diff --git a/internal/adapter/converter/ffmpeg/ffmpeg.go b/internal/adapter/converter/ffmpeg/ffmpeg.go index f1a4c24..c71a66f 100644 --- a/internal/adapter/converter/ffmpeg/ffmpeg.go +++ b/internal/adapter/converter/ffmpeg/ffmpeg.go @@ -1,6 +1,7 @@ package ffmpeg import ( + "context" "fmt" "os" "os/exec" @@ -15,7 +16,7 @@ func NewFfmpegConverter() *FfmpegConverter { return &FfmpegConverter{} } -func (c *FfmpegConverter) Convert(src, dest string) error { +func (c *FfmpegConverter) Convert(ctx context.Context, src, dest string) error { // Проверяем существование исходного файла if _, err := os.Stat(src); os.IsNotExist(err) { return fmt.Errorf("input file does not exist: %s", src) @@ -26,8 +27,9 @@ func (c *FfmpegConverter) Convert(src, dest string) error { return fmt.Errorf("ffmpeg not found in PATH: %w", err) } - // Создаем команду ffmpeg для конвертации в OGG - cmd := exec.Command(ffmpegExecutable, + // Команда заводится с контекстом: отменённый контекст убивает процесс, а не + // оставляет его дожёвывать чужую запись после остановки воркера. + cmd := exec.CommandContext(ctx, ffmpegExecutable, "-i", src, // входной файл "-c:a", "libvorbis", // кодек Vorbis для OGG "-q:a", "4", // качество аудио (0-10, где 4 - хорошее качество) diff --git a/internal/adapter/metaviewer/ffmpeg/ffmpeg.go b/internal/adapter/metaviewer/ffmpeg/ffmpeg.go index 1671074..bdd79d9 100644 --- a/internal/adapter/metaviewer/ffmpeg/ffmpeg.go +++ b/internal/adapter/metaviewer/ffmpeg/ffmpeg.go @@ -1,6 +1,7 @@ package ffmpeg import ( + "context" "encoding/json" "fmt" "os" @@ -26,7 +27,7 @@ func NewFfmpegMetaViewer() *FfmpegMetaViewer { return &FfmpegMetaViewer{} } -func (m *FfmpegMetaViewer) GetInfo(src string) (*contract.AudioInfo, error) { +func (m *FfmpegMetaViewer) GetInfo(ctx context.Context, src string) (*contract.AudioInfo, error) { // Проверяем существование исходного файла if _, err := os.Stat(src); os.IsNotExist(err) { return nil, fmt.Errorf("input file does not exist: %s", src) @@ -37,8 +38,9 @@ func (m *FfmpegMetaViewer) GetInfo(src string) (*contract.AudioInfo, error) { return nil, fmt.Errorf("ffprobe not found in PATH: %w", err) } - // Создаем команду ffprobe для получения метаданных - cmd := exec.Command(ffprobeExecutable, + // Команда заводится с контекстом: отправитель, закрывший соединение, не + // оставляет за собой чтение метаданных чужого файла. + cmd := exec.CommandContext(ctx, ffprobeExecutable, "-v", "quiet", // тихий режим (без лишнего вывода) "-print_format", "json", // вывод в формате JSON "-show_format", // показать информацию о формате diff --git a/internal/adapter/recognizer/memory.go b/internal/adapter/recognizer/memory.go index f7036cf..7f20f84 100644 --- a/internal/adapter/recognizer/memory.go +++ b/internal/adapter/recognizer/memory.go @@ -1,6 +1,7 @@ package recognizer import ( + "context" "io" "git.vakhrushev.me/av/transcriber/internal/entity" @@ -9,14 +10,14 @@ import ( type MemoryAudioRecognizer struct{} -func (r *MemoryAudioRecognizer) Recognize(file io.Reader, fileName string) (operationID string, err error) { +func (r *MemoryAudioRecognizer) Recognize(ctx context.Context, file io.Reader, fileName string) (operationID string, err error) { return uuid.NewString(), nil } -func (r *MemoryAudioRecognizer) GetRecognitionText(operationID string) (string, error) { +func (r *MemoryAudioRecognizer) GetRecognitionText(ctx context.Context, operationID string) (string, error) { return "Foo bar, Baz.", nil } -func (r *MemoryAudioRecognizer) CheckRecognitionStatus(operationID string) (*entity.RecognitionResult, error) { +func (r *MemoryAudioRecognizer) CheckRecognitionStatus(ctx context.Context, operationID string) (*entity.RecognitionResult, error) { return entity.NewCompletedResult(), nil } diff --git a/internal/adapter/recognizer/yandex/recognizer.go b/internal/adapter/recognizer/yandex/recognizer.go index 2185757..8db2113 100644 --- a/internal/adapter/recognizer/yandex/recognizer.go +++ b/internal/adapter/recognizer/yandex/recognizer.go @@ -1,8 +1,10 @@ package yandex import ( + "context" "fmt" "io" + "time" "git.vakhrushev.me/av/transcriber/internal/entity" ) @@ -54,16 +56,31 @@ func (s *YandexAudioRecognizerService) Close() error { return s.sttService.Close() } -func (s *YandexAudioRecognizerService) Recognize(file io.Reader, fileName string) (string, error) { +// startRecognitionTimeout — сколько ждём принятия операции, когда нас уже +// остановили. Число меньше жёсткого предела остановки: иначе процесс убьют +// прежде, чем ответ дойдёт, и защита ничего не даст. +const startRecognitionTimeout = 10 * time.Second - err := s.s3Sevice.uploadFile(file, fileName) +func (s *YandexAudioRecognizerService) Recognize(ctx context.Context, file io.Reader, fileName string) (string, error) { + + // Заливка отменяется штатно: она дорога по времени, а повтор её бесплатен — + // объект ложится под тем же ключом. + err := s.s3Sevice.uploadFile(ctx, file, fileName) if err != nil { return "", err } uri := s.s3Sevice.fileUrl(fileName) - opId, err := s.sttService.recognizeFileFromS3(uri) + // А вот принятие операции от отмены защищено. Окно короткое и дорогое: + // SpeechKit может операцию принять и начать считать деньги, а ответ до нас + // не доедет — идентификатор потеряется навсегда, и повтор оплатит ту же + // запись второй раз. Свой предел вызову оставлен, чтобы остановка не ждала + // вечно. + startCtx, cancel := protectFromCancel(ctx, startRecognitionTimeout) + defer cancel() + + opId, err := s.sttService.recognizeFileFromS3(startCtx, uri) if err != nil { return "", err } @@ -71,12 +88,19 @@ func (s *YandexAudioRecognizerService) Recognize(file io.Reader, fileName string return opId, nil } -func (s *YandexAudioRecognizerService) GetRecognitionText(operationID string) (string, error) { - return s.sttService.getRecognitionText(operationID) +// protectFromCancel отвязывает вызов от отмены родителя, оставляя ему значения +// родителя и собственный предел по времени. Употребляется там, где обрыв стоит +// дороже ожидания: у платной операции, чей результат нельзя переспросить. +func protectFromCancel(ctx context.Context, timeout time.Duration) (context.Context, context.CancelFunc) { + return context.WithTimeout(context.WithoutCancel(ctx), timeout) } -func (s *YandexAudioRecognizerService) CheckRecognitionStatus(operationID string) (*entity.RecognitionResult, error) { - operation, err := s.sttService.checkOperationStatus(operationID) +func (s *YandexAudioRecognizerService) GetRecognitionText(ctx context.Context, operationID string) (string, error) { + return s.sttService.getRecognitionText(ctx, operationID) +} + +func (s *YandexAudioRecognizerService) CheckRecognitionStatus(ctx context.Context, operationID string) (*entity.RecognitionResult, error) { + operation, err := s.sttService.checkOperationStatus(ctx, operationID) if err != nil { return nil, err } diff --git a/internal/adapter/recognizer/yandex/recognizer_test.go b/internal/adapter/recognizer/yandex/recognizer_test.go new file mode 100644 index 0000000..dbf1ff8 --- /dev/null +++ b/internal/adapter/recognizer/yandex/recognizer_test.go @@ -0,0 +1,38 @@ +package yandex + +import ( + "context" + "testing" + "time" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +// Принятие операции распознавания защищено от отмены: остановка сервиса не +// должна обрывать вызов, который уже мог начать стоить денег и чей результат +// нельзя переспросить. Проверяется само средство защиты — проводка к нему +// оракула не имеет: клиент SpeechKit подставить нечем, а прогон на реальных +// ключах запрещён (CLAUDE.md, «Запреты»). +func TestProtectedContextSurvivesParentCancel(t *testing.T) { + parent, cancel := context.WithCancel(t.Context()) + + protected, release := protectFromCancel(parent, time.Minute) + defer release() + + cancel() + + require.Error(t, parent.Err(), "родитель отменён — иначе проверка судит не то") + assert.NoError(t, protected.Err(), "защищённый вызов пережил отмену родителя") +} + +// Защита не бессрочна: у вызова свой предел, иначе остановка ждала бы вечно. +func TestProtectedContextKeepsItsOwnDeadline(t *testing.T) { + protected, release := protectFromCancel(t.Context(), time.Minute) + defer release() + + deadline, ok := protected.Deadline() + + require.True(t, ok, "у защищённого вызова обязан быть свой предел") + assert.WithinDuration(t, time.Now().Add(time.Minute), deadline, 5*time.Second) +} diff --git a/internal/adapter/recognizer/yandex/s3.go b/internal/adapter/recognizer/yandex/s3.go index 1942482..96a2940 100644 --- a/internal/adapter/recognizer/yandex/s3.go +++ b/internal/adapter/recognizer/yandex/s3.go @@ -67,8 +67,8 @@ func newYandexS3Service(cfg s3Config) (*yandexS3Service, error) { }, nil } -func (s *yandexS3Service) uploadFile(file io.Reader, fileName string) error { - _, err := s.uploader.Upload(context.Background(), &s3.PutObjectInput{ +func (s *yandexS3Service) uploadFile(ctx context.Context, file io.Reader, fileName string) error { + _, err := s.uploader.Upload(ctx, &s3.PutObjectInput{ Bucket: aws.String(s.bucketName), Key: aws.String(fileName), Body: file, diff --git a/internal/adapter/recognizer/yandex/speechkit.go b/internal/adapter/recognizer/yandex/speechkit.go index 6a19bb2..f9f2b9d 100644 --- a/internal/adapter/recognizer/yandex/speechkit.go +++ b/internal/adapter/recognizer/yandex/speechkit.go @@ -94,9 +94,7 @@ func (s *speechKitService) Close() error { } // recognizeFileFromS3 запускает асинхронное распознавание файла из S3 -func (s *speechKitService) recognizeFileFromS3(s3URI string) (string, error) { - ctx := context.Background() - +func (s *speechKitService) recognizeFileFromS3(ctx context.Context, s3URI string) (string, error) { // Добавляем авторизацию и folder_id в контекст ctx = metadata.AppendToOutgoingContext(ctx, "authorization", "Api-Key "+s.apiKey) ctx = metadata.AppendToOutgoingContext(ctx, "x-folder-id", s.folderID) @@ -137,9 +135,7 @@ func (s *speechKitService) recognizeFileFromS3(s3URI string) (string, error) { } // GetRecognitionResult получает результат распознавания по ID операции -func (s *speechKitService) getRecognitionText(operationID string) (string, error) { - ctx := context.Background() - +func (s *speechKitService) getRecognitionText(ctx context.Context, operationID string) (string, error) { // Добавляем авторизацию и folder_id в контекст ctx = metadata.AppendToOutgoingContext(ctx, "authorization", "Api-Key "+s.apiKey) ctx = metadata.AppendToOutgoingContext(ctx, "x-folder-id", s.folderID) @@ -180,9 +176,7 @@ func (s *speechKitService) getRecognitionText(operationID string) (string, error } // checkOperationStatus проверяет статус операции распознавания -func (s *speechKitService) checkOperationStatus(operationID string) (*operation.Operation, error) { - ctx := context.Background() - +func (s *speechKitService) checkOperationStatus(ctx context.Context, operationID string) (*operation.Operation, error) { ctx = metadata.AppendToOutgoingContext(ctx, "authorization", "Api-Key "+s.apiKey) ctx = metadata.AppendToOutgoingContext(ctx, "x-folder-id", s.folderID) diff --git a/internal/adapter/telegram/bot.go b/internal/adapter/telegram/bot.go new file mode 100644 index 0000000..e56d878 --- /dev/null +++ b/internal/adapter/telegram/bot.go @@ -0,0 +1,101 @@ +package telegram + +import ( + "errors" + "fmt" + "log/slog" + "net/http" + "net/url" + "strings" + + tgbotapi "github.com/go-telegram-bot-api/telegram-bot-api/v5" +) + +// ErrEmptyToken — токен бота не задан. Отдельным значением, потому что подъём +// без Telegram — законный исход: сервис продолжает работать с HTTP API. +var ErrEmptyToken = errors.New("telegram bot token is empty") + +// NewBot заводит клиента Bot API — и это **единая точка**, через которую с +// библиотекой разговаривают оба пакета: адаптер отправки и транспорт бота. +// +// Точка нужна ради инварианта «секрет не покидает конфиг». Токен живёт в пути +// каждого обращения к Bot API (`https://api.telegram.org/bot/getFile`), +// а `http.Client` кладёт адрес запроса в `*url.Error` целиком. Библиотека +// отдаёт этот отказ вызывающему как есть, поэтому чистка на месте употребления +// закрывает ровно один вызов из пяти: остаются `getFile`, `sendMessage`, +// `getMe` из конструктора и длинный опрос. Здесь закрыты все. +func NewBot(token string, logger *slog.Logger) (*tgbotapi.BotAPI, error) { + return newBot(token, tgbotapi.APIEndpoint, logger) +} + +// newBot принимает адрес отдельно — иначе проверка утечки токена ходила бы за +// подтверждением в живой Telegram, а боевым токеном запускаться запрещено. +func newBot(token, endpoint string, logger *slog.Logger) (*tgbotapi.BotAPI, error) { + if token == "" { + return nil, ErrEmptyToken + } + + // Длинный опрос живёт внутри библиотеки и печатает свой отказ пакетным + // логгером в stderr (`GetUpdatesChan`), минуя и наш `slog`, и чистку выше. + // Это самый частый путь: опрос идёт непрерывно, а скачивание — только когда + // кто-то прислал запись. Логгер пакетный, поэтому и подменяется один раз. + if err := tgbotapi.SetLogger(&redactingLogger{token: token, logger: logger}); err != nil { + return nil, fmt.Errorf("failed to set telegram logger: %w", err) + } + + return tgbotapi.NewBotAPIWithClient(token, endpoint, &safeClient{inner: &http.Client{}}) +} + +// safeClient — клиент, чей отказ не несёт адреса. Библиотека объявляет +// зависимость интерфейсом `HTTPClient` и возвращает наш отказ вызывающему +// нетронутым, поэтому чистка отсюда доходит до каждого вызова Bot API. +type safeClient struct { + inner *http.Client +} + +func (c *safeClient) Do(req *http.Request) (*http.Response, error) { + resp, err := c.inner.Do(req) + if err != nil { + return nil, WithoutURL(err) + } + return resp, nil +} + +// WithoutURL снимает с отказа адрес запроса, сохраняя причину. Стандартный +// клиент кладёт в `*url.Error` полный URL, а в ссылке Telegram стоит токен +// бота: без этой чистки первый же сбой сети печатает секрет в журнал. +// Причина остаётся и узнаётся `errors.Is` по-прежнему. +func WithoutURL(err error) error { + var urlErr *url.Error + if errors.As(err, &urlErr) { + return urlErr.Err + } + return err +} + +// redactingLogger отдаёт сообщения библиотеки нашему журналу, вычеркнув токен. +// Здесь чистится текст, а не ошибка: библиотека печатает уже отформатированную +// строку, и разбирать в ней `*url.Error` нечего. Замена точная — токен известен. +type redactingLogger struct { + token string + logger *slog.Logger +} + +const redactedToken = "«токен»" + +func (l *redactingLogger) Println(v ...any) { + l.write(strings.TrimSuffix(fmt.Sprintln(v...), "\n")) +} + +func (l *redactingLogger) Printf(format string, v ...any) { + l.write(fmt.Sprintf(format, v...)) +} + +// write пишет на WARN: это сбой фонового цикла со штатным повтором, а не +// событие, требующее разбора (docs/conventions/logging.md, «Уровень — +// это адресат»). Сообщение нейтрально: тем же логгером библиотека печатает и +// отладку, если её включить, а разделить их она не даёт. +func (l *redactingLogger) write(message string) { + l.logger.Warn("Telegram library log", + "message", strings.ReplaceAll(message, l.token, redactedToken)) +} diff --git a/internal/adapter/telegram/bot_test.go b/internal/adapter/telegram/bot_test.go new file mode 100644 index 0000000..e4b58fd --- /dev/null +++ b/internal/adapter/telegram/bot_test.go @@ -0,0 +1,140 @@ +package telegram + +import ( + "bytes" + "errors" + "log/slog" + "net/http" + "net/http/httptest" + "net/url" + "strings" + "testing" + + tgbotapi "github.com/go-telegram-bot-api/telegram-bot-api/v5" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +// Токен виден в пути каждого обращения к Bot API, а `http.Client` кладёт путь в +// `*url.Error` целиком. Проверки ниже судят по тексту: секрет не должен +// встречаться ни в отказе, ни в строке журнала. Утечка необратима — утёкший +// токен отзывают руками (CLAUDE.md, «Инварианты», critical). +const probeToken = "7654321:AAHsecretBOTtokenVALUE" + +// getMeResponse — ответ, которым подставной Telegram пускает конструктор +// дальше: `NewBotAPIWithClient` ходит за `getMe` прежде, чем отдать клиента. +const getMeResponse = `{"ok":true,"result":{"id":1,"is_bot":true,"first_name":"probe","username":"probe_bot"}}` + +func newProbeBot(t *testing.T, handler http.HandlerFunc) (*tgbotapi.BotAPI, *httptest.Server) { + t.Helper() + + server := httptest.NewServer(handler) + t.Cleanup(server.Close) + + bot, err := newBot(probeToken, server.URL+"/bot%s/%s", slog.New(slog.DiscardHandler)) + require.NoError(t, err) + + return bot, server +} + +// Отказ транспорта на любом вызове Bot API не несёт токена: чистка стоит на +// границе клиента, а не у места употребления, поэтому закрыты все вызовы разом. +func TestBotAPIFailureDoesNotCarryToken(t *testing.T) { + bot, server := newProbeBot(t, func(w http.ResponseWriter, _ *http.Request) { + if _, err := w.Write([]byte(getMeResponse)); err != nil { + t.Errorf("подставной Telegram не смог ответить: %v", err) + } + }) + + // Собеседник исчез — так выглядит обрыв сети, DNS-сбой и недоступность + // api.telegram.org. + server.Close() + + t.Run("getFile", func(t *testing.T) { + _, err := bot.GetFile(tgbotapi.FileConfig{FileID: "any"}) + require.Error(t, err) + assert.NotContains(t, err.Error(), probeToken, "токен уехал в отказ: %v", err) + }) + + t.Run("sendMessage", func(t *testing.T) { + _, err := bot.Send(tgbotapi.NewMessage(1, "текст")) + require.Error(t, err) + assert.NotContains(t, err.Error(), probeToken, "токен уехал в отказ: %v", err) + }) +} + +// Отказ конструктора несёт тот же путь: `NewBotAPIWithClient` ходит за `getMe`, +// и контейнер, стартующий раньше сети, печатал бы токен в первую же секунду. +func TestBotConstructionFailureDoesNotCarryToken(t *testing.T) { + server := httptest.NewServer(http.HandlerFunc(func(http.ResponseWriter, *http.Request) {})) + server.Close() + + _, err := newBot(probeToken, server.URL+"/bot%s/%s", slog.New(slog.DiscardHandler)) + + require.Error(t, err) + assert.NotContains(t, err.Error(), probeToken, "токен уехал в отказ конструктора: %v", err) +} + +// Длинный опрос печатает свои отказы пакетным логгером самой библиотеки, минуя +// наш `slog`. Логгер подменён — значит, и эта строка идёт через вычистку. +func TestLibraryLoggerRedactsToken(t *testing.T) { + journal := &bytes.Buffer{} + logger := slog.New(slog.NewTextHandler(journal, nil)) + + redacting := &redactingLogger{token: probeToken, logger: logger} + redacting.Println(errors.New(`Post "https://api.telegram.org/bot` + probeToken + `/getUpdates": dial tcp: refused`)) + redacting.Printf("Failed to get updates from %s", "https://api.telegram.org/bot"+probeToken+"/getUpdates") + + written := journal.String() + assert.NotContains(t, written, probeToken, "токен уехал в журнал: %s", written) + assert.Equal(t, 2, strings.Count(written, redactedToken), "вместо токена стоит пометка") + assert.Contains(t, written, "dial tcp", "причина отказа осталась") +} + +// Пустой токен — законный исход подъёма без Telegram, и узнаётся он по смыслу. +func TestEmptyTokenIsRecognizedByValue(t *testing.T) { + _, err := NewBot("", slog.New(slog.DiscardHandler)) + + require.ErrorIs(t, err, ErrEmptyToken) +} + +// WithoutURL снимает адрес, но не причину: `errors.Is` по цепочке продолжает +// работать, иначе чистка стоила бы узнаваемости отказа. +func TestWithoutURLKeepsCause(t *testing.T) { + cause := errors.New("dial tcp: connection refused") + wrapped := &url.Error{Op: "Post", URL: "https://api.telegram.org/bot" + probeToken + "/getMe", Err: cause} + + cleaned := WithoutURL(wrapped) + + assert.NotContains(t, cleaned.Error(), probeToken) + require.ErrorIs(t, cleaned, cause) + assert.Equal(t, cause, WithoutURL(cause), "отказ без адреса не трогают") +} + +// Стык, которого не сторожил никто: подмена пакетного логгера держится одной +// строкой в `NewBot`, а снятие этой строки не роняло ни одной проверки. Оракул +// косвенный по необходимости — библиотека не отдаёт установленный логгер +// обратно, — поэтому он смотрит на исход: её собственная строка обязана +// оказаться в нашем журнале. +func TestLibraryLoggerIsActuallyInstalled(t *testing.T) { + journal := &bytes.Buffer{} + + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + if _, err := w.Write([]byte(getMeResponse)); err != nil { + t.Errorf("подставной Telegram не смог ответить: %v", err) + } + })) + t.Cleanup(server.Close) + + bot, err := newBot(probeToken, server.URL+"/bot%s/%s", slog.New(slog.NewTextHandler(journal, nil))) + require.NoError(t, err) + + // Отладку библиотека печатает тем же логгером, что и отказы: включаем её, + // чтобы строка появилась без обрыва сети. + bot.Debug = true + _, err = bot.GetMe() + require.NoError(t, err) + + assert.Contains(t, journal.String(), "Telegram library log", + "строка библиотеки прошла мимо нашего журнала: логгер не подменён") +} diff --git a/internal/adapter/telegram/sender.go b/internal/adapter/telegram/sender.go index c21b9b8..ab6976c 100644 --- a/internal/adapter/telegram/sender.go +++ b/internal/adapter/telegram/sender.go @@ -16,7 +16,9 @@ type TelegramMessageSender struct { } func NewTelegramMessageSender(botToken string, logger *slog.Logger) (*TelegramMessageSender, error) { - bot, err := tgbotapi.NewBotAPI(botToken) + // Клиент заводится единой точкой: её отказ не несёт токена, а отказ + // конструктора несёт — `NewBotAPI` зовёт `getMe`. + bot, err := NewBot(botToken, logger) if err != nil { return nil, err } diff --git a/internal/contract/contract.go b/internal/contract/contract.go index 894e110..4802621 100644 --- a/internal/contract/contract.go +++ b/internal/contract/contract.go @@ -1,6 +1,7 @@ package contract import ( + "context" "io" "git.vakhrushev.me/av/transcriber/internal/entity" @@ -10,18 +11,23 @@ type AudioInfo struct { Seconds int // Длина аудиофайла в секундах } +// Контекст первым доводом несут все интерфейсы, за которыми стоит внешний +// собеседник — процесс `ffmpeg`, S3, SpeechKit. Он здесь не украшение: остановка +// сервиса обязана доходить до чужой работы, а не оставлять её сиротой. Без него +// конвертация шестичасовой записи переживает остановку воркера, а запрос к +// платному распознаванию висит до собственного таймаута библиотеки. type AudioMetaViewer interface { - GetInfo(src string) (*AudioInfo, error) + GetInfo(ctx context.Context, src string) (*AudioInfo, error) } type AudioFileConverter interface { - Convert(src, dest string) error + Convert(ctx context.Context, src, dest string) error } type AudioRecognizer interface { - Recognize(file io.Reader, fileName string) (operationID string, err error) - GetRecognitionText(operationID string) (string, error) - CheckRecognitionStatus(operationID string) (*entity.RecognitionResult, error) + Recognize(ctx context.Context, file io.Reader, fileName string) (operationID string, err error) + GetRecognitionText(ctx context.Context, operationID string) (string, error) + CheckRecognitionStatus(ctx context.Context, operationID string) (*entity.RecognitionResult, error) } type TelegramMessageSender interface { diff --git a/internal/controller/http/transcribe.go b/internal/controller/http/transcribe.go index 4d73e28..083094d 100644 --- a/internal/controller/http/transcribe.go +++ b/internal/controller/http/transcribe.go @@ -1,6 +1,7 @@ package http import ( + "context" "log/slog" "net/http" "time" @@ -71,7 +72,14 @@ func (h *TranscribeHandler) CreateTranscribeJob(e *core.RequestEvent) error { } }() - job, err := h.trsService.CreateJobFromApi(file, header.Filename) + // Запись доехала целиком, поэтому задача заводится независимо от того, + // дождётся ли отправитель ответа: на контексте запроса приём терял бы + // полностью загруженную запись от одного обрыва соединения, а забрать + // результат он может и позже — по `GET /status/{id}`. Значения контекста + // (журнал запроса, сессия) при этом сохраняются, теряется только отмена. + ctx := context.WithoutCancel(e.Request.Context()) + + job, err := h.trsService.CreateJobFromApi(ctx, file, header.Filename) if err != nil { // Второй раз отказ не логируем: приём назван конвенцией логирующей // границей и уже написал о нём. Транспорт переводит ошибку в ответ. diff --git a/internal/controller/http/transcribe_test.go b/internal/controller/http/transcribe_test.go index 45c8f8b..1d70c5d 100644 --- a/internal/controller/http/transcribe_test.go +++ b/internal/controller/http/transcribe_test.go @@ -2,6 +2,7 @@ package http import ( "bytes" + "context" "encoding/json" "errors" "fmt" @@ -39,7 +40,12 @@ type stubMetaViewer struct { err error } -func (m *stubMetaViewer) GetInfo(string) (*contract.AudioInfo, error) { +func (m *stubMetaViewer) GetInfo(ctx context.Context, _ string) (*contract.AudioInfo, error) { + // Настоящий `ffprobe` заведён с контекстом и по отмене умирает; стаб, + // который контекст игнорирует, сделал бы проверку приёма неспособной упасть. + if err := ctx.Err(); err != nil { + return nil, err + } if m.err != nil { return nil, m.err } @@ -50,7 +56,7 @@ func (m *stubMetaViewer) GetInfo(string) (*contract.AudioInfo, error) { // этого файла её не зовёт. type stubConverter struct{} -func (c *stubConverter) Convert(string, string) error { return nil } +func (c *stubConverter) Convert(context.Context, string, string) error { return nil } // TestTgSender: приём по HTTP в Telegram не отвечает, но сервису отправитель нужен. type TestTgSender struct{} @@ -727,3 +733,27 @@ func TestGetTranscribeJobStatus_NotFound(t *testing.T) { assert.Equal(t, "Job not found", response["error"]) } + +// Отправитель, у которого соединение оборвалось после полной загрузки, задачу +// всё равно получает: запись доехала целиком, а результат он заберёт позже по +// `GET /status/{id}`. Приём на контексте запроса терял бы такую запись молча — +// решение владельца от 2026-08-13. +func TestAcceptedRecordSurvivesSenderDisconnect(t *testing.T) { + env := setupTestEnv(t, readableMetaViewer()) + + req := createMultipartRequest(t, "sample.m4a", []byte("аудио")) + // Так выглядит ушедший отправитель: контекст запроса отменяется сервером, + // когда соединение закрылось. + ctx, cancel := context.WithCancel(req.Context()) + cancel() + req = req.WithContext(ctx) + + w := httptest.NewRecorder() + env.serve(w, req) + + require.Equal(t, http.StatusCreated, w.Result().StatusCode, "тело ответа: %s", w.Body.String()) + + jobs, err := env.app.FindAllRecords(migrations.JobsCollection) + require.NoError(t, err) + assert.Len(t, jobs, 1, "задача заведена, несмотря на ушедшего отправителя") +} diff --git a/internal/controller/tg/download_test.go b/internal/controller/tg/download_test.go new file mode 100644 index 0000000..8a76465 --- /dev/null +++ b/internal/controller/tg/download_test.go @@ -0,0 +1,111 @@ +package tg + +import ( + "io" + "log/slog" + "net/http" + "strings" + "testing" + + tgbotapi "github.com/go-telegram-bot-api/telegram-bot-api/v5" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +// probeClient подменяет клиента бота и запоминает, кого спрашивали. Через него +// проверяется стык: скачивание обязано идти клиентом бота, а не общим +// `http.DefaultClient` — чистку отказа от адреса с токеном несёт именно клиент +// (`internal/adapter/telegram`). Подмена на общий клиент правил гейта не +// нарушает, поэтому сторожить стык может только проверка. +type probeClient struct { + seen []string + download func(w *probeResponse) +} + +type probeResponse struct { + status int + body string +} + +func (c *probeClient) Do(req *http.Request) (*http.Response, error) { + c.seen = append(c.seen, req.URL.Path) + + switch { + case strings.Contains(req.URL.Path, "/getMe"): + return jsonResponse(`{"ok":true,"result":{"id":1,"is_bot":true,"username":"probe_bot"}}`), nil + case strings.Contains(req.URL.Path, "/getFile"): + return jsonResponse(`{"ok":true,"result":{"file_id":"x","file_path":"voice/file_1.ogg"}}`), nil + } + + answer := &probeResponse{status: http.StatusOK, body: "аудио"} + if c.download != nil { + c.download(answer) + } + return &http.Response{ + StatusCode: answer.status, + Body: io.NopCloser(strings.NewReader(answer.body)), + Header: make(http.Header), + }, nil +} + +func jsonResponse(body string) *http.Response { + header := make(http.Header) + header.Set("Content-Type", "application/json") + return &http.Response{ + StatusCode: http.StatusOK, + Body: io.NopCloser(strings.NewReader(body)), + Header: header, + } +} + +func newProbeController(t *testing.T, client *probeClient) *TelegramController { + t.Helper() + + // Клиент подставной, поэтому адрес значения не имеет — важно лишь, что + // библиотека соберёт из него разбираемый URL. + bot, err := tgbotapi.NewBotAPIWithClient( + "7654321:AAHsecretBOTtokenVALUE", + "http://telegram.probe/bot%s/%s", + client, + ) + require.NoError(t, err) + + return &TelegramController{ + bot: bot, + logger: slog.New(slog.DiscardHandler), + } +} + +// Скачивание идёт клиентом бота: иначе отказ пойдёт мимо чистки и унесёт токен. +func TestDownloadGoesThroughBotClient(t *testing.T) { + client := &probeClient{} + controller := newProbeController(t, client) + + body, name, err := controller.downloadAudioFile(t.Context(), "file-id") + require.NoError(t, err) + t.Cleanup(func() { + if err := body.Close(); err != nil { + t.Errorf("не удалось закрыть тело: %v", err) + } + }) + + assert.Equal(t, "voice/file_1.ogg", name) + require.Len(t, client.seen, 3, "клиент бота видел все обращения: getMe, getFile и скачивание") + assert.Contains(t, client.seen[2], "voice/file_1.ogg", "скачивание ушло мимо клиента бота") +} + +// Отказ выдачи файла — это не запись: тело такого ответа не должно доехать до +// хранилища и умереть на `ffprobe`, уведя диагностику к чужой причине. +func TestDownloadRejectsNonOKStatus(t *testing.T) { + client := &probeClient{download: func(w *probeResponse) { + w.status = http.StatusUnauthorized + w.body = `{"ok":false,"error_code":401,"description":"Unauthorized"}` + }} + controller := newProbeController(t, client) + + body, _, err := controller.downloadAudioFile(t.Context(), "file-id") + + require.Error(t, err) + assert.Nil(t, body, "тело отказа наружу не отдают") + assert.Contains(t, err.Error(), "401", "код ответа назван — по нему видно, что отказал Telegram") +} diff --git a/internal/controller/tg/tg.go b/internal/controller/tg/tg.go index 59ebbbd..cba0f49 100644 --- a/internal/controller/tg/tg.go +++ b/internal/controller/tg/tg.go @@ -1,6 +1,8 @@ package tg import ( + "context" + "errors" "fmt" "io" "log/slog" @@ -8,6 +10,12 @@ import ( "slices" "strings" + // Транспорт знает адаптер Telegram ровно ради единой точки чистки отказа: + // второй экземпляр той же функции здесь был бы вторым способом делать одно + // и то же, а секрет в журнале — необратим. Направление «транспорт не знает + // адаптера» правилом не держится и уже нарушено HTTP-поверхностью + // (docs/conventions/go-linters.md, «Что остаётся прозой»). + "git.vakhrushev.me/av/transcriber/internal/adapter/telegram" "git.vakhrushev.me/av/transcriber/internal/contract" "git.vakhrushev.me/av/transcriber/internal/service" tgbotapi "github.com/go-telegram-bot-api/telegram-bot-api/v5" @@ -25,25 +33,22 @@ type TelegramController struct { } type TelegramConfig struct { - BotToken string UpdateTimeout int UserWhiteList []string } +// NewTelegramController принимает готового клиента, а не токен: клиента заводит +// единая точка `internal/adapter/telegram`, и только её отказ не несёт секрета. +// Токен сюда не приезжает вовсе — значит, и утечь отсюда ему неоткуда. func NewTelegramController( config TelegramConfig, + bot *tgbotapi.BotAPI, transcribeService *service.TranscribeService, jobRepo contract.TranscriptJobRepository, logger *slog.Logger, ) (*TelegramController, error) { - botToken := config.BotToken - if botToken == "" { - return nil, &EmptyBotTokenError{} - } - - bot, err := tgbotapi.NewBotAPI(botToken) - if err != nil { - return nil, err + if bot == nil { + return nil, errors.New("telegram bot is not created") } controller := &TelegramController{ @@ -58,7 +63,11 @@ func NewTelegramController( return controller, nil } -func (c *TelegramController) Start() { +// Start принимает контекст жизни процесса и отдаёт его каждому обработчику: +// скачивание записи и разбор её метаданных — работа с внешним собеседником, и +// остановка сервиса обязана до неё доходить. Приём обновлений контекстом не +// правится: его прекращает Stop. +func (c *TelegramController) Start(ctx context.Context) { c.logger.Info("Telegram bot started", "username", c.bot.Self.UserName) u := tgbotapi.NewUpdate(0) @@ -94,11 +103,11 @@ func (c *TelegramController) Start() { // Handle audio messages and files if update.Message.Audio != nil { - c.handleAudioMessage(update.Message) + c.handleAudioMessage(ctx, update.Message) } else if update.Message.Voice != nil { - c.handleVoiceMessage(update.Message) + c.handleVoiceMessage(ctx, update.Message) } else if update.Message.Document != nil { - c.handleDocumentMessage(update.Message) + c.handleDocumentMessage(ctx, update.Message) } } } @@ -148,7 +157,7 @@ func (c *TelegramController) handleHelpCommand(message *tgbotapi.Message) { c.send(msg) } -func (c *TelegramController) handleAudioMessage(message *tgbotapi.Message) { +func (c *TelegramController) handleAudioMessage(ctx context.Context, message *tgbotapi.Message) { // Отправляем сообщение о начале обработки progressMsg := tgbotapi.NewMessage(message.Chat.ID, "Обрабатываю аудиофайл...") progressMsg.ReplyToMessageID = message.MessageID @@ -159,7 +168,7 @@ func (c *TelegramController) handleAudioMessage(message *tgbotapi.Message) { } // Скачиваем файл - fileReader, fileName, err := c.downloadAudioFile(message.Audio.FileID) + fileReader, fileName, err := c.downloadAudioFile(ctx, message.Audio.FileID) if err != nil { c.logger.Error("Failed to download audio file", "error", err) errorMsg := tgbotapi.NewMessage(message.Chat.ID, "Ошибка при скачивании аудиофайла. Попробуйте еще раз.") @@ -169,7 +178,7 @@ func (c *TelegramController) handleAudioMessage(message *tgbotapi.Message) { defer fileReader.Close() // Обрабатываем файл - job, err := c.transcribeService.CreateJobFromTelegram(fileReader, fileName, message.Chat.ID, sentProgressMsg.MessageID) + job, err := c.transcribeService.CreateJobFromTelegram(ctx, fileReader, fileName, message.Chat.ID, sentProgressMsg.MessageID) if err != nil { c.logger.Error("Failed to create transcribe job", "error", err) errorMsg := tgbotapi.NewMessage(message.Chat.ID, "Ошибка при создании задачи на расшифровку. Попробуйте еще раз.") @@ -183,7 +192,7 @@ func (c *TelegramController) handleAudioMessage(message *tgbotapi.Message) { c.send(successMsg) } -func (c *TelegramController) handleVoiceMessage(message *tgbotapi.Message) { +func (c *TelegramController) handleVoiceMessage(ctx context.Context, message *tgbotapi.Message) { // Отправляем сообщение о начале обработки progressMsg := tgbotapi.NewMessage(message.Chat.ID, "Обрабатываю голосовое сообщение...") progressMsg.ReplyToMessageID = message.MessageID @@ -194,7 +203,7 @@ func (c *TelegramController) handleVoiceMessage(message *tgbotapi.Message) { } // Скачиваем файл - fileReader, fileName, err := c.downloadAudioFile(message.Voice.FileID) + fileReader, fileName, err := c.downloadAudioFile(ctx, message.Voice.FileID) if err != nil { c.logger.Error("Failed to download voice file", "error", err) errorMsg := tgbotapi.NewMessage(message.Chat.ID, "Ошибка при скачивании голосового сообщения. Попробуйте еще раз.") @@ -204,7 +213,7 @@ func (c *TelegramController) handleVoiceMessage(message *tgbotapi.Message) { defer fileReader.Close() // Обрабатываем файл - job, err := c.transcribeService.CreateJobFromTelegram(fileReader, fileName, message.Chat.ID, sentProgressMsg.MessageID) + job, err := c.transcribeService.CreateJobFromTelegram(ctx, fileReader, fileName, message.Chat.ID, sentProgressMsg.MessageID) if err != nil { c.logger.Error("Failed to create transcribe job", "error", err) errorMsg := tgbotapi.NewMessage(message.Chat.ID, "Ошибка при создании задачи на расшифровку. Попробуйте еще раз.") @@ -218,7 +227,7 @@ func (c *TelegramController) handleVoiceMessage(message *tgbotapi.Message) { c.send(successMsg) } -func (c *TelegramController) handleDocumentMessage(message *tgbotapi.Message) { +func (c *TelegramController) handleDocumentMessage(ctx context.Context, message *tgbotapi.Message) { // Проверяем, является ли документ аудиофайлом if !c.isAudioDocument(message.Document) { return @@ -234,7 +243,7 @@ func (c *TelegramController) handleDocumentMessage(message *tgbotapi.Message) { } // Скачиваем файл - fileReader, fileName, err := c.downloadAudioFile(message.Document.FileID) + fileReader, fileName, err := c.downloadAudioFile(ctx, message.Document.FileID) if err != nil { c.logger.Error("Failed to download document file", "error", err) errorMsg := tgbotapi.NewMessage(message.Chat.ID, "Ошибка при скачивании аудиофайла. Попробуйте еще раз.") @@ -244,7 +253,7 @@ func (c *TelegramController) handleDocumentMessage(message *tgbotapi.Message) { defer fileReader.Close() // Обрабатываем файл - job, err := c.transcribeService.CreateJobFromTelegram(fileReader, fileName, message.Chat.ID, sentProgressMsg.MessageID) + job, err := c.transcribeService.CreateJobFromTelegram(ctx, fileReader, fileName, message.Chat.ID, sentProgressMsg.MessageID) if err != nil { c.logger.Error("Failed to create transcribe job", "error", err) errorMsg := tgbotapi.NewMessage(message.Chat.ID, "Ошибка при создании задачи на расшифровку. Попробуйте еще раз.") @@ -258,20 +267,41 @@ func (c *TelegramController) handleDocumentMessage(message *tgbotapi.Message) { c.send(successMsg) } -func (c *TelegramController) downloadAudioFile(fileID string) (io.ReadCloser, string, error) { +func (c *TelegramController) downloadAudioFile(ctx context.Context, fileID string) (io.ReadCloser, string, error) { // Получаем информацию о файле file, err := c.bot.GetFile(tgbotapi.FileConfig{FileID: fileID}) if err != nil { return nil, "", fmt.Errorf("failed to get file info: %w", err) } - // Скачиваем файл + // Скачиваем файл. Запрос заводится с контекстом: скачивание шестичасовой + // записи иначе продолжается и после остановки сервиса, а ссылка на файл + // несёт токен бота — держать её живой дольше нужного незачем. + // + // Клиент берётся у бота, а не `http.DefaultClient`: у бота он свой, и его + // отказ уже не несёт адреса (`internal/adapter/telegram`, единая точка). fileURL := file.Link(c.bot.Token) - resp, err := http.Get(fileURL) + request, err := http.NewRequestWithContext(ctx, http.MethodGet, fileURL, nil) + if err != nil { + return nil, "", fmt.Errorf("failed to build download request: %w", telegram.WithoutURL(err)) + } + + resp, err := c.bot.Client.Do(request) if err != nil { return nil, "", fmt.Errorf("failed to download file: %w", err) } + // Отказ выдачи файла — это не запись. Без проверки телом «записи» станет + // JSON вида `{"ok":false,…}`: он доедет до хранилища, ляжет рабочей копией + // и умрёт на `ffprobe`, а отправитель получит жалобу на свой файл вместо + // правды о протухшей ссылке. + if resp.StatusCode != http.StatusOK { + if err := resp.Body.Close(); err != nil { + c.logger.Error("Failed to close download response", "error", err) + } + return nil, "", fmt.Errorf("failed to download file: unexpected status %d", resp.StatusCode) + } + // Получаем имя файла из URL fileName := file.FilePath if fileName == "" { diff --git a/internal/controller/worker/worker.go b/internal/controller/worker/worker.go index c43d7a7..744a90a 100644 --- a/internal/controller/worker/worker.go +++ b/internal/controller/worker/worker.go @@ -17,21 +17,31 @@ type Worker interface { Name() string } +// pollInterval — пауза между прогонами шага. Полем, а не константой по месту: +// проверке нужен второй прогон, чтобы остановить воркер **после** того, как он +// рассудил об исходе первого. Отменять контекст изнутри шага она не может — +// отменённый контекст теперь и значит «нас остановили». +const pollInterval = time.Second + type CallbackWorker struct { - name string - f func() error - logger *slog.Logger + name string + // Шаг принимает контекст воркера: остановка обязана доходить до чужой + // работы, которую шаг завёл, а не только прерывать цикл между шагами. + f func(ctx context.Context) error + logger *slog.Logger + interval time.Duration } -func NewCallbackWorker(name string, f func() error, logger *slog.Logger) *CallbackWorker { +func NewCallbackWorker(name string, f func(ctx context.Context) error, logger *slog.Logger) *CallbackWorker { if logger == nil { logger = slog.Default() } return &CallbackWorker{ - name: name, - f: f, - logger: logger, + name: name, + f: f, + logger: logger, + interval: pollInterval, } } @@ -48,25 +58,34 @@ func (w *CallbackWorker) Start(ctx context.Context) { w.logger.Info("Worker received shutdown signal", "worker", w.Name()) return default: - err := w.f() + err := w.f(ctx) // Признак узнаётся по смыслу, а не по точной форме значения: // приведение типа видело только вершину цепочки и сломалось бы от // первой же обёртки `%w`, которая в проекте — умолчание. var noop *contract.NoopJobError isNoop := errors.As(err, &noop) - if !isNoop { + // Остановка — не отказ шага: контекст отменили мы сами. Считать её + // в метрику и писать владельцу «Worker error» значит красить каждую + // выкладку как поломку — по тому же доводу, по которому не считается + // `NoopJobError`. Судит контекст, а не текст ошибки: убитый процесс + // отдаёт «signal: killed», и `errors.Is` его с отменой не свяжет. + stopped := err != nil && !isNoop && ctx.Err() != nil + if !isNoop && !stopped { metrics.WorkerJobCounter.WithLabelValues(w.Name(), strconv.FormatBool(err != nil)).Inc() } - if err != nil && !isNoop { + if err != nil && !isNoop && !stopped { w.logger.Error("Worker error", "worker", w.Name(), "error", err) } + if stopped { + w.logger.Info("Worker step interrupted by shutdown", "worker", w.Name()) + } - // Ждем 1 секунду перед следующей итерацией + // Ждем перед следующей итерацией select { case <-ctx.Done(): w.logger.Info("Worker received shutdown signal during sleep", "worker", w.Name()) return - case <-time.After(1 * time.Second): + case <-time.After(w.interval): // Продолжаем работу } } diff --git a/internal/controller/worker/worker_test.go b/internal/controller/worker/worker_test.go index d4816db..ce73e53 100644 --- a/internal/controller/worker/worker_test.go +++ b/internal/controller/worker/worker_test.go @@ -40,10 +40,13 @@ func (b *journalBuffer) String() string { } // runOnce прогоняет воркер ровно один раз и возвращает журнал этого прогона. -// Цикл воркера бесконечен и спит секунду между прогонами, поэтому контекст -// отменяется сразу после первого вызова работы: ждать второго прогона нечего, а -// секунда сна на проверку — цена ни за что. -func runOnce(t *testing.T, name string, work func() error) string { +// +// Воркер останавливает **второй** прогон, а не первый: отменённый контекст +// теперь и значит «нас остановили», и отмена изнутри первого шага сделала бы +// его исход неотличимым от остановки — проверка судила бы не то, что заявляет. +// Пауза между прогонами на время проверки укорочена до миллисекунды: ждать +// секунду ради второго вызова незачем. +func runOnce(t *testing.T, name string, work func(ctx context.Context) error) string { t.Helper() journal := &journalBuffer{} @@ -54,15 +57,21 @@ func runOnce(t *testing.T, name string, work func() error) string { var once sync.Once done := make(chan struct{}) + calls := 0 - w := NewCallbackWorker(name, func() error { - err := work() - once.Do(func() { - cancel() - close(done) - }) - return err + w := NewCallbackWorker(name, func(ctx context.Context) error { + calls++ + if calls > 1 { + // Первый прогон уже рассужен: журнал написан, счётчик сдвинут. + once.Do(func() { + cancel() + close(done) + }) + return &contract.NoopJobError{State: "stopping"} + } + return work(ctx) }, logger) + w.interval = time.Millisecond finished := make(chan struct{}) go func() { @@ -147,7 +156,7 @@ func TestWrappedNoopIsNotAFailure(t *testing.T) { before := jobCount(t, name, "false") beforeErr := jobCount(t, name, "true") - journal := runOnce(t, name, func() error { + journal := runOnce(t, name, func(context.Context) error { return fmt.Errorf("find and acquire job: %w", &contract.NoopJobError{State: "created"}) }) @@ -171,7 +180,7 @@ func TestFailureIsLoggedAndCounted(t *testing.T) { before := jobCount(t, name, "true") - journal := runOnce(t, name, func() error { + journal := runOnce(t, name, func(context.Context) error { return errors.New("database is gone") }) @@ -191,7 +200,7 @@ func TestSuccessIsCounted(t *testing.T) { before := jobCount(t, name, "false") - journal := runOnce(t, name, func() error { + journal := runOnce(t, name, func(context.Context) error { return nil }) @@ -202,3 +211,50 @@ func TestSuccessIsCounted(t *testing.T) { t.Errorf("успешный прогон записан отказом: журнал %q", journal) } } + +// Остановка сервиса — не отказ шага: контекст отменили мы сами. Без этой +// развилки каждая выкладка красит журнал владельца отказами и накручивает +// счётчик сбоев, которых не было, — тот же довод, по которому не считается +// `NoopJobError`. Судит контекст, а не текст ошибки: убитый по контексту +// процесс отдаёт «signal: killed», и `errors.Is` его с отменой не свяжет. +func TestShutdownIsNotAFailure(t *testing.T) { + const name = "stopped_worker" + + beforeErr := jobCount(t, name, "true") + beforeOk := jobCount(t, name, "false") + + journal := &journalBuffer{} + logger := slog.New(slog.NewTextHandler(journal, nil)) + + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + + w := NewCallbackWorker(name, func(context.Context) error { + // Так выглядит шаг, которого застала остановка. + cancel() + return errors.New("ffmpeg conversion failed: signal: killed") + }, logger) + w.interval = time.Millisecond + + finished := make(chan struct{}) + go func() { + w.Start(ctx) + close(finished) + }() + + select { + case <-finished: + case <-time.After(5 * time.Second): + t.Fatal("воркер не остановился по отмене контекста") + } + + if got := journal.String(); strings.Contains(got, "Worker error") { + t.Errorf("остановка записана отказом: журнал %q", got) + } + if got := jobCount(t, name, "true"); got != beforeErr { + t.Errorf("остановка засчитана отказом: было %v, стало %v", beforeErr, got) + } + if got := jobCount(t, name, "false"); got != beforeOk { + t.Errorf("остановка засчитана успешным прогоном: было %v, стало %v", beforeOk, got) + } +} diff --git a/internal/service/pipeline_test.go b/internal/service/pipeline_test.go index fe75e2f..0af6c72 100644 --- a/internal/service/pipeline_test.go +++ b/internal/service/pipeline_test.go @@ -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) diff --git a/internal/service/recognition_test.go b/internal/service/recognition_test.go index 8cd4008..c3fb8f9 100644 --- a/internal/service/recognition_test.go +++ b/internal/service/recognition_test.go @@ -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) diff --git a/internal/service/shutdown_test.go b/internal/service/shutdown_test.go new file mode 100644 index 0000000..63b7b69 --- /dev/null +++ b/internal/service/shutdown_test.go @@ -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")) +} diff --git a/internal/service/transcribe.go b/internal/service/transcribe.go index 8f177da..db4362f 100644 --- a/internal/service/transcribe.go +++ b/internal/service/transcribe.go @@ -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 { diff --git a/main.go b/main.go index b9ff958..1794762 100644 --- a/main.go +++ b/main.go @@ -19,6 +19,7 @@ import ( pbrepo "git.vakhrushev.me/av/transcriber/internal/adapter/repo/pocketbase" "git.vakhrushev.me/av/transcriber/internal/adapter/telegram" "git.vakhrushev.me/av/transcriber/internal/config" + "git.vakhrushev.me/av/transcriber/internal/contract" httpcontroller "git.vakhrushev.me/av/transcriber/internal/controller/http" tgcontroller "git.vakhrushev.me/av/transcriber/internal/controller/tg" "git.vakhrushev.me/av/transcriber/internal/controller/worker" @@ -136,13 +137,13 @@ func main() { var wg sync.WaitGroup tgConfig := tgcontroller.TelegramConfig{ - BotToken: cfg.Telegram.BotToken, UpdateTimeout: cfg.Telegram.UpdateTimeout, UserWhiteList: cfg.Server.UsersWhiteList, } - // Создаем Telegram бот - tgController, err := tgcontroller.NewTelegramController(tgConfig, transcribeService, jobRepo, logger) + // Клиента бота заводит единая точка: её отказ не несёт токена, тогда как + // отказ `NewBotAPI` несёт — он ходит за `getMe`. + tgController, err := newTelegramController(cfg.Telegram.BotToken, tgConfig, transcribeService, jobRepo, logger) if err != nil { logger.Error("Failed to create Telegram controller", "error", err) // Не останавливаем приложение, если Telegram бот не создан @@ -152,7 +153,7 @@ func main() { go func() { defer wg.Done() logger.Info("Starting Telegram bot") - tgController.Start() + tgController.Start(ctx) logger.Info("Telegram bot stopped gracefully") }() } @@ -327,3 +328,21 @@ func main() { logger.Info("Transcriber service stopped") } + +// newTelegramController собирает бота и транспорт вокруг него. Токен доходит +// до единой точки `internal/adapter/telegram` и дальше не идёт: транспорт его +// не видит вовсе, а отказ, который увидит журнал, адреса с токеном не несёт. +func newTelegramController( + botToken string, + cfg tgcontroller.TelegramConfig, + transcribeService *service.TranscribeService, + jobRepo contract.TranscriptJobRepository, + logger *slog.Logger, +) (*tgcontroller.TelegramController, error) { + bot, err := telegram.NewBot(botToken, logger) + if err != nil { + return nil, err + } + + return tgcontroller.NewTelegramController(cfg, bot, transcribeService, jobRepo, logger) +}