Skip to content

Flow 开发教程

Flow 是跨 Agent 的多阶段工作流编排,多个角色(Agent)并行/串行协作完成复杂任务。本文覆盖主链路源码、从简单到复杂的实例。


什么是 Flow

概念与 Recipe 的区别
Recipe单 Agent 内,自动化一个技能
Flow多 Agent 协作,编排跨角色工作流
🔀 Flow 执行模型 —— 创意工坊预览
┌──────────────────────────────────────────────────────────────────────┐
│  FlowTemplate (JSON)                                                 │
│    roles: {role_A → agent_1, role_B → agent_2, role_C → agent_3}    │
└───────────────────────────────┬──────────────────────────────────────┘


┌──────────────────────────────────────────────────────────────────────┐
│                    FlowEngine(并发执行)                             │
│                                                                      │
│   owner-thread (协调)                                                │
│     │                                                                │
│     ├──▶ role_A worker ──▶ agent_1 执行                              │
│     │                                                                │
│     ├──▶ role_B worker ──▶ agent_2 执行  (并行)                    │
│     │                                                                │
│     └──▶ role_C worker ──▶ agent_3 执行  (并行)                    │
│                                                                      │
│   事件队列 ──▶ SSE 推送到前端                                        │
└──────────────────────────────────────────────────────────────────────┘


┌──────────────────────────────────────────────────────────────────────┐
│  断点续跑 · src/flow/resume.py:rehydrate()                           │
│    重放 mutation log + hash 校验                                      │
└──────────────────────────────────────────────────────────────────────┘

📄 src/flow/engine.py — FlowEngine 主循环 📄 src/flow/resume.py — 断点续跑(rehydrate) 📄 src/flow/models/graph.py — 图定义


主链路:从模板到执行

1. 编译 template.json → FlowTemplate 对象
   src/flow/compiler.py:compile_template()

2. 创建 Flow 实例
   POST /api/flows
   → 分配 flow_id, 初始化 mutation log

3. 启动执行
   FlowEngine.run()
   ├── owner-thread: 协调节点调度
   ├── per-role worker: 各角色并行执行
   └── 事件队列: 推送 SSE 到前端

4. 断点续跑(可选)
   src/flow/resume.py:rehydrate()
   → 重放 mutation log + 校验 hash

关键源码

文件职责
src/flow/engine.py:FlowEngine主引擎:owner-thread + worker-threads
src/flow/engine.py:_execute_node()单节点执行
src/flow/resume.py:rehydrate()从 DB 恢复中断的 flow
src/flow/models/graph.py:GraphNode节点定义
src/flow/compiler.py:compile_template()编译期校验
src/bridge/routes/flows.pyREST API

节点类型

1. starter(输入定义)

yaml
node:
  id: starter
  type: starter
  config:
    inputs:
      - name: user_request
        type: string
        required: true
      - name: priority
        type: enum
        values: [low, medium, high]
        default: medium

2. llm_call(LLM 调用)

yaml
node:
  id: analyze_request
  type: llm_call
  role: analyst        # 绑定到 analyst 角色
  config:
    system_prompt: "分析用户请求,提取关键需求。"
    user_prompt_template: "{{user_request}}"
    output_var: analysis

3. for_each(循环执行)

yaml
node:
  id: process_items
  type: for_each
  config:
    items_var: analysis.items      # 要遍历的列表
    item_var: current_item         # 循环变量名
    body_nodes:                    # 循环体内的节点
      - id: process_one
        type: llm_call
        role: worker
        config:
          system_prompt: "处理单个项目"
          user_prompt_template: "项目:{{current_item}}"
          output_var: result
    output_var: all_results        # 收集所有结果

4. branch(条件分支)

yaml
node:
  id: route_by_priority
  type: branch
  config:
    cases:
      - condition: "priority == 'high'"
        target: urgent_handler
      - condition: "priority == 'low'"
        target: batch_handler
    else: normal_handler

⚠️ conditionAST 白名单src/recipes/ast_guard.py),语法同 Recipe。

5. summary(输出汇总)

yaml
node:
  id: summary
  type: summary
  config:
    outputs:
      - name: final_report
        from: merge_results
      - name: status
        value: "completed"

