feat(parallel): 并发安全改为声明式,并审计标注 37 个工具

把"能不能并发"从内核硬编码名单改成**工具自己的声明项**,形态照 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() 缺陷
This commit is contained in:
JianFeeeee
2026-09-27 15:15:57 +08:00
parent d36b4f6d34
commit 43bc25525d
16 changed files with 907 additions and 52 deletions

View File

@ -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
}

View File

@ -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 必须胜出")
}
}

View File

@ -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) }

View File

@ -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{}{

View File

@ -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 {

View File

@ -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)
})

View File

@ -112,6 +112,8 @@ func (p *Plugin) Start(s *sdk.PluginSDK) error {
},
"required": []string{"prompt"},
},
// 外部调用,插件内无共享可变状态
ParallelSafe: true,
}, p.handleGenerate)
return nil

View File

@ -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")

View File

@ -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 == "" {

View File

@ -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] + "..."
}

View File

@ -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

View File

@ -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

View File

@ -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 == "" {

View File

@ -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)
}

View File

@ -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)

View File

@ -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.