mirror of
https://gitcode.com/JianFeeeee/HomeAgent.git
synced 2026-09-29 22:12:52 +00:00
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:
242
internal/plugins/seq/stress_extreme_test.go
Normal file
242
internal/plugins/seq/stress_extreme_test.go
Normal 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)
|
||||
}
|
||||
Reference in New Issue
Block a user