实例一:简单 Flow —— 双人翻译审校

目标:Agent A 翻译 → Agent B 审校 → 输出终稿。

🔀 translation_review · 创意工坊
┌──────────────────────┐
│ ● starter            │
│   inputs: text,      │
│     target_lang      │
└──────────┬───────────┘


┌──────────────────────┐
│ ● llm_call           │  role: translator
│   translate          │  → agent_1
│   system: 专业翻译    │
│   → draft            │
└──────────┬───────────┘


┌──────────────────────┐
│ ● llm_call           │  role: reviewer
│   review             │  → agent_2
│   system: 审校专家    │
│   → review_result    │
└──────────┬───────────┘


┌──────────────────────┐
│ ● summary            │
│   done               │
│   outputs:           │
│     translation,     │
│     fixes            │
└──────────────────────┘

完整 template.json

json
{
  "id": "translation_review",
  "name": "双人翻译审校",
  "version": "1.0.0",
  "roles": {
    "translator": {"agent_id": 1, "description": "翻译员"},
    "reviewer": {"agent_id": 2, "description": "审校员"}
  },
  "nodes": [
    {
      "id": "starter",
      "type": "starter",
      "config": {
        "inputs": [
          {"name": "text", "type": "string", "required": true},
          {"name": "target_lang", "type": "string", "default": "en"}
        ]
      }
    },
    {
      "id": "translate",
      "type": "llm_call",
      "role": "translator",
      "config": {
        "system_prompt": "你是专业翻译。只输出翻译结果。",
        "user_prompt_template": "翻译为{{target_lang}}:\n{{text}}",
        "output_var": "draft"
      }
    },
    {
      "id": "review",
      "type": "llm_call",
      "role": "reviewer",
      "config": {
        "system_prompt": "你是审校专家。检查翻译质量,修正错误。",
        "user_prompt_template": "原文:{{text}}\n译文:{{draft}}",
        "output_format": {
          "type": "object",
          "properties": {
            "final": {"type": "string"},
            "issues_fixed": {"type": "array", "items": {"type": "string"}}
          }
        },
        "output_var": "review_result"
      }
    },
    {
      "id": "summary",
      "type": "summary",
      "config": {
        "outputs": [
          {"name": "translation", "from": "review_result.final"},
          {"name": "fixes", "from": "review_result.issues_fixed"}
        ]
      }
    }
  ],
  "edges": [
    {"from": "starter", "to": "translate"},
    {"from": "translate", "to": "review"},
    {"from": "review", "to": "summary"}
  ]
}

运行

bash
# 创建 flow 实例
curl -X POST http://127.0.0.1:18711/api/flows \
  -H "X-Bridge-Token: $TOKEN" \
  -d '{"template_id": "translation_review", "inputs": {"text": "你好世界"}}'

# 查看执行状态
curl http://127.0.0.1:18711/api/flows/{flow_id}

# SSE 实时事件
curl http://127.0.0.1:18711/api/flows/{flow_id}/events

实例二:并行 Flow —— 多源信息采集

目标:3 个 Agent 并行搜索不同来源 → 汇总去重 → 输出报告。

🔀 multi_source_collect · 创意工坊
                        ┌─────────────────────────────────────────────┐
                        │ ● parallel_group                           │
                        │   branches:                                 │
                        │     ┌─────────────────────────────────────┐ │
                        │     │ branch 1 · role: web_searcher       │ │
                        │     │   ● llm_call: web                   │ │
                        │     │     tools: browser_search           │ │
                        │     └─────────────────────────────────────┘ │
                        │     ┌─────────────────────────────────────┐ │
            ┌──split ───┤     │ branch 2 · role: db_searcher        │ ├──→ merger ──▶ done
            │           │     │   ● llm_call: db                    │ │      (role_C)
            │           │     │     tools: search_internal_db       │ │
            │           │     └─────────────────────────────────────┘ │
            │           │     ┌─────────────────────────────────────┐ │
            │           │     │ branch 3 · role: api_fetcher        │ │
            │           │     │   ● llm_call: api                   │ │
            │           │     │     tools: fetch_api                │ │
            │           │     └─────────────────────────────────────┘ │
            │           └─────────────────────────────────────────────┘

