package service import ( "context" "errors" "testing" "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" ) // killedConverter ведёт себя как настоящий `ffmpeg`, убитый по контексту: // дожидается отмены и отдаёт отказ, который **не** несёт `context.Canceled`. // Это не упрощение, а суть проверки: `exec` отдаёт `*exec.ExitError` с текстом // «signal: killed», и отличить остановку от негодной записи по самой ошибке // нельзя — только по контексту шага. type killedConverter struct { cancel func() } func (c *killedConverter) Convert(ctx context.Context, _, _ string) error { c.cancel() <-ctx.Done() return errors.New("ffmpeg conversion failed: signal: killed") } // Остановка сервиса посреди конвертации не выносит записи приговора: задача // остаётся пригодной к повтору, попытку не тратит и отправителю о сбое, // которого не было, не сообщает. Прежде любой отказ `Convert` уводил задачу в // терминальное `failed`, откуда её возвращает только владелец правкой в панели. func TestShutdownDuringConversionKeepsJobRetryable(t *testing.T) { ctx, cancel := context.WithCancel(t.Context()) defer cancel() converter := &killedConverter{cancel: cancel} env := newPipelineEnv(t, &okMetaViewer{}, converter) job := newTelegramJob(t, env) err := env.service.FindAndRunConversionJob(ctx) require.Error(t, err, "шаг обязан сообщить об обрыве наверх") require.ErrorIs(t, err, context.Canceled, "обрыв узнаётся по смыслу, а не по тексту") after, err := env.jobRepo.GetByID(job.Id) require.NoError(t, err) assert.Equal(t, entity.StateCreated, after.State, "задача осталась на повтор, а не похоронена") assert.Nil(t, after.AcquisitionID, "захват снят: задачу возьмёт следующий прогон") assert.Equal(t, 0, after.Attempts, "остановка попытки не тратит") assert.Nil(t, after.ErrorText, "приговора не выносили") assert.Empty(t, env.sender.messages, "отправителю о несуществующем сбое не сообщают") } // Задача, которую шаг не успел взять, потому что нас уже остановили, остаётся // нетронутой: захват не случился, попытка не потрачена. func TestShutdownBeforeStepLeavesJobUntouched(t *testing.T) { env := newPipelineEnv(t, &okMetaViewer{}, &failingConverter{}) job := newTelegramJob(t, env) ctx, cancel := context.WithCancel(t.Context()) cancel() err := env.service.FindAndRunConversionJob(ctx) // Исход «шаг не сделал ничего» — это `NoopJobError`: воркер не пишет о нём // владельцу и не считает его в метрику. require.Error(t, err) var noop *contract.NoopJobError require.ErrorAs(t, err, &noop) after, err := env.jobRepo.GetByID(job.Id) require.NoError(t, err) assert.Equal(t, entity.StateCreated, after.State) assert.Equal(t, 0, after.Attempts, "захвата не было — попытке взяться неоткуда") record, err := env.app.FindRecordById(migrations.JobsCollection, job.Id) require.NoError(t, err) assert.Empty(t, record.GetString("acquisition_id")) }