diff --git a/docs/zh/resident-subagent-design.md b/docs/zh/resident-subagent-design.md index 8da0e3b..a1be4be 100644 --- a/docs/zh/resident-subagent-design.md +++ b/docs/zh/resident-subagent-design.md @@ -341,28 +341,40 @@ - **创建/销毁/回收/查看/发送**是**父可调用的原语(工具)**;**决策**(压还是收、收哪些) 在父的模型手里 —— 内核不替父决定。 -### 7.1 例外:积压任务自动转投(内核主动拉起)[已定 · 唯一例外] +### 7.1 积压任务的**及时反馈**(内核主动拉起分诊助手)[已定] -**背景(2026-09-19 线上实测)**:主 agent 被一条长任务占住时(当天现场:12 分 8 秒、 -69 次工具调用),后来的 QQ 消息全部以 `level insufficient` 排进中断队列干等 —— +**背景(2026-09-19 线上实测)**:主 agent 被一条长任务占住时(当天现场:13 分 5 秒、 +8 次 cmd_run),后来的 QQ 消息全部以 `level insufficient` 排进中断队列干等 —— 同级中断不能抢占同级运行任务(`scheduler.canPreempt`),只能等前一个跑完。 -而内核明明有驻留子(独立 agent + 独立调度器 + 共享输出通道视图)可以并行干活。 +用户在这十几分钟里**收不到任何回复**。 + +**定性(用户明确)**:这不是"内核替父决定",而是**及时反馈** —— +主 agent 忙时不该让用户干等。分诊助手的职责是: +- **简单的、不需主 agent 介入的** → 直接处理并回复; +- **需要主 agent 介入的** → 立刻回「主 agent 忙碌中,请稍候」,**不勉强作答**。 **行为**:当运行任务已持续超过 `core.agent.offload_busy_after`(默认 5m) **且**排队输入积到 `offload_min_pending`(默认 3)条时,内核: -1. 拉起(或在 `offload_max_residents` 内复用一个)**转投专用驻留子**; -2. 把积压的**纯排队输入**转投给它; -3. 在原队列位置留下一条说明:`[系统] N 条积压任务已转投给驻留子 agent X 处理…`。 +1. 拉起(或在 `offload_max_residents` 内复用一个)**分诊助手**(`OffloadOwned`); +2. 把积压的**纯排队输入**交给它先行分诊; +3. 在原队列位置留下一条说明:`[系统] N 条积压消息已在主 agent 忙期间交由临时助手 X 先行分诊…`。 -**转投驻留子的通道配置(刻意与人工创建的子不同)**: +**分诊助手的通道配置(刻意与人工创建的子不同)**: - **不配 inputch**:它是内核的干活 agent,不接收任何插件的用户输入; - **持有全部输出通道**(`AllowedOutputs` 为空 = 完整授权):它必须能把结果发回 qq/webui 等正确通道(否则干活结果无处可去)。 -**为什么这是对上述原则的例外,且可接受**:父此刻正忙(物理上无法做决策), -而积压任务**本来就是空的** —— 转投只是把「排队干等」换成「有人在做」, -不改变任何已提交决策的语义。若不做例外,这个能力就只能由父的模型发起, -而它恰恰是忙不过来的那个。 +**为什么这套机制自然(用户观察)**:分诊助手就在**同一张登记表**里 —— +父能 `inspect` 它的处理表与轮次、能 `send`、能按需 `compress`/`reclaim`/`destroy`。 +控制面 6 个动作均按 id 生效、不区分来源,因此回收策略对它自动适用。 +状态面额外暴露 `offload_owned`,让父能分清"我建的子"与"内核临时拉的助手"。 + +**残余任务由父显式决定**(用户 2026-09-19 要求):回收/销毁一个分诊助手时, +它手头可能还有尚未处理的消息。内核**不自己决定**这些消息的命运,而是: +- `residual=keep`(默认):逐条转回父自己的队列,父稍后处理; +- `residual=drop`:明确丢弃,**逐条记日志**(不可追溯的丢弃是不允许的); +- 两种路径都仍要给 `ResponseCh` 补终态,否则 cli/a2a 这类无超时同步调用方 + 会永久挂起(设计 §7 I5)。 **默认关闭**(`core.agent.offload_enabled=false`):它改变的是系统行为而非修 bug, 按「显式才是特权」(与 `scheduler.DefaultLevel` 同一条理由)由部署方打开。 @@ -373,6 +385,14 @@ - 只对**根 agent** 生效:子再去拉孙子会形成无界增殖,而积压的源头是根那条链。 - 实现上检查跑在**独立 goroutine**:`schedulerLoop` 是同步执行的, 放在那里在「正忙」期间根本不会回到循环顶部(等于永不触发)。 +- 转投时**推原事件**(`DeliverRouted`)而不是重建:重建会丢掉 `ResponseCh`, + 使同步调用方永久挂起。 + +💡 **两条容易重犯的坑(都已在实现里修掉并写进测试)**: +1. 积压可能全堆在 `io.inputCh`(因为忙时 `pumpInbox` 没被调用), + 只数 `sched.queue` 会得到 0 ⇒ 永不触发。 +2. 分诊助手必须**继承父的 SystemPrompt**:它里面写着「面向 qq 等异步通道时 + 必须显式 `output_send`,纯文本会被静默丢弃」。缺了它,子处理完却发不出去。 - **默认完整授权**[已定]:子默认拿到全部插件与工具(含输出门); 父可在创建时**收窄**(收窄工具子集、收窄可用输出通道集合)。 diff --git a/internal/agent/core/eventloop.go b/internal/agent/core/eventloop.go index 4aa8924..c30ca78 100644 --- a/internal/agent/core/eventloop.go +++ b/internal/agent/core/eventloop.go @@ -313,15 +313,28 @@ func (a *Agent) mediaToBlocks(payload map[string]interface{}, mediaType string, return blocks, alt } -// emitSkippedReply 给被跳过任务的**同步**调用方一个终态。 +// emitSkippedReply 给被跳过任务的调用方一个终态。 // // 为什么要单独一条路径而不是复用 emitResponse:跳过意味着“我们没有处理这条输入”, // 不应对外发 agent_output 事件(否则 WebUI 聊天记录会凭空多出一条空消息), // 但必须写 ResponseCh——否则 cli/clawhub 这类无超时的同步注入会永久挂起。 // +// ❗异步来源(qq / wechat / rss 等)**没有 ResponseCh**,于是这里以前是直接 return。 +// 后果是任务被丢弃时**完全无声**:用户什么都没收到、日志里也没痕迹, +// 他只会以为消息丢了。转投子被回收/销毁时队列里的积压正落在这个盲区里 +// (父可随时对子 reclaim/destroy,而子手上可能还握着几条 QQ 消息)。 +// 现在至少留一条带来源与通道的日志,让“这条消息为什么没回”可被追溯。 +// // 非阻塞写:ResponseCh 由同步调用方以 cap=1 创建,调用方超时离开后仍可写入。 func (a *Agent) emitSkippedReply(evt *agentIO.InputEvent, reason string) { - if evt == nil || evt.ResponseCh == nil { + if evt == nil { + return + } + if evt.ResponseCh == nil { + // 无可回执的通道:不静默。异步来源本就靠 agent 主动 output_send, + // 丢弃后没有任何东西会告诉用户,因此这条日志是唯一的线索。 + log.Printf("[agent] %s: 丢弃一条无回执通道的输入(source=%s channel=%s request=%s reason=%s)", + a.id, evt.Source, evt.OutputChannel, evt.RequestID, reason) return } ch := evt.OutputChannel diff --git a/internal/agent/core/offload.go b/internal/agent/core/offload.go index 37616a1..21e48b9 100644 --- a/internal/agent/core/offload.go +++ b/internal/agent/core/offload.go @@ -77,12 +77,16 @@ func (o OffloadOptions) normalized() OffloadOptions { // // 它必须**自己说清是系统做的**:用户看到队列里出现一条没人发过的消息时, // 唯一能解释这件事的就是这句话本身。 +// +// 同时要说清"不必重复处理":那些消息已由子 agent 回复(或已回复"忙碌中"), +// 主 agent 再处理一遍会让用户收到重复回复。 func offloadNotice(count int, residentID string) string { return fmt.Sprintf( - "[系统] %d 条积压任务已转投给驻留子 agent %s 处理(主 agent 正忙于长任务,"+ - "内核为它们拉起了独立 agent 并行执行)。它们的回复会由 %s 直接发到对应通道;"+ - "本提示仅用于说明「那几条消息不会再由你处理」,无需为它们采取任何行动。", - count, residentID, residentID) + "[系统] %d 条积压消息已在主 agent 忙期间交由临时助手 %s 先行分诊"+ + "(简单的已直接处理并回复,需要你的那些已告知用户「忙碌中,请稍候」)。"+ + "它们**不需要你再处理**了;若其中有需要你后续跟进的,请查看上述通道的会话记录。"+ + "本提示仅用于说明情况,无需回复。", + count, residentID) } // offloadCandidate 是一条可被转投的排队任务。 @@ -283,13 +287,14 @@ func (a *Agent) ensureOffloadResident(opts OffloadOptions) (string, error) { info, err := a.SpawnResident(ResidentOptions{ ID: id, // 不配 inputch:它是内核的**干活** agent,不接收任何插件的用户输入 - // (用户要求"不配输入通道")。它只由父经 sub/ 投喂任务。 + // (用户要求"不配输入通道")。它只由父经转投拿到任务。 InputChs: nil, // 全部输出通道:它要能把结果发回 qq/webui 等正确通道 // (用户要求"持有全部输出通道")。nil = 完整授权。 AllowedOutputs: nil, TempPath: a.residentTempPath(id), OffloadOwned: true, + TaskPrompt: offloadTaskPrompt(), }) if err != nil { return "", err @@ -297,6 +302,42 @@ func (a *Agent) ensureOffloadResident(opts OffloadOptions) (string, error) { return info.ID, nil } +// offloadTaskPrompt 是转投专用驻留子的**分诊职责**说明。 +// +// 为什么必须给:不给的话子完全不知道自己为什么存在(只知道自己是"小宅"), +// 拿到一条转投消息时不知道它是"用户正在等回复的请求", +// 也不知道自己只有两条路可走(直接办 / 报忙碌)。 +// +// 用户的定位(2026-09-19 明确):这不是"内核替父决定",而是**及时反馈** —— +// 主 agent 忙时不该让用户干等(实测有 13 分钟的现场)。 +// 子的职责是**分诊**(triage): +// - 简单、不需主 agent 介入的 → 直接办完并回复; +// - 需要主 agent 介入的 → 立刻回「忙碌中,请稍候」,**不要勉强做**。 +func offloadTaskPrompt() string { + return `你是主 agent 的临时助手,负责在主 agent 忙不过来时**分诊**它的积压消息。 + +背景:主 agent 正在执行一个长任务,短时间无法处理新消息。你被临时拉起, +专门承接这些积压的请求,**避免用户干等**(此前用户可能要等十几分钟)。 + +对每一条消息,你只有两条路: + +1. 【直接办】如果这件事简单、明确、不需要主 agent 的全局上下文或长期规划 + (例如:查个信息、跑个小命令、读个文件、简单问答)—— + **直接做完,并把结果发回原通道**。 + +2. 【报忙碌】如果这件事需要主 agent 介入(需要它的长期记忆、正在进行的任务上下文、 + 需要它做多步决策,或你无法确定怎么做)—— + **不要勉强尝试**。立刻回复用户:主 agent 当前忙碌中,请稍候。 + +重要约束: +- **必须把回复发到用户原本的通道**。面向 qq、wechat 等异步通道时, + 纯文本返回会被丢弃 —— 必须显式调用 output_send__{通道名},否则用户收不到, + 而你会以为已经回过了。 +- 不要向用户暴露"我是被临时拉起的助手"这类内部细节,用主 agent 的口吻回复。 +- 拿不准属于哪一类时,选【报忙碌】。宁可让用户稍后得到准确答复, + 也不要给出错误的直接回答。` +} + // residentTempPath 计算某个驻留子的 temp 图记忆路径(与既有约定一致)。 func (a *Agent) residentTempPath(id string) string { return strings.TrimRight(a.dataDir, "/") + "/residents/" + id + "/graph.db" diff --git a/internal/agent/core/offload_test.go b/internal/agent/core/offload_test.go index 37abc58..6516b1a 100644 --- a/internal/agent/core/offload_test.go +++ b/internal/agent/core/offload_test.go @@ -137,7 +137,8 @@ func TestRequeueFrontKeepsAllTasks(t *testing.T) { // 唯一能解释这件事的就是这句话本身。 func TestOffloadNoticeExplainsItself(t *testing.T) { msg := offloadNotice(3, "offload-123") - for _, want := range []string{"系统", "3 条", "offload-123", "转投"} { + // 用词按用户口径:这是**分诊**(及时反馈),不是"内核替父决定"。 + for _, want := range []string{"系统", "3 条", "offload-123", "分诊", "不需要你再处理"} { if !strings.Contains(msg, want) { t.Errorf("说明缺少 %q:%s", want, msg) } @@ -232,7 +233,7 @@ func TestOffloadMovesTasksToResidentEndToEnd(t *testing.T) { t.Fatalf("留下的应是内核说明,实际 %+v", notice.Event) } content, _ := notice.Event.Payload["content"].(string) - for _, want := range []string{"2 条", resident.ID, "转投"} { + for _, want := range []string{"2 条", resident.ID, "分诊"} { if !strings.Contains(content, want) { t.Errorf("说明缺少 %q:%s", want, content) } @@ -408,3 +409,174 @@ func TestForwardKeepsResponseCh(t *testing.T) { t.Fatal("★ 同步调用方没收到回执:转投丢了 ResponseCh(线上表现为 HTTP 挂起 504)") } } + +// ★★ 回归:转投子必须拿到**分诊职责**提示词。 +// +// 我第一版没给 TaskPrompt,于是子完全不知道自己为什么存在(只知道自己叫小宅)。 +// 用户对这个特性的定位是**及时反馈**:主 agent 忙时不能让用户干等十几分钟 +// (实测现场 785,951ms)。子的职责是分诊 —— 简单的直接办,需要主 agent 的 +// 立刻回「忙碌中,请稍候」,而不是勉强作答。 +func TestOffloadResidentGetsTriagePrompt(t *testing.T) { + p := offloadTaskPrompt() + for _, want := range []string{"分诊", "直接办", "忙碌中", "output_send"} { + if !strings.Contains(p, want) { + t.Errorf("分诊提示词缺少 %q", want) + } + } + // 拿不准时的默认动作必须是保守的那条(报忙碌),不能是"勉强作答" + if !strings.Contains(p, "选【报忙碌】") { + t.Error("必须写明拿不准时选报忙碌(避免给用户错误答复)") + } +} + +// ★★ 回归:驻留子必须**继承父的 SystemPrompt**。 +// +// 我第一版没传 SystemPrompt,子只能用 buildSystemPrompt 的一句兜底文案。 +// 而父的提示词里有「回复投递规则」:面向 qq/wechat 等**异步**通道时, +// 纯文本返回会被静默丢弃,必须显式 output_send__{通道名}。 +// 缺了它,子处理完 QQ 积压却发不出去,且自己不会意识到(实测:webui 这类 +// **同步**通道能回是因为走 ResponseCh,掩盖了这个缺陷)。 +func TestResidentInheritsParentSystemPrompt(t *testing.T) { + root, main := newRootWithoutSchedulerLoop(t) + defer main.Close() + root.systemPrompt = "父的提示词:异步通道必须显式 output_send" + + info, err := root.SpawnResident(ResidentOptions{ + ID: "inherit-test", TempPath: root.residentTempPath("inherit-test"), + }) + if err != nil { + t.Fatalf("创建驻留子失败: %v", err) + } + root.residentMu.Lock() + rc := root.residents[info.ID] + root.residentMu.Unlock() + if rc == nil || rc.agent == nil { + t.Fatal("驻留子不可用") + } + if rc.agent.systemPrompt != root.systemPrompt { + t.Fatalf("★ 驻留子未继承父的 SystemPrompt:子=%q 父=%q", + rc.agent.systemPrompt, root.systemPrompt) + } + // 真正要看的是提示词里确实带上了投递规则 + built := rc.agent.buildSystemPrompt("", "x") + if !strings.Contains(built, "output_send") { + t.Errorf("子拼出的系统提示词里没有投递规则:%s", built) + } +} + +// ★★ 残余任务必须由父**显式**决定保留还是丢弃(用户 2026-09-19 要求)。 +// +// 现场问题:回收/销毁驻留子时它手头可能还有没处理的消息。异步通道(qq) +// 没有 ResponseCh,静默丢弃时用户零反馈、日志也无痕迹 —— 消息就像没发过一样。 +// 所以内核只负责「把残余任务列清楚」,处置由父的模型决定(设计 §7)。 +func TestResidualKeepReturnsTasksToParent(t *testing.T) { + root, main := newRootWithoutSchedulerLoop(t) + defer main.Close() + + info, err := root.SpawnResident(ResidentOptions{ + ID: "res-keep", TempPath: root.residentTempPath("res-keep"), + }) + if err != nil { + t.Fatalf("创建驻留子失败: %v", err) + } + root.residentMu.Lock() + child := root.residents[info.ID].agent + root.residentMu.Unlock() + + // 给子塞两条尚未处理的残余任务(一条带同步回执、一条不带=模拟 qq) + syncCh := make(chan *agentIO.OutputEvent, 1) + child.sched.enqueue(newInputTask(&agentIO.InputEvent{ + RequestID: "r1", Source: "webui", OutputChannel: "webui", + Payload: map[string]interface{}{"content": "a"}, ResponseCh: syncCh, + })) + child.sched.enqueue(newInputTask(&agentIO.InputEvent{ + RequestID: "r2", Source: "qq", OutputChannel: "qq", + Payload: map[string]interface{}{"content": "b"}, + })) + + n, msg, err := root.ApplyResidual(info.ID, ResidualKeep) + if err != nil { + t.Fatalf("ApplyResidual(keep): %v", err) + } + if n != 2 { + t.Fatalf("应处置 2 条,实际 %d(msg=%s)", n, msg) + } + // keep = 转回父自己:两条都要出现在父的队列里,且**带 ResponseCh 的那条仍带** + if len(root.sched.queue) != 2 { + t.Fatalf("两条残余任务应转回父队列,实际 %d", len(root.sched.queue)) + } + var hasResponseCh bool + for _, task := range root.sched.queue { + if task.Event != nil && task.Event.ResponseCh != nil { + hasResponseCh = true + } + } + if !hasResponseCh { + t.Error("★ keep 丢了 ResponseCh:同步调用方会永久挂起") + } + // keep 后子队列应清空(已交出去):再取一次应为空 + if left, _ := root.TakeResidual(info.ID); len(left) != 0 { + t.Errorf("交接后子队列应清空,实际剩 %d", len(left)) + } +} + +// drop 也必须给同步调用方一个终态,否则 cli/a2a 会永久挂起。 +func TestResidualDropNotifiesSyncCaller(t *testing.T) { + root, main := newRootWithoutSchedulerLoop(t) + defer main.Close() + + info, err := root.SpawnResident(ResidentOptions{ + ID: "res-drop", TempPath: root.residentTempPath("res-drop"), + }) + if err != nil { + t.Fatalf("创建驻留子失败: %v", err) + } + root.residentMu.Lock() + child := root.residents[info.ID].agent + root.residentMu.Unlock() + + respCh := make(chan *agentIO.OutputEvent, 1) + child.sched.enqueue(newInputTask(&agentIO.InputEvent{ + RequestID: "d1", Source: "cli", OutputChannel: "cli", + Payload: map[string]interface{}{"content": "x"}, ResponseCh: respCh, + })) + + n, _, err := root.ApplyResidual(info.ID, ResidualDrop) + if err != nil { + t.Fatalf("ApplyResidual(drop): %v", err) + } + if n != 1 { + t.Fatalf("应处置 1 条,实际 %d", n) + } + select { + case out := <-respCh: + if out == nil || !out.Done { + t.Error("drop 应给同步调用方一个终态") + } + case <-time.After(3 * time.Second): + t.Fatal("★ drop 未通知同步调用方:cli/a2a 会永久挂起") + } + // drop 后父队列不应多出东西 + if len(root.sched.queue) != 0 { + t.Errorf("drop 不应把任务转回父队列,实际 %d", len(root.sched.queue)) + } +} + +// 无残余任务时应明确说"无",而不是让父以为丢了什么。 +func TestResidualEmptyIsReported(t *testing.T) { + root, main := newRootWithoutSchedulerLoop(t) + defer main.Close() + info, err := root.SpawnResident(ResidentOptions{ + ID: "res-empty", TempPath: root.residentTempPath("res-empty"), + }) + if err != nil { + t.Fatalf("创建驻留子失败: %v", err) + } + n, msg, err := root.ApplyResidual(info.ID, ResidualKeep) + if err != nil { + t.Fatalf("ApplyResidual: %v", err) + } + if n != 0 || msg != "无残余任务" { + t.Errorf("应报告无残余任务,实际 n=%d msg=%q", n, msg) + } +} diff --git a/internal/agent/core/resident.go b/internal/agent/core/resident.go index aa508a0..620c034 100644 --- a/internal/agent/core/resident.go +++ b/internal/agent/core/resident.go @@ -14,6 +14,7 @@ package core import ( "encoding/json" "fmt" + "log" "os" "path/filepath" "sort" @@ -58,15 +59,17 @@ type ResidentOptions struct { // ResidentInfo 是父对某个驻留子的可查询状态(登记表条目 + 状态面摘要)。 type ResidentInfo struct { - ID string `json:"id"` - State string `json:"state"` - InputChs []string `json:"inputchs"` - AllowedOutputs []string `json:"allowed_outputs"` - Rounds int `json:"rounds"` - ContextFull bool `json:"context_full"` - CreatedAt time.Time `json:"created_at"` - TableSize int `json:"table_size"` - Table []InputchRecord `json:"table,omitempty"` + ID string `json:"id"` + State string `json:"state"` + InputChs []string `json:"inputchs"` + AllowedOutputs []string `json:"allowed_outputs"` + Rounds int `json:"rounds"` + ContextFull bool `json:"context_full"` + CreatedAt time.Time `json:"created_at"` + TableSize int `json:"table_size"` + // OffloadOwned 标记这是内核为承接积压而拉起的**临时分诊助手**(见 offload.go)。 + OffloadOwned bool `json:"offload_owned,omitempty"` + Table []InputchRecord `json:"table,omitempty"` // Sched* 是这个驻留子**自己的**输入调度器积压摘要(排队 / 待处理中断 / // 中断栈 / 四级中断队列)。 @@ -89,7 +92,7 @@ type residentChild struct { dir string inputChs []string allowed []string - createdAt time.Time + createdAt time.Time // offloadOwned 标记这是内核为转投拉起的子(见 offload.go)。 offloadOwned bool @@ -165,7 +168,19 @@ func (a *Agent) SpawnResident(opts ResidentOptions) (ResidentInfo, error) { parentID := string(a.id) child := New(AgentConfig{ - ID: types.AgentID(opts.ID), + ID: types.AgentID(opts.ID), + // SystemPrompt 必须**继承父的**: + // + // 拿不到它时子只能用 buildSystemPrompt 里的兜底句(一句人格描述)。 + // 而父的提示词里写着**回复投递规则** ——「面向 qq、wechat、a2a、acp 等异步 + // 通道时,必须显式调用 output_send__{通道名};只返回纯文本会被直接丢弃, + // 用户永远收不到,而你会误以为已经回复过了」,以及事实性约束、命令与文件 + // 操作策略等。缺了这些,子于异步通道(如 QQ 积压)处理完却发不出去, + // 而且它自己不会意识到(提示词没告诉它)。 + // + // 实测(2026-09-19):webui 这类**同步**通道能回,是因为走 ResponseCh; + // 而 QQ 是异步通道、必须子主动 output_send —— 这正是本字段必须接上的理由。 + SystemPrompt: a.systemPrompt, Provider: a.provider, ProviderManager: a.providerManager, IO: childIO, @@ -263,6 +278,88 @@ func (a *Agent) DestroyResident(id string) error { return nil } +// ResidualPolicy 是父对**残余任务**的显式决定。 +// +// 为什么需要它(用户 2026-09-19 明确要求):回收/销毁驻留子时, +// 它手头可能还有**尚未处理的消息**。这些消息的处置不能由内核悄悄决定: +// - 直接丢弃 → 用户消息无声消失(异步通道更是零反馈:qq 无 ResponseCh, +// 用户不知道发生了什么,系统里也没任何痕迹); +// - 无条件转回父 → 父本来就很忙,把一堆活重新塞回去可能反而加剧积压。 +// +// 因此与其它控制面动作一致:**原语在内核,决定在父的模型**(设计 §7)。 +// 内核负责把残余任务**列清楚**(来源、通道、内容),父选 keep(转回自己)/ drop。 +// +// 无论选哪个,内核都会**逐条留日志**:丢弃必须可追溯。 +// 默认(不传 policy)取 keep:宁可多做一件,不可默默丢一条。 +type ResidualPolicy string + +const ( + // ResidualKeep 把残余任务转回父自己的队列,父稍后处理。 + ResidualKeep ResidualPolicy = "keep" + // ResidualDrop 明确丢弃残余任务(父已看过清单并确认)。 + ResidualDrop ResidualPolicy = "drop" +) + +// TakeResidual 取走一个驻留子手头**尚未处理**的残余任务(不处置)。 +// +// 调它**会**从子的队列里移除这些任务,所以父必须先看返回值再决定: +// 一旦取走,不再调 ApplyResidual 就等于把它们丢了。控制面工具走 ApplyResidual +// (取+处置一步完成)以避免这个误用。 +func (a *Agent) TakeResidual(id string) ([]*agentIO.InputEvent, error) { + a.residentMu.Lock() + rc := a.residents[id] + a.residentMu.Unlock() + if rc == nil || rc.agent == nil { + return nil, fmt.Errorf("驻留子 %s 不存在", id) + } + return rc.agent.sched.takeAllPendingEvents(), nil +} + +// ApplyResidual 按父给出的策略处置一个驻留子的残余任务。 +// +// keep:逐条转回父自己的队列(保留原事件,含 ResponseCh 与来源/通道); +// drop:逐条记录日志后丢弃(不可追溯的丢弃是不允许的)。 +// +// 无论哪种,带 ResponseCh 的都要给终态,否则 cli/a2a 这类无超时同步调用方 +// 会永久挂起(设计 §7 I5)。 +func (a *Agent) ApplyResidual(id string, policy ResidualPolicy) (int, string, error) { + events, err := a.TakeResidual(id) + if err != nil { + return 0, "", err + } + if len(events) == 0 { + return 0, "无残余任务", nil + } + if policy == ResidualDrop { + for _, evt := range events { + // 丢弃必须留痕:异步通道(qq/wechat)没有 ResponseCh, + // 不记日志的话“这条消息为什么没人回”永远查不出来。 + log.Printf("[resident] %s 残余任务按父的决定丢弃:source=%s channel=%s request=%s content=%s", + id, evt.Source, evt.OutputChannel, evt.RequestID, truncateStr(inputTextOf(evt), 80)) + a.emitSkippedReply(evt, "residual_dropped_by_parent") + } + return len(events), fmt.Sprintf("已按父的决定丢弃 %d 条残余任务(逐条已记日志)", len(events)), nil + } + // 默认 keep:转回父自己 + for _, evt := range events { + if evt == nil { + continue + } + a.sched.enqueue(newInputTask(evt)) + } + a.sched.signalWake() + return len(events), fmt.Sprintf("已把 %d 条残余任务转回主 agent 队列", len(events)), nil +} + +// inputTextOf 取一条输入事件的正文(供日志描述残余任务用)。 +func inputTextOf(evt *agentIO.InputEvent) string { + if evt == nil || evt.Payload == nil { + return "" + } + s, _ := evt.Payload["content"].(string) + return s +} + // teardownResident 停内核、放通道、丢 temp(销毁与回收共用)。 func (a *Agent) teardownResident(rc *residentChild) { rc.mu.Lock() diff --git a/internal/agent/core/resident_tools.go b/internal/agent/core/resident_tools.go index 8836947..bc4125a 100644 --- a/internal/agent/core/resident_tools.go +++ b/internal/agent/core/resident_tools.go @@ -123,17 +123,31 @@ func (a *Agent) executeResidentAgents(tc agentAPI.ToolCall) string { return fmt.Sprintf("已压缩子 agent 上下文(丢弃 %d 条旧事件,并发清理其 inputch 处理表);子继续存在", n) case "reclaim": - info, err := a.ReclaimResident(strArg(tc, "id"), reclaimKeepAll) + id := strArg(tc, "id") + // 残余任务先由父**显式**决定保留还是丢弃(用户 2026-09-19 要求)。 + // 不传 residual 时默认 keep:宁可多做一件,不可默默丢一条。 + residual, rmsg, rerr := a.ApplyResidual(id, residualPolicyOf(tc)) + if rerr != nil { + return fmt.Sprintf("处置残余任务失败: %v", rerr) + } + info, err := a.ReclaimResident(id, reclaimKeepAll) if err != nil { return fmt.Sprintf("回收失败: %v", err) } - return "已回收(temp 中选中的记录已合入主记忆,该驻留子已取消): " + MarshalResidentInfo(info) + return fmt.Sprintf("残余任务(%d 条):%s\n已回收(temp 中选中的记录已合入主记忆,该驻留子已取消): %s", + residual, rmsg, MarshalResidentInfo(info)) case "destroy": - if err := a.DestroyResident(strArg(tc, "id")); err != nil { + id := strArg(tc, "id") + // 同上:销毁前先把残余任务交出去,否则它们会随子一起无声消失。 + residual, rmsg, rerr := a.ApplyResidual(id, residualPolicyOf(tc)) + if rerr != nil { + return fmt.Sprintf("处置残余任务失败: %v", rerr) + } + if err := a.DestroyResident(id); err != nil { return fmt.Sprintf("销毁失败: %v", err) } - return "已销毁并移除该驻留子" + return fmt.Sprintf("残余任务(%d 条):%s\n已销毁并移除该驻留子", residual, rmsg) default: return fmt.Sprintf("未知 action=%q;可用:list | create | send | inspect | compress | reclaim | destroy", action) @@ -144,6 +158,19 @@ func (a *Agent) executeResidentAgents(tc agentAPI.ToolCall) string { // ("哪些纳入"由父的模型决定——这里给的是"全要"这一档)。 func reclaimKeepAll(_ []InputchRecord, triples []memory.Triple) []memory.Triple { return triples } +// residualPolicyOf 从工具参数读残余任务的处置策略(keep/drop)。 +// +// 默认 keep:父没明确说丢时,一律转回自己而不是丢弃。 +// "宁可多做一件,不可默默丢一条"——吞掉一条输入比多处理一条更糟。 +func residualPolicyOf(tc agentAPI.ToolCall) ResidualPolicy { + switch strings.ToLower(strArg(tc, "residual")) { + case "drop", "discard": + return ResidualDrop + default: + return ResidualKeep + } +} + func intArg(tc agentAPI.ToolCall, key string) int { switch v := tc.Arguments[key].(type) { case float64: diff --git a/internal/agent/core/scheduler.go b/internal/agent/core/scheduler.go index 324ab35..bd418f4 100644 --- a/internal/agent/core/scheduler.go +++ b/internal/agent/core/scheduler.go @@ -1011,6 +1011,41 @@ func (s *scheduler) pendingEvents() []*agentIO.InputEvent { return out } +// takeAllPendingEvents 取走**全部尚未执行**的输入事件(含无回执通道的异步输入), +// 并从队列中移除它们。返回的事件不再会被本调度器执行。 +// +// 与 pendingEvents 的区别(两者用途完全不同,不要混用): +// - pendingEvents 只挑**带 ResponseCh** 的,用途是给同步调用方补终态; +// - 本函数**不筛通道**,因为回收/销毁驻留子时要向父交代的是"手头还有哪些活", +// 而 QQ/微信这类异步消息本来就没有 ResponseCh —— 它们恰恰是最容易被无声丢掉的。 +// +// 只取 TaskQueued/KindInput 与中断队列里的输入事件;self 任务是内核内部记账 +// (记忆整理),换 agent 没有意义,原地丢弃即可(不计入返回)。 +func (s *scheduler) takeAllPendingEvents() []*agentIO.InputEvent { + s.mu.Lock() + defer s.mu.Unlock() + var out []*agentIO.InputEvent + for _, t := range s.queue { + if t != nil && t.Kind == TaskKindInput && t.Event != nil { + out = append(out, t.Event) + } + } + s.queue = nil + for lv := LevelBackground; lv <= LevelCritical; lv++ { + for _, t := range s.interruptQueues[lv] { + if t != nil && t.Event != nil { + out = append(out, t.Event) + } + } + s.interruptQueues[lv] = nil + } + if s.immediate != nil && s.immediate.Event != nil { + out = append(out, s.immediate.Event) + } + s.immediate = nil + return out +} + // done 标记任务执行结束。 func (s *scheduler) done(t *Task) { s.mu.Lock() diff --git a/internal/agent/core/status.go b/internal/agent/core/status.go index a63a228..912f0b3 100644 --- a/internal/agent/core/status.go +++ b/internal/agent/core/status.go @@ -331,6 +331,7 @@ func residentStatuses(list []ResidentInfo) []ResidentStatus { InputChs: r.InputChs, AllowedOutputs: r.AllowedOutputs, InputChTable: r.TableSize, + OffloadOwned: r.OffloadOwned, // 子自己的调度器积压:per-agent 负载图靠这四项,缺了就只能画根。 ReadyQueueDepth: r.SchedReady, PendingInterrupts: r.SchedPending, diff --git a/internal/agent/core/tooldefs.go b/internal/agent/core/tooldefs.go index 2cd1310..bf5b4a4 100644 --- a/internal/agent/core/tooldefs.go +++ b/internal/agent/core/tooldefs.go @@ -508,7 +508,9 @@ func (a *Agent) buildToolDefs() []interface{} { tools = append(tools, toolDef("resident_agents", "管理驻留子 agent(长期派驻的下属):list 列出 / create 创建(划入 inputch + "+ "授权输出通道 + 注入任务提示词)/ send 发送消息(对子而言是 L4 中断,取消其当前状态并插入新消息)"+ "/ inspect 查看其 inputch 处理表(不打断它)/ compress 压缩其上下文(保留语义,子继续存在)"+ - "/ reclaim 回收(父选哪些纳入主记忆,然后取消该子)/ destroy 立刻销毁并移除。", map[string]interface{}{ + "/ reclaim 回收(父选哪些纳入主记忆,然后取消该子)/ destroy 立刻销毁并移除。"+ + "reclaim/destroy 时它手头**尚未处理的消息**(残余任务)由你用 residual 决定:"+ + "keep=转回你自己的队列(默认)/ drop=明确丢弃,两者都会逐条记日志。", map[string]interface{}{ "action": map[string]interface{}{ "type": "string", "enum": []string{"list", "create", "send", "inspect", "compress", "reclaim", "destroy"}, @@ -520,6 +522,13 @@ func (a *Agent) buildToolDefs() []interface{} { "capacity": map[string]interface{}{"type": "number", "description": "create:划入 inputch 的队列容量"}, "temp_path": map[string]interface{}{"type": "string", "description": "create:temp 图记忆路径(留空则用 data_dir/residents//graph.db)"}, "text": map[string]interface{}{"type": "string", "description": "send:要发给子 agent 的消息"}, + "residual": map[string]interface{}{ + "type": "string", + "description": "reclaim/destroy:该子手头未处理消息的处置。" + + "keep=转回主 agent 队列(默认,宁可多做一件不可默默丢一条)/ drop=丢弃(已确认不要)。" + + "无论哪种都会逐条记日志;异步通道(如 qq)被丢时用户收不到任何回复,请慎重选 drop。", + "enum": []string{"keep", "drop"}, + }, }, "action")) } diff --git a/internal/sdk/status.go b/internal/sdk/status.go index 79c628d..c5f58ba 100644 --- a/internal/sdk/status.go +++ b/internal/sdk/status.go @@ -122,6 +122,13 @@ type ResidentStatus struct { InputChTable int `json:"input_ch_table"` CreatedAt string `json:"created_at,omitempty"` + // OffloadOwned 标记这是**内核为承接积压而拉起**的临时助手(见 core/offload.go)。 + // + // 为什么要暴露给状态面/工具面:父的模型需要分清"我建的子"与 + // "内核临时拉的分诊助手" —— 前者该按需回收,后者在空闲、无积压时 + // 可以安全回收,而**正忙时不该回收**(会让用户的消息再次无声丢失)。 + OffloadOwned bool `json:"offload_owned,omitempty"` + // 以下四项是该驻留子**自己的**调度器积压摘要,用于 per-agent 负载环形图。 // // 根 agent 的积压看 KernelStatus.Scheduler;每个驻留子是独立 agent、