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/migrations" "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(migrations.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, "отправителю ничего не отправлено") }