Разбор файла · по блокам

workflows.py

Главный файл репозитория: здесь живёт «дирижёр» — класс DevPipeline, который решает, кого запустить и что делать дальше.

Оригинал: temporal/worker/workflows.py (836 строк) · ← на главную · про спеку (ТЗ) · про ЧЕЛОВЕК · про Gate

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

Сначала карта файла

Файл отвечает на один вопрос: «что происходит от старта до финиша». В нём пять частей:

1. Что такое «шаг»

описание шага: промпт, какая джоба, пишет ли репозиторий.

2. Главный сценарий

run() — точка входа: подготовка спеки (ТЗ) → шаги DEV-AGENT.

3. Как выполняется шаг

_run_step() — запустить SPEC-AGENT или DEV-AGENT, подождать, забрать результат.

4. Входы извне

сигналы agent_done и human_verdictSPEC-AGENT/DEV-AGENT или ЧЕЛОВЕК «стучат в дверь».

5. Gate спеки (ТЗ)

_run_spec_review_gate_or_none()ЧЕЛОВЕК проверяет спеку (ТЗ) до разработки.

Блок кода 1. Что такое «шаг»

Шаг — это «один подход агента» (SPEC-AGENT или DEV-AGENT)

Класс Step — просто паспорт одного запуска. Он ничего не делает сам, только хранит данные: какой промпт дать агенту, какую джобу запустить, и можно ли этому агенту писать в репозиторий (коммиты, MR).

«Писать в репозиторий» — важно: DEV-AGENT, который пишет код, получает временный ключ доступа. SPEC-AGENT, который только читает, — не получает.

Важно про значения полей — что из них ограниченный набор, а что свободная строка:

  • step_nameсвободное имя шага, а не тип. Это просто строка-метка. Сейчас в пайплайне три имени ("spec", "step-1-intro", "step-2-followup"), но код не ограничивает — можно назвать шаг как угодно.
  • job_nameимя зарегистрированной Nomad-джобы. Сейчас пайплайн запускает только две (spec-agent и dev-agent), но значение — просто имя: запустить можно любую джобу, которая реально зарегистрирована в Nomad (в nomad/ лежат и mock-agent, и goose-agent, и другие).
  • consumerстрогий enum (ограниченный набор ролей). Код принимает только значения из перечисления Consumer: SPEC_AGENT, DEV_AGENT, REVIEW_AGENT и другие — произвольную строку сюда не передать.
  • factoryстрогий enum из двух значений: DEV или PRODUCT. Это «какая фабрика» выполняет шаг. Сейчас весь пайплайн работает в фабрике DEV (dev-фабрика агентов-разработчиков); PRODUCT зарезервировано под будущую продуктовую линию. Подробнее — ниже.
  • writes_repo / uses_task_package — просто флаги True/False.

То есть «spec» в step_name — это не тип спеки и не ограничение «только spec и task»: это имя конкретного шага (здесь — шага, который пишет спеку), и таких имён может быть сколько угодно.

Что такое factory — подробнее

factory отвечает на вопрос: «в рамках какой фабрики работает этот шаг?». Одна и та же роль (например, агент-разработчик) может работать в разных фабриках, поэтому принадлежность к фабрике хранится отдельно от роли consumer:

ПолеВопросПример
consumerкто именно потребляет токеныSPEC-AGENT, DEV-AGENT, REVIEW-AGENT
factoryв рамках какой фабрики (линии)DEV или PRODUCT

Сейчас в перечислении Factory ровно два значения:

  • DEV — dev-фабрика: агенты-разработчики и спека-агенты, которые делают разработку. Вся текущая работа идёт здесь.
  • PRODUCT — продуктовая фабрика: зарезервирована под будущую линию (продуктовые агенты), пока ни один шаг её не использует.

Зачем это нужно. factory вместе с consumer уходит в телеметрию и учёт токенов: код собирает «паспорт шага» (correlation), куда входят workflow, шаг, роль и фабрика, и прокидывает его наружу — в OTel-атрибуты и в заголовки запросов к LLM-шлюзу. Благодаря этому по логам и метрикам видно не только «какой агент жёг токены», но и «в рамках какой фабрики». Если бы фабрики не было, нельзя было бы разделить стоимость разработки и будущей продуктовой линии — всё слилось бы в одну кучу под consumer.

