Files
transcriber/internal/service/pipeline_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

365 lines
16 KiB
Go

package service
import (
"context"
"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/adapter/repo/pocketbase/migrations"
"git.vakhrushev.me/av/transcriber/internal/contract"
"git.vakhrushev.me/av/transcriber/internal/entity"
)
// Проверки конвейера идут против настоящего хранилища: захват, число попыток и
// переход в «мертва» держатся на запросе, и подставной репозиторий проверял бы
// собственную заглушку, а не то, что делает база.
// 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("запись не читается")
}
// 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.DiscardHandler),
)
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(t.Context(), 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(migrations.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(migrations.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(t.Context())
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(t.Context()))
record, err := env.app.FindRecordById(migrations.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(migrations.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(t.Context()))
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(t.Context()))
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(t.Context(), 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(t.Context(), 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(t.Context()))
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(t.Context(), 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(t.Context(), 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(t.Context()))
leftovers, err = filepath.Glob(filepath.Join(tempDir, "transcriber-*"))
require.NoError(t, err)
assert.Empty(t, leftovers, "ни исходной копии, ни копии под результат не осталось")
}