Перейти к содержанию

Проектирование workflow

Фреймворк отвечает только за механику: как выполняется шаг, как сцепляются сессии, что делать при сбое, как экономить контекст. Workflow пишете вы — фреймворк не знает, какой у вас проект и какой язык, и знать не должен. Эта страница о том, как спроектировать workflow; полные таблицы полей Step и Workflow — в Python API.

Какую задачу это решает

Один long-horizon-запуск не описывается одной строкой prompt: сначала выяснить требования, потом исследовать, потом реализовать, потом перепроверить — у каждого отрезка своя роль, свой контекст, свои условия приёмки. Если запихнуть всё в одну подсказку, модель сама решит, какой отрезок пропустить; если оформить это как workflow, порядок, условия выхода и передача состояния становятся кодом на Python — читаемым, тестируемым, с возможностью перезапустить только сломавшийся шаг.

Workflow делает ровно три вещи:

  • выполняет последовательность шагов по порядку;
  • определяет, что каждый шаг видит из предыдущих (три способа сцепки сессий + словарь ctx);
  • определяет, когда повторять попытку, а когда выходить досрочно.

Никаких предметных допущений внутри нет. Где резать, что принимать на каждом шаге, что делать при непрохождении — вот эти четыре вещи и есть «проектирование workflow».

Как этим пользоваться (минимальный код)

# flows.py
from flower import AgentSpec, Step, Workflow

terse = AgentSpec(
    name="terse",
    instructions="回答极简,一行以内,不解释不寒暄。",
    allowed_tools=["Read", "Glob"],
    max_turns=4,
)


def main() -> Workflow:
    return Workflow([
        # Новая сессия: видит только то, что передано в prompt
        Step("取词", terse, "读 seed.txt,只回文件里那个词。"),
        # Тоже новая сессия, результат предыдущего шага впрыснут в prompt (дёшево, защищает от загрязнения)
        Step("造句", terse, lambda ctx: f"用「{ctx['取词']}」造一个五字短句,只回短句。"),
    ])
flower run flows.py:main -w /path/to/repo

Аргумент flower run — это модуль:атрибут или путь_к_файлу:атрибут. Если полученный объект вызываемый, он сначала вызывается один раз; затем берётся Workflow и запускается. По завершении в терминал печатается суммарная стоимость и путь к манифесту запуска.

Можно написать и свой драйвер — первый аргумент Workflow.run это Runtime:

ctx = await wf.run(rt, on_step=lambda step, r: print(f"{step.name} ok={r.ok} ${r.cost_usd:.4f}"))

Что оно делает на самом деле

Что получает Step и что обязан вернуть

Step — это не функция, а декларация. Реально выполняется Runtime.run(step.spec, отрендеренный prompt, ...)один шаг = один вызов Runtime.run = одна сессия.

Первые три поля — позиционные аргументы, Step(name, spec, prompt):

  • name — имя шага. Оно же ключ в ctx, имя строки в runs/manifest.json и ключ межпроцессной родословной.
  • spec — какой AgentSpec использовать. Он задаёт для этого шага белый список инструментов, модель и бюджет.
  • promptstr либо (ctx) -> str. В вызываемом варианте получает текущий ctx; это самый дешёвый способ передать результат предыдущего шага (второй способ — сцепка сессий, см. ниже).

«Возвращает» шаг объект StepResult, но внутри workflow вы получаете две вещи:

  • ctx[step.name] — по умолчанию result.text, а если задан reduce — то, что вернул reduce;
  • ctx["_results"][step.name] — полный StepResult (стоимость, число ходов, число попыток, session_id).

result.text собирает только основной текст main thread: реплики subagent лежат в его собственном transcript, выданное ему задание имеет kind="prompt", а синтетическая ошибка обрыва — kind="error"; ни одно из трёх сюда не попадает.

reduce: это не синтаксический сахар

По умолчанию дальше передаётся дословный ответ модели. У некоторых шагов дословный ответ передавать дальше нельзя:

Step("确认需求", spec=确认者, prompt="帮我做一个 X",
     reduce=lambda r, ctx: ctx["_brief"].prompt_block())