Значение строгое (enum), а не свободная строка: код принимает только DEV или PRODUCT, и не даст случайно написать factory="marketing". Так телеметрия остаётся чистой — неожиданное значение не просочится в метрики.

Step(
    step_name = "spec",                    # свободное имя шага (строка-метка, не тип!)
    prompt    = "Сгенерируй спеку (ТЗ)…",      # задание SPEC-AGENT
    job_name  = "spec-agent",              # имя Nomad-джобы: сейчас spec-agent или dev-agent
    consumer  = SPEC_AGENT,                # строгий enum Consumer (SPEC_AGENT/DEV_AGENT/…)
    factory   = DEV,                       # строгий enum Factory: DEV или PRODUCT
    writes_repo = False,                   # НЕ пишет в репозиторий → нет ключа
    uses_task_package = True,              # вход — пакет задачи, а не спека
)
настоящий код
@dataclass(frozen=True, slots=True)
class Step:
    step_name: str
    prompt: str
    job_name: str
    consumer: Consumer
    factory: Factory
    writes_repo: bool = True
    uses_task_package: bool = False

SPEC_STEP = Step(
    step_name="spec",
    prompt="Generate a technical specification from the provided task package.",
    job_name="spec-agent",
    consumer=Consumer.SPEC_AGENT,
    factory=Factory.DEV,
    writes_repo=False,
    uses_task_package=True,
)

Блок кода 2. Главный сценарий run()

Точка входа: «сначала спека (ТЗ), потом работа»

Всё начинается здесь. run() — метод, который Temporal запускает, когда приходит задача. У него два варианта:

  • Пришёл только номер спеки — спека уже готова и одобрена, сразу к разработке.
  • Пришёл номер задачи — сначала надо сделать спеку (ТЗ): SPEC-AGENT пишет черновик, ЧЕЛОВЕК его проверяет, и только потом код едет к DEV-AGENT.

Дальше — единый путь: проверить спеку, узнать «от какой версии кода отталкиваться» (base_sha), и прогнать шаги DEV-AGENT по очереди.

Почему «сейчас их два» и что это за шаги

У DEV-AGENT работа идёт не одним большим запуском, а серией маленьких шагов: агент запускается отдельно на каждый шаг, каждый раз с новым промптом. Эти шаги перечислены в run() списком steps — и workflow просто проходит по нему по очереди.

Сейчас в этом списке ровно два шага, и оба — демонстрационные:

ШагПромпт (задание)
step-1-intro«Ответь одним предложением: кто ты и какая ты модель?»
step-2-followup«Назови одно преимущество durable execution для пайплайнов агентов»

Задания не про настоящую разработку — они учебные, чтобы проверить сам механизм: workflow умеет запустить DEV-AGENT, дождаться его, сохранить вывод и перейти к следующему шагу. В бою эти демо-промпты заменят настоящими шагами, собранными из одобренной спеки (ТЗ): «сделай X», «покрой тестами», «открой MR» и так далее. Число шагов тогда вырастет — список не ограничен двумя, их будет столько, сколько нужно для задачи.

Каждый шаг из списка проходит одинаковый цикл: запустить DEV-AGENT → сохранить его вывод (артефакт-результат) → проверить на Gate → PASS значит «шаг готов», дальше следующий; не-PASS — пайплайн останавливается с понятной причиной.

Что именно «приходит» в Temporal

На входе у run() — не документ и не пакет, а два строковых аргумента: идентификаторы (UUID) артефактов из artifacts-service. Один из них может быть пустым — по нему workflow понимает, какой из двух путей запускать:

АргументТипЧто этоКогда заполнен
spec_idстрока (UUID)номер готовой спеки (ТЗ), артефакт kind=specпуть 1 — спека уже одобрена
task_idстрока (UUID)номер задачи, артефакт kind=task; из неё будет собрана спекапуть 2 — спеки ещё нет

Workflow выбирает путь просто: если task_id непустой — надо сделать спеку (ТЗ) с нуля; если пустой — используем spec_id, который уже пришёл готовым. Никакой другой структуры на вход не передаётся — вся «тяжёлая» информация (тело спеки, пакет задачи) подтягивается самим workflow из artifacts-service уже внутри прогона.

