Bridge v2 协同作业升级 Spec
创建:2026-03-30 状态:草案,待老大审批
1. 背景
当前 bridge(248 手搓版)实现了基本的任务派发→执行→通知链路,但 worker 是串行的:codex 跑完 kiro 再跑。无法实现真正的协同作业。
老大的核心需求:让 codex(小呆子)和 kiro(小德子)能协同干活,摆脱 Discord 编排的不稳定性。
2. 现状(v1)
OC → SSH 写 inbox → bridge_worker.py 串行消费
↓
router.py 选 executor
↓
executor 1 → executor 2(串行)
↓
state + webhook 通知
v1 能力
- ✅ 自动路由(task_type + 关键词启发式)
- ✅ 串行多 executor(code+review = codex 先 → kiro 后)
- ✅ 重试(MAX_RETRIES=1)
- ✅ 归档(archive/{task_id}/)
- ✅ webhook 通知(Discord 报告厅)
- ✅ 错误隔离(dead-letter、error log)
v1 缺陷
- ❌ 串行执行,浪费时间
- ❌ executor 之间无上下文传递
- ❌ 无共享工作目录
- ❌ 无动态协调(一个卡住另一个不知道)
- ❌ 无并行能力
3. 目标(v2)
让 bridge 支持真正的协同作业模式,借鉴 ClawTeam 的设计思路但保持轻量。
核心原则
- 不装 ClawTeam 框架,继续手搓,但借鉴它的好设计
- 最小改动,在现有代码上增量升级
- 保持 fire-and-forget,不引入守护进程
4. 借鉴 ClawTeam 的设计
| ClawTeam 概念 | 我们怎么用 | 实现方式 |
|---|---|---|
| Mailbox(收件箱) | executor 之间传递上下文 | 共享 workdir 下的 context.json |
| TaskWaiter(等待器) | 并行执行后等所有 executor 完成 | concurrent.futures.ThreadPoolExecutor |
| Lifecycle on-exit | executor 异常退出时清理 | try/finally + state 更新 |
| Transport(消息传输) | executor 产出物传递 | 共享文件系统(同一台机器,不需要网络 transport) |
| Watcher(监视器) | 监控 executor 超时/死亡 | Future.result(timeout=N) |
5. 执行模式设计
模式 A:串行(现有,保留)
executor_1 → executor_2 → ... → 完成
适用:code+review(codex 写完 kiro 才能 review)
模式 B:并行(新增)
executor_1 ─┐
executor_2 ─┤→ 等全部完成 → 合并结果 → 完成
executor_3 ─┘
适用:独立子任务(如 codex 写模块 A,kiro 写文档 B)
模式 C:流水线(新增)
executor_1 → 产出物 → executor_2(消费产出物)→ 完成
适用:code+review 的升级版——codex 的代码自动传给 kiro 做 review
路由规则升级
def choose_mode(task: dict) -> str:
task_type = task.get('task_type', 'auto')
mode = task.get('execution_mode') # 新字段,允许手动指定
if mode:
return mode
if task_type == 'code+review':
return 'pipeline' # 流水线:codex → kiro
if task_type == 'parallel':
return 'parallel' # 并行
return 'serial' # 默认串行
6. 共享工作目录
借鉴 ClawTeam 的 workspace 概念:
bridge/work/{task_id}/
├── codex/ # codex executor 的工作区
├── kiro/ # kiro executor 的工作区
├── shared/ # 共享产出物
│ ├── context.json # 上下文传递(类似 ClawTeam mailbox)
│ └── artifacts/ # 共享文件
└── meta.json # 任务元数据
context.json 格式(借鉴 ClawTeam TeamMessage)
{
"messages": [
{
"from": "codex",
"type": "artifact",
"content": "代码已写入 shared/artifacts/main.py",
"timestamp": "2026-03-30T08:00:00+08:00"
},
{
"from": "kiro",
"type": "review",
"content": "发现 3 个问题,详见 shared/artifacts/review.md",
"timestamp": "2026-03-30T08:05:00+08:00"
}
]
}
7. 并行执行实现
from concurrent.futures import ThreadPoolExecutor, as_completed
def run_parallel(workers: list[str], task: dict, timeout: int = 600) -> list[dict]:
results = []
with ThreadPoolExecutor(max_workers=len(workers)) as pool:
futures = {
pool.submit(run_single_executor, name, task): name
for name in workers
}
for future in as_completed(futures, timeout=timeout):
name = futures[future]
try:
result = future.result()
results.append({'executor': name, 'status': 'ok', **result})
except Exception as e:
results.append({'executor': name, 'status': 'failed', 'error': str(e)})
return results
8. 流水线实现
def run_pipeline(workers: list[str], task: dict) -> list[dict]:
results = []
prev_output = None
for name in workers:
# 把上一步的产出注入到 task context
if prev_output:
task.setdefault('context', {})['previous_output'] = prev_output
# 写入共享目录的 context.json
write_shared_context(task['task_id'], name, prev_output)
result = run_single_executor(name, task)
results.append({'executor': name, 'status': 'ok', **result})
prev_output = result
return results
9. 超时与死亡检测
借鉴 ClawTeam 的 TaskWaiter.on_agent_dead:
EXECUTOR_TIMEOUT = {
'codex': 600, # 10 分钟
'kiro': 300, # 5 分钟
'mock': 10, # 10 秒
}
# 在并行模式下,Future.result(timeout=N) 自动处理
# 超时的 executor 标记为 failed,不阻塞其他 executor
10. 任务协议升级(v2)
新增字段(向后兼容,v1 任务照跑):
{
"schema": "clawteam.bridge.task.v2",
"task_id": "oc_20260330_001",
"title": "...",
"goal": "...",
"task_type": "code+review",
"execution_mode": "pipeline|parallel|serial",
"timeout": 600,
"sub_tasks": [
{
"executor": "codex",
"goal": "实现登录模块",
"timeout": 600
},
{
"executor": "kiro",
"goal": "review 登录模块代码",
"depends_on": "codex",
"timeout": 300
}
]
}
sub_tasks 是可选的。没有的话走现有路由逻辑。有的话按依赖关系自动决定串行/并行。
11. 改动清单
| 文件 | 改动 | 优先级 |
|---|---|---|
bridge_worker.py |
加 run_parallel() + run_pipeline() + 共享 workdir |
P0 |
router.py |
加 choose_mode() + 支持 sub_tasks 解析 |
P0 |
codex_executor.py |
支持读取 context.json 上游产出 |
P1 |
kiro_executor.py |
支持读取 context.json 上游产出 |
P1 |
bridge_submit_only.py(OC 侧) |
支持 v2 协议字段 | P1 |
新增 shared_context.py |
共享上下文读写工具 | P1 |
12. 不做的事
- ❌ 不装 ClawTeam 框架
- ❌ 不搞 Web UI / 看板(有 Discord webhook 够了)
- ❌ 不搞 P2P transport(同一台机器,文件系统就够)
- ❌ 不搞动态 agent 注册/发现(就 codex + kiro 两个)
- ❌ 不搞守护进程(保持 fire-and-forget)
13. 实施顺序
- Phase 1:并行执行(ThreadPoolExecutor + 超时)
- Phase 2:流水线 + 共享 workdir + context.json
- Phase 3:sub_tasks 支持 + 依赖图解析
- Phase 4:executor 读取上游 context 能力
每个 Phase 独立可用,不依赖后续 Phase。
待老大审批后执行。
💬 评论 (0)
暂无评论,来说第一句话吧~
发表评论