На практике шаг уточнения требований помимо четырёх разделов вклеивает в ответ весь код целиком. Дальше по цепочке должны идти только разобранные четыре раздела, иначе вся эта груда кода попадёт в prompt следующего шага. clarify_step держится именно на этом поле.

reduce обязан быть синхронной функцией; gate / when / on_reject могут быть async.

Как состояние течёт через ctx

ctx — это dict[str, Any], сам объект Workflow.context. После каждого шага запись идёт по этой таблице:

Ситуация ctx[имя шага] Прочее
when(ctx) вернул False не пишется, весь шаг пропущен result не создаётся, в _results тоже не попадает
Прошёл reduce(result, ctx), а если не задан — result.text
Провал + on_fail="stop" (по умолчанию) не пишется пишется ctx["_failed_at"], весь workflow останавливается на этом шаге
Провал + on_fail="skip" не пишется выполнение продолжается дальше
Провал + on_fail="continue" result.text (неполный, reduce не применяется) выполнение продолжается дальше

Независимо от прохождения ctx["_results"][имя шага] пишется всегда; если result.session_id непустой, он дополнительно пишется в ctx["_sessions"] и заносится в родословную.

Успешен ли был workflow, определяется по ctx.get("_failed_at"), а не по тому, выдал ли последний шаг что-нибудь.

Все ключи, начинающиеся с подчёркивания, кладёт сам Workflow.run: _runtime, _on_event, _sessions, _results, _lineage, _woke, _aborted, _failed_at — не используйте их как имена своих шагов. Отдельные механизмы кладут свои (_brief / _goal / _verdict и т. д.), полный список — в Python API.

Из них _runtime и _on_event предназначены для gate: gate может сам отправить агента вынести вердикт, и процесс вынесения вердикта при этом всё равно уходит в UI — иначе на десяток с лишним секунд интерфейс чернеет и выглядит зависшим. Страж цели реализован именно так.

ctx — один и тот же dict: при втором запуске того же объекта Workflow ключи от прошлого раза остаются на месте. Нужен чистый старт — создайте новый объект или явно передайте context={}.

При on_fail=skip ctx[имя шага] не пишется

Ниже по цепочке lambda ctx: ctx["某步"] даст прямой KeyError. Чтобы идти дальше с неполным результатом, используйте on_fail="continue"; если действительно нужно пропустить, нижестоящие шаги должны сами подстраховываться через ctx.get(...).

Вердикт и возврат на доработку: gate, on_reject, StepAbort

gate(result, ctx) -> bool решает вопрос «отработал, но годится ли». Две детали, которые надо знать обязательно:

  • если result.ok ложно, gate вообще не вызывается (короткое замыкание);
  • на каждую попытку он вызывается ровно один раз, результат сохраняется для дальнейшего использования — у него могут быть побочные эффекты. Gate у clarify_step записывает бриф на диск, повторный вызов — повторная запись.

Как всё повторяется после непрохождения gate, зависит от того, задан ли on_reject:

Как идёт следующий круг Имя в манифесте
Только retries Прогон с нуля, исходный prompt, исходный resume_from X#retry1
Плюс on_reject Продолжение той самой отклонённой сессии, prompt заменяется на возвращённое on_reject, fork принудительно False X#round2

Второй вариант — это «вернуть на доработку, указать, чего не хватает, и дать дописать»: уже сделанная работа и контекст остаются на месте. Если on_reject вернул пустую строку или та попытка вообще не получила session_id, всё вырождается в прогон с нуля.

gate может также бросить StepAbort, что означает пробовать дальше бесполезно, не тратьте оставшиеся круги:

from flower import StepAbort

def gate(result, ctx):
    if "这个环境装不了依赖" in result.text:
        raise StepAbort("环境缺依赖,再跑几轮也一样")
    return "验收通过" in result.text

После броска: причина заносится в ctx["_aborted"], шаг считается провалившимся и обрабатывается по on_fail (по умолчанию "stop"), цикл повторов прерывается на месте, ни одна из оставшихся retries не тратится.

Запомните разницу: вернуть False — это «в этот раз не вышло, ещё круг»; StepAbort — это «ещё круг не поможет». Типичный случай: цель признана невыполнимой в этой среде и спросить некого — продолжать крутиться вхолостую самый дорогой вариант.

