Отмена доходит до внешнего собеседника, а токен не покидает единой точки
- контекст проложен от воркера и обоих входов до внешних вызовов: 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, уводя диагностику
This commit is contained in:
@@ -1,6 +1,7 @@
|
||||
package service
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"io"
|
||||
"log/slog"
|
||||
@@ -29,19 +30,19 @@ import (
|
||||
// failingConverter отказывает на каждой попытке.
|
||||
type failingConverter struct{}
|
||||
|
||||
func (c *failingConverter) Convert(string, string) error {
|
||||
func (c *failingConverter) Convert(context.Context, string, string) error {
|
||||
return errors.New("конвертация не удалась")
|
||||
}
|
||||
|
||||
type okMetaViewer struct{}
|
||||
|
||||
func (m *okMetaViewer) GetInfo(string) (*contract.AudioInfo, error) {
|
||||
func (m *okMetaViewer) GetInfo(context.Context, string) (*contract.AudioInfo, error) {
|
||||
return &contract.AudioInfo{Seconds: 1}, nil
|
||||
}
|
||||
|
||||
type failingMetaViewer struct{}
|
||||
|
||||
func (m *failingMetaViewer) GetInfo(string) (*contract.AudioInfo, error) {
|
||||
func (m *failingMetaViewer) GetInfo(context.Context, string) (*contract.AudioInfo, error) {
|
||||
return nil, errors.New("запись не читается")
|
||||
}
|
||||
|
||||
@@ -101,7 +102,7 @@ func newTelegramJob(t *testing.T, env *pipelineEnv) *entity.TranscribeJob {
|
||||
t.Helper()
|
||||
|
||||
chatId := int64(100)
|
||||
job, err := env.service.CreateJobFromTelegram(strings.NewReader("запись"), "voice.ogg", chatId, 1)
|
||||
job, err := env.service.CreateJobFromTelegram(t.Context(), strings.NewReader("запись"), "voice.ogg", chatId, 1)
|
||||
require.NoError(t, err)
|
||||
return job
|
||||
}
|
||||
@@ -149,7 +150,7 @@ func TestJobDiesAfterAttemptLimit(t *testing.T) {
|
||||
rotAcquisition(t, env, job.Id)
|
||||
|
||||
// Следующий захват видит перебор и хоронит задачу.
|
||||
err := env.service.FindAndRunConversionJob()
|
||||
err := env.service.FindAndRunConversionJob(t.Context())
|
||||
|
||||
var noop *contract.NoopJobError
|
||||
require.ErrorAs(t, err, &noop, "мёртвая задача шагу не отдаётся")
|
||||
@@ -178,7 +179,7 @@ func TestDeadJobReturnsAfterStateEdit(t *testing.T) {
|
||||
_, err := env.jobRepo.FindAndAcquire(entity.StateCreated, "holder", time.Now().Add(time.Hour))
|
||||
require.NoError(t, err)
|
||||
}
|
||||
require.Error(t, env.service.FindAndRunConversionJob())
|
||||
require.Error(t, env.service.FindAndRunConversionJob(t.Context()))
|
||||
|
||||
record, err := env.app.FindRecordById(migrations.JobsCollection, job.Id)
|
||||
require.NoError(t, err)
|
||||
@@ -209,7 +210,7 @@ func TestFailedStepSchedulesRetryWithGrowingDelay(t *testing.T) {
|
||||
require.NoError(t, env.app.Save(record))
|
||||
|
||||
// Первый отказ.
|
||||
require.Error(t, env.service.FindAndRunConversionJob())
|
||||
require.Error(t, env.service.FindAndRunConversionJob(t.Context()))
|
||||
|
||||
after, err := env.jobRepo.GetByID(job.Id)
|
||||
require.NoError(t, err)
|
||||
@@ -220,7 +221,7 @@ func TestFailedStepSchedulesRetryWithGrowingDelay(t *testing.T) {
|
||||
|
||||
// Второй отказ — с той же задачи, пауза снята вручную.
|
||||
clearDelay(t, env, job.Id)
|
||||
require.Error(t, env.service.FindAndRunConversionJob())
|
||||
require.Error(t, env.service.FindAndRunConversionJob(t.Context()))
|
||||
|
||||
after, err = env.jobRepo.GetByID(job.Id)
|
||||
require.NoError(t, err)
|
||||
@@ -247,7 +248,7 @@ func TestWorkFileRemovedAfterIntakeFailure(t *testing.T) {
|
||||
|
||||
env := newPipelineEnv(t, &failingMetaViewer{}, &failingConverter{})
|
||||
|
||||
_, err := env.service.CreateJobFromApi(strings.NewReader("запись"), "sample.mp3")
|
||||
_, err := env.service.CreateJobFromApi(t.Context(), strings.NewReader("запись"), "sample.mp3")
|
||||
require.Error(t, err, "отказ источника метаданных роняет приём")
|
||||
|
||||
leftovers, err := filepath.Glob(filepath.Join(tempDir, "transcriber-*"))
|
||||
@@ -263,7 +264,7 @@ func TestWorkFileRemovedAfterSuccessfulIntake(t *testing.T) {
|
||||
|
||||
env := newPipelineEnv(t, &okMetaViewer{}, &failingConverter{})
|
||||
|
||||
_, err := env.service.CreateJobFromApi(strings.NewReader("запись"), "sample.mp3")
|
||||
_, err := env.service.CreateJobFromApi(t.Context(), strings.NewReader("запись"), "sample.mp3")
|
||||
require.NoError(t, err)
|
||||
|
||||
leftovers, err := filepath.Glob(filepath.Join(tempDir, "transcriber-*"))
|
||||
@@ -280,7 +281,7 @@ func TestJobNeverPointsToMissingFile(t *testing.T) {
|
||||
|
||||
// Конвертация отказывает — задача уходит в `failed`, но ссылка остаётся на
|
||||
// исходную запись, а не на несозданный результат.
|
||||
require.NoError(t, env.service.FindAndRunConversionJob())
|
||||
require.NoError(t, env.service.FindAndRunConversionJob(t.Context()))
|
||||
|
||||
after, err := env.jobRepo.GetByID(job.Id)
|
||||
require.NoError(t, err)
|
||||
@@ -298,7 +299,7 @@ func TestStoredContentSurvivesRoundTrip(t *testing.T) {
|
||||
|
||||
content := strings.Repeat("запись ", 1000)
|
||||
|
||||
job, err := env.service.CreateJobFromApi(strings.NewReader(content), "sample.mp3")
|
||||
job, err := env.service.CreateJobFromApi(t.Context(), strings.NewReader(content), "sample.mp3")
|
||||
require.NoError(t, err)
|
||||
require.NotNil(t, job.FileID)
|
||||
|
||||
@@ -321,7 +322,7 @@ func TestStoredContentSurvivesRoundTrip(t *testing.T) {
|
||||
func TestLocalizeGivesReadableCopy(t *testing.T) {
|
||||
env := newPipelineEnv(t, &okMetaViewer{}, &failingConverter{})
|
||||
|
||||
job, err := env.service.CreateJobFromApi(strings.NewReader("содержимое"), "sample.mp3")
|
||||
job, err := env.service.CreateJobFromApi(t.Context(), strings.NewReader("содержимое"), "sample.mp3")
|
||||
require.NoError(t, err)
|
||||
require.NotNil(t, job.FileID)
|
||||
|
||||
@@ -355,7 +356,7 @@ func TestWorkFilesRemovedAfterConversionFailure(t *testing.T) {
|
||||
require.Empty(t, leftovers, "приём убрал свою рабочую копию")
|
||||
|
||||
// Конвертация отказывает — задача уходит в `failed`, копии убраны.
|
||||
require.NoError(t, env.service.FindAndRunConversionJob())
|
||||
require.NoError(t, env.service.FindAndRunConversionJob(t.Context()))
|
||||
|
||||
leftovers, err = filepath.Glob(filepath.Join(tempDir, "transcriber-*"))
|
||||
require.NoError(t, err)
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
package service
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"io"
|
||||
"strings"
|
||||
@@ -32,7 +33,7 @@ type scriptedRecognizer struct {
|
||||
lastObjectKey string
|
||||
}
|
||||
|
||||
func (r *scriptedRecognizer) Recognize(file io.Reader, fileName string) (string, error) {
|
||||
func (r *scriptedRecognizer) Recognize(_ context.Context, file io.Reader, fileName string) (string, error) {
|
||||
r.recognizeCalls++
|
||||
r.lastObjectKey = fileName
|
||||
if r.recognizeErr != nil {
|
||||
@@ -45,11 +46,11 @@ func (r *scriptedRecognizer) Recognize(file io.Reader, fileName string) (string,
|
||||
return "operation-id", nil
|
||||
}
|
||||
|
||||
func (r *scriptedRecognizer) GetRecognitionText(string) (string, error) {
|
||||
func (r *scriptedRecognizer) GetRecognitionText(context.Context, string) (string, error) {
|
||||
return r.text, nil
|
||||
}
|
||||
|
||||
func (r *scriptedRecognizer) CheckRecognitionStatus(string) (*entity.RecognitionResult, error) {
|
||||
func (r *scriptedRecognizer) CheckRecognitionStatus(context.Context, string) (*entity.RecognitionResult, error) {
|
||||
return r.result, nil
|
||||
}
|
||||
|
||||
@@ -92,7 +93,7 @@ func TestTranscribeJobHandsRecordOverAndMovesOn(t *testing.T) {
|
||||
rec := &scriptedRecognizer{result: entity.NewInProgressResult()}
|
||||
svc := withRecognizer(env, rec)
|
||||
|
||||
require.NoError(t, svc.FindAndRunTranscribeJob())
|
||||
require.NoError(t, svc.FindAndRunTranscribeJob(t.Context()))
|
||||
|
||||
assert.Equal(t, 1, rec.recognizeCalls, "содержимое отдано распознавателю")
|
||||
assert.NotEmpty(t, rec.lastObjectKey, "ключ объекта назван")
|
||||
@@ -119,7 +120,7 @@ func TestTranscribeJobKeepsJobRetryableOnRecognizerFailure(t *testing.T) {
|
||||
rec := &scriptedRecognizer{recognizeErr: errors.New("распознаватель недоступен")}
|
||||
svc := withRecognizer(env, rec)
|
||||
|
||||
require.Error(t, svc.FindAndRunTranscribeJob())
|
||||
require.Error(t, svc.FindAndRunTranscribeJob(t.Context()))
|
||||
|
||||
after, err := env.jobRepo.GetByID(job.Id)
|
||||
require.NoError(t, err)
|
||||
@@ -134,7 +135,7 @@ func transcribingJob(t *testing.T, env *pipelineEnv, rec contract.AudioRecognize
|
||||
t.Helper()
|
||||
|
||||
job := convertedJob(t, env)
|
||||
require.NoError(t, withRecognizer(env, rec).FindAndRunTranscribeJob())
|
||||
require.NoError(t, withRecognizer(env, rec).FindAndRunTranscribeJob(t.Context()))
|
||||
|
||||
clearDelay(t, env, job.Id)
|
||||
return job
|
||||
@@ -150,7 +151,7 @@ func TestCheckJobWaitsWithoutSpendingAttempts(t *testing.T) {
|
||||
svc := withRecognizer(env, rec)
|
||||
|
||||
for i := 0; i < 3; i++ {
|
||||
require.NoError(t, svc.FindAndRunTranscribeCheckJob())
|
||||
require.NoError(t, svc.FindAndRunTranscribeCheckJob(t.Context()))
|
||||
|
||||
after, err := env.jobRepo.GetByID(job.Id)
|
||||
require.NoError(t, err)
|
||||
@@ -174,7 +175,7 @@ func TestCheckJobFailsJobAndTellsSender(t *testing.T) {
|
||||
rec.result = entity.NewFailedResult("операция отклонена")
|
||||
svc := withRecognizer(env, rec)
|
||||
|
||||
require.NoError(t, svc.FindAndRunTranscribeCheckJob())
|
||||
require.NoError(t, svc.FindAndRunTranscribeCheckJob(t.Context()))
|
||||
|
||||
after, err := env.jobRepo.GetByID(job.Id)
|
||||
require.NoError(t, err)
|
||||
@@ -196,7 +197,7 @@ func TestCheckJobCompletesAndAnswersOnce(t *testing.T) {
|
||||
rec.text = "расшифровка записи"
|
||||
svc := withRecognizer(env, rec)
|
||||
|
||||
require.NoError(t, svc.FindAndRunTranscribeCheckJob())
|
||||
require.NoError(t, svc.FindAndRunTranscribeCheckJob(t.Context()))
|
||||
|
||||
after, err := env.jobRepo.GetByID(job.Id)
|
||||
require.NoError(t, err)
|
||||
@@ -224,7 +225,7 @@ func TestCheckJobCompletesEmptyTextWithExplanation(t *testing.T) {
|
||||
rec.text = ""
|
||||
svc := withRecognizer(env, rec)
|
||||
|
||||
require.NoError(t, svc.FindAndRunTranscribeCheckJob())
|
||||
require.NoError(t, svc.FindAndRunTranscribeCheckJob(t.Context()))
|
||||
|
||||
after, err := env.jobRepo.GetByID(job.Id)
|
||||
require.NoError(t, err)
|
||||
@@ -255,7 +256,7 @@ func TestCheckJobWritesNothingWhenAcquisitionLost(t *testing.T) {
|
||||
require.NoError(t, env.app.Save(record))
|
||||
|
||||
svc := withRecognizer(env, rec)
|
||||
err = svc.checkTranscribeJob(acquired, "mine")
|
||||
err = svc.checkTranscribeJob(t.Context(), acquired, "mine")
|
||||
|
||||
var lost *contract.LostAcquisitionError
|
||||
require.ErrorAs(t, err, &lost)
|
||||
|
||||
@@ -0,0 +1,83 @@
|
||||
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"))
|
||||
}
|
||||
@@ -1,6 +1,7 @@
|
||||
package service
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
@@ -75,7 +76,7 @@ func NewTranscribeService(
|
||||
}
|
||||
}
|
||||
|
||||
func (s *TranscribeService) CreateJobFromTelegram(file io.Reader, fileName string, chatId int64, replyMsgId int) (*entity.TranscribeJob, error) {
|
||||
func (s *TranscribeService) CreateJobFromTelegram(ctx context.Context, file io.Reader, fileName string, chatId int64, replyMsgId int) (*entity.TranscribeJob, error) {
|
||||
job := &entity.TranscribeJob{
|
||||
State: entity.StateCreated,
|
||||
Source: entity.SourceTelegram,
|
||||
@@ -83,19 +84,19 @@ func (s *TranscribeService) CreateJobFromTelegram(file io.Reader, fileName strin
|
||||
TgReplyMessageId: &replyMsgId,
|
||||
}
|
||||
|
||||
return s.createTranscribeJob(job, file, fileName)
|
||||
return s.createTranscribeJob(ctx, job, file, fileName)
|
||||
}
|
||||
|
||||
func (s *TranscribeService) CreateJobFromApi(file io.Reader, fileName string) (*entity.TranscribeJob, error) {
|
||||
func (s *TranscribeService) CreateJobFromApi(ctx context.Context, file io.Reader, fileName string) (*entity.TranscribeJob, error) {
|
||||
job := &entity.TranscribeJob{
|
||||
State: entity.StateCreated,
|
||||
Source: entity.SourceApi,
|
||||
}
|
||||
|
||||
return s.createTranscribeJob(job, file, fileName)
|
||||
return s.createTranscribeJob(ctx, job, file, fileName)
|
||||
}
|
||||
|
||||
func (s *TranscribeService) createTranscribeJob(job *entity.TranscribeJob, file io.Reader, fileName string) (*entity.TranscribeJob, error) {
|
||||
func (s *TranscribeService) createTranscribeJob(ctx context.Context, job *entity.TranscribeJob, file io.Reader, fileName string) (*entity.TranscribeJob, error) {
|
||||
// Определяем расширение файла
|
||||
ext := filepath.Ext(fileName)
|
||||
if ext == "" {
|
||||
@@ -122,7 +123,7 @@ func (s *TranscribeService) createTranscribeJob(job *entity.TranscribeJob, file
|
||||
// строка журнала вместе с идентификатором записи собрала бы её целиком.
|
||||
s.logger.Info("Creating transcribe job", "file_ext", ext)
|
||||
|
||||
info, err := s.metaviewer.GetInfo(work.Path())
|
||||
info, err := s.metaviewer.GetInfo(ctx, work.Path())
|
||||
if err != nil {
|
||||
s.logger.Error("Failed to get file info", "error", err, "file_ext", ext)
|
||||
return nil, err
|
||||
@@ -160,28 +161,41 @@ func (s *TranscribeService) createTranscribeJob(job *entity.TranscribeJob, file
|
||||
return job, nil
|
||||
}
|
||||
|
||||
func (s *TranscribeService) FindAndRunConversionJob() error {
|
||||
return s.runStep(entity.StateCreated, conversionAcquireTimeout, s.convertJob)
|
||||
func (s *TranscribeService) FindAndRunConversionJob(ctx context.Context) error {
|
||||
return s.runStep(ctx, entity.StateCreated, conversionAcquireTimeout, s.convertJob)
|
||||
}
|
||||
|
||||
func (s *TranscribeService) FindAndRunTranscribeJob() error {
|
||||
return s.runStep(entity.StateConverted, transcribeAcquireTimeout, s.transcribeJob)
|
||||
func (s *TranscribeService) FindAndRunTranscribeJob(ctx context.Context) error {
|
||||
return s.runStep(ctx, entity.StateConverted, transcribeAcquireTimeout, s.transcribeJob)
|
||||
}
|
||||
|
||||
func (s *TranscribeService) FindAndRunTranscribeCheckJob() error {
|
||||
return s.runStep(entity.StateTranscribe, checkAcquireTimeout, s.checkTranscribeJob)
|
||||
func (s *TranscribeService) FindAndRunTranscribeCheckJob(ctx context.Context) error {
|
||||
return s.runStep(ctx, entity.StateTranscribe, checkAcquireTimeout, s.checkTranscribeJob)
|
||||
}
|
||||
|
||||
// runStep забирает задачу и отдаёт её шагу. Отказ шага не оставляет задачу
|
||||
// захваченной до конца срока: захват снимается, и задача ждёт нарастающую паузу
|
||||
// — иначе повтор наступал бы через восемь часов, а не через секунду.
|
||||
func (s *TranscribeService) runStep(state string, expiration time.Duration, step func(job *entity.TranscribeJob, holder string) error) error {
|
||||
//
|
||||
// Контекст доходит до шага, а через него — до внешнего собеседника: остановка
|
||||
// сервиса убивает `ffmpeg` и обрывает запрос к распознаванию. Прерванный шаг
|
||||
// приговора не выносит: задача остаётся пригодной к повтору, попытки не тратит
|
||||
// и отправителю о несуществующем сбое не сообщает — исход остановки отличается
|
||||
// от исхода отказа на каждом шаге.
|
||||
func (s *TranscribeService) runStep(ctx context.Context, state string, expiration time.Duration, step func(ctx context.Context, job *entity.TranscribeJob, holder string) error) error {
|
||||
// Нас уже остановили — задачу не забираем: захват стоил бы ей попытки, а
|
||||
// работы всё равно не будет. Исход «шаг не сделал ничего» — это `NoopJobError`
|
||||
// по смыслу, и он же не поднимает уровень и не считается в метрику.
|
||||
if ctx.Err() != nil {
|
||||
return &contract.NoopJobError{State: state}
|
||||
}
|
||||
|
||||
job, holder, err := s.findJob(state, expiration)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
if err := step(job, holder); err != nil {
|
||||
if err := step(ctx, job, holder); err != nil {
|
||||
s.scheduleRetry(job, holder, err)
|
||||
return err
|
||||
}
|
||||
@@ -189,7 +203,7 @@ func (s *TranscribeService) runStep(state string, expiration time.Duration, step
|
||||
return nil
|
||||
}
|
||||
|
||||
func (s *TranscribeService) convertJob(job *entity.TranscribeJob, holder string) error {
|
||||
func (s *TranscribeService) convertJob(ctx context.Context, job *entity.TranscribeJob, holder string) error {
|
||||
s.logger.Info("Starting conversion job", "job_id", job.Id)
|
||||
|
||||
if job.FileID == nil {
|
||||
@@ -227,13 +241,27 @@ func (s *TranscribeService) convertJob(job *entity.TranscribeJob, holder string)
|
||||
|
||||
// Измеряем время конвертации
|
||||
startTime := clock.Start()
|
||||
err = s.converter.Convert(src.Path(), dest.Path())
|
||||
err = s.converter.Convert(ctx, src.Path(), dest.Path())
|
||||
conversionDuration := time.Since(startTime)
|
||||
|
||||
// Записываем метрику времени конвертации
|
||||
metrics.ObserveConversionDuration(srcExt, "ogg", err != nil, conversionDuration.Seconds())
|
||||
|
||||
if err != nil {
|
||||
// Остановка сервиса — не приговор записи. Убитый по контексту `ffmpeg`
|
||||
// отдаёт `signal: killed`, и от настоящего отказа конвертации
|
||||
// (`exit status N`) эта ошибка неотличима ни типом, ни `errors.Is`:
|
||||
// различает их только контекст шага. Без этой развилки каждый деплой
|
||||
// хоронил бы конвертируемую запись в `failed` — состояние терминальное,
|
||||
// и вернуть её оттуда может только владелец правкой в панели, — да ещё
|
||||
// и сообщал бы отправителю о сбое, которого не было.
|
||||
if ctxErr := ctx.Err(); ctxErr != nil {
|
||||
s.logger.Info("File conversion interrupted by shutdown",
|
||||
"job_id", job.Id,
|
||||
"duration", conversionDuration)
|
||||
return fmt.Errorf("conversion interrupted: %w", ctxErr)
|
||||
}
|
||||
|
||||
s.logger.Error("File conversion failed",
|
||||
"error", err,
|
||||
"job_id", job.Id,
|
||||
@@ -276,7 +304,7 @@ func (s *TranscribeService) convertJob(job *entity.TranscribeJob, holder string)
|
||||
return nil
|
||||
}
|
||||
|
||||
func (s *TranscribeService) transcribeJob(job *entity.TranscribeJob, holder string) error {
|
||||
func (s *TranscribeService) transcribeJob(ctx context.Context, job *entity.TranscribeJob, holder string) error {
|
||||
s.logger.Info("Starting transcribe job", "job_id", job.Id)
|
||||
|
||||
if job.FileID == nil {
|
||||
@@ -304,8 +332,12 @@ func (s *TranscribeService) transcribeJob(job *entity.TranscribeJob, holder stri
|
||||
s.logger.Info("Starting recognition", "job_id", job.Id, "file_id", *job.FileID)
|
||||
|
||||
// Запускаем асинхронное распознавание
|
||||
operationID, err := s.recognizer.Recognize(content, fileRecord.FileName)
|
||||
operationID, err := s.recognizer.Recognize(ctx, content, fileRecord.FileName)
|
||||
if err != nil {
|
||||
if ctxErr := ctx.Err(); ctxErr != nil {
|
||||
s.logger.Info("Recognition interrupted by shutdown", "job_id", job.Id)
|
||||
return fmt.Errorf("recognition interrupted: %w", ctxErr)
|
||||
}
|
||||
s.logger.Error("Failed to start recognition", "error", err, "job_id", job.Id)
|
||||
return err
|
||||
}
|
||||
@@ -335,7 +367,7 @@ func (s *TranscribeService) transcribeJob(job *entity.TranscribeJob, holder stri
|
||||
return nil
|
||||
}
|
||||
|
||||
func (s *TranscribeService) checkTranscribeJob(job *entity.TranscribeJob, holder string) error {
|
||||
func (s *TranscribeService) checkTranscribeJob(ctx context.Context, job *entity.TranscribeJob, holder string) error {
|
||||
if job.RecognitionOpID == nil {
|
||||
s.logger.Error("Recognition operation ID not found", "job_id", job.Id)
|
||||
return fmt.Errorf("recognition opId not found for job: %s", job.Id)
|
||||
@@ -345,8 +377,12 @@ func (s *TranscribeService) checkTranscribeJob(job *entity.TranscribeJob, holder
|
||||
|
||||
// Проверяем статус операции
|
||||
s.logger.Info("Checking operation status", "job_id", job.Id, "operation_id", opId)
|
||||
recResult, err := s.recognizer.CheckRecognitionStatus(opId)
|
||||
recResult, err := s.recognizer.CheckRecognitionStatus(ctx, opId)
|
||||
if err != nil {
|
||||
if ctxErr := ctx.Err(); ctxErr != nil {
|
||||
s.logger.Info("Status check interrupted by shutdown", "job_id", job.Id)
|
||||
return fmt.Errorf("status check interrupted: %w", ctxErr)
|
||||
}
|
||||
s.logger.Error("Failed to check recognition status", "error", err, "operation_id", opId)
|
||||
return err
|
||||
}
|
||||
@@ -375,8 +411,12 @@ func (s *TranscribeService) checkTranscribeJob(job *entity.TranscribeJob, holder
|
||||
}
|
||||
|
||||
// Операция завершена, получаем результат
|
||||
transcriptionText, err := s.recognizer.GetRecognitionText(opId)
|
||||
transcriptionText, err := s.recognizer.GetRecognitionText(ctx, opId)
|
||||
if err != nil {
|
||||
if ctxErr := ctx.Err(); ctxErr != nil {
|
||||
s.logger.Info("Text fetch interrupted by shutdown", "job_id", job.Id)
|
||||
return fmt.Errorf("text fetch interrupted: %w", ctxErr)
|
||||
}
|
||||
s.logger.Error("Failed to get recognition text", "error", err, "operation_id", opId)
|
||||
return err
|
||||
}
|
||||
@@ -450,6 +490,16 @@ func (s *TranscribeService) scheduleRetry(job *entity.TranscribeJob, holder stri
|
||||
return
|
||||
}
|
||||
|
||||
// Остановка попытки не тратит: задача не виновата в том, что нас
|
||||
// перезапустили. Счётчик растёт при захвате, поэтому здесь его возвращают
|
||||
// назад — иначе пять выкладок подряд уводят живую запись в «мертва» с
|
||||
// приговором «попытки исчерпаны».
|
||||
if errors.Is(stepErr, context.Canceled) || errors.Is(stepErr, context.DeadlineExceeded) {
|
||||
if job.Attempts > 0 {
|
||||
job.Attempts--
|
||||
}
|
||||
}
|
||||
|
||||
job.RetryAfter(clock.Now().Add(retryDelay(job.Attempts)))
|
||||
|
||||
if err := s.jobRepo.Save(job, holder); err != nil {
|
||||
|
||||
Reference in New Issue
Block a user