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

Как работает движок прогонов

Проблема: исполнение пайплайна — долгий, склонный к отказам процесс, трогающий LLM-провайдеров, базу и сборщик сайтов. Если его состояние живёт в памяти того процесса, которому довелось его исполнять, краш теряет правду о случившемся, а два процесса, пишущие «состояние», дают две правды.

Слой 1 — один вход, стримы между, один писатель

Заголовок раздела «Слой 1 — один вход, стримы между, один писатель»

Создание прогона — единый шов: дубль-предохранитель, бюджетный preflight, резолв версии из каталога, строка pipeline_run — затем джоб-сообщение в Redis Stream. Исполнение отвечает только событиями.

SSE clientsPostgreSQLProjectorruns:events (stream)Native workerruns:jobs (stream)Platform APISSE clientsPostgreSQLProjectorruns:events (stream)Native workerruns:jobs (stream)Platform APIenqueue job (XADD)consumer group deliverstyped events: run_started,step_completed, run_succeeded…consumer group deliversapply_event (sole DB writer)fan-out via Redis Pub/Sub

Вес несут два решения:

  • Типизированные фабрики событий с JSON-Schema контрактом. Воркеры никогда не собирают event-словари руками — фабрики штампуют версию схемы и валидируют поля, а схема едет закоммиченным контрактом с фикстурами на каждый тип события. Мок-эмулятор веб-UI проигрывает те же фикстуры, так что фронтенд разрабатывается против точного wire-формата.
  • Проектор — единственный писатель. Всё остальное эмитит; применяет один потребитель. Идемпотентность структурна: счётчики шагов считаются как distinct терминальные шаги, а не инкременты, поэтому at-least-once редоставка не может посчитать дважды. Отравленное сообщение подтверждается и захватывается, а не ретраится вечно.

Слой 2 — состояния прогона и кто вправе их двигать

Заголовок раздела «Слой 2 — состояния прогона и кто вправе их двигать»
create_runworker picks upcancelstep error (fail-loud)watchdog: heartbeat lostqueuedrunningcancelledsucceededfailedinterrupted

Честно: сегодня эти переходы держатся конвенцией между писателями (проектор, вачдог, ретрай), а не формальной таблицей предохранителей — явный FSM с chaos-тестами измеренного восстановления это спроектированная, ещё не принятая эволюция. Что живо — самолечащаяся пара:

  • вачдог помечает прогоны с протухшим heartbeat как interrupted (упавший воркер не может оставить прогон «running» навсегда);
  • рипер возвращает брошенные джоб-сообщения из consumer-группы, переставляет их в очередь до ретрай-кэпа, затем — в dead-letter. Двое никогда не пересекаются: прогон, упавший по собственным заслугам, терминален и не трогается.

Слой 3 — инпуты замораживаются до первого шага

Заголовок раздела «Слой 3 — инпуты замораживаются до первого шага»

Инпуты прогона (профиль бренда, профиль автора, конфиг сбора) собираются один раз на старте в замороженный типизированный контракт с отредактированными секретами — шаги читают ctx.inputs и никогда не перечитывают базу посреди прогона. Правило, которое это обеспечивает: чистая граница между что тебе дали и что ты произвёл. Оно существует потому, что оба провала реально случались: бренд-гайдлайны, тихо обрезанные ad-hoc чтением, и конфиг-поле, мутированное шагом 1 и съеденное шагом 6 как «инпут».

Слой 4 — конкурентность, чувствующая провайдера

Заголовок раздела «Слой 4 — конкурентность, чувствующая провайдера»

Глобальный кэп LLM-конкурентности — не константа, а контур управления:

telemetryset capread capLLM adapterreads x-ratelimit-* headersRedisAutoscaler tickC* = min of provider ·DB pool · host · ceilingWorkers acquire slots

Sense — каждый живой LLM-ответ роняет rate-limit телеметрию в Redis (fail-soft: телеметрия не может сломать вызов). Decide — периодический тик считает кэп как минимум нескольких лимитеров и применяет AIMD (поднимай медленно, режь резко, мгновенный срез на 429). Actuate — кэп это Redis-ключ с TTL; умрёт контур управления — воркеры сами откатятся к статической настройке.

Почему несколько лимитеров, а не только провайдер: живой эксперимент нашёл тихую стену — выше некоторой конкурентности исчерпывался пул БД, пока LLM-провайдер и CPU выглядели зелёными, и пропускная способность падала в ноль. Связывающее ограничение на практике — пул, а не провайдер. Гигиена границ: контексту прогонов import-линтером запрещено лезть во внутренности интеграций, поэтому Sense и Decide общаются строго через Redis — контур управления уважает те же модульные границы, что и всё остальное.

Cancel-на-бегу сегодня фактически no-op (канал управления существует, но его никто не потребляет: queued-прогоны отменяются мгновенно, бегущие добегают или падают). Exactly-once — по построению, а не по формальному replay-тесту. Оба — известные, заложенные в дизайн зазоры, не сюрпризы.

POST /api/v1/runs запускает прогон; смотри его вживую по SSE-каналу прогонов или читай его строки pipeline_run / run_step — каждое утверждение выше видно в этих данных.

Спеки: SPEC-007 (фреймворк прогонов), SPEC-055 (типизированные события + мок-эмулятор), SPEC-062 (идемпотентный проектор), SPEC-071 (адаптивная конкурентность), SPEC-069/070 (замороженные контракты инпутов), SPEC-104 (формальный FSM — спроектирован).