Files
transcriber/internal/service/recognition_test.go
T
av f494dcb83e Отмена доходит до внешнего собеседника, а токен не покидает единой точки
- контекст проложен от воркера и обоих входов до внешних вызовов: ffmpeg и
  ffprobe заводятся через exec.CommandContext, SpeechKit и Object Storage
  принимают ctx вместо context.Background, скачивание записи идёт запросом с
  контекстом. Прежде остановка сервиса не доходила до чужой работы вовсе
- прерванный шаг приговора не выносит: убитый по контексту ffmpeg отдаёт
  «signal: killed», от настоящего отказа неотличимо ни типом, ни errors.Is, и
  различает их только ctx.Err(). Задача остаётся на повтор, попытку не тратит и
  отправителю о несуществующем сбое не сообщает; воркер не считает остановку
  отказом, а задача не забирается вовсе, если нас уже остановили
- клиента Bot API заводит единая точка internal/adapter/telegram: токен стоит в
  пути каждого обращения, а http.Client кладёт адрес в *url.Error целиком.
  Чистка на месте употребления закрывала один вызов из пяти — теперь свой Do
  чистит отказ, подменённый логгер вычищает токен из строк самой библиотеки, а
  транспорт бота токена не получает вовсе
- принятие операции распознавания защищено от отмены своим пределом: SpeechKit
  мог её принять и начать считать деньги, а потерянный идентификатор заставил
  бы повтор оплатить ту же запись второй раз
- приём по HTTP доводит запись до задачи независимо от отправителя: на
  контексте запроса один обрыв соединения терял полностью загруженную запись
- ответ Telegram с не-2xx кодом больше не становится записью: прежде тело
  отказа доезжало до хранилища и умирало на ffprobe, уводя диагностику
2026-08-13 10:27:54 +03:00

269 lines
12 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
package service
import (
"context"
"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(_ context.Context, 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(context.Context, string) (string, error) {
return r.text, nil
}
func (r *scriptedRecognizer) CheckRecognitionStatus(context.Context, 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(t.Context()))
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(t.Context()))
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(t.Context()))
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(t.Context()))
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(t.Context()))
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(t.Context()))
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(t.Context()))
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(t.Context(), 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, "отправителю ничего не отправлено")
}