猫咖
首页博客工具
搜索
语言
切换网站风格
选择主题颜色
点击特效
主题

猫咖 · 持续更新中,源码见 GitHub。

Index 15: Agent Teams — MessageBus / 收件箱 / 权限冒泡

2026年8月13日
AI智能体Kotlin后端

系列

使用Kotlin从0开发一个ClaudeCode

系列

使用Kotlin从0开发一个ClaudeCode

进度 15 / 21

使用Kotlin从0开发一个ClaudeCode

上一篇

Index 14: Cron Scheduler - 持久化调度 / 会话级触发

下一篇

Index 16: Team Protocols — 关机握手 / 计划审批

目标

s06 的子智能体是「匿名、fire-and-forget」模型:主智能体派发一个 prompt,子智能体跑完返回,单向、无命名、互不通信。有一类场景它覆盖不了:需要多个命名智能体长期协作——研究员做调研、审查员审报告、编译员跑构建,彼此交换中间结果、互相提问。这类场景需要的是「团队」,不是「一次性任务」。

Index 15 引入命名智能体团队,对齐 Claude Code 的 Agent(spawn 命名 teammate)+ SendMessage(消息路由)+ 后台 teammate 的权限上浮机制:

  1. MessageBus — 按名字路由消息到接收方收件箱
  2. TeamInbox — 每智能体一个挂起式消息队列(kotlinx Channel),自动投递
  3. PermissionBubble + PermissionBroker — teammate 命中 ASK 权限时冒泡到主用户决策
  4. TeamFactory — 隔离组装 + 消费协程("resume from transcript")
  5. team_spawn / send_message 工具 + /team 命令 — 主智能体与人类的操作入口

完成后的效果——主智能体可以像 Claude Code 一样派生命名 teammate、互相发消息、权限请求自动冒泡:

TEXT
$ cat-code
> 派生一个 researcher 调查 s14 的 cron 实现,再让 reviewer 审阅它的报告
🤖 spawn teammate 'researcher' ...         ← LLM 调 team_spawn
🤖 spawn teammate 'reviewer' ...
> /team
Team members (2):
  researcher  [IDLE]
  reviewer    [IDLE]
> /team messages                            ← 回合间 drain 主收件箱
📨 researcher -> main: s14 用会话级轮询+持久化存储,nextFire 缓存推进...
> /team permissions                         ← teammate 冒泡的权限请求
Permission request from 'researcher':
  tool=bash  rule=DangerousCommandRule
  input: {"command":"rm -rf build/"}
  [a]llow / [d]eny / [A]llways / [D]eny-always? d
📨 researcher -> main: 被拒绝,改用 ./gradlew clean
> exit
Stopping team scope and closing main inbox
Goodbye! 🐾

为什么需要

现实问题

s06 的三个缺口,s15 逐一补上:

  • 子智能体无名字 — TaskTool 返回的是 task_id(task_xxx),主智能体只能用轮询追踪状态。团队需要的是稳定的身份:"让 researcher 继续","给 reviewer 发最新报告"——名字即通信地址。
  • 子智能体单向返回 — s06 的模型是「派发 → 跑完 → 轮询取结果」,子智能体之间无法对话。团队协作天然是多轮、多向的:研究员发现问题问审查员,审查员让研究员补充证据。
  • 后台无权限模型 — s06 有意不给子智能体注入权限管线(后台无界面,弹窗会死锁)。但长寿命 teammate 可能执行需要审批的危险操作(删文件、跑网络命令),完全放开太危险。需要一个「后台请求 → 上浮主用户」的桥。

依赖图表明 s15 是阶段四的起点,s06 提供隔离组装范式,s12 提供任务状态机范式:

TEXT
s06 Subagent ──→ s15 AgentTeams ──→ s16 TeamProtocols ──→ s17 AutonomousAgents
     (隔离组装)      (通信+冒泡)       (关机握手/计划审批)    (空闲循环/自动认领)

设计原则

原则取舍依据
名字即路由键MessageBus 按名字查 inbox 表,TeamStore 保证名字唯一。名字持续有效直到停止——支撑"names keep working"
消息用 Channel 而非队列s13 NotificationQueue 的 drain 模式是被动展示副作用;TeamInbox 是控制流——消费协程需挂起/唤醒,Channel.receive() 天然满足
消费循环驱动 resumeteammate 跑完首轮进入 IDLE,挂起在 receive();消息到达即唤醒,AgentLoop.run(msg) 把消息 append 进历史续跑
冒泡经独立 Broker权限请求不走 MessageBus(后者是 LLM 级消息),走 PermissionBroker 专用队列——决策是同步的、控制流的,不是对话文本
回合间 drain(对齐 s13)主智能体正忙时 teammate 挂起等待;权限请求在主 REPL 回合间 drain 询问用户,实时打断留后续 Index
每 teammate 完全隔离独立 TodoStore/ToolRegistry/HookManager/ApprovalStore——一个 teammate 的 always-allow 不污染其他,与 s06 同构
禁止递归 spawnteammate 的工具集不注册 team_spawn,与 s06 禁 task 递归同理,防资源失控

