Files
transcriber/internal/controller/worker/worker.go
T
av 1576d06735 внутренняя модель перестроена вокруг аудиозаписи
- audiorecords вместо transcribe_jobs: приложения (texts, structures,
  recognitions, record_events, topics) живут своими коллекциями, ссылки на
  исходник и на приведённую копию перестали переставляться
- рубеж называет достигнутое, отказ стал признаком остановки с причиной, а
  сторожей стало двое: число отказов и время в рубеже
- воркеры потеряли специализацию, их число задаётся [pipeline] workers, шаг
  выбирается по рубежу, а захват отдаёт идентификатор и признак захвата
2026-08-14 20:20:33 +03:00

141 lines
5.6 KiB
Go

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()
}