Проектирование 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 — это модуль:атрибут или путь_к_файлу:атрибут. Если полученный объект вызываемый, он сначала вызывается один раз; затем берётся Workflow и запускается. По завершении в терминал печатается суммарная стоимость и путь к манифесту запуска.
Можно написать и свой драйвер — первый аргумент Workflow.run это Runtime:
Что оно делает на самом деле¶
Что получает Step и что обязан вернуть¶
Step — это не функция, а декларация. Реально выполняется Runtime.run(step.spec, отрендеренный prompt, ...) — один шаг = один вызов Runtime.run = одна сессия.
Первые три поля — позиционные аргументы, Step(name, spec, prompt):
name— имя шага. Оно же ключ вctx, имя строки вruns/manifest.jsonи ключ межпроцессной родословной.spec— какойAgentSpecиспользовать. Он задаёт для этого шага белый список инструментов, модель и бюджет.prompt—strлибо(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: это не синтаксический сахар¶
По умолчанию дальше передаётся дословный ответ модели. У некоторых шагов дословный ответ передавать дальше нельзя:
На практике шаг уточнения требований помимо четырёх разделов вклеивает в ответ весь код целиком. Дальше по цепочке должны идти только разобранные четыре раздела, иначе вся эта груда кода попадёт в 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 — без тихой деградации до новой сессии, потому что это молча ломало бы допущение о «связной памяти».
Несколько выводов, за которые уже не раз платили:
- Один шаг — одна принимаемая цель. Граница шага — это граница контекста: там, где
resume_from=None, все прежние результаты инструментов окончательно перестают висеть в контексте. См. Экономику контекста. - Сомневаетесь — начните с
clarify_step. В long-horizon «неверно понял цель» самая дорогая ошибка, и как раз она относится к тому классу, который те несколько слоёв экономии контекста вычистить не могут. См. clarify. - Выдаваемая задача должна быть самодостаточной. У subagent чистый контекст, он не знает того, что знает координатор. Нужный фон пишите в задание либо указывайте, какой артефакт прочитать.
- Длинный вывод идёт на диск, а не в ответ. Это уже записано в
WORKER_RULES— не отменяйте это своимиinstructions(«вставь мне сюда полный лог»). gateв первую очередь проверяет жёсткие условия. Есть ли файл, равен ли код выхода 0 — то, что решается одной строкой Python, не надо отдавать модели. Если решать должна модель, берите готовыйwith_goal— он подменяет gate реализацией, запускающей отдельного судью; не собирайте это вручную внутри gate.- Параллельные правки одного репозитория — только
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: пройти тот же путь ещё раз¶
Три способа сцепки выше — про шаги внутри одного запуска. Между процессами это другая ось:
Повторный запуск в том же рабочем каталоге — и каждый шаг продолжает говорить в той же сессии, что и в прошлый раз, за счёт пар «имя шага → 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; термины — в глоссарии.