package sqlite import ( "context" "database/sql" "errors" "fmt" "strings" "git.vakhrushev.me/av/transcriber/internal/clock" "git.vakhrushev.me/av/transcriber/internal/contract" "git.vakhrushev.me/av/transcriber/internal/entity" "git.vakhrushev.me/av/transcriber/internal/ident" ) // recordsTable — таблица аудиозаписей. const recordsTable = "audio_records" type AudioRecordRepository struct { db *DB } func NewAudioRecordRepository(db *DB) *AudioRecordRepository { return &AudioRecordRepository{db: db} } // Create заводит аудиозапись. // // Идентификатор приходит готовым, когда его назначил вызывающий: приём знает его // раньше, чем кладёт файл, — копии записи лежат подкаталогом под этим самым // идентификатором. Пустой заполняется единой точкой выдачи. func (repo *AudioRecordRepository) Create(r *entity.AudioRecord) error { if r.Id == "" { r.Id = ident.New() } now := clock.Now() if r.CreatedAt.IsZero() { r.CreatedAt = now } r.UpdatedAt = now if r.StateEnteredAt.IsZero() { r.StateEnteredAt = now } query, args := insertSQL(recordsTable, writeRecord(r)) if _, err := repo.db.Writer().ExecContext(context.Background(), query, args...); err != nil { return fmt.Errorf("failed to insert audio record: %w", err) } return nil } // Save сохраняет запись, захват которой держит holder. // // Сверка захвата и запись идут **одним запросом**: значение признака стоит // условием правки, поэтому между проверкой и записью не остаётся окна. Сверяется // именно значение, а не занятость записи — захват, перевыданный другому по // протуханию срока или после того, как человек вернул запись в работу, обязан // обратить запись первого в отказ; условие по непустоте признака пропустило бы // обоих, и два шага записали бы в одну запись по очереди, портя её результат. // // Пустой holder снимает условность и в конвейере не употребляется: все его шаги // получают признак захвата от FindAndAcquire. func (repo *AudioRecordRepository) Save(r *entity.AudioRecord, holder string) error { r.UpdatedAt = clock.Now() where := "id = :id" whereArgs := []any{sql.Named("id", r.Id)} if holder != "" { where += " AND acquisition_id = :holder" whereArgs = append(whereArgs, sql.Named("holder", holder)) } query, args := updateSQL(recordsTable, writeOwnedByPipeline(r), where, whereArgs) result, err := repo.db.Writer().ExecContext(context.Background(), query, args...) if err != nil { return fmt.Errorf("failed to update audio record: %w", err) } affected, err := result.RowsAffected() if err != nil { return fmt.Errorf("failed to read the outcome of an audio record update: %w", err) } if affected > 0 { return nil } // Строк не тронуто по одной из двух причин, и различить их можно только // чтением: записи нет вовсе либо захват достался другому. Разница несущая — // первая означает поломку, вторая штатный исход шага, потерявшего запись. if _, err := repo.Get(r.Id); err != nil { return err } return &contract.LostAcquisitionError{JobID: r.Id} } // 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.read("id = ? AND owner_id = ?", id, ownerID) if err != nil { return nil, err } return record, nil } // Get отдаёт запись без сужения владельцем: им пользуется конвейер, чья выборка // владельцем не сужается. func (repo *AudioRecordRepository) Get(id string) (*entity.AudioRecord, error) { return repo.read("id = ?", id) } // read читает одну запись по условию. func (repo *AudioRecordRepository) read(where string, args ...any) (*entity.AudioRecord, error) { row := &recordRow{} columns, targets := selectList(readRecordColumns(row), "") query := "SELECT " + columns + " FROM " + recordsTable + " WHERE " + where + " LIMIT 1" if err := repo.db.Reader().QueryRowContext(context.Background(), query, args...).Scan(targets...); 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) } record := rowToAudioRecord(row) topics, err := repo.topicsOf([]string{record.Id}) if err != nil { return nil, err } record.TopicIDs = topics[record.Id] return record, nil } // FindAndAcquire забирает пригодную к работе запись одним неделимым шагом: // выбор подходящей и пометка её захваченной идут вместе, одним оператором с // возвратом. // // Возвращается **идентификатор и признак этого захвата**, а не перечень колонок. // Колонки шаг читает обычным чтением: иначе всякая новая колонка записи попадала // бы под инвариант проекта о колонках очереди, а забытая приезжала бы нулевой, и // первое же сохранение писало бы этот ноль поверх сохранённого значения. // // Срок протухания захвата приезжает **с рубежом**, а не с воркером: воркер не // привязан к шагу и не знает заранее, что вытянет. Перечень рубежей и их сроков // приходит одним дескриптором — перечислять их порознь нельзя: рубеж, забытый в // отборе, не выдаётся ни одному воркеру никогда, а пустой прогон по инварианту // проекта не пишется в журнал и не считается в метрику. // // Запрос идёт по **пишущему** соединению: он читает состояние, которое сам же // меняет, а транзакцию, начатую на читающем соединении, SQLite до пишущей не // повышает. func (repo *AudioRecordRepository) FindAndAcquire(stages []entity.Stage) (*contract.AcquiredRecord, error) { if len(stages) == 0 { return nil, &contract.JobNotFoundError{Message: "no working stages declared"} } now := clock.Now() holder := ident.New() args := []any{ sql.Named("holder", holder), sql.Named("now", formatTime(now)), } // Срок протухания у каждого рубежа свой, поэтому он выбирается по рубежу // самой записи прямо в запросе: воркер, ещё не знающий, что вытянет, // подставить его не может. var expiry strings.Builder expiry.WriteString("CASE state") states := make([]string, 0, len(stages)) for i, stage := range stages { stateKey := fmt.Sprintf("state%d", i) expiryKey := fmt.Sprintf("expiry%d", i) fmt.Fprintf(&expiry, " WHEN :%s THEN :%s", stateKey, expiryKey) args = append(args, sql.Named(stateKey, stage.Name), sql.Named(expiryKey, formatTime(now.Add(stage.AcquireTimeout))), ) states = append(states, ":"+stateKey) } expiry.WriteString(" END") // Порядок выборки определён однозначно: время заведения плюс ключ записи. // Сравнения по неуникальному значению для этого мало — порядок обработки // стал бы невоспроизводимым. query := ` UPDATE ` + recordsTable + ` SET acquisition_id = :holder, acquire_expires_at = ` + expiry.String() + `, attempts = attempts + 1, updated_at = :now WHERE id = ( SELECT id FROM ` + recordsTable + ` WHERE state IN (` + strings.Join(states, ", ") + `) AND halted_at IS NULL AND (delay_time IS NULL OR delay_time < :now) AND (acquisition_id IS NULL OR acquire_expires_at IS NULL OR acquire_expires_at < :now) ORDER BY created_at, id LIMIT 1 ) RETURNING id` var id string if err := repo.db.Writer().QueryRowContext(context.Background(), query, args...).Scan(&id); 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: id, Holder: holder}, nil }