Перейти к содержимому

Планировщик

Планировщик выполняет процессы по принципу work-stealing. Воркеры поддерживают локальные деки и заимствуют работу друг у друга при простое.

Планировщик работает с любым типом, реализующим интерфейс Process:

type Process interface {
Init(ctx context.Context, method string, input payload.Payloads) error
Step(events []Event, out *StepOutput) error
Close()
}
МетодНазначение
InitПодготовка процесса с именем метода и входными аргументами
StepПродвижение автомата состояний входящими событиями, запись yields в output
CloseОсвобождение ресурсов

Параметр method в Init указывает, какую точку входа вызывать. Экземпляр процесса может предоставлять несколько точек входа, и вызывающий выбирает нужную. Это также служит проверкой корректности инициации процесса планировщиком.

Планировщик вызывает Step() многократно, передавая события (завершения yield’ов, сообщения) и собирая yield’ы (команды для диспатчинга). Процесс записывает свой статус и yield’ы в буфер StepOutput.

type Event struct {
Type EventType // EventYieldComplete или EventMessage
Tag uint64 // Корреляционный тег для завершений yield
Data any // Данные результата или payload сообщения
Error error // Ошибка, если yield завершился неудачей
}

Планировщик порождает GOMAXPROCS воркеров по умолчанию. У каждого воркера локальный дек для кеш-дружественного LIFO-доступа. Глобальная FIFO-очередь обрабатывает новые submissions и межворкерные переносы. Процессы отслеживаются по PID для маршрутизации сообщений.

flowchart TD
W[Воркер ищет работу] --> L{Локальный дек?}
L -->|есть элементы| LP[Pop со дна LIFO]
L -->|пуст| G{Глобальная очередь?}
G -->|есть элементы| GP[Pop + пакетный перенос до 16]
G -->|пуста| S[Заимствовать у случайной жертвы]
S --> SH[StealHalfInto дека жертвы]

Воркеры проверяют источники в порядке приоритета:

ПриоритетИсточникПаттерн
1Локальный декLIFO pop, lock-free, кеш-дружественный
2Глобальная очередьFIFO pop с пакетным переносом
3Другие воркерыЗаимствовать половину из дека жертвы

При pop из глобальной очереди воркеры берут один элемент и пакетно переносят до 16 дополнительных в свой локальный дек.

У каждого воркера дек Chase-Lev для work-stealing:

type Deque struct {
buffer atomic.Pointer[dequeBuffer]
top atomic.Int64 // Воры заимствуют отсюда (CAS)
bottom atomic.Int64 // Владелец push/pop отсюда
}

Владелец push’ит и pop’ит со дна (LIFO) без синхронизации. Воры заимствуют с верха (FIFO) через CAS. Это даёт владельцу кеш-дружественный доступ к недавно добавленным элементам, распределяя более старую работу ворам.

StealHalfInto забирает половину элементов одной CAS-операцией, снижая конкуренцию.

Перед блокировкой на condition variable воркеры адаптивно вращаются:

Счётчик вращенийДействие
< 4Тесный цикл
4-15Уступить поток (runtime.Gosched)
>= 16Заблокироваться на condition variable
stateDiagram-v2
[*] --> Ready: Submit
Ready --> Running: CAS воркером
Running --> Complete: done
Running --> Blocked: yields commands
Running --> Idle: ожидание сообщений
Blocked --> Ready: CompleteYield
Idle --> Ready: Send прибыл
СостояниеОписание
ReadyВ очереди на выполнение
RunningВоркер выполняет Step()
BlockedОжидание завершения yield
IdleОжидание сообщений
CompleteВыполнение завершено

Флаг wakeup обрабатывает гонки: если обработчик вызывает CompleteYield, пока воркер ещё владеет процессом (Running), он устанавливает флаг. Воркер проверяет флаг после диспатчинга и перепланирует, если установлен.

У каждого процесса MPSC (multi-producer, single-consumer) очередь событий:

  • Producers: обработчики команд (CompleteYield), отправители сообщений (Send)
  • Consumer: воркер сливает события в Step()

Планировщик реализует relay.Receiver для маршрутизации сообщений процессам. При вызове Send() он находит целевой PID в карте byPID, push’ит сообщение как событие в очередь процесса и будит процесс, если тот Idle, push’ом в глобальную очередь.

При завершении планировщик отправляет события отмены всем выполняющимся процессам и ждёт их завершения или таймаута. Воркеры завершаются, когда работы больше нет.