хранилище переехало с PocketBase на SQLite со своим каталогом файлов

- база своя: два пула, захват одним UPDATE ... RETURNING, шаги схемы на goose
  под файловым замком, одна миграция начальной схемы вместо семи прежних
- транспорт переписан на net/http: свои слои, свой ограничитель частоты,
  отдача файла с проверкой владельца; панель /_/ и пространство /api/ исчезли
- по находкам ревью: журнал не пишет путь под корнем приложения, ключ бюджета
  читается справа налево, узнавание известного идёт читающим пулом
This commit is contained in:
av
2026-08-23 08:06:04 +03:00
parent 1edf8cb225
commit c9b7765646
118 changed files with 11668 additions and 6679 deletions
+6 -6
View File
@@ -18,11 +18,11 @@ func TestWorkerTakesRecordsOfEveryOwner(t *testing.T) {
env := newPipelineEnv(t, &okMetaViewer{}, &failingConverter{})
first, err := env.service.CreateJobFromApi(t.Context(),
strings.NewReader("первая"), "one.mp3", newOwner(t, env.app))
strings.NewReader("первая"), "one.mp3", newOwner(t, env))
require.NoError(t, err)
second, err := env.service.CreateJobFromApi(t.Context(),
strings.NewReader("вторая"), "two.mp3", newOwner(t, env.app))
strings.NewReader("вторая"), "two.mp3", newOwner(t, env))
require.NoError(t, err)
// Третья — ещё одного владельца: воркер не сужается ни одним из них.
@@ -52,7 +52,7 @@ func TestWorkerTakesRecordsOfEveryOwner(t *testing.T) {
func TestAcquireReturnsIdentifierAndHolder(t *testing.T) {
env := newPipelineEnv(t, &okMetaViewer{}, &failingConverter{})
owner := newOwner(t, env.app)
owner := newOwner(t, env)
record, err := env.service.CreateJobFromApi(t.Context(), strings.NewReader("запись"), "one.mp3", owner)
require.NoError(t, err)
@@ -74,7 +74,7 @@ func TestAcquireReturnsIdentifierAndHolder(t *testing.T) {
func TestPipelineStepKeepsOwner(t *testing.T) {
env := newPipelineEnv(t, &okMetaViewer{}, &failingConverter{})
owner := newOwner(t, env.app)
owner := newOwner(t, env)
record, err := env.service.CreateJobFromApi(t.Context(), strings.NewReader("запись"), "one.mp3", owner)
require.NoError(t, err)
@@ -102,8 +102,8 @@ func TestCreateJobFromApiRequiresOwner(t *testing.T) {
func TestGetByIDHidesForeignRecords(t *testing.T) {
env := newPipelineEnv(t, &okMetaViewer{}, &failingConverter{})
owner := newOwner(t, env.app)
stranger := newOwner(t, env.app)
owner := newOwner(t, env)
stranger := newOwner(t, env)
record, err := env.service.CreateJobFromApi(t.Context(), strings.NewReader("запись"), "one.mp3", owner)
require.NoError(t, err)
+86 -86
View File
@@ -12,19 +12,21 @@ import (
"testing"
"time"
"github.com/google/uuid"
"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/adapter/repo/pocketbase/migrations"
sqliterepo "git.vakhrushev.me/av/transcriber/internal/adapter/repo/sqlite"
"git.vakhrushev.me/av/transcriber/internal/clock"
"git.vakhrushev.me/av/transcriber/internal/contract"
"git.vakhrushev.me/av/transcriber/internal/entity"
"git.vakhrushev.me/av/transcriber/internal/ident"
)
// timeLayout — вид времени в колонках базы. Проверки, двигающие сроки, пишут его
// тем же видом, каким пишет хранилище: сравнение там побайтово.
const timeLayout = "2006-01-02T15:04:05Z"
// Проверки конвейера идут против настоящего хранилища: захват, число отказов и
// остановка держатся на запросе, и подставной репозиторий проверял бы
// собственную заглушку, а не то, что делает база.
@@ -61,11 +63,13 @@ func (m *failingMetaViewer) GetInfo(context.Context, string) (*contract.AudioInf
}
type pipelineEnv struct {
app core.App
db *sqliterepo.DB
store *sqliterepo.Store
users *sqliterepo.UserRepository
service *TranscribeService
repos Repositories
recordRepo *pbrepo.AudioRecordRepository
fileRepo *pbrepo.FileRepository
recordRepo *sqliterepo.AudioRecordRepository
fileRepo *sqliterepo.FileRepository
}
// testLimits — пределы простоя проверок. Числа боевые; проверка застревания
@@ -98,33 +102,34 @@ func newPipelineEnvWithLogger(
) *pipelineEnv {
t.Helper()
app, err := pbrepo.New(t.TempDir())
dir := t.TempDir()
db, err := sqliterepo.Open(dir, sqliterepo.Settings{BusyTimeoutMs: 5000, ReadConnections: 4})
require.NoError(t, err)
t.Cleanup(func() {
if err := app.ResetBootstrapState(); err != nil {
t.Logf("не удалось закрыть хранилище: %v", err)
if err := db.Close(); err != nil {
t.Logf("не удалось закрыть базу: %v", err)
}
})
// Правила панели вешаются и здесь: конфигурация под проверкой обязана
// совпадать с боевой, иначе утверждения говорят про прод то, чего в проде
// нет.
pbrepo.BindPanelRules(app)
require.NoError(t, sqliterepo.Migrate(t.Context(), db, dir, slog.New(slog.DiscardHandler)))
recordRepo := pbrepo.NewAudioRecordRepository(app)
fileRepo := pbrepo.NewFileRepository(app)
store := sqliterepo.NewStore(dir)
recordRepo := sqliterepo.NewAudioRecordRepository(db)
fileRepo := sqliterepo.NewFileRepository(db, store)
repos := Repositories{
Records: recordRepo,
Files: fileRepo,
Texts: pbrepo.NewTextRepository(app),
Structures: pbrepo.NewStructureRepository(app),
Recognitions: pbrepo.NewRecognitionRepository(app),
Events: pbrepo.NewRecordEventRepository(app),
Texts: sqliterepo.NewTextRepository(db),
Structures: sqliterepo.NewStructureRepository(db),
Recognitions: sqliterepo.NewRecognitionRepository(db, store),
Events: sqliterepo.NewRecordEventRepository(db),
}
svc := NewTranscribeService(repos, metaviewer, converter, rec, testLimits, logger)
return &pipelineEnv{
app: app,
db: db,
store: store,
users: sqliterepo.NewUserRepository(db),
service: svc,
repos: repos,
recordRepo: recordRepo,
@@ -132,13 +137,22 @@ func newPipelineEnvWithLogger(
}
}
// exec выполняет запрос к базе от имени проверки: фикстуры двигают колонки
// напрямую там, где домен такого перехода не делает.
func (e *pipelineEnv) exec(t *testing.T, query string, args ...any) {
t.Helper()
_, err := e.db.Writer().ExecContext(context.Background(), query, args...)
require.NoError(t, err)
}
// newRecord заводит запись — так, как её заводит приём по HTTP: от имени
// вошедшего, потому что ничьей записи в хранилище не бывает.
func newRecord(t *testing.T, env *pipelineEnv) *entity.AudioRecord {
t.Helper()
record, err := env.service.CreateJobFromApi(
t.Context(), strings.NewReader("запись"), "voice.ogg", newOwner(t, env.app))
t.Context(), strings.NewReader("запись"), "voice.ogg", newOwner(t, env))
require.NoError(t, err)
return record
}
@@ -147,10 +161,7 @@ func newRecord(t *testing.T, env *pipelineEnv) *entity.AudioRecord {
func clearDelay(t *testing.T, env *pipelineEnv, recordID string) {
t.Helper()
record, err := env.app.FindRecordById(migrations.RecordsCollection, recordID)
require.NoError(t, err)
record.Set("delay_time", "")
require.NoError(t, env.app.Save(record))
env.exec(t, "UPDATE audio_records SET delay_time = NULL WHERE id = ?", recordID)
}
// enteredStateAt отодвигает время входа записи в рубеж: так это выглядит, когда
@@ -158,13 +169,8 @@ func clearDelay(t *testing.T, env *pipelineEnv, recordID string) {
func enteredStateAt(t *testing.T, env *pipelineEnv, recordID string, moment time.Time) {
t.Helper()
record, err := env.app.FindRecordById(migrations.RecordsCollection, recordID)
require.NoError(t, err)
record.Set("state_entered_at", types.DateTime{}.Add(0))
stamp, err := types.ParseDateTime(moment)
require.NoError(t, err)
record.Set("state_entered_at", stamp)
require.NoError(t, env.app.Save(record))
env.exec(t, "UPDATE audio_records SET state_entered_at = ? WHERE id = ?",
moment.UTC().Format(timeLayout), recordID)
}
// setAttempts ставит записи число отказов: так она выглядит, когда шаг отказал
@@ -172,10 +178,7 @@ func enteredStateAt(t *testing.T, env *pipelineEnv, recordID string, moment time
func setAttempts(t *testing.T, env *pipelineEnv, recordID string, attempts int) {
t.Helper()
record, err := env.app.FindRecordById(migrations.RecordsCollection, recordID)
require.NoError(t, err)
record.Set("attempts", attempts)
require.NoError(t, env.app.Save(record))
env.exec(t, "UPDATE audio_records SET attempts = ? WHERE id = ?", attempts, recordID)
}
// drain крутит конвейер, пока он двигает записи. Паузы опроса снимаются: они
@@ -253,10 +256,8 @@ func TestRecordHaltsAfterAttemptLimit(t *testing.T) {
func expireAcquisition(t *testing.T, env *pipelineEnv, recordID string) {
t.Helper()
record, err := env.app.FindRecordById(migrations.RecordsCollection, recordID)
require.NoError(t, err)
record.Set("acquire_expires_at", types.NowDateTime().Add(-time.Hour))
require.NoError(t, env.app.Save(record))
env.exec(t, "UPDATE audio_records SET acquire_expires_at = ? WHERE id = ?",
clock.Now().Add(-time.Hour).Format(timeLayout), recordID)
}
// Критерий приёмки 1. Остановленная на шаге запись перезапускается снятием
@@ -508,11 +509,10 @@ func TestEveryHaltReasonRecordsItsCause(t *testing.T) {
events := recordEvents(t, env, record.Id)
require.Len(t, events, 1, "остановка оставила строку журнала событий")
outcome := events[0].GetString("outcome")
assert.Contains(t,
[]string{entity.EventOutcomeHalted, entity.EventOutcomeFailed}, outcome,
[]string{entity.EventOutcomeHalted, entity.EventOutcomeFailed}, events[0].Outcome,
"строка журнала называет исход остановкой")
assert.NotEmpty(t, events[0].GetString("outcome_text"), "и несёт причину")
assert.NotEmpty(t, events[0].OutcomeText, "и несёт причину")
})
}
}
@@ -536,9 +536,9 @@ func TestRepeatedStoreKeepsSingleText(t *testing.T) {
require.NoError(t, err)
assert.Equal(t, "второй разбор", stored.Contents)
count, err := env.app.CountRecords(migrations.TextsCollection)
require.NoError(t, err)
assert.Equal(t, int64(1), count, "второго комплекта строк не завелось")
var count int
require.NoError(t, env.db.Reader().QueryRowContext(context.Background(), "SELECT COUNT(*) FROM texts").Scan(&count))
assert.Equal(t, 1, count, "второго комплекта строк не завелось")
}
// Пауза растёт с числом отказов и упирается в потолок.
@@ -559,19 +559,12 @@ func TestFailedStepSchedulesRetryWithGrowingDelay(t *testing.T) {
// Ссылка переставляется на запись о файле без содержимого: шаг отказывает на
// получении рабочей копии — то есть отказом, а не приговором записи.
files, err := env.app.FindCollectionByNameOrId(migrations.FilesCollection)
require.NoError(t, err)
empty := core.NewRecord(files)
empty.Set("location", entity.LocationLocal)
empty.Set("size", 1)
emptyID := ident.New()
// Владелец обязателен и у файла: схема ничьих не принимает.
empty.Set("owner", record.OwnerID)
require.NoError(t, env.app.Save(empty))
stored, err := env.app.FindRecordById(migrations.RecordsCollection, record.Id)
require.NoError(t, err)
stored.Set("original_file", empty.Id)
require.NoError(t, env.app.Save(stored))
env.exec(t, `INSERT INTO files (id, owner_id, record_id, file_name, size_bytes, created_at)
VALUES (?, ?, ?, ?, ?, ?)`,
emptyID, record.OwnerID, record.Id, "missing.mp3", 1, clock.Now().Format(timeLayout))
env.exec(t, "UPDATE audio_records SET original_file_id = ? WHERE id = ?", emptyID, record.Id)
require.Error(t, env.service.RunStep(t.Context()))
@@ -596,7 +589,7 @@ func TestWorkFileRemovedAfterIntakeFailure(t *testing.T) {
env := newPipelineEnv(t, &failingMetaViewer{}, &failingConverter{})
_, err := env.service.CreateJobFromApi(t.Context(), strings.NewReader("запись"), "sample.mp3", newOwner(t, env.app))
_, err := env.service.CreateJobFromApi(t.Context(), strings.NewReader("запись"), "sample.mp3", newOwner(t, env))
require.Error(t, err, "отказ источника метаданных роняет приём")
leftovers, err := filepath.Glob(filepath.Join(tempDir, "transcriber-*"))
@@ -612,7 +605,7 @@ func TestWorkFileRemovedAfterSuccessfulIntake(t *testing.T) {
env := newPipelineEnv(t, &okMetaViewer{}, &failingConverter{})
_, err := env.service.CreateJobFromApi(t.Context(), strings.NewReader("запись"), "sample.mp3", newOwner(t, env.app))
_, err := env.service.CreateJobFromApi(t.Context(), strings.NewReader("запись"), "sample.mp3", newOwner(t, env))
require.NoError(t, err)
leftovers, err := filepath.Glob(filepath.Join(tempDir, "transcriber-*"))
@@ -666,7 +659,7 @@ func TestStoredContentSurvivesRoundTrip(t *testing.T) {
content := strings.Repeat("запись ", 1000)
record, err := env.service.CreateJobFromApi(t.Context(), strings.NewReader(content), "sample.mp3", newOwner(t, env.app))
record, err := env.service.CreateJobFromApi(t.Context(), strings.NewReader(content), "sample.mp3", newOwner(t, env))
require.NoError(t, err)
require.NotNil(t, record.OriginalFileID)
@@ -689,7 +682,7 @@ func TestStoredContentSurvivesRoundTrip(t *testing.T) {
func TestLocalizeGivesReadableCopy(t *testing.T) {
env := newPipelineEnv(t, &okMetaViewer{}, &failingConverter{})
record, err := env.service.CreateJobFromApi(t.Context(), strings.NewReader("содержимое"), "sample.mp3", newOwner(t, env.app))
record, err := env.service.CreateJobFromApi(t.Context(), strings.NewReader("содержимое"), "sample.mp3", newOwner(t, env))
require.NoError(t, err)
require.NotNil(t, record.OriginalFileID)
@@ -707,24 +700,17 @@ func TestLocalizeGivesReadableCopy(t *testing.T) {
// newOwner заводит учётную запись и отдаёт её идентификатор.
//
// Владелец — связь с коллекцией пользователей, и хранилище проверяет, что такая
// запись есть: выдуманный идентификатор запись завести не даст.
func newOwner(t *testing.T, app core.App) string {
// Владелец — связь с таблицей учётных записей, и база проверяет, что такая
// строка есть: выдуманный идентификатор запись завести не даст.
func newOwner(t *testing.T, env *pipelineEnv) string {
t.Helper()
users, err := app.FindCollectionByNameOrId(migrations.UsersCollection)
// Логин у провайдера — ключ учётной записи, и он уникален: две записи с
// одним ключом схема не примет.
account, _, err := env.users.EnsureUser(contract.Identity{Login: ident.New()})
require.NoError(t, err)
record := core.NewRecord(users)
// Логин у провайдера — ключ учётной записи, и он уникален: две записи с
// пустым ключом схема не примет.
record.Set(migrations.ProviderLoginField, uuid.NewString())
record.Set("email", uuid.NewString()+"@example.test")
record.Set("verified", true)
record.Set("password", uuid.NewString())
require.NoError(t, app.Save(record))
return record.Id
return account.ID
}
// Остановка приговором шага засчитывается **отказом**, а не успехом, и пишет в
@@ -752,7 +738,7 @@ func TestHaltIsCountedAsFailureAndLoggedOnce(t *testing.T) {
events := recordEvents(t, env, record.Id)
require.Len(t, events, 1, "одна строка журнала, а не две")
assert.Equal(t, entity.EventOutcomeFailed, events[0].GetString("outcome"),
assert.Equal(t, entity.EventOutcomeFailed, events[0].Outcome,
"исход назван приговором, а не сделанной работой")
}
@@ -780,18 +766,32 @@ func TestPostponeWritesNoEvent(t *testing.T) {
"пять откладываний не оставили в журнале ни строки")
}
// recordEventRow — строка журнала событий записи, какой её видит проверка.
type recordEventRow struct {
Origin string
Step string
Outcome string
OutcomeText string
}
// recordEvents читает журнал событий одной записи в порядке заведения.
func recordEvents(t *testing.T, env *pipelineEnv, recordID string) []*core.Record {
func recordEvents(t *testing.T, env *pipelineEnv, recordID string) []recordEventRow {
t.Helper()
all, err := env.app.FindAllRecords(migrations.RecordEventsCollection)
rows, err := env.db.Reader().QueryContext(context.Background(),
"SELECT origin, step, outcome, outcome_text FROM record_events WHERE record_id = ? ORDER BY id",
recordID,
)
require.NoError(t, err)
defer func() { require.NoError(t, rows.Close()) }()
var own []*core.Record
for _, event := range all {
if event.GetString("record") == recordID {
own = append(own, event)
}
var own []recordEventRow
for rows.Next() {
var event recordEventRow
require.NoError(t, rows.Scan(&event.Origin, &event.Step, &event.Outcome, &event.OutcomeText))
own = append(own, event)
}
require.NoError(t, rows.Err())
return own
}
+15 -13
View File
@@ -11,11 +11,10 @@ import (
"strings"
"time"
"github.com/google/uuid"
"git.vakhrushev.me/av/transcriber/internal/clock"
"git.vakhrushev.me/av/transcriber/internal/contract"
"git.vakhrushev.me/av/transcriber/internal/entity"
"git.vakhrushev.me/av/transcriber/internal/ident"
"git.vakhrushev.me/av/transcriber/internal/metrics"
)
@@ -155,9 +154,12 @@ func (s *TranscribeService) CreateJobFromApi(ctx context.Context, file io.Reader
return nil, contract.ErrOwnerRequired
}
// Идентификатор назначается здесь, до укладки файла: копии записи лежат её
// подкаталогом, названным этим идентификатором, и знать его надо раньше, чем
// класть первую копию.
record := &entity.AudioRecord{
Id: ident.New(),
State: entity.StateUploaded,
Source: entity.SourceApi,
OwnerID: ownerID,
}
@@ -176,10 +178,10 @@ func (s *TranscribeService) createRecord(ctx context.Context, r *entity.AudioRec
ext = fmt.Sprintf(".%s", defaultAudioExt)
}
// Собственное имя записи: идентификатор с расширением. Имя, данное
// отправителем, в хранилище не попадает — от него взято только расширение.
fileId := uuid.NewString()
storageFileName := fmt.Sprintf("%s%s", fileId, ext)
// Собственное имя копии: идентификатор с расширением. Имя, данное
// отправителем, в каталог данных не попадает — от него взято только
// расширение.
storageFileName := fmt.Sprintf("%s%s", ident.New(), ext)
// Содержимое ложится в рабочую копию потоком: в память запись целиком не
// читается, расчётный потолок — шесть часов.
@@ -217,7 +219,7 @@ func (s *TranscribeService) createRecord(ctx context.Context, r *entity.AudioRec
DurationMs: int64(info.Seconds) * 1000,
}
fileRecord, err := s.repos.Files.Create(storageFileName, work, meta, r.OwnerID)
fileRecord, err := s.repos.Files.Create(r.Id, storageFileName, work, meta, r.OwnerID)
if err != nil {
s.logger.Error("Failed to create file record", "error", err, "file_ext", ext)
return nil, err
@@ -474,9 +476,9 @@ func (s *TranscribeService) normalize(ctx context.Context, r *entity.AudioRecord
metrics.OutputFileSizeHistogram.WithLabelValues("ogg").Observe(float64(destSize))
destFileName := fmt.Sprintf("%s%s", uuid.NewString(), ".ogg")
destFileName := fmt.Sprintf("%s%s", ident.New(), ".ogg")
destMeta := contract.FileMeta{Format: "ogg", DurationMs: srcFile.DurationMs}
destFileRecord, err := s.repos.Files.Create(destFileName, dest, destMeta, r.OwnerID)
destFileRecord, err := s.repos.Files.Create(r.Id, destFileName, dest, destMeta, r.OwnerID)
if err != nil {
s.logger.Error("Failed to create normalized file record", "error", err, "record_id", r.Id)
return outcomeDone, err
@@ -780,8 +782,8 @@ func (s *TranscribeService) storeOutcome(r *entity.AudioRecord, outcome *entity.
}
// finish доводит запись до конечного рубежа. Наружу шаг не обращается: доставки
// ответа отправителю у сервиса нет, и свой исход отправитель узнаёт опросом
// готовности.
// ответа отправителю у сервиса нет, и свой исход владелец записи узнаёт её
// карточкой.
func (s *TranscribeService) finish(ctx context.Context, r *entity.AudioRecord, holder string) (stepOutcome, error) {
r.MoveToState(entity.StateDone)
if err := s.repos.Records.Save(r, holder); err != nil {
@@ -806,7 +808,7 @@ func (s *TranscribeService) failStep(r *entity.AudioRecord, holder, step string,
// событий.
//
// Отправителю отсюда ничего не уходит: инвариант проекта «Принятая запись не
// теряется молча» держится теперь опросом готовности — остановка видна там
// теряется молча» держится теперь карточкой записи — остановка видна там
// признаком — и журналом владельца, где у неё стоит причина.
//
// Счётчик растит **сама остановка**, а не воркер, и это не стилистика.