- контекст проложен от воркера и обоих входов до внешних вызовов: 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, уводя диагностику
94 lines
3.5 KiB
Go
94 lines
3.5 KiB
Go
package worker
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"log/slog"
|
|
"strconv"
|
|
"time"
|
|
|
|
"git.vakhrushev.me/av/transcriber/internal/contract"
|
|
"git.vakhrushev.me/av/transcriber/internal/metrics"
|
|
)
|
|
|
|
// Worker представляет базовый интерфейс для всех воркеров
|
|
type Worker interface {
|
|
Start(ctx context.Context)
|
|
Name() string
|
|
}
|
|
|
|
// pollInterval — пауза между прогонами шага. Полем, а не константой по месту:
|
|
// проверке нужен второй прогон, чтобы остановить воркер **после** того, как он
|
|
// рассудил об исходе первого. Отменять контекст изнутри шага она не может —
|
|
// отменённый контекст теперь и значит «нас остановили».
|
|
const pollInterval = time.Second
|
|
|
|
type CallbackWorker struct {
|
|
name string
|
|
// Шаг принимает контекст воркера: остановка обязана доходить до чужой
|
|
// работы, которую шаг завёл, а не только прерывать цикл между шагами.
|
|
f func(ctx context.Context) error
|
|
logger *slog.Logger
|
|
interval time.Duration
|
|
}
|
|
|
|
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,
|
|
interval: pollInterval,
|
|
}
|
|
}
|
|
|
|
func (w *CallbackWorker) Name() string {
|
|
return w.name
|
|
}
|
|
|
|
func (w *CallbackWorker) Start(ctx context.Context) {
|
|
w.logger.Info("Worker started", "worker", w.Name())
|
|
|
|
for {
|
|
select {
|
|
case <-ctx.Done():
|
|
w.logger.Info("Worker received shutdown signal", "worker", w.Name())
|
|
return
|
|
default:
|
|
err := w.f(ctx)
|
|
// Признак узнаётся по смыслу, а не по точной форме значения:
|
|
// приведение типа видело только вершину цепочки и сломалось бы от
|
|
// первой же обёртки `%w`, которая в проекте — умолчание.
|
|
var noop *contract.NoopJobError
|
|
isNoop := errors.As(err, &noop)
|
|
// Остановка — не отказ шага: контекст отменили мы сами. Считать её
|
|
// в метрику и писать владельцу «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 && !stopped {
|
|
w.logger.Error("Worker error", "worker", w.Name(), "error", err)
|
|
}
|
|
if stopped {
|
|
w.logger.Info("Worker step interrupted by shutdown", "worker", w.Name())
|
|
}
|
|
|
|
// Ждем перед следующей итерацией
|
|
select {
|
|
case <-ctx.Done():
|
|
w.logger.Info("Worker received shutdown signal during sleep", "worker", w.Name())
|
|
return
|
|
case <-time.After(w.interval):
|
|
// Продолжаем работу
|
|
}
|
|
}
|
|
}
|
|
}
|