Не путайте два уровня повторов

Step.retries Runtime(resilience=...)
За что отвечает Прикладной провал: gate не пройден, result.ok ложно Инфраструктура: дрожание сети, обрыв, 5xx
Как повторяет Весь шаг заново, тот же prompt и тот же resume_from resume с места обрыва, прежние траты не пропадают
Что делает до этого Ничего Держит DNS + TCP-пробы и ждёт возврата сети (без HTTP, без учётных данных — проба обязана быть бесплатной)
Что не повторяется Ошибка учётных данных, ошибка параметров — немедленная остановка, без бесконечного ожидания

Prompt для продолжения намеренно не содержит никаких деталей ошибки: модели нужно знать «тебя прервали, продолжай», а не то, был ли это ENOTFOUND или 503.

Как сцеплять шаги

Передавать состояние между шагами можно тремя способами; выбор определяет, что увидит следующий шаг:

Запись Что видит следующий шаг Где применять
resume_from=None (по умолчанию) + впрыск в prompt Только те слова, которые вы впрыснули Шаги независимы. Дёшево, защищает от загрязнения
resume_from="имя предыдущего шага" Полную историю сессии Когда нужна связная память
resume_from="имя предыдущего шага" + fork=True Полную историю, но в отдельной ветке Перепроверка / параллельные варианты / повтор без загрязнения исходной линии

Шаг, на который указывает resume_from, обязан действительно создать сессию. Если он был пропущен через when или не выполнялся вовсе, Workflow.run сразу бросает ValueError — без тихой деградации до новой сессии, потому что это молча ломало бы допущение о «связной памяти».

Несколько выводов, за которые уже не раз платили:

  1. Один шаг — одна принимаемая цель. Граница шага — это граница контекста: там, где resume_from=None, все прежние результаты инструментов окончательно перестают висеть в контексте. См. Экономику контекста.
  2. Сомневаетесь — начните с clarify_step. В long-horizon «неверно понял цель» самая дорогая ошибка, и как раз она относится к тому классу, который те несколько слоёв экономии контекста вычистить не могут. См. clarify.
  3. Выдаваемая задача должна быть самодостаточной. У subagent чистый контекст, он не знает того, что знает координатор. Нужный фон пишите в задание либо указывайте, какой артефакт прочитать.
  4. Длинный вывод идёт на диск, а не в ответ. Это уже записано в WORKER_RULES — не отменяйте это своими instructions («вставь мне сюда полный лог»).
  5. gate в первую очередь проверяет жёсткие условия. Есть ли файл, равен ли код выхода 0 — то, что решается одной строкой Python, не надо отдавать модели. Если решать должна модель, берите готовый with_goal — он подменяет gate реализацией, запускающей отдельного судью; не собирайте это вручную внутри gate.
  6. Параллельные правки одного репозитория — только worker(isolate=True). Завершение (слияние, очистка worktree, открытие PR) сейчас остаётся на ваш workflow, harness гарантирует лишь то, что правки лягут каждая в свой worktree.

Верстак должен висеть на Workflow

Везде, где workflow пишет файлы в верстак — типовой случай clarify_step(brief_path=...) — нужно самому создать Workbench и одновременно повесить его на Workflow.workbench и передать в Runtime:

from pathlib import Path

from flower import (HumanChannel, Runtime, Step, Workbench, Workflow,
                    clarify_step, coordinator, worker)

wb = Workbench(Path.cwd()).ensure()
ch = HumanChannel(log_path=wb.notes / "问答记录.md", timeout_s=1800.0)

主控 = coordinator("协调者", "", {
    "coder": worker("写代码与测试。要动手实现的活派给它。",
                    "你负责实现。每改一处就跑一次验证,别攒到最后。"),
}, channel=ch)

wf = Workflow(
    [
        clarify_step(ch, brief_path=wb.notes / "需求.md", prompt="帮我做一个 X"),
        Step("干活", spec=主控, prompt=lambda ctx: f"照这份需求做:\n\n{ctx['确认需求']}"),
    ],
    channel=ch,
    workbench=wb,
)

