- audiorecords вместо transcribe_jobs: приложения (texts, structures, recognitions, record_events, topics) живут своими коллекциями, ссылки на исходник и на приведённую копию перестали переставляться - рубеж называет достигнутое, отказ стал признаком остановки с причиной, а сторожей стало двое: число отказов и время в рубеже - воркеры потеряли специализацию, их число задаётся [pipeline] workers, шаг выбирается по рубежу, а захват отдаёт идентификатор и признак захвата
130 lines
3.9 KiB
Go
130 lines
3.9 KiB
Go
package worker
|
|
|
|
import (
|
|
"context"
|
|
"log/slog"
|
|
"sync/atomic"
|
|
"testing"
|
|
"time"
|
|
)
|
|
|
|
// Пул — предмет этой работы: воркеров стало сколько угодно одинаковых вместо
|
|
// трёх именованных. Проверки ниже судят саму обвязку — подъём, остановку и
|
|
// нулевой размер, — а не шаг, который она крутит: шаг судят проверки конвейера.
|
|
|
|
// countingStep считает свои вызовы и отпускает проверку, когда их набралось
|
|
// достаточно.
|
|
func countingStep(t *testing.T, enough int64) (func(context.Context) error, <-chan struct{}, *atomic.Int64) {
|
|
t.Helper()
|
|
|
|
var calls atomic.Int64
|
|
done := make(chan struct{})
|
|
var closed atomic.Bool
|
|
|
|
return func(context.Context) error {
|
|
if calls.Add(1) >= enough && closed.CompareAndSwap(false, true) {
|
|
close(done)
|
|
}
|
|
return nil
|
|
}, done, &calls
|
|
}
|
|
|
|
// Пул поднимает столько воркеров, сколько ему назвали, и все они крутят шаг.
|
|
func TestPoolRunsEveryWorker(t *testing.T) {
|
|
const size = 4
|
|
|
|
step, done, calls := countingStep(t, size)
|
|
pool := NewPool(size, step, slog.New(slog.DiscardHandler))
|
|
for _, w := range pool.workers {
|
|
w.interval = time.Millisecond
|
|
}
|
|
|
|
if pool.Size() != size {
|
|
t.Fatalf("в пуле %d воркеров вместо %d", pool.Size(), size)
|
|
}
|
|
|
|
ctx, cancel := context.WithCancel(context.Background())
|
|
defer cancel()
|
|
|
|
finished := make(chan struct{})
|
|
go func() {
|
|
pool.Start(ctx)
|
|
close(finished)
|
|
}()
|
|
|
|
select {
|
|
case <-done:
|
|
case <-time.After(5 * time.Second):
|
|
t.Fatalf("шаг позвали %d раз вместо %d: не все воркеры поднялись", calls.Load(), size)
|
|
}
|
|
|
|
cancel()
|
|
|
|
select {
|
|
case <-finished:
|
|
case <-time.After(5 * time.Second):
|
|
t.Fatal("пул не дождался остановки воркеров: горутина осталась висеть")
|
|
}
|
|
}
|
|
|
|
// Нулевой пул — законный режим, а не поломка: сервис поднимается, записи
|
|
// принимаются и не двигаются. Проверка судит именно это: шаг не зовётся ни разу,
|
|
// а подъём не блокируется.
|
|
func TestZeroPoolRunsNothingAndReturns(t *testing.T) {
|
|
var calls atomic.Int64
|
|
pool := NewPool(0, func(context.Context) error {
|
|
calls.Add(1)
|
|
return nil
|
|
}, slog.New(slog.DiscardHandler))
|
|
|
|
if pool.Size() != 0 {
|
|
t.Fatalf("пустой пул завёл %d воркеров", pool.Size())
|
|
}
|
|
|
|
finished := make(chan struct{})
|
|
go func() {
|
|
pool.Start(context.Background())
|
|
close(finished)
|
|
}()
|
|
|
|
select {
|
|
case <-finished:
|
|
case <-time.After(5 * time.Second):
|
|
t.Fatal("пустой пул не вернул управление: подъём сервиса заблокирован")
|
|
}
|
|
|
|
if got := calls.Load(); got != 0 {
|
|
t.Errorf("шаг позвали %d раз при нулевом пуле", got)
|
|
}
|
|
}
|
|
|
|
// Отменённый контекст останавливает **всех** воркеров пула: забытая горутина не
|
|
// падает и не пишет, а держит процесс и продолжает опрашивать базу после
|
|
// остановки сервиса.
|
|
func TestPoolStopsEveryWorkerOnCancel(t *testing.T) {
|
|
const size = 3
|
|
|
|
step, done, _ := countingStep(t, size)
|
|
pool := NewPool(size, step, slog.New(slog.DiscardHandler))
|
|
for _, w := range pool.workers {
|
|
w.interval = time.Millisecond
|
|
}
|
|
|
|
ctx, cancel := context.WithCancel(context.Background())
|
|
|
|
finished := make(chan struct{})
|
|
go func() {
|
|
pool.Start(ctx)
|
|
close(finished)
|
|
}()
|
|
|
|
<-done
|
|
cancel()
|
|
|
|
select {
|
|
case <-finished:
|
|
case <-time.After(5 * time.Second):
|
|
t.Fatal("пул не остановился по отмене контекста")
|
|
}
|
|
}
|