Как это выглядит в коде запуска

Так workflow стартуют тесты и реальные вызовы (клиент Temporal → start_workflow):

# путь 1: спека уже есть — передаём только spec_id
handle = await client.start_workflow(
    DevPipeline.run,
    "11111111-1111-1111-1111-111111111111",   # spec_id (UUID спеки)
    id="dev-pipeline-1",
    task_queue="dev-factory",
)

# путь 2: спеки нет — передаём только task_id
handle = await client.start_workflow(
    DevPipeline.run,
    args=[
        "",                                     # spec_id пустой — не используется
        "33333333-3333-3333-3333-333333333333", # task_id (UUID задачи)
    ],
    id="dev-pipeline-2",
    task_queue="dev-factory",
)
run(spec_id, task_id=""):
    ЕСЛИ есть task_id:                       # надо сначала сделать спеку (ТЗ)
        SPEC-AGENT пишет черновик спеки (ТЗ)  # _run_step(SPEC_STEP)
        человек проверяет спеку (ТЗ) (Gate)    # approve? → едем дальше
                                           # reject? → SPEC-AGENT переделывает
                                           # возвратов уже было 2 (потолок)?
                                           #   → эскалация человеку, цикл стоп

    проверить спеку                          # существует? одобрена?
    узнать base_sha                          # от какой версии кода работаем

    # Дальше — сам процесс разработки: список шагов DEV-AGENT.
    # Сейчас в списке ДВА демо-шага (step-1-intro, step-2-followup) —
    # просто чтобы проверить, что пайплайн шагает по кругу.
    # В бою эти демо-промпты заменят настоящими шагами из ТЗ.
    ДЛЯ каждого шага DEV-AGENT:
        запустить DEV-AGENT                  # _run_step
        сохранить его вывод                  # persist_result
        проверить результат (Gate)           # PASS → дальше, иначе стоп
        вернуть {"status": "success", "completed_steps": [...]}
настоящий код
@workflow.run
async def run(self, spec_id: str, task_id: str = "") -> dict:
    if task_id:
        self._task_id = task_id
        spec_step_result = await self._run_step(SPEC_STEP)
        if spec_step_result.exit_code != "0":
            return {"status": "failed", "failed_step": SPEC_STEP.step_name,
                    "exit_code": spec_step_result.exit_code,
                    "reason": spec_step_result.reason, "completed_steps": []}

        spec_gate_failure = await self._run_spec_review_gate_or_none(spec_step_result.result)
        if spec_gate_failure is not None:
            return spec_gate_failure
        spec_id = self._spec_id

    validation_failure = await self._validate_spec_or_none(spec_id)
    if validation_failure is not None:
        return validation_failure
    self._spec_id = spec_id

    base_sha_failure = await self._read_base_sha_or_none(spec_id)
    if base_sha_failure is not None:
        return base_sha_failure

    steps = [ Step(step_name="step-1-intro", ...), Step(step_name="step-2-followup", ...) ]

    completed = []
    for step in steps:
        result = await self._run_step(step)
        persist_output = await workflow.execute_activity("persist_result", ...)
        verdict = await run_gate(GateSpec(name=f"{step_name}-exit", ...), ...)
        if verdict.outcome != Outcome.PASS:
            return {"status": "failed", "failed_step": step_name, ...}
        completed.append({"step": step_name, "result": result.result})
    return {"status": "success", "completed_steps": completed}

Блок кода 3. Как выполняется один шаг

Сердце системы: _run_step()

Это самый важный механизм. Один шаг = запустить агента (SPEC-AGENT или DEV-AGENT), дождаться, забрать результат. Хитрость в том, что «ждать» умеет сам Temporal: пока запущенный агент работает (минуты!), код ничего не делает, но и не забывает — это и есть durable. Когда агент закончил, он шлёт сигнал agent_done.

Плюс две страховки:

  • Если запущенный агент «молчит» — код сам сходит в Nomad и спросит: задача вообще запускалась? завершилась? зависла?
  • Если агент пишет в репозиторий (сейчас это DEV-AGENT) — перед запуском выдаётся временный ключ, а после шага он всегда отзывается (даже если шаг упал).
