package worker import ( "context" "errors" "fmt" "log/slog" "sync" "time" "git.vakhrushev.me/av/transcriber/internal/contract" ) // pollInterval — пауза между прогонами шага. Полем, а не константой по месту: // проверке нужен второй прогон, чтобы остановить воркер **после** того, как он // рассудил об исходе первого. Отменять контекст изнутри шага она не может — // отменённый контекст теперь и значит «нас остановили». const pollInterval = time.Second // CallbackWorker крутит один и тот же шаг, опрашивая очередь. // // Специализации у него нет: шаг сам берёт любую пригодную к работе запись и // выбирает работу по её рубежу. Раньше воркеров было три именованных, по одному // на состояние, и каждый новый рубеж требовал четвёртого. 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 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): // Продолжаем работу } } } } // Pool — пул одинаковых воркеров конвейера. // // Число задаётся настройкой, и ноль — законное значение: сервис поднимается, // записи принимаются и не двигаются. Это режим, а не поломка: он нужен местному // запуску и выкладке, где конвейер надо остановить, не роняя приём. type Pool struct { workers []*CallbackWorker logger *slog.Logger } // NewPool собирает пул из size одинаковых воркеров, крутящих один и тот же шаг. func NewPool(size int, step func(ctx context.Context) error, logger *slog.Logger) *Pool { if logger == nil { logger = slog.Default() } workers := make([]*CallbackWorker, 0, size) for i := range size { workers = append(workers, NewCallbackWorker(fmt.Sprintf("pipeline_worker_%d", i+1), step, logger)) } return &Pool{workers: workers, logger: logger} } // Size — сколько воркеров в пуле. func (p *Pool) Size() int { return len(p.workers) } // Start поднимает всех воркеров пула и ждёт их остановки. func (p *Pool) Start(ctx context.Context) { if len(p.workers) == 0 { // Молчать нельзя: пустой пул неотличим от поломки, а объявленный режим // обязан быть назван. p.logger.Info("Pipeline workers are disabled by configuration") return } var wg sync.WaitGroup for _, w := range p.workers { wg.Add(1) go func(worker *CallbackWorker) { defer wg.Done() worker.Start(ctx) p.logger.Info("Worker stopped gracefully", "worker", worker.Name()) }(w) } wg.Wait() }