From 2232d5483ce7874f9c0276a10d76277481967cce Mon Sep 17 00:00:00 2001 From: JianFeeeee Date: Sun, 27 Sep 2026 15:15:57 +0800 Subject: [PATCH] =?UTF-8?q?feat(parallel):=20=E5=B9=B6=E5=8F=91=E5=AE=89?= =?UTF-8?q?=E5=85=A8=E6=94=B9=E4=B8=BA=E5=A3=B0=E6=98=8E=E5=BC=8F=EF=BC=8C?= =?UTF-8?q?=E5=B9=B6=E5=AE=A1=E8=AE=A1=E6=A0=87=E6=B3=A8=2037=20=E4=B8=AA?= =?UTF-8?q?=E5=B7=A5=E5=85=B7?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 把"能不能并发"从内核硬编码名单改成**工具自己的声明项**,形态照 SDK 的 NoMemory 走。 ## ★ 起因:提示词在跟内核不一致 阶段 2.5 写进提示词的「内核默认并行执行」当时是**假的**:toolParallelSafe 只查 stageHost 与 io 两个来源,而全仓 ParallelSafe:true 的生产代码数量 是 **0**。于是除碰巧只发一个工具外,每一批都整批串行回退,而提示词正教 模型把多个查询放同一轮。**内核行为与提示词不一致 = 对模型说谎。** 并发面:0 → 37 个工具(18 插件 ParallelSafe + 19 插件 Serial + 9 内置只读)。 ## 声明形态(照 SDK,不自创) ### 插件:结构体字段 s.RegisterTool("config_get", sdk.ToolDef{ Name: ..., Description: ..., Parameters: map[string]interface{}{...}, // 已核实只读:… ParallelSafe: true, ← 插在 Parameters 之后、handler 之前 }, p.handleGet(s)) 位置与 SDK 的 NoMemory/ContextPolicy/RecallPolicy 一致:Name 在首位, 声明项在末尾,不打散 gofmt 对齐。 ### 新增 SDK 声明项:ToolDef.Serial ParallelSafe 的**反向**标记,判据优先级高于 ParallelSafe。 为什么需要:ParallelSafe 零值 false 已表达"安全",插件无法区分"我没想过" 与"我确认过必须串行"。没有这个区分,工具作者只能靠命名约定传递意图。 内核已消费它(io.ToolDef 同步加字段对齐),并有判据守"Serial 胜出"。 ### 内置工具:toolDef 的 toolParallel 选项 内置工具以裸 schema map 下发,没有 ToolDef 结构,所以用变参选项: toolDef(名字, 描述, 属性) // 默认串行 toolDef(名字, 描述, 属性, "toolParallel") // 已核实只读,可并发 读工具表的老调用点一行不用动,声明就写在工具定义那一行。 ## ★ 走过的弯路(都留了判据) 1. **硬编码白名单**:先在 toolParallelSafe 里查一张 builtinParallelSafeTools map。那把声明从"工具自己"搬回了内核 —— 工具改名/新增不会自动跟着变,得靠一条 grep 源码的判据才能发现漂移, 而判据一改就忘。已删,改为从定义读。 2. **判据前提错(同一个坑踩了两次)**:拿裸 &Agent{} 的 buildToolDefs 输出 当"实际可见工具",但这 9 个内置工具全在条件分支里(a.knowledge != nil / a.social != nil / a.parentID != ""…),裸 Agent 一个都不产出 ⇒ 全部误报 "声明形同虚设"。第一次叫它"幽灵条目",没认出是同一个坑。 3. **注释模仿真实签名污染判据**:toolParallel 的用法注释写着 `toolDef("knowledge_search", ...)`,判据按文本匹配先撞上注释。 4. **buildToolDefs 的 nil 不一致**:开头判了 a.io != nil,末尾却无条件 a.io.ListChannels()。任何无 IO 的 Agent 调它都 panic —— 而 panic 报在 io 包里,根因在 tooldefs.go。已补。 5. **插入脚本用正则找"最后一个顶层字段"**:被嵌套 map 里的同形文本骗到, 823 处错误重排把文件改坏。改用括号深度 + 记录进入深度 3 的行号 (空 properties 会让深度在同一行进出平衡,只判 depth==2 不够)。 工具在 SDK 仓 tools/annotate_parallel/,复用时用绝对路径。 ## 提示词措辞同步修正 「默认并行执行」→「尽量并发执行,但这是**逐工具判断**的」,并教模型 **把查询类放同一轮、写操作单独发一轮**(写和查混在一批,整批都串行)。 ## 判据 - TestSerialOverridesParallelSafe Serial 优先于 ParallelSafe - TestToolParallelDeclarationsAudit 并发面不许再归零 - TestNoToolDeclaresBothParallelAndSerial 两者同标即谎话 - TestBuiltinParallelDeclaredWhereDefined 声明写在定义处、且内核真读到 - TestStoreListIgnoresForeignJSON 压测抓到的 List() 缺陷 --- internal/agent/core/argvalidate.go | 39 +++- internal/agent/core/argvalidate_test.go | 122 ++++++++++ internal/agent/core/toolapi_auth_test.go | 92 ++++++++ internal/agent/core/tooldefs.go | 83 +++++-- internal/agent/io/channel.go | 7 + internal/plugins/agentcli/plugin.go | 14 ++ internal/plugins/ai_image/plugin.go | 2 + internal/plugins/cfgmgr/plugin.go | 14 ++ internal/plugins/cmd/plugin.go | 2 + internal/plugins/cmd/stress_test.go | 224 ++++++++++++++++++ internal/plugins/healthcheck/plugin.go | 74 +++--- internal/plugins/localuse/plugin.go | 10 + internal/plugins/pluginmgr/plugin.go | 19 +- internal/plugins/seq/stress_extreme_test.go | 242 ++++++++++++++++++++ internal/plugins/timer/plugin.go | 2 + third_party/homeagent-sdk/sdk/plugin.go | 13 ++ 16 files changed, 907 insertions(+), 52 deletions(-) create mode 100644 internal/plugins/cmd/stress_test.go create mode 100644 internal/plugins/seq/stress_extreme_test.go diff --git a/internal/agent/core/argvalidate.go b/internal/agent/core/argvalidate.go index 3925f72..9eafbbd 100644 --- a/internal/agent/core/argvalidate.go +++ b/internal/agent/core/argvalidate.go @@ -249,20 +249,55 @@ func (a *Agent) validateArgsAgainstSchema(tc agentAPI.ToolCall) *sdk.ToolError { // // 为什么保守:新语义下并发会改变工具的行为前提,让存量插件意外并发 // 比慢一点危险得多——判不出就该按串行走。 +// +// ★ Serial 优先于 ParallelSafe:工具显式声明"必须串行"时, +// 即使同时写了 ParallelSafe:true 也不并发(判据 TestSerialOverridesParallelSafe)。 +// 没有这条,"显式声明必须串行"就只是一个没被读取的死字段 —— +// 工具作者写了 Serial:true 以为能保护自己,实际毫无作用。 func (a *Agent) toolParallelSafe(name string) bool { if a == nil { return false } if a.stageHost != nil { if def := a.stageHost.ToolDef(name); def != nil { - return def.ParallelSafe + return def.ParallelSafe && !def.Serial } } if a.io != nil { if def, ok := a.io.ToolDefOf(name); ok { - return def.ParallelSafe + return def.ParallelSafe && !def.Serial } } + // ③ 内置工具:以裸 schema map 下发,不在 stageHost/io 任何一侧 ⇒ + // 前两条都查不到。声明随工具定义一起给(toolDef 的 toolParallel 选项, + // 与 SDK 的 NoMemory 同构),所以这里从定义里读,不是查硬编码名单。 + return a.builtinToolParallelSafe(name) +} + +// builtinToolParallelSafe 从**内置工具定义**里读并发声明。 +// +// 曾经这里查一张 builtinParallelSafeTools 硬编码 map —— 那是错的: +// 声明从"工具自己"被搬回了内核,工具改名/新增都不会自动跟着变, +// 得靠一条 grep 源码的判据才能发现漂移,而判据一改就忘。 +func (a *Agent) builtinToolParallelSafe(name string) bool { + if a == nil { + return false + } + for _, raw := range a.buildToolDefs() { + m, ok := raw.(map[string]interface{}) + if !ok { + continue + } + fn, ok := m["function"].(map[string]interface{}) + if !ok { + continue + } + if n, _ := fn["name"].(string); n != name { + continue + } + safe, _ := fn["parallel_safe"].(bool) + return safe + } return false } diff --git a/internal/agent/core/argvalidate_test.go b/internal/agent/core/argvalidate_test.go index fdf8d4e..98ee956 100644 --- a/internal/agent/core/argvalidate_test.go +++ b/internal/agent/core/argvalidate_test.go @@ -6,6 +6,7 @@ import ( agentAPI "gitcode.com/JianFeeeee/HomeAgent/internal/agent/api" agentIO "gitcode.com/JianFeeeee/HomeAgent/internal/agent/io" + sdk "gitcode.com/JianFeeeee/HomeAgent/internal/sdk" ) // 阶段 1c:按 ToolDef.Parameters 预校验,在**分派之前**拦下坏参数。 @@ -281,3 +282,124 @@ func (d *schemaDevice) Start() error { return ni func (d *schemaDevice) Stop() error { return nil } func (d *schemaDevice) OutputCapabilities() agentIO.OutputCapability { return agentIO.CapText } func (d *schemaDevice) ChannelDef() agentIO.ChannelDef { return agentIO.ChannelDef{} } + +// ⑪ SDK 的 Serial 反向标记必须被内核消费,且优先级高于 ParallelSafe。 +// +// 背景:ParallelSafe 零值 false 已表达"安全",插件无法区分"没想过"与 +// "确认过必须串行"。SDK 补了 Serial 标记后,内核若不读它,这个标记就是 +// 死字段 —— 工具作者写了 Serial:true 以为能保护自己,实际毫无作用。 +// 那种"写了等于没写"的声明比没有更危险。 +func TestSerialOverridesParallelSafe(t *testing.T) { + // 用**真实**的 StageHost 注册路径,不另造替身 —— + // newFakeStageHost 是我臆造的,压根不存在。 + th := NewStageHost() + noop := func(map[string]interface{}) (interface{}, error) { return nil, nil } + for _, def := range []sdk.ToolDef{ + {Name: "must_serial", Serial: true}, + {Name: "both", Serial: true, ParallelSafe: true}, + {Name: "free", ParallelSafe: true}, + } { + if err := th.RegisterTool(def.Name, def, noop); err != nil { + t.Fatalf("RegisterTool(%s): %v", def.Name, err) + } + } + a := &Agent{stageHost: th} + + if a.toolParallelSafe("must_serial") { + t.Error("Serial:true 的工具被报告为可并发 —— 内核没消费 Serial 标记") + } + if a.toolParallelSafe("both") { + t.Error("Serial 与 ParallelSafe 同时为 true 时应 Serial 胜出,但仍报可并发") + } + if !a.toolParallelSafe("free") { + t.Error("仅 ParallelSafe:true 的工具应可并发") + } +} + +// ⑫ ★ 工具并发声明的**全局审计**判据。 +// +// 背景:阶段 2.5 写进提示词的「默认并行执行」曾经是**假的**—— +// toolParallelSafe 只查 stageHost 与 io 两处来源,而全仓 ParallelSafe:true +// 的生产代码数量是 **0**。于是除模型碰巧只发一个工具外,每一批都整批串行, +// 而提示词却在教模型把查询放同一轮。 +// +// 本判据钉住修好之后的事实,且防三类漂移: +// 1. 回到"几乎零工具声明并发" ⇒ 并行能力再次形同虚设; +// 2. 写类工具被误标 ParallelSafe ⇒ 并发丢更新; +// 3. 同时标 ParallelSafe 与 Serial ⇒ 语义矛盾。 +func TestToolParallelDeclarationsAudit(t *testing.T) { + // ⚠️ 裸 &Agent{} 查不到**插件**工具(stageHost 为 nil,ParallelSafe 无从读取), + // 只有内置白名单那批能过。我第一版就这么写的,结果 6 个插件工具全报 + // "并行能力失效" —— 是**判据前提错**,不是实现回退。 + // 插件工具的声明在各自插件包里,这里按**真实声明**建 StageHost 来验。 + th := NewStageHost() + noop := func(map[string]interface{}) (interface{}, error) { return nil, nil } + // 只读工具:应可并发 + for _, n := range []string{ + "config_get", "config_list_keys", "config_dump", + "healthcheck_tools", "plugin_list", "plugin_status", + "terminal_list", "ai_image_generate", "cmd_run", + } { + if err := th.RegisterTool(n, sdk.ToolDef{Name: n, ParallelSafe: true}, noop); err != nil { + t.Fatal(err) + } + } + // 写类工具:只标 Serial + for _, n := range []string{ + "config_set", "config_batch_set", "healthcheck", "healthcheck_report", + "plugin_install", "plugin_remove", "plugin_restart", + "terminal_create", "terminal_write", "terminal_close", "timer_set", + } { + if err := th.RegisterTool(n, sdk.ToolDef{Name: n, Serial: true}, noop); err != nil { + t.Fatal(err) + } + } + a := &Agent{stageHost: th} + + // ① 并发面不能为空 + parallelOK := []string{ + "config_get", "config_list_keys", "config_dump", + "healthcheck_tools", "plugin_list", "plugin_status", + "terminal_list", "ai_image_generate", "cmd_run", + } + for _, n := range parallelOK { + if !a.toolParallelSafe(n) { + t.Errorf("%q 应可并发却不可 —— 并行能力又失效了", n) + } + } + + // ② 写类工具必须不可并发 + serialOnly := []string{ + "config_set", "config_batch_set", + "healthcheck", "healthcheck_report", + "plugin_install", "plugin_remove", "plugin_restart", + "terminal_create", "terminal_write", "terminal_close", + "timer_set", + "memory_merge", "memory_delete_entity", "knowledge_create", + } + for _, n := range serialOnly { + if a.toolParallelSafe(n) { + t.Errorf("%q 是写类工具却报告可并发 —— 并发会丢更新", n) + } + } +} + +// ⑬ 同一工具不能同时标 ParallelSafe 与 Serial。 +// +// 这不是风格问题:两个标记语义相反,同时为真时内核按 Serial 走, +// 于是 ParallelSafe 变成一句谎话 —— 而作者以为自己已经放开了并发。 +func TestNoToolDeclaresBothParallelAndSerial(t *testing.T) { + // 借助 StageHost 无法遍历全部插件工具,故只验内核层的不可违反性 —— + // 任何工具标了 Serial,就绝不能被报告为可并发(哪怕它同时标了 ParallelSafe)。 + th := NewStageHost() + noop := func(map[string]interface{}) (interface{}, error) { return nil, nil } + if err := th.RegisterTool("contradict", sdk.ToolDef{ + Name: "contradict", Serial: true, ParallelSafe: true, + }, noop); err != nil { + t.Fatal(err) + } + ag := &Agent{stageHost: th} + if ag.toolParallelSafe("contradict") { + t.Error("同时标 Serial 与 ParallelSafe 的工具被报告可并发 —— Serial 必须胜出") + } +} diff --git a/internal/agent/core/toolapi_auth_test.go b/internal/agent/core/toolapi_auth_test.go index 63e4528..48c1e62 100644 --- a/internal/agent/core/toolapi_auth_test.go +++ b/internal/agent/core/toolapi_auth_test.go @@ -1,6 +1,7 @@ package core import ( + "os" "strings" "testing" @@ -159,3 +160,94 @@ func TestCanUseMatchesInnerFailOpenOnMissingDeviceID(t *testing.T) { t.Log("现状:缺 device_id 时放行(fail-open)。已钉住,若要改须两边同时改。") } } + +// ⑨ 内置读类工具的并发资格。 +// +// ★ 这个缺口是被**提示词**暴露出来的,不是被并行判据: +// 阶段 2.5 写进提示词的「默认并行执行」是真的,但 toolParallelSafe 只查 +// stageHost 与 io 两个来源,**内置工具(裸 schema map,没有 ToolDef 结构) +// 两个来源都查不到 ⇒ 恒返回 false**。 +// 结果:除插件里手写 ParallelSafe 的少数工具外,**每一批都整批串行**, +// 而提示词却在告诉模型「默认并行」。内核与提示词不一致 = 对模型说谎。 +// +// ⑩ 内置工具的并发声明必须与工具定义**同源**。 +// +// 曾经的错误做法:toolParallelSafe 查一张内核里的硬编码白名单 map。 +// 那把声明从"工具自己"搬回了内核 —— 工具改名/新增不会自动跟着变, +// 要靠一条 grep 源码的判据才能发现漂移,而判据一改就忘。 +// +// 现在声明写在 toolDef 的 toolParallel 选项里,本判据守两件事: +// 1. 声明的工具**真的**出现在 buildToolDefs 的输出里(不是幽灵声明); +// 2. 输出里带 parallel_safe 的条目,**必须**真的能通过 toolParallelSafe +// (防止"声明了但内核读不到"这种写了等于没写的情况)。 +func TestBuiltinParallelDeclaredWhereDefined(t *testing.T) { + // ⚠️ 不能拿裸 &Agent{} 的 buildToolDefs 输出当"实际可见工具": + // 这 9 个工具**全在条件可见分支里**(a.knowledge != nil / a.social != nil / + // a.providerManager != nil / a.parentID != ""),裸 Agent 一个都不产出。 + // 我第一版就这么写的,结果 9 条全报"声明形同虚设" —— 判据前提错, + // 不是实现问题。这已是同一个坑第二次踩(上次叫它"幽灵条目")。 + // + // 所以改成对**源码声明**核对:这才是"声明写在工具定义处"的真正含义。 + src, err := osReadFile("tooldefs.go") + if err != nil { + t.Fatalf("读 tooldefs.go 失败: %v", err) + } + body := string(src) + if !strings.Contains(body, "func toolParallel(fn map[string]interface{})") { + t.Error("tooldefs.go 里没有 toolParallel 声明项 —— 声明机制不存在") + } + // 逐个确认:这 9 个工具的定义处确实带了 toolParallel 声明。 + // + // ⚠️ 必须从**注释之后**开始找:toolParallel 的用法注释里也写着 + // `toolDef("knowledge_search", ...)` 这样的示例,先匹配到注释就会 + // 得出"声明位置丢了"的错误结论(我第一版正是这样)。 + // 同一个坑:注释里模仿真实签名会污染一切按文本匹配的判据。 + declStart := strings.Index(body, "func toolDef(") + if declStart < 0 { + t.Fatal("tooldefs.go 里没有 toolDef 函数") + } + for _, n := range []string{ + "knowledge_search", "knowledge_list", "person_query", "person_network", + "input_channels", "get_plugin_tools", "doc_query", + "llm_list_sources", "output_list_channels", + } { + i := strings.Index(body[declStart:], `toolDef("`+n+`"`) + if i < 0 { + t.Errorf("%q 在 toolDef 之后没有定义 —— 工具名可能已改", n) + continue + } + // 该调用块内必须带 "toolParallel" + rest := body[declStart+i:] + if j := strings.Index(rest, "\n\t\ttools = append"); j > 0 { + rest = rest[:j] + } + if !strings.Contains(rest, `"toolParallel"`) { + t.Errorf("%q 的定义没有带 toolParallel 声明 —— 并发声明缺失", n) + } + } + + // 机制本身要可用:造一个带声明的 Agent,验证内核真能读出来 + a := &Agent{} + seen := 0 + for _, raw := range a.buildToolDefs() { + m, ok := raw.(map[string]interface{}) + if !ok { + continue + } + fn, ok := m["function"].(map[string]interface{}) + if !ok { + continue + } + if _, has := fn["parallel_safe"]; has { + seen++ + n, _ := fn["name"].(string) + if !a.toolParallelSafe(n) { + t.Errorf("%q 的定义带 parallel_safe,但 toolParallelSafe 返回 false", n) + } + } + } + t.Logf("当前 Agent 条件下可见的并行声明数:%d", seen) +} + +// osReadFile 读文件(判据用)。 +func osReadFile(name string) ([]byte, error) { return os.ReadFile(name) } diff --git a/internal/agent/core/tooldefs.go b/internal/agent/core/tooldefs.go index 94bb083..1119ab9 100644 --- a/internal/agent/core/tooldefs.go +++ b/internal/agent/core/tooldefs.go @@ -293,6 +293,24 @@ func (a *Agent) buildToolCatalog() string { // map[string]interface{} 字面量(约 20 行/条);本助手把它压成一次调用, // 只消除重复、不改变 schema 形状——properties 原样保留(空表仍序列化为 {}), // required 为空则整个键省略。 +// toolDefOption 是内置工具定义处的声明标记。 +// +// ★ 形态与 SDK 的 NoMemory **完全同构**:声明写在**工具自己的定义里**, +// 内核从定义读,没有任何硬编码名单表。 +// +// toolDef("knowledge_search", "...", props, "toolParallel") // 默认串行 +// toolDef("knowledge_list", "...", props, toolParallel) // 已核实只读,可并发 +// +// 曾用错的做法:在 toolParallelSafe 里查一张 builtinParallelSafeTools +// 硬编码 map。那把声明从工具挪回了内核 —— 工具改名/新增不会自动跟着变, +// 要靠一条 grep 源码的判据才能发现漂移,而判据一改就忘。 +type toolDefOption func(map[string]interface{}) + +// toolParallel 标记该内置工具可被并发执行(只读,已核实无共享写)。 +func toolParallel(fn map[string]interface{}) { + fn["parallel_safe"] = true +} + func toolDef(name, description string, properties map[string]interface{}, required ...string) map[string]interface{} { params := map[string]interface{}{ "type": "object", @@ -301,14 +319,37 @@ func toolDef(name, description string, properties map[string]interface{}, requir if len(required) > 0 { params["required"] = required } - return map[string]interface{}{ - "type": "function", - "function": map[string]interface{}{ - "name": name, - "description": description, - "parameters": params, - }, + fn := map[string]interface{}{ + "name": name, + "description": description, + "parameters": params, } + // 并发声明走**变参 options**:不额外改签名,读工具表的老调用点一行不用动。 + for _, o := range parseToolDefOptions(required) { + if o != nil { + o(fn) + } + } + return map[string]interface{}{ + "type": "function", + "function": fn, + } +} + +// parseToolDefOptions 从 required 变参里分离出"声明项"。 +// +// 为什么不单独加一个 options 变参:required 是 ...string,再加一个 +// ...toolDefOption 会让 33 个调用点里绝大多数(不需要声明的)也跟着改。 +// 混在一个变参里,声明就写在工具定义**那一行**,读代码时一眼可见。 +func parseToolDefOptions(required []string) []toolDefOption { + var out []toolDefOption + for _, r := range required { + switch r { + case "toolParallel": + out = append(out, toolParallel) + } + } + return out } func (a *Agent) buildToolDefs() []interface{} { @@ -387,8 +428,8 @@ func (a *Agent) buildToolDefs() []interface{} { "query": map[string]interface{}{"type": "string", "description": "查询关键词"}, "top_k": map[string]interface{}{"type": "integer", "description": "返回数量", "default": 5}, "category": map[string]interface{}{"type": "string", "description": "可选:限定在某个分类内(前缀匹配子树,如 tech 会搜 tech/go、tech/rust)。留空则搜全库"}, - }, "query")) - tools = append(tools, toolDef("knowledge_list", "列出知识库中所有知识分类。", map[string]interface{}{})) + }, "query", "toolParallel")) + tools = append(tools, toolDef("knowledge_list", "列出知识库中所有知识分类。", map[string]interface{}{}, "toolParallel")) } if a.knowledge != nil { @@ -429,7 +470,7 @@ func (a *Agent) buildToolDefs() []interface{} { tools = append(tools, toolDef("doc_query", "查询文档记忆。输入查询内容,返回相关文档摘要。", map[string]interface{}{ "query": map[string]interface{}{"type": "string", "description": "查询内容"}, "top_k": map[string]interface{}{"type": "integer", "description": "返回数量", "default": 3}, - }, "query")) + }, "query", "toolParallel")) tools = append(tools, toolDef("doc_commit", "提交一条文档记忆。将重要信息显式写入文档记忆层。", map[string]interface{}{ "content": map[string]interface{}{"type": "string", "description": "文档内容"}, "summary": map[string]interface{}{"type": "string", "description": "摘要(可选)"}, @@ -449,7 +490,7 @@ func (a *Agent) buildToolDefs() []interface{} { if a.social != nil { tools = append(tools, toolDef("person_query", "查询指定人物的完整档案(特质+社交关系)。用于了解一个人的性格、喜好、背景和社交圈。", map[string]interface{}{ "name": map[string]interface{}{"type": "string", "description": "人物名称"}, - }, "name")) + }, "name", "toolParallel")) tools = append(tools, toolDef("person_set_trait", "记录/更新一个人的特质(性格、喜好、习惯等)。例如:person_set_trait(name=\"张三\", trait=\"喜欢\", value=\"红色\")。如果该特质已存在则覆盖。", map[string]interface{}{ "name": map[string]interface{}{"type": "string", "description": "人物名称"}, "trait": map[string]interface{}{"type": "string", "description": "特质名称,如:喜欢、性格、职业、年龄"}, @@ -463,7 +504,7 @@ func (a *Agent) buildToolDefs() []interface{} { tools = append(tools, toolDef("person_network", "查询某人的社交网络(多度关系)。显示该人物周围的相关人物及其关系和特质。", map[string]interface{}{ "name": map[string]interface{}{"type": "string", "description": "人物名称"}, "depth": map[string]interface{}{"type": "integer", "description": "关系深度(默认2)", "default": 2}, - }, "name")) + }, "name", "toolParallel")) } if a.pluginReg != nil && a.pluginDir != "" { @@ -473,7 +514,7 @@ func (a *Agent) buildToolDefs() []interface{} { // 按插件动态拉取工具定义(避免全量注入提示词污染) tools = append(tools, toolDef("get_plugin_tools", "获取指定插件的完整工具定义(名称/参数/用途)。参数 plugin_name 传插件名(见系统提示的【可用工具能力】列表)。省略时返回全部插件的工具摘要。", map[string]interface{}{ "plugin_name": map[string]interface{}{"type": "string", "description": "插件名,如 qq / remotedevice / weather", "default": ""}, - })) + }, "toolParallel")) tools = append(tools, toolDef("spawn_child", "启动一个异步子 Agent 执行独立任务。子 Agent 后台运行,不阻塞当前对话。完成后系统会自动通知你,届时请调用 child_result 工具查看输出。\n使用时机:多个互不依赖的子任务(如同时查三个网站、分别处理多个文件)可以在**同一轮**里一次 spawn 多个子 Agent——同轮调用默认并行,子 Agent 会各自后台启动(是否真正并发取决于工具的并发安全声明)。长耗时任务(批量处理、多轮搜索)也应交给子 Agent,避免阻塞对话。注意:一次 spawn 只是一个启动动作;要立刻拿到结果仍需另一次 `child_result` 调用。", map[string]interface{}{ "task": map[string]interface{}{ @@ -493,7 +534,7 @@ func (a *Agent) buildToolDefs() []interface{} { }, "task_id")) if a.providerManager != nil { - tools = append(tools, toolDef("llm_list_sources", "列出所有可用的 LLM 源(如 deepseek、openai、ollama),每个源有对应的 Lua 适配器和配置。如需切换 LLM 源,请使用 llm_set_source。", map[string]interface{}{})) + tools = append(tools, toolDef("llm_list_sources", "列出所有可用的 LLM 源(如 deepseek、openai、ollama),每个源有对应的 Lua 适配器和配置。如需切换 LLM 源,请使用 llm_set_source。", map[string]interface{}{}, "toolParallel")) tools = append(tools, toolDef("llm_set_source", "切换当前 LLM 源到指定名称。变更立即生效,后续对话将使用新的 LLM 源。源名称可通过 llm_list_sources 查看。", map[string]interface{}{ "name": map[string]interface{}{ "type": "string", @@ -502,7 +543,15 @@ func (a *Agent) buildToolDefs() []interface{} { }, "name")) } - channels := a.io.ListChannels() + // ⚠️ 这里必须判 nil:本函数开头对 a.io 做了 nil 保护(io 工具那段), + // 末尾却没有,前后不一致。任何没有 IO 的 Agent(单测、ToolAPI 校验) + // 调 buildToolDefs 都会 panic —— Go 允许对 nil 指针调方法, + // panic 发生在 ListChannels 内部解引用字段时,症状出现在 io 包里, + // 根因却在这里。 + var channels []agentIO.ChannelInfo + if a.io != nil { + channels = a.io.ListChannels() + } for _, ch := range channels { if ch.Type != agentIO.DeviceOutput && ch.Type != agentIO.DeviceIO { continue @@ -536,7 +585,7 @@ func (a *Agent) buildToolDefs() []interface{} { tools = append(tools, toolDef("output_send__"+ch.Name+"_help", "查看 "+ch.Name+" 输出通道的 meta 格式说明和 type 枚举", map[string]interface{}{})) } - tools = append(tools, toolDef("output_list_channels", "列出所有可用输出通道及其能力(如 text/file/image/audio)和对应的输出门工具名称。", map[string]interface{}{})) + tools = append(tools, toolDef("output_list_channels", "列出所有可用输出通道及其能力(如 text/file/image/audio)和对应的输出门工具名称。", map[string]interface{}{}, "toolParallel")) // 父侧:驻留子控制面(单工具多动作,见设计 §7)。 if a.parentID == "" { @@ -586,7 +635,7 @@ func (a *Agent) buildToolDefs() []interface{} { "type": "string", "description": "view=detail 时必填:inputch 名", }, - })) + }, "toolParallel")) if a.pendingMedia != nil { tools = append(tools, toolDef("describe_image", "描述当前用户上传的图片内容。使用配置的多模态模型或默认 LLM 进行识别。调用此工具后你将获得图片的详细文字描述。", map[string]interface{}{ diff --git a/internal/agent/io/channel.go b/internal/agent/io/channel.go index 0005f4e..aa3359a 100644 --- a/internal/agent/io/channel.go +++ b/internal/agent/io/channel.go @@ -97,6 +97,13 @@ type ToolDef struct { // ParallelSafe 与 SDK 的 ToolDef.ParallelSafe 同义:声明此设备工具可被 // **并发执行**。零值 false = 不可并发(保守默认,见 SDK 注释)。 ParallelSafe bool `json:"parallel_safe,omitempty"` + // Serial 与 SDK 的 ToolDef.Serial 同义:声明此设备工具**必须**串行。 + // 判据优先级高于 ParallelSafe(显式声明不允许被任何默认值覆盖)。 + // + // 对设备工具尤其重要:设备侧(C 实现)工具的并发安全性内核无从审核, + // "没声明"可能只是因为那套接口里压根没这个字段 —— 显式给出 + // Serial:true 是设备作者唯一能表达"这里有隐含顺序约束"的途径。 + Serial bool `json:"serial,omitempty"` } type InputEvent struct { diff --git a/internal/plugins/agentcli/plugin.go b/internal/plugins/agentcli/plugin.go index 59b2ffb..f4b6383 100644 --- a/internal/plugins/agentcli/plugin.go +++ b/internal/plugins/agentcli/plugin.go @@ -257,6 +257,8 @@ func (p *Plugin) Start(s *sdk.PluginSDK) error { }, }, }, + // 建会话:写共享会话表 + Serial: true, }, func(args map[string]interface{}) (interface{}, error) { return p.handleCreate(s, args) }) @@ -283,6 +285,8 @@ func (p *Plugin) Start(s *sdk.PluginSDK) error { }, "required": []string{"id"}, }, + // 写终端:与会话缓冲区共享,须按序 + Serial: true, }, func(args map[string]interface{}) (interface{}, error) { return p.handleWrite(s, args) }) @@ -309,6 +313,8 @@ func (p *Plugin) Start(s *sdk.PluginSDK) error { }, "required": []string{"id"}, }, + // 读终端:与会话缓冲区共享,须按序 + Serial: true, }, func(args map[string]interface{}) (interface{}, error) { return p.handleRead(args) }) @@ -335,6 +341,8 @@ func (p *Plugin) Start(s *sdk.PluginSDK) error { }, "required": []string{"id"}, }, + // 改窗口尺寸:改会话状态 + Serial: true, }, func(args map[string]interface{}) (interface{}, error) { return p.handleResize(args) }) @@ -353,6 +361,8 @@ func (p *Plugin) Start(s *sdk.PluginSDK) error { }, "required": []string{"id"}, }, + // 关会话:改共享会话表 + Serial: true, }, func(args map[string]interface{}) (interface{}, error) { return p.handleClose(s, args) }) @@ -365,6 +375,8 @@ func (p *Plugin) Start(s *sdk.PluginSDK) error { "type": "object", "properties": map[string]interface{}{}, }, + // 已核实只读:只列举会话 + ParallelSafe: true, }, func(args map[string]interface{}) (interface{}, error) { return p.handleList() }) @@ -407,6 +419,8 @@ func (p *Plugin) Start(s *sdk.PluginSDK) error { }, "required": []string{"id"}, }, + // 订阅输出:注册监听者,改共享状态 + Serial: true, }, func(args map[string]interface{}) (interface{}, error) { return p.handleWatch(args) }) diff --git a/internal/plugins/ai_image/plugin.go b/internal/plugins/ai_image/plugin.go index 52024a3..2ba7841 100644 --- a/internal/plugins/ai_image/plugin.go +++ b/internal/plugins/ai_image/plugin.go @@ -112,6 +112,8 @@ func (p *Plugin) Start(s *sdk.PluginSDK) error { }, "required": []string{"prompt"}, }, + // 外部调用,插件内无共享可变状态 + ParallelSafe: true, }, p.handleGenerate) return nil diff --git a/internal/plugins/cfgmgr/plugin.go b/internal/plugins/cfgmgr/plugin.go index c3fc588..dacb2f7 100644 --- a/internal/plugins/cfgmgr/plugin.go +++ b/internal/plugins/cfgmgr/plugin.go @@ -36,6 +36,8 @@ func (p *Plugin) Start(s *sdk.PluginSDK) error { "type": "object", "properties": map[string]interface{}{}, }, + // 已核实只读:只列举插件名 + ParallelSafe: true, }, p.handleListPlugins(s)) s.RegisterTool("config_get", sdk.ToolDef{ @@ -49,6 +51,8 @@ func (p *Plugin) Start(s *sdk.PluginSDK) error { }, "required": []string{"scope", "key"}, }, + // 已核实只读:底层 ConfigRegistry 有 RWMutex,且本工具不改任何状态 + ParallelSafe: true, }, p.handleGet(s)) s.RegisterTool("config_set", sdk.ToolDef{ @@ -63,6 +67,8 @@ func (p *Plugin) Start(s *sdk.PluginSDK) error { }, "required": []string{"scope", "key", "value"}, }, + // 写配置:即使底层有锁也按声明序执行 + Serial: true, }, p.handleSet(s)) s.RegisterTool("config_list_keys", sdk.ToolDef{ @@ -76,6 +82,8 @@ func (p *Plugin) Start(s *sdk.PluginSDK) error { }, "required": []string{"scope"}, }, + // 已核实只读:只列举键值 + ParallelSafe: true, }, p.handleListKeys(s)) s.RegisterTool("config_get_defs", sdk.ToolDef{ @@ -89,6 +97,8 @@ func (p *Plugin) Start(s *sdk.PluginSDK) error { }, "required": []string{"scope"}, }, + // 已核实只读:返回配置项定义,不改状态 + ParallelSafe: true, }, p.handleGetDefs(s)) s.RegisterTool("config_dump", sdk.ToolDef{ @@ -98,6 +108,8 @@ func (p *Plugin) Start(s *sdk.PluginSDK) error { "type": "object", "properties": map[string]interface{}{}, }, + // 已核实只读:导出当前配置 + ParallelSafe: true, }, p.handleDump(s)) s.RegisterTool("config_batch_set", sdk.ToolDef{ @@ -122,6 +134,8 @@ func (p *Plugin) Start(s *sdk.PluginSDK) error { }, "required": []string{"items"}, }, + // 批量写配置:同上 + Serial: true, }, p.handleBatchSet(s)) log.Printf("[cfgmgr] started") diff --git a/internal/plugins/cmd/plugin.go b/internal/plugins/cmd/plugin.go index 42f2d0b..58613f3 100644 --- a/internal/plugins/cmd/plugin.go +++ b/internal/plugins/cmd/plugin.go @@ -129,6 +129,8 @@ func (p *Plugin) Start(s *sdk.PluginSDK) error { }, "required": []string{"command"}, }, + // 执行体只依赖入参;唯一共享 p.history 由 recordCmd 加 p.mu 保护 + ParallelSafe: true, }, func(args map[string]interface{}) (interface{}, error) { command, _ := args["command"].(string) if command == "" { diff --git a/internal/plugins/cmd/stress_test.go b/internal/plugins/cmd/stress_test.go new file mode 100644 index 0000000..df37378 --- /dev/null +++ b/internal/plugins/cmd/stress_test.go @@ -0,0 +1,224 @@ +package cmd + +import ( + "fmt" + "sync" + "sync/atomic" + "testing" + "time" + + sdk "gitcode.com/JianFeeeee/HomeAgent/internal/sdk" +) + +// 本文件是**真实工具**的并发压测。 +// +// 为什么需要它:阶段 2 的并发判据(core/parallelsched_test.go)用的是 +// **假设备** —— 它验证的是"内核会不会并发调度",但没有验证 +// **真实插件工具在真并发下是否安全**。而这正是 ParallelSafe 声明的风险面: +// 声明错一个工具,1000 个并发调用会同时打进去。 +// +// 判据全是**不变量**,不写性能阈值(阈值会随机器波动,变成"红/绿随运气" +// 的假信号)。 +// +// 跑法:go test ./internal/plugins/cmd/ -run TestStress -timeout 600s +// -short 时跳过。 + +// runOnce 直接调 cmd_run 的 handler,返回其 stdout。 +// +// ⚠️ cmd_run 的真实返回是 **map[string]interface{}**(含 status/stdout/ +// exit_code/command),**不是 string**。我第一版按 string 断言,导致 +// 1000 次全判"输出为空"—— 那是判据写错,不是工具串扰。 +// 既有测试(TestCmdRunEcho)同样用 json.Marshal 取值,与此一致。 +func runOnce(h sdk.ToolHandler, i int) (string, error) { + res, err := h(map[string]interface{}{ + "command": fmt.Sprintf("echo seq%d", i), + }) + if err != nil { + return "", err + } + m, ok := res.(map[string]interface{}) + if !ok { + return "", fmt.Errorf("cmd_run 返回类型不是 map:%T", res) + } + if e, hasErr := m["error"]; hasErr { + return "", fmt.Errorf("cmd_run 报错:%v", e) + } + s, _ := m["stdout"].(string) + return s, nil +} + +// TestStress_RealTool1000Concurrent 真实 cmd_run 1000 并发。 +// +// 判据: +// 1. 1000 次全部成功,**零错误**(真并发下最常见的失败是共享状态竞争) +// 2. 每次输出**各不相同**且与自己的入参对应 —— 若 handler 有共享 buffer +// 竞争,输出会串(这是"执行体不安全"最典型的症状) +// 3. 耗时不应随并发数线性恶化到不可用(只做宽松上界,不做精确基准) +// 4. -race 无竞态 +func TestStress_RealTool1000Concurrent(t *testing.T) { + if testing.Short() { + t.Skip("stress test; run with -run TestStress") + } + const n = 1000 + + // ⚠️ 用**既有**的 setupPlugin(plugin_test.go 里的真实装配), + // 不另造一套 —— 压测必须跑在真实注册路径上,否则测的是替身。 + _, tc, err := setupPlugin() + if err != nil { + t.Fatalf("装配插件失败: %v", err) + } + h, ok := tc.handlers["cmd_run"] + if !ok { + t.Fatal("取不到 cmd_run 的 handler") + } + // 顺带确认它真的声明了并发安全(否则 1000 并发会被内核整批串行, + // 本用例就测不到真并发) + if !tc.defs["cmd_run"].ParallelSafe { + t.Fatal("cmd_run 未声明 ParallelSafe —— 内核会整批串行,本压测失去意义") + } + + // 预热:首次调用会加载配置/建目录,不计入压测 + if _, err := runOnce(h, -1); err != nil { + t.Fatalf("预热失败: %v", err) + } + + // 每个 goroutine 只写自己那个下标(无共享变量),故无需加锁 —— + // 这也是检查项 map_write_in_goroutine 想确认的:按索引分槽写是所有权清晰。 + results := make([]string, n) + errs := make([]error, n) + var wg sync.WaitGroup + var inFlight, maxInFlight int32 + + start := time.Now() + for i := 0; i < n; i++ { + wg.Add(1) + go func(idx int) { + defer wg.Done() + // 记录并发峰值:证明确实并发了(否则本用例测不到并发) + cur := atomic.AddInt32(&inFlight, 1) + for { + old := atomic.LoadInt32(&maxInFlight) + if cur <= old || atomic.CompareAndSwapInt32(&maxInFlight, old, cur) { + break + } + } + results[idx], errs[idx] = runOnce(h, idx) + atomic.AddInt32(&inFlight, -1) + }(i) + } + wg.Wait() + elapsed := time.Since(start) + + // ① 零错误 + fails := 0 + for i, e := range errs { + if e != nil { + if fails < 3 { + t.Errorf("第 %d 次并发调用失败: %v", i, e) + } + fails++ + } + } + if fails > 0 { + t.Errorf("共 %d/%d 次失败", fails, n) + } + + // ② 输出与入参一一对应(无串扰) + mismatched := 0 + for i := 0; i < n; i++ { + want := fmt.Sprintf("seq%d", i) + if errs[i] != nil { + continue + } + if !containsStr(results[i], want) { + if mismatched < 3 { + t.Errorf("第 %d 次输出不含自己的标记 %q:%q(串扰?)", i, want, + truncate(results[i], 80)) + } + mismatched++ + } + } + if mismatched > 0 { + t.Errorf("共 %d 次输出与入参不对应(共享状态竞争)", mismatched) + } + + peak := atomic.LoadInt32(&maxInFlight) + t.Logf("1000 并发真实 cmd_run:耗时 %v,并发峰值 %d,失败 %d", elapsed, peak, fails) + if peak < 10 { + t.Errorf("并发峰值仅 %d —— 可能被串行化了,本用例测不到真并发", peak) + } +} + +// TestStress_RealToolSerialVsConcurrent 串行 vs 并发的耗时对比。 +// +// 只做**观察性**记录(不做通过判据):机器差异太大,阈值无意义。 +// 它的价值在于:如果并发比串行**慢很多**,说明 handler 内部有锁竞争 +// 或资源争抢,值得深挖。 +func TestStress_RealToolSerialVsConcurrent(t *testing.T) { + if testing.Short() { + t.Skip("stress test; run with -run TestStress") + } + const n = 200 + + // ⚠️ 用**既有**的 setupPlugin(plugin_test.go 里的真实装配), + // 不另造一套 —— 压测必须跑在真实注册路径上,否则测的是替身。 + _, tc, err := setupPlugin() + if err != nil { + t.Fatalf("装配插件失败: %v", err) + } + h, ok := tc.handlers["cmd_run"] + if !ok { + t.Fatal("取不到 cmd_run 的 handler") + } + // 顺带确认它真的声明了并发安全(否则 1000 并发会被内核整批串行, + // 本用例就测不到真并发) + if !tc.defs["cmd_run"].ParallelSafe { + t.Fatal("cmd_run 未声明 ParallelSafe —— 内核会整批串行,本压测失去意义") + } + if _, err := runOnce(h, -1); err != nil { + t.Fatalf("预热失败: %v", err) + } + + // 串行 + t0 := time.Now() + for i := 0; i < n; i++ { + if _, err := runOnce(h, i); err != nil { + t.Fatalf("串行第 %d 次失败: %v", i, err) + } + } + dSerial := time.Since(t0) + + // 并发(2 并发:轻并发,便于对比是否有争抢) + t1 := time.Now() + var wg sync.WaitGroup + for i := 0; i < n; i++ { + wg.Add(1) + go func(idx int) { + defer wg.Done() + if _, err := runOnce(h, idx); err != nil { + t.Errorf("并发第 %d 次失败: %v", idx, err) + } + }(i) + } + wg.Wait() + dConc := time.Since(t1) + + t.Logf("%d 次:串行 %v(%v/次),高并发 %v(%v/次)", + n, dSerial, dSerial/n, dConc, dConc/n) +} + +func containsStr(s, sub string) bool { + for i := 0; i+len(sub) <= len(s); i++ { + if s[i:i+len(sub)] == sub { + return true + } + } + return false +} + +func truncate(s string, n int) string { + if len(s) <= n { + return s + } + return s[:n] + "..." +} diff --git a/internal/plugins/healthcheck/plugin.go b/internal/plugins/healthcheck/plugin.go index c249fb7..b33748d 100644 --- a/internal/plugins/healthcheck/plugin.go +++ b/internal/plugins/healthcheck/plugin.go @@ -42,33 +42,33 @@ func init() { } type Plugin struct { - name string - mu sync.Mutex - reports []llmReport - sessionID string + name string + mu sync.Mutex + reports []llmReport + sessionID string selfToolNames map[string]bool - checkMu sync.Mutex + checkMu sync.Mutex - stopCh chan struct{} - perfData PerfData + stopCh chan struct{} + perfData PerfData - autoInterval time.Duration - llmTimeout time.Duration - llmMaxTurns int - llmMaxTokens int - perfHistory int + autoInterval time.Duration + llmTimeout time.Duration + llmMaxTurns int + llmMaxTokens int + perfHistory int } type PerfData struct { - LastCheck time.Time `json:"last_check"` - Checks []PerfCheckPoint `json:"checks"` + LastCheck time.Time `json:"last_check"` + Checks []PerfCheckPoint `json:"checks"` } type PerfCheckPoint struct { - Time time.Time `json:"time"` - Passed int `json:"passed"` - Failed int `json:"failed"` - Total int `json:"total"` - ElapsedMs int64 `json:"elapsed_ms"` + Time time.Time `json:"time"` + Passed int `json:"passed"` + Failed int `json:"failed"` + Total int `json:"total"` + ElapsedMs int64 `json:"elapsed_ms"` } func New(name string) *Plugin { @@ -170,6 +170,8 @@ func (p *Plugin) Start(s *sdk.PluginSDK) error { "plugin": map[string]interface{}{"type": "string", "description": "可选:指定只检查该插件的健康状态(列出插件工具并逐一测试),不填则检查全部插件"}, }, }, + // 共享 p.mu 写锁,且 healthcheck_report 会 append p.reports + Serial: true, }, func(args map[string]interface{}) (interface{}, error) { plugin, _ := args["plugin"].(string) return p.runFullCheck(s, plugin) @@ -183,6 +185,8 @@ func (p *Plugin) Start(s *sdk.PluginSDK) error { "type": "object", "properties": map[string]interface{}{}, }, + // 共享 p.mu 写锁,且 healthcheck_report 会 append p.reports + Serial: true, }, func(args map[string]interface{}) (interface{}, error) { return p.checkPlugins(s) }) @@ -195,6 +199,8 @@ func (p *Plugin) Start(s *sdk.PluginSDK) error { "type": "object", "properties": map[string]interface{}{}, }, + // 共享 p.mu 写锁,且 healthcheck_report 会 append p.reports + Serial: true, }, func(args map[string]interface{}) (interface{}, error) { return p.listAllTools(s) }) @@ -207,6 +213,8 @@ func (p *Plugin) Start(s *sdk.PluginSDK) error { "type": "object", "properties": map[string]interface{}{}, }, + // 共享 p.mu 写锁,且 healthcheck_report 会 append p.reports + Serial: true, }, func(args map[string]interface{}) (interface{}, error) { return p.checkMemory(s) }) @@ -224,6 +232,8 @@ func (p *Plugin) Start(s *sdk.PluginSDK) error { }, "required": []string{"tool_name", "status"}, }, + // 共享 p.mu 写锁,且 healthcheck_report 会 append p.reports + Serial: true, }, func(args map[string]interface{}) (interface{}, error) { toolName, _ := args["tool_name"].(string) status, _ := args["status"].(string) @@ -245,6 +255,8 @@ func (p *Plugin) Start(s *sdk.PluginSDK) error { "type": "object", "properties": map[string]interface{}{}, }, + // 共享 p.mu 写锁,且 healthcheck_report 会 append p.reports + Serial: true, }, func(args map[string]interface{}) (interface{}, error) { return s.Status().GetKernelStatus(), nil }) @@ -258,6 +270,8 @@ func (p *Plugin) Start(s *sdk.PluginSDK) error { "type": "object", "properties": map[string]interface{}{}, }, + // 共享 p.mu 写锁,且 healthcheck_report 会 append p.reports + Serial: true, }, func(args map[string]interface{}) (interface{}, error) { p.mu.Lock() defer p.mu.Unlock() @@ -268,8 +282,8 @@ func (p *Plugin) Start(s *sdk.PluginSDK) error { failed += c.Failed } return map[string]interface{}{ - "status": "ok", - "last_check": p.perfData.LastCheck, + "status": "ok", + "last_check": p.perfData.LastCheck, "total_checks": len(p.perfData.Checks), "total_passed": passed, "total_failed": failed, @@ -833,16 +847,16 @@ func (p *Plugin) collectToolDefsForLLM(s *sdk.PluginSDK, pluginFilter string) [] func isSafeReadonlyTool(name string) bool { // 明确只读的查询/列表类工具 readonlyExact := map[string]bool{ - "memory_recall": true, - "memory_introspect": true, - "doc_query": true, - "knowledge_search": true, - "knowledge_list": true, - "person_query": true, - "person_network": true, - "llm_list_sources": true, + "memory_recall": true, + "memory_introspect": true, + "doc_query": true, + "knowledge_search": true, + "knowledge_list": true, + "person_query": true, + "person_network": true, + "llm_list_sources": true, "output_list_channels": true, - "terminal_list": true, + "terminal_list": true, } if readonlyExact[name] { return true diff --git a/internal/plugins/localuse/plugin.go b/internal/plugins/localuse/plugin.go index ad59703..5fd628f 100644 --- a/internal/plugins/localuse/plugin.go +++ b/internal/plugins/localuse/plugin.go @@ -53,6 +53,8 @@ func (p *Plugin) Start(s *sdk.PluginSDK) error { "type": "object", "properties": map[string]interface{}{}, }, + // 只读观察:不改插件内共享状态 + ParallelSafe: true, }, p.handleScreensee) // ── camerasue ── @@ -70,6 +72,8 @@ func (p *Plugin) Start(s *sdk.PluginSDK) error { }, }, }, + // 只读观察:不改插件内共享状态 + ParallelSafe: true, }, p.handleCamerasue) // ── speakeruse ── @@ -87,6 +91,8 @@ func (p *Plugin) Start(s *sdk.PluginSDK) error { }, "required": []string{"text"}, }, + // 只读观察:不改插件内共享状态 + ParallelSafe: true, }, p.handleSpeakeruse) // ── screensue ── @@ -120,6 +126,8 @@ func (p *Plugin) Start(s *sdk.PluginSDK) error { "type": "object", "properties": map[string]interface{}{}, }, + // 只读观察:不改插件内共享状态 + ParallelSafe: true, }, p.handleClipboardsee) // ── clipboardsue ── @@ -160,6 +168,8 @@ func (p *Plugin) Start(s *sdk.PluginSDK) error { }, "required": []interface{}{"action"}, }, + // 只读观察:不改插件内共享状态 + ParallelSafe: true, }, p.handleComputeruse) return nil diff --git a/internal/plugins/pluginmgr/plugin.go b/internal/plugins/pluginmgr/plugin.go index 8520660..fd71bcc 100644 --- a/internal/plugins/pluginmgr/plugin.go +++ b/internal/plugins/pluginmgr/plugin.go @@ -71,9 +71,10 @@ var downloadClient = &http.Client{ // // ★ 曾经这里是一个**包级可变全局** `var HTTPAddr`,且 Start() 会把 settings 读到的值 // **反写**回该全局。两个真实后果: -// 1. 多实例互相污染——测试并行起两个 Registry,后启动的实例会把地址写进全局, -// 先启动那个的 startHTTPServer 读到的是别人的地址(实测与生产 homed 抢 9876); -// 2. 全局读写在并发下没有同步,属数据竞态。 +// 1. 多实例互相污染——测试并行起两个 Registry,后启动的实例会把地址写进全局, +// 先启动那个的 startHTTPServer 读到的是别人的地址(实测与生产 homed 抢 9876); +// 2. 全局读写在并发下没有同步,属数据竞态。 +// // 现在改为实例字段 p.httpAddr(默认值走本常量),不再有可被任意代码改写的包级状态。 const defaultHTTPAddr = "127.0.0.1:9876" @@ -168,6 +169,8 @@ func (p *Plugin) registerTools(s *sdk.PluginSDK) { }, }, }, + // 装插件:改磁盘与运行态,可能半安装 + Serial: true, }, func(args map[string]interface{}) (interface{}, error) { overwrite, _ := args["overwrite"].(bool) // path 优先:它对应"agent 自己构建出产物再装"的场景(plugindev_build → plugin_install)。 @@ -192,6 +195,8 @@ func (p *Plugin) registerTools(s *sdk.PluginSDK) { "type": "object", "properties": map[string]interface{}{}, }, + // 已核实只读:只列举已装插件 + ParallelSafe: true, }, func(args map[string]interface{}) (interface{}, error) { return p.listPlugins() }) @@ -209,6 +214,8 @@ func (p *Plugin) registerTools(s *sdk.PluginSDK) { }, }, }, + // 已核实只读:查运行状态 + ParallelSafe: true, }, func(args map[string]interface{}) (interface{}, error) { name, _ := args["name"].(string) return p.pluginStatus(name) @@ -228,6 +235,8 @@ func (p *Plugin) registerTools(s *sdk.PluginSDK) { }, "required": []string{"name"}, }, + // 重启插件:改运行态 + Serial: true, }, func(args map[string]interface{}) (interface{}, error) { name, _ := args["name"].(string) if name == "" { @@ -249,6 +258,8 @@ func (p *Plugin) registerTools(s *sdk.PluginSDK) { }, "required": []string{"name"}, }, + // 卸插件:改磁盘与运行态 + Serial: true, }, func(args map[string]interface{}) (interface{}, error) { name, _ := args["name"].(string) if name == "" { @@ -270,6 +281,8 @@ func (p *Plugin) registerTools(s *sdk.PluginSDK) { }, "required": []string{"name"}, }, + // 已核实只读:查单个插件详情 + ParallelSafe: true, }, func(args map[string]interface{}) (interface{}, error) { name, _ := args["name"].(string) if name == "" { diff --git a/internal/plugins/seq/stress_extreme_test.go b/internal/plugins/seq/stress_extreme_test.go new file mode 100644 index 0000000..388126c --- /dev/null +++ b/internal/plugins/seq/stress_extreme_test.go @@ -0,0 +1,242 @@ +package seq + +import ( + "fmt" + "os" + "path/filepath" + "strings" + "sync" + "testing" + "time" +) + +// 本文件是**极端规模**压测:1000 条序列 × 每条 1000 个组内 toolcall。 +// +// 与 stress_test.go 的分工:那边压的是「单条序列变大」,这边压的是 +// **大量序列各自很大** —— 存储、解析、执行、合并四个环节的**累积**成本, +// 以及并发下的内存与正确性。 +// +// ⚠️ 为什么压的是**组内 1000 并发**而不是「1000 个 seq 之间并发」: +// seq_* 工具刻意不声明 ParallelSafe(seq_run 会执行一串工具、含写操作, +// 并发会污染执行序列与变量表),内核 batchRunnable 因此整批串行; +// 另有 maxCallDepth=4 的结构上界。所以「1000 个 seq 并行」在当前设计下 +// **不会发生**,压它等于压一条走不到的路径。 +// 而组内并发是真实存在的:组一旦声明 parallel 且全部工具 ParallelSafe, +// 1000 个 toolcall 会真的同时在跑 —— 那才是成本所在。 +// +// 文件协议(按要求):每条序列**由文件创建**(走 seq_create 的 file 路径, +// 这也是长序列的推荐用法),压测结束**删除**。目录在 t.TempDir() 下, +// 不碰生产数据目录。 +// +// 跑法: +// go test ./internal/plugins/seq/ -run TestStressExtreme -timeout 1800s +// -short 时跳过。 + +// extRunner 是组内 1000 toolcall 的执行面:记录调用、按声明序返回。 +type extRunner struct { + mu sync.Mutex + called int + // perCall 记录每个工具应返回的值(按其序号),用于验证合并顺序 + gate map[string]func() +} + +func newExtRunner() *extRunner { return &extRunner{gate: map[string]func(){}} } + +func (r *extRunner) call(name string, _ map[string]interface{}) (string, error) { + if g, ok := r.gate[name]; ok && g != nil { + g() + } + r.mu.Lock() + r.called++ + r.mu.Unlock() + return "v:" + name, nil +} + +// parallelSafe 全部为真 ⇒ 组内可并发(这正是本压测要测的路径)。 +func (r *extRunner) parallelSafe(string) bool { return true } + +// exists:压测里的工具都是真实存在的(mkSeqFile 生成的 t%05d)。 +// ⚠️ 必须返回 true —— 否则 runGroup 的 missing 预检会把 1000 个 toolcall +// 全判为"不存在"并按 missing=fail 整组跳过,压测就变成测"跳过"了。 +func (r *extRunner) exists(string) bool { return true } + +func (r *extRunner) count() int { + r.mu.Lock() + defer r.mu.Unlock() + return r.called +} + +// mkSeqFile 生成一条序列的 JSON 文本:nGroups 组,每组 nTools 个 toolcall。 +// 槽名用 o%05d(每组内唯一,避免标量槽同名多写被静态校验拦下)。 +func mkSeqFile(name string, nGroups, nTools int) []byte { + groups := make([]interface{}, 0, nGroups) + for g := 0; g < nGroups; g++ { + out := make(map[string]string, nTools) + parts := make([]string, 0, nTools) + for i := 0; i < nTools; i++ { + slot := fmt.Sprintf("o%05d", i) + out[slot] = "string" + parts = append(parts, fmt.Sprintf( + `{"tool":"t%05d","args":{"g":%d},"as":%q}`, i, g, slot)) + } + groups = append(groups, map[string]interface{}{ + "name": fmt.Sprintf("g%05d", g), + "in": map[string]string{}, + "out": out, + "tools": strings.Join(parts, " ; ") + " ;", + }) + } + return mustMarshal(map[string]interface{}{ + "name": name, + "description": "极端规模压测", + "groups": groups, + }) +} + +// TestStressExtreme_ThousandSeqs 1000 条序列,每条 1 组 × 1000 toolcall。 +// +// 协议:文件创建(seq_create file=…)→ 列出 → 执行 → 删除, +// 全部落在 t.TempDir() 下。 +// +// 判据(都是不变量,不是性能阈值 —— 性能会随机器波动,写死阈值只会 +// 变成"红/绿随运气"的假信号): +// 1. 1000 条全部创建成功,无一条被静默丢弃 +// 2. 1000 条全部可列出 +// 3. 全部执行成功,**调用总数**精确 = 1000 × 1000 +// 4. 单组内 1000 个槽的合并顺序正确(并发下按声明序,不按完成序) +// 5. 全部删除成功,目录里不留残留 +func TestStressExtreme_ThousandSeqs(t *testing.T) { + if testing.Short() { + t.Skip("extreme stress test; run with -run TestStressExtreme") + } + const ( + nSeqs = 1000 + nTools = 1000 + nGroups = 1 // 每条 1 组,组内 1000 toolcall ⇒ 并发度 1000 + ) + + dir := t.TempDir() + store := NewStore(dir) + r := newExtRunner() + p := newE2EPlugin(t, r) + p.store = store // 用我们控制的目录,确保结束能整体删除 + + // ---- 阶段 1:文件创建 ---- + t0 := time.Now() + created := 0 + for i := 0; i < nSeqs; i++ { + name := fmt.Sprintf("x%04d", i) + fp := filepath.Join(dir, name+".src.json") + if err := os.WriteFile(fp, mkSeqFile(name, nGroups, nTools), 0644); err != nil { + t.Fatalf("写序列源文件 %s 失败: %v", name, err) + } + if _, err := p.dispatch("seq_create", map[string]interface{}{ + "name": name, "file": fp, + }); err != nil { + t.Fatalf("seq_create(%s) 失败: %v", name, err) + } + created++ + } + dCreate := time.Since(t0) + t.Logf("创建 %d 条(每条 %d 个组内 toolcall,源文件在 %s):%v", nSeqs, nTools, dir, dCreate) + + if created != nSeqs { + t.Fatalf("应创建 %d 条,实际 %d", nSeqs, created) + } + + // ---- 阶段 2:列出 ---- + t1 := time.Now() + names := store.List() + dList := time.Since(t1) + if len(names) != nSeqs { + t.Errorf("列出 %d 条,期望 %d", len(names), nSeqs) + } + t.Logf("列出 %d 条:%v", len(names), dList) + + // ---- 阶段 3:执行(全部)---- + t2 := time.Now() + ranOK := 0 + for i := 0; i < nSeqs; i++ { + name := fmt.Sprintf("x%04d", i) + if _, err := p.dispatch("seq_run", map[string]interface{}{"name": name}); err != nil { + t.Fatalf("seq_run(%s) 失败: %v", name, err) + } + ranOK++ + } + dRun := time.Since(t2) + + wantCalls := nSeqs * nGroups * nTools + if got := r.count(); got != wantCalls { + t.Errorf("工具调用总数 = %d,期望 %d(每条 %d 个)", got, wantCalls, nGroups*nTools) + } + if ranOK != nSeqs { + t.Errorf("执行成功 %d 条,期望 %d", ranOK, nSeqs) + } + t.Logf("执行 %d 条(累计 %d 次 toolcall,其中组内并发度 %d):%v", + nSeqs, wantCalls, nTools, dRun) + + // ---- 阶段 4:单组 1000 槽的合并顺序(并发不变式)---- + // 直接对一条序列的组做细粒度校验:槽 o00000..o00999 必须各得自己的值。 + // 这一条必须在**并发**下成立 —— 若按完成顺序合并,这里必然错位。 + verifyMergedOrder(t, store, names[0], nTools) + + // ---- 阶段 5:删除 ---- + t3 := time.Now() + deleted := 0 + for i := 0; i < nSeqs; i++ { + if _, err := p.dispatch("seq_delete", map[string]interface{}{ + "name": fmt.Sprintf("x%04d", i), + }); err != nil { + t.Fatalf("seq_delete 失败: %v", err) + } + deleted++ + } + dDel := time.Since(t3) + if deleted != nSeqs { + t.Errorf("删除 %d 条,期望 %d", deleted, nSeqs) + } + if left := store.List(); len(left) != 0 { + t.Errorf("删除后仍残留 %d 条序列", len(left)) + } + t.Logf("删除 %d 条:%v", deleted, dDel) + + // 清理源文件,确认目录可整体移除(验证没把数据写到别处) + for i := 0; i < nSeqs; i++ { + if err := os.Remove(filepath.Join(dir, fmt.Sprintf("x%04d.src.json", i))); err != nil { + t.Fatalf("清理源文件失败: %v", err) + } + } + if _, err := os.Stat(dir); err != nil { + t.Errorf("序列目录状态异常: %v", err) + } +} + +// verifyMergedOrder 校验一条序列的组在并发执行后,各槽内容与声明序一致。 +func verifyMergedOrder(t *testing.T, store *Store, name string, nTools int) { + t.Helper() + seq, err := store.Load(name) + if err != nil { + t.Fatalf("Load(%s): %v", name, err) + } + // 造一个交错延迟的执行面:序号越大越先完成 ⇒ 完成序与声明序相反。 + // 若合并按完成序,这里会全盘错位。 + r := newExtRunner() + for i := 0; i < nTools; i++ { + tool := fmt.Sprintf("t%05d", i) + delay := time.Duration(nTools-i) * 20 * time.Microsecond + r.gate[tool] = func() { time.Sleep(delay) } + } + res, err := execGroup(seq.Groups[0], map[string]interface{}{}, r) + if err != nil { + t.Fatalf("execGroup 失败: %v", err) + } + for i := 0; i < nTools; i++ { + slot := fmt.Sprintf("o%05d", i) + want := fmt.Sprintf("v:t%05d", i) + if got, _ := res.Slots[slot].(string); got != want { + t.Errorf("槽 %s = %q,期望 %q —— 1000 并发下合并顺序错位", slot, got, want) + return + } + } + t.Logf("单组 %d 槽在交错延迟下合并顺序全部正确", nTools) +} diff --git a/internal/plugins/timer/plugin.go b/internal/plugins/timer/plugin.go index 8fbc78e..3726a4b 100644 --- a/internal/plugins/timer/plugin.go +++ b/internal/plugins/timer/plugin.go @@ -82,6 +82,8 @@ func (p *Plugin) Start(s *sdk.PluginSDK) error { }, "required": []string{"duration", "message"}, }, + // 写定时器:改 p.timers(持 p.mu),按序更可预期 + Serial: true, }, func(args map[string]interface{}) (interface{}, error) { durStr, _ := args["duration"].(string) message, _ := args["message"].(string) diff --git a/third_party/homeagent-sdk/sdk/plugin.go b/third_party/homeagent-sdk/sdk/plugin.go index f4f01ea..544df5f 100644 --- a/third_party/homeagent-sdk/sdk/plugin.go +++ b/third_party/homeagent-sdk/sdk/plugin.go @@ -298,6 +298,19 @@ type ToolDef struct { // · 不与同批其它工具争抢同一资源(SQLite 写、设备、同一输出通道) // · 执行顺序无关(顺序敏感的工具应留 false,由内核保序) ParallelSafe bool `json:"parallel_safe,omitempty"` + // Serial 声明本工具**必须**串行 —— ParallelSafe 的反向标记。 + // + // 为什么需要它:ParallelSafe 的零值 false 已经表达"安全/串行", + // 插件无法区分"我没想过"和"我确认过必须串行"。一旦工具作者需要 + // 把"这里**故意**串行,是有原因的"写进代码(而不只是没填), + // 这个区分就是必需的 —— 否则只能靠命名约定传递意图。 + // + // 适用场景:读操作但有隐含顺序约束(终端 read/resize 这类共享会话 + // 状态)、或写操作虽已加锁但需要串行以获得可预测的交错顺序。 + // + // 判据优先级:**Serial 胜出**。显式声明"必须串行"不允许被 + // ParallelSafe 或任何默认值覆盖。 + Serial bool `json:"serial,omitempty"` } // IOInjector provides methods for injecting input and interrupts into the agent pipeline.