核心设计与实现

架构全景

TEXT
                          ┌──────────────────────────────────────┐
                          │            ReplLoop (主)              │
                          │  ┌─────────────┐  ┌───────────────┐  │
                          │  │ main AgentLoop│  │ PermissionBroker│ │
                          │  │  (REPL 驱动)  │  │  (drain 给用户) │ │
                          │  └──────┬───────┘  └───────▲───────┘  │
                          │         │                  │ pending  │
                          │  回合间 drain            resolveAll      │
                          │  (主 inbox 通知)        (userPrompter)   │
                          └─────────┼──────────────────┼───────────┘
                                    │                  │
            ┌───────────────────────┼──────────────────┼────────┐
            │              MessageBus (name -> inbox)   │        │
            │   send(to,msg) 路由到 recipient.inbox    │        │
            └──────────────┬─────────────────────────┬─┘        │
                           │                          │          │
                  ┌────────▼─────────┐       ┌─────────▼──────┐  │
                  │  researcher      │       │  reviewer       │  │
                  │  TeamMember      │       │  TeamMember     │  │
                  │  ┌────────────┐  │       │  ┌────────────┐ │  │
                  │  │ TeamInbox  │◄─┼───────┼──│ TeamInbox  │◄┼──┘ send_message
                  │  │ (Channel)  │  │ send  │  │ (Channel)  │ │
                  │  └─────▲──────┘  │       │  └─────▲──────┘ │
                  │        │receive  │       │        │receive │
                  │  ┌─────┴──────┐  │       │  ┌─────┴──────┐ │
                  │  │ AgentLoop  │  │       │  │ AgentLoop   │ │
                  │  │ (隔离)     │  │       │  │ (隔离)      │ │
                  │  │ PermBubble │──┼───────┼──┘             │ │
                  │  └────────────┘  │ ASK   │  ┌────────────┐ │
                  │                  │       │  │PermBubble  │ │
                  └──────────────────┘       └──┴────────────┘─┘
                                                  │ submit
                                                  ▼
                                          PermissionBroker.pending
                                          (→ 冒泡到主用户)

模块依赖(符合 5.2 规则 team -> subagent, task):team 依赖 agent(AgentLoop)、llm、tool、permission(Pipeline + UserPrompter)、system(PromptSegment)、todo、hooks;team 不被 agent 反向依赖(增强层可插拔)。

新文件一览(src/main/kotlin/com/sepcai/code/team/ 14 个 + repl/commands/TeamCommand.kt):

文件职责
TeamMessage.kt消息 data class(from/to/content/summary/timestamp)
TeamStatus.kt状态机:PENDING → RUNNING ↔ IDLE → COMPLETED/FAILED
TeamMember.ktmember 不可变快照记录(仅可观测状态,不持运行态引用)
TeamInbox.kt挂起式消息队列(Channel<TeamMessage>(UNLIMITED))
MessageBus.ktname→inbox 路由表 + SendResult
TeamStore.ktmember 注册表 + 状态机(ConcurrentHashMap + computeIfPresent)
TeamFactory.kt隔离组装 + 消费协程 + SpawnResult + stop
TeamPrompts.ktteammate 系统提示(名字、通信工具、冒泡说明)
TeamSpawnTool.ktLLM 工具:派生命名 teammate
SendMessageTool.ktLLM 工具:发消息给 teammate 或 "main"
PermissionRequest.kt冒泡请求(id/memberName/toolCall/ruleName/deferred)
PermissionBubble.ktteammate 的 UserPrompter 实现(冒泡+挂起)
PermissionBroker.kt主侧请求队列(submit/drainPending/resolveAll)
repl/commands/TeamCommand.kt/team list/messages/permissions/stop

接线变更:ReplLoop(teamScope/messageBus/teamStore/broker/主 inbox/工具注册/回合间 drain/退出清理)、ReplContext(+5 字段)、AgentConfig(+maxTeamMembers)、CronCreateTool(注入 clock,修 pre-existing flaky)。


逐类拆解

TeamMessage / TeamStatus / TeamMember(模型层)

KOTLIN
// TeamMessage.kt
data class TeamMessage(
    val from: String,        // 发送方名字("main" 或 teammate 名字)
    val to: String,          // 接收方名字
    val content: String,     // 消息正文,作为 user 消息传给接收方 AgentLoop
    val summary: String = "",// 可选摘要——接收方为 main 时 drain 通知展示用
    val timestamp: Long = System.currentTimeMillis()
)
KOTLIN
// TeamStatus.kt
enum class TeamStatus {
    PENDING, RUNNING, IDLE, COMPLETED, FAILED
}

关键设计点——IDLE 是 s15 与 s06 的本质差异:s06 SubagentStatus 只有 PENDING/RUNNING/COMPLETED/FAILED,subagent 跑完即终态销毁。team member 跑完进入 IDLE 而非终态——名字持续有效,下一条消息从历史续跑。TeamStore.markRunning 因此接受 PENDING 或 IDLE 作为前置(resume)。

KOTLIN
// TeamMember.kt —— 只存可观测快照,不持运行态引用
data class TeamMember(
    val name: String,
    val status: TeamStatus,
    val createdAt: Long,
    val completedAt: Long? = null,
    val error: String? = null
)

运行态的 AgentLoop / TeamInbox 引用由 TeamFactory 的协程闭包持有,不进 record——与 SubagentRecord 一致的「不可变快照」设计,读取方拿到的总是某一时刻的完整状态。

TeamInbox(挂起式队列)

KOTLIN
// TeamInbox.kt
class TeamInbox {
    private val channel = Channel<TeamMessage>(Channel.UNLIMITED)
    private val counter = AtomicInteger(0)

    fun deliver(message: TeamMessage): Boolean {          // trySend + 计数
        val result = channel.trySend(message)
        if (result.isSuccess) { counter.incrementAndGet(); return true }
        return false                                       // inbox 已关闭,消息丢弃
    }

    suspend fun receive(): TeamMessage? {                 // 挂起到有消息
        val result = channel.receiveCatching()
        if (result.isSuccess) counter.decrementAndGet()
        return result.getOrNull()                          // null = 已关闭且取空
    }
    // poll(): 非阻塞取队首;drain(): 批量清空;close(): 关闭;pendingCount()
}

为何用 Channel 而非 s13 的 ConcurrentLinkedQueue:NotificationQueue 的 KDoc 明确写了"为何不用 Channel/Flow——通知是尽力展示的副作用,不参与控制流"。TeamInbox 恰恰相反,它是控制流:消费协程需要"无消息时挂起、有消息时唤醒",Channel.receive() 天然满足。两种队列并存是刻意的——一个被动 drain,一个主动挂起,语义不同(详见开发过程记录的设计矛盾)。

关闭语义:close() 后已缓冲消息仍可取出,取空后 receive() 返回 null——消费循环据此退出(receive() ?: break),实现 /team stop 优雅停止单个 teammate 而不取消整个 teamScope。

MessageBus(路由)

KOTLIN
// MessageBus.kt
class MessageBus {
    private val inboxes = ConcurrentHashMap<String, TeamInbox>()

    fun register(name: String, inbox: TeamInbox) { inboxes[name] = inbox }
    fun unregister(name: String) { inboxes.remove(name) }
    fun hasRecipient(name: String): Boolean = inboxes.containsKey(name)

    fun send(message: TeamMessage): SendResult {
        val inbox = inboxes[message.to]
            ?: return SendResult.UnknownRecipient(message.to)
        if (inbox.deliver(message)) return SendResult.Delivered
        inboxes.remove(message.to)   // 投递到死信箱 -> 自愈注销
        return SendResult.UnknownRecipient(message.to)
    }
}

sealed class SendResult {
    object Delivered : SendResult()
    data class UnknownRecipient(val name: String) : SendResult()
}

关键设计点:

  • bus 不依赖 TeamStore——两者按名字各自索引,spawn 时由 TeamFactory 同时向 bus 注册 inbox、向 store 注册 member record。避免循环依赖,单一职责。
  • 投递失败自愈:若 recipient 已注册但 inbox 已关闭(teammate 已停止但未清理),deliver 返回 false,bus 自动注销该名字并报告未知——后续发送快速失败,而非反复投递到死信箱。
  • "main" 是保留名:主智能体的 inbox 在 ReplLoop 初始化时 messageBus.register("main", mainInbox) 注册,teammate 可用 send_message(to="main") 向主智能体汇报。

TeamStore(注册表 + 状态机)

仿 SubagentStore:ConcurrentHashMap + computeIfPresent 状态校验 + AtomicInteger 计数。

KOTLIN
// TeamStore.kt(节选)
fun markRunning(name: String) {
    members.computeIfPresent(name) { _, m ->
        if (m.status == TeamStatus.PENDING || m.status == TeamStatus.IDLE) {
            runningCount.incrementAndGet()
            m.copy(status = TeamStatus.RUNNING)
        } else m
    }
}

计数语义(易错点):runningCount 只统计 RUNNING。markRunning 从 PENDING/IDLE 进入时 +1,markIdle 从 RUNNING 进入时 -1,markCompleted/markFailed 仅当从 RUNNING 转入才 -1(从 IDLE 转入不计,避免下溢)。activeCount() 统计非终态成员(PENDING/RUNNING/IDLE),用于 team_spawn 的 maxTeamMembers 上限。

TeamFactory(隔离组装 + 消费协程)—— s15 的心脏

仿 SubagentFactory 的隔离模式,但增强三点:注册 SendMessageTool、接入 PermissionPipeline+PermissionBubble、启动消费协程。

KOTLIN
// TeamFactory.kt(节选)
class TeamFactory(
    private val llmProvider: LLMProvider,
    private val config: AgentConfig,
    private val messageBus: MessageBus,
    private val permissionBroker: PermissionBroker,
    private val scope: CoroutineScope
) {
    fun spawn(name: String, prompt: String, store: TeamStore): SpawnResult {
        if (!store.register(name)) {
            return SpawnResult.NameTaken(name)      // 名字唯一性由 putIfAbsent 原子保证
        }
        val inbox = TeamInbox()
        messageBus.register(name, inbox)
        inboxes[name] = inbox
        val loop = buildIsolatedLoop(name)
        scope.launch { runConsumerLoop(name, prompt, loop, inbox, store) }
        return SpawnResult.Spawned(name)
    }

    fun stop(name: String): Boolean {              // /team stop 优雅停止
        val inbox = inboxes.remove(name) ?: return false
        inbox.close()                               // receive() 返回 null,循环退出
        return true
    }
}

消费协程("resume from transcript")——这是「names keep working after an agent completes」的实现核心:

KOTLIN
// TeamFactory.kt(节选)
private suspend fun runConsumerLoop(
    name: String, prompt: String, loop: AgentLoop, inbox: TeamInbox, store: TeamStore
) {
    try {
        store.markRunning(name)
        loop.run(prompt)                            // 首轮:执行初始任务
        store.markIdle(name)
        while (coroutineContext.isActive) {
            val msg = inbox.receive() ?: break      // null = inbox 已关闭
            store.markRunning(name)
            loop.run(formatMessage(msg))            // 消息作为新 user 消息续跑历史
            store.markIdle(name)
        }
        store.markCompleted(name)
    } catch (e: CancellationException) {
        throw e                                     // 结构化并发:scope 取消即传播
    } catch (e: Exception) {
        store.markFailed(name, e.message ?: "Unknown error")
    } finally {
        messageBus.unregister(name)                 // 停止后名字不再路由
    }
}

AgentLoop.run(userInput) 本身就把输入 append 进 messagesHistory 并跑到 END_TURN——所以 loop.run(formatMessage(msg)) 天然实现了「从历史续跑」。teammate END_TURN 后进入 IDLE 等消息,收到消息即续跑。

隔离组装(buildIsolatedLoop):

KOTLIN
// TeamFactory.kt(节选)
private fun buildIsolatedLoop(name: String): AgentLoop {
    val todoStore = TodoStore()
    val bubble = PermissionBubble(permissionBroker, memberName = name)
    val pipeline = PermissionPipeline(
        rules = listOf(DangerousCommandRule(), PathAllowlistRule()),
        approvalStore = ApprovalStore(),           // 每 teammate 独立,隔离 always-allow
        userPrompter = bubble                      // ASK -> 冒泡,非 ReadLinePrompter
    )
    val toolRegistry = ToolRegistry().apply {
        register(ReadFileTool()); register(WriteFileTool()); register(BashTool())
        register(TodoWriteTool(todoStore))
        register(SendMessageTool(messageBus, from = name))   // 双向通信
        // 不注册 team_spawn —— 禁止递归派生
    }
    val hookManager = HookManager().apply {
        LoggingHook().registerTo(this)             // 不注册 TodoDisplayHook,避免干扰主 REPL
    }
    return AgentLoop(
        llmProvider = llmProvider,
        systemPromptBuilder = SystemPromptBuilder().apply {
            register(PromptSegment("base", 0) { TeamPrompts.forName(name) })
        },
        config = config,
        hooks = AgentLoopHooks(
            onBeforeToolExecute = { tc -> pipeline.approve(tc) },
            onPreToolUse = { tc -> hookManager.firePreToolUse(tc) },
            onPostToolUse = { tc, result -> hookManager.firePostToolUse(tc, result) }
        ),
        toolRegistry = toolRegistry
    )
}

与 s06 SubagentFactory 的权限差异(刻意):s06 不注入 PermissionPipeline(无界面会死锁);s15 注入,但用 PermissionBubble 替代 ReadLinePrompter——DENY 由本地规则直接拦截(安全规则本地执行),只有 ASK 才冒泡。这是「后台 agent 也能安全执行」的正确模型。

PermissionBubble + PermissionBroker(权限冒泡)

KOTLIN
// PermissionBubble.kt —— teammate 的 UserPrompter
class PermissionBubble(
    private val broker: PermissionBroker,
    private val memberName: String
) : UserPrompter {
    override suspend fun ask(toolCall: ToolCall, ruleName: String): PermissionDecision {
        val request = PermissionRequest(
            id = "perm_${UUID.randomUUID()}",
            memberName = memberName,
            toolCall = toolCall,
            ruleName = ruleName,
            deferred = CompletableDeferred()
        )
        return broker.submitAndAwait(request)   // teammate 协程在此挂起
    }
}
KOTLIN
// PermissionBroker.kt —— 主侧队列
class PermissionBroker {
    private val pending = ConcurrentLinkedQueue<PermissionRequest>()

    suspend fun submitAndAwait(request: PermissionRequest): PermissionDecision {
        pending.add(request)
        return request.deferred.await()
    }

    fun drainPending(): List<PermissionRequest> = pending.toList()   // 快照,不改队列

    suspend fun resolveAll(userPrompter: UserPrompter): Int {
        var resolved = 0
        while (true) {
            val req = pending.poll() ?: break
            val decision = userPrompter.ask(req.toolCall, req.ruleName)
            req.deferred.complete(decision)     // 唤醒挂起的 teammate
            resolved++
        }
        return resolved
    }
}

端到端冒泡流程:

TEXT
teammate 命中 ASK
    │ PermissionBubble.ask()
    ▼
PermissionBroker.submitAndAwait() ── 请求入队,teammate 挂起在 deferred.await()
    │
    ▼  主 REPL 回合间 drainPermissionRequests()
PermissionBroker.resolveAll(ReadLinePrompter)
    │ 逐个询问用户
    ▼
userPrompter.ask(toolCall, ruleName) ── 用户 [a]llow/[d]eny/[A]lways/[D]eny-always
    │
    ▼
deferred.complete(decision) ── 唤醒 teammate,PermissionPipeline.approve 返回决策

取舍(与 s13 对齐):冒泡请求在主 REPL 回合间 drain。主智能体正忙时 teammate 等到当前 run() 返回。实时打断式提醒留后续 Index(与 s13 drainNotifications 注释同理)。不超时——teammate 无限期等主用户决策,s16 协议层可加关机/超时。

工具层

team_spawn(TeamSpawnTool,仿 TaskTool):

KOTLIN
// TeamSpawnTool.kt(节选)
override suspend fun execute(input: JsonObject): ToolResult {
    val name = input["name"]?.jsonPrimitive?.content ?: return error("name is required")
    val prompt = input["prompt"]?.jsonPrimitive?.content ?: return error("prompt is required")
    if (store.activeCount() >= maxMembers) {
        return ToolResult("", "Error: max team members ($maxMembers) reached; stop an existing teammate first", isError = true)
    }
    return when (val result = factory.spawn(name, prompt, store)) {
        is SpawnResult.Spawned -> ToolResult("", "Spawned teammate: ${result.name}\nUse send_message(...) to communicate with it.")
        is SpawnResult.NameTaken -> ToolResult("", "Error: name '${result.name}' is already in use", isError = true)
    }
}

send_message(SendMessageTool):

KOTLIN
// SendMessageTool.kt(节选)
return when (val result = bus.send(TeamMessage(from = from, to = to, content = message, summary = summary))) {
    SendResult.Delivered -> ToolResult("", "Message delivered to $to")
    is SendResult.UnknownRecipient -> ToolResult("", "Error: unknown recipient: ${result.name}", isError = true)
}

与 s06 TaskTool 的对比:

维度s06 tasks15 team_spawn
身份匿名,task_id命名,name 即通信地址
通信单向(query_task 轮询)双向(send_message)
生命周期跑完即终态跑完进入 IDLE,可续跑
权限不注入管线注入 + PermissionBubble 冒泡

/team 命令

KOTLIN
// TeamCommand.kt(节选)
return when (subcommand) {
    "", "list" -> listMembers(ctx)            // 列出成员 + 状态 + 待处理权限数
    "messages" -> drainMessages(ctx)          // drain 主 inbox 消息
    "permissions" -> drainPermissions(ctx)    // resolveAll(逐个询问用户)
    "stop" -> stopMember(name, ctx)           // factory.stop(name) 优雅停止
    else -> { println("Unknown subcommand: $subcommand"); ... }
}

错误处理

失败场景行为
send_message 到未知名字UnknownRecipient,工具返回 isError,LLM 可见
team_spawn 重名NameTaken,工具返回 isError(名字唯一性由 putIfAbsent 原子保证)
team_spawn 超 maxTeamMembers工具返回 isError,提示先 stop 一个
teammate LLM 抛异常消费协程 catch → markFailed,TeamMember.error 填充
teamScope 取消(会话退出)CancellationException 重新抛出(结构化并发),finally 注销 inbox
/team stop 未知名字命令打印 "No active teammate named 'x'"
inbox 已关闭仍投递bus 自动注销该名字,返回 UnknownRecipient(自愈)
teammate 挂起等权限无超时;主回合间 resolveAll 唤醒(文档记录,s16 可加超时)

与 s06 Subagent 的衔接

s15 复用 s06 的全部隔离范式(独立 TodoStore/ToolRegistry/HookManager、SubagentFactory.create 的组装结构),新增了「命名 + 双向 + 冒泡」三件事。TeamFactory 与 SubagentFactory 几乎同构,差异集中在三处:SendMessageTool 注册、PermissionPipeline+PermissionBubble 注入、消费协程循环。


端到端:一次完整团队生命周期

以「调查 s14 代码并审阅报告」为例,串联 s15 的全部机制。每个步骤标注对应的真实代码路径。

TEXT
Step 1  用户:调查 s14 的 cron 实现,再让一个审查员审阅报告
        │
        ▼ 主 AgentLoop.run(userInput) → LLM 决定派生 2 个 teammate
Step 2  主 LLM 调 team_spawn(name="researcher", prompt="阅读 cron 包…")
        │ TeamSpawnTool.execute → TeamFactory.spawn(TeamSpawnTool.kt:80)
        │   ├─ TeamStore.register("researcher")        → PENDING(TeamStore.kt:72)
        │   ├─ MessageBus.register("researcher", inbox)(TeamFactory.kt:87)
        │   ├─ buildIsolatedLoop("researcher")         → 独立 TodoStore/ToolRegistry/PermBubble
        │   └─ scope.launch { runConsumerLoop(...) }   → markRunning → loop.run(prompt)
        │
        ▼ 主 LLM 调 team_spawn(name="reviewer", prompt="审阅报告并给建议")
        │   同样的派生流程,teammate 各自在独立协程跑初始 prompt
        │
Step 3  researcher 读完 s14 源码,END_TURN → markIdle(TeamFactory.kt:117)
        │  消费协程挂起在 inbox.receive()(TeamInbox.kt:52)
        │
Step 4  researcher 向 main 汇报(send_message(to="main", message=…, summary="s14 用会话级轮询…"))
        │ SendMessageTool.execute → MessageBus.send → 主 inbox.deliver(TeamInbox.kt:39)
        │
        ▼ 主 REPL 回合间 drainTeamMessages()(ReplLoop.kt:455)
        │  📨 researcher -> main: s14 用会话级轮询+持久化存储…
        │
Step 5  主 LLM 让 researcher 深入(send_message(to="researcher", message="检查 nextFire 缓存…"))
        │ MessageBus.send → researcher 的 inbox.deliver
        │
        ▼ researcher 消费协程被唤醒(receive() 返回)
        │  markRunning(IDLE→RUNNING)→ loop.run("[Message from main]…")  ← 从历史续跑
        │  → markIdle(TeamFactory.kt:118-120)
        │
Step 6  reviewer 命中危险命令 BashTool("rm -rf build/")
        │ DangerCommandRule → ASK → PermissionPipeline.approve
        │ → PermissionBubble.ask → PermissionBroker.submitAndAwait(PermissionBroker.kt:40)
        │   reviewer 协程挂起在 deferred.await()
        │
        ▼ 主 REPL 回合间 drainPermissionRequests()(ReplLoop.kt:468)
        │  PermissionBroker.resolveAll(ReadLinePrompter)
        │  → userPrompter.ask(toolCall, ruleName) → 用户输入 "d"(拒绝)
        │  → deferred.complete(DENY) → reviewer 协程恢复
        │  → BashTool 返回 isError → reviewer 改用 ./gradlew clean
        │
Step 7  用户 /team 查看状态
        │  /team list → researcher [IDLE]  reviewer [IDLE]
        │  /team stop reviewer → TeamFactory.stop → inbox.close → 消费循环 break → COMPLETED
        │
        ▼ 会话退出
Step 8  ReplLoop finally:teamScope.cancel()(所有消费协程取消)+ mainInbox.close()(ReplLoop.kt:361)
        │  每个消费协程的 finally 里 messageBus.unregister(name)(TeamFactory.kt:133)
        ▼
      session 结束,团队状态已清理

这个生命周期把三个独立机制串成一条线:Step 2-3 的 spawn + 隔离消费循环 → Step 4-5 的双向消息 + resume → Step 6 的权限冒泡。三者都挂在一个时间维度上:主 REPL 的回合间 drain(Step 4、6)让后台 teammate 的异步行为被同步感知,这是 s15 与 s13 共享的核心节奏。


测试策略

测试类覆盖点
TeamInboxTest (11)deliver/receive 顺序、poll 非阻塞、drain 批量、pendingCount、receive 挂起、close 语义
MessageBusTest (9)路由到正确 inbox、未知 recipient、register/unregister、"main" 保留名、死信箱自愈、re-register 覆盖
TeamStoreTest (14)名字唯一性、PENDING→RUNNING→IDLE 状态机、resume(IDLE→RUNNING)、计数不下溢、非法转换忽略
PermissionBrokerTest (6)submitAndAwait 挂起/唤醒、drainPending 快照不改队列、resolveAll FIFO、多轮 resolve
PermissionBubbleTest (4)请求字段(memberName/toolCall/ruleName)、互异 id、全链路冒泡→唤醒
TeamFactoryTest (7)spawn 注册+IDLE、重名、消息续跑、stop 优雅停止、LLM 异常→FAILED、多成员隔离
TeamSpawnToolTest (6)派生成功、缺参、重名、max 上限、停止释放槽位
SendMessageToolTest (7)投递、未知 recipient、缺参、from 字段、summary 默认/显式
TeamCommandTest (9)list/messages/permissions/stop、未知子命令
回归改造5 个既有命令测试补 ReplContext 新字段;AgentLoopTest.FakeLLMProvider.callCount 加 @Volatile

测试设计上的坑(记录如下,供团队复盘):

  1. runBlocking 挂起协程 = 死锁。runBlocking { launch { broker.submitAndAwait(req) } }——runBlocking 会等所有子协程完成,而 submitAndAwait 挂起在 deferred.await() 永不返回。全量测试一度"卡死"超时。修复:把 resolveAll 放进同一个 runBlocking(先 launch 挂起,再 execute 触发 resolveAll,再 job.join())。这个坑在 TeamCommandTest 的 permissions 用例上踩实。

  2. IDLE 态在 resume 前后无法区分。awaitStatus(store, name, IDLE) 在消息投递后立即返回——因为消费协程还没醒,状态仍是首轮的 IDLE。可靠信号是 callCount(每次 loop.run 触发一次 chat)。改用 awaitCondition { llm.callCount >= 2 }。且 callCount 由后台线程写、测试线程读,须 @Volatile(否则 JMM 可见性问题)。

  3. finally 块晚于状态变更。markCompleted 在 try 末尾,messageBus.unregister 在 finally——awaitStatus(COMPLETED) 可能先于 unregister 返回。测试用 awaitCondition { !bus.hasRecipient(name) } 轮询,而非立即断言。

  4. pre-existing flaky 修复(顺带):CronCreateToolTest 的 one-shot 用例断言 fires at: 2026-08-11 09:00,但 CronCreateTool 用 LocalDateTime.now()(真实时间)计算 fireAt,而 fixture 的 scheduler clock 固定在 2026-08-10。测试在 8 月 10 日能过,8 月 11 日就挂(next 9am 变 8/12)。修复:给 CronCreateTool 注入 clock 参数(与 CronScheduler 同构),测试传固定时钟。这是"测试假设真实时间≈fixture 日期"的经典 flaky。


开发过程记录

设计矛盾 1:Channel vs ConcurrentLinkedQueue

s13 的 NotificationQueue 明确不用 Channel。s15 的 TeamInbox 却用了。是不是不一致?

审查后的结论:两者语义不同,并存是刻意设计。

  • NotificationQueue:后台任务结束是副作用,需要"展示出来"但不控制流程 → 队列 + 回合间 drain。
  • TeamInbox:消息是控制流,消费协程必须"无消息挂起、有消息唤醒" → Channel.receive()。

写进代码注释:两个队列的 KDoc 互相对照(NotificationQueue 解释为何不用 Channel,TeamInbox 解释为何用)。s13 的设计决策不适用于 s15,因为目的不同——这是"设计原则服从于目的"的一个好例子。

设计矛盾 2:主智能体收到 teammate 消息后是否自动续跑?

Claude Code 的 SendMessage 语义是"messages delivered automatically; you don't check an inbox"——teammate 自动续跑。那主智能体呢?

取舍:主智能体本期不在消息上自动续跑,只做回合间 drain 通知。原因:

  • 主循环由用户输入驱动(REPL 读 stdin),没有"空闲检查 inbox"的天然挂点。
  • 自动续跑意味着主智能体的 AgentLoop.run 要能在无用户输入时被消息唤醒——这其实是 s17 自治(idle loop)的范畴。

方案:主 inbox 回合间 drain 打印 📨 <from> -> main: <summary>,让用户感知 teammate 汇报;主 LLM 是否处理由用户决定。清晰标注"s17 自治智能体将接管主 inbox 的自动续跑"。

设计矛盾 3:权限冒泡走 MessageBus 还是独立 Broker?

初版设想是权限请求作为一种特殊 TeamMessage 经 MessageBus 发给 "main"。审查发现有问题:

  • MessageBus 是 LLM 级对话消息的载体,teammate 的消费协程会把消息作为 user 输入续跑 AgentLoop。
  • 权限决策是 同步、控制流的(必须等待用户决策才能继续执行工具),不是对话文本。
  • 若走 MessageBus,main 的消费协程(不存在)或 drain 逻辑要把权限请求重新解析成决策,语义混乱。

最终:权限请求走独立的 PermissionBroker(ConcurrentLinkedQueue + CompletableDeferred),与 MessageBus 完全解耦。这是一个"通道语义"的取舍——消息通道承载对话,控制通道承载决策。

设计矛盾 4:权限规则集选择

teammate 的 PermissionPipeline 用 [DangerousCommandRule, PathAllowlistRule],不放 ToolCategoryRule。

理由:ToolCategoryRule 对工具类别做整体判断(如 read_file 一律 ALLOW),这在主智能体有交互界面时可接受;teammate 无界面,规则应尽量"本地可判定"——DangerousCommandRule 拦截危险命令、PathAllowlistRule 拦路径,两者都返回确定性的 DENY/ALLOW;只有真正需要用户判断的才 ASK 冒泡。若放 ToolCategoryRule,它会把大量工具判为 ASK,全部冒泡轰炸用户。

开发中发现并修复的问题

  1. isActive 接收者错误:runConsumerLoop 是普通 suspend 函数(不是 CoroutineScope 扩展),while (isActive) 编译失败("receiver type mismatch")。改为 while (coroutineContext.isActive)——kotlinx.coroutines.isActive 有 CoroutineContext 扩展,suspend 函数内可用。
  2. TeamFactoryTest 的 PENDING 断言竞态:spawn 后立即断言 status == PENDING,但消费协程在 Dispatchers.Default 上立即推进到 RUNNING/IDLE。PENDING 是瞬态,删除该断言,只等终态。
  3. PermissionBrokerTest 手动 complete 不出队:request.deferred.complete(...) 不会把请求移出队列(只有 resolveAll 的 poll 会)。测试手动 complete 后 pendingCount 仍为 1,断言错误。改为用 resolveAll 解析(既回填又出队),并加 KDoc 说明"resolveAll 是唯一出队路径"。

关键取舍汇总

决策选择理由
消息队列kotlinx Channel(UNLIMITED)单消费者天然适配、挂起/唤醒、缓冲不丢消息
主 inbox 投递回合间 drain 通知(非注入 LLM)与 s13 对齐;自动续跑属 s17
权限冒泡独立 Broker(非 MessageBus)消息通道承载对话,控制通道承载决策
teammate 规则集Danger + Path(无 ToolCategory)本地可判定,避免 ASK 轰炸
teammate ApprovalStore每 teammate 独立隔离语义,一个的 always-allow 不污染其他
递归 spawn禁止(不注册 team_spawn)与 s06 禁 task 递归同理
挂起无超时teammate 无限期等 broker resolves15 不引入超时配置;s16 协议层可加

下一站

s15 为阶段四后续 Index 留下了什么基础:

  1. s16 Team Protocols — PermissionBroker 已有"请求-回填"通道,s16 可在此之上加 ShutdownHandshake(teammate 关机时向主智能体确认待办)+ PlanApproval(teammate 提出计划请求主审批,复用 broker 的挂起/唤醒机制)。TeamFactory.stop 的优雅停止为握手提供终止入口。
  2. s17 Autonomous Agents — 主 inbox 的回合间 drain 现在是"展示通知",s17 的 IdleLoop 可接管:主智能体空闲时从 inbox 取消息自动续跑 AgentLoop.run,加上 AutoClaimer(认领 TeamStore.all() 中的 IDLE 成员)。TeamStatus.IDLE 就是为"可认领"设计的。
  3. s18 Worktree Isolation — teammate 的隔离模式(独立 TodoStore/ToolRegistry)延伸为文件系统级隔离:每个 member 绑定一个 git worktree。
  4. s20 Comprehensive Agent — TeamFactory 的隔离组装 + 消费循环将成为"全机制归到一个循环"的骨架参考。

未完事项:权限冒泡无超时(teammate 可能长期挂起);主 inbox 自动续跑留待 s17;teammate 的消息进历史无 token 上限(长会话可能需压缩,可复用 s08 ContextCompactor)。

目录

当前章节:目标

  • 1. 目标
  • 2. 为什么需要
  • 3. 现实问题
  • 4. 设计原则
  • 5. 核心设计与实现
  • 6. 架构全景
  • 7. 逐类拆解
  • 8. 错误处理
  • 9. 与 s06 Subagent 的衔接
  • 10. 端到端:一次完整团队生命周期
  • 11. 测试策略
  • 12. 开发过程记录
  • 13. 设计矛盾 1:Channel vs ConcurrentLinkedQueue
  • 14. 设计矛盾 2:主智能体收到 teammate 消息后是否自动续跑?
  • 15. 设计矛盾 3:权限冒泡走 MessageBus 还是独立 Broker?
  • 16. 设计矛盾 4:权限规则集选择
  • 17. 开发中发现并修复的问题
  • 18. 关键取舍汇总
  • 19. 下一站
回到顶部

相关推荐

查看全部文章
Index 21: Terminal UX —— Claude Code 风格终端交互

Index 21: Terminal UX —— Claude Code 风格终端交互

2026年8月21日

Index 21 将 cat-code 终端交互升级为 Claude Code 风格,解决旧 REPL 黑箱、无中断及输入体验差的问题。通过 JLine3 与 Mordant 实现历史补全、实时工具可见性、Spinner 状态及 Esc 中断。核心采用 UI 与 AgentLoop 解耦的事件流架构,支持内联权限菜单与优雅降级。同时配置 logback 收敛控制台日志,确保 TUI 清爽且功能无损,显著提升可用性。

Index 20: Comprehensive Agent —— 全机制集成(收口)

Index 20: Comprehensive Agent —— 全机制集成(收口)

2026年8月13日

作为收官之作,把 s01-s19 的二十个核心机制整合为一个全面智能体:统一的状态查询与运行周期、自治认领与工作区隔离的贯通,以及 /status 命令的全局可视化。全机制协同运转,标志着 Cat-Code 从零完整复刻 Claude Code 核心能力的收官。

Index 19: MCP Plugin —— 多传输 / 通道路由 / 工具池组装

Index 19: MCP Plugin —— 多传输 / 通道路由 / 工具池组装

2026年8月13日

引入 MCP(Model Context Protocol)插件机制,通过多传输适配与通道路由,把外部 MCP 服务器的工具动态接入智能体的工具池。工具注册从静态编译期扩展为运行时动态组装,让智能体能力随外部服务即插即用。