diff --git a/docs/zh/resident-subagent-design.md b/docs/zh/resident-subagent-design.md index 7ff552c..4c2ad23 100644 --- a/docs/zh/resident-subagent-design.md +++ b/docs/zh/resident-subagent-design.md @@ -125,6 +125,28 @@ - **输出通道 → 目标 agent 的 inputch** 的解析; - **授权过滤**:`output_list_channels` 只列该 agent 被授权的通道。 +### 4.6 inputch 总览:**单工具多视图** [已定] + +父 agent 必须能看清两件事:**有哪些 inputch 已注册(谁注册的)**、**它们是怎么划分的**。 +按用户要求构筑为**单工具多视图**(一个工具 + 一个 `view` 参数),而不是一堆小工具 —— +视图切换比工具增殖更好用,也更省提示词预算。 + +工具:**`input_channels`** + +| view | 内容 | +|---|---| +| `all`(默认) | 全部已注册 inputch:名字 / **归属插件** / 归属 agent / 容量 / 记忆策略标记 | +| `mine` | 划给**本 agent** 的 | +| `unassigned` | **尚未划出**的(可按需分配) | +| `by_agent` | **划分情况总览**:按归属 agent 分组列出各自拥有哪些 inputch | +| `detail`(需 `name`) | 单个 inputch 的全字段(注册插件 / 归属 / 容量 / 默认回程 / 记忆策略) | + +- 未知 `view` **必须报错并列出可用值**(拼错不得被静默当成默认视图)。 +- 登记表是**可共享对象**(`*ChannelRegistry`):根 agent 与它的驻留子共用同一份, + 这样"划入/授权"才有意义(默认每个 agent 自带一份,向后兼容)。 +- **插件重载不得抹掉划分**:重复登记只更新「归属插件 + 策略」, + 保留已有的 Owner / Capacity / Output。 + --- ## 5. 记忆模型:两级空间 [已定] @@ -375,6 +397,8 @@ | S16 | 父退出清理 | 父 `Stop()`(多个子、有子在工作中) | 全部子被销毁;登记表空;**无孤儿 goroutine / 无悬空通道 / 无残留 temp** | | S17 | L4 独占 | 子内部(子自己的输入/工具/定时器)试图产生 L4 | 被夹到 **≤L3**;只有父的消息是 L4 | | S21 | 一个插件多个 inputch 可分别路由 | 同一插件注册 2 个 inputch,分别划给父与子后各投一条输入 | 各自只到被划给的 agent,互不串台 | +| S22 | inputch 总览(单工具多视图) | 一个插件注册 2 个 inputch、另一插件 1 个;把其中若干划给本 agent | `view=all` 列出全部(带归属插件);`mine`/`unassigned` 各自正确;`by_agent` 给出划分总览;`detail` 给出单条全字段;未知 view 报错并列出可用值 | +| S23 | 插件重载不抹划分 | 先划分 inputch,再重复登记(模拟插件重载) | Owner/Capacity 保留,仅策略被更新 | | S19 | 父消息必能打断子 | 子在长任务中(LLM 流式段)时父发送消息 | 子按 L4 被打断;若子在不可抢占临界区,则在安全点生效 | | S20 | 子的内核级事件上报 | 子 panic / 子 contextfull | 在**父的阶梯上以 L4(带子标识)**出现;子侧不自己产生 L4 | | S18 | 内核不持有"当前通道" | 抢占/中断后被打断任务恢复并发响应 | 提示词与事件标签都**只来自输入事件**(不再有被覆盖的字段) | @@ -388,6 +412,7 @@ | **R1** | 驻留子的内核形态 | **[已定]独立轻量内核**(自己的调度器/中断栈/上下文/记忆作用域) | | **R2** | 插件与工具 | **[已定]父授权,默认完整授权**(可收窄) | | **R3** | 记忆模型 | **[已定]两级空间**:读 temp∪main,写 temp | +| **R13** | inputch 总览的形态 | **[已定]单工具多视图**(`input_channels` + `view`) | | **R11** | 父消息的级别 | **[已定]对子 = L4**(子的阶梯上唯一 L4 来源;子内部一律 ≤L3) | | **R12** | 子的内核级事件(panic / contextfull)上报级别 | **[默认/推论]在父的阶梯上以 L4(带子标识)上报** | | **R4** | 父→子消息落地 | **[默认]挂起/恢复**(现场不丢);一处开关可改"直接丢弃" | @@ -405,7 +430,8 @@ | 里程碑 | 内容 | 验收 | |---|---|---| | **N0** | **无状态化**:删 `Agent.currentOutputChannel`、删提示词里的通道预设 | S18;既有全部测试通过(这是纯收敛,不含新能力) | -| **N1** | **通道层**:inputch 一等化(归属 / 容量 / 授权 / 目标解析)+ 输出通道授权过滤 | S1–S3 | +| **N1a** | **通道登记层**:inputch 一等化(归属插件 / 归属 agent / 容量 / 共享登记表)+ **单工具多视图总览** | S21–S23 | +| **N1b** | 输出通道授权过滤 + 目标解析(outputch → 目标 agent 的 inputch) | S1–S3 | | **N2** | **轻量内核 + 记忆作用域**:记忆子系统加 space 维度;`AgentConfig` 加作用域参数(写 temp、读 temp∪main) | S13–S15 | | **N3** | **驻留子生命周期**:创建 / 销毁 / 登记表 / 父退出清理 | S12、S16 | | **N4** | **跨 agent 投递**:子→父 L3、父→子 L4(+ `isKernelLevelSource` 分层) | S4、S6、S17 | diff --git a/internal/agent/core/inputch.go b/internal/agent/core/inputch.go new file mode 100644 index 0000000..dcfbccf --- /dev/null +++ b/internal/agent/core/inputch.go @@ -0,0 +1,167 @@ +package core + +// inputch 总览:**单工具多视图**。 +// +// inputch 是**最基本的输入路由单位**(由插件注册,一个插件可注册多个)。 +// 父 agent 需要能看清两件事: +// 1. 有哪些 inputch 已注册(谁注册的); +// 2. 它们是怎么划分的(各自划给了哪个 agent、容量多少)。 +// +// 按用户要求做成**单工具多视图**(一个 `input_channels` 工具 + `view` 参数), +// 而不是一堆小工具 —— 视图切换比工具增殖更好用,也更省提示词预算。 + +import ( + "fmt" + "sort" + "strings" + + agentAPI "gitcode.com/JianFeeeee/HomeAgent/internal/agent/api" + agentIO "gitcode.com/JianFeeeee/HomeAgent/internal/agent/io" +) + +func (a *Agent) executeInputChannels(tc agentAPI.ToolCall) string { + view, _ := tc.Arguments["view"].(string) + view = strings.TrimSpace(view) + if view == "" { + view = "all" + } + name, _ := tc.Arguments["name"].(string) + + all := a.io.InputChannels() + if len(all) == 0 { + return "没有任何已注册的 inputch。" + } + + switch view { + case "all": + return a.renderInputChannels(all, "全部已注册 inputch") + case "mine": + return a.renderInputChannels(a.channelRegistry().ListByOwner(string(a.id)), + "划给本 agent("+string(a.id)+")的 inputch") + case "unassigned": + return a.renderInputChannels(a.channelRegistry().ListByOwner(""), + "尚未划出的 inputch(可按需分配)") + case "by_agent": + return a.renderInputChannelsByAgent(all) + case "detail": + if name == "" { + return "view=detail 需要 name 参数(inputch 名)" + } + ch, ok := a.io.LookupInputChannel(name) + if !ok { + return fmt.Sprintf("inputch %q 未注册", name) + } + return renderInputChannelDetail(ch) + default: + return fmt.Sprintf("未知 view=%q;可用:all | mine | unassigned | by_agent | detail", view) + } +} + +func (a *Agent) channelRegistry() *agentIO.ChannelRegistry { return a.io.ChannelRegistry() } + +// renderInputChannels 渲染一组 inputch 的一行式概览。 +func (a *Agent) renderInputChannels(list []agentIO.InputChannel, title string) string { + if len(list) == 0 { + return title + ":无" + } + var b strings.Builder + fmt.Fprintf(&b, "%s(%d 个):", title, len(list)) + for _, ch := range list { + fmt.Fprintf(&b, "\n - %s%s%s", ch.Name, pluginSuffix(ch), policySuffix(ch)) + fmt.Fprintf(&b, "\n 归属: %s", ownerLabel(ch.Owner)) + if ch.Capacity > 0 { + fmt.Fprintf(&b, " | 容量: %d", ch.Capacity) + } + if ch.Output != "" { + fmt.Fprintf(&b, " | 默认回程: %s", ch.Output) + } + } + return b.String() +} + +// renderInputChannelsByAgent 按归属分组("划分情况"总览)。 +func (a *Agent) renderInputChannelsByAgent(all []agentIO.InputChannel) string { + byOwner := map[string][]agentIO.InputChannel{} + for _, ch := range all { + byOwner[ch.Owner] = append(byOwner[ch.Owner], ch) + } + owners := make([]string, 0, len(byOwner)) + for o := range byOwner { + owners = append(owners, o) + } + sort.Strings(owners) + + var b strings.Builder + fmt.Fprintf(&b, "inputch 划分情况(共 %d 个):", len(all)) + for _, o := range owners { + names := make([]string, 0, len(byOwner[o])) + for _, ch := range byOwner[o] { + names = append(names, ch.Name) + } + sort.Strings(names) + fmt.Fprintf(&b, "\n - %s: %s", ownerLabel(o), strings.Join(names, ", ")) + } + return b.String() +} + +// renderInputChannelDetail 渲染单个 inputch 的全部字段。 +func renderInputChannelDetail(ch agentIO.InputChannel) string { + var b strings.Builder + fmt.Fprintf(&b, "inputch: %s\n", ch.Name) + fmt.Fprintf(&b, " 注册插件: %s\n", orDash(ch.Plugin)) + fmt.Fprintf(&b, " 归属 agent: %s\n", ownerLabel(ch.Owner)) + if ch.Capacity > 0 { + fmt.Fprintf(&b, " 容量: %d\n", ch.Capacity) + } else { + fmt.Fprintf(&b, " 容量: 内核默认\n") + } + fmt.Fprintf(&b, " 默认回程输出通道: %s\n", orDash(ch.Output)) + fmt.Fprintf(&b, " 记忆策略: %s\n", policyLabel(ch)) + return b.String() +} + +func pluginSuffix(ch agentIO.InputChannel) string { + if ch.Plugin == "" { + return "" + } + return "(插件 " + ch.Plugin + ")" +} + +// policySuffix 用短标记提示策略(详见 view=detail)。 +func policySuffix(ch agentIO.InputChannel) string { + var m []string + if ch.Def.NoMemory { + m = append(m, "无记忆") + } + if ch.Def.Cleaner != nil { + m = append(m, "清洗") + } + if ch.Def.ContextPolicy != "" && ch.Def.ContextPolicy != "none" { + m = append(m, "裁剪:"+ch.Def.ContextPolicy) + } + if len(m) == 0 { + return "" + } + return " [" + strings.Join(m, "/") + "]" +} + +func policyLabel(ch agentIO.InputChannel) string { + if s := policySuffix(ch); s != "" { + return strings.Trim(s, " []") + } + return "默认(记入记忆、不裁剪)" +} + +func ownerLabel(owner string) string { + if owner == "" { + return "未分配(根 agent/内核默认)" + } + return owner +} + +func orDash(s string) string { + if s == "" { + return "-" + } + return s +} diff --git a/internal/agent/core/inputch_test.go b/internal/agent/core/inputch_test.go new file mode 100644 index 0000000..2835507 --- /dev/null +++ b/internal/agent/core/inputch_test.go @@ -0,0 +1,104 @@ +package core + +// inputch 总览工具(单工具多视图)的验收。 +// +// 需求:父 agent 能看到**所有已注册的 inputch**以及**它们的划分情况**。 +// 构筑方式:单工具(input_channels)+ 多视图(view=all|mine|unassigned|by_agent|detail)。 + +import ( + "strings" + "testing" + + agentAPI "gitcode.com/JianFeeeee/HomeAgent/internal/agent/api" + agentIO "gitcode.com/JianFeeeee/HomeAgent/internal/agent/io" +) + +func inputChToolCall(args map[string]interface{}) agentAPI.ToolCall { + return agentAPI.ToolCall{ID: "ic1", Name: "input_channels", Arguments: args} +} + +func TestInputChannelsTool_SingleToolMultipleViews(t *testing.T) { + a := newPreemptAgent(t, newPreemptProvider()) + me := string(a.id) + + // 同一个插件注册多个 inputch(最基本的输入路由单位);另一个插件再注册一个。 + if err := a.io.RegisterInputChannelFrom("qq", "qq", agentIO.ChannelDef{NoMemory: true}); err != nil { + t.Fatal(err) + } + if err := a.io.RegisterInputChannelFrom("qq", "qq/device-2", agentIO.ChannelDef{}); err != nil { + t.Fatal(err) + } + if err := a.io.RegisterInputChannelFrom("sub", "sub/in", agentIO.ChannelDef{}); err != nil { + t.Fatal(err) + } + // 划分:qq 与 sub/in 归本 agent,qq/device-2 留未分配。 + if err := a.io.AssignInputChannel("qq", me, 64); err != nil { + t.Fatal(err) + } + if err := a.io.AssignInputChannel("sub/in", me, 0); err != nil { + t.Fatal(err) + } + + // view=all(默认):全部已注册,且带归属插件。 + all := a.executeInputChannels(inputChToolCall(nil)) + for _, want := range []string{"qq", "qq/device-2", "sub/in", "插件 qq", "插件 sub", "无记忆"} { + if !strings.Contains(all, want) { + t.Fatalf("view=all 缺少 %q:\n%s", want, all) + } + } + + // view=mine:只有划给本 agent 的。 + mine := a.executeInputChannels(inputChToolCall(map[string]interface{}{"view": "mine"})) + if !strings.Contains(mine, "qq") || !strings.Contains(mine, "sub/in") { + t.Fatalf("view=mine 应含 qq 与 sub/in:\n%s", mine) + } + if strings.Contains(mine, "qq/device-2") { + t.Fatalf("view=mine 不应含未分配的 qq/device-2:\n%s", mine) + } + + // view=unassigned:只有尚未划出的。 + un := a.executeInputChannels(inputChToolCall(map[string]interface{}{"view": "unassigned"})) + if !strings.Contains(un, "qq/device-2") { + t.Fatalf("view=unassigned 应含 qq/device-2:\n%s", un) + } + if strings.Contains(un, "sub/in") { + t.Fatalf("view=unassigned 不应含已划分的 sub/in:\n%s", un) + } + + // view=by_agent:划分情况总览(按归属分组)。 + byAgent := a.executeInputChannels(inputChToolCall(map[string]interface{}{"view": "by_agent"})) + if !strings.Contains(byAgent, me+":") { + t.Fatalf("view=by_agent 应列出归属 %q:\n%s", me, byAgent) + } + if !strings.Contains(byAgent, "未分配") { + t.Fatalf("view=by_agent 应列出未分配一组:\n%s", byAgent) + } + + // view=detail:单个 inputch 的全字段(容量 / 策略 / 回程 / 归属)。 + detail := a.executeInputChannels(inputChToolCall(map[string]interface{}{"view": "detail", "name": "qq"})) + for _, want := range []string{"inputch: qq", "注册插件: qq", "容量: 64", "无记忆", me} { + if !strings.Contains(detail, want) { + t.Fatalf("view=detail 缺少 %q:\n%s", want, detail) + } + } + if miss := a.executeInputChannels(inputChToolCall(map[string]interface{}{"view": "detail", "name": "nope"})); !strings.Contains(miss, "未注册") { + t.Fatalf("detail 查未注册的 inputch 应明确报错:%s", miss) + } + if noName := a.executeInputChannels(inputChToolCall(map[string]interface{}{"view": "detail"})); !strings.Contains(noName, "需要 name") { + t.Fatalf("detail 缺 name 应提示:%s", noName) + } + + // 未知视图必须报错并列出可用值(不要把拼错静默当成默认视图)。 + bad := a.executeInputChannels(inputChToolCall(map[string]interface{}{"view": "whatever"})) + if !strings.Contains(bad, "未知 view") || !strings.Contains(bad, "by_agent") { + t.Fatalf("未知 view 应报错并列出可用值:%s", bad) + } +} + +// 一个 inputch 都没有时应给出明确说明,而不是空串。 +func TestInputChannelsTool_NoChannels(t *testing.T) { + a := newPreemptAgent(t, newPreemptProvider()) + if out := a.executeInputChannels(inputChToolCall(nil)); !strings.Contains(out, "没有任何已注册") { + t.Fatalf("空登记表应明确说明:%q", out) + } +} diff --git a/internal/agent/core/toolcall.go b/internal/agent/core/toolcall.go index 1e06d2e..cca2ff9 100644 --- a/internal/agent/core/toolcall.go +++ b/internal/agent/core/toolcall.go @@ -61,6 +61,8 @@ func (a *Agent) executeToolCallInner(tc agentAPI.ToolCall, channel string) strin return a.executeOutputSendTool(tc) case tc.Name == "output_list_channels": return a.executeOutputListChannels() + case tc.Name == "input_channels": + return a.executeInputChannels(tc) case tc.Name == "plgreload": return a.executePluginReload() case tc.Name == "get_plugin_tools": diff --git a/internal/agent/core/tooldefs.go b/internal/agent/core/tooldefs.go index 2a0ecda..1aeb8bc 100644 --- a/internal/agent/core/tooldefs.go +++ b/internal/agent/core/tooldefs.go @@ -637,6 +637,30 @@ func (a *Agent) buildToolDefs() []interface{} { }, }) + tools = append(tools, map[string]interface{}{ + "type": "function", + "function": map[string]interface{}{ + "name": "input_channels", + "description": "查看 inputch(最基本的输入路由单位):哪些已注册、谁注册的、" + + "各自划给了哪个 agent、容量与记忆策略。单工具多视图。", + "parameters": map[string]interface{}{ + "type": "object", + "properties": map[string]interface{}{ + "view": map[string]interface{}{ + "type": "string", + "description": "all=全部已注册(默认)| mine=划给本 agent 的 | " + + "unassigned=尚未划出的 | by_agent=按归属分组的划分总览 | detail=单个详情", + "enum": []string{"all", "mine", "unassigned", "by_agent", "detail"}, + }, + "name": map[string]interface{}{ + "type": "string", + "description": "view=detail 时必填:inputch 名", + }, + }, + }, + }, + }) + if a.pendingMedia != nil { tools = append(tools, map[string]interface{}{ "type": "function", diff --git a/internal/agent/io/channel.go b/internal/agent/io/channel.go index bb96225..16e89a0 100644 --- a/internal/agent/io/channel.go +++ b/internal/agent/io/channel.go @@ -96,13 +96,13 @@ type OutputEvent struct { } type IOManager struct { - mu sync.RWMutex - devices map[string]Device - inputCh chan *InputEvent - interruptCh chan *InputEvent - outputCh chan *OutputEvent - nextReqID int64 - inputChannels map[string]ChannelDef + mu sync.RWMutex + devices map[string]Device + inputCh chan *InputEvent + interruptCh chan *InputEvent + outputCh chan *OutputEvent + nextReqID int64 + channelReg *ChannelRegistry // toolBlocks:插件工具注入多模态内容块,process.go 在下一条 tool message 时消费。 // 用 interface{}[] 避免 import api.ContentBlock 导致的循环依赖。 @@ -112,11 +112,11 @@ type IOManager struct { func NewIOManager() *IOManager { return &IOManager{ - devices: make(map[string]Device), - inputCh: make(chan *InputEvent, 256), - interruptCh: make(chan *InputEvent, 64), - outputCh: make(chan *OutputEvent, 256), - inputChannels: make(map[string]ChannelDef), + devices: make(map[string]Device), + inputCh: make(chan *InputEvent, 256), + interruptCh: make(chan *InputEvent, 64), + outputCh: make(chan *OutputEvent, 256), + channelReg: NewChannelRegistry(), } } @@ -444,26 +444,57 @@ func (m *IOManager) EmitTextTo(target, outputChannel, text string) { func (m *IOManager) InputChan() <-chan *InputEvent { return m.inputCh } func (m *IOManager) OutputChan() <-chan *OutputEvent { return m.outputCh } -// RegisterInputChannel 注册输入通道的记忆行为 +// RegisterInputChannel 注册一个 inputch(不带插件归属,兼容旧调用)。 +// +// inputch 是**最基本的输入路由单位**;一个插件可以注册多个。 +// 新代码请用 RegisterInputChannelFrom 以便登记归属插件(可追溯)。 func (m *IOManager) RegisterInputChannel(name string, def ChannelDef) { - m.mu.Lock() - defer m.mu.Unlock() - m.inputChannels[name] = def + _ = m.RegisterInputChannelFrom("", name, def) } -// UnregisterInputChannel 注销输入通道 +// RegisterInputChannelFrom 注册一个 inputch 并登记归属插件。 +func (m *IOManager) RegisterInputChannelFrom(plugin, name string, def ChannelDef) error { + return m.channelReg.Register(InputChannel{Name: name, Plugin: plugin, Def: def}) +} + +// UnregisterInputChannel 注销一个 inputch。 func (m *IOManager) UnregisterInputChannel(name string) { - m.mu.Lock() - defer m.mu.Unlock() - delete(m.inputChannels, name) + m.channelReg.Unregister(name) } -// GetInputChannelDef 查询输入通道的记忆行为定义 +// AssignInputChannel 把一个 inputch 划给某个 agent(见 ChannelRegistry.Assign)。 +func (m *IOManager) AssignInputChannel(name, agentID string, capacity int) error { + return m.channelReg.Assign(name, agentID, capacity) +} + +// SetChannelRegistry 注入一份**共享的**登记表(根 agent 与驻留子共用同一份)。 +func (m *IOManager) SetChannelRegistry(r *ChannelRegistry) { + if r == nil { + return + } + m.mu.Lock() + defer m.mu.Unlock() + m.channelReg = r +} + +// ChannelRegistry 返回底层登记表(只读用途;可直接读 Views)。 +func (m *IOManager) ChannelRegistry() *ChannelRegistry { return m.channelReg } + +// InputChannels 返回全部已注册 inputch(按名字排序)。 +func (m *IOManager) InputChannels() []InputChannel { return m.channelReg.List() } + +// LookupInputChannel 查询单个 inputch 的完整登记记录。 +func (m *IOManager) LookupInputChannel(name string) (InputChannel, bool) { + return m.channelReg.Lookup(name) +} + +// GetInputChannelDef 查询 inputch 的记忆行为定义。 func (m *IOManager) GetInputChannelDef(name string) (ChannelDef, bool) { - m.mu.RLock() - defer m.mu.RUnlock() - def, ok := m.inputChannels[name] - return def, ok + ch, ok := m.channelReg.Lookup(name) + if !ok { + return ChannelDef{}, false + } + return ch.Def, true } func (m *IOManager) GetAllTools() []ToolDef { diff --git a/internal/agent/io/inputch.go b/internal/agent/io/inputch.go new file mode 100644 index 0000000..69bfbb5 --- /dev/null +++ b/internal/agent/io/inputch.go @@ -0,0 +1,152 @@ +package io + +// inputch 是**最基本的输入路由单位**(设计见 docs/zh/resident-subagent-design.md §4)。 +// +// 一个插件可以注册多个 inputch;每个 inputch 是彼此独立的路由单位: +// 可以被划给不同的 agent、可以分别限额。中断输入与排队输入**两类都从 inputch 进出**, +// 而"中断 vs 排队"是每条输入自己的类别 —— 不是 inputch 的属性。 +// +// 本文件只放**登记表**(谁是注册者、划给了谁、容量多少、记忆策略是什么)。 +// 路由本身发生在进内核之前:投递方决定"这条输入投给哪个 inputch"。 + +import ( + "errors" + "sort" + "sync" +) + +var ( + // ErrInputChannelUnknown 表示引用了未注册的 inputch。 + ErrInputChannelUnknown = errors.New("inputch 未注册") + // ErrInputChannelNameEmpty 表示 inputch 名为空。 + ErrInputChannelNameEmpty = errors.New("inputch 名不能为空") +) + +// InputChannel 是一个 inputch 的完整登记记录。 +type InputChannel struct { + // Name 是路由单位 id(全局唯一,如 "qq"、"webui"、"qq/device-2")。 + Name string `json:"name"` + // Plugin 是注册它的插件名("插件可注册多个 inputch",归属可追溯)。 + Plugin string `json:"plugin,omitempty"` + // Owner 是**被划给的 agent id**("" = 未分配,归根 agent/内核默认)。 + Owner string `json:"owner,omitempty"` + // Capacity 是该 inputch 的队列容量(0 = 用内核默认值)。 + Capacity int `json:"capacity,omitempty"` + // Output 是该 inputch 的默认回程输出通道("" = 由来源/调用方决定)。 + // + // 注意:这不是"内核路由"—— 输出仍然是 agent 的主动调用;这里只是登记 + // "这个 inputch 的回复默认该往哪个 outputch 走"的映射依据。 + Output string `json:"output,omitempty"` + // Def 是记忆/上下文策略(沿用 ChannelDef:NoMemory / Cleaner / ContextPolicy)。 + // Cleaner 是函数,故本字段不可序列化(json:"-")。 + Def ChannelDef `json:"-"` +} + +// ChannelRegistry 是 inputch 的登记表。 +// +// 它被设计成**可共享对象**(`*ChannelRegistry`):根 agent 与它的驻留子共用同一份, +// 这样"划入/授权"才有意义;每个 IOManager 默认自带一份(向后兼容)。 +type ChannelRegistry struct { + mu sync.RWMutex + channels map[string]InputChannel +} + +// NewChannelRegistry 构造一个空的 inputch 登记表。 +func NewChannelRegistry() *ChannelRegistry { + return &ChannelRegistry{channels: make(map[string]InputChannel)} +} + +// Register 登记/更新一个 inputch。 +// +// 重复登记(插件重载)时**保留已有的 Owner/Capacity/Output**,只更新 +// Plugin 与 Def —— 否则一次插件重载就会把父 agent 做的划分抹掉。 +func (r *ChannelRegistry) Register(ch InputChannel) error { + if ch.Name == "" { + return ErrInputChannelNameEmpty + } + r.mu.Lock() + defer r.mu.Unlock() + if r.channels == nil { + r.channels = make(map[string]InputChannel) + } + if old, ok := r.channels[ch.Name]; ok { + if ch.Owner == "" { + ch.Owner = old.Owner + } + if ch.Capacity == 0 { + ch.Capacity = old.Capacity + } + if ch.Output == "" { + ch.Output = old.Output + } + if ch.Plugin == "" { + ch.Plugin = old.Plugin + } + } + r.channels[ch.Name] = ch + return nil +} + +// Unregister 注销一个 inputch。 +func (r *ChannelRegistry) Unregister(name string) { + r.mu.Lock() + defer r.mu.Unlock() + delete(r.channels, name) +} + +// Lookup 查询一个 inputch。 +func (r *ChannelRegistry) Lookup(name string) (InputChannel, bool) { + r.mu.RLock() + defer r.mu.RUnlock() + ch, ok := r.channels[name] + return ch, ok +} + +// Assign 把一个 inputch **划给**某个 agent(可同时给定容量)。 +// +// 语义(默认取值,见设计文档 R6):**读写授权**,不转移所有权 —— +// 登记表仍记录 Plugin(谁注册的)与 Owner(划给了谁)两件事。 +func (r *ChannelRegistry) Assign(name, agentID string, capacity int) error { + r.mu.Lock() + defer r.mu.Unlock() + ch, ok := r.channels[name] + if !ok { + return ErrInputChannelUnknown + } + ch.Owner = agentID + if capacity > 0 { + ch.Capacity = capacity + } + r.channels[name] = ch + return nil +} + +// List 返回全部已注册 inputch(按名字稳定排序)。 +func (r *ChannelRegistry) List() []InputChannel { + r.mu.RLock() + defer r.mu.RUnlock() + out := make([]InputChannel, 0, len(r.channels)) + for _, ch := range r.channels { + out = append(out, ch) + } + sort.Slice(out, func(i, j int) bool { return out[i].Name < out[j].Name }) + return out +} + +// ListByOwner 返回划给某个 agent 的 inputch(owner == "" 时返回**未分配**的)。 +func (r *ChannelRegistry) ListByOwner(agentID string) []InputChannel { + var out []InputChannel + for _, ch := range r.List() { + if ch.Owner == agentID { + out = append(out, ch) + } + } + return out +} + +// Count 返回已注册 inputch 数量。 +func (r *ChannelRegistry) Count() int { + r.mu.RLock() + defer r.mu.RUnlock() + return len(r.channels) +} diff --git a/internal/agent/io/inputch_test.go b/internal/agent/io/inputch_test.go new file mode 100644 index 0000000..dce1de5 --- /dev/null +++ b/internal/agent/io/inputch_test.go @@ -0,0 +1,134 @@ +package io + +// inputch 登记表与划分(单工具多视图背后的数据面)。 +// +// 设计依据 docs/zh/resident-subagent-design.md §4: +// inputch 是**最基本的输入路由单位**,由插件注册,一个插件可注册多个。 + +import "testing" + +func TestChannelRegistry_RegisterKeepsAttribution(t *testing.T) { + r := NewChannelRegistry() + + // 一个插件注册多个 inputch(最基本的输入路由单位)。 + if err := r.Register(InputChannel{Name: "qq", Plugin: "qq", Def: ChannelDef{NoMemory: true}}); err != nil { + t.Fatal(err) + } + if err := r.Register(InputChannel{Name: "qq/device-2", Plugin: "qq"}); err != nil { + t.Fatal(err) + } + if n := r.Count(); n != 2 { + t.Fatalf("注册数=%d,期望 2(同一插件的多个 inputch 彼此独立)", n) + } + + ch, ok := r.Lookup("qq/device-2") + if !ok || ch.Plugin != "qq" { + t.Fatalf("归属插件未记录:%+v ok=%v", ch, ok) + } + if _, ok := r.Lookup("nope"); ok { + t.Fatal("未注册的 inputch 不应查得到") + } + if err := r.Register(InputChannel{Name: ""}); err != ErrInputChannelNameEmpty { + t.Fatalf("空名应报错,实际 %v", err) + } +} + +// 插件重载(重复登记)**不得抹掉划分**:Owner/Capacity/Output 必须保留。 +func TestChannelRegistry_ReRegisterKeepsAllocation(t *testing.T) { + r := NewChannelRegistry() + _ = r.Register(InputChannel{Name: "qq", Plugin: "qq"}) + if err := r.Assign("qq", "agent-1", 128); err != nil { + t.Fatal(err) + } + + // 插件重载:只带 Plugin/Def,不带 Owner/Capacity。 + if err := r.Register(InputChannel{Name: "qq", Plugin: "qq", Def: ChannelDef{NoMemory: true}}); err != nil { + t.Fatal(err) + } + ch, _ := r.Lookup("qq") + if ch.Owner != "agent-1" || ch.Capacity != 128 { + t.Fatalf("重载后划分被抹掉:owner=%q capacity=%d", ch.Owner, ch.Capacity) + } + if !ch.Def.NoMemory { + t.Fatal("重载应更新 Def") + } +} + +func TestChannelRegistry_AssignAndViews(t *testing.T) { + r := NewChannelRegistry() + _ = r.Register(InputChannel{Name: "qq", Plugin: "qq"}) + _ = r.Register(InputChannel{Name: "cli", Plugin: "cli"}) + _ = r.Register(InputChannel{Name: "sub/in", Plugin: "sub"}) + + if err := r.Assign("qq", "root", 0); err != nil { + t.Fatal(err) + } + if err := r.Assign("sub/in", "child-1", 32); err != nil { + t.Fatal(err) + } + if err := r.Assign("nope", "x", 0); err != ErrInputChannelUnknown { + t.Fatalf("划分未注册的 inputch 应报错,实际 %v", err) + } + + if got := len(r.ListByOwner("root")); got != 1 { + t.Fatalf("root 的 inputch 数=%d,期望 1", got) + } + child := r.ListByOwner("child-1") + if len(child) != 1 || child[0].Name != "sub/in" || child[0].Capacity != 32 { + t.Fatalf("child-1 的划分=%+v", child) + } + // 未分配的:cli。 + un := r.ListByOwner("") + if len(un) != 1 || un[0].Name != "cli" { + t.Fatalf("未分配的 inputch=%+v,期望只有 cli", un) + } + // List 按名排序(视图输出稳定)。 + all := r.List() + if len(all) != 3 || all[0].Name != "cli" || all[2].Name != "sub/in" { + t.Fatalf("List 未按名排序:%+v", all) + } +} + +// 共享登记表:根 agent 与驻留子共用同一份,划分才有意义。 +func TestChannelRegistry_SharedBetweenManagers(t *testing.T) { + shared := NewChannelRegistry() + root := NewIOManager() + child := NewIOManager() + root.SetChannelRegistry(shared) + child.SetChannelRegistry(shared) + + if err := root.RegisterInputChannelFrom("sub", "sub/in", ChannelDef{}); err != nil { + t.Fatal(err) + } + if err := root.AssignInputChannel("sub/in", "child-1", 8); err != nil { + t.Fatal(err) + } + + // 子在**自己**的 io 上就能看到这份划分。 + ch, ok := child.LookupInputChannel("sub/in") + if !ok { + t.Fatal("共享登记表后,子应看得到根注册的 inputch") + } + if ch.Owner != "child-1" || ch.Capacity != 8 { + t.Fatalf("子看到的划分=%+v", ch) + } + + // 注销也要跨 manager 生效。 + root.UnregisterInputChannel("sub/in") + if _, ok := child.LookupInputChannel("sub/in"); ok { + t.Fatal("注销应跨 manager 生效") + } +} + +// 记忆策略查询保持向后兼容(原 GetInputChannelDef 的语义)。 +func TestChannelRegistry_DefLookupCompat(t *testing.T) { + m := NewIOManager() + m.RegisterInputChannel("qq", ChannelDef{NoMemory: true, ContextPolicy: "prune"}) + def, ok := m.GetInputChannelDef("qq") + if !ok || !def.NoMemory || def.ContextPolicy != "prune" { + t.Fatalf("策略查询=%+v ok=%v", def, ok) + } + if _, ok := m.GetInputChannelDef("nope"); ok { + t.Fatal("未注册的 inputch 不应有策略") + } +} diff --git a/internal/plugin/registry.go b/internal/plugin/registry.go index e05c54d..1559a2f 100644 --- a/internal/plugin/registry.go +++ b/internal/plugin/registry.go @@ -311,7 +311,9 @@ func (r *Registry) buildSDK(name string) *sdk.PluginSDK { if r.iom == nil { return nil } - r.iom.RegisterInputChannel(chName, agentIO.ChannelDef(def)) + // inputch 是最基本的输入路由单位:登记**归属插件**,便于父 agent 看清 + // "哪个插件的哪个 inputch 划给了谁"(一个插件可注册多个 inputch)。 + _ = r.iom.RegisterInputChannelFrom(name, chName, agentIO.ChannelDef(def)) r.noteChannel(name, chName, false) return nil } diff --git a/internal/plugins/deepsearch_e2e_test.go b/internal/plugins/deepsearch_e2e_test.go new file mode 100644 index 0000000..2e65885 --- /dev/null +++ b/internal/plugins/deepsearch_e2e_test.go @@ -0,0 +1,118 @@ +//go:build linux || darwin + +package plugins + +import ( + "fmt" + "path/filepath" + "strings" + "testing" +) + +// deepsearch 插件的真实调用往返:内核 ExecuteTool → RPC → 插件子进程 → 本地 SearXNG → 结果回传。 +// +// 与 TestRealPlugin_ToolInvokeRoundTrip 的区别:那条只断言「链路通(拿到结果或拿到插件侧的错误)」, +// 这条断言**内容形状**——返回里必须有「摘要」与「引擎覆盖度」。这正是旧实现(抓 Bing HTML) +// 拿不到的东西,也是「搜索能力不行」的根因,所以它必须成为回归判据。 +// +// 前置:本机 127.0.0.1:8888 上跑着 SearXNG(部署见 /root/searxng-agent)。 +// 未启动时 fail 并给出可操作提示,而不是 skip —— 否则这条判据会在环境退化时静默失效。 +func TestRealPlugin_DeepSearchInvoke(t *testing.T) { + env := setupIntegration(t) + defer env.cleanup() + + plgDir := filepath.Join(env.tmpDir, "plugins") + installRealPlugin(t, plgDir, "deepsearch") + + if err := env.pluginReg.Load(plgDir); err != nil { + t.Fatalf("加载插件: %v", err) + } + if env.pluginReg.Get("deepsearch") == nil { + t.Fatal("deepsearch 未经 proc 通道加载") + } + + var toolName string + for _, def := range env.stageHost.GetToolDefs() { + if strings.HasPrefix(def.Name, "deepsearch") && strings.HasSuffix(def.Name, "_search") { + toolName = def.Name + break + } + } + if toolName == "" { + t.Fatal("deepsearch 未注册检索工具") + } + t.Logf("调用工具 %s", toolName) + + // 用当初失败的那条查询做判据 + res, err := env.stageHost.ExecuteTool(toolName, map[string]interface{}{ + "query": "深度科技 deepin 开发者 被开除", + "count": float64(3), + }) + if err != nil { + if strings.Contains(err.Error(), "not found in any plugin") { + t.Fatalf("工具未注册到 stageHost: %v", err) + } + t.Fatalf("工具调用失败(检查本机 SearXNG 是否在 127.0.0.1:8888 运行): %v", err) + } + if res == nil { + t.Fatal("工具返回 nil 且无错误") + } + + text := fmt.Sprintf("%v", res) + t.Logf("工具返回前 500 字:\n%s", truncRunes(text, 500)) + + if !strings.Contains(text, "摘要:") { + t.Errorf("返回内容缺少摘要——这正是旧实现拿不到的部分:\n%s", truncRunes(text, 800)) + } + if !strings.Contains(text, "覆盖:") { + t.Errorf("返回内容缺少引擎覆盖度(模型据此判断可信度):\n%s", truncRunes(text, 800)) + } + if !strings.Contains(text, "http") { + t.Errorf("返回内容缺少结果链接:\n%s", truncRunes(text, 800)) + } +} + +// deepsearch_status 也走一遍真实调用:它把「后端是否可用、哪些引擎在出结果」暴露成工具, +// 出问题时 agent 可以先自检,而不是盲目换词重搜。 +func TestRealPlugin_DeepSearchStatusInvoke(t *testing.T) { + env := setupIntegration(t) + defer env.cleanup() + + plgDir := filepath.Join(env.tmpDir, "plugins") + installRealPlugin(t, plgDir, "deepsearch") + if err := env.pluginReg.Load(plgDir); err != nil { + t.Fatalf("加载插件: %v", err) + } + + var toolName string + for _, def := range env.stageHost.GetToolDefs() { + if strings.HasPrefix(def.Name, "deepsearch") && strings.HasSuffix(def.Name, "_status") { + toolName = def.Name + break + } + } + if toolName == "" { + t.Fatal("deepsearch 未注册自检工具") + } + + res, err := env.stageHost.ExecuteTool(toolName, map[string]interface{}{"probe": "test"}) + if err != nil { + t.Fatalf("自检调用失败(检查本机 SearXNG 是否运行): %v", err) + } + text := fmt.Sprintf("%v", res) + t.Logf("自检返回:%s", truncRunes(text, 400)) + + for _, want := range []string{"healthz", "search_ok", "engines_returning_results"} { + if !strings.Contains(text, want) { + t.Errorf("自检结果缺少字段 %q:%s", want, truncRunes(text, 500)) + } + } +} + +func truncRunes(s string, n int) string { + r := []rune(s) + if len(r) <= n { + return s + } + return string(r[:n]) + "…" +}