package service import ( "context" "errors" "io" "log/slog" "os" "path/filepath" "strings" "sync" "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" "git.vakhrushev.me/av/transcriber/internal/contract" "git.vakhrushev.me/av/transcriber/internal/entity" ) // Проверки конвейера идут против настоящего хранилища: захват, число отказов и // остановка держатся на запросе, и подставной репозиторий проверял бы // собственную заглушку, а не то, что делает база. // okConverter переливает исходник в результат — так это выглядит у настоящего // приведения к рабочему формату. type okConverter struct{} func (c *okConverter) Convert(_ context.Context, src, dest string) error { content, err := os.ReadFile(src) if err != nil { return err } return os.WriteFile(dest, content, 0o600) } // failingConverter отказывает на каждой попытке. type failingConverter struct{} func (c *failingConverter) Convert(context.Context, string, string) error { return errors.New("конвертация не удалась") } type okMetaViewer struct{} func (m *okMetaViewer) GetInfo(context.Context, string) (*contract.AudioInfo, error) { return &contract.AudioInfo{Seconds: 1}, nil } type failingMetaViewer struct{} func (m *failingMetaViewer) GetInfo(context.Context, string) (*contract.AudioInfo, error) { return nil, errors.New("запись не читается") } type pipelineEnv struct { app core.App service *TranscribeService repos Repositories recordRepo *pbrepo.AudioRecordRepository fileRepo *pbrepo.FileRepository } // testLimits — пределы простоя проверок. Числа боевые; проверка застревания // двигает не их, а время входа записи в рубеж: сторож меряет именно его. var testLimits = entity.StuckLimits{Own: time.Hour, Foreign: 24 * time.Hour} func newPipelineEnv(t *testing.T, metaviewer contract.AudioMetaViewer, converter contract.AudioFileConverter) *pipelineEnv { t.Helper() return newPipelineEnvWith(t, metaviewer, converter, &recognizer.MemoryAudioRecognizer{}) } func newPipelineEnvWith( t *testing.T, metaviewer contract.AudioMetaViewer, converter contract.AudioFileConverter, rec contract.AudioRecognizer, ) *pipelineEnv { t.Helper() return newPipelineEnvWithLogger(t, metaviewer, converter, rec, slog.New(slog.DiscardHandler)) } // newPipelineEnvWithLogger — то же с подменённым журналом: проверке, судящей // строку журнала, нужен свой, а не общий. func newPipelineEnvWithLogger( t *testing.T, metaviewer contract.AudioMetaViewer, converter contract.AudioFileConverter, rec contract.AudioRecognizer, logger *slog.Logger, ) *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) recordRepo := pbrepo.NewAudioRecordRepository(app) fileRepo := pbrepo.NewFileRepository(app) repos := Repositories{ Records: recordRepo, Files: fileRepo, Texts: pbrepo.NewTextRepository(app), Structures: pbrepo.NewStructureRepository(app), Recognitions: pbrepo.NewRecognitionRepository(app), Events: pbrepo.NewRecordEventRepository(app), } svc := NewTranscribeService(repos, metaviewer, converter, rec, testLimits, logger) return &pipelineEnv{ app: app, service: svc, repos: repos, recordRepo: recordRepo, fileRepo: fileRepo, } } // 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)) require.NoError(t, err) return record } // clearDelay снимает паузу, чтобы следующий прогон взял запись сразу. 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)) } // enteredStateAt отодвигает время входа записи в рубеж: так это выглядит, когда // запись простояла в нём дольше предела. 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)) } // setAttempts ставит записи число отказов: так она выглядит, когда шаг отказал // столько раз подряд, что предел исчерпан. 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)) } // drain крутит конвейер, пока он двигает записи. Паузы опроса снимаются: они // проверяются отдельно, а здесь мешают дойти до конца. func drain(t *testing.T, env *pipelineEnv, recordID string) { t.Helper() for range 20 { clearDelay(t, env, recordID) err := env.service.RunStep(t.Context()) var noop *contract.NoopJobError if errors.As(err, &noop) { return } require.NoError(t, err) after := readRecord(t, env, recordID) if after.State == entity.StateDone || after.IsHalted() { return } } t.Fatal("конвейер не дошёл до исхода за отведённое число прогонов") } // readRecord читает запись мимо сужения владельцем. // // Читающий метод сервиса отдаёт запись только её владельцу, а проверки конвейера // смотрят записи, принятые ботом: владельца у таких нет вовсе. Проверке нужно не // право, а состояние записи после шага. func readRecord(t *testing.T, env *pipelineEnv, id string) *entity.AudioRecord { t.Helper() record, err := env.recordRepo.Get(id) require.NoError(t, err) return record } // Запись, отказывающая на каждой попытке, останавливается признаком: из выборки // исчезает, рубеж сохраняет, а отправитель узнаёт о неудаче. Инвариант «Принятая // запись не теряется молча» допускает два исхода, и молчаливая остановка не // подходит ни под один. func TestRecordHaltsAfterAttemptLimit(t *testing.T) { env := newPipelineEnv(t, &okMetaViewer{}, &failingConverter{}) record := newRecord(t, env) // Захват без выполнения шага — так это выглядит при гибели процесса: отказа // шаг объявить не успевает, а попытка засчитана. for range maxAttempts { _, err := env.recordRepo.FindAndAcquire(entity.WorkingStages()) require.NoError(t, err) expireAcquisition(t, env, record.Id) } err := env.service.RunStep(t.Context()) var noop *contract.NoopJobError require.ErrorAs(t, err, &noop, "остановленная запись шагу не отдаётся") after := readRecord(t, env, record.Id) assert.True(t, after.IsHalted(), "запись остановлена признаком") require.NotNil(t, after.HaltReason) assert.Equal(t, entity.HaltReasonAttempts, *after.HaltReason) assert.Equal(t, entity.StateUploaded, after.State, "рубеж пережил остановку") assert.Greater(t, after.Attempts, maxAttempts, "число отказов сохранено") // И из выборки она исчезла. _, err = env.recordRepo.FindAndAcquire(entity.WorkingStages()) var missing *contract.JobNotFoundError assert.ErrorAs(t, err, &missing) } // expireAcquisition отодвигает срок протухания захвата в прошлое: так это // выглядит, когда шаг оборвался вместе с процессом. 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)) } // Критерий приёмки 1. Остановленная на шаге запись перезапускается снятием // признака и продолжает с того рубежа, где стояла, — следующим идёт отправка на // распознавание, а не повторное приведение. func TestHaltedRecordResumesFromItsStage(t *testing.T) { env := newPipelineEnv(t, &okMetaViewer{}, &okConverter{}) record := newRecord(t, env) // Доводим до рубежа приведения и останавливаем на нём. require.NoError(t, env.service.RunStep(t.Context())) after := readRecord(t, env, record.Id) require.Equal(t, entity.StateNormalized, after.State) after.Halt(entity.HaltReasonStepFailed, "проверочная остановка") require.NoError(t, env.recordRepo.Save(after, "")) // Остановленная запись захвату не выдаётся. _, err := env.recordRepo.FindAndAcquire(entity.WorkingStages()) var missing *contract.JobNotFoundError require.ErrorAs(t, err, &missing, "остановленная запись из выборки исчезла") // Снятие признака возвращает её в работу с сохранённого рубежа. halted := readRecord(t, env, record.Id) halted.Resume() require.NoError(t, env.recordRepo.Save(halted, "")) resumed := readRecord(t, env, record.Id) assert.Equal(t, entity.StateNormalized, resumed.State, "рубеж сохранён") assert.Equal(t, 0, resumed.Attempts, "отказы сброшены") // Следующим идёт отправка на распознавание, а не повторное приведение. require.NoError(t, env.service.RunStep(t.Context())) next := readRecord(t, env, record.Id) assert.Equal(t, entity.StateSubmitted, next.State, "продолжили с места остановки") assert.NotNil(t, next.RecognitionID, "отправка состоялась") } // Критерий приёмки 2. У прошедшей конвейер записи ссылки на исходник и на // приведённую копию ведут на разные существующие файлы: прежде ссылка была одна, // и каждый шаг переставлял её на свой результат. func TestBothFileLinksSurvivePipeline(t *testing.T) { env := newPipelineEnv(t, &okMetaViewer{}, &okConverter{}) record := newRecord(t, env) drain(t, env, record.Id) after := readRecord(t, env, record.Id) require.Equal(t, entity.StateDone, after.State) require.NotNil(t, after.OriginalFileID, "ссылка на принятую копию заполнена") require.NotNil(t, after.NormalizedFileID, "ссылка на приведённую копию заполнена") assert.NotEqual(t, *after.OriginalFileID, *after.NormalizedFileID, "копии разные") for _, id := range []string{*after.OriginalFileID, *after.NormalizedFileID} { reader, err := env.fileRepo.Open(id) require.NoError(t, err, "обе копии открываются") content, err := io.ReadAll(reader) require.NoError(t, err) assert.NotEmpty(t, content) require.NoError(t, reader.Close()) } } // Критерий приёмки 3. Поведение записи не зависит от числа воркеров: при одном // и при нескольких она доходит до конечного рубежа. func TestOutcomeDoesNotDependOnWorkerCount(t *testing.T) { for _, workers := range []int{1, 4} { env := newPipelineEnv(t, &okMetaViewer{}, &okConverter{}) record := newRecord(t, env) for range 20 { clearDelay(t, env, record.Id) var wg sync.WaitGroup for range workers { wg.Add(1) go func() { defer wg.Done() // Исход прогона здесь не судится: воркеров несколько, и всем, // кроме одного, работы не достаётся. Судится состояние записи. if err := env.service.RunStep(t.Context()); err != nil { t.Log(err) } }() } wg.Wait() if readRecord(t, env, record.Id).State == entity.StateDone { break } } after := readRecord(t, env, record.Id) assert.Equalf(t, entity.StateDone, after.State, "запись дошла до конца при %d воркерах", workers) assert.Falsef(t, after.IsHalted(), "запись не остановлена при %d воркерах", workers) } } // Тот же критерий, вторая половина: нулевое число воркеров — законное значение. // Записи принимаются и не двигаются, и это режим, а не поломка. func TestZeroWorkersLeaveRecordUntouched(t *testing.T) { env := newPipelineEnv(t, &okMetaViewer{}, &okConverter{}) record := newRecord(t, env) // Пул нулевого размера ни одного прогона не делает — потому запись остаётся // там, где её оставил приём. after := readRecord(t, env, record.Id) assert.Equal(t, entity.StateUploaded, after.State, "запись принята и стоит на первом рубеже") assert.False(t, after.IsHalted(), "и не потеряна") assert.Equal(t, record.Id, after.Id) } // Критерий приёмки 5, первая половина. Откладывание опроса не двигает время // входа в рубеж и не обнуляет отсчёт застревания: иначе запись, чью чужую // операцию опрашивают раз в несколько секунд, не достигла бы предела никогда. func TestPostponeKeepsStuckCountdown(t *testing.T) { entered := time.Date(2026, 8, 14, 10, 0, 0, 0, time.UTC) record := &entity.AudioRecord{State: entity.StateSubmitted, StateEnteredAt: entered} for range 100 { record.Postpone(time.Now().Add(5 * time.Second)) } assert.Equal(t, entered, record.StateEnteredAt, "сотня откладываний не двигает отсчёт") assert.Equal(t, entity.StateSubmitted, record.State, "и рубежа не трогает") assert.Nil(t, record.AcquisitionID, "захват при этом снят") assert.NotNil(t, record.DelayTime, "а пауза поставлена") } // Критерий приёмки 5, вторая половина. Запись, простоявшая в рубеже дольше // предела, останавливается признаком с причиной «застряла». func TestStuckRecordIsHalted(t *testing.T) { env := newPipelineEnv(t, &okMetaViewer{}, &okConverter{}) record := newRecord(t, env) enteredStateAt(t, env, record.Id, time.Now().Add(-2*testLimits.Own)) err := env.service.RunStep(t.Context()) var noop *contract.NoopJobError require.ErrorAs(t, err, &noop, "застрявшая запись шагу не отдаётся") after := readRecord(t, env, record.Id) require.True(t, after.IsHalted()) require.NotNil(t, after.HaltReason) assert.Equal(t, entity.HaltReasonStuck, *after.HaltReason, "причина названа") assert.Equal(t, entity.StateUploaded, after.State, "рубеж сохранён") } // Конечный рубеж под предел простоя не подпадает: стоять в нём запись будет // вечно по построению, и сторож остановил бы всякую доведённую запись. func TestDoneRecordIsNeverStuck(t *testing.T) { env := newPipelineEnv(t, &okMetaViewer{}, &okConverter{}) record := newRecord(t, env) drain(t, env, record.Id) enteredStateAt(t, env, record.Id, time.Now().Add(-10*24*time.Hour)) err := env.service.RunStep(t.Context()) var noop *contract.NoopJobError require.ErrorAs(t, err, &noop, "доведённая запись в работу не берётся") after := readRecord(t, env, record.Id) assert.Equal(t, entity.StateDone, after.State) assert.False(t, after.IsHalted(), "признака остановки у доведённой записи нет") } // Критерий приёмки 6. Держатель захвата отличим значением, а не занятостью // записи: шаг, чей захват достался другому, результата не пишет и отправителю // ничего не шлёт. func TestOnlyHolderWritesResult(t *testing.T) { env := newPipelineEnv(t, &okMetaViewer{}, &okConverter{}) record := newRecord(t, env) first, err := env.recordRepo.FindAndAcquire(entity.WorkingStages()) require.NoError(t, err) require.Equal(t, record.Id, first.ID) // Человек снял признак остановки — панель чистит признак захвата, — и запись // достаётся другому воркеру. expireAcquisition(t, env, record.Id) second, err := env.recordRepo.FindAndAcquire(entity.WorkingStages()) require.NoError(t, err) require.NotEqual(t, first.Holder, second.Holder, "признак нового захвата отличается") // Первый доходит до записи результата и получает отказ. stale := readRecord(t, env, record.Id) stale.MoveToState(entity.StateNormalized) err = env.recordRepo.Save(stale, first.Holder) var lost *contract.LostAcquisitionError require.ErrorAs(t, err, &lost, "потерявший захват не пишет") after := readRecord(t, env, record.Id) assert.Equal(t, entity.StateUploaded, after.State, "рубеж не сдвинут потерявшим захват") } // Критерий приёмки 7. Всякий способ вывести запись из работы оставляет причину // остановки: причин больше одной, и обязанность у них общая. Отправитель узнаёт // исход опросом готовности, а владелец сервиса — журналом событий записи. func TestEveryHaltReasonRecordsItsCause(t *testing.T) { reasons := []struct { name string halt func(t *testing.T, env *pipelineEnv, recordID string) }{ { name: "приговор шага", halt: func(t *testing.T, env *pipelineEnv, recordID string) { require.NoError(t, env.service.RunStep(t.Context())) }, }, { name: "застревание", halt: func(t *testing.T, env *pipelineEnv, recordID string) { enteredStateAt(t, env, recordID, time.Now().Add(-2*testLimits.Own)) var noop *contract.NoopJobError require.ErrorAs(t, env.service.RunStep(t.Context()), &noop) }, }, { name: "исчерпанные отказы", halt: func(t *testing.T, env *pipelineEnv, recordID string) { setAttempts(t, env, recordID, maxAttempts+1) var noop *contract.NoopJobError require.ErrorAs(t, env.service.RunStep(t.Context()), &noop) }, }, } for _, reason := range reasons { t.Run(reason.name, func(t *testing.T) { env := newPipelineEnv(t, &okMetaViewer{}, &failingConverter{}) record := newRecord(t, env) reason.halt(t, env, record.Id) after := readRecord(t, env, record.Id) require.True(t, after.IsHalted(), "запись остановлена") // Признак остановки и причина ставятся одним движением, поэтому // вторым утверждением берётся **журнал событий**: он пишется // отдельной строкой, отдельным сохранением, и упасть может сам по // себе. Прежде эту роль играл счёт ответов отправителю; ответы ушли // вместе с входом Telegram, и без замены проверка вывелась бы из // собственной предыдущей строки. 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, "строка журнала называет исход остановкой") assert.NotEmpty(t, events[0].GetString("outcome_text"), "и несёт причину") }) } } // Критерий приёмки 8. Повтор шага не создаёт второго приложения: пара «запись и // вид» уникальна, и вопрос «какой текст отдавать человеку» не становится // вопросом порядка записи. func TestRepeatedStoreKeepsSingleText(t *testing.T) { env := newPipelineEnv(t, &okMetaViewer{}, &okConverter{}) record := newRecord(t, env) first, err := env.repos.Texts.Put(record.Id, entity.TextKindTranscript, "первый разбор") require.NoError(t, err) second, err := env.repos.Texts.Put(record.Id, entity.TextKindTranscript, "второй разбор") require.NoError(t, err) assert.Equal(t, first.Id, second.Id, "строка та же, а не вторая") stored, err := env.repos.Texts.GetByID(first.Id) 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, "второго комплекта строк не завелось") } // Пауза растёт с числом отказов и упирается в потолок. 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 TestFailedStepSchedulesRetryWithGrowingDelay(t *testing.T) { env := newPipelineEnv(t, &okMetaViewer{}, &failingConverter{}) record := newRecord(t, env) // Ссылка переставляется на запись о файле без содержимого: шаг отказывает на // получении рабочей копии — то есть отказом, а не приговором записи. files, err := env.app.FindCollectionByNameOrId(migrations.FilesCollection) require.NoError(t, err) empty := core.NewRecord(files) empty.Set("location", entity.LocationLocal) empty.Set("size", 1) // Владелец обязателен и у файла: схема ничьих не принимает. 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)) require.Error(t, env.service.RunStep(t.Context())) after := readRecord(t, env, record.Id) require.Nil(t, after.AcquisitionID, "захват снят: запись пригодна к повтору") require.NotNil(t, after.DelayTime, "пауза поставлена") firstDelay := time.Until(*after.DelayTime) clearDelay(t, env, record.Id) require.Error(t, env.service.RunStep(t.Context())) after = readRecord(t, env, record.Id) require.NotNil(t, after.DelayTime) assert.Greater(t, time.Until(*after.DelayTime), firstDelay, "вторая пауза длиннее первой") } // Рабочая копия убирается на любом исходе, включая отказ. Забытая копия — это // шестичасовая запись во временном каталоге, и узнать о ней неоткуда. func TestWorkFileRemovedAfterIntakeFailure(t *testing.T) { tempDir := t.TempDir() t.Setenv("TMPDIR", tempDir) env := newPipelineEnv(t, &failingMetaViewer{}, &failingConverter{}) _, err := env.service.CreateJobFromApi(t.Context(), strings.NewReader("запись"), "sample.mp3", newOwner(t, env.app)) 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(t.Context(), strings.NewReader("запись"), "sample.mp3", newOwner(t, env.app)) require.NoError(t, err) leftovers, err := filepath.Glob(filepath.Join(tempDir, "transcriber-*")) require.NoError(t, err) assert.Empty(t, leftovers, "рабочей копии после успеха не остаётся") } // Уборка проверяется и на шаге приведения: репозиторий даёт единственный способ // убрать копию, но зовёт его шаг. Копий здесь две — исходник и результат. func TestWorkFilesRemovedAfterConversionFailure(t *testing.T) { tempDir := t.TempDir() t.Setenv("TMPDIR", tempDir) env := newPipelineEnv(t, &okMetaViewer{}, &failingConverter{}) newRecord(t, env) leftovers, err := filepath.Glob(filepath.Join(tempDir, "transcriber-*")) require.NoError(t, err) require.Empty(t, leftovers, "приём убрал свою рабочую копию") require.NoError(t, env.service.RunStep(t.Context())) leftovers, err = filepath.Glob(filepath.Join(tempDir, "transcriber-*")) require.NoError(t, err) assert.Empty(t, leftovers, "ни исходной копии, ни копии под результат не осталось") } // Запись не остаётся ссылающейся на файл, которого нет: ссылка ставится только // после того, как запись о файле существует. func TestRecordNeverPointsToMissingFile(t *testing.T) { env := newPipelineEnv(t, &okMetaViewer{}, &failingConverter{}) record := newRecord(t, env) require.NoError(t, env.service.RunStep(t.Context())) after := readRecord(t, env, record.Id) assert.True(t, after.IsHalted(), "приговор шага остановил запись") require.NotNil(t, after.OriginalFileID) assert.Nil(t, after.NormalizedFileID, "ссылки на несозданный результат не появилось") file, err := env.fileRepo.GetByID(*after.OriginalFileID) require.NoError(t, err, "ссылка ведёт на существующую запись о файле") assert.NotEmpty(t, file.FileName) } // Содержимое доезжает до хранилища целиком и читается обратно тем же. func TestStoredContentSurvivesRoundTrip(t *testing.T) { env := newPipelineEnv(t, &okMetaViewer{}, &failingConverter{}) content := strings.Repeat("запись ", 1000) record, err := env.service.CreateJobFromApi(t.Context(), strings.NewReader(content), "sample.mp3", newOwner(t, env.app)) require.NoError(t, err) require.NotNil(t, record.OriginalFileID) reader, err := env.fileRepo.Open(*record.OriginalFileID) 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(*record.OriginalFileID) require.NoError(t, err) assert.Equal(t, int64(len(content)), file.Size) assert.Equal(t, "mp3", file.Format, "формат копии записан") } // Рабочая копия хранимого файла отдаётся именем на диске — так её получают шаги, // отдающие файл внешней программе. 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)) require.NoError(t, err) require.NotNil(t, record.OriginalFileID) work, err := env.fileRepo.Localize(*record.OriginalFileID) 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), "закрытая копия убрана") } // newOwner заводит учётную запись и отдаёт её идентификатор. // // Владелец — связь с коллекцией пользователей, и хранилище проверяет, что такая // запись есть: выдуманный идентификатор запись завести не даст. func newOwner(t *testing.T, app core.App) string { t.Helper() users, err := app.FindCollectionByNameOrId(migrations.UsersCollection) require.NoError(t, err) record := core.NewRecord(users) 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 } // Остановка приговором шага засчитывается **отказом**, а не успехом, и пишет в // журнал записи одну строку, а не две. // // Пока исход шага выводился из «шаг не вернул ошибку», остановленная запись // получала событием `done` вслед за `halted` и растила счётчик успехов — то есть // единственный канал владельца молчал ровно там, где запись встала. func TestHaltIsCountedAsFailureAndLoggedOnce(t *testing.T) { env := newPipelineEnv(t, &okMetaViewer{}, &failingConverter{}) beforeOk := stageCount(t, entity.StateUploaded, "false") beforeErr := stageCount(t, entity.StateUploaded, "true") record := newRecord(t, env) require.NoError(t, env.service.RunStep(t.Context())) after := readRecord(t, env, record.Id) require.True(t, after.IsHalted(), "приговор шага остановил запись") assert.InDelta(t, beforeErr+1, stageCount(t, entity.StateUploaded, "true"), 0, "остановка засчитана отказом") assert.InDelta(t, beforeOk, stageCount(t, entity.StateUploaded, "false"), 0, "и успехом не засчитана") events := recordEvents(t, env, record.Id) require.Len(t, events, 1, "одна строка журнала, а не две") assert.Equal(t, entity.EventOutcomeFailed, events[0].GetString("outcome"), "исход назван приговором, а не сделанной работой") } // Откладывание опроса в журнал событий не пишется: часовое ожидание чужой // операции дало бы там сотни строк ни о чём, и журнал, заведённый для человека, // стал бы нечитаемым ровно у долгой записи. func TestPostponeWritesNoEvent(t *testing.T) { rec := &countingRecognizer{} rec.inProgress.Store(5) env := newPipelineEnvWith(t, &okMetaViewer{}, &okConverter{}, rec) record := newRecord(t, env) require.NoError(t, env.service.RunStep(t.Context())) // приведение require.NoError(t, env.service.RunStep(t.Context())) // отправка require.Equal(t, entity.StateSubmitted, readRecord(t, env, record.Id).State) before := len(recordEvents(t, env, record.Id)) for range 5 { clearDelay(t, env, record.Id) require.NoError(t, env.service.RunStep(t.Context())) } assert.Len(t, recordEvents(t, env, record.Id), before, "пять откладываний не оставили в журнале ни строки") } // recordEvents читает журнал событий одной записи в порядке заведения. func recordEvents(t *testing.T, env *pipelineEnv, recordID string) []*core.Record { t.Helper() all, err := env.app.FindAllRecords(migrations.RecordEventsCollection) require.NoError(t, err) var own []*core.Record for _, event := range all { if event.GetString("record") == recordID { own = append(own, event) } } return own }