- audiorecords вместо transcribe_jobs: приложения (texts, structures, recognitions, record_events, topics) живут своими коллекциями, ссылки на исходник и на приведённую копию перестали переставляться - рубеж называет достигнутое, отказ стал признаком остановки с причиной, а сторожей стало двое: число отказов и время в рубеже - воркеры потеряли специализацию, их число задаётся [pipeline] workers, шаг выбирается по рубежу, а захват отдаёт идентификатор и признак захвата
141 lines
5.6 KiB
Go
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()
|
|
}
|