rt = Runtime(workspace=Path.cwd(), run_dir="runs", workbench=wb)

Держать channel на workflow нужно по двум причинам: run() подключит его on_event к тому же выходу событий (только если channel.on_event всё ещё None), а драйвер по этому же полю узнаёт, кому отвечать.

Самодельный путь к верстаку отваливается молча

По умолчанию Runtime(workbench=True) кладёт верстак в <run_dir>/workbench, а Workbench(ws) — в <ws>/.flower: это не один и тот же каталог. Когда workflow вызывается из CLI, он не видит run_dir, так что самодельный путь соберётся куда-то не туда — бриф запишется в каталог A, а впрыснутый индекс просканирует каталог B, и ошибки при этом не будет. Создайте один объект и используйте его в обоих местах — проблемы не возникнет. Если Workflow.workbench задан, ключ -W командной строки игнорируется, приоритет у него.

continuous=True: пройти тот же путь ещё раз

Три способа сцепки выше — про шаги внутри одного запуска. Между процессами это другая ось:

Workflow([...], continuous=True)     # значение по умолчанию

Повторный запуск в том же рабочем каталоге — и каждый шаг продолжает говорить в той же сессии, что и в прошлый раз, за счёт пар «имя шага → session_id» в <run_dir>/lineage.json. При загрузке каждая запись проверяется через runtime.has_session() — жива ли ещё сессия в хранилище, и используется только живая: файл родословной может пережить sessions.db, а resume несуществующей сессии рванёт только после запуска дочернего процесса.

Три следствия:

  • resume_from=None не значит «совершенно новая сессия». В первый запуск — да, во второй уже нет. Если новая сессия нужна каждый раз, пишите явно Workflow(..., continuous=False).
  • При преемственности контекст будет только расти. Если при продолжении нужно сказать что-то другое, задайте Step.resume_prompt — то, что уже есть в контексте собеседника, повторно слать не нужно.
  • Шаги с явно заданным resume_from это не затрагивает, у него приоритет.

Имя шага — межпроцессный ключ

Переименовать шаг значит оборвать родословную этого шага: в следующий раз продолжения не будет, и ошибки не будет тоже. Имена повторов с суффиксами #retry1 / #round2 в родословную не попадают (записывается всегда исходное имя) — это один из способов, которым обеспечивается «судья всегда в новой сессии».

Полный дизайн и --new — в преемственности.

Когда этим пользоваться не нужно

  • Один агент и никаких вердиктов — не нужен Workflow. Просто await rt.run(spec, "…") или командная строка flower once "读一眼这个仓库".
  • Форма ровно такая: «выяснить требования → задать цель → работать» — берите готовый starter_flow(), писать своё не нужно:

    from flower import starter_flow
    
    wf = starter_flow("帮我做一个 X", workspace=".", run_dir="runs",
                      rounds=3, timeout_s=1800.0, isolate=False)
    

    Это три шага: 确认需求设定目标干活 (с циклом вердикта, шаг вердикта называется 干活·判定#N). При goal=False нет второго шага и цикла вердикта, при clarify_only=True остаётся только первый шаг. Он сам создаёт HumanChannel и Workbench и вешает их на workflow, поэтому Runtime(workbench=wf.workbench) можно брать напрямую, не собирая ещё один.

    Код писать тоже не обязательно: flower "帮我做一个 X" в каталоге проекта запускает именно его. Это не «рекомендуемый дизайн workflow», а всего лишь способ стартовать без конфигурации.

  • Шаги нарезаны мельче, чем «одна принимаемая цель» — чистый убыток. Каждый шаг поднимает новую сессию, а у новой сессии есть стартовый пол (у координатора на практике около 34k контекста), и размазать его не выйдет.

  • Хочется потом откатиться к конкретному сообщению — этим путём Workflow не ходит, он никогда не передаёт resume_at. Вызывайте напрямую Runtime.run(spec, "从这里重来", resume=sid, resume_at=uuid).

Как выбирать роль (coordinator / worker / clarify / judge / oracle) и семантика Step и Workflow по каждому полю — см. Python API; термины — в глоссарии.