_run_step(step):
    сбросить прошлый результат
    вычислить step_id (номер шага в пайплайне)

    ЕСЛИ шаг пишет в репозиторий:
        выпустить временный ключ доступа (deploy key)

    запустить SPEC-AGENT или DEV-AGENT через activity  # dispatch в Nomad
    ЖДАТЬ сигнал agent_done                  # минуты, durable-ожидание

    ЕСЛИ запущенный агент завершился успешно И пишет в репо:
        создать MR (merge request)

    ВСЕГДА (даже при ошибке):
        отозвать временный ключ

    вернуть: {exit_code, result}
настоящий код (сокращённо)
async def _run_step(self, step: Step, max_turns: str = "8") -> AgentResult:
    self._agent_result = None
    step_id = self._last_step_id = correlation.step_id

    key_id = None
    try:
        if step.writes_repo:                    # временный ключ для записи
            key_id = (await workflow.execute_activity("mint_deploy_key", ...))["key_id"]

        await workflow.execute_activity("prepare_agent_run", PrepareAgentRunInput(...))

        try:                                    # durable-ожидание сигнала
            await workflow.wait_condition(lambda: self._agent_result is not None,
                                          timeout=_FETCH_FAIL_PROBE_TIMEOUT)
        except asyncio.TimeoutError:
            alloc = await workflow.execute_activity("read_alloc_status", ...)
            if _is_infra_delivery_failure(alloc):
                self._agent_result = _resolve_timeout_fallback(alloc)
            else:                               # ждём дальше, до 30 минут всего
                try:
                    await workflow.wait_condition(..., timeout=_STEP_TIMEOUT - _FETCH_FAIL_PROBE_TIMEOUT)
                except asyncio.TimeoutError:
                    alloc = await workflow.execute_activity("read_alloc_status", ...)
                    self._agent_result = _resolve_timeout_fallback(alloc)

        if step.writes_repo and self._agent_result.exit_code == "0":
            await workflow.execute_activity("create_mr", ...)   # создать MR

    finally:
        if step.writes_repo and key_id is not None:              # отозвать ключ ВСЕГДА
            await workflow.execute_activity("revoke_deploy_key", ...)
    return self._agent_result

Блок кода 4. «Стук в дверь» — сигналы извне

Как внешний мир сообщает коду результат

Код сам ничего не опрашивает. Вместо этого к нему приходят сигналы — Temporal записывает их в историю и «будит» сценарий. Два таких входа:

agent_done(exit_code, result, step_id):      # ← SPEC-AGENT или DEV-AGENT закончил шаг
    ЕСЛИ step_id не совпадает с ожидаемым:    # опоздавший/чужой сигнал
        проигнорировать (с предупреждением)
    ИНАЧЕ:
        запомнить результат агента (его result.json)

human_verdict(artifact_id, step_id, feedback): # ← человек вынес вердикт
    # только «будит» код — статус меняет сам artifacts-сервис через POST /status
    ЕСЛИ нет Gate, который ждёт вердикт:  проигнорировать
    ЕСЛИ artifact_id не тот:               проигнорировать
    ЕСЛИ такой же step_id уже был:         проигнорировать (дубликат)
    ИНАЧЕ:
        запомнить feedback (текст замечаний)
        разбудить ожидающий Gate
настоящий код
@workflow.signal
async def agent_done(self, exit_code: str, result: dict, step_id: str = "") -> None:
    if step_id != self._last_step_id:
        workflow.logger.warning("signal agent_done: stale/unexpected step_id, ignoring: ...")
        return
    self._agent_result = AgentResult(exit_code=exit_code, result=result,
                                     reason=FailureReason.AGENT_REPORT.value)

@workflow.signal
async def human_verdict(self, artifact_id="", step_id="", feedback="") -> None:
    ctx = self._active_gate_ctx
    if ctx is None:                        # нет Gate → нечего будить
        return
    expected = ctx.subject.get("artifact_id")
    if artifact_id and artifact_id != expected:   # не тот артефакт
        return
    if step_id and step_id == ctx.last_verdict_step_id:  # дубликат
        return
    if step_id:
        ctx.last_verdict_step_id = step_id
    ctx.human_feedback = feedback          # текст замечаний для rework
    ctx.notify_human_verdict()