关键节点

json
{
  "id": "parallel_search",
  "type": "parallel_group",
  "config": {
    "branches": [
      {
        "role": "web_searcher",
        "nodes": [
          {"id": "web", "type": "llm_call", "config": {"tool_whitelist": ["browser_search"]}}
        ]
      },
      {
        "role": "db_searcher",
        "nodes": [
          {"id": "db", "type": "llm_call", "config": {"tool_whitelist": ["search_internal_db"]}}
        ]
      },
      {
        "role": "api_fetcher",
        "nodes": [
          {"id": "api", "type": "llm_call", "config": {"tool_whitelist": ["fetch_api"]}}
        ]
      }
    ]
  }
}

实例三:带循环的 Flow —— 迭代优化写作

目标:Writer 写初稿 → Reviewer 评分 → 低于阈值则重写(最多 3 轮)。

🔀 iterative_writing · 创意工坊
┌──────────────────────┐
│ ● starter            │
│   inputs: topic,     │
│     min_score        │
└──────────┬───────────┘


┌──────────────────────┐
│ ● llm_call           │  role: writer
│   write              │  → agent_1
│   system: 专业写手    │  loop_context.feedback
│   → draft            │  接收 reviewer 反馈
└──────────┬───────────┘


┌──────────────────────┐
│ ● llm_call           │  role: reviewer
│   review             │  → agent_2
│   system: 质量评审    │
│   → review_result    │
│     (score, feedback)│
└──────────┬───────────┘


┌──────────────────────┐
│ ● branch             │
│   route_by_score     │
│                      │
│   score >= min_score ─────yes────▶ ┌──────────────────────┐
│                      │             │ ● summary            │
│   score < min_score  │             │   done               │
│   max_loops: 3       │             │   outputs: final     │
│   loop_context:      │             └──────────────────────┘
│     feedback:        │
│       review_result  │
│       .feedback      │
│                      │
│      no ─────────▶ (回到 write,携带 feedback)
└──────────────────────┘

循环实现

json
{
  "id": "check_score",
  "type": "branch",
  "config": {
    "cases": [
      {
        "condition": "review_result.score >= min_score",
        "target": "summary"
      },
      {
        "condition": "review_result.score < min_score",
        "target": "write",
        "loop_context": {"feedback": "review_result.feedback"}
      }
    ],
    "max_loops": 3
  }
}

编程助手与可视化编辑

Flow 编辑器前端(electron/renderer/src/components/flow-builder/)提供可视化拖拽:

flow-builder/
├── FlowCanvas.vue          # VueFlow 画布
├── GraphNode.vue           # 节点组件
├── GraphEdge.vue           # 边组件
├── RolePanel.vue           # 角色绑定面板
├── NodeConfigModal.vue     # 节点配置弹窗
└── stores/
    └── flowBuilder.ts      # Pinia store

📄 electron/renderer/src/components/flow-builder/ 📄 src/flow/_builtin/hello_world/ — 内置示例


调试与运维

bash
# 列出所有 flow
python main.py flows list --pretty

# 查看 flow 详情(含 trace)
python main.py flows get <flow_id>

# 取消运行中的 flow
python main.py flows cancel <flow_id>

# 断点续跑(bridge 崩溃恢复)
python -c "from src.flow.resume import rehydrate; rehydrate()"

常见坑

问题原因解决
并行组死锁节点等待不存在的依赖检查 branch edges
分支体丢失存放位置错误branch_bodies[case.id]
续跑拒绝模板被改过hash 校验失败,需重新创建
VueFlow 点击失效组件被 reactive 包裹markRaw()

更多

  • 源码src/flow/
  • 内置模板src/flow/_builtin/hello_world/
  • 编辑器前端electron/renderer/src/components/flow-builder/
  • 完整来源FLOW_DEVELOPMENT.md

✨ Familiars · 多 Agent AI 桌面应用 —— 对话树记忆 · 长期记忆 · 工具系统 · 数字人 · Broker 生态