хранилище, файлы записей и очередь переведены на встроенную PocketBase

- записи, метаданные и файлы съехались под один каталог данных; появилась
  панель владельца, а gin, goqu, goose и требование CGO ушли
- захват задачи стал одним запросом с RETURNING; заведены число попыток,
  состояние dead и нарастающая пауза вместо признака is_error
- имя файла в хранилище задаёт сервис и в журнал не идёт: вместе с
  идентификатором записи оно собирало бы ссылку на скачивание
This commit is contained in:
av
2026-08-12 08:31:59 +03:00
parent 09cedc4e61
commit 01cc31d45f
55 changed files with 5238 additions and 1235 deletions
+12 -1
View File
@@ -2,6 +2,7 @@ package yandex
import (
"context"
"errors"
"fmt"
"io"
"strings"
@@ -11,6 +12,7 @@ import (
"github.com/aws/aws-sdk-go-v2/credentials"
"github.com/aws/aws-sdk-go-v2/feature/s3/manager"
"github.com/aws/aws-sdk-go-v2/service/s3"
"github.com/aws/smithy-go"
)
type s3Config struct {
@@ -72,7 +74,16 @@ func (s *yandexS3Service) uploadFile(file io.Reader, fileName string) error {
Body: file,
})
if err != nil {
return fmt.Errorf("failed to upload file to S3: %w", err)
// Отказ SDK несёт полный URL объекта, то есть имя файла в хранилище, а
// оно — последняя часть ссылки на скачивание: цепочка `%w` уехала бы в
// журнал вместе с ключом. Наружу идёт класс отказа и только он — по
// нему «ключи отозваны» отличимо от «бакета нет» и от «сети нет», а
// адреса в коде отказа SDK не бывает.
var apiErr smithy.APIError
if errors.As(err, &apiErr) {
return fmt.Errorf("failed to upload file to S3: %s", apiErr.ErrorCode())
}
return errors.New("failed to upload file to S3")
}
return nil
+56
View File
@@ -0,0 +1,56 @@
// Package pocketbase — хранилище задач и файлов поверх встроенной PocketBase.
//
// Приложение поднимается библиотекой, а не её набором команд: разбор флагов и
// мягкая остановка остаются нашими, а ключ `-c config.toml` — объявленный
// контракт запуска.
package pocketbase
import (
"fmt"
pb "github.com/pocketbase/pocketbase"
"github.com/pocketbase/pocketbase/core"
)
// Имена коллекций. Они же — часть пути к файлу в раскладке хранилища и часть
// адреса ссылки на него, поэтому меняются только новым шагом схемы.
const (
FilesCollection = "files"
JobsCollection = "transcribe_jobs"
)
// New создаёт приложение хранилища на заданном каталоге данных и приводит его в
// рабочее состояние: открывает базу, читает настройки и накатывает непринятые
// шаги схемы.
//
// Схема накатывается **здесь**, а не оставляется серверу, хотя тот и гоняет
// непринятые шаги сам. Причина в порядке: воркеры стартуют раньше сервера, и на
// чистом каталоге их первые опросы приходились бы на несуществующую таблицу —
// отказ в журнале и в счётчике на каждую секунду до конца накатки.
func New(dataDir string) (*pb.PocketBase, error) {
app := pb.NewWithConfig(pb.Config{
DefaultDataDir: dataDir,
HideStartBanner: true,
})
if err := app.Bootstrap(); err != nil {
return nil, fmt.Errorf("failed to bootstrap storage: %w", err)
}
if err := app.RunAllMigrations(); err != nil {
return nil, fmt.Errorf("failed to apply storage schema: %w", err)
}
return app, nil
}
// MustFindCollection достаёт коллекцию по имени. Отсутствие коллекции здесь —
// не отказ окружения, а несделанный шаг схемы: сервис до этой точки не доходит,
// потому что Serve накатывает схему прежде, чем поднять сервер.
func findCollection(app core.App, name string) (*core.Collection, error) {
collection, err := app.FindCollectionByNameOrId(name)
if err != nil {
return nil, fmt.Errorf("failed to find collection %s: %w", name, err)
}
return collection, nil
}
@@ -0,0 +1,269 @@
package pocketbase
import (
"errors"
"fmt"
"io"
"os"
"path/filepath"
"github.com/pocketbase/pocketbase/core"
"github.com/pocketbase/pocketbase/tools/filesystem"
"git.vakhrushev.me/av/transcriber/internal/contract"
"git.vakhrushev.me/av/transcriber/internal/entity"
)
// workFile — рабочая копия файла на диске. Живёт во временном каталоге
// системы, а не в каталоге данных: последний смонтирован на сервере, и
// временному там не место.
type workFile struct {
path string
}
func (w *workFile) Path() string { return w.path }
func (w *workFile) Size() (int64, error) {
info, err := os.Stat(w.path)
if err != nil {
return 0, fmt.Errorf("failed to stat work file: %w", err)
}
return info.Size(), nil
}
// Close убирает копию. Отсутствие файла отказом не считается: шаг мог не дойти
// до его создания, и повторный Close тоже законен.
func (w *workFile) Close() error {
if err := os.Remove(w.path); err != nil && !os.IsNotExist(err) {
return fmt.Errorf("failed to remove work file: %w", err)
}
return nil
}
type FileRepository struct {
app core.App
}
func NewFileRepository(app core.App) *FileRepository {
return &FileRepository{app: app}
}
// newWorkFile заводит пустую копию во временном каталоге. Расширение сохраняется
// в имени: `ffprobe` и `ffmpeg` по нему выбирают разбор.
func newWorkFile(ext string) (*workFile, error) {
f, err := os.CreateTemp("", "transcriber-*"+ext)
if err != nil {
return nil, fmt.Errorf("failed to create work file: %w", err)
}
path := f.Name()
if err := f.Close(); err != nil {
_ = os.Remove(path)
return nil, fmt.Errorf("failed to close work file: %w", err)
}
return &workFile{path: path}, nil
}
func (repo *FileRepository) StageEmpty(ext string) (contract.WorkFile, error) {
return newWorkFile(ext)
}
func (repo *FileRepository) Stage(ext string, content io.Reader) (contract.WorkFile, error) {
work, err := newWorkFile(ext)
if err != nil {
return nil, err
}
if err := writeTo(work.path, content); err != nil {
// Отказ уборки не подменяет отказ записи, но и не теряется.
return nil, errors.Join(err, work.Close())
}
return work, nil
}
func (repo *FileRepository) Localize(fileID string) (contract.WorkFile, error) {
record, err := repo.app.FindRecordById(FilesCollection, fileID)
if err != nil {
return nil, fmt.Errorf("failed to find file %s: %w", fileID, err)
}
name := firstFileName(record)
if name == "" {
return nil, fmt.Errorf("file %s has no content in storage", fileID)
}
work, err := newWorkFile(filepath.Ext(name))
if err != nil {
return nil, err
}
src, err := repo.openStored(record, name)
if err != nil {
return nil, errors.Join(err, work.Close())
}
defer src.Close()
if err := writeTo(work.path, src); err != nil {
return nil, errors.Join(err, work.Close())
}
return work, nil
}
// CreateLocal кладёт рабочую копию в хранилище. Имя задаём мы: умолчание
// библиотеки строит его из имени, данного отправителем, а имя отправителя в
// хранилище не попадает — путь к файлу читается в журнале, и инвариант
// приватности этого не допускает. Свой суффикс хранилище допишет само.
func (repo *FileRepository) CreateLocal(name string, work contract.WorkFile) (*entity.File, error) {
collection, err := findCollection(repo.app, FilesCollection)
if err != nil {
return nil, err
}
stored, err := filesystem.NewFileFromPath(work.Path())
if err != nil {
return nil, fmt.Errorf("failed to read work file: %w", err)
}
stored.Name = name
record := core.NewRecord(collection)
record.Set("file", stored)
record.Set("location", entity.LocationLocal)
record.Set("size", stored.Size)
if err := repo.app.Save(record); err != nil {
// Отказ укладки называет имя файла — то самое, из которого строится
// ссылка на скачивание. В цепочку оно не идёт по той же причине, что и
// ключ при чтении.
return nil, errors.New("failed to store file")
}
return recordToFile(record), nil
}
func (repo *FileRepository) CreateRemote(objectKey string, size int64) (*entity.File, error) {
collection, err := findCollection(repo.app, FilesCollection)
if err != nil {
return nil, err
}
record := core.NewRecord(collection)
record.Set("location", entity.LocationS3)
record.Set("object_key", objectKey)
record.Set("size", size)
if err := repo.app.Save(record); err != nil {
return nil, fmt.Errorf("failed to store remote file record: %w", err)
}
return recordToFile(record), nil
}
func (repo *FileRepository) GetByID(id string) (*entity.File, error) {
record, err := repo.app.FindRecordById(FilesCollection, id)
if err != nil {
return nil, fmt.Errorf("failed to get file: %w", err)
}
return recordToFile(record), nil
}
func (repo *FileRepository) Open(fileID string) (io.ReadCloser, error) {
record, err := repo.app.FindRecordById(FilesCollection, fileID)
if err != nil {
return nil, fmt.Errorf("failed to find file %s: %w", fileID, err)
}
name := firstFileName(record)
if name == "" {
return nil, fmt.Errorf("file %s has no content in storage", fileID)
}
return repo.openStored(record, name)
}
// openStored открывает содержимое файла в хранилище потоком.
func (repo *FileRepository) openStored(record *core.Record, name string) (io.ReadCloser, error) {
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() + "/" + name)
if err != nil {
// Отказ хранилища несёт ключ файла целиком, а ключ — последняя часть
// ссылки `/api/files/...`, по которой запись скачивают. Наружу отдаётся
// идентификатор записи, и только он: цепочка `%w` уехала бы в журнал и
// стала бы там бессрочным ключом к чужому аудио.
return nil, errors.Join(
fmt.Errorf("failed to read stored file of record %s", record.Id),
fsys.Close(),
)
}
return &storedReader{reader: reader, fsys: fsys}, nil
}
// storedReader держит открытой файловую систему хранилища на всё время чтения:
// закрытая раньше времени, она обрывает поток на середине записи.
type storedReader struct {
reader io.ReadCloser
fsys io.Closer
}
func (r *storedReader) Read(p []byte) (int, error) { return r.reader.Read(p) }
func (r *storedReader) Close() error {
readerErr := r.reader.Close()
fsysErr := r.fsys.Close()
switch {
case readerErr != nil && fsysErr != nil:
return errors.New("failed to close stored file and its filesystem")
case readerErr != nil:
return errors.New("failed to close stored file")
default:
return fsysErr
}
}
// writeTo переливает содержимое в файл потоком. В память запись целиком не
// читается: расчётный потолок — шесть часов.
func writeTo(path string, content io.Reader) error {
dst, err := os.Create(path)
if err != nil {
return fmt.Errorf("failed to open work file: %w", err)
}
if _, err := io.Copy(dst, content); err != nil {
_ = dst.Close()
return fmt.Errorf("failed to write work file: %w", err)
}
if err := dst.Close(); err != nil {
return fmt.Errorf("failed to close work file: %w", err)
}
return nil
}
func firstFileName(record *core.Record) string {
names := record.GetStringSlice("file")
if len(names) == 0 {
return ""
}
return names[0]
}
func recordToFile(record *core.Record) *entity.File {
name := firstFileName(record)
if name == "" {
name = record.GetString("object_key")
}
return &entity.File{
Id: record.Id,
Location: record.GetString("location"),
FileName: name,
Size: int64(record.GetInt("size")),
CreatedAt: record.GetDateTime("created").Time(),
}
}
@@ -0,0 +1,38 @@
package pocketbase
import (
"strings"
"testing"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"git.vakhrushev.me/av/transcriber/internal/entity"
)
// Потолок размера у поля файла задан числом, а не нулём: нулём библиотека читает
// собственное умолчание в 5 МиБ, и на нём отвергалось бы всё длиннее примерно
// пяти минут — то есть штатная запись сервиса. Проверка судит запись, которая
// заведомо больше этого умолчания: обновление библиотеки, вернувшее умолчание,
// иначе прошло бы молча.
func TestCreateLocal_AcceptsRecordLargerThanLibraryDefault(t *testing.T) {
app := newTestApp(t)
repo := NewFileRepository(app)
const libraryDefault = 5 << 20
// Ровно на байт больше умолчания: проверка судит границу, а не пропускную
// способность — лишние мегабайты стоили бы секунд на каждом прогоне.
work, err := repo.Stage(".mp3", strings.NewReader(strings.Repeat("a", libraryDefault+1)))
require.NoError(t, err)
defer func() { require.NoError(t, work.Close()) }()
size, err := work.Size()
require.NoError(t, err)
require.Greater(t, size, int64(libraryDefault), "запись заведомо больше умолчания библиотеки")
file, err := repo.CreateLocal("big.mp3", work)
require.NoError(t, err, "запись длиннее умолчания библиотеки ложится в хранилище")
assert.Equal(t, size, file.Size)
assert.Greater(t, entity.MaxRecordSize, size, "объявленный потолок выше проверяемого размера")
}
@@ -0,0 +1,203 @@
package pocketbase
import (
"database/sql"
"time"
"github.com/pocketbase/pocketbase/core"
"github.com/pocketbase/pocketbase/tools/types"
"git.vakhrushev.me/av/transcriber/internal/entity"
)
// Отображение задачи в запись коллекции и обратно живёт одним местом. Прежде
// список колонок был переписан четырежды — в каждом запросе своего слоя, — и
// расхождение проявлялось как потерянное при сохранении поле.
// applyOwnedByPipeline кладёт в запись только те поля, которыми распоряжается
// конвейер. Поля, которые он не меняет никогда — куда отвечать отправителю и
// каким входом пришла запись, — не трогаются вовсе.
//
// Разрез нужен потому, что шаг держит задачу снимком с момента захвата и до
// своего сохранения, а это до восьми часов. Всё, что владелец правил в панели за
// это время, безусловная запись снимка стёрла бы молча: ни строки в журнале, ни
// отказа в панели — владелец видел бы успешное сохранение и был бы уверен, что
// правка на месте.
func applyOwnedByPipeline(record *core.Record, job *entity.TranscribeJob) {
record.Set("state", job.State)
record.Set("file", derefString(job.FileID))
record.Set("error_text", derefString(job.ErrorText))
record.Set("acquisition_id", derefString(job.AcquisitionID))
record.Set("acquire_time", dateOrEmpty(job.AcquireTime))
record.Set("delay_time", dateOrEmpty(job.DelayTime))
record.Set("attempts", job.Attempts)
record.Set("recognition_op_id", derefString(job.RecognitionOpID))
record.Set("transcription_text", derefString(job.TranscriptionText))
}
// applyToRecord кладёт задачу в запись целиком — это заведение, и спорить за
// поля здесь не с кем.
func applyToRecord(record *core.Record, job *entity.TranscribeJob) {
applyOwnedByPipeline(record, job)
record.Set("source", job.Source)
record.Set("tg_chat_id", derefInt64(job.TgChatId))
record.Set("tg_reply_message_id", derefInt(job.TgReplyMessageId))
}
func recordToJob(record *core.Record) *entity.TranscribeJob {
return &entity.TranscribeJob{
Id: record.Id,
State: record.GetString("state"),
Source: record.GetString("source"),
FileID: nilIfEmpty(record.GetString("file")),
ErrorText: nilIfEmpty(record.GetString("error_text")),
AcquisitionID: nilIfEmpty(record.GetString("acquisition_id")),
AcquireTime: timeOrNil(record.GetDateTime("acquire_time")),
DelayTime: timeOrNil(record.GetDateTime("delay_time")),
Attempts: record.GetInt("attempts"),
RecognitionOpID: nilIfEmpty(record.GetString("recognition_op_id")),
TranscriptionText: nilIfEmpty(record.GetString("transcription_text")),
TgChatId: nilIfZero64(int64(record.GetInt("tg_chat_id"))),
TgReplyMessageId: nilIfZeroInt(record.GetInt("tg_reply_message_id")),
CreatedAt: record.GetDateTime("created").Time(),
UpdatedAt: record.GetDateTime("updated").Time(),
}
}
// acquiredRow — задача, прочитанная сырым запросом захвата. Колонки читаются
// именно так, потому что запрос идёт мимо записей коллекции; связь с их
// перечнем держит константа acquireColumns и тест захвата, читающий задачу
// целиком.
type acquiredRow struct {
Id string `db:"id"`
State string `db:"state"`
Source string `db:"source"`
FileID sql.NullString `db:"file"`
ErrorText sql.NullString `db:"error_text"`
AcquisitionID sql.NullString `db:"acquisition_id"`
AcquireTime sql.NullString `db:"acquire_time"`
DelayTime sql.NullString `db:"delay_time"`
Attempts int `db:"attempts"`
RecognitionOpID sql.NullString `db:"recognition_op_id"`
TranscriptionText sql.NullString `db:"transcription_text"`
TgChatId sql.NullInt64 `db:"tg_chat_id"`
TgReplyMessageId sql.NullInt64 `db:"tg_reply_message_id"`
Created sql.NullString `db:"created"`
Updated sql.NullString `db:"updated"`
}
func (r *acquiredRow) toJob() *entity.TranscribeJob {
job := &entity.TranscribeJob{
Id: r.Id,
State: r.State,
Source: r.Source,
FileID: nullToPtr(r.FileID),
ErrorText: nullToPtr(r.ErrorText),
AcquisitionID: nullToPtr(r.AcquisitionID),
AcquireTime: parseTimeOrNil(r.AcquireTime),
DelayTime: parseTimeOrNil(r.DelayTime),
Attempts: r.Attempts,
RecognitionOpID: nullToPtr(r.RecognitionOpID),
TranscriptionText: nullToPtr(r.TranscriptionText),
}
if r.TgChatId.Valid && r.TgChatId.Int64 != 0 {
chatId := r.TgChatId.Int64
job.TgChatId = &chatId
}
if r.TgReplyMessageId.Valid && r.TgReplyMessageId.Int64 != 0 {
msgId := int(r.TgReplyMessageId.Int64)
job.TgReplyMessageId = &msgId
}
if created := parseTimeOrNil(r.Created); created != nil {
job.CreatedAt = *created
}
if updated := parseTimeOrNil(r.Updated); updated != nil {
job.UpdatedAt = *updated
}
return job
}
func derefString(v *string) string {
if v == nil {
return ""
}
return *v
}
func derefInt64(v *int64) int64 {
if v == nil {
return 0
}
return *v
}
func derefInt(v *int) int {
if v == nil {
return 0
}
return *v
}
// dateOrEmpty отдаёт пустое значение вместо нулевой даты: пустая колонка даты в
// хранилище это пустая строка, и она же значит «времени нет».
func dateOrEmpty(v *time.Time) any {
if v == nil {
return ""
}
date, err := types.ParseDateTime(*v)
if err != nil {
return ""
}
return date
}
func nilIfEmpty(v string) *string {
if v == "" {
return nil
}
return &v
}
func timeOrNil(v types.DateTime) *time.Time {
if v.IsZero() {
return nil
}
t := v.Time()
return &t
}
func nilIfZero64(v int64) *int64 {
if v == 0 {
return nil
}
return &v
}
func nilIfZeroInt(v int) *int {
if v == 0 {
return nil
}
return &v
}
func nullToPtr(v sql.NullString) *string {
if !v.Valid || v.String == "" {
return nil
}
s := v.String
return &s
}
func parseTimeOrNil(v sql.NullString) *time.Time {
if !v.Valid || v.String == "" {
return nil
}
date, err := types.ParseDateTime(v.String)
if err != nil || date.IsZero() {
return nil
}
t := date.Time()
return &t
}
@@ -0,0 +1,122 @@
package pocketbase
import (
"github.com/pocketbase/pocketbase/core"
"github.com/pocketbase/pocketbase/migrations"
"git.vakhrushev.me/av/transcriber/internal/entity"
)
// Схема заводится версионированными шагами, и применённый шаг не переписывается
// — только новым шагом. Инвариант проекта перенесён дословно: хранилище считает
// применённое по имени файла шага.
//
// Шаг регистрируется в списке приложения при загрузке пакета, а накатывает его
// `apis.Serve` прежде, чем поднять сервер.
func init() {
migrations.Register(up202608110001, down202608110001, "202608110001_init.go")
}
func up202608110001(app core.App) error {
files := core.NewBaseCollection(FilesCollection)
files.Fields.Add(
// Сам файл. Защищённым поле не помечено намеренно: право прочитать
// запись даёт знание её идентификатора, и файл встаёт вровень с опросом
// готовности задачи, а не ниже.
//
// Потолок задан **числом**: нулём библиотека читает не «без предела», а
// своё умолчание в 5 МиБ, и на нём отваливалось бы всё длиннее пяти
// минут. Число выведено из расчётного потолка записи в шесть часов с
// запасом на видео; оно же стоит строкой в docs/database.md.
&core.FileField{Name: "file", MaxSelect: 1, MaxSize: entity.MaxRecordSize},
// Где лежит копия. Поле названо `location`, а не `storage`: последним
// словом зовут само хранилище, и третий смысл развёл бы одно слово по
// разным вещам.
&core.SelectField{
Name: "location",
Values: []string{entity.LocationLocal, entity.LocationS3},
MaxSelect: 1,
Required: true,
},
// Ключ объекта во внешнем хранилище; у местной копии пуст.
&core.TextField{Name: "object_key"},
&core.NumberField{Name: "size", OnlyInt: true},
&core.AutodateField{Name: "created", OnCreate: true},
&core.AutodateField{Name: "updated", OnCreate: true, OnUpdate: true},
)
if err := app.Save(files); err != nil {
return err
}
jobs := core.NewBaseCollection(JobsCollection)
jobs.Fields.Add(
// Перечень состояний закрыт схемой: задача, заведённая в панели руками,
// не должна попасть в выборку с состоянием, которого конвейер не знает.
&core.SelectField{
Name: "state",
Values: []string{
entity.StateCreated,
entity.StateConverted,
entity.StateTranscribe,
entity.StateDone,
entity.StateFailed,
entity.StateDead,
},
MaxSelect: 1,
Required: true,
},
&core.SelectField{
Name: "source",
Values: []string{entity.SourceUnknown, entity.SourceApi, entity.SourceTelegram},
MaxSelect: 1,
Required: true,
},
// Текущий файл задачи: шаг конвейера переставляет ссылку на свой
// результат.
// Обязательна: задача без записи не может пройти ни одного шага, и
// заведённая в панели руками она дошла бы до шага только затем, чтобы
// отказать. Компилятор этого не держит — держит схема.
&core.RelationField{
Name: "file",
CollectionId: files.Id,
MaxSelect: 1,
Required: true,
},
&core.TextField{Name: "error_text"},
&core.TextField{Name: "acquisition_id"},
&core.DateField{Name: "acquire_time"},
&core.DateField{Name: "delay_time"},
// Число попыток: растёт при каждом захвате, обнуляется на шаге,
// завершившемся без отказа.
&core.NumberField{Name: "attempts", OnlyInt: true, Min: ptr(0.0)},
&core.TextField{Name: "recognition_op_id"},
&core.EditorField{Name: "transcription_text"},
&core.NumberField{Name: "tg_chat_id", OnlyInt: true},
&core.NumberField{Name: "tg_reply_message_id", OnlyInt: true},
&core.AutodateField{Name: "created", OnCreate: true},
&core.AutodateField{Name: "updated", OnCreate: true, OnUpdate: true},
)
// Выборка воркера идёт по состоянию, паузе и сроку захвата — индекс по
// состоянию снимает полный перебор, который был у прежней таблицы.
jobs.AddIndex("idx_transcribe_jobs_state", false, "state", "")
return app.Save(jobs)
}
func down202608110001(app core.App) error {
// Порядок обратный порядку заведения: задачи ссылаются на файлы.
for _, name := range []string{JobsCollection, FilesCollection} {
collection, err := app.FindCollectionByNameOrId(name)
if err != nil {
continue
}
if err := app.Delete(collection); err != nil {
return err
}
}
return nil
}
func ptr[T any](v T) *T { return &v }
+36
View File
@@ -0,0 +1,36 @@
package pocketbase
import (
"github.com/pocketbase/pocketbase/core"
)
// BindPanelRules подчиняет правку задачи в панели тем же правилам перехода, что
// и правку из кода.
//
// Панель — вход в задачу наравне с конвейером, а не окно просмотра: ради правки
// она и покупалась, мёртвая задача оживляется сменой состояния. Но правка полем
// идёт мимо кода, который чистит служебные поля прошлого состояния, и владелец,
// «вернувший задачу в работу», получил бы задачу с прежним признаком захвата
// (захвату она не выдастся до конца срока) и с числом попыток на пределе (умрёт
// от первого же отказа). Узнать об этом ему неоткуда.
//
// Хук стоит на правке **запросом**, а не на всяком сохранении записи. Модельное
// событие не различает, кто пишет, и срабатывало бы на каждом переходе
// конвейера: тогда задержка, поставленная шагом вместе со сменой состояния,
// стиралась бы тем же сохранением, а число попыток мёртвой задачи — которое
// переход хранит намеренно — приходило бы владельцу нулём.
func BindPanelRules(app core.App) {
app.OnRecordUpdateRequest(JobsCollection).BindFunc(func(e *core.RecordRequestEvent) error {
original := e.Record.Original()
if original == nil || original.GetString("state") == e.Record.GetString("state") {
return e.Next()
}
e.Record.Set("acquisition_id", "")
e.Record.Set("acquire_time", "")
e.Record.Set("delay_time", "")
e.Record.Set("attempts", 0)
return e.Next()
})
}
@@ -0,0 +1,147 @@
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
}
@@ -0,0 +1,410 @@
package pocketbase
import (
"net/http"
"net/http/httptest"
"strings"
"sync"
"testing"
"time"
"github.com/pocketbase/pocketbase/apis"
"github.com/pocketbase/pocketbase/core"
"github.com/pocketbase/pocketbase/tools/types"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"git.vakhrushev.me/av/transcriber/internal/contract"
"git.vakhrushev.me/av/transcriber/internal/entity"
)
// newTestApp поднимает хранилище на пустом каталоге и накатывает схему — тем же
// путём, каким это делает сервис при старте.
func newTestApp(t *testing.T) core.App {
t.Helper()
app, err := New(t.TempDir())
require.NoError(t, err)
t.Cleanup(func() {
if err := app.ResetBootstrapState(); err != nil {
t.Logf("не удалось закрыть хранилище: %v", err)
}
})
return app
}
// newFile заводит запись о файле: ссылка на неё у задачи обязательна схемой.
func newFile(t *testing.T, app core.App) *entity.File {
t.Helper()
repo := NewFileRepository(app)
work, err := repo.Stage(".mp3", strings.NewReader("запись"))
require.NoError(t, err)
defer func() { require.NoError(t, work.Close()) }()
file, err := repo.CreateLocal("sample.mp3", work)
require.NoError(t, err)
return file
}
func newJob(t *testing.T, repo *TranscriptJobRepository, state string) *entity.TranscribeJob {
t.Helper()
file := newFile(t, repo.app)
job := &entity.TranscribeJob{State: state, Source: entity.SourceApi, FileID: &file.Id}
require.NoError(t, repo.Create(job))
return job
}
// Захват неделим: выбор подходящей задачи и пометка её захваченной идут вместе.
// Двум вызывающим, пришедшим за одним состоянием разом, запись достаётся
// одному — на этом стоит инвариант «Принятая запись не теряется молча».
func TestFindAndAcquire_OnlyOneOfThreeGetsTheJob(t *testing.T) {
app := newTestApp(t)
repo := NewTranscriptJobRepository(app)
job := newJob(t, repo, entity.StateCreated)
const racers = 3
var (
wg sync.WaitGroup
mu sync.Mutex
got []*entity.TranscribeJob
notFound int
)
start := make(chan struct{})
for i := 0; i < racers; i++ {
wg.Add(1)
go func(n int) {
defer wg.Done()
<-start
acquired, err := repo.FindAndAcquire(entity.StateCreated, "holder", time.Now().Add(-time.Hour))
mu.Lock()
defer mu.Unlock()
if err != nil {
var missing *contract.JobNotFoundError
if assert.ErrorAs(t, err, &missing) {
notFound++
}
return
}
got = append(got, acquired)
}(i)
}
close(start)
wg.Wait()
require.Len(t, got, 1, "запись получает ровно один из трёх захватов")
assert.Equal(t, job.Id, got[0].Id)
assert.Equal(t, racers-1, notFound, "остальные получают признак «работы нет»")
}
// Захваченная задача второй раз не выдаётся, пока срок захвата не истёк.
func TestFindAndAcquire_AcquiredJobIsNotHandedOutAgain(t *testing.T) {
app := newTestApp(t)
repo := NewTranscriptJobRepository(app)
newJob(t, repo, entity.StateCreated)
first, err := repo.FindAndAcquire(entity.StateCreated, "first", time.Now().Add(-time.Hour))
require.NoError(t, err)
require.NotNil(t, first)
_, err = repo.FindAndAcquire(entity.StateCreated, "second", time.Now().Add(-time.Hour))
var missing *contract.JobNotFoundError
assert.ErrorAs(t, err, &missing, "захваченная задача второму не выдаётся")
}
// Захват протухает, и задача достаётся снова. Время захвата кладётся **не**
// нашим кодом, а тем же путём, что и `created`: проверка, кладущая его своим
// форматом, была бы зелена и тогда, когда сравнение вида сломано.
func TestFindAndAcquire_RottenAcquisitionIsHandedOutAgain(t *testing.T) {
app := newTestApp(t)
repo := NewTranscriptJobRepository(app)
job := newJob(t, repo, entity.StateCreated)
_, err := repo.FindAndAcquire(entity.StateCreated, "first", time.Now().Add(-time.Hour))
require.NoError(t, err)
// Задним числом — записью коллекции, то есть тем же слоем, который пишет
// собственные времена хранилища.
record, err := app.FindRecordById(JobsCollection, job.Id)
require.NoError(t, err)
record.Set("acquire_time", types.NowDateTime().Add(-2*time.Hour))
require.NoError(t, app.Save(record))
again, err := repo.FindAndAcquire(entity.StateCreated, "second", time.Now().Add(-time.Hour))
require.NoError(t, err, "протухший захват не мешает выдать задачу следующему")
assert.Equal(t, job.Id, again.Id)
}
// Пауза держит задачу от выдачи, пока не кончится.
func TestFindAndAcquire_DelayedJobIsNotHandedOut(t *testing.T) {
app := newTestApp(t)
repo := NewTranscriptJobRepository(app)
job := newJob(t, repo, entity.StateCreated)
delay := time.Now().Add(time.Hour)
job.DelayTime = &delay
require.NoError(t, repo.Save(job, ""))
_, err := repo.FindAndAcquire(entity.StateCreated, "holder", time.Now().Add(-time.Hour))
var missing *contract.JobNotFoundError
assert.ErrorAs(t, err, &missing, "задача не выдаётся, пока пауза не кончилась")
}
// Число попыток растёт при каждом захвате: только так попытка засчитывается и
// задаче, брошенной вместе с процессом.
func TestFindAndAcquire_AttemptsGrowOnEveryAcquisition(t *testing.T) {
app := newTestApp(t)
repo := NewTranscriptJobRepository(app)
newJob(t, repo, entity.StateCreated)
for expected := 1; expected <= 3; expected++ {
acquired, err := repo.FindAndAcquire(entity.StateCreated, "holder", time.Now().Add(time.Hour))
require.NoError(t, err)
assert.Equal(t, expected, acquired.Attempts)
}
}
// Захват отдаёт задачу целиком, а не только её ключ: сырой запрос идёт мимо
// записей коллекции, и расхождение перечня колонок иначе проявилось бы как
// потерянное поле.
func TestFindAndAcquire_ReturnsWholeJob(t *testing.T) {
app := newTestApp(t)
repo := NewTranscriptJobRepository(app)
chatId := int64(4242)
replyId := 17
opId := "operation-id"
text := "расшифровка"
file := newFile(t, app)
job := &entity.TranscribeJob{
State: entity.StateTranscribe,
Source: entity.SourceTelegram,
FileID: &file.Id,
TgChatId: &chatId,
TgReplyMessageId: &replyId,
RecognitionOpID: &opId,
TranscriptionText: &text,
}
require.NoError(t, repo.Create(job))
acquired, err := repo.FindAndAcquire(entity.StateTranscribe, "holder", time.Now().Add(-time.Hour))
require.NoError(t, err)
assert.Equal(t, job.Id, acquired.Id)
assert.Equal(t, entity.StateTranscribe, acquired.State)
assert.Equal(t, entity.SourceTelegram, acquired.Source)
require.NotNil(t, acquired.TgChatId)
assert.Equal(t, chatId, *acquired.TgChatId)
require.NotNil(t, acquired.TgReplyMessageId)
assert.Equal(t, replyId, *acquired.TgReplyMessageId)
require.NotNil(t, acquired.RecognitionOpID)
assert.Equal(t, opId, *acquired.RecognitionOpID)
require.NotNil(t, acquired.TranscriptionText)
assert.Equal(t, text, *acquired.TranscriptionText)
assert.False(t, acquired.CreatedAt.IsZero(), "время заведения доехало")
}
// Шаг, потерявший захват за время работы, результата не пишет: иначе два
// воркера пишут в одну задачу по очереди, а отправитель получает два ответа.
func TestSave_RefusesWriteFromLostAcquisition(t *testing.T) {
app := newTestApp(t)
repo := NewTranscriptJobRepository(app)
newJob(t, repo, entity.StateCreated)
mine, err := repo.FindAndAcquire(entity.StateCreated, "mine", time.Now().Add(-time.Hour))
require.NoError(t, err)
// Задача досталась другому, пока шаг работал.
record, err := app.FindRecordById(JobsCollection, mine.Id)
require.NoError(t, err)
record.Set("acquisition_id", "someone-else")
require.NoError(t, app.Save(record))
mine.MoveToState(entity.StateConverted)
err = repo.Save(mine, "mine")
var lost *contract.LostAcquisitionError
require.ErrorAs(t, err, &lost)
// И состояние не поехало.
after, err := repo.GetByID(mine.Id)
require.NoError(t, err)
assert.Equal(t, entity.StateCreated, after.State)
}
// Пустой держатель значит «задача не захватывалась» — так её сохраняет приём.
func TestSave_WithoutHolderWritesAnyway(t *testing.T) {
app := newTestApp(t)
repo := NewTranscriptJobRepository(app)
job := newJob(t, repo, entity.StateCreated)
job.MoveToState(entity.StateConverted)
require.NoError(t, repo.Save(job, ""))
after, err := repo.GetByID(job.Id)
require.NoError(t, err)
assert.Equal(t, entity.StateConverted, after.State)
}
// Правка состояния **запросом** — то есть из панели — чистит служебные поля
// прошлого состояния: те же, что чистит переход из кода. Иначе владелец,
// вернувший мёртвую задачу в работу, получил бы задачу, которая не выдаётся
// захвату и умирает от первого же отказа, и не узнал бы об этом.
func TestPanelRules_StateChangeByRequestClearsAcquisition(t *testing.T) {
app := newTestApp(t)
BindPanelRules(app)
repo := NewTranscriptJobRepository(app)
job := newJob(t, repo, entity.StateCreated)
acquired, err := repo.FindAndAcquire(entity.StateCreated, "holder", time.Now().Add(-time.Hour))
require.NoError(t, err)
require.NotNil(t, acquired.AcquisitionID)
record, err := app.FindRecordById(JobsCollection, job.Id)
require.NoError(t, err)
record.Set("attempts", 5)
record.Set("state", entity.StateDead)
require.NoError(t, app.Save(record))
// Владелец возвращает задачу в работу правкой состояния в панели — то есть
// запросом к записи, а не сохранением из кода.
patchRecord(t, app, job.Id, `{"state":"`+entity.StateCreated+`"}`)
after, err := repo.GetByID(job.Id)
require.NoError(t, err)
assert.Nil(t, after.AcquisitionID, "признак захвата снят")
assert.Nil(t, after.AcquireTime, "время захвата снято")
assert.Nil(t, after.DelayTime, "пауза снята")
assert.Equal(t, 0, after.Attempts, "число попыток обнулено")
// И ближайший захват задачу выдаёт.
again, err := repo.FindAndAcquire(entity.StateCreated, "next", time.Now().Add(-time.Hour))
require.NoError(t, err)
assert.Equal(t, job.Id, again.Id)
}
// Обратная сторона того же правила, и она дороже: правила панели MUST не
// трогать записи, которые правит сам конвейер. Модельный хук их не различал, и
// пауза, поставленная шагом вместе со сменой состояния, стиралась тем же
// сохранением, а число попыток мёртвой задачи приходило владельцу нулём.
func TestPanelRules_DoNotTouchPipelineWrites(t *testing.T) {
app := newTestApp(t)
BindPanelRules(app)
repo := NewTranscriptJobRepository(app)
job := newJob(t, repo, entity.StateConverted)
acquired, err := repo.FindAndAcquire(entity.StateConverted, "holder", time.Now().Add(-time.Hour))
require.NoError(t, err)
// Шаг ставит задержку опроса вместе со сменой состояния.
delay := time.Now().Add(10 * time.Second)
acquired.MoveToStateAndDelay(entity.StateTranscribe, &delay)
require.NoError(t, repo.Save(acquired, "holder"))
after, err := repo.GetByID(job.Id)
require.NoError(t, err)
require.NotNil(t, after.DelayTime, "задержка, поставленная шагом, пережила сохранение")
// Переход в «мертва» хранит число попыток намеренно: по нему владелец видит,
// сколько раз мы пробовали.
after.Attempts = 6
after.Die("attempts exhausted: 6")
require.NoError(t, repo.Save(after, ""))
dead, err := repo.GetByID(job.Id)
require.NoError(t, err)
assert.Equal(t, entity.StateDead, dead.State)
assert.Equal(t, 6, dead.Attempts, "число попыток мёртвой задачи сохранено")
}
// patchRecord правит запись тем же путём, каким её правит панель: запросом к
// API от имени владельца.
func patchRecord(t *testing.T, app core.App, recordID, body string) {
t.Helper()
superusers, err := app.FindCollectionByNameOrId(core.CollectionNameSuperusers)
require.NoError(t, err)
owner := core.NewRecord(superusers)
owner.Set("email", "owner@example.com")
owner.Set("password", "ownerpassword123")
require.NoError(t, app.Save(owner))
token, err := owner.NewStaticAuthToken(time.Hour)
require.NoError(t, err)
router, err := apis.NewRouter(app)
require.NoError(t, err)
mux, err := router.BuildMux()
require.NoError(t, err)
req := httptest.NewRequest(
http.MethodPatch,
"/api/collections/"+JobsCollection+"/records/"+recordID,
strings.NewReader(body),
)
req.Header.Set("Content-Type", "application/json")
req.Header.Set("Authorization", token)
w := httptest.NewRecorder()
mux.ServeHTTP(w, req)
require.Equal(t, http.StatusOK, w.Code, "правка записи владельцем: %s", w.Body.String())
}
// Правка владельца в панели переживает сохранение шага. Шаг держит задачу
// снимком с момента захвата и до своего сохранения — до восьми часов, — и
// безусловная запись снимка стёрла бы правку молча: ни строки в журнале, ни
// отказа в панели.
func TestSave_KeepsOwnerEditMadeWhileStepHeldTheJob(t *testing.T) {
app := newTestApp(t)
BindPanelRules(app)
repo := NewTranscriptJobRepository(app)
file := newFile(t, app)
chatId := int64(111)
job := &entity.TranscribeJob{
State: entity.StateCreated,
Source: entity.SourceTelegram,
FileID: &file.Id,
TgChatId: &chatId,
}
require.NoError(t, repo.Create(job))
// Шаг захватил задачу и работает.
acquired, err := repo.FindAndAcquire(entity.StateCreated, "holder", time.Now().Add(-time.Hour))
require.NoError(t, err)
// Владелец правит в панели поле, которого конвейер не касается.
patchRecord(t, app, job.Id, `{"tg_chat_id":999999}`)
// Шаг доработал и сохраняет свой снимок.
acquired.MoveToState(entity.StateConverted)
require.NoError(t, repo.Save(acquired, "holder"))
after, err := repo.GetByID(job.Id)
require.NoError(t, err)
assert.Equal(t, entity.StateConverted, after.State, "шаг свой результат записал")
require.NotNil(t, after.TgChatId)
assert.Equal(t, int64(999999), *after.TgChatId, "правка владельца пережила сохранение шага")
}
-56
View File
@@ -1,56 +0,0 @@
package sqlite
import (
"database/sql"
"fmt"
"git.vakhrushev.me/av/transcriber/internal/entity"
"github.com/doug-martin/goqu/v9"
)
type FileRepository struct {
db *sql.DB
gq *goqu.Database
}
func NewFileRepository(conn *sql.DB, gq *goqu.Database) *FileRepository {
return &FileRepository{conn, gq}
}
func (repo *FileRepository) Create(file *entity.File) error {
record := goqu.Record{
"id": file.Id,
"storage": file.Storage,
"file_name": file.FileName,
"size": file.Size,
"created_at": file.CreatedAt,
}
query := repo.gq.Insert("files").Rows(record)
sql, args, err := query.ToSQL()
if err != nil {
return fmt.Errorf("failed to build query: %w", err)
}
_, err = repo.db.Exec(sql, args...)
if err != nil {
return fmt.Errorf("failed to insert file: %w", err)
}
return nil
}
func (repo *FileRepository) GetByID(id string) (*entity.File, error) {
query := repo.gq.From("files").Select("id", "storage", "file_name", "size", "created_at").Where(goqu.C("id").Eq(id))
sql, args, err := query.ToSQL()
if err != nil {
return nil, fmt.Errorf("failed to build query: %w", err)
}
var file entity.File
err = repo.db.QueryRow(sql, args...).Scan(&file.Id, &file.Storage, &file.FileName, &file.Size, &file.CreatedAt)
if err != nil {
return nil, fmt.Errorf("failed to get file: %w", err)
}
return &file, nil
}
@@ -1,230 +0,0 @@
package sqlite
import (
"database/sql"
"fmt"
"time"
"git.vakhrushev.me/av/transcriber/internal/contract"
"git.vakhrushev.me/av/transcriber/internal/entity"
goqu "github.com/doug-martin/goqu/v9"
)
type TranscriptJobRepository struct {
db *sql.DB
gq *goqu.Database
}
func NewTranscriptJobRepository(db *sql.DB, gq *goqu.Database) *TranscriptJobRepository {
return &TranscriptJobRepository{db, gq}
}
func (repo *TranscriptJobRepository) Create(job *entity.TranscribeJob) error {
record := goqu.Record{
"id": job.Id,
"state": job.State,
"source": job.Source,
"file_id": job.FileID,
"is_error": job.IsError,
"error_text": job.ErrorText,
"acquisition_id": job.AcquisitionID,
"acquire_time": job.AcquireTime,
"delay_time": job.DelayTime,
"recognition_op_id": job.RecognitionOpID,
"transcription_text": job.TranscriptionText,
"tg_chat_id": job.TgChatId,
"tg_reply_message_id": job.TgReplyMessageId,
"created_at": job.CreatedAt,
"updated_at": job.UpdatedAt,
}
query := repo.gq.Insert("transcribe_jobs").Rows(record)
sql, args, err := query.ToSQL()
if err != nil {
return fmt.Errorf("failed to build query: %w", err)
}
_, err = repo.db.Exec(sql, args...)
if err != nil {
return fmt.Errorf("failed to insert transcribe job: %w", err)
}
return nil
}
func (repo *TranscriptJobRepository) Save(job *entity.TranscribeJob) error {
record := goqu.Record{
"state": job.State,
"source": job.Source,
"file_id": job.FileID,
"is_error": job.IsError,
"error_text": job.ErrorText,
"acquisition_id": job.AcquisitionID,
"acquire_time": job.AcquireTime,
"delay_time": job.DelayTime,
"recognition_op_id": job.RecognitionOpID,
"transcription_text": job.TranscriptionText,
"tg_chat_id": job.TgChatId,
"tg_reply_message_id": job.TgReplyMessageId,
"updated_at": job.UpdatedAt,
}
query := repo.gq.Update("transcribe_jobs").Set(record).Where(goqu.C("id").Eq(job.Id))
sql, args, err := query.ToSQL()
if err != nil {
return fmt.Errorf("failed to build query: %w", err)
}
_, err = repo.db.Exec(sql, args...)
if err != nil {
return fmt.Errorf("failed to update transcribe job: %w", err)
}
return nil
}
func (repo *TranscriptJobRepository) GetByID(id string) (*entity.TranscribeJob, error) {
query := repo.gq.From("transcribe_jobs").Select(
"id",
"state",
"source",
"file_id",
"is_error",
"error_text",
"acquisition_id",
"acquire_time",
"delay_time",
"recognition_op_id",
"transcription_text",
"tg_chat_id",
"tg_reply_message_id",
"created_at",
"updated_at",
).Where(goqu.C("id").Eq(id))
sql, args, err := query.ToSQL()
if err != nil {
return nil, fmt.Errorf("failed to build query: %w", err)
}
var job entity.TranscribeJob
err = repo.db.QueryRow(sql, args...).Scan(
&job.Id,
&job.State,
&job.Source,
&job.FileID,
&job.IsError,
&job.ErrorText,
&job.AcquisitionID,
&job.AcquireTime,
&job.DelayTime,
&job.RecognitionOpID,
&job.TranscriptionText,
&job.TgChatId,
&job.TgReplyMessageId,
&job.CreatedAt,
&job.UpdatedAt,
)
if err != nil {
return nil, fmt.Errorf("failed to get transcribe job: %w", err)
}
return &job, nil
}
func (repo *TranscriptJobRepository) FindAndAcquire(state, acquisitionId string, rottingTime time.Time) (*entity.TranscribeJob, error) {
updateQuery := repo.gq.Update("transcribe_jobs").
Set(
goqu.Record{
"acquisition_id": acquisitionId,
"acquire_time": time.Now(),
},
).
Where(
goqu.C("id").Eq(
repo.gq.From("transcribe_jobs").Select("id").
Where(
goqu.And(
goqu.C("state").Eq(state),
goqu.C("is_error").Eq(0),
goqu.Or(
goqu.C("delay_time").IsNull(),
goqu.C("delay_time").Lt(time.Now()),
),
goqu.Or(
goqu.C("acquisition_id").IsNull(),
goqu.C("acquire_time").Lt(rottingTime),
),
),
).
Limit(1),
),
)
sql, args, err := updateQuery.ToSQL()
if err != nil {
return nil, fmt.Errorf("failed to build query: %w", err)
}
// log.Printf("aquire sql: %s", sql)
result, err := repo.db.Exec(sql, args...)
if err != nil {
return nil, fmt.Errorf("failed to aquire job with state %s: %w", state, err)
}
rowsAffected, err := result.RowsAffected()
if err != nil {
return nil, fmt.Errorf("failed check affected rows: %w", err)
}
if rowsAffected == 0 {
e := contract.JobNotFoundError{State: state, Message: "appropriate job not found"}
return nil, &e
}
if rowsAffected != 1 {
return nil, fmt.Errorf("unexpected affected rows count: %d", rowsAffected)
}
selectQuery := repo.gq.From("transcribe_jobs").Select(
"id",
"state",
"source",
"file_id",
"is_error",
"error_text",
"acquisition_id",
"acquire_time",
"delay_time",
"recognition_op_id",
"transcription_text",
"tg_chat_id",
"tg_reply_message_id",
"created_at",
"updated_at",
).Where(goqu.C("acquisition_id").Eq(acquisitionId))
sql, args, err = selectQuery.ToSQL()
if err != nil {
return nil, fmt.Errorf("failed to build query: %w", err)
}
var job entity.TranscribeJob
err = repo.db.QueryRow(sql, args...).Scan(
&job.Id,
&job.State,
&job.Source,
&job.FileID,
&job.IsError,
&job.ErrorText,
&job.AcquisitionID,
&job.AcquireTime,
&job.DelayTime,
&job.RecognitionOpID,
&job.TranscriptionText,
&job.TgChatId,
&job.TgReplyMessageId,
&job.CreatedAt,
&job.UpdatedAt,
)
if err != nil {
return nil, fmt.Errorf("failed to get transcribe job: %w", err)
}
return &job, nil
}
+4 -10
View File
@@ -9,7 +9,6 @@ import (
type Config struct {
Server ServerConfig `toml:"server"`
Database DatabaseConfig `toml:"database"`
Storage StorageConfig `toml:"storage"`
Yandex YandexConfig `toml:"yandex"`
Telegram TelegramConfig `toml:"telegram"`
@@ -22,12 +21,10 @@ type ServerConfig struct {
UsersWhiteList []string `toml:"users_while_list"`
}
type DatabaseConfig struct {
Path string `toml:"path"`
}
// StorageConfig — единственный каталог данных: под ним лежат и база, и файлы
// записей. Двух путей, как было раньше, у хранилища не бывает.
type StorageConfig struct {
Path string `toml:"path"`
DataDir string `toml:"data_dir"`
}
type YandexConfig struct {
@@ -53,11 +50,8 @@ func defaultConfig() *Config {
ShutdownTimeout: 5,
ForceShutdownTimeout: 20,
},
Database: DatabaseConfig{
Path: "data/transcriber.db",
},
Storage: StorageConfig{
Path: "data/files",
DataDir: "data",
},
Yandex: YandexConfig{
FolderID: "",
+12
View File
@@ -11,6 +11,18 @@ func (e *JobNotFoundError) Error() string {
return fmt.Sprintf("%s - %s", e.State, e.Message)
}
// LostAcquisitionError — захват задачи за время работы шага достался другому.
// Шаг, получивший его, завершается без записи результата и без ответа
// отправителю: иначе два воркера пишут в одну задачу по очереди, а отправитель
// получает два ответа на одну запись.
type LostAcquisitionError struct {
JobID string
}
func (e *LostAcquisitionError) Error() string {
return fmt.Sprintf("%s: job acquisition lost", e.JobID)
}
type NoopJobError struct {
State string
}
+41 -2
View File
@@ -1,19 +1,58 @@
package contract
import (
"io"
"time"
"git.vakhrushev.me/av/transcriber/internal/entity"
)
// WorkFile — рабочая копия файла на диске: её просят шаги, отдающие файл
// внешней программе, потому что `ffmpeg` и `ffprobe` принимают имя аргументом.
//
// Заводится копия одним способом — репозиторием файлов, — и убирает её за собой
// Close. Каждый шаг, заводящий копию сам, повторял бы и обязанность прибрать, а
// забытая копия это шестичасовая запись во временном каталоге, о которой не
// узнает никто.
type WorkFile interface {
// Path — имя копии на диске, годное для внешней программы.
Path() string
// Size — длина копии в байтах на момент вызова.
Size() (int64, error)
// Close убирает копию. Зовётся на любом исходе, включая отказ.
Close() error
}
type FileRepository interface {
Create(file *entity.File) error
// Stage принимает содержимое потоком в рабочую копию с заданным
// расширением: по нему внешняя программа выбирает разбор. В память запись
// целиком не читается — расчётный потолок шесть часов.
Stage(ext string, content io.Reader) (WorkFile, error)
// StageEmpty заводит пустую рабочую копию с заданным расширением — под
// результат внешней программы, которая пишет по имени.
StageEmpty(ext string) (WorkFile, error)
// Localize выдаёт рабочую копию хранимого файла.
Localize(fileID string) (WorkFile, error)
// CreateLocal кладёт рабочую копию в хранилище под именем name и заводит
// запись о файле. Имя задаёт сервис: умолчание хранилища, строящее его из
// имени отправителя, не применяется.
CreateLocal(name string, work WorkFile) (*entity.File, error)
// CreateRemote заводит запись о копии, лежащей во внешнем хранилище.
CreateRemote(objectKey string, size int64) (*entity.File, error)
GetByID(id string) (*entity.File, error)
// Open отдаёт содержимое хранимого файла потоком.
Open(fileID string) (io.ReadCloser, error)
}
type TranscriptJobRepository interface {
Create(job *entity.TranscribeJob) error
Save(job *entity.TranscribeJob) error
// Save сохраняет задачу, захват которой держит holder. Захват, доставшийся
// за время работы другому, даёт LostAcquisitionError и запись не проводит.
// Пустой holder снимает эту условность и в конвейере не употребляется: все
// его шаги получают признак захвата от FindAndAcquire.
Save(job *entity.TranscribeJob, holder string) error
GetByID(id string) (*entity.TranscribeJob, error)
// FindAndAcquire забирает задачу одним неделимым шагом и увеличивает число
// её попыток. Работы в состоянии нет — JobNotFoundError.
FindAndAcquire(state, acquisitionId string, rottingTime time.Time) (*entity.TranscribeJob, error)
}
+41 -51
View File
@@ -1,22 +1,30 @@
package http
import (
"log"
"log/slog"
"net/http"
"time"
"github.com/pocketbase/pocketbase/apis"
"github.com/pocketbase/pocketbase/core"
"github.com/pocketbase/pocketbase/tools/router"
"git.vakhrushev.me/av/transcriber/internal/contract"
"git.vakhrushev.me/av/transcriber/internal/entity"
"git.vakhrushev.me/av/transcriber/internal/service"
"github.com/gin-gonic/gin"
)
type TranscribeHandler struct {
jobRepo contract.TranscriptJobRepository
trsService *service.TranscribeService
logger *slog.Logger
}
func NewTranscribeHandler(jobRepo contract.TranscriptJobRepository, trsService *service.TranscribeService) *TranscribeHandler {
return &TranscribeHandler{jobRepo: jobRepo, trsService: trsService}
func NewTranscribeHandler(jobRepo contract.TranscriptJobRepository, trsService *service.TranscribeService, logger *slog.Logger) *TranscribeHandler {
if logger == nil {
logger = slog.Default()
}
return &TranscribeHandler{jobRepo: jobRepo, trsService: trsService, logger: logger}
}
type CreateTranscribeJobResponse struct {
@@ -31,74 +39,56 @@ type GetTranscribeJobResponse struct {
TranscriptionText *string `json:"transcription_text,omitempty"`
}
func (h *TranscribeHandler) CreateTranscribeJob(c *gin.Context) {
// Register вешает маршруты сервиса на роутер хранилища. Порт у сервиса и у
// панели один, поэтому и роутер один; имена полей ответа и коды при переезде
// сохранены — публичный контракт API объявлен необратимым.
func (h *TranscribeHandler) Register(r *router.Router[*core.RequestEvent]) {
api := r.Group("/api")
// Умолчание роутера хранилища — 32 МиБ на тело, и оно отсекало бы запись
// раньше обработчика, без строки в журнале приёма. Приём размеру не судья,
// поэтому предел тела равен потолку самой записи.
api.POST("/audio", h.CreateTranscribeJob).Bind(apis.BodyLimit(entity.MaxRecordSize))
api.GET("/status/{id}", h.GetTranscribeJobStatus)
}
func (h *TranscribeHandler) CreateTranscribeJob(e *core.RequestEvent) error {
// Получаем файл из формы
file, header, err := c.Request.FormFile("audio")
file, header, err := e.Request.FormFile("audio")
if err != nil {
c.JSON(http.StatusBadRequest, gin.H{"error": "No audio file provided"})
return
return e.JSON(http.StatusBadRequest, map[string]string{"error": "No audio file provided"})
}
defer file.Close()
defer func() {
if err := file.Close(); err != nil {
h.logger.Error("Failed to close uploaded file", "error", err)
}
}()
job, err := h.trsService.CreateJobFromApi(file, header.Filename)
if err != nil {
log.Printf("Err: %v", err)
c.JSON(http.StatusInternalServerError, gin.H{"error": "Failed to create transcibe job"})
return
// Второй раз отказ не логируем: приём назван конвенцией логирующей
// границей и уже написал о нём. Транспорт переводит ошибку в ответ.
return e.JSON(http.StatusInternalServerError, map[string]string{"error": "Failed to create transcibe job"})
}
// Возвращаем успешный ответ
response := CreateTranscribeJobResponse{
return e.JSON(http.StatusCreated, CreateTranscribeJobResponse{
JobID: job.Id,
State: job.State,
}
c.JSON(http.StatusCreated, response)
})
}
func (h *TranscribeHandler) GetTranscribeJobStatus(c *gin.Context) {
jobID := c.Param("id")
func (h *TranscribeHandler) GetTranscribeJobStatus(e *core.RequestEvent) error {
jobID := e.Request.PathValue("id")
job, err := h.jobRepo.GetByID(jobID)
if err != nil {
c.JSON(http.StatusNotFound, gin.H{"error": "Job not found"})
return
return e.JSON(http.StatusNotFound, map[string]string{"error": "Job not found"})
}
c.JSON(http.StatusOK, GetTranscribeJobResponse{
return e.JSON(http.StatusOK, GetTranscribeJobResponse{
JobID: job.Id,
State: job.State,
CreatedAt: job.CreatedAt,
TranscriptionText: job.TranscriptionText,
})
}
func (h *TranscribeHandler) RunConversionJob(c *gin.Context) {
err := h.trsService.FindAndRunConversionJob()
if err != nil {
c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()})
return
}
c.Status(http.StatusOK)
}
func (h *TranscribeHandler) RunTranscribeJob(c *gin.Context) {
err := h.trsService.FindAndRunTranscribeJob()
if err != nil {
c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()})
return
}
c.Status(http.StatusOK)
}
func (h *TranscribeHandler) RunRecognitionCheckJob(c *gin.Context) {
err := h.trsService.FindAndRunTranscribeCheckJob()
if err != nil {
c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()})
return
}
c.Status(http.StatusOK)
}
+188 -187
View File
@@ -2,7 +2,6 @@ package http
import (
"bytes"
"database/sql"
"encoding/json"
"errors"
"fmt"
@@ -10,30 +9,22 @@ import (
"mime/multipart"
"net/http"
"net/http/httptest"
"os"
"path"
"path/filepath"
"regexp"
"runtime"
"strings"
"sync"
"testing"
"time"
"github.com/pocketbase/pocketbase/apis"
"github.com/pocketbase/pocketbase/core"
"github.com/prometheus/client_golang/prometheus"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"git.vakhrushev.me/av/transcriber/internal/adapter/recognizer"
"git.vakhrushev.me/av/transcriber/internal/adapter/repo/sqlite"
pbrepo "git.vakhrushev.me/av/transcriber/internal/adapter/repo/pocketbase"
"git.vakhrushev.me/av/transcriber/internal/contract"
"git.vakhrushev.me/av/transcriber/internal/entity"
"git.vakhrushev.me/av/transcriber/internal/service"
"github.com/doug-martin/goqu/v9"
_ "github.com/doug-martin/goqu/v9/dialect/sqlite3"
"github.com/gin-gonic/gin"
_ "github.com/mattn/go-sqlite3"
"github.com/pressly/goose/v3"
"github.com/prometheus/client_golang/prometheus"
sloggin "github.com/samber/slog-gin"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
// Подставные адаптеры вместо ffprobe и ffmpeg. Проверки судят приём — что
@@ -72,21 +63,19 @@ func readableMetaViewer() *stubMetaViewer {
return &stubMetaViewer{seconds: 42}
}
// testEnv — собранное окружение одной проверки. Каталог хранения свой у
// каждой: рабочий каталог процесса проверки не трогают.
// testEnv — собранное окружение одной проверки. Каталог данных свой у каждой:
// рабочий каталог процесса проверки не трогают.
type testEnv struct {
router *gin.Engine
handler *TranscribeHandler
db *sql.DB
storageDir string
journal *journalBuffer
mux http.Handler
handler *TranscribeHandler
app core.App
journal *journalBuffer
}
// journalBuffer — перехваченный журнал одной проверки. Свой на случай: общий на
// пакет сделал бы исход функцией от соседних случаев — «поля на месте» прошло бы
// на чужой строке, а «маркера нет» покраснело бы от чужой. Замок нужен потому,
// что пишущих в него потоков три: логгер сервиса, стандартный `log` транспорта
// и middleware запроса.
// что пишущих в него потоков два: логгер сервиса и логгер обработчика.
type journalBuffer struct {
mu sync.Mutex
text strings.Builder
@@ -106,44 +95,29 @@ func (b *journalBuffer) String() string {
return b.text.String()
}
func setupTestDB(t *testing.T) (*sql.DB, *goqu.Database) {
db, err := sql.Open("sqlite3", ":memory:")
// newTestStorage поднимает хранилище на пустом каталоге и накатывает схему —
// ровно тем же путём, каким это делает сервис при старте.
func newTestStorage(t *testing.T) core.App {
t.Helper()
app, err := pbrepo.New(t.TempDir())
require.NoError(t, err)
t.Cleanup(func() { db.Close() })
t.Cleanup(func() {
if err := app.ResetBootstrapState(); err != nil {
t.Logf("не удалось закрыть хранилище: %v", err)
}
})
// Каждому новому соединению с `:memory:` драйвер выдаёт свою базу, и
// второй потребитель пула не увидел бы накатанных миграций. Одно
// соединение снимает класс целиком.
db.SetMaxOpenConns(1)
gq := goqu.New("sqlite3", db)
err = goose.SetDialect("sqlite3")
require.NoError(t, err)
_, b, _, _ := runtime.Caller(0)
migpath, err := filepath.Abs(path.Join(b, "../../../../migrations"))
require.NoError(t, err)
goose.SetLogger(goose.NopLogger())
err = goose.Up(db, migpath)
require.NoError(t, err)
return db, gq
return app
}
func setupTestEnv(t *testing.T, metaviewer contract.AudioMetaViewer) *testEnv {
gin.SetMode(gin.TestMode)
app := newTestStorage(t)
db, gq := setupTestDB(t)
pbrepo.BindPanelRules(app)
storageDir := filepath.Join(t.TempDir(), "files")
require.NoError(t, os.MkdirAll(storageDir, 0o755))
fileRepo := sqlite.NewFileRepository(db, gq)
jobRepo := sqlite.NewTranscriptJobRepository(db, gq)
fileRepo := pbrepo.NewFileRepository(app)
jobRepo := pbrepo.NewTranscriptJobRepository(app)
// Журнал уходит в буфер, а не в никуда: по нему судит проверка запрета на
// имя отправителя. Вывод прогона от этого не меняется — ERROR-строки ветки
@@ -152,20 +126,6 @@ func setupTestEnv(t *testing.T, metaviewer contract.AudioMetaViewer) *testEnv {
journal := &journalBuffer{}
logger := slog.New(slog.NewTextHandler(journal, nil))
// Второй писатель журнала приёма — HTTP-транспорт: он пишет через стандартный
// `log` (расхождение записано в docs/conventions/logging.md). В бою `main.go`
// зовёт `slog.SetDefault`, и такая запись садится в `msg` строки `slog`;
// повторяем это здесь, чтобы оракул видел ту же цепочку, что и прод, а не
// свою. Без перехвата оракул был бы уже требования, которое накрывает все
// журнальные записи приёма.
//
// Подмена процессная, а не своя у случая: `t.Parallel()` в этом файле
// запрещён. При параллельных случаях вывод указывал бы на буфер соседа, и
// проверка запрета прошла бы, ничего не прочитав.
prevDefault := slog.Default()
slog.SetDefault(logger)
t.Cleanup(func() { slog.SetDefault(prevDefault) })
trsService := service.NewTranscribeService(
jobRepo,
fileRepo,
@@ -173,27 +133,21 @@ func setupTestEnv(t *testing.T, metaviewer contract.AudioMetaViewer) *testEnv {
&stubConverter{},
&recognizer.MemoryAudioRecognizer{},
&TestTgSender{},
storageDir,
logger,
)
handler := NewTranscribeHandler(jobRepo, trsService)
handler := NewTranscribeHandler(jobRepo, trsService, logger)
// Роутер собирается той же цепочкой, что и боевой (main.go): у приёма три
// пишущих в журнал потока, и middleware — третий. Без него требование «ни
// одна журнальная запись приёма» проверялось бы шире, чем оракул смотрит.
router := gin.New()
router.Use(sloggin.New(logger))
router.Use(gin.Recovery())
router.MaxMultipartMemory = 32 << 20 // 32 MiB
// Роутер собирается тем же способом, что и боевой: маршруты вешает сам
// обработчик, и проверка судит ту же цепочку, что и прод.
r, err := apis.NewRouter(app)
require.NoError(t, err)
handler.Register(r)
api := router.Group("/api")
{
api.POST("/audio", handler.CreateTranscribeJob)
api.GET("/status/:id", handler.GetTranscribeJobStatus)
}
mux, err := r.BuildMux()
require.NoError(t, err)
return &testEnv{router: router, handler: handler, db: db, storageDir: storageDir, journal: journal}
return &testEnv{mux: mux, handler: handler, app: app, journal: journal}
}
// createMultipartRequest собирает запрос из имени и содержимого. Файла на диске
@@ -217,35 +171,67 @@ func createMultipartRequestWithField(t *testing.T, field, fileName string, conte
err = writer.Close()
require.NoError(t, err)
req, err := http.NewRequest("POST", "/api/audio", &buf)
require.NoError(t, err)
req := httptest.NewRequest("POST", "/api/audio", &buf)
req.Header.Set("Content-Type", writer.FormDataContentType())
return req
}
// storedFiles отдаёт содержимое каталога хранения.
func storedFiles(t *testing.T, env *testEnv) []string {
files, err := filepath.Glob(filepath.Join(env.storageDir, "*"))
// storedFileNames отдаёт имена, под которыми файлы легли в хранилище.
func storedFileNames(t *testing.T, env *testEnv) []string {
records, err := env.app.FindAllRecords(pbrepo.FilesCollection)
require.NoError(t, err)
return files
var names []string
for _, record := range records {
names = append(names, record.GetStringSlice("file")...)
}
return names
}
// countFiles считает записи о файлах.
func countFiles(t *testing.T, env *testEnv) int {
records, err := env.app.FindAllRecords(pbrepo.FilesCollection)
require.NoError(t, err)
return len(records)
}
// countJobs считает заведённые задачи расшифровки.
func countJobs(t *testing.T, env *testEnv) int {
var count int
err := env.db.QueryRow("SELECT COUNT(*) FROM transcribe_jobs").Scan(&count)
records, err := env.app.FindAllRecords(pbrepo.JobsCollection)
require.NoError(t, err)
return count
return len(records)
}
// storedFileName отдаёт имя файла, записанное в учёте под данным идентификатором.
func storedFileName(t *testing.T, env *testEnv, fileID string) string {
var name string
err := env.db.QueryRow("SELECT file_name FROM files WHERE id = ?", fileID).Scan(&name)
// jobWithFile заводит задачу вместе с её записью: ссылка на файл обязательна
// схемой, потому что без неё задача не пройдёт ни одного шага.
func jobWithFile(t *testing.T, env *testEnv) *entity.TranscribeJob {
t.Helper()
repo := pbrepo.NewFileRepository(env.app)
work, err := repo.Stage(".mp3", strings.NewReader("запись"))
require.NoError(t, err)
return name
defer func() { require.NoError(t, work.Close()) }()
file, err := repo.CreateLocal("sample.mp3", work)
require.NoError(t, err)
job := &entity.TranscribeJob{State: entity.StateCreated, Source: entity.SourceApi, FileID: &file.Id}
require.NoError(t, env.handler.jobRepo.Create(job))
return job
}
// storedContent читает содержимое файла из хранилища.
func storedContent(t *testing.T, env *testEnv, fileID string) []byte {
repo := pbrepo.NewFileRepository(env.app)
reader, err := repo.Open(fileID)
require.NoError(t, err)
defer reader.Close()
var buf bytes.Buffer
_, err = buf.ReadFrom(reader)
require.NoError(t, err)
return buf.Bytes()
}
func TestCreateTranscribeJob_Success(t *testing.T) {
@@ -255,7 +241,7 @@ func TestCreateTranscribeJob_Success(t *testing.T) {
req := createMultipartRequest(t, "sample.m4a", content)
w := httptest.NewRecorder()
env.router.ServeHTTP(w, req)
env.mux.ServeHTTP(w, req)
require.Equal(t, http.StatusCreated, w.Code)
@@ -285,17 +271,9 @@ func TestCreateTranscribeJob_Success(t *testing.T) {
require.NotNil(t, job.FileID)
assert.NotEmpty(t, *job.FileID)
// Содержимое лежит в каталоге хранения одним файлом и целиком.
files := storedFiles(t, env)
require.Len(t, files, 1)
stored, err := os.ReadFile(files[0])
require.NoError(t, err)
assert.Equal(t, content, stored)
// Учёт указывает на этот самый файл, а не на какой-то другой: дальше по
// конвейеру путь берётся только из учёта, и разъезд убил бы задачу молча.
assert.Equal(t, filepath.Base(files[0]), storedFileName(t, env, *job.FileID))
// Содержимое лежит в хранилище одним файлом и целиком.
require.Equal(t, 1, countFiles(t, env))
assert.Equal(t, content, storedContent(t, env, *job.FileID))
}
func TestCreateTranscribeJob_NoFile(t *testing.T) {
@@ -308,9 +286,9 @@ func TestCreateTranscribeJob_NoFile(t *testing.T) {
{
name: "no body at all",
req: func(t *testing.T) *http.Request {
req, err := http.NewRequest("POST", "/api/audio", nil)
require.NoError(t, err)
return req
// Запрос строится так, как его видит сервер: у пришедшего по
// проводу тело не бывает пустым указателем.
return httptest.NewRequest("POST", "/api/audio", http.NoBody)
},
},
{
@@ -326,7 +304,7 @@ func TestCreateTranscribeJob_NoFile(t *testing.T) {
env := setupTestEnv(t, readableMetaViewer())
w := httptest.NewRecorder()
env.router.ServeHTTP(w, tc.req(t))
env.mux.ServeHTTP(w, tc.req(t))
require.Equal(t, http.StatusBadRequest, w.Code)
@@ -335,7 +313,7 @@ func TestCreateTranscribeJob_NoFile(t *testing.T) {
require.NoError(t, err)
assert.Equal(t, "No audio file provided", response["error"])
assert.Empty(t, storedFiles(t, env))
assert.Equal(t, 0, countFiles(t, env))
assert.Equal(t, 0, countJobs(t, env))
})
}
@@ -349,7 +327,7 @@ func TestCreateTranscribeJob_EmptyFile(t *testing.T) {
req := createMultipartRequest(t, "empty.m4a", nil)
w := httptest.NewRecorder()
env.router.ServeHTTP(w, req)
env.mux.ServeHTTP(w, req)
require.Equal(t, http.StatusCreated, w.Code)
@@ -396,28 +374,52 @@ func TestCreateTranscribeJob_DifferentFileExtensions(t *testing.T) {
req := createMultipartRequest(t, tc.fileName, []byte("запись"))
w := httptest.NewRecorder()
env.router.ServeHTTP(w, req)
env.mux.ServeHTTP(w, req)
require.Equal(t, http.StatusCreated, w.Code)
files := storedFiles(t, env)
require.Len(t, files, 1)
names := storedFileNames(t, env)
require.Len(t, names, 1)
// Имя отправителя в хранилище не попадает: имя файла — свой
// идентификатор, от отправителя взято только расширение.
assert.Equal(t, tc.expectExt, filepath.Ext(files[0]))
assert.NotContains(t, filepath.Base(files[0]), tc.fileName)
// идентификатор, от отправителя взято только расширение. Суффикс
// дописывает само хранилище, поэтому сверяем хвост, а не Ext.
assert.True(t, strings.HasSuffix(names[0], tc.expectExt),
"имя в хранилище %q оканчивается на %q", names[0], tc.expectExt)
assert.NotContains(t, names[0], strings.TrimSuffix(tc.fileName, tc.expectExt))
})
}
}
// Имя, данное отправителем, в хранилище не попадает целиком — умолчание
// библиотеки, строящее имя файла из него, не применяется. Проверка отдельная от
// перебора расширений: там сверяется хвост, здесь — что основы имени нет.
func TestCreateTranscribeJob_SenderFileNameNotStored(t *testing.T) {
env := setupTestEnv(t, readableMetaViewer())
req := createMultipartRequest(t, "секретное-слово.mp3", []byte("запись"))
w := httptest.NewRecorder()
env.mux.ServeHTTP(w, req)
require.Equal(t, http.StatusCreated, w.Code)
names := storedFileNames(t, env)
require.Len(t, names, 1)
assert.NotContains(t, names[0], "секретное-слово",
"имя, данное отправителем, в хранилище не попадает")
assert.True(t, strings.HasSuffix(names[0], ".mp3"),
"расширение при этом сохраняется: %q", names[0])
}
func TestCreateTranscribeJob_MetaViewerFailure(t *testing.T) {
env := setupTestEnv(t, &stubMetaViewer{err: errors.New("не удалось прочитать запись")})
req := createMultipartRequest(t, "broken.m4a", []byte("не запись вовсе"))
w := httptest.NewRecorder()
env.router.ServeHTTP(w, req)
env.mux.ServeHTTP(w, req)
require.Equal(t, http.StatusInternalServerError, w.Code)
@@ -439,14 +441,14 @@ func TestCreateTranscribeJob_MetaViewerFailure(t *testing.T) {
// пропустил бы.
const senderNameMarker = "SENDERNAMELEAKMARKER7Q2"
// Тексты, по которым проверки находят журнальные строки. Оба — записанный долг
// `docs/conventions/logging.md`: `msg` обязан стать короткой категорией, а
// транспорту не положено логировать вовсе. Когда долг закроют, правка будет
// здесь и одна, а смысл утверждений менять не придётся.
// Тексты, по которым проверки находят журнальные строки. Первый — записанный
// долг `docs/conventions/logging.md`: `msg` обязан стать короткой категорией.
// Когда долг закроют, правка будет здесь и одна.
const (
msgIntake = "Creating transcribe job"
msgTransportErr = "Err:"
msgMiddleware = "Incoming request"
msgIntake = "Creating transcribe job"
// Отказ пишет доменная граница — приём, — а не транспорт: конвенция просит
// логировать ошибку один раз, и повторная запись транспорта снята.
msgIntakeErr = "Failed to get file info"
)
func TestCreateTranscribeJob_SenderFileNameNotLogged(t *testing.T) {
@@ -457,17 +459,16 @@ func TestCreateTranscribeJob_SenderFileNameNotLogged(t *testing.T) {
req := createMultipartRequest(t, senderNameMarker+".mp3", []byte("запись"))
w := httptest.NewRecorder()
env.router.ServeHTTP(w, req)
env.mux.ServeHTTP(w, req)
require.Equal(t, http.StatusCreated, w.Code)
journal := env.journal.String()
// Сперва — что поток middleware вообще перехвачен. Он третий писатель
// журнала приёма, и без этого утверждения снятие его из тестового роутера
// сузило бы оракул молча.
require.Contains(t, journal, msgMiddleware,
"строка middleware о запросе попадает в перехваченный журнал")
// Сперва — что журнал приёма вообще перехвачен. Без этого утверждения
// пустой буфер сделал бы проверку запрета зелёной, ничего не прочитав.
require.Contains(t, journal, msgIntake,
"строка приёма попадает в перехваченный журнал")
assert.NotContains(t, journal, senderNameMarker,
"имя, данное отправителем, не пишется в журнал: инвариант приватности")
@@ -481,22 +482,45 @@ func TestCreateTranscribeJob_SenderFileNameNotLoggedOnFailure(t *testing.T) {
req := createMultipartRequest(t, senderNameMarker+".mp3", []byte("не запись вовсе"))
w := httptest.NewRecorder()
env.router.ServeHTTP(w, req)
env.mux.ServeHTTP(w, req)
require.Equal(t, http.StatusInternalServerError, w.Code)
journal := env.journal.String()
// Сперва — что второй поток журнала вообще перехвачен. Без этого
// утверждения снятие `slog.SetDefault` из окружения оставило бы проверку
// зелёной, а оракул критического инварианта молча сузился бы вдвое.
require.Contains(t, journal, msgTransportErr,
"строка транспорта, идущая мимо slog, попадает в перехваченный журнал")
// Сперва — что журнал ветки отказа вообще перехвачен: обработчик пишет свою
// строку, и без неё оракул молча сузился бы вдвое.
require.Contains(t, journal, msgIntakeErr,
"строка приёма об отказе попадает в перехваченный журнал")
assert.NotContains(t, journal, senderNameMarker,
"имя отправителя не пишется в журнал и на пути отказа")
}
// Имя, под которым файл лёг в хранилище, — это последняя часть ссылки
// `/api/files/...`, по которой запись скачивают. Попав в журнал, строка стала бы
// бессрочным ключом к чужому аудио, поэтому в журнал идёт имя, заданное
// сервисом, а суффикс хранилища — нет.
func TestCreateTranscribeJob_StorageFileNameNotLogged(t *testing.T) {
env := setupTestEnv(t, readableMetaViewer())
req := createMultipartRequest(t, "sample.mp3", []byte("запись"))
w := httptest.NewRecorder()
env.mux.ServeHTTP(w, req)
require.Equal(t, http.StatusCreated, w.Code)
names := storedFileNames(t, env)
require.Len(t, names, 1)
journal := env.journal.String()
require.Contains(t, journal, msgIntake, "журнал приёма перехвачен")
assert.NotContains(t, journal, names[0],
"имени файла в хранилище в журнале нет: по нему собирается ссылка на скачивание")
}
func TestCreateTranscribeJob_JournalTracesRecord(t *testing.T) {
env := setupTestEnv(t, readableMetaViewer())
@@ -504,7 +528,7 @@ func TestCreateTranscribeJob_JournalTracesRecord(t *testing.T) {
req := createMultipartRequest(t, "sample.mp3", content)
w := httptest.NewRecorder()
env.router.ServeHTTP(w, req)
env.mux.ServeHTTP(w, req)
require.Equal(t, http.StatusCreated, w.Code)
@@ -569,7 +593,7 @@ func TestCreateTranscribeJob_MetricLabelCarriesNoSenderName(t *testing.T) {
req := createMultipartRequest(t, "sample."+senderNameMarker, []byte("запись"))
w := httptest.NewRecorder()
env.router.ServeHTTP(w, req)
env.mux.ServeHTTP(w, req)
require.Equal(t, http.StatusCreated, w.Code)
@@ -582,41 +606,30 @@ func TestCreateTranscribeJob_MetricLabelCarriesNoSenderName(t *testing.T) {
assert.Contains(t, values, "other",
"незнакомое расширение приведено к общему значению")
// А на диске расширение остаётся пришедшим: раскладка каталога записей
// объявлена необратимой, и приведение сюда не распространяется.
files := storedFiles(t, env)
require.Len(t, files, 1)
assert.Equal(t, "."+senderNameMarker, filepath.Ext(files[0]))
// А в хранилище расширение остаётся пришедшим: приведение сюда не
// распространяется.
names := storedFileNames(t, env)
require.Len(t, names, 1)
assert.True(t, strings.HasSuffix(names[0], "."+senderNameMarker),
"имя в хранилище сохраняет пришедшее расширение: %q", names[0])
}
func TestGetTranscribeJobStatus_Success(t *testing.T) {
env := setupTestEnv(t, readableMetaViewer())
job := &entity.TranscribeJob{
Id: "test-job-id",
State: entity.StateCreated,
Source: entity.SourceApi,
FileID: nil,
IsError: false,
CreatedAt: time.Now(),
}
job := jobWithFile(t, env)
err := env.handler.jobRepo.Create(job)
require.NoError(t, err)
req, err := http.NewRequest("GET", "/api/status/test-job-id", nil)
require.NoError(t, err)
req := httptest.NewRequest("GET", "/api/status/"+job.Id, http.NoBody)
w := httptest.NewRecorder()
env.router.ServeHTTP(w, req)
env.mux.ServeHTTP(w, req)
require.Equal(t, http.StatusOK, w.Code)
var response GetTranscribeJobResponse
err = json.Unmarshal(w.Body.Bytes(), &response)
require.NoError(t, err)
require.NoError(t, json.Unmarshal(w.Body.Bytes(), &response))
assert.Equal(t, "test-job-id", response.JobID)
assert.Equal(t, job.Id, response.JobID)
assert.Equal(t, entity.StateCreated, response.State)
assert.NotZero(t, response.CreatedAt)
}
@@ -624,21 +637,12 @@ func TestGetTranscribeJobStatus_Success(t *testing.T) {
func TestGetTranscribeJobStatus_NoTranscriptionText(t *testing.T) {
env := setupTestEnv(t, readableMetaViewer())
job := &entity.TranscribeJob{
Id: "job-without-text",
State: entity.StateCreated,
Source: entity.SourceApi,
CreatedAt: time.Now(),
}
job := jobWithFile(t, env)
err := env.handler.jobRepo.Create(job)
require.NoError(t, err)
req, err := http.NewRequest("GET", "/api/status/job-without-text", nil)
require.NoError(t, err)
req := httptest.NewRequest("GET", "/api/status/"+job.Id, http.NoBody)
w := httptest.NewRecorder()
env.router.ServeHTTP(w, req)
env.mux.ServeHTTP(w, req)
require.Equal(t, http.StatusOK, w.Code)
@@ -646,8 +650,7 @@ func TestGetTranscribeJobStatus_NoTranscriptionText(t *testing.T) {
// читается клиентом как «расшифровка пуста», и разобранная структура
// эти два случая не различает.
var raw map[string]json.RawMessage
err = json.Unmarshal(w.Body.Bytes(), &raw)
require.NoError(t, err)
require.NoError(t, json.Unmarshal(w.Body.Bytes(), &raw))
assert.Contains(t, raw, "job_id")
assert.Contains(t, raw, "status")
@@ -658,17 +661,15 @@ func TestGetTranscribeJobStatus_NoTranscriptionText(t *testing.T) {
func TestGetTranscribeJobStatus_NotFound(t *testing.T) {
env := setupTestEnv(t, readableMetaViewer())
req, err := http.NewRequest("GET", "/api/status/non-existent-id", nil)
require.NoError(t, err)
req := httptest.NewRequest("GET", "/api/status/non-existent-id", http.NoBody)
w := httptest.NewRecorder()
env.router.ServeHTTP(w, req)
env.mux.ServeHTTP(w, req)
require.Equal(t, http.StatusNotFound, w.Code)
var response map[string]string
err = json.Unmarshal(w.Body.Bytes(), &response)
require.NoError(t, err)
require.NoError(t, json.Unmarshal(w.Body.Bytes(), &response))
assert.Equal(t, "Job not found", response["error"])
}
+22 -14
View File
@@ -4,25 +4,33 @@ import (
"time"
)
// Где лежит копия файла. Поле названо `location`, а не `storage`: последним
// словом зовут само хранилище, и третий смысл у одного слова развёл бы по
// разным вещам запись о файле и хранилище, в котором она лежит.
const (
StorageLocal = "local"
StorageS3 = "s3"
LocationLocal = "local"
LocationS3 = "s3"
)
// MaxRecordSize — потолок размера одного файла записи. Выведен из расчётного
// потолка записи в шесть часов с запасом на видео, а не из замера.
//
// Число нужно назвать **явно** в двух местах сразу: у поля файла в хранилище
// нулевой потолок значит не «без предела», а умолчание библиотеки в 5 МиБ, а у
// тела запроса приёма умолчание роутера отсекало бы запись раньше, чем она
// дойдёт до обработчика — без строки в журнале приёма.
const MaxRecordSize int64 = 8 << 30 // 8 ГиБ
// File — одна физическая копия: исходник, результат конвертации и копия во
// внешнем хранилище — три разные записи.
type File struct {
Id string
Storage string
Id string
Location string
// FileName — имя, под которым файл лежит: у местной копии это имя, заданное
// сервисом, у внешней — ключ объекта. Своего суффикса хранилище к заданному
// имени не дописывает: суффикс появляется только у имён, которые оно строит
// само из имени отправителя, а это умолчание не применяется.
FileName string
Size int64
CreatedAt time.Time
}
func (f *File) CopyWithStorage(newId, storage string) *File {
return &File{
Id: newId,
Storage: storage,
FileName: f.FileName,
Size: f.Size,
CreatedAt: time.Now(),
}
}
+30 -2
View File
@@ -9,11 +9,11 @@ type TranscribeJob struct {
State string
Source string
FileID *string
IsError bool
ErrorText *string
AcquisitionID *string
AcquireTime *time.Time
DelayTime *time.Time
Attempts int // Число попыток: растёт при захвате, обнуляется на шаге без отказа
RecognitionOpID *string // ID операции распознавания в Yandex Cloud
TranscriptionText *string // Результат распознавания
TgChatId *int64 // Telegram: в какой чат отправить результат распознавания
@@ -28,6 +28,11 @@ const (
StateTranscribe = "transcribe"
StateDone = "done"
StateFailed = "failed"
// StateDead — задача, которую мы повторяли и перестали. От `failed` она
// отличается тем, чей это приговор: в `failed` задачу переводит шаг,
// рассудивший об этой записи окончательно, а сюда она уходит без такого
// суждения. Ни один шаг конвейера в неё не переводит сам.
StateDead = "dead"
)
const (
@@ -43,6 +48,10 @@ func (j *TranscribeJob) MoveToState(state string) {
j.DelayTime = nil
j.AcquisitionID = nil
j.AcquireTime = nil
// Шаг, дошедший до перехода, завершился без отказа, а попытки считают
// именно отказавшие: иначе задача, прошедшая конвейер целиком, накопила бы
// их поштучно и умерла бы здоровой.
j.Attempts = 0
j.UpdatedAt = time.Now()
}
@@ -59,6 +68,25 @@ func (j *TranscribeJob) Done(transcriptionText string) {
func (j *TranscribeJob) Fail(errText string) {
j.MoveToState(StateFailed)
j.IsError = true
j.ErrorText = &errText
}
// RetryAfter освобождает отказавшую задачу для повтора: захват снимается,
// пауза ставится, а число попыток сохраняется — по нему растёт пауза и
// наступает предел.
func (j *TranscribeJob) RetryAfter(delay time.Time) {
j.AcquisitionID = nil
j.AcquireTime = nil
j.DelayTime = &delay
j.UpdatedAt = time.Now()
}
// Die переводит задачу, исчерпавшую попытки, в состояние «мертва». Число
// попыток при этом сохраняется: по нему видно, сколько раз мы пробовали, а
// возвращает задачу в работу владелец правкой состояния.
func (j *TranscribeJob) Die(errText string) {
attempts := j.Attempts
j.MoveToState(StateDead)
j.Attempts = attempts
j.ErrorText = &errText
}
+5 -5
View File
@@ -24,8 +24,8 @@ type stubJobRepo struct {
err error
}
func (r *stubJobRepo) Create(*entity.TranscribeJob) error { return nil }
func (r *stubJobRepo) Save(*entity.TranscribeJob) error { return nil }
func (r *stubJobRepo) Create(*entity.TranscribeJob) error { return nil }
func (r *stubJobRepo) Save(*entity.TranscribeJob, string) error { return nil }
func (r *stubJobRepo) GetByID(string) (*entity.TranscribeJob, error) {
return nil, errors.New("не зовётся этими проверками")
@@ -37,7 +37,7 @@ func (r *stubJobRepo) FindAndAcquire(string, string, time.Time) (*entity.Transcr
func serviceWithRepo(repo contract.TranscriptJobRepository) *TranscribeService {
logger := slog.New(slog.NewTextHandler(io.Discard, nil))
return NewTranscribeService(repo, nil, nil, nil, nil, nil, "", logger)
return NewTranscribeService(repo, nil, nil, nil, nil, nil, logger)
}
// Репозиторий вправе добавить своему отказу пояснение — соседние ветки того же
@@ -50,7 +50,7 @@ func TestFindJobTranslatesWrappedNotFoundToNoop(t *testing.T) {
&contract.JobNotFoundError{State: "created", Message: "appropriate job not found"}),
})
_, err := svc.findJob("created", time.Minute)
_, _, err := svc.findJob("created", time.Minute)
var noop *contract.NoopJobError
if !errors.As(err, &noop) {
@@ -66,7 +66,7 @@ func TestFindJobTranslatesWrappedNotFoundToNoop(t *testing.T) {
func TestFindJobKeepsRealFailure(t *testing.T) {
svc := serviceWithRepo(&stubJobRepo{err: errors.New("database is gone")})
_, err := svc.findJob("created", time.Minute)
_, _, err := svc.findJob("created", time.Minute)
var noop *contract.NoopJobError
if errors.As(err, &noop) {
+362
View File
@@ -0,0 +1,362 @@
package service
import (
"errors"
"io"
"log/slog"
"os"
"path/filepath"
"strings"
"testing"
"time"
"github.com/pocketbase/pocketbase/core"
"github.com/pocketbase/pocketbase/tools/types"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"git.vakhrushev.me/av/transcriber/internal/adapter/recognizer"
pbrepo "git.vakhrushev.me/av/transcriber/internal/adapter/repo/pocketbase"
"git.vakhrushev.me/av/transcriber/internal/contract"
"git.vakhrushev.me/av/transcriber/internal/entity"
)
// Проверки конвейера идут против настоящего хранилища: захват, число попыток и
// переход в «мертва» держатся на запросе, и подставной репозиторий проверял бы
// собственную заглушку, а не то, что делает база.
// failingConverter отказывает на каждой попытке.
type failingConverter struct{}
func (c *failingConverter) Convert(string, string) error {
return errors.New("конвертация не удалась")
}
type okMetaViewer struct{}
func (m *okMetaViewer) GetInfo(string) (*contract.AudioInfo, error) {
return &contract.AudioInfo{Seconds: 1}, nil
}
type failingMetaViewer struct{}
func (m *failingMetaViewer) GetInfo(string) (*contract.AudioInfo, error) {
return nil, errors.New("запись не читается")
}
// recordingSender запоминает, что и куда отправлено.
type recordingSender struct {
messages []string
}
func (s *recordingSender) Send(text string, chatId int64, replyMsgId *int) error {
s.messages = append(s.messages, text)
return nil
}
type pipelineEnv struct {
app core.App
service *TranscribeService
jobRepo *pbrepo.TranscriptJobRepository
fileRepo *pbrepo.FileRepository
sender *recordingSender
}
func newPipelineEnv(t *testing.T, metaviewer contract.AudioMetaViewer, converter contract.AudioFileConverter) *pipelineEnv {
t.Helper()
app, err := pbrepo.New(t.TempDir())
require.NoError(t, err)
t.Cleanup(func() {
if err := app.ResetBootstrapState(); err != nil {
t.Logf("не удалось закрыть хранилище: %v", err)
}
})
// Правила панели вешаются и здесь: конфигурация под проверкой обязана
// совпадать с боевой, иначе утверждения говорят про прод то, чего в проде
// нет.
pbrepo.BindPanelRules(app)
jobRepo := pbrepo.NewTranscriptJobRepository(app)
fileRepo := pbrepo.NewFileRepository(app)
sender := &recordingSender{}
svc := NewTranscribeService(
jobRepo,
fileRepo,
metaviewer,
converter,
&recognizer.MemoryAudioRecognizer{},
sender,
slog.New(slog.NewTextHandler(io.Discard, nil)),
)
return &pipelineEnv{app: app, service: svc, jobRepo: jobRepo, fileRepo: fileRepo, sender: sender}
}
// newTelegramJob заводит задачу с записью — так, как её завёл бы приём.
func newTelegramJob(t *testing.T, env *pipelineEnv) *entity.TranscribeJob {
t.Helper()
chatId := int64(100)
job, err := env.service.CreateJobFromTelegram(strings.NewReader("запись"), "voice.ogg", chatId, 1)
require.NoError(t, err)
return job
}
// clearDelay снимает паузу, чтобы следующий прогон взял задачу сразу: проверка
// судит счётчик попыток, а не то, умеет ли она ждать.
func clearDelay(t *testing.T, env *pipelineEnv, jobID string) {
t.Helper()
record, err := env.app.FindRecordById(pbrepo.JobsCollection, jobID)
require.NoError(t, err)
record.Set("delay_time", "")
require.NoError(t, env.app.Save(record))
}
// rotAcquisition отодвигает время захвата так, чтобы он протух: так это
// выглядит, когда шаг оборвался вместе с процессом.
func rotAcquisition(t *testing.T, env *pipelineEnv, jobID string) {
t.Helper()
record, err := env.app.FindRecordById(pbrepo.JobsCollection, jobID)
require.NoError(t, err)
record.Set("acquire_time", types.NowDateTime().Add(-24*time.Hour))
require.NoError(t, env.app.Save(record))
}
// Задача, падающая на каждой попытке, уходит в «мертва»: из выборки исчезает,
// видна отбором по состоянию, а отправитель узнаёт о неудаче. Инвариант
// «Принятая запись не теряется молча» допускает два исхода, и молчаливая смерть
// не подходит ни под один.
func TestJobDiesAfterAttemptLimit(t *testing.T) {
env := newPipelineEnv(t, &okMetaViewer{}, &failingConverter{})
job := newTelegramJob(t, env)
// Отказ конвертации переводит задачу в `failed` сразу, поэтому предел
// попыток проверяем на шаге, который отказывает *не* приговором: подменяем
// его отказом источника метаданных внутри самого шага конвертации нельзя, и
// вместо этого гоняем захват без выполнения шага — так же, как это выглядит
// при гибели процесса.
for i := 0; i < maxAttempts; i++ {
_, err := env.jobRepo.FindAndAcquire(entity.StateCreated, "holder", time.Now().Add(time.Hour))
require.NoError(t, err)
}
rotAcquisition(t, env, job.Id)
// Следующий захват видит перебор и хоронит задачу.
err := env.service.FindAndRunConversionJob()
var noop *contract.NoopJobError
require.ErrorAs(t, err, &noop, "мёртвая задача шагу не отдаётся")
after, err := env.jobRepo.GetByID(job.Id)
require.NoError(t, err)
assert.Equal(t, entity.StateDead, after.State, "задача видна отбором по состоянию")
assert.Greater(t, after.Attempts, maxAttempts, "число попыток сохранено")
require.Len(t, env.sender.messages, 1, "отправитель узнал о неудаче")
assert.Contains(t, env.sender.messages[0], "попытки исчерпаны")
// И из выборки она исчезла.
_, err = env.jobRepo.FindAndAcquire(entity.StateCreated, "next", time.Now().Add(time.Hour))
var missing *contract.JobNotFoundError
assert.ErrorAs(t, err, &missing)
}
// Мёртвая задача возвращается в работу правкой состояния.
func TestDeadJobReturnsAfterStateEdit(t *testing.T) {
env := newPipelineEnv(t, &okMetaViewer{}, &failingConverter{})
job := newTelegramJob(t, env)
for i := 0; i < maxAttempts; i++ {
_, err := env.jobRepo.FindAndAcquire(entity.StateCreated, "holder", time.Now().Add(time.Hour))
require.NoError(t, err)
}
require.Error(t, env.service.FindAndRunConversionJob())
record, err := env.app.FindRecordById(pbrepo.JobsCollection, job.Id)
require.NoError(t, err)
record.Set("state", entity.StateCreated)
require.NoError(t, env.app.Save(record))
again, err := env.jobRepo.FindAndAcquire(entity.StateCreated, "next", time.Now().Add(time.Hour))
require.NoError(t, err, "снятое состояние возвращает задачу в работу")
assert.Equal(t, job.Id, again.Id)
}
// Отказ шага не оставляет задачу захваченной до конца срока: захват снимается,
// и задача ждёт нарастающую паузу. Иначе повтор наступал бы через восемь часов.
func TestFailedStepSchedulesRetryWithGrowingDelay(t *testing.T) {
// Источник метаданных отказывает — это отказ шага, а не приговор записи.
env := newPipelineEnv(t, &okMetaViewer{}, &failingConverter{})
job := newTelegramJob(t, env)
// Ссылку переставляем на запись без содержимого: шаг отказывает на получении
// рабочей копии — то есть отказом, а не приговором записи.
empty, err := env.fileRepo.CreateRemote("object-key", 1)
require.NoError(t, err)
record, err := env.app.FindRecordById(pbrepo.JobsCollection, job.Id)
require.NoError(t, err)
record.Set("file", empty.Id)
require.NoError(t, env.app.Save(record))
// Первый отказ.
require.Error(t, env.service.FindAndRunConversionJob())
after, err := env.jobRepo.GetByID(job.Id)
require.NoError(t, err)
require.Nil(t, after.AcquisitionID, "захват снят: задача пригодна к повтору")
require.NotNil(t, after.DelayTime, "пауза поставлена")
firstDelay := time.Until(*after.DelayTime)
// Второй отказ — с той же задачи, пауза снята вручную.
clearDelay(t, env, job.Id)
require.Error(t, env.service.FindAndRunConversionJob())
after, err = env.jobRepo.GetByID(job.Id)
require.NoError(t, err)
require.NotNil(t, after.DelayTime)
secondDelay := time.Until(*after.DelayTime)
assert.Greater(t, secondDelay, firstDelay, "вторая пауза длиннее первой")
}
// Пауза растёт с числом попыток и упирается в потолок.
func TestRetryDelayGrowsAndCaps(t *testing.T) {
assert.Equal(t, retryDelayBase, retryDelay(1))
assert.Equal(t, 2*retryDelayBase, retryDelay(2))
assert.Greater(t, retryDelay(3), retryDelay(2))
assert.Equal(t, retryDelayCap, retryDelay(100), "пауза упирается в потолок")
assert.Equal(t, retryDelayBase, retryDelay(0), "нулевая попытка не даёт нулевой паузы")
}
// Рабочая копия убирается на любом исходе, включая отказ. Забытая копия — это
// шестичасовая запись во временном каталоге, и узнать о ней неоткуда.
func TestWorkFileRemovedAfterIntakeFailure(t *testing.T) {
tempDir := t.TempDir()
t.Setenv("TMPDIR", tempDir)
env := newPipelineEnv(t, &failingMetaViewer{}, &failingConverter{})
_, err := env.service.CreateJobFromApi(strings.NewReader("запись"), "sample.mp3")
require.Error(t, err, "отказ источника метаданных роняет приём")
leftovers, err := filepath.Glob(filepath.Join(tempDir, "transcriber-*"))
require.NoError(t, err)
assert.Empty(t, leftovers, "рабочей копии после отказа не остаётся")
}
// Успешный приём тоже за собой убирает: копия нужна была только на время
// укладки в хранилище.
func TestWorkFileRemovedAfterSuccessfulIntake(t *testing.T) {
tempDir := t.TempDir()
t.Setenv("TMPDIR", tempDir)
env := newPipelineEnv(t, &okMetaViewer{}, &failingConverter{})
_, err := env.service.CreateJobFromApi(strings.NewReader("запись"), "sample.mp3")
require.NoError(t, err)
leftovers, err := filepath.Glob(filepath.Join(tempDir, "transcriber-*"))
require.NoError(t, err)
assert.Empty(t, leftovers, "рабочей копии после успеха не остаётся")
}
// Задача не остаётся ссылающейся на файл, которого нет: ссылка переставляется
// только после того, как запись о новом файле существует.
func TestJobNeverPointsToMissingFile(t *testing.T) {
env := newPipelineEnv(t, &okMetaViewer{}, &failingConverter{})
job := newTelegramJob(t, env)
// Конвертация отказывает — задача уходит в `failed`, но ссылка остаётся на
// исходную запись, а не на несозданный результат.
require.NoError(t, env.service.FindAndRunConversionJob())
after, err := env.jobRepo.GetByID(job.Id)
require.NoError(t, err)
assert.Equal(t, entity.StateFailed, after.State)
require.NotNil(t, after.FileID)
file, err := env.fileRepo.GetByID(*after.FileID)
require.NoError(t, err, "ссылка задачи ведёт на существующую запись о файле")
assert.NotEmpty(t, file.FileName)
}
// Содержимое доезжает до хранилища целиком и читается обратно тем же.
func TestStoredContentSurvivesRoundTrip(t *testing.T) {
env := newPipelineEnv(t, &okMetaViewer{}, &failingConverter{})
content := strings.Repeat("запись ", 1000)
job, err := env.service.CreateJobFromApi(strings.NewReader(content), "sample.mp3")
require.NoError(t, err)
require.NotNil(t, job.FileID)
reader, err := env.fileRepo.Open(*job.FileID)
require.NoError(t, err)
defer reader.Close()
stored, err := io.ReadAll(reader)
require.NoError(t, err)
assert.Equal(t, content, string(stored))
// И длина в учёте совпадает с длиной принятого.
file, err := env.fileRepo.GetByID(*job.FileID)
require.NoError(t, err)
assert.Equal(t, int64(len(content)), file.Size)
}
// Рабочая копия хранимого файла отдаётся именем на диске — так её получают
// шаги, отдающие файл внешней программе.
func TestLocalizeGivesReadableCopy(t *testing.T) {
env := newPipelineEnv(t, &okMetaViewer{}, &failingConverter{})
job, err := env.service.CreateJobFromApi(strings.NewReader("содержимое"), "sample.mp3")
require.NoError(t, err)
require.NotNil(t, job.FileID)
work, err := env.fileRepo.Localize(*job.FileID)
require.NoError(t, err)
content, err := os.ReadFile(work.Path())
require.NoError(t, err)
assert.Equal(t, "содержимое", string(content))
require.NoError(t, work.Close())
_, err = os.Stat(work.Path())
assert.True(t, os.IsNotExist(err), "закрытая копия убрана")
}
// Уборка рабочей копии проверяется и на шаге конвертации: репозиторий даёт
// единственный способ убрать копию, но зовёт его шаг, и норма держится
// проверкой, а не построением. Копий здесь две — исходник и результат.
func TestWorkFilesRemovedAfterConversionFailure(t *testing.T) {
tempDir := t.TempDir()
t.Setenv("TMPDIR", tempDir)
env := newPipelineEnv(t, &okMetaViewer{}, &failingConverter{})
newTelegramJob(t, env)
// Приём уже отработал — убеждаемся, что за собой он прибрал, иначе остаток
// от него зачёлся бы шагу конвертации.
leftovers, err := filepath.Glob(filepath.Join(tempDir, "transcriber-*"))
require.NoError(t, err)
require.Empty(t, leftovers, "приём убрал свою рабочую копию")
// Конвертация отказывает — задача уходит в `failed`, копии убраны.
require.NoError(t, env.service.FindAndRunConversionJob())
leftovers, err = filepath.Glob(filepath.Join(tempDir, "transcriber-*"))
require.NoError(t, err)
assert.Empty(t, leftovers, "ни исходной копии, ни копии под результат не осталось")
}
+267
View File
@@ -0,0 +1,267 @@
package service
import (
"errors"
"io"
"strings"
"testing"
"time"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"git.vakhrushev.me/av/transcriber/internal/adapter/repo/pocketbase"
"git.vakhrushev.me/av/transcriber/internal/contract"
"git.vakhrushev.me/av/transcriber/internal/entity"
)
// Шаги распознавания переписаны переездом на новое хранилище целиком: они берут
// содержимое по записи, заводят запись о копии во внешнем хранилище и пишут
// результат условием по держателю захвата. Подставной распознаватель проекта
// умеет только «завершено с фиксированным текстом», поэтому ветки ожидания,
// отказа операции и пустого текста изобразить нечем — для них нужен управляемый
// двойник.
// scriptedRecognizer отдаёт заданный исход проверки операции и заданный текст.
type scriptedRecognizer struct {
result *entity.RecognitionResult
text string
recognizeErr error
recognizeCalls int
lastObjectKey string
}
func (r *scriptedRecognizer) Recognize(file io.Reader, fileName string) (string, error) {
r.recognizeCalls++
r.lastObjectKey = fileName
if r.recognizeErr != nil {
return "", r.recognizeErr
}
// Содержимое обязано быть читаемым: шаг отдаёт его наружу потоком.
if _, err := io.Copy(io.Discard, file); err != nil {
return "", err
}
return "operation-id", nil
}
func (r *scriptedRecognizer) GetRecognitionText(string) (string, error) {
return r.text, nil
}
func (r *scriptedRecognizer) CheckRecognitionStatus(string) (*entity.RecognitionResult, error) {
return r.result, nil
}
// convertedJob доводит задачу до состояния, с которого работает шаг
// распознавания: запись принята и сконвертирована.
func convertedJob(t *testing.T, env *pipelineEnv) *entity.TranscribeJob {
t.Helper()
job := newTelegramJob(t, env)
acquired, err := env.jobRepo.FindAndAcquire(entity.StateCreated, "setup", time.Now().Add(-time.Hour))
require.NoError(t, err)
acquired.MoveToState(entity.StateConverted)
require.NoError(t, env.jobRepo.Save(acquired, "setup"))
return job
}
// withRecognizer пересобирает сервис с управляемым распознавателем поверх того
// же хранилища.
func withRecognizer(env *pipelineEnv, rec contract.AudioRecognizer) *TranscribeService {
return NewTranscribeService(
env.jobRepo,
env.fileRepo,
&okMetaViewer{},
&failingConverter{},
rec,
env.sender,
env.service.logger,
)
}
// Шаг распознавания отдаёт содержимое наружу, заводит запись о копии во внешнем
// хранилище и переставляет на неё ссылку задачи — только после того, как запись
// о копии существует.
func TestTranscribeJobHandsRecordOverAndMovesOn(t *testing.T) {
env := newPipelineEnv(t, &okMetaViewer{}, &failingConverter{})
job := convertedJob(t, env)
rec := &scriptedRecognizer{result: entity.NewInProgressResult()}
svc := withRecognizer(env, rec)
require.NoError(t, svc.FindAndRunTranscribeJob())
assert.Equal(t, 1, rec.recognizeCalls, "содержимое отдано распознавателю")
assert.NotEmpty(t, rec.lastObjectKey, "ключ объекта назван")
after, err := env.jobRepo.GetByID(job.Id)
require.NoError(t, err)
assert.Equal(t, entity.StateTranscribe, after.State)
require.NotNil(t, after.RecognitionOpID)
assert.Equal(t, "operation-id", *after.RecognitionOpID)
require.NotNil(t, after.DelayTime, "задержка перед первой проверкой поставлена")
// Ссылка задачи ведёт на существующую запись о копии, а не на несозданную.
require.NotNil(t, after.FileID)
copyRecord, err := env.fileRepo.GetByID(*after.FileID)
require.NoError(t, err)
assert.Equal(t, entity.LocationS3, copyRecord.Location)
}
// Отказ распознавателя не двигает задачу: она остаётся пригодной к повтору.
func TestTranscribeJobKeepsJobRetryableOnRecognizerFailure(t *testing.T) {
env := newPipelineEnv(t, &okMetaViewer{}, &failingConverter{})
job := convertedJob(t, env)
rec := &scriptedRecognizer{recognizeErr: errors.New("распознаватель недоступен")}
svc := withRecognizer(env, rec)
require.Error(t, svc.FindAndRunTranscribeJob())
after, err := env.jobRepo.GetByID(job.Id)
require.NoError(t, err)
assert.Equal(t, entity.StateConverted, after.State, "задача осталась на своём шаге")
assert.Nil(t, after.AcquisitionID, "захват снят: задача пригодна к повтору")
assert.NotNil(t, after.DelayTime, "пауза перед повтором поставлена")
assert.Empty(t, env.sender.messages, "отправителю про повторимый отказ не пишут")
}
// transcribingJob доводит задачу до состояния ожидания операции.
func transcribingJob(t *testing.T, env *pipelineEnv, rec contract.AudioRecognizer) *entity.TranscribeJob {
t.Helper()
job := convertedJob(t, env)
require.NoError(t, withRecognizer(env, rec).FindAndRunTranscribeJob())
clearDelay(t, env, job.Id)
return job
}
// Ожидание чужой операции попытку не тратит и опрос не учащает: шаг отработал
// без отказа, и задержка у него своя, числом.
func TestCheckJobWaitsWithoutSpendingAttempts(t *testing.T) {
env := newPipelineEnv(t, &okMetaViewer{}, &failingConverter{})
rec := &scriptedRecognizer{result: entity.NewInProgressResult()}
job := transcribingJob(t, env, rec)
svc := withRecognizer(env, rec)
for i := 0; i < 3; i++ {
require.NoError(t, svc.FindAndRunTranscribeCheckJob())
after, err := env.jobRepo.GetByID(job.Id)
require.NoError(t, err)
assert.Equal(t, entity.StateTranscribe, after.State)
assert.Equal(t, 0, after.Attempts, "ожидание операции попытку не тратит")
require.NotNil(t, after.DelayTime)
assert.InDelta(t, nextCheckDelay.Seconds(), time.Until(*after.DelayTime).Seconds(), 2,
"задержка опроса не выродилась в наименьшую паузу повтора")
clearDelay(t, env, job.Id)
}
}
// Отказ операции распознавания — приговор записи: задача уходит в `failed`, а
// отправитель узнаёт причину человеческим текстом.
func TestCheckJobFailsJobAndTellsSender(t *testing.T) {
env := newPipelineEnv(t, &okMetaViewer{}, &failingConverter{})
rec := &scriptedRecognizer{result: entity.NewInProgressResult()}
job := transcribingJob(t, env, rec)
rec.result = entity.NewFailedResult("операция отклонена")
svc := withRecognizer(env, rec)
require.NoError(t, svc.FindAndRunTranscribeCheckJob())
after, err := env.jobRepo.GetByID(job.Id)
require.NoError(t, err)
assert.Equal(t, entity.StateFailed, after.State)
require.Len(t, env.sender.messages, 1, "отправитель узнал об отказе")
assert.Contains(t, env.sender.messages[0], "сбой при распознавании файла")
assert.NotContains(t, env.sender.messages[0], "операция отклонена",
"машинная причина отправителю не идёт")
}
// Готовая операция завершает задачу и отдаёт текст отправителю ровно один раз.
func TestCheckJobCompletesAndAnswersOnce(t *testing.T) {
env := newPipelineEnv(t, &okMetaViewer{}, &failingConverter{})
rec := &scriptedRecognizer{result: entity.NewInProgressResult()}
job := transcribingJob(t, env, rec)
rec.result = entity.NewCompletedResult()
rec.text = "расшифровка записи"
svc := withRecognizer(env, rec)
require.NoError(t, svc.FindAndRunTranscribeCheckJob())
after, err := env.jobRepo.GetByID(job.Id)
require.NoError(t, err)
assert.Equal(t, entity.StateDone, after.State)
require.NotNil(t, after.TranscriptionText)
assert.Equal(t, "расшифровка записи", *after.TranscriptionText)
require.Len(t, env.sender.messages, 1, "отправитель получил ровно один ответ")
assert.Equal(t, "расшифровка записи", env.sender.messages[0])
// И задача из выборки исчезла: второй ответ отправителю неоткуда взяться.
_, err = env.jobRepo.FindAndAcquire(entity.StateTranscribe, "next", time.Now().Add(-time.Hour))
var missing *contract.JobNotFoundError
assert.ErrorAs(t, err, &missing)
}
// Пустая расшифровка — не отказ: задача завершается, а отправителю уходит
// объяснение вместо пустого сообщения.
func TestCheckJobCompletesEmptyTextWithExplanation(t *testing.T) {
env := newPipelineEnv(t, &okMetaViewer{}, &failingConverter{})
rec := &scriptedRecognizer{result: entity.NewInProgressResult()}
job := transcribingJob(t, env, rec)
rec.result = entity.NewCompletedResult()
rec.text = ""
svc := withRecognizer(env, rec)
require.NoError(t, svc.FindAndRunTranscribeCheckJob())
after, err := env.jobRepo.GetByID(job.Id)
require.NoError(t, err)
assert.Equal(t, entity.StateDone, after.State)
require.Len(t, env.sender.messages, 1)
assert.Contains(t, strings.ToLower(env.sender.messages[0]), "нет текста")
}
// Шаг, потерявший захват за время работы, результата не пишет и отправителю не
// отвечает: иначе два воркера пишут в одну задачу, а отправитель получает два
// ответа на одну запись.
func TestCheckJobWritesNothingWhenAcquisitionLost(t *testing.T) {
env := newPipelineEnv(t, &okMetaViewer{}, &failingConverter{})
rec := &scriptedRecognizer{result: entity.NewInProgressResult()}
job := transcribingJob(t, env, rec)
rec.result = entity.NewCompletedResult()
rec.text = "расшифровка записи"
// Захват задачи достался другому, пока шаг работал.
acquired, err := env.jobRepo.FindAndAcquire(entity.StateTranscribe, "mine", time.Now().Add(-time.Hour))
require.NoError(t, err)
record, err := env.app.FindRecordById(pocketbase.JobsCollection, job.Id)
require.NoError(t, err)
record.Set("acquisition_id", "someone-else")
require.NoError(t, env.app.Save(record))
svc := withRecognizer(env, rec)
err = svc.checkTranscribeJob(acquired, "mine")
var lost *contract.LostAcquisitionError
require.ErrorAs(t, err, &lost)
after, err := env.jobRepo.GetByID(job.Id)
require.NoError(t, err)
assert.Equal(t, entity.StateTranscribe, after.State, "результат не записан")
assert.Empty(t, env.sender.messages, "отправителю ничего не отправлено")
}
+262 -187
View File
@@ -5,7 +5,7 @@ import (
"fmt"
"io"
"log/slog"
"os"
"math"
"path/filepath"
"strings"
"time"
@@ -18,17 +18,39 @@ import (
const (
defaultAudioExt = "audio"
// Предел попыток. Число обратимо и живёт здесь одним местом; счётчик растёт
// при захвате и обнуляется на шаге, завершившемся без отказа.
maxAttempts = 5
// Пауза перед повтором отказавшей задачи растёт с числом попыток до
// потолка. Ожидание чужой операции этой паузой не выражается — у него своя
// задержка числом, и попытки оно не тратит.
retryDelayBase = time.Second
retryDelayCap = 5 * time.Minute
// Сроки захвата. Каждый не меньше того, что его шаг может занять на самом
// длинном допустимом входе: расчётный потолок записи — шесть часов, и
// конвертация такой записи идёт дольше часа по построению.
conversionAcquireTimeout = 8 * time.Hour
transcribeAcquireTimeout = 8 * time.Hour
checkAcquireTimeout = time.Hour
// Задержки опроса операции распознавания. Числа, а не функция числа попыток:
// счётчик на ожидании обнулён, и выведенная из него пауза выродилась бы в
// своё наименьшее значение, учащая опрос платного сервиса.
firstCheckDelay = 10 * time.Second
nextCheckDelay = 5 * time.Second
)
type TranscribeService struct {
jobRepo contract.TranscriptJobRepository
fileRepo contract.FileRepository
metaviewer contract.AudioMetaViewer
converter contract.AudioFileConverter
recognizer contract.AudioRecognizer
tgSender contract.TelegramMessageSender
storagePath string
logger *slog.Logger
jobRepo contract.TranscriptJobRepository
fileRepo contract.FileRepository
metaviewer contract.AudioMetaViewer
converter contract.AudioFileConverter
recognizer contract.AudioRecognizer
tgSender contract.TelegramMessageSender
logger *slog.Logger
}
func NewTranscribeService(
@@ -38,172 +60,172 @@ func NewTranscribeService(
converter contract.AudioFileConverter,
recognizer contract.AudioRecognizer,
tgSender contract.TelegramMessageSender,
storagePath string,
logger *slog.Logger,
) *TranscribeService {
return &TranscribeService{
jobRepo: jobRepo,
fileRepo: fileRepo,
metaviewer: metaviewer,
converter: converter,
recognizer: recognizer,
tgSender: tgSender,
storagePath: storagePath,
logger: logger,
jobRepo: jobRepo,
fileRepo: fileRepo,
metaviewer: metaviewer,
converter: converter,
recognizer: recognizer,
tgSender: tgSender,
logger: logger,
}
}
func (s *TranscribeService) CreateJobFromTelegram(file io.Reader, fileName string, chatId int64, replyMsgId int) (*entity.TranscribeJob, error) {
jobId := uuid.NewString()
now := time.Now()
job := &entity.TranscribeJob{
Id: jobId,
State: entity.StateCreated,
Source: entity.SourceTelegram,
TgChatId: &chatId,
TgReplyMessageId: &replyMsgId,
IsError: false,
CreatedAt: now,
UpdatedAt: now,
}
return s.createTranscribeJob(job, file, fileName)
}
func (s *TranscribeService) CreateJobFromApi(file io.Reader, fileName string) (*entity.TranscribeJob, error) {
jobId := uuid.NewString()
now := time.Now()
job := &entity.TranscribeJob{
Id: jobId,
State: entity.StateCreated,
Source: entity.SourceApi,
IsError: false,
CreatedAt: now,
UpdatedAt: now,
State: entity.StateCreated,
Source: entity.SourceApi,
}
return s.createTranscribeJob(job, file, fileName)
}
func (s *TranscribeService) createTranscribeJob(job *entity.TranscribeJob, file io.Reader, fileName string) (*entity.TranscribeJob, error) {
// Генерируем UUID для файла
fileId := uuid.NewString()
// Определяем расширение файла
ext := filepath.Ext(fileName)
if ext == "" {
ext = fmt.Sprintf(".%s", defaultAudioExt) // fallback если расширение не определено
}
// Создаем путь для сохранения файла
// Собственное имя записи: идентификатор с расширением. Имя, данное
// отправителем, в хранилище не попадает — от него взято только расширение.
fileId := uuid.NewString()
storageFileName := fmt.Sprintf("%s%s", fileId, ext)
storageFilePath := filepath.Join(s.storagePath, storageFileName)
// Имя, данное отправителем, в журнал не идёт: инвариант приватности.
// Расширение из него уже стоит в собственном имени файла на диске.
s.logger.Info("Creating transcribe job",
"file_id", fileId,
"storage_path", storageFilePath)
// Создаем файл на диске
dst, err := os.Create(storageFilePath)
// Содержимое ложится в рабочую копию потоком: в память запись целиком не
// читается, расчётный потолок — шесть часов.
work, err := s.fileRepo.Stage(ext, file)
if err != nil {
s.logger.Error("Failed to create file", "error", err, "path", storageFilePath)
s.logger.Error("Failed to stage uploaded file", "error", err)
return nil, err
}
defer dst.Close()
defer s.closeWork(work)
// Копируем содержимое загруженного файла
size, err := io.Copy(dst, file)
// В журнал идёт расширение, и только оно. Имя, данное отправителем, не
// пишется по инварианту приватности; имя, под которым файл ложится в
// хранилище, — потому что оно последняя часть ссылки на скачивание, и
// строка журнала вместе с идентификатором записи собрала бы её целиком.
s.logger.Info("Creating transcribe job", "file_ext", ext)
info, err := s.metaviewer.GetInfo(work.Path())
if err != nil {
s.logger.Error("Failed to copy file content", "error", err)
s.logger.Error("Failed to get file info", "error", err, "file_ext", ext)
return nil, err
}
if err := dst.Close(); err != nil {
s.logger.Error("Failed to close file", "error", err)
size, err := work.Size()
if err != nil {
s.logger.Error("Failed to measure uploaded file", "error", err)
return nil, err
}
info, err := s.metaviewer.GetInfo(storageFilePath)
fileRecord, err := s.fileRepo.CreateLocal(storageFileName, work)
if err != nil {
s.logger.Error("Failed to get file info", "error", err, "path", storageFilePath)
s.logger.Error("Failed to create file record", "error", err, "file_ext", ext)
return nil, err
}
s.logger.Info("File uploaded successfully",
"file_id", fileId,
"file_id", fileRecord.Id,
"size", size,
"duration_seconds", info.Seconds)
metrics.InputFileDurationHistogram.WithLabelValues().Observe(float64(info.Seconds))
metrics.ObserveInputFileSize(ext, size)
// Создаем запись в таблице files
fileRecord := &entity.File{
Id: fileId,
Storage: entity.StorageLocal,
FileName: storageFileName,
Size: size,
CreatedAt: time.Now(),
}
if err := s.fileRepo.Create(fileRecord); err != nil {
// Удаляем файл если не удалось создать запись в БД
os.Remove(storageFilePath)
s.logger.Error("Failed to create file record", "error", err, "file_id", fileId)
return nil, err
}
job.FileID = &fileId
job.FileID = &fileRecord.Id
if err := s.jobRepo.Create(job); err != nil {
s.logger.Error("Failed to create job record", "error", err, "job_id", job.Id)
s.logger.Error("Failed to create job record", "error", err, "file_id", fileRecord.Id)
return nil, err
}
s.logger.Info("Transcribe job created successfully", "job_id", job.Id, "file_id", fileId)
s.logger.Info("Transcribe job created successfully", "job_id", job.Id, "file_id", fileRecord.Id)
return job, nil
}
func (s *TranscribeService) FindAndRunConversionJob() error {
job, err := s.findJob(entity.StateCreated, time.Hour)
return s.runStep(entity.StateCreated, conversionAcquireTimeout, s.convertJob)
}
func (s *TranscribeService) FindAndRunTranscribeJob() error {
return s.runStep(entity.StateConverted, transcribeAcquireTimeout, s.transcribeJob)
}
func (s *TranscribeService) FindAndRunTranscribeCheckJob() error {
return s.runStep(entity.StateTranscribe, checkAcquireTimeout, s.checkTranscribeJob)
}
// runStep забирает задачу и отдаёт её шагу. Отказ шага не оставляет задачу
// захваченной до конца срока: захват снимается, и задача ждёт нарастающую паузу
// — иначе повтор наступал бы через восемь часов, а не через секунду.
func (s *TranscribeService) runStep(state string, expiration time.Duration, step func(job *entity.TranscribeJob, holder string) error) error {
job, holder, err := s.findJob(state, expiration)
if err != nil {
return err
}
if err := step(job, holder); err != nil {
s.scheduleRetry(job, holder, err)
return err
}
return nil
}
func (s *TranscribeService) convertJob(job *entity.TranscribeJob, holder string) error {
s.logger.Info("Starting conversion job", "job_id", job.Id)
if job.FileID == nil {
s.logger.Error("Job has no file", "job_id", job.Id)
return s.failJob(job, holder, errors.New("job has no file"), "у задачи нет записи")
}
srcFile, err := s.fileRepo.GetByID(*job.FileID)
if err != nil {
s.logger.Error("Failed to get source file", "error", err, "file_id", *job.FileID)
return err
}
srcFilePath := filepath.Join(s.storagePath, srcFile.FileName)
destFileId := uuid.NewString()
destFileName := fmt.Sprintf("%s%s", destFileId, ".ogg")
destFilePath := filepath.Join(s.storagePath, destFileName)
// Получаем расширение исходного файла для метрики
srcExt := strings.TrimPrefix(filepath.Ext(srcFile.FileName), ".")
if srcExt == "" {
srcExt = defaultAudioExt
}
s.logger.Info("Converting file",
"job_id", job.Id,
"src_path", srcFilePath,
"dest_path", destFilePath,
"src_format", srcExt)
src, err := s.fileRepo.Localize(*job.FileID)
if err != nil {
s.logger.Error("Failed to localize source file", "error", err, "file_id", *job.FileID)
return err
}
defer s.closeWork(src)
dest, err := s.fileRepo.StageEmpty(".ogg")
if err != nil {
s.logger.Error("Failed to stage converted file", "error", err, "job_id", job.Id)
return err
}
defer s.closeWork(dest)
s.logger.Info("Converting file", "job_id", job.Id, "src_format", srcExt)
// Измеряем время конвертации
startTime := time.Now()
err = s.converter.Convert(srcFilePath, destFilePath)
err = s.converter.Convert(src.Path(), dest.Path())
conversionDuration := time.Since(startTime)
// Записываем метрику времени конвертации
@@ -214,43 +236,36 @@ func (s *TranscribeService) FindAndRunConversionJob() error {
"error", err,
"job_id", job.Id,
"duration", conversionDuration)
return s.failJob(job, err, "сбой конвертации файла")
return s.failJob(job, holder, err, "сбой конвертации файла")
}
stat, err := os.Stat(destFilePath)
destSize, err := dest.Size()
if err != nil {
s.logger.Error("Failed to stat converted file", "error", err, "path", destFilePath)
s.logger.Error("Failed to measure converted file", "error", err, "job_id", job.Id)
return err
}
s.logger.Info("File conversion completed",
"job_id", job.Id,
"duration", conversionDuration,
"output_size", stat.Size())
"output_size", destSize)
// Записываем метрику размера выходного файла
metrics.OutputFileSizeHistogram.WithLabelValues("ogg").Observe(float64(stat.Size()))
metrics.OutputFileSizeHistogram.WithLabelValues("ogg").Observe(float64(destSize))
// Создаем запись в таблице files
destFileRecord := &entity.File{
Id: destFileId,
Storage: entity.StorageLocal,
FileName: destFileName,
Size: stat.Size(),
CreatedAt: time.Now(),
}
job.FileID = &destFileId
job.MoveToState(entity.StateConverted)
err = s.fileRepo.Create(destFileRecord)
destFileName := fmt.Sprintf("%s%s", uuid.NewString(), ".ogg")
destFileRecord, err := s.fileRepo.CreateLocal(destFileName, dest)
if err != nil {
s.logger.Error("Failed to create converted file record", "error", err, "file_id", destFileId)
s.logger.Error("Failed to create converted file record", "error", err, "job_id", job.Id)
return err
}
err = s.jobRepo.Save(job)
if err != nil {
// Ссылка переставляется только после того, как запись о новом файле есть:
// иначе повтор оставил бы задачу указывающей на файл, которого нет.
job.FileID = &destFileRecord.Id
job.MoveToState(entity.StateConverted)
if err := s.jobRepo.Save(job, holder); err != nil {
s.logger.Error("Failed to save job", "error", err, "job_id", job.Id)
return err
}
@@ -259,36 +274,35 @@ func (s *TranscribeService) FindAndRunConversionJob() error {
return nil
}
func (s *TranscribeService) FindAndRunTranscribeJob() error {
job, err := s.findJob(entity.StateConverted, time.Hour)
if err != nil {
return err
}
func (s *TranscribeService) transcribeJob(job *entity.TranscribeJob, holder string) error {
s.logger.Info("Starting transcribe job", "job_id", job.Id)
if job.FileID == nil {
s.logger.Error("Job has no file", "job_id", job.Id)
return s.failJob(job, holder, errors.New("job has no file"), "у задачи нет записи")
}
fileRecord, err := s.fileRepo.GetByID(*job.FileID)
if err != nil {
s.logger.Error("Failed to get file record", "error", err, "file_id", *job.FileID)
return err
}
filePath := filepath.Join(s.storagePath, fileRecord.FileName)
file, err := os.Open(filePath)
content, err := s.fileRepo.Open(*job.FileID)
if err != nil {
s.logger.Error("Failed to open file", "error", err, "path", filePath)
s.logger.Error("Failed to open file", "error", err, "file_id", *job.FileID)
return err
}
defer file.Close()
defer func() {
if err := content.Close(); err != nil {
s.logger.Error("Failed to close file", "error", err, "file_id", *job.FileID)
}
}()
destFileId := uuid.NewString()
destFileRecord := fileRecord.CopyWithStorage(destFileId, entity.StorageS3)
s.logger.Info("Starting recognition", "job_id", job.Id, "file_path", filePath)
s.logger.Info("Starting recognition", "job_id", job.Id, "file_id", *job.FileID)
// Запускаем асинхронное распознавание
operationID, err := s.recognizer.Recognize(file, destFileRecord.FileName)
operationID, err := s.recognizer.Recognize(content, fileRecord.FileName)
if err != nil {
s.logger.Error("Failed to start recognition", "error", err, "job_id", job.Id)
return err
@@ -298,20 +312,19 @@ func (s *TranscribeService) FindAndRunTranscribeJob() error {
"job_id", job.Id,
"operation_id", operationID)
// Обновляем задачу с ID операции распознавания
job.FileID = &destFileId
job.RecognitionOpID = &operationID
delayTime := time.Now().Add(10 * time.Second)
job.MoveToStateAndDelay(entity.StateTranscribe, &delayTime)
err = s.fileRepo.Create(destFileRecord)
destFileRecord, err := s.fileRepo.CreateRemote(fileRecord.FileName, fileRecord.Size)
if err != nil {
s.logger.Error("Failed to create S3 file record", "error", err, "file_id", destFileId)
s.logger.Error("Failed to create S3 file record", "error", err, "job_id", job.Id)
return err
}
err = s.jobRepo.Save(job)
if err != nil {
// Обновляем задачу с ID операции распознавания
job.FileID = &destFileRecord.Id
job.RecognitionOpID = &operationID
delayTime := time.Now().Add(firstCheckDelay)
job.MoveToStateAndDelay(entity.StateTranscribe, &delayTime)
if err := s.jobRepo.Save(job, holder); err != nil {
s.logger.Error("Failed to save job", "error", err, "job_id", job.Id)
return err
}
@@ -320,12 +333,7 @@ func (s *TranscribeService) FindAndRunTranscribeJob() error {
return nil
}
func (s *TranscribeService) FindAndRunTranscribeCheckJob() error {
job, err := s.findJob(entity.StateTranscribe, 24*time.Hour)
if err != nil {
return err
}
func (s *TranscribeService) checkTranscribeJob(job *entity.TranscribeJob, holder string) error {
if job.RecognitionOpID == nil {
s.logger.Error("Recognition operation ID not found", "job_id", job.Id)
return fmt.Errorf("recogniton opId not found for job: %s", job.Id)
@@ -342,12 +350,13 @@ func (s *TranscribeService) FindAndRunTranscribeCheckJob() error {
}
if recResult.IsInProgress() {
// Операция еще не завершена, оставляем в статусе обработки
// Операция ещё не завершена. Шаг отработал без отказа, поэтому задержка
// здесь своя, числом, а число попыток обнуляется переходом: ожидание
// чужой операции попытку не тратит.
s.logger.Info("Operation in progress", "job_id", job.Id, "operation_id", opId)
delayTime := time.Now().Add(5 * time.Second)
delayTime := time.Now().Add(nextCheckDelay)
job.MoveToStateAndDelay(entity.StateTranscribe, &delayTime)
err := s.jobRepo.Save(job)
if err != nil {
if err := s.jobRepo.Save(job, holder); err != nil {
s.logger.Error("Failed to save job", "error", err, "job_id", job.Id)
return err
}
@@ -360,7 +369,7 @@ func (s *TranscribeService) FindAndRunTranscribeCheckJob() error {
"job_id", job.Id,
"operation_id", opId,
"error_message", errorText)
return s.failJob(job, errors.New(errorText), "сбой при распознавании файла")
return s.failJob(job, holder, errors.New(errorText), "сбой при распознавании файла")
}
// Операция завершена, получаем результат
@@ -376,86 +385,152 @@ func (s *TranscribeService) FindAndRunTranscribeCheckJob() error {
"text_length", len(transcriptionText))
if len(transcriptionText) == 0 {
return s.completeJob(job, "Ой, кажется, на аудиозаписи нет текста.")
return s.completeJob(job, holder, "Ой, кажется, на аудиозаписи нет текста.")
}
// Завершаем задачу
return s.completeJob(job, transcriptionText)
return s.completeJob(job, holder, transcriptionText)
}
func (s *TranscribeService) findJob(state string, expiration time.Duration) (job *entity.TranscribeJob, err error) {
// findJob забирает задачу и отдаёт её вместе с признаком захвата, который шаг
// держит. Задача, захваченная сверх предела попыток, до шага не доходит: её
// переводят в «мертва» и сообщают об этом отправителю.
func (s *TranscribeService) findJob(state string, expiration time.Duration) (*entity.TranscribeJob, string, error) {
acquisitionId := uuid.NewString()
rottingTime := time.Now().Add(-1 * expiration)
job, err = s.jobRepo.FindAndAcquire(state, acquisitionId, rottingTime)
job, err := s.jobRepo.FindAndAcquire(state, acquisitionId, rottingTime)
if err != nil {
// Признак узнаётся по смыслу: репозиторий вправе обернуть свой отказ
// пояснением, и приведение типа от этого сломалось бы молча.
var notFound *contract.JobNotFoundError
if errors.As(err, &notFound) {
return nil, &contract.NoopJobError{State: state}
return nil, "", &contract.NoopJobError{State: state}
}
s.logger.Error("Failed to find and acquire job", "state", state, "error", err)
return nil, fmt.Errorf("failed find and acquire job: %s, %w", state, err)
return nil, "", fmt.Errorf("failed find and acquire job: %s, %w", state, err)
}
return job, nil
if job.Attempts > maxAttempts {
s.killJob(job, acquisitionId)
return nil, "", &contract.NoopJobError{State: state}
}
return job, acquisitionId, nil
}
func (s *TranscribeService) completeJob(job *entity.TranscribeJob, transcriptionText string) error {
// killJob переводит исчерпавшую попытки задачу в «мертва» и сообщает об этом
// отправителю. Инвариант «Принятая запись не теряется молча» допускает два
// исхода — задача пригодна к повтору либо об отказе сказано, — и молчаливая
// смерть не подходит ни под один.
func (s *TranscribeService) killJob(job *entity.TranscribeJob, holder string) {
s.logger.Error("Job exhausted its attempts",
"job_id", job.Id,
"state", job.State,
"attempts", job.Attempts)
job.Die(fmt.Sprintf("attempts exhausted: %d", job.Attempts))
if err := s.jobRepo.Save(job, holder); err != nil {
s.logger.Error("Failed to save dead job", "error", err, "job_id", job.Id)
return
}
s.notify(job, "Не удалось обработать запись: попытки исчерпаны.\nПожалуйста, попробуйте еще раз.")
}
// scheduleRetry снимает захват с отказавшей задачи и ставит нарастающую паузу.
// Захват, оставленный до конца срока, отложил бы повтор на часы.
func (s *TranscribeService) scheduleRetry(job *entity.TranscribeJob, holder string, stepErr error) {
// Шаг, потерявший захват, задачу уже не трогает: ею занят другой.
var lost *contract.LostAcquisitionError
if errors.As(stepErr, &lost) {
return
}
job.RetryAfter(time.Now().Add(retryDelay(job.Attempts)))
if err := s.jobRepo.Save(job, holder); err != nil {
var lostOnSave *contract.LostAcquisitionError
if errors.As(err, &lostOnSave) {
return
}
s.logger.Error("Failed to schedule job retry", "error", err, "job_id", job.Id)
}
}
// retryDelay растит паузу с числом попыток до потолка.
func retryDelay(attempts int) time.Duration {
if attempts < 1 {
attempts = 1
}
delay := time.Duration(math.Pow(2, float64(attempts-1))) * retryDelayBase
if delay > retryDelayCap || delay <= 0 {
return retryDelayCap
}
return delay
}
func (s *TranscribeService) completeJob(job *entity.TranscribeJob, holder string, transcriptionText string) error {
// Обновляем задачу с результатом
job.Done(transcriptionText)
// Сохраняем задачу в базу
err := s.jobRepo.Save(job)
if err != nil {
if err := s.jobRepo.Save(job, holder); err != nil {
s.logger.Error("Failed to save job", "error", err, "job_id", job.Id)
return fmt.Errorf("failed to save job: %w", err)
}
// Отправляем распознанный текст обратно пользователю
switch job.Source {
case entity.SourceTelegram:
if job.TgChatId == nil {
s.logger.Error("Telegram chat not specified", "job_id", job.Id)
return fmt.Errorf("tg chat id not specified, job id: %s", job.Id)
}
err := s.tgSender.Send(transcriptionText, *job.TgChatId, job.TgReplyMessageId)
if err != nil {
s.logger.Error("Failed to sent transcription text to client", "job_id", job.Id)
return fmt.Errorf("failed to sent message to client, job id: %s, err: %w", job.Id, err)
}
}
return nil
return s.send(job, transcriptionText)
}
func (s *TranscribeService) failJob(job *entity.TranscribeJob, jobErr error, humanErrorText string) error {
func (s *TranscribeService) failJob(job *entity.TranscribeJob, holder string, jobErr error, humanErrorText string) error {
// Обновляем задачу с результатом
job.Fail(jobErr.Error())
// Сохраняем задачу в базу
err := s.jobRepo.Save(job)
if err != nil {
if err := s.jobRepo.Save(job, holder); err != nil {
s.logger.Error("Failed to save job", "error", err, "job_id", job.Id)
return fmt.Errorf("failed to save job: %w", err)
}
// Отправляем текст об ошибке пользователю
switch job.Source {
case entity.SourceTelegram:
if job.TgChatId == nil {
s.logger.Error("Telegram chat not specified", "job_id", job.Id)
return fmt.Errorf("tg chat id not specified, job id: %s", job.Id)
}
errorMessage := fmt.Sprintf("При обработке задачи произошла ошибка: %s.\nПожалуйста, попробуйте еще раз.", humanErrorText)
return s.send(job, errorMessage)
}
errorMessage := fmt.Sprintf("При обработке задачи произошла ошибка: %s.\nПожалуйста, попробуйте еще раз.", humanErrorText)
err := s.tgSender.Send(errorMessage, *job.TgChatId, job.TgReplyMessageId)
if err != nil {
s.logger.Error("Failed to sent message to client", "job_id", job.Id)
return fmt.Errorf("failed to sent message to client, job id: %s, err: %w", job.Id, err)
}
// send отвечает отправителю там, откуда пришла запись, и отказ отправки
// поднимает вверх: он принадлежит шагу.
func (s *TranscribeService) send(job *entity.TranscribeJob, text string) error {
if job.Source != entity.SourceTelegram {
return nil
}
if job.TgChatId == nil {
s.logger.Error("Telegram chat not specified", "job_id", job.Id)
return fmt.Errorf("tg chat id not specified, job id: %s", job.Id)
}
if err := s.tgSender.Send(text, *job.TgChatId, job.TgReplyMessageId); err != nil {
s.logger.Error("Failed to sent message to client", "job_id", job.Id)
return fmt.Errorf("failed to sent message to client, job id: %s, err: %w", job.Id, err)
}
return nil
}
// notify отвечает отправителю там, где поднимать отказ некуда: задача уже
// доведена до конца, и отказ отправки остаётся записью в журнале владельца.
func (s *TranscribeService) notify(job *entity.TranscribeJob, text string) {
if err := s.send(job, text); err != nil {
s.logger.Error("Failed to notify sender", "error", err, "job_id", job.Id)
}
}
// closeWork убирает рабочую копию. Отказ уборки не роняет шаг, но и не
// проглатывается: забытая копия это шестичасовая запись во временном каталоге.
func (s *TranscribeService) closeWork(work contract.WorkFile) {
if err := work.Close(); err != nil {
s.logger.Error("Failed to remove work file", "error", err)
}
}