Блок кода 5. Gate спеки (ТЗ): ЧЕЛОВЕК проверяет черновик

Что такое спека (ТЗ), откуда она берётся и как добирается до DEV-AGENT — на отдельной странице: про спеку (ТЗ) →
Петля «вернули → переделал → показали снова»

Когда SPEC-AGENT написал черновик спеки (ТЗ), его нельзя просто взять и использовать. Он отправляется на ревью ЧЕЛОВЕК. Дальше — цикл:

  • Одобрено → берём эту спеку (ТЗ) и едем к DEV-AGENT.
  • Вернули с замечаниямиSPEC-AGENT переделывает (замечания добавляются в его задание), показываем снова.
  • Слишком много возвратов (потолок = 2) → стоп, эскалация ЧЕЛОВЕК: дальше решает ЧЕЛОВЕК, а не цикл.

Важная деталь: каждая новая версия спеки (ТЗ) — новый артефакт. Старую не переписывают, её «хоронят» и создают новую. Счётчик возвратов (rework_count) можно спросить снаружи.

_run_spec_review_gate_or_none(результат SPEC-AGENT):
    вытащить spec_id из результата; ЕСЛИ его нет → fail
    rework_count = 0

    ЦИКЛ:
        отправить спеку на ревью человеку     # pending_review
        ЖДАТЬ вердикт человека (Gate)

        ЕСЛИ одобрено:
            вернуть "ок, спеку можно использовать"

        # иначе — reject
        ЕСЛИ уже было 2 возврата (rework_count >= MAX_REWORK):
            вернуть {status: "escalated",
                     reason: "спека отклонялась уже MAX_REWORK раз — дальше решает человек"}

        rework_count += 1
        SPEC-AGENT переделывает с замечаниями в промпте
        новая версия → снова в начало ЦИКЛА
настоящий код
async def _run_spec_review_gate_or_none(self, spec_step_result):
    self._rework_count = 0
    spec_id_result = extract_spec_id(spec_step_result)
    if isinstance(spec_id_result, Err):
        return {"status": "failed", "failed_step": "spec-review", "error": ...}
    self._spec_id = str(spec_id_result.ok_value)

    while True:
        await workflow.execute_activity("submit_spec_for_review",
                                        SubmitSpecForReviewInput(spec_id=self._spec_id), ...)
        verdict = await run_gate(GateSpec(name="spec-review",
                                          verifier=HumanVerifier(), max_rework=0),
                                 GateContext(subject={"artifact_id": self._spec_id}))
        if verdict.outcome == Outcome.PASS:
            return None                       # одобрено — можно работать

        if self._rework_count >= MAX_REWORK:  # уже 2 возврата (потолок) — стоп, дальше человек
            return {"status": "escalated", "failed_step": "spec-review",
                    "reason": f"spec rework ceiling reached: max_rework={MAX_REWORK}",
                    "rework_count": self._rework_count, "completed_steps": []}

        self._rework_count += 1
        rework_step = replace(SPEC_STEP, prompt=f"{SPEC_STEP.prompt}\n\nРевьюер вернул спеку (ТЗ):\n{feedback}")
        rework_result = await self._run_step(rework_step)
        ...
        self._spec_id = новый spec_id         # новая версия → в начало цикла

Соберём всё вместе

Один запуск целиком, от начала до конца:

  1. Есть задача → запускается SPEC-AGENT, пишет черновик спеки (ТЗ).
  2. Черновик уходит ЧЕЛОВЕК на проверку. Не понравилось — возврат, переделка (до 2 раз).
  3. Спека (ТЗ) одобрена → код узнаёт base_sha (от какой версии кода работать).
  4. Шаг за шагом запускается DEV-AGENT: запустил → подождал → проверил → дальше.
  5. Каждый результат сохраняется; провал любого шага останавливает всё с понятной причиной.
  6. Финал: {"status": "success", "completed_steps": [...]}.
Три идеи, которые стоит унести с собой: (1) код «спит» и просыпается от сигналов, а не опрашивает — это durable; (2) ЧЕЛОВЕК встроен в цикл там, где нужно решение (Gate); (3) временные ключи живут ровно один шаг — всегда выдаются и всегда отзываются.