package pocketbase import ( "database/sql" "errors" "fmt" "strings" "github.com/google/uuid" "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" "git.vakhrushev.me/av/transcriber/internal/adapter/repo/pocketbase/migrations" "git.vakhrushev.me/av/transcriber/internal/clock" ) type AudioRecordRepository struct { app core.App } func NewAudioRecordRepository(app core.App) *AudioRecordRepository { return &AudioRecordRepository{app: app} } func (repo *AudioRecordRepository) Create(r *entity.AudioRecord) error { collection, err := findCollection(repo.app, migrations.RecordsCollection) if err != nil { return err } record := core.NewRecord(collection) if r.Id != "" { record.Id = r.Id } applyToRecord(record, r) if err := repo.app.Save(record); err != nil { return fmt.Errorf("failed to insert audio record: %w", err) } r.Id = record.Id r.CreatedAt = record.GetDateTime("created").Time() r.UpdatedAt = record.GetDateTime("updated").Time() return nil } // Save сохраняет запись, захват которой держит holder. Проверка и запись идут // одной транзакцией: шаг, потерявший запись за время работы, получает // LostAcquisitionError и результата не пишет. // // Сверяется **значение** признака захвата, а не занятость записи. Захват, // перевыданный другому — по протуханию срока или после того, как человек снял // признак остановки в панели, — обязан обратить запись первого в отказ; условие // по непустоте признака пропустило бы обоих, и два шага записали бы в одну // запись по очереди, портя её результат. func (repo *AudioRecordRepository) Save(r *entity.AudioRecord, holder string) error { return repo.app.RunInTransaction(func(txApp core.App) error { record, err := txApp.FindRecordById(migrations.RecordsCollection, r.Id) if err != nil { return fmt.Errorf("failed to find audio record: %w", err) } if holder != "" && record.GetString("acquisition_id") != holder { return &contract.LostAcquisitionError{JobID: r.Id} } // Кладём только то, чем распоряжается конвейер: правку владельца в // панели снимок шага стирать не должен. applyOwnedByPipeline(record, r) if err := txApp.Save(record); err != nil { return fmt.Errorf("failed to update audio record: %w", err) } r.UpdatedAt = record.GetDateTime("updated").Time() return nil }) } // GetByID отдаёт запись, только если её владелец — ownerID. // // Чужая запись, запись без владельца и несуществующая дают одну и ту же ошибку: // по разнице ответов иначе перебирается список заведённых записей, а // идентификатор записи и есть то, что разграничение прячет. // // Пустой ownerID отсекается **до** чтения и не совпадает ни с чем. Правило это // не стало избыточным с обязательностью колонки: схема запрещает **заводить** // ничью запись, а здесь запрещено **спрашивать** ничьим именем — иначе // вызывающий без учётной записи получил бы выборку вместо отказа. func (repo *AudioRecordRepository) GetByID(id, ownerID string) (*entity.AudioRecord, error) { if ownerID == "" { return nil, &contract.JobNotFoundError{Message: "record not found"} } record, err := repo.find(id) if err != nil { return nil, err } if record.GetString("owner") != ownerID { return nil, &contract.JobNotFoundError{Message: "record not found"} } return recordToAudioRecord(record), nil } // Get отдаёт запись без сужения владельцем: им пользуется конвейер, чья выборка // владельцем не сужается. func (repo *AudioRecordRepository) Get(id string) (*entity.AudioRecord, error) { record, err := repo.find(id) if err != nil { return nil, err } return recordToAudioRecord(record), nil } func (repo *AudioRecordRepository) find(id string) (*core.Record, error) { record, err := repo.app.FindRecordById(migrations.RecordsCollection, id) if err != nil { // «Такой записи нет» переводится в доменную ошибку **здесь**, у // источника, как велит конвенция об ошибках. Иначе три исхода, которые // разграничение обязано сделать неразличимыми, разъезжаются: чужая и // ничья записи дают доменную ошибку, а несуществующая — отказ базы, // неотличимый от настоящей аварии хранилища. if errors.Is(err, sql.ErrNoRows) { return nil, &contract.JobNotFoundError{Message: "record not found"} } return nil, fmt.Errorf("failed to get audio record: %w", err) } return record, nil } // FindAndAcquire забирает пригодную к работе запись одним неделимым шагом: // выбор подходящей и пометка её захваченной идут вместе. // // Возвращается **идентификатор и признак этого захвата**, а не перечень колонок. // Колонки шаг читает обычным чтением: иначе всякая новая колонка записи попадала // бы под инвариант проекта о колонках очереди, а забытая приезжала бы нулевой, и // первое же сохранение писало бы этот ноль поверх сохранённого значения. // // Срок протухания захвата приезжает **с рубежом**, а не с воркером: воркер не // привязан к шагу и не знает заранее, что вытянет. Перечень рубежей и их сроков // приходит одним дескриптором — перечислять их порознь нельзя: рубеж, забытый в // отборе, не выдаётся ни одному воркеру никогда, а пустой прогон по инварианту // проекта не пишется в журнал и не считается в метрику. // // Запрос идёт сырым, мимо записей коллекции: `app.DB()` направляет всё, кроме // выборок, в пул с единственным соединением, и захваты выстраиваются в очередь. // Хуки коллекции на нём не срабатывают, поэтому время изменения проставляет сам // запрос. // // Все времена кладутся и сравниваются тем же видом, каким хранилище пишет свои // `created`/`updated`: сравнение строк побайтово, и вид, разошедшийся хоть // разделителем, обратил бы условие срока в постоянную истину или постоянную // ложь — молча. func (repo *AudioRecordRepository) FindAndAcquire(stages []entity.Stage) (*contract.AcquiredRecord, error) { if len(stages) == 0 { return nil, &contract.JobNotFoundError{Message: "no working stages declared"} } // Метка времени берётся единой точкой, а не `types.NowDateTime()`: обёртка // хранилища читает часы сама, и запрет линтера её не видит — новая метка в // этом запросе обошла бы единую точку молча. now, err := types.ParseDateTime(clock.Now()) if err != nil { return nil, fmt.Errorf("failed to parse current time: %w", err) } holder := uuid.NewString() params := dbx.Params{ "holder": holder, "now": now.String(), } // Срок протухания у каждого рубежа свой, поэтому он выбирается по рубежу // самой записи прямо в запросе: воркер, ещё не знающий, что вытянет, // подставить его не может. var expiry strings.Builder expiry.WriteString("CASE state") var states []string for i, stage := range stages { stateKey := fmt.Sprintf("state%d", i) expiryKey := fmt.Sprintf("expiry%d", i) deadline, err := types.ParseDateTime(clock.Now().Add(stage.AcquireTimeout)) if err != nil { return nil, fmt.Errorf("failed to parse acquire deadline: %w", err) } fmt.Fprintf(&expiry, " WHEN {:%s} THEN {:%s}", stateKey, expiryKey) params[stateKey] = stage.Name params[expiryKey] = deadline.String() states = append(states, "{:"+stateKey+"}") } expiry.WriteString(" END") table := "{{" + migrations.RecordsCollection + "}}" query := repo.app.DB().NewQuery(` UPDATE ` + table + ` SET acquisition_id = {:holder}, acquire_expires_at = ` + expiry.String() + `, attempts = attempts + 1, updated = {:now} WHERE id = ( SELECT id FROM ` + table + ` WHERE state IN (` + strings.Join(states, ", ") + `) AND (halted_at = '' OR halted_at IS NULL) AND (delay_time = '' OR delay_time IS NULL OR delay_time < {:now}) AND (acquisition_id = '' OR acquisition_id IS NULL OR acquire_expires_at = '' OR acquire_expires_at IS NULL OR acquire_expires_at < {:now}) ORDER BY created, id LIMIT 1 ) RETURNING id`) query.Bind(params) var row struct { Id string `db:"id"` } if err := query.One(&row); err != nil { if errors.Is(err, sql.ErrNoRows) { return nil, &contract.JobNotFoundError{Message: "no record is ready for work"} } return nil, fmt.Errorf("failed to acquire an audio record: %w", err) } return &contract.AcquiredRecord{ID: row.Id, Holder: holder}, nil }