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): // Продолжаем работу } } } }