- audiorecords вместо transcribe_jobs: приложения (texts, structures, recognitions, record_events, topics) живут своими коллекциями, ссылки на исходник и на приведённую копию перестали переставляться - рубеж называет достигнутое, отказ стал признаком остановки с причиной, а сторожей стало двое: число отказов и время в рубеже - воркеры потеряли специализацию, их число задаётся [pipeline] workers, шаг выбирается по рубежу, а захват отдаёт идентификатор и признак захвата
162 lines
6.3 KiB
Go
162 lines
6.3 KiB
Go
package pocketbase
|
|
|
|
import (
|
|
"errors"
|
|
"fmt"
|
|
"io"
|
|
|
|
"github.com/pocketbase/pocketbase/core"
|
|
"github.com/pocketbase/pocketbase/tools/filesystem"
|
|
|
|
"git.vakhrushev.me/av/transcriber/internal/adapter/repo/pocketbase/migrations"
|
|
"git.vakhrushev.me/av/transcriber/internal/clock"
|
|
"git.vakhrushev.me/av/transcriber/internal/entity"
|
|
)
|
|
|
|
type RecognitionRepository struct {
|
|
app core.App
|
|
}
|
|
|
|
func NewRecognitionRepository(app core.App) *RecognitionRepository {
|
|
return &RecognitionRepository{app: app}
|
|
}
|
|
|
|
// Create заводит строку попытки **до** обращения к провайдеру.
|
|
//
|
|
// Порядок здесь несущий: окно между ответом провайдера и записью идентификатора
|
|
// операции — то место, где теряется оплаченное. Заведённая заранее строка даёт
|
|
// повторному шагу, чем проверить сделанное прежде, чем платить второй раз.
|
|
func (repo *RecognitionRepository) Create(r *entity.Recognition) error {
|
|
collection, err := findCollection(repo.app, migrations.RecognitionsCollection)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
started := clock.Now()
|
|
record := core.NewRecord(collection)
|
|
record.Set("record", r.RecordID)
|
|
record.Set("provider", r.Provider)
|
|
record.Set("model", r.Model)
|
|
record.Set("external_id", r.ExternalID)
|
|
record.Set("source_uri", r.SourceURI)
|
|
record.Set("started_at", dateOrEmpty(&started))
|
|
|
|
if err := repo.app.Save(record); err != nil {
|
|
return fmt.Errorf("failed to create recognition attempt for record %s: %w", r.RecordID, err)
|
|
}
|
|
|
|
r.Id = record.Id
|
|
r.StartedAt = &started
|
|
return nil
|
|
}
|
|
|
|
// Submitted сохраняет адрес аудио и идентификатор заведённой операции. По
|
|
// последнему повторный шаг узнаёт, что за эту запись уже заплачено, и второй раз
|
|
// наружу не платит.
|
|
func (repo *RecognitionRepository) Submitted(id, sourceURI, externalID string) error {
|
|
record, err := repo.app.FindRecordById(migrations.RecognitionsCollection, id)
|
|
if err != nil {
|
|
return fmt.Errorf("failed to find recognition attempt %s: %w", id, err)
|
|
}
|
|
record.Set("source_uri", sourceURI)
|
|
record.Set("external_id", externalID)
|
|
if err := repo.app.Save(record); err != nil {
|
|
return fmt.Errorf("failed to store operation id of attempt %s: %w", id, err)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// Finish кладёт сырой ответ провайдера вложением и отмечает завершение.
|
|
//
|
|
// Вложением, а не колонкой: шаг опроса читает эту строку раз в несколько секунд,
|
|
// а хранилище читает запись целиком — ответ на многочасовую запись ехал бы в
|
|
// память при каждом опросе. Хранится он потому, что результат операции у
|
|
// провайдера не переспрашивается.
|
|
func (repo *RecognitionRepository) Finish(id string, raw []byte) error {
|
|
record, err := repo.app.FindRecordById(migrations.RecognitionsCollection, id)
|
|
if err != nil {
|
|
return fmt.Errorf("failed to find recognition attempt %s: %w", id, err)
|
|
}
|
|
|
|
if len(raw) > 0 {
|
|
// Имя вложения задаём мы: умолчание хранилища строит его из имени
|
|
// исходного файла, а имя, данное отправителем, в хранилище не попадает.
|
|
payload, err := filesystem.NewFileFromBytes(raw, id+".payload")
|
|
if err != nil {
|
|
return fmt.Errorf("failed to prepare provider payload of attempt %s", id)
|
|
}
|
|
record.Set("payload", payload)
|
|
}
|
|
|
|
finished := clock.Now()
|
|
record.Set("finished_at", dateOrEmpty(&finished))
|
|
|
|
if err := repo.app.Save(record); err != nil {
|
|
// Отказ хранилища несёт имя файла вложения целиком, а оно — последняя
|
|
// часть ссылки: цепочка `%w` уехала бы в журнал вместе с ним.
|
|
return fmt.Errorf("failed to store provider payload of attempt %s", id)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (repo *RecognitionRepository) GetByID(id string) (*entity.Recognition, error) {
|
|
record, err := repo.app.FindRecordById(migrations.RecognitionsCollection, id)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("failed to get recognition attempt %s: %w", id, err)
|
|
}
|
|
|
|
return &entity.Recognition{
|
|
Id: record.Id,
|
|
RecordID: record.GetString("record"),
|
|
Provider: record.GetString("provider"),
|
|
Model: record.GetString("model"),
|
|
ExternalID: record.GetString("external_id"),
|
|
SourceURI: record.GetString("source_uri"),
|
|
StartedAt: timeOrNil(record.GetDateTime("started_at")),
|
|
FinishedAt: timeOrNil(record.GetDateTime("finished_at")),
|
|
}, nil
|
|
}
|
|
|
|
// ReadRaw отдаёт сохранённый ответ провайдера. Зовётся только тогда, когда ответ
|
|
// нужен: шаг опроса читает строку попытки без него.
|
|
func (repo *RecognitionRepository) ReadRaw(id string) ([]byte, error) {
|
|
record, err := repo.app.FindRecordById(migrations.RecognitionsCollection, id)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("failed to find recognition attempt %s: %w", id, err)
|
|
}
|
|
|
|
names := record.GetStringSlice("payload")
|
|
if len(names) == 0 {
|
|
return nil, fmt.Errorf("recognition attempt %s has no stored payload", id)
|
|
}
|
|
|
|
fsys, err := repo.app.NewFilesystem()
|
|
if err != nil {
|
|
return nil, fmt.Errorf("failed to open storage filesystem: %w", err)
|
|
}
|
|
|
|
reader, err := fsys.GetReader(record.BaseFilesPath() + "/" + names[0])
|
|
if err != nil {
|
|
// Отказ хранилища несёт имя вложения целиком, а имя — последняя часть
|
|
// ссылки на скачивание: наружу идёт идентификатор попытки, и только он.
|
|
return nil, errors.Join(
|
|
fmt.Errorf("failed to read stored payload of attempt %s", id),
|
|
fsys.Close(),
|
|
)
|
|
}
|
|
|
|
raw, readErr := io.ReadAll(reader)
|
|
closeErr := errors.Join(reader.Close(), fsys.Close())
|
|
if readErr != nil {
|
|
return nil, errors.Join(
|
|
fmt.Errorf("failed to read stored payload of attempt %s", id),
|
|
closeErr,
|
|
)
|
|
}
|
|
if closeErr != nil {
|
|
return nil, closeErr
|
|
}
|
|
|
|
return raw, nil
|
|
}
|