把 Grok、Gemini、ChatGPT/Codex 和 Claude Code 串成一条自动化流水线,关键不在接入了多少个模型,而在于状态由谁推进、门禁由谁把守、非受信内容在哪里止步。本文给出实测可运行的 watchdog 编排器、带鉴权的 FastAPI Broker 与交接协议模板,逐一修正常见写法中的坑:状态卡死在 IN_PROGRESS、审查者被赋予写权限、`git diff HEAD~1` 漏看提交、TOTP 密钥错放进 Redis;并讨论跨 Agent 的提示注入、`ANTHROPIC_API_KEY` 的凭据优先级陷阱,以及为什么“第二个账号”带不来第二种视角。
适用人群:软件工程师、系统架构师、AI DevTools 开发者、技术极客
核心主题:异构 AI Agent 协作模式、跨账号状态解耦、上下文交接协议、本地中转服务(Broker)与自动化 Pipeline 搭建
本文中的 Claude Code 参数均对照 Claude Code 2.1.278 的
claude --help与官方文档核对。第 4 节的编排器、Broker 与轮询脚本都实际运行过:我们用一个模拟的claude可执行文件跑了五种情形(一次通过、测试先失败后通过、审查先打回后通过、无进展、连续打回三次),状态迁移全部符合预期。各家 CLI 迭代很快,读到本文时请以你本机--help的输出为准。
1. 背景与架构设计原则
1.1 单 Agent 的局限性
在日益复杂的软件工程中,只依赖单一 LLM 或单个 AI Agent 终端,往往会遇到以下瓶颈:
- 上下文窗口污染与遗忘(Context Pollution):长对话之后,Agent 会逐渐丢掉早期的架构约束,出现语义漂移乃至幻觉。长期有效的约束应该写进每次启动都会重新读取的项目指令文件(
CLAUDE.md、AGENTS.md),而不是寄希望于聊天记录。 - 配额与计费边界(Rate Limits & Quota):长时间密集生成代码,很容易撞上单个平台的用量上限。把不同性质的工作分散到不同平台,本身就是一种负载拆分——但这不等于“多开几个同一服务的账号轮换用”,这一点在 1.3 节单独讨论。
- 单一模型的思维盲区:不同模型在推理、代码重构、海量日志分析和实时资讯检索上各有所长,单一模型很难兼顾。
- 审计与执行同体(Self-Audit Risk):让同一个 Agent 既写代码又做安全审计,它很容易放过自己的盲点。
1.2 跨账号/跨平台的管道解耦架构
构建跨账号、多平台 Agent Pipeline 的核心思想,是把“隐性会话记忆”转化为“显性状态资产”:
1[1] 调研 ── Grok:实时 Web / X 检索;Gemini:长文档与多模态资料消化
2 │ 产物:docs/research/*.md(每条结论附来源 URL)
3 ▼
4[2] 规格 ── ChatGPT / Codex:推理、接口定义、测试清单
5 │ 产物:接口文件、测试清单、.pipeline/HANDOVER.md
6 │ ★ 人工确认点:架构决策在这里拍板
7 ▼
8[3] 实现 ── Claude Code(profile acc_a):读写仓库、运行测试
9 │ 门禁:编排器亲自运行 npm test,通过后才提交
10 ▼
11[4] 审查 ── Claude Code(profile acc_b,只读)或另一家厂商的模型
12 APPROVE → 完成;CHANGES_REQUESTED → 退回 [3]
13
14状态总线:Git 仓库(代码与提交历史)+ .pipeline/(HANDOVER.md、pipeline_state.json,不入库)通过引入状态总线(State Bus)与交接协议(Handover Protocol),不同账号、不同平台上的 Agent 可以在完全不共享登录凭据与 API 密钥的前提下,完成异步的任务接力。真正的“总线”其实有两条:代码和提交历史走 Git,流程状态走 .pipeline/ 目录(记得把它写进 .gitignore)。
1.3 “跨账号”到底解决什么问题
一种常见的做法是用“第二个账号”运行审查者,以期获得独立的视角。这里需要把两件事拆开:
- 独立上下文不需要第二个账号。 每次
claude -p调用都是一个全新的会话,看不到实现者的对话历史;Claude Code 内置的 subagent(.claude/agents/*.md,可以用tools:限定工具)也有自己独立的上下文窗口。 - 第二个账号也带不来第二种视角。 同一个模型换个账号,训练出来的偏好和盲点不会变。想要真正的“交叉审计”,审查阶段就应该换一家厂商的模型(见第 5 节 Step 4 末尾的替代方案)。
那么独立的配置目录(CLAUDE_CONFIG_DIR)和独立的账号还有没有意义?有,但理由是隔离而不是扩容:工作组织和个人账号分开、流水线使用单独计费并可设置支出上限的 API key、每个角色的凭据泄露影响面互不波及、审计日志可以按角色区分。
还有一条必须写明的边界:Anthropic 在 Claude Code 的法务与合规页面中说明,Pro/Max 套餐公布的用量额度以“普通的个人使用”为前提;构建产品或服务的开发者应当使用 Claude Console 的 API key。把多个订阅账号当成配额池、给 7×24 小时无人值守的流水线轮换使用,既违背这一前提,也会让你的流水线建立在随时可能失效的基础上。要跑无人值守的自动化,就用 API key,并在 Console 里设置支出上限。其他厂商也有各自的条款,上线前值得逐一读一遍。
2. 异构 Agent 矩阵与能力分工
为了最大化利用各个模型的优势,需要按照它们的原生能力定义明确的管道角色:
| Agent 平台 | 核心优势 | 定位与典型产物 | 自动化接入方式 |
|---|---|---|---|
| Grok | 实时 Web 与 X 平台检索、技术动态跟踪 | 情报员:技术选型报告、最新 API 变动说明(附来源) | Grok Build CLI(grok -p,Beta);或 xAI API(兼容 OpenAI 格式)+服务端 web_search/x_search 工具 |
| Gemini | 超长上下文、PDF/图片/视频等多模态解析 | 分析师:遗留仓库总结、海量日志与长文档摘要 | Gemini CLI(gemini -p) |
| ChatGPT/Codex | 逻辑推理、复杂算法设计、规范制定 | 架构师:接口规范、测试清单、HANDOVER.md | Codex CLI(codex exec)或 OpenAI API |
| Claude Code(acc_a) | 终端操控、多文件修改、测试驱动的迭代 | 主程序员:业务代码、单元测试 | claude -p |
| Claude Code(acc_b)或其他厂商 | 独立上下文、只读权限 | 审查员:审查报告与明确的通过/打回结论 | claude -p 只读配置;或 codex exec --sandbox read-only |
两点补充:
- 本文不写死模型型号。 o1/o3、GPT-4o、grok-3 这些在同类教程里常见的型号,到 2026 年都已不是各家的当前主力——grok-3 甚至已在 2026 年 5 月退役,请求会被自动转到新模型上。流水线配置里应该把型号当作参数,而不是写进架构图。
- “长上下文”已经不是 Gemini 独有的卖点。 主流旗舰模型的上下文窗口普遍达到数十万乃至百万 token。把长文档交给 Gemini,更实际的理由是它的多模态解析能力、独立的配额,以及“换一个模型来读”本身带来的视角差异。
3. 上下文对齐协议(Context Handover Standard)
要实现 Agent 之间的无缝交接,必须制定标准的“交接文件协议”,并且同时满足人类可读和机器可解析。
先分清两类上下文:
- 静态约束(代码规范、目录边界、禁止事项):写进项目根目录的
AGENTS.md。Codex、Cursor、Jules 等工具原生读取它;Gemini CLI 默认读GEMINI.md,但可以在.gemini/settings.json里设置"context": {"fileName": "AGENTS.md"};Claude Code 则在CLAUDE.md里写一行@AGENTS.md导入即可。一份文件,所有 Agent 共用。 - 动态交接(这一次任务做到哪了、下一步做什么):写进
.pipeline/HANDOVER.md,每个阶段结束时更新。
3.1 结构化 HANDOVER.md 规范
每个 Agent 结束自己的工作节点时,都要输出(或更新).pipeline/HANDOVER.md:
1# Agent Handover Protocol v1.1
2
3## 1. 任务基本信息 (Task Metadata)
4
5- **Source Agent**: codex(规格阶段)
6- **Target Agent**: claude-code:acc_a(实现阶段)
7- **Timestamp**: 2026-09-21T11:30:00+09:00
8- **Pipeline ID**: pipe_feat_auth_v2_88f9a
9- **Base Commit**: 3f2a9c1e7b4d
10
11## 2. 核心目标与范围 (Goal & Scope)
12
13- **Goal**: 为现有的 JWT 登录流程增加基于 TOTP 的双因素认证(2FA)。
14- **In-Scope**: `src/auth/`, `src/services/totpService.ts`, `src/middleware/auth.ts`, `tests/auth/`
15- **Out-of-Scope**: UI 前端组件;`src/config/jwt.ts` 中的签名算法与密钥轮换逻辑。
16
17## 3. 已完成的工作 (Completed Actions)
18
19- [x] 定义了 `IAuthService` 接口(见 `docs/specs/auth_spec.md`)。
20- [x] 选定 `otplib` v13 作为 TOTP 库。
21
22## 4. 上下文与约束条件 (Context & Constraints)
23
24- **关键文件**:
25 - `docs/specs/auth_spec.md`: 必须严格遵循的接口定义。
26 - `src/config/jwt.ts`: 只读参考,不得修改。
27- **已知技术坑点 (Pitfalls)**:
28 - TOTP 密钥是长期凭据,必须加密后持久化到数据库,不能只放在 Redis 或进程内存里。
29 Redis 只负责两件事:记录每个用户最近一次通过校验的时间步(RFC 6238 §5.2
30 要求同一时间步内的 OTP 不得被第二次接受),以及失败次数限流。
31 - otplib v13 是完全重写的版本:`authenticator` 导出已被移除,`verify()` 改为异步并返回
32 对象(读 `result.valid`)。网上大量 v12 示例不能照抄;防重放可用 `afterTimeStep` 选项。
33 - `auth.ts` 中间件依赖自定义的 `AppError` 类,不要直接抛出原生 `Error`。
34
35## 5. 验收标准 (Acceptance Criteria)
36
37- `npm test` 全部通过,`totpService` 行覆盖率 ≥ 90%。
38- 同一 OTP 在同一时间步内第二次提交时返回 401。
39
40## 6. 目标 Agent 下一步指令 (Next Actions for Receiver)
41
421. 运行 `npm install otplib@^13` 安装依赖。
432. 按照 `docs/specs/auth_spec.md` 实现 `src/services/totpService.ts`。
443. 为 `totpService` 编写单元测试。
454. 不要执行 git commit——提交由编排器在测试通过后完成。
46
47## 7. 未决问题 (Open Questions)
48
49- 恢复码(recovery codes)的生成与存储方案尚未确定,本阶段不实现。这份模板里有几处设计值得展开说明,它们各自对应交接文档里的一种常见错误:
- 范围与指令必须对齐:如果 In-Scope 里没有
src/services/,下一步指令却要求在那里新建文件,Agent 要么越界,要么卡住。 - 不要让交接文档传递错误的设计决定:一种常见的写法是“多实例部署下 TOTP 密钥必须写入 Redis”。但 TOTP 密钥和密码一样是长期凭据,应该加密存进数据库;Redis 作为缓存,可能被驱逐、可能没有开启持久化,适合放的是防重放记录和限流计数。交接文档里的错误约束,会被下游 Agent 当作“已经拍板的决定”忠实地执行。
- “基于 JWT 的双因素认证”是个混淆的说法:JWT 负责会话,TOTP 负责第二因素,两者是正交的,目标应当写成“为 JWT 登录流程增加 TOTP”。
- 验收标准和未决问题不可省略:没有验收标准,“完成”就只能由 Agent 自己宣布;没有未决问题列表,Agent 会自己替你做决定。
- 时间戳使用 ISO 8601,并记录
Base Commit,供第 3.3 节的 diff 使用。
3.2 机器状态机 pipeline_state.json
除了 Markdown,还需要一个 JSON 文件记录原子化的状态变化,供编排脚本(Watcher/Broker)判断下一步:
1{
2 "pipeline_id": "pipe_feat_auth_v2_88f9a",
3 "current_stage": "IMPLEMENTATION",
4 "status": "AWAITING_EXECUTION",
5 "base_commit": "3f2a9c1e7b4d",
6 "retry_count": 0,
7 "updated_at": "2026-09-21T11:30:00+09:00",
8 "history": [
9 {
10 "stage": "RESEARCH",
11 "agent": "grok",
12 "status": "COMPLETED",
13 "output_artifacts": ["docs/research/2fa.md"]
14 },
15 {
16 "stage": "SPECIFICATION",
17 "agent": "codex",
18 "status": "COMPLETED",
19 "output_artifacts": ["docs/specs/auth_spec.md", ".pipeline/HANDOVER.md"]
20 }
21 ]
22}status 的迁移路径是 AWAITING_EXECUTION → IN_PROGRESS →(下一阶段的 AWAITING_EXECUTION / COMPLETED / HUMAN_INTERVENTION_REQUIRED)。这里刻意没有 next_agent_target 这类字段:“哪个阶段由哪个 Agent 执行”属于编排器的配置,再写进状态文件只会制造两个互相矛盾的事实来源。
更关键的是谁有权写这个文件。上游(人,或第 4 节的 Broker 轮询脚本)只负责把它置为 AWAITING_EXECUTION;此后的每一次迁移都由编排器完成,Agent 本身不写状态文件。LLM 可能写出非法 JSON、跳过阶段,或者在测试没过的时候宣称“已完成”——状态机的推进应该建立在退出码和测试结果上,而不是 Agent 的自我汇报上。
3.3 基于 Git 的代码语义感知交接
对于修改代码的阶段,代码改动本身就是最强的上下文。以规格阶段为例,结束时提交:
1git add docs/specs/auth_spec.md src/auth/IAuthService.ts
2git commit -m "feat(auth): define 2FA spec and IAuthService interface
3
4- Add IAuthService interface
5- Store TOTP secrets encrypted in the DB; Redis only for replay guard and rate limits
6
7Pipeline-ID: pipe_feat_auth_v2_88f9a
8Handover-To: claude-code:acc_a"几个细节:
- 显式列出要提交的路径,而不是
git add .:后者会把.env、构建产物和.pipeline/一并扫进去。 - 把
Pipeline-ID写成 Git trailer(放在提交信息最后一段)。之后可以用git log --grep="Pipeline-ID: pipe_feat_auth_v2_88f9a"找出一条流水线的全部提交。 - 下游用
git diff <base_commit>...HEAD,而不是git diff HEAD~1。 实现阶段完全可能产生多个提交,经过“审查打回、再实现”的循环之后更是如此,HEAD~1只能看到最后一个。三点语法对比的是合并基点到HEAD的全部改动。
diff 告诉下游改了什么,提交信息和 HANDOVER.md 告诉它为什么改,两者缺一不可。但 diff 不包含没有被改动的调用方,所以审查者仍然需要读取相关文件——这也是第 4 节的审查配置里保留 Read、Grep 的原因。
4. 跨账号与跨平台通信总线(Inter-Agent Communication Bus)
要让彼此独立的 Agent 互相触发,就需要一条自动化的通信总线。先澄清一个前提:网页版的 ChatGPT、Gemini、Grok 既不能被脚本唤起,也没法把结果写进你的仓库。 想让它们进入自动化流水线,就要换成各家的 CLI 或 API——OpenAI 的 Codex CLI(codex exec)、Google 的 Gemini CLI(gemini -p)、xAI 的 Grok Build CLI(grok -p,目前为 Beta)或 xAI API。做不到全自动的阶段,保留为“人把产物放进仓库”的半自动节点,完全合理。
方案 A:共享文件系统与轻量级 Watchdog 监听器
这是最适合单机、多终端场景的方案:Python 脚本监听 .pipeline/pipeline_state.json 的变化,自动唤起下一个 Agent。
自动化调度脚本 agent_orchestrator.py
1#!/usr/bin/env python3
2"""
3agent_orchestrator.py — 监听 .pipeline/pipeline_state.json,按阶段调度本地 Claude Code。
4
5三条设计原则:
61. 状态迁移只由编排器执行。Agent 只产出工件然后退出,不碰状态文件。
72. 以退出码和编排器亲自跑的测试作为门禁,不采信 Agent 的自我汇报。
83. 状态文件一律原子替换写入,任何读取方都不会读到半截 JSON。
9假设同一时刻只有一个编排器实例在运行。
10"""
11import json
12import os
13import subprocess
14import tempfile
15import threading
16import time
17from datetime import datetime
18from pathlib import Path
19
20from watchdog.events import FileSystemEventHandler
21from watchdog.observers import Observer
22
23PIPELINE_DIR = Path(".pipeline").resolve()
24STATE_FILE = PIPELINE_DIR / "pipeline_state.json"
25# 在代码里把 "~" 展开成绝对路径:shell 不会展开引号或 Python 字符串里的 "~"
26PROFILES = Path.home() / ".claude_profiles"
27MAX_RETRIES = 3
28AGENT_TIMEOUT_SEC = 30 * 60
29
30STAGES = {
31 "IMPLEMENTATION": {
32 "profile": "acc_a",
33 "prompt": (
34 "读取 .pipeline/HANDOVER.md,在其限定的范围内完成实现并补齐单元测试。"
35 "如果存在 .pipeline/feedback.md,先逐条处理其中列出的问题。"
36 "不要修改 .pipeline/ 目录下的任何文件,也不要执行 git commit。"
37 ),
38 "flags": [
39 "--max-turns", "60",
40 "--max-budget-usd", "5",
41 "--permission-mode", "acceptEdits",
42 "--allowedTools", "Bash(npm test *)", "Bash(npx tsc *)", "Bash(git diff *)", "Bash(git status)",
43 ],
44 },
45 "REVIEW": {
46 "profile": "acc_b",
47 "prompt": (
48 "你是独立的安全审查员。标准输入是本次改动的完整 diff。"
49 "可以读取仓库中的文件补充上下文,但不得修改任何文件。"
50 "重点检查:连接泄漏、未处理的 Promise rejection、EVAL 中拼接的 Lua 脚本、"
51 "跨 slot 的多 key 操作,以及 package.json scripts、CI 配置等会改变执行行为的改动。"
52 "输出 Markdown 格式的审查报告,"
53 "最后一行必须是 VERDICT: APPROVE 或 VERDICT: CHANGES_REQUESTED。"
54 ),
55 "flags": [
56 "--max-turns", "30",
57 "--max-budget-usd", "2",
58 # dontAsk:凡是需要确认的操作一律自动拒绝,只放行下面白名单里的只读工具
59 "--permission-mode", "dontAsk",
60 "--allowedTools", "Read", "Grep", "Glob",
61 "--disallowedTools", "Edit", "Write", "NotebookEdit", "Bash", "WebFetch", "WebSearch",
62 ],
63 },
64}
65
66wake = threading.Event()
67
68
69class StateFileHandler(FileSystemEventHandler):
70 # 只订阅“写”类事件。watchdog 在 Linux 上连读取都会发 opened/closed 事件,
71 # 用 on_any_event 会让编排器每读一次状态就把自己再唤醒一次。
72 # 很多工具以“写临时文件 + rename”的方式保存,目标文件只收到 moved,收不到 modified。
73 def on_modified(self, event):
74 self._check(event.src_path)
75
76 def on_created(self, event):
77 self._check(event.src_path)
78
79 def on_moved(self, event):
80 self._check(event.dest_path)
81
82 def _check(self, path):
83 if Path(path).resolve() == STATE_FILE:
84 wake.set()
85
86
87def load_state():
88 for _ in range(10):
89 try:
90 return json.loads(STATE_FILE.read_text(encoding="utf-8"))
91 except FileNotFoundError:
92 return None
93 except json.JSONDecodeError:
94 time.sleep(0.2) # 外部写入方可能还没写完
95 raise RuntimeError(f"{STATE_FILE} 持续无法解析")
96
97
98def save_state(state):
99 state["updated_at"] = datetime.now().astimezone().isoformat(timespec="seconds")
100 fd, tmp = tempfile.mkstemp(dir=PIPELINE_DIR, suffix=".tmp")
101 with os.fdopen(fd, "w", encoding="utf-8") as f:
102 json.dump(state, f, ensure_ascii=False, indent=2)
103 os.replace(tmp, STATE_FILE) # 同一文件系统内的原子替换
104
105
106def git(*args):
107 return subprocess.run(
108 ["git", *args], capture_output=True, text=True, check=True
109 ).stdout
110
111
112def run_claude(stage, stdin_text=None):
113 spec = STAGES[stage]
114 # -p 模式下只要环境里有 ANTHROPIC_API_KEY 就一定优先使用它,配置目录隔离随之失效。
115 # 想让流水线走 API 计费,就删掉这个过滤,改为显式传入专用的 key
116 env = {k: v for k, v in os.environ.items() if k not in ("ANTHROPIC_API_KEY", "ANTHROPIC_AUTH_TOKEN")}
117 env["CLAUDE_CONFIG_DIR"] = str(PROFILES / spec["profile"])
118 # 提示词紧跟 -p,--allowedTools 这类变长参数统一放在末尾
119 cmd = ["claude", "-p", spec["prompt"], "--output-format", "json", *spec["flags"]]
120 # 没有输入时也给一个空 stdin:否则子进程继承编排器的 stdin,在非 TTY 环境下可能一直等输入
121 proc = subprocess.run(
122 cmd, env=env, input=stdin_text or "", capture_output=True, text=True,
123 timeout=AGENT_TIMEOUT_SEC,
124 )
125 try:
126 result = json.loads(proc.stdout)
127 except json.JSONDecodeError:
128 result = {"is_error": True, "result": proc.stderr[-2000:]}
129 result["exit_code"] = proc.returncode
130 return result
131
132
133def record(state, stage, status, result):
134 state.setdefault("history", []).append({
135 "stage": stage,
136 "agent": f"claude-code:{STAGES[stage]['profile']}",
137 "status": status,
138 "session_id": result.get("session_id"),
139 "num_turns": result.get("num_turns"),
140 "cost_usd": result.get("total_cost_usd"),
141 "at": datetime.now().astimezone().isoformat(timespec="seconds"),
142 })
143
144
145def advance(state, stage, next_stage, result):
146 record(state, stage, "COMPLETED", result)
147 state["current_stage"] = next_stage
148 state["status"] = "COMPLETED" if next_stage == "DONE" else "AWAITING_EXECUTION"
149 save_state(state)
150
151
152def fail(state, stage, feedback, result, back_to=None, give_up=False):
153 # retry_count 统计整条流水线的失败次数:测试没过、审查打回、超时都算
154 state["retry_count"] = state.get("retry_count", 0) + 1
155 record(state, stage, "FAILED", result)
156 (PIPELINE_DIR / "feedback.md").write_text(feedback, encoding="utf-8")
157 if give_up or state["retry_count"] >= MAX_RETRIES:
158 state["status"] = "HUMAN_INTERVENTION_REQUIRED"
159 notify(f"Pipeline {state['pipeline_id']} 在 {stage} 阶段停止:{feedback[:200]}")
160 else:
161 state["current_stage"] = back_to or stage
162 state["status"] = "AWAITING_EXECUTION"
163 save_state(state)
164
165
166def notify(message):
167 print(f"[!] {message}", flush=True) # 生产环境换成 webhook 告警,见 6.3 节
168
169
170def implement(state):
171 result = run_claude("IMPLEMENTATION")
172 if result.get("is_error") or result["exit_code"] != 0:
173 return fail(state, "IMPLEMENTATION", str(result.get("result", "")), result)
174 tests = subprocess.run(["npm", "test"], capture_output=True, text=True)
175 if tests.returncode != 0:
176 return fail(state, "IMPLEMENTATION", tests.stdout[-4000:] + tests.stderr[-2000:], result)
177
178 git("add", "-A", "--", ".", ":(exclude).pipeline")
179 tree = git("write-tree").strip()
180 # 与 HEAD 或上一轮完全相同的树:Agent 在原地打转,再跑也只是烧 token
181 if tree in (state.get("last_tree"), git("rev-parse", "HEAD^{tree}").strip()):
182 return fail(state, "IMPLEMENTATION", "本轮没有产生任何新改动,判定为无进展。", result, give_up=True)
183 state["last_tree"] = tree
184 git("commit", "-m", f"feat: implement {state['pipeline_id']}\n\nPipeline-ID: {state['pipeline_id']}")
185 (PIPELINE_DIR / "feedback.md").unlink(missing_ok=True)
186 advance(state, "IMPLEMENTATION", "REVIEW", result)
187
188
189def review(state):
190 diff = git("diff", f"{state['base_commit']}...HEAD", "--", ".", ":(exclude).pipeline")
191 result = run_claude("REVIEW", stdin_text=diff)
192 report = str(result.get("result", ""))
193 (PIPELINE_DIR / "audit_report.md").write_text(report, encoding="utf-8")
194 if result.get("is_error") or result["exit_code"] != 0:
195 return fail(state, "REVIEW", report, result)
196 if report.rstrip().endswith("VERDICT: APPROVE"):
197 return advance(state, "REVIEW", "DONE", result)
198 # 缺少判定行也按“打回”处理:宁可多跑一轮,也不放过没审完的改动
199 fail(state, "REVIEW", report, result, back_to="IMPLEMENTATION")
200
201
202HANDLERS = {"IMPLEMENTATION": implement, "REVIEW": review}
203
204
205def tick():
206 state = load_state()
207 while (
208 state
209 and state.get("status") == "AWAITING_EXECUTION"
210 and state.get("current_stage") in HANDLERS
211 ):
212 stage = state["current_stage"]
213 state["status"] = "IN_PROGRESS" # 进程若在此后崩溃,需人工把状态改回 AWAITING_EXECUTION
214 save_state(state)
215 print(f"[>] {state['pipeline_id']}: {stage}", flush=True)
216 try:
217 HANDLERS[stage](state)
218 except Exception as exc: # 超时、git 失败等,同样计入重试次数
219 fail(state, stage, f"orchestrator error: {exc!r}", {})
220 state = load_state()
221
222
223if __name__ == "__main__":
224 PIPELINE_DIR.mkdir(exist_ok=True)
225 observer = Observer()
226 observer.schedule(StateFileHandler(), str(PIPELINE_DIR), recursive=False)
227 observer.start()
228 print(f"[+] Orchestrator 已启动,监听 {STATE_FILE}", flush=True)
229 wake.set() # 启动时先检查一次,接住停机期间积压的任务
230 try:
231 while True:
232 wake.wait(timeout=30) # 以文件事件唤醒为主,30 秒轮询兜底
233 wake.clear()
234 tick()
235 except KeyboardInterrupt:
236 pass
237 finally:
238 observer.stop()
239 observer.join()网上常见的最小示例通常只有几十行:在 on_modified 回调里读 JSON、改状态,然后直接 subprocess.run 一个开着 --dangerously-skip-permissions 的 Agent。思路没错,但照搬会踩到下面这些坑,上面的实现逐一处理了:
- 文件事件比想象的复杂。 我们在 watchdog 6.0.0(Linux)上实测:单纯读取文件就会触发
FileOpenedEvent和FileClosedNoWriteEvent;以“写临时文件再 rename”方式保存时,目标文件只收到一个FileMovedEvent,on_modified根本不会触发;一次普通写入也可能触发多次on_modified。所以这里只订阅写类事件,并用一个threading.Event把连续事件合并成一次处理。 - 调度不能阻塞事件线程。 在 watchdog 的回调里直接运行一个可能跑半小时的 Agent,会把事件分发整个卡住;这里回调只负责“叫醒”,真正的工作放在主循环里。
- 状态不能只进不出。 朴素写法把状态改成
IN_PROGRESS之后,既不检查退出码,也不负责推进状态,完全指望 Agent 自己去改 JSON——Agent 只要漏掉这一步,流水线就会永远卡住,而且不报任何错。 - 审查者根本不需要写权限。 diff 通过 stdin 交给它,报告从 stdout 拿回来,由编排器写入
audit_report.md。一边要求审查者只读、一边用--dangerously-skip-permissions启动它,是这类示例里很常见的自相矛盾。 --output-format json会返回session_id、num_turns、total_cost_usd等字段,编排器把它们记进history,事后可以追溯每一轮花了多少钱、用了几轮。
两个配置目录需要各自先交互式登录一次:CLAUDE_CONFIG_DIR="$HOME/.claude_profiles/acc_a" claude,然后执行 /login。按照官方文档,设置了 CLAUDE_CONFIG_DIR 之后,凭据文件(macOS 上则是钥匙串条目)都按目录区分。需要完全无人值守时,可以用 claude setup-token 生成有效期一年的 OAuth token,通过 CLAUDE_CODE_OAUTH_TOKEN 传入;或者直接用 ANTHROPIC_API_KEY 走 API 计费。注意优先级:在 -p 模式下,只要环境里存在 ANTHROPIC_API_KEY,它就一定会被使用,排在 OAuth token 和订阅登录之前——这就是上面的代码要把它从子进程环境中过滤掉的原因。
方案 B:轻量级 Python HTTP Broker(跨机器的中转站)
如果 Agent 分布在不同的物理机器或 CI 环境中,可以部署一个 FastAPI Broker:
1# broker_server.py — 跨机器交接的最小中转站(单进程,演示用)
2import os
3import secrets
4from collections import defaultdict, deque
5
6from fastapi import Depends, FastAPI, Header, HTTPException
7from pydantic import BaseModel
8
9app = FastAPI(title="Multi-Agent Handover Broker")
10TOKEN = os.environ["BROKER_TOKEN"] # 没有设置 token 就拒绝启动
11
12
13def require_token(authorization: str = Header(default="")):
14 if not secrets.compare_digest(authorization.encode(), f"Bearer {TOKEN}".encode()):
15 raise HTTPException(status_code=401)
16
17
18class HandoverPayload(BaseModel):
19 pipeline_id: str
20 sender_agent: str
21 receiver_agent: str
22 stage: str
23 handover_markdown: str
24 git_commit: str # 已推送到共享远端的提交:跨机器只传引用,不传本地路径
25
26
27# 每个接收方一条 FIFO 队列。只存在进程内存里:重启即丢失,也不能开多个 worker
28queues: dict[str, deque] = defaultdict(deque)
29
30
31@app.post("/api/v1/handover", dependencies=[Depends(require_token)])
32async def enqueue(payload: HandoverPayload):
33 queues[payload.receiver_agent].append(payload)
34 return {"status": "ACK", "queued": len(queues[payload.receiver_agent])}
35
36
37@app.get("/api/v1/poll/{agent_id}", dependencies=[Depends(require_token)])
38async def poll(agent_id: str):
39 queue = queues.get(agent_id)
40 if not queue:
41 return {"has_task": False, "data": None}
42 return {"has_task": True, "data": queue.popleft()}
43
44
45if __name__ == "__main__":
46 import uvicorn
47
48 # 只监听本机。跨机器访问请放在 VPN / Tailscale 或带 TLS 的反向代理之后
49 uvicorn.run(app, host="127.0.0.1", port=8080)在跑 Claude Code 的机器上,再配一个轮询脚本,把领到的任务写成方案 A 的状态文件,后面的事交给编排器:
1# poll_worker.py — 运行在 Claude Code 所在的机器上,把 Broker 的任务转交给方案 A 的编排器
2import json
3import os
4import subprocess
5import time
6from pathlib import Path
7
8import httpx
9
10BROKER = os.environ.get("BROKER_URL", "http://127.0.0.1:8080")
11HEADERS = {"Authorization": f"Bearer {os.environ['BROKER_TOKEN']}"}
12AGENT_ID = "claude-code-acc-a"
13STATE_FILE = Path(".pipeline/pipeline_state.json") # .pipeline/ 应写进 .gitignore
14
15
16def busy():
17 if not STATE_FILE.exists():
18 return False
19 return json.loads(STATE_FILE.read_text("utf-8"))["status"] in ("AWAITING_EXECUTION", "IN_PROGRESS")
20
21
22while True:
23 time.sleep(10)
24 if busy(): # 编排器手上还有活,先不领新任务
25 continue
26 try:
27 resp = httpx.get(f"{BROKER}/api/v1/poll/{AGENT_ID}", headers=HEADERS, timeout=10)
28 resp.raise_for_status()
29 except httpx.HTTPError as exc:
30 print(f"[worker] broker 不可达:{exc}", flush=True)
31 continue
32 task = resp.json()
33 if not task["has_task"]:
34 continue
35
36 data = task["data"]
37 subprocess.run(["git", "fetch", "origin"], check=True)
38 subprocess.run(["git", "checkout", "-B", f"pipeline/{data['pipeline_id']}", data["git_commit"]], check=True)
39 STATE_FILE.parent.mkdir(exist_ok=True)
40 (STATE_FILE.parent / "HANDOVER.md").write_text(data["handover_markdown"], encoding="utf-8")
41 state = {
42 "pipeline_id": data["pipeline_id"],
43 "current_stage": data["stage"],
44 "status": "AWAITING_EXECUTION",
45 "base_commit": data["git_commit"],
46 "retry_count": 0,
47 "history": [],
48 }
49 tmp = STATE_FILE.with_suffix(".tmp")
50 tmp.write_text(json.dumps(state, ensure_ascii=False, indent=2), encoding="utf-8")
51 tmp.replace(STATE_FILE) # 原子替换,方案 A 的编排器随即被唤醒这个 Broker 很短,但有四处是刻意为之:
- 鉴权与监听地址。 很多示例监听
0.0.0.0且没有任何鉴权,而 Broker 投递的任务最终会交给一个能执行命令的 Agent——等于把一台机器的 shell 暴露给了整个网段。 - 队列必须真的是队列。 用
message_queue[receiver] = payload实现的“队列”,会让第二个任务直接覆盖第一个。这里每个接收方一条 FIFO 队列。 - 跨机器不要传本地路径。 在 payload 里传
docs/spec.md这样的本地路径,到了另一台机器上毫无意义;这里改为传git_commit,工件本身通过 Git 远端同步。 - 投递语义要心里有数。
poll取出即删除,属于“至多一次”:worker 在取走任务之后崩溃,这个任务就丢了。生产环境应换成带确认机制的队列,例如 Redis Streams 的消费者组(XREADGROUP+XACK)或 SQS 的可见性超时。
调用方式:所谓的“ChatGPT/Gemini 节点”,实际上是你自己写的脚本——它调用 Codex CLI、Gemini CLI 或各家 API 完成工作,把结果推送到 Git 远端,再向 /api/v1/handover 发送 POST。网页版的 ChatGPT 是没法往你内网的 Broker 发请求的。
方案 C:CLI 标准输入/输出(Stdio/Pipe)重定向
在 Shell 层面,可以用 Unix 管道把一个平台的输出直接交给另一个 Agent:
1#!/usr/bin/env bash
2# pipeline_pipe.sh
3set -euo pipefail
4mkdir -p .pipeline
5
6echo "=== 步骤 1:用 Gemini CLI 分析长文档 ==="
7gemini -p "@docs/legacy_architecture.pdf 详细总结其中的核心性能瓶颈,输出为简洁的 Markdown" \
8 > .pipeline/gemini_summary.md
9
10echo "=== 步骤 2:把总结通过 stdin 交给 Claude Code 实施重构 ==="
11claude -p "标准输入是一份架构瓶颈总结。据此重构 src/legacy_module.js,完成后运行 npm test。" \
12 --permission-mode acceptEdits \
13 --allowedTools "Bash(npm test *)" \
14 --max-turns 40 \
15 < .pipeline/gemini_summary.md两个容易写错的细节:Google 官方 CLI 的可执行文件名是 gemini,而不是仓库名 gemini-cli;文件通过提示词里的 @路径 引入。set -euo pipefail 让第一步失败时整条管道立即停止,而不是把一份空的总结交给下游。
第二步刻意没有使用 --dangerously-skip-permissions,而是“自动接受编辑+只放行测试命令”。原因不只是权限大小:管道会把上游模型的输出原封不动地变成下游 Agent 的指令。 Gemini 读的是一份 PDF,PDF 里的任何一段文字,都可能经过这条管道变成你机器上的一条命令。这正是 6.2 节要讨论的问题。
5. 实战演练:端到端跨平台 Agent 联合开发 Pipeline
下面是一个真实场景的半自动工作流:前两个阶段由人触发并确认,后两个阶段交给第 4 节的编排器。
场景:企业级 Redis 缓存层重构与独立安全审查
Step 1:调研阶段(Grok)
- Prompt:
"调研 Node.js 在 Redis Cluster 模式下的客户端选型与重连策略。每条结论附来源链接,并注明各客户端库当前的维护状态。" - 产物:
docs/research/redis_client.md
提示词刻意没有写成“查询 ioredis 的最佳重连策略”——那样已经预设了答案。先问选型,调研会带回一条关键事实:ioredis 的 README 写明它的维护是“best-effort”(尽力而为),并且新项目推荐使用 node-redis。如果跳过调研直接让编码 Agent 凭记忆写,它大概率会选训练数据里出现最多的那个库。
要求“每条结论附来源”也不是形式主义:检索型模型同样会编造 API 变动,没有链接的结论不应该进入下一阶段。
Step 2:规范与用例设计阶段(ChatGPT/Codex)
- 输入:
docs/research/redis_client.md - Prompt:
"作为首席架构师,根据调研报告定义 CacheManager 的 TypeScript 接口。接口不得暴露具体客户端库的类型。输出单元测试用例清单和 .pipeline/HANDOVER.md。" - 产物:
src/cache/ICacheManager.tstests/specs/cache_spec.md.pipeline/HANDOVER.md
“接口不暴露客户端库类型”这一条,让 node-redis 和 ioredis 之争变成一个可以推迟、可以撤回的实现细节。这里也是整条流水线里最值得安排人工确认的地方:架构决策错了,后面每个阶段都会把错误执行得非常认真。确认之后提交,并把这个提交的 SHA 作为 base_commit 写进状态文件,状态置为 AWAITING_EXECUTION,编排器随即接手。
Step 3:核心编码阶段(Claude Code,profile acc_a)
编排器执行的命令等价于:
1CLAUDE_CONFIG_DIR="$HOME/.claude_profiles/acc_a" claude -p \
2 "读取 .pipeline/HANDOVER.md 和 src/cache/ICacheManager.ts,实现 CacheManager 并补齐单元测试。不要执行 git commit。" \
3 --output-format json \
4 --max-turns 60 --max-budget-usd 5 \
5 --permission-mode acceptEdits \
6 --allowedTools "Bash(npm test *)" "Bash(npx tsc *)"注意这里写的是 $HOME,而不是 CLAUDE_CONFIG_DIR="~/.claude_acc_a":双引号里的 ~ 不会被 shell 展开,程序收到的就是字面量 ~/.claude_acc_a。能不能用,取决于每个读取这个变量的程序是否自己处理 ~——与其赌这一点,不如直接写 $HOME。另外,权限参数不能省:在 -p 模式下没有人能点“允许”,不预先放行的话,编辑文件和运行测试的请求都会被拒绝,Agent 实际上什么也改不了。
执行过程:
- Agent 读取
HANDOVER.md中的范围与约束。 - 生成
src/cache/CacheManager.ts及对应的单元测试,并自行运行npm test迭代到通过。 - Agent 退出后,编排器亲自再跑一遍
npm test——这一次的结果才算数。 - 测试通过,编排器提交改动(附
Pipeline-IDtrailer),把状态推进到REVIEW。
为什么不让 Agent 自己执行 git add . && git commit、自己修改 pipeline_state.json?因为只要提示词里漏掉一句,流水线就会就此停住,而且没有任何报错;git add . 还会把不该提交的文件一并扫进去。
Step 4:交叉安全审查阶段(Claude Code,profile acc_b)
1git diff "$BASE_COMMIT"...HEAD | CLAUDE_CONFIG_DIR="$HOME/.claude_profiles/acc_b" claude -p \
2 "你是独立的安全审查员。标准输入是本次改动的完整 diff。重点检查连接泄漏、未监听的 error 事件与未处理的 Promise rejection、EVAL 中拼接的 Lua 脚本、跨 slot 的多 key 操作。最后一行输出 VERDICT: APPROVE 或 VERDICT: CHANGES_REQUESTED。" \
3 --permission-mode dontAsk \
4 --allowedTools "Read" "Grep" "Glob" \
5 --disallowedTools "Edit" "Write" "Bash" \
6 --max-turns 30 > .pipeline/audit_report.md- 产物:
.pipeline/audit_report.md。最后一行是APPROVE时流水线结束;是CHANGES_REQUESTED,或者根本没有判定行时,报告被写入feedback.md,任务退回 Step 3。
审查重点里没有写“Redis 注入”,这是有意的。RESP 协议按参数数组传输命令,并不存在 SQL 注入那种意义上的拼接漏洞;Redis 真正的风险点是用字符串拼接生成的 Lua 脚本(EVAL)和缺少命名空间前缀的 key。Cluster 模式下还要额外检查多 key 操作:涉及的 key 不在同一个 slot 时会直接报 CROSSSLOT 错误,需要用 {user:42} 这样的 hash tag 让它们落到同一个 slot。
想要真正的模型多样性,可以把这一步换成另一家厂商的模型。例如用 Codex CLI 在只读沙箱中审查:codex exec --sandbox read-only "审查 git diff $BASE_COMMIT...HEAD 的改动……"。同一个模型换一个账号,只能带来一个干净的上下文,带不来第二种视角。
6. 工程安全、沙箱隔离与容错机制
多 Agent 流水线越复杂,容错与安全控制就越关键。
6.1 权限与账号沙箱隔离
-
配置目录隔离(Profile Sandboxing) 用环境变量把不同角色的认证状态分开,不要跨角色共用凭据:
alias claude-dev='CLAUDE_CONFIG_DIR="$HOME/.claude_profiles/acc_a" claude' alias claude-audit='CLAUDE_CONFIG_DIR="$HOME/.claude_profiles/acc_b" claude'alias 用单引号定义,变量会在每次使用时才展开。
-
按角色的最小权限
角色 权限模式 放行 禁止 实现者 acceptEdits文件编辑; Bash(npm test *)、Bash(npx tsc *)其余 Bash 命令( -p模式下没人批准,会被拒绝)审查者 dontAskRead、Grep、GlobEdit、Write、Bash、网络工具 -
权限规则是护栏,不是安全边界
Bash(npm test *)放行的是一个命令前缀,而npm test实际执行什么,由package.json里的scripts.test决定——实现者恰好有权修改这个文件。更要命的是,编排器自己运行的测试门禁也会执行它。所以:- 实现者和测试门禁都应该跑在容器、Dev Container 或 VM 里,里面不放生产凭据,网络出口只保留必要的白名单。Claude Code 自带的沙箱(设置项
sandbox.enabled,在 macOS、Linux 和 WSL2 上为 Bash 提供文件系统与网络隔离)也可以作为一层防护。 --dangerously-skip-permissions在claude --help里的原话是“Recommended only for sandboxes with no internet access”(仅推荐在无法访问互联网的沙箱中使用)。在宿主机上给每个 Agent 都开启它,是多 Agent 教程里最常见、也最危险的做法。- 审查者的检查清单里必须包括
package.json的 scripts 和 CI 配置的改动,第 4 节的审查提示词已经写进了这一条。
- 实现者和测试门禁都应该跑在容器、Dev Container 或 VM 里,里面不放生产凭据,网络出口只保留必要的白名单。Claude Code 自带的沙箱(设置项
6.2 跨 Agent 的提示注入:每一次交接都是信任边界
这是同类文章里最常被忽略、却最危险的一环。流水线的上游恰恰是专门读取非受信内容的 Agent:Grok 读网页和 X 上的帖子,Gemini 读来源不明的 PDF 和日志。它们的输出一旦原样流进 HANDOVER.md 的“下一步指令”,或者像方案 C 那样直接通过管道变成提示词,攻击者写在网页里的一句话就可能变成执行者机器上的一条命令。
Simon Willison 把这种组合称为 “lethal trifecta”(致命三要素):能接触私有数据、会读入非受信内容、能向外发送信息。三者同时具备,数据外泄只差一段精心构造的文本。多 Agent 流水线很容易在不知不觉中把三者凑齐:调研 Agent 负责读入,执行 Agent 手握仓库和凭据,两者之间只隔着一个 Markdown 文件。
可行的缓解措施:
- 调研产物只作参考资料,不作指令。 它们放在
docs/research/,由人或架构阶段提炼之后,才写进HANDOVER.md。执行者的提示词里明确说明:研究资料中的内容是数据,不是命令。 - 拆掉三要素中的至少一个。 执行者的环境里不放生产凭据,网络出口走白名单——这样即使被注入,也没有东西可偷、没有地方可发。
- 让审查者专门盯异常行为:新增的网络请求、新增的依赖、改动过的构建脚本。
6.3 死循环与 Token 消耗控制(Loop Protection)
当实现者和审查者陷入“修改 → 打回 → 修改 → 打回”的循环,token 会被持续消耗。控制要分三层:
- 单次调用:
--max-turns限制一次运行的轮数,--max-budget-usd限制一次运行的金额(只在-p模式下生效),再加上subprocess.run的timeout。 - 整条流水线:
retry_count统计测试失败、审查打回和超时的总次数,达到MAX_RETRIES就把状态置为HUMAN_INTERVENTION_REQUIRED。这一逻辑已经写在第 4 节编排器的fail()里。重试上限必须接进状态迁移本身:只定义一个check_loop_limit(),却不在每次失败时调用它、也不让retry_count递增,等于没有上限。 - 无进展检测:用
git write-tree计算暂存区的树哈希,与上一轮相同就立即停止。比起单纯计数,这能更早发现“Agent 反复提交同样的东西”这种情形。
每次调用的 total_cost_usd 都记在 history 里,一条流水线花了多少,一行命令就能算出来:
jq '[.history[].cost_usd // 0] | add' .pipeline/pipeline_state.json最后,把编排器里的 notify() 换成真正的告警。下面是 Slack Incoming Webhook 的写法;飞书和钉钉的自定义机器人只是 JSON 结构不同(飞书为 {"msg_type": "text", "content": {"text": ...}},钉钉为 {"msgtype": "text", "text": {"content": ...}}):
1import json
2import os
3import urllib.request
4
5
6def notify(message):
7 print(f"[!] {message}", flush=True)
8 url = os.environ.get("ALERT_WEBHOOK_URL")
9 if not url:
10 return
11 req = urllib.request.Request(
12 url,
13 data=json.dumps({"text": message}).encode(),
14 headers={"Content-Type": "application/json"},
15 )
16 try:
17 urllib.request.urlopen(req, timeout=10)
18 except OSError as exc: # 告警发不出去,不能反过来拖垮编排器
19 print(f"[!] 告警发送失败:{exc}", flush=True)7. 总结与未来的 Agent Mesh 趋势
跨账号、跨平台的 AI Agent 联合工作流,标志着 AI 辅助开发从**“单打独斗的 AI 助手时代”迈入“多 Agent 协作的流水线时代”**。
把 Grok(实时情报)、Gemini(长文档与多模态理解)、ChatGPT/Codex(推理与规范定义) 和 Claude Code(终端执行与重构) 组合起来,再配上标准化的交接协议(HANDOVER.md)和轻量级调度总线,开发者可以搭起一条高吞吐、强隔离、能互相审查的自动化软件生产线。但本文反复验证的一点是:这条生产线的可靠性,不取决于接入了多少个模型,而取决于状态由谁推进、门禁由谁把守、非受信内容在哪里止步。
至于未来,互联标准已经相当具体,用不着自己去发明 Unix Domain Socket 上的私有协议:
- MCP(Model Context Protocol) 负责把工具和上下文接进 Agent。Claude Code 自身就可以通过
claude mcp serve作为 MCP 服务器,被其他 Agent 当作工具调用。 - A2A(Agent2Agent)协议 负责 Agent 之间的通信,2026 年 8 月加入了 Linux Foundation 旗下的 Agentic AI Foundation(AAIF),规范版本为 1.0。Agent 通过
/.well-known/agent-card.json发布自己的能力名片。 - AGENTS.md 同样由 AAIF 管理,正在成为跨工具共享项目约束的事实标准。
回头看,本文手写的 pipeline_state.json 状态机和 Broker,本质上就是 A2A 中任务(Task)生命周期与工件(Artifact)传递的手工版本;HUMAN_INTERVENTION_REQUIRED 对应的,正是 A2A 中的 input-required 状态。先亲手搭一遍简化版,再迁移到标准协议,你会清楚地知道每一个字段为什么存在。
落地清单:
- 把
.pipeline/写进.gitignore,并为每个配置目录先交互式/login一次。 - 状态迁移只由编排器执行,门禁看退出码和测试结果。
- 审查者只读,diff 走 stdin,报告走 stdout。
- 执行者跑在沙箱里,不接触生产凭据,网络出口走白名单。
- 每次调用都设置
--max-turns和--max-budget-usd,整条流水线设置重试上限。 - 在规格阶段之后安排一次人工确认。
- 无人值守的流水线使用 API key,而不是在多个订阅账号之间轮换。

Comments
Comments (0)