внутренняя модель перестроена вокруг аудиозаписи
- audiorecords вместо transcribe_jobs: приложения (texts, structures, recognitions, record_events, topics) живут своими коллекциями, ссылки на исходник и на приведённую копию перестали переставляться - рубеж называет достигнутое, отказ стал признаком остановки с причиной, а сторожей стало двое: число отказов и время в рубеже - воркеры потеряли специализацию, их число задаётся [pipeline] workers, шаг выбирается по рубежу, а захват отдаёт идентификатор и признак захвата
This commit is contained in:
@@ -3,26 +3,25 @@ package worker
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"log/slog"
|
||||
"strconv"
|
||||
"sync"
|
||||
"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
|
||||
|
||||
// CallbackWorker крутит один и тот же шаг, опрашивая очередь.
|
||||
//
|
||||
// Специализации у него нет: шаг сам берёт любую пригодную к работе запись и
|
||||
// выбирает работу по её рубежу. Раньше воркеров было три именованных, по одному
|
||||
// на состояние, и каждый новый рубеж требовал четвёртого.
|
||||
type CallbackWorker struct {
|
||||
name string
|
||||
// Шаг принимает контекст воркера: остановка обязана доходить до чужой
|
||||
@@ -64,15 +63,12 @@ func (w *CallbackWorker) Start(ctx context.Context) {
|
||||
// первой же обёртки `%w`, которая в проекте — умолчание.
|
||||
var noop *contract.NoopJobError
|
||||
isNoop := errors.As(err, &noop)
|
||||
// Остановка — не отказ шага: контекст отменили мы сами. Считать её
|
||||
// в метрику и писать владельцу «Worker error» значит красить каждую
|
||||
// выкладку как поломку — по тому же доводу, по которому не считается
|
||||
// Остановка — не отказ шага: контекст отменили мы сами. Писать
|
||||
// владельцу «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)
|
||||
}
|
||||
@@ -80,6 +76,9 @@ func (w *CallbackWorker) Start(ctx context.Context) {
|
||||
w.logger.Info("Worker step interrupted by shutdown", "worker", w.Name())
|
||||
}
|
||||
|
||||
// Счётчик работы растит сам шаг: только он знает рубеж, с которого
|
||||
// взята запись, а воркер к рубежу больше не привязан.
|
||||
|
||||
// Ждем перед следующей итерацией
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
@@ -91,3 +90,51 @@ func (w *CallbackWorker) Start(ctx context.Context) {
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// 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()
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user