PicoBot/docs/ASYNC_SUBAGENT_DESIGN.md
oudecheng f05f09b636 docs: 异步子代理设计文档
记录 v8 设计方案的三个核心不变量、wait 模式、取消机制与崩溃恢复策略
2026-08-13 22:13:00 +08:00

20 KiB
Raw Permalink Blame History

异步子 Agent 方案设计

版本: v8 (最终版) 日期: 2026-08-12 状态: 设计定稿

一、核心思路

子 agent 执行与父 agent 流程解耦,通过双通道实现异步通信:

  • 子代理结果走独立队列 sub_done_q(直达 wait不绕 SQLite → bus → reload
  • 用户消息走现有 bus(保持现有机制,零 channel 适配器改动)

wait 工具通过 释放锁 + select! 等待 + 重新获取锁 实现真等待bus 消息能在 wait 期间通过 is_waiting 分流注入 history 并唤醒 wait。

核心机制

sub_done_q独立队列  ← 后台子代理完成时注入 → wait select! 直达消费
bus现有机制         ← ws.rs 用户消息注入    → process_one 分流注入 + wakeup 唤醒 wait

二、关键设计决策

决策点 选择 理由
子 agent 执行方式 异步tokio::spawn 消除最大阻塞源
子代理结果通道 sub_done_q 独立队列 直达 wait不绕 SQLite → bus → reload
用户消息通道 复用 bus 现有机制channel 零改动
wait 语义 释放锁 + select! + 重获取锁 真等待,无死锁,无 TOCTOU
is_waiting 判断 持锁后判断 原子操作,无 TOCTOU
完成回调内容 当前结果 + 未完成子代理 id 列表 LLM 据此决策下一步
task tool_result running + 引导提示(只提 wait LLM 知道应调 wait
嵌套层 仅顶层异步 避免复杂度爆炸
模式1中断 不做 复杂度高,收益低

三、数据结构

1. SessionHistory 扩展(核心)

// src/gateway/session_history.rs
pub(crate) struct SessionHistory {
    // 现有字段
    topic_histories: HashMap<String, Vec<ChatMessage>>,
    chat_topic_ids: HashMap<String, String>,
    compression_in_flight: HashSet<String>,
    topic_serial_locks: HashMap<String, Arc<tokio::sync::Mutex<()>>>,
    conversations: Arc<dyn ConversationRepository>,
    skill_events: Arc<dyn SkillEventRepository>,

    // 新增per-topic 子代理完成队列
    sub_done_queues: HashMap<String, mpsc::Sender<SubagentResult>>,
    // 新增per-topic wait 唤醒信号
    wait_wakeups: HashMap<String, Arc<tokio::sync::Notify>>,
    // 新增per-topic 等待状态
    waiting_flags: HashMap<String, bool>,
}

生命周期与 topic_serial_locks 完全一致

  • 创建:topic_serial_lock(topic_id) 时同步创建
  • 驱逐:evict_inactive_if_needed 时同步移除
  • 清理topic 删除时同步清理

2. SubagentResult 结构

struct SubagentResult {
    task_id: String,
    status: SubagentStatus,        // completed/failed/timeout/cancelled
    output: String,                // 当前子代理的输出(与 task 工具返回格式一致)
    pending_task_ids: Vec<String>, // 未完成的子代理 id 列表
}

关键设计:每个完成回调都带上 pending_task_idsLLM 据此判断全局进度。

3. pending_subagents 表(新增)

CREATE TABLE pending_subagents (
    task_id TEXT PRIMARY KEY,
    parent_session_id TEXT NOT NULL,
    parent_topic_id TEXT NOT NULL,
    parent_chat_id TEXT NOT NULL,
    parent_channel TEXT NOT NULL,
    def_name TEXT,
    spawned_at INTEGER NOT NULL,
    status TEXT NOT NULL DEFAULT 'running'  -- running/completed/failed/interrupted/cancelled
);

四、执行流程

阶段 ① 发起 task

用户消息 → bus → process_one → lock(topic) → agent.process()
  → LLM 调用 task 工具
  → TaskTool::execute():
      ① INSERT pending_subagents (status='running')
      ② tokio::spawn 子 agent后台执行
      ③ 立即返回 tool_result:
         "running, task_id=t1。
          调用 wait_for_subagents(timeout_secs) 等待子代理完成并获取结果。"
  → LLM 看到 running + 引导,自主决策下一步

LLM 的 3 种可能行为

  • 调用 wait_for_subagents主动等待— 主路径
  • 调用其他工具(继续工作)— 罕见
  • 生成最终回答(不调 wait— 后端兜底

阶段 ② 后台子 agent 执行

tokio::spawn 独立运行:
  → SubAgentEmitter 实时推消息流给前端
  → Semaphore 限流(默认 8 并发)
  → 执行完成或失败或超时

阶段 ③ 子 agent 完成回调(直达 sub_done_q

后台 task 结束:
  ① 查询 pending_subagents WHERE parent_topic_id=? AND status='running'
     → 得到未完成列表 ["t2", "t3"]
  ② sub_done_q.send(SubagentResult{
       task_id: "t1",
       status: "completed",
       output: "{与 task 工具一致的输出格式}",
       pending_task_ids: ["t2", "t3"]
     })
  ③ UPDATE pending_subagents SET status='completed' WHERE task_id='t1'

关键:子代理结果不写 SQLite、不入 bus,直接通过 sub_done_q 传递给 wait。

顺序保证:查询未完成 → send → UPDATE

阶段 ④ wait 工具(释放锁 + select! + 重获取锁)

LLM 调用 wait_for_subagents(timeout_secs):
  ① 查询 pending_subagents
     → 无 running → 返回 "无需等待",继续循环
     → 有 running → 进入等待

  ② 释放锁,标记等待
     drop(lock_guard)
     session.set_waiting(topic_id, true)

  ③ select! 等待(不持锁)
     select! {
       result = sub_done_q.recv() => Event::Subagent(result),
       _ = wakeup.notified()       => Event::UserMessage,
       _ = timeout(60s)            => Event::Timeout,
     }

  ④ 重新获取锁,清除标记(顺序重要:先获取锁,再清除)
     let lock_guard = lock(topic_id).await
     session.set_waiting(topic_id, false)

  ⑤ 返回 tool_result
     Subagent(r) => "子代理 {task_id} 完成: {output}。未完成: {pending}"
     UserMessage => "有新用户消息到达"
     Timeout     => "等待超时。未完成子代理: {pending}"

阶段 ⑤ process_one 分流wait 期间注入用户消息)

async fn process_one(&self, inbound: InboundMessage) {
    let guard = lock(topic_id).await;
    
    if session.is_waiting(topic_id) {
        // wait 期间:只注入消息到 history不启动 process
        session.inject_to_history(inbound).await;  // INSERT + 更新内存 history
        drop(guard);
        session.wakeup(topic_id);  // 唤醒 wait
        return;
    }
    
    // 正常处理topic 空闲或 agent 已退出)
    agent.process().await;
}

阶段 ⑥ process 退出兜底

process 退出前:
  if waited=true:
    → 不发 ExecutionCompleted等 bus 消息触发新轮)
  else if pending_subagents 有 running:
    → 不发 ExecutionCompletedLLM 没调 wait 但有 pending
  else:
    → 发 ExecutionCompleted

五、wait 工具的完整设计

工具定义

// 工具名wait_for_subagents
// 参数timeout_secs可选默认 60s
// 描述:等待子代理完成或用户消息。用于有 pending 子代理时进入等待状态。

wait_for_subagents(timeout_secs: Option<u64>) -> ToolResult {
    let timeout = Duration::from_secs(timeout_secs.unwrap_or(60));
    let pending = query_pending_task_ids(topic_id);
    
    if pending.is_empty() {
        return ToolResult {
            output: "无 pending 子代理,无需等待".to_string(),
            metadata: { wait_marker: false }
        };
    }
    
    // 释放锁,标记等待
    drop(lock_guard);
    session.set_waiting(topic_id, true);
    
    // 等待(不持锁)
    let event = select! {
        result = sub_done_q.recv() => Event::Subagent(result),
        _ = wakeup.notified() => Event::UserMessage,
        _ = sleep(timeout) => Event::Timeout,
    };
    
    // 重新获取锁,清除标记(顺序重要)
    let lock_guard = lock(topic_id).await;
    session.set_waiting(topic_id, false);
    
    // 先 drain sub_done_q 积压(处理多子代理同时完成)
    match event {
        Event::Subagent(result) => {
            format_subagent_result(result)
        }
        Event::UserMessage => {
            "有新用户消息到达".to_string()
        }
        Event::Timeout => {
            let pending = query_pending_task_ids(topic_id);
            format!("等待超时。未完成子代理: {:?}", pending)
        }
    }
}

wait 的三种返回场景

事件 返回内容 LLM 行为
子代理完成 "子代理 t1 完成: {output}。未完成: [t2, t3]" 处理结果,看 pending 决定继续 wait 或综合
用户消息到达 "有新用户消息到达" LLM 下一轮看到注入的 user 消息
超时 "等待超时。未完成子代理: [t1, t2]" 决定继续等或放弃

六、两个通道的职责(核心设计)

sub_done_qwait 专属消费

后台子代理完成 → sub_done_q.send(result)  ← 不获取锁,直达 wait
                    ↓
wait select! → 立即收到 → 返回给 LLM

不经过 SQLite、不经过 bus,直接队列传递。

bus用户消息wait 期间分流注入

用户消息 → bus.publish_inbound(UserMessage)
  → process_one → lock(topic)
      ├─ is_waiting=true → 注入 history + wakeup 唤醒 wait
      └─ is_waiting=false → 正常启动 process

wait 期间:用户消息持锁注入 historywakeup 唤醒 wait。 非 wait 期间:正常走 process_one。

七、LLM 看到的信息流

时机 LLM 看到的内容
task 调用后 tool_result("running, task_id=t1。调用 wait_for_subagents 等待...")
wait 调用后 tool_result("已进入等待2 个子代理运行中: [t1,t2]")
t1 完成wait 返回) tool_result("子代理 t1 完成: {output}。未完成: [t2]")
用户消息wait 返回) tool_result("有新用户消息到达") + 下轮 user("用户消息")
t2 完成wait 返回) tool_result("子代理 t2 完成: {output}。未完成: []")

LLM 通过 pending_task_ids 能判断:

  • [] 空列表 → 全部完成,综合回答
  • ["t3"] 非空 → 还有未完成的,继续调 wait

八、消息的完整路径

子代理结果

后台 task 完成 → sub_done_q.send(SubagentResult)
  → wait select! 立即收到 → 返回 tool_result 给 LLM

直达,不绕路

用户消息

所有 channel 消息 → bus.publish_inbound(UserMessage)
  → process_one → lock(topic)
      ├─ is_waiting=true → 注入 history + wakeup 唤醒 wait
      │                   → wait 返回 "有新用户消息到达"
      │                   → LLM 下一轮看到注入的 user 消息
      └─ is_waiting=false → 正常启动 process

两条路径职责清晰,互不干扰

九、三个边界问题的解法

边界 1wakeup vs timeout 竞态

timeout 和 wakeup 同时触发 → select! 随机选一个
如果 timeout 赢:
  - 用户消息已注入 historyprocess_one 在锁内完成注入)
  - wait 返回 "超时"
  - LLM 下一轮 reload history 会看到用户消息
  - 不会丢失

解法无需特殊处理消息已持久化reload 能读到。

边界 2is_waiting 清除时机

错误顺序:清除 is_waiting → 获取锁
  → 中间窗口 process_one 看到 is_waiting=false → 启动 process → 死锁

正确顺序:获取锁 → 清除 is_waiting
  → 持锁后才清除process_one 在锁外看到 is_waiting=true → 走注入路径

解法let guard = lock().await; set_waiting(false); 顺序保证。

边界 3多条用户消息

wait 释放锁后3 条用户消息依次到达:
  msg1 → lock → 注入 → drop → wakeup
  msg2 → lock → 注入 → drop → wakeup
  msg3 → lock → 注入 → drop → wakeup

wakeup 是 Notify3 次 notify 只存储 1 个
wait 被唤醒 1 次,但 history 已有 3 条消息
wait 返回 "有新用户消息到达"
LLM 下一轮看到 3 条 user 消息

解法Notify 合并是正确行为history 完整。

十、时序场景表现

场景 1多子代理并发完成

① task(t1, t2, t3) → running
② wait → 释放锁 → select!
③ t1, t2, t3 几乎同时完成 → sub_done_q: [t1, t2, t3]
④ wait select 收到 t1 → 返回 "t1 完成。未完成: [t2, t3]"
⑤ LLM 处理 t1 → 再调 wait
   wait 先 try_recv → 立即拿到 t2 → 返回 "t2 完成。未完成: [t3]"
⑥ LLM 处理 t2 → 再调 wait
   wait 先 try_recv → 立即拿到 t3 → 返回 "t3 完成。未完成: []"
⑦ LLM 综合 t1+t2+t3 回答

3 轮 wait但无延迟try_recv 立即返回积压)

场景 2wait 期间用户发消息

① task(t1) → running → wait → 释放锁 → select!
② 用户发消息 → bus → process_one → lockwait 已释放,获取成功)
   → is_waiting=true → 注入 history + wakeup
③ wait 被 wakeup 唤醒 → 重新获取锁 → 返回 "有新用户消息到达"
④ LLM 看到 [wait_result, user("用户消息")] → 处理用户消息

用户消息即时响应wait 被唤醒

场景 3子代理完成 + 用户消息同时到达

① task(t1) → running → wait → 释放锁 → select!
② t1 完成 → sub_done_q
   用户发消息 → bus → process_one → 注入 + wakeup
③ select! 随机选一个:
   - 选 sub_done_q → 返回 "t1 完成" → LLM 下一轮看到 user 消息
   - 选 wakeup → 返回 "有新用户消息" → LLM 下一轮看到 t1 结果try_recv

两种情况都不丢失消息

场景 4父 agent 生成最终回答(不调 wait

① task(t1) → running
② LLM 生成最终回答(无 tool_calls→ process 退出
③ 退出前检查: pending 有 running → 不发 ExecutionCompleted
④ t1 完成 → sub_done_q.send无人消费留在队列
   → 触发新轮 process_one通过 bus 或定期检查)
⑤ 新轮 LLM 看到 [自己之前的回答] → try_recv sub_done_q → 处理 t1 结果

注意:此场景需要额外机制触发新轮 processsub_done_q 有残留时)。可通过:

  • 后台 task 完成后同时 publish_inbound(bus) 作为触发信号
  • 或 wait 退出兜底检查 sub_done_q

场景 5wait 超时

① task(t1) → running → wait(60s) → 释放锁 → select!
② 60s 内无事件 → timeout 触发
③ wait 重新获取锁 → 返回 "等待超时。未完成: [t1]"
④ LLM 决定继续等(再调 wait或放弃

十一、关键机制

1. 队列与状态生命周期管理

// session_history.rs

/// 获取或创建该 topic 的 sub_done_q sender
pub(crate) fn sub_done_queue(&mut self, topic_id: &str) -> mpsc::Sender<SubagentResult> {
    self.sub_done_queues
        .entry(topic_id.to_string())
        .or_insert_with(|| {
            let (tx, _rx) = mpsc::channel(32);
            tx
        })
        .clone()
}

/// 获取或创建该 topic 的 wakeup
pub(crate) fn wait_wakeup(&mut self, topic_id: &str) -> Arc<Notify> {
    self.wait_wakeups
        .entry(topic_id.to_string())
        .or_insert_with(|| Arc::new(Notify::new()))
        .clone()
}

/// 设置等待状态
pub(crate) fn set_waiting(&mut self, topic_id: &str, waiting: bool) {
    self.waiting_flags.insert(topic_id.to_string(), waiting);
}

/// 检查等待状态
pub(crate) fn is_waiting(&self, topic_id: &str) -> bool {
    self.waiting_flags.get(topic_id).copied().unwrap_or(false)
}

// evict_inactive_if_needed 中同步清理
fn evict_inactive_if_needed(&mut self) {
    // ... 现有逻辑
    if let Some(tid) = to_evict.cloned() {
        self.topic_histories.remove(&tid);
        self.topic_serial_locks.remove(&tid);
        self.sub_done_queues.remove(&tid);      // 新增
        self.wait_wakeups.remove(&tid);          // 新增
        self.waiting_flags.remove(&tid);         // 新增
    }
}

2. pending_subagents 状态管理

状态 含义
running 后台执行中
completed 正常完成
failed 执行失败
interrupted 进程崩溃,启动扫描标记
cancelled 用户取消

3. 崩溃恢复

启动扫描

UPDATE pending_subagents SET status='interrupted' WHERE status='running';

history 加载时

  • 遇到 tool_result("running, task_id=...") → 查 pending_subagents
  • status=interrupted → 替换为"子代理因重启中断,请决定是否重新发起"

两步操作之间崩溃(查询未完成 → send → UPDATE

  • 崩溃在任何点:启动扫描标记 interruptedhistory 加载时替换占位
  • 最坏情况sub_done_q 消息丢失wait 未消费),但 pending 表状态可恢复

4. 取消传播

  • task_id → CancellationToken 映射
  • 用户取消 → 遍历 pending → 触发 cancel_token
  • 子代理退出 → sub_done_q.send(SubagentResult{status:"cancelled"})

5. 去重

  • DefaultSubAgentRuntime 维护 completed_tasks: HashSet<task_id>
  • 入队前检查,已入队的不再入

6. topic_serial_lock 一致性

  • 用户消息走 bus用 topic_id 作 lock_key
  • 子代理结果走 sub_done_q不需要 lock直达 wait
  • wait 释放锁后bus 的 process_one 能获取锁注入消息

十二、改动范围

文件 改动
session_history.rs 新增 sub_done_q + wait_wakeup + waiting_flagsper-topic+ 生命周期
runtime.rs spawn 异步化 + 完成回调(含 pending_task_ids注入 sub_done_q
task_tool.rs(新) TaskTool 返回 running + 引导提示 + spawn
wait_tool.rs(新) wait 释放锁 + select! + 重获取锁 + 返回 pending 信息
agent_loop.rs 传递 lock_guard 给 wait 工具(支持释放/重获取)
processor.rs process_one 分流is_waiting 检查 + 注入 + wakeup
db.rs pending_subagents 表 + 启动扫描
history.rs running 占位替换
cancel.rs 取消传播

预估总改动:~800 行

十三、与之前方案的对比

维度 v6bus + break v8双通道 + wait 释放锁)
子代理结果通道 bus绕 SQLite → reload sub_done_q 直达
用户消息通道 bus bus
wait 语义 break 退出 select! 真等待
wait 期间用户消息 等锁agent 退出后) 即时注入 + wakeup
is_waiting 判断 不需要 需要(持锁后判断,无 TOCTOU
TOCTOU 风险 (持锁后判断)
死锁风险 wait 释放锁)
LLM 调用次数 2batch/ N+2 2-N
process 退出 每次都退出 reload wait 不退出,连续推理
上下文连贯 每轮 reload 同 process 内连续
实现复杂度 中(~20 行新增 + 边界解法)
channel 适配器改动

十四、不做的事

  • 模式1中断插入
  • user_msg_q用户消息走 bus
  • 嵌套层异步
  • 子代理结果走 bus改回 sub_done_q
  • wait 不持锁(改为释放锁 + select + 重获取)

十五、方案定位总结

机制 定位 通道
wait 释放锁 + select! 主路径(真等待,子代理直达) sub_done_q
process_one 分流 + wakeup 用户消息路径(注入 + 唤醒) bus
process 退出兜底pending 检查) 安全网LLM 不调 wait 时) -

三机制分工清晰:

  • wait select! 是核心(子代理结果直达 + 用户消息 wakeup
  • process_one 分流是用户消息路径wait 期间注入 + 唤醒)
  • 退出兜底是安全网(防止 LLM 不调 wait 时丢消息)

十六、配置说明

{
  "subagent": {
    "max_concurrent": 8,
    "spawn_timeout_secs": 300,
    "wait_default_timeout_secs": 60
  }
}
配置项 默认值 说明
max_concurrent 8 并发子 agent 数量Semaphore
spawn_timeout_secs 300 子 agent 执行超时
wait_default_timeout_secs 60 wait 工具默认超时