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. 实施顺序

  1. Phase 1:并行执行(ThreadPoolExecutor + 超时)
  2. Phase 2:流水线 + 共享 workdir + context.json
  3. Phase 3:sub_tasks 支持 + 依赖图解析
  4. Phase 4:executor 读取上游 context 能力

每个 Phase 独立可用,不依赖后续 Phase。


待老大审批后执行。