Files
transcriber/internal/adapter/repo/pocketbase/transcript_job_repo.go
T
av 01cc31d45f хранилище, файлы записей и очередь переведены на встроенную PocketBase
- записи, метаданные и файлы съехались под один каталог данных; появилась
  панель владельца, а gin, goqu, goose и требование CGO ушли
- захват задачи стал одним запросом с RETURNING; заведены число попыток,
  состояние dead и нарастающая пауза вместо признака is_error
- имя файла в хранилище задаёт сервис и в журнал не идёт: вместе с
  идентификатором записи оно собирало бы ссылку на скачивание
2026-08-12 08:31:59 +03:00

148 lines
5.5 KiB
Go

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
}