package pocketbase import ( "database/sql" "errors" "fmt" "time" "github.com/pocketbase/dbx" "github.com/pocketbase/pocketbase/core" "github.com/pocketbase/pocketbase/tools/types" "git.vakhrushev.me/av/transcriber/internal/contract" "git.vakhrushev.me/av/transcriber/internal/entity" ) type TranscriptJobRepository struct { app core.App } func NewTranscriptJobRepository(app core.App) *TranscriptJobRepository { return &TranscriptJobRepository{app: app} } func (repo *TranscriptJobRepository) Create(job *entity.TranscribeJob) error { collection, err := findCollection(repo.app, JobsCollection) if err != nil { return err } record := core.NewRecord(collection) if job.Id != "" { record.Id = job.Id } applyToRecord(record, job) if err := repo.app.Save(record); err != nil { return fmt.Errorf("failed to insert transcribe job: %w", err) } job.Id = record.Id job.CreatedAt = record.GetDateTime("created").Time() job.UpdatedAt = record.GetDateTime("updated").Time() return nil } // Save сохраняет задачу, захват которой держит holder. Проверка и запись идут // одной транзакцией: шаг, потерявший задачу за время работы, получает // LostAcquisitionError и результата не пишет. func (repo *TranscriptJobRepository) Save(job *entity.TranscribeJob, holder string) error { err := repo.app.RunInTransaction(func(txApp core.App) error { record, err := txApp.FindRecordById(JobsCollection, job.Id) if err != nil { return fmt.Errorf("failed to find transcribe job: %w", err) } if holder != "" && record.GetString("acquisition_id") != holder { return &contract.LostAcquisitionError{JobID: job.Id} } // Кладём только то, чем распоряжается конвейер: правку владельца в // панели снимок шага стирать не должен. applyOwnedByPipeline(record, job) if err := txApp.Save(record); err != nil { return fmt.Errorf("failed to update transcribe job: %w", err) } job.UpdatedAt = record.GetDateTime("updated").Time() return nil }) if err != nil { return err } return nil } func (repo *TranscriptJobRepository) GetByID(id string) (*entity.TranscribeJob, error) { record, err := repo.app.FindRecordById(JobsCollection, id) if err != nil { return nil, fmt.Errorf("failed to get transcribe job: %w", err) } return recordToJob(record), nil } // Колонки, которые читает захват. Список нужен запросу дословно: `RETURNING *` // отдал бы и порядок, зависящий от схемы. const acquireColumns = `id, state, source, file, error_text, acquisition_id, ` + `acquire_time, delay_time, attempts, recognition_op_id, transcription_text, ` + `tg_chat_id, tg_reply_message_id, created, updated` // FindAndAcquire забирает задачу одним неделимым шагом: выбор подходящей и // пометка её захваченной идут вместе, и захваченная возвращается тем же // запросом. Двум вызывающим, пришедшим за одним состоянием, запись достаётся // одному — на этом стоит инвариант «Принятая запись не теряется молча». // // Запрос идёт сырым, мимо записей коллекции: `app.DB()` направляет всё, кроме // выборок, в пул с единственным соединением, и захваты выстраиваются в очередь. // Хуки коллекции на нём не срабатывают, поэтому время изменения проставляет сам // запрос. // // Все времена кладутся и сравниваются тем же видом, каким хранилище пишет свои // `created`/`updated`: сравнение строк побайтово, и вид, разошедшийся хоть // разделителем, обратил бы условие срока в постоянную истину или постоянную // ложь — молча. func (repo *TranscriptJobRepository) FindAndAcquire(state, acquisitionId string, rottingTime time.Time) (*entity.TranscribeJob, error) { now := types.NowDateTime() query := repo.app.DB().NewQuery(` UPDATE {{` + JobsCollection + `}} SET acquisition_id = {:acquisition_id}, acquire_time = {:now}, attempts = attempts + 1, updated = {:now} WHERE id = ( SELECT id FROM {{` + JobsCollection + `}} WHERE state = {:state} AND (delay_time = '' OR delay_time IS NULL OR delay_time < {:now}) AND (acquisition_id = '' OR acquisition_id IS NULL OR acquire_time < {:rotting}) ORDER BY created, id LIMIT 1 ) RETURNING ` + acquireColumns) rotting, err := types.ParseDateTime(rottingTime) if err != nil { return nil, fmt.Errorf("failed to parse rotting time: %w", err) } query.Bind(dbx.Params{ "acquisition_id": acquisitionId, "now": now.String(), "state": state, "rotting": rotting.String(), }) var row acquiredRow if err := query.One(&row); err != nil { if errors.Is(err, sql.ErrNoRows) { return nil, &contract.JobNotFoundError{State: state, Message: "appropriate job not found"} } return nil, fmt.Errorf("failed to aquire job with state %s: %w", state, err) } return row.toJob(), nil }