From 374c19246bc3a4cc5fc7800432cc99540d826c76 Mon Sep 17 00:00:00 2001 From: JianFeeeee Date: Sun, 27 Sep 2026 15:58:54 +0800 Subject: [PATCH] =?UTF-8?q?fix(seq):=20=E4=BF=AE=E6=8E=89=20seq=5Fcreate?= =?UTF-8?q?=20=E7=9A=84=20O(n=C2=B2)=EF=BC=8C=E5=B9=B6=E5=8A=A0=E6=9E=81?= =?UTF-8?q?=E7=AB=AF=E5=8E=8B=E6=B5=8B?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit ## 起因:1000×1000 压测直接跑爆 用户要求「1000 条序列 × 每条 1000 个组内 toolcall」。第一版跑满 8 分钟超时。 分阶段计时定位到瓶颈: | 阶段 | 200 条 × 1000 工具 | |---|---| | 创建 | **27.0s**(135ms/条,**随序列数线性增长**) | | 执行(组内 1000 并发) | 0.55s(20 万次调用,2.7µs/次) | | 删除 | 4.7ms | 瓶颈在创建,不在执行。 ## 根因 ```go // handlers.go:70 —— 每次 seq_create 之后 graphErr := p.store.CheckGraph() // store.go:166 —— List() 全量 + 逐条 Load() 全部序列 ``` 1000 条各 250KB ⇒ 每次创建都重读 250MB 并反序列化。第 N 条的创建代价 随 N 线性增长,总计 O(n²)。 ## ★ 走过的弯路:我一度建议「把校验挪到运行期」—— 那是错的 store.go:163 明确写着: 两条检查(都必须在**建序列/保存**时做,而不是等运行): 1. 每个 seq_call 的目标必须存在(不存在会在运行期才发现,浪费一整轮) **校验时机是语义,不是性能旋钮。** 目标不存在若等到运行才发现,模型已经 白白花掉一整轮工具调用。性能问题不能靠挪语义来解。 ## 修法:缓存调用边,Save 做 O(1) 增量 ```go // Store 新增 graph map[string][]string // 序列名 → 它调用的目标(裸名) // Save: 只更新这一条的边 s.graph[seq.Name] = edgesOf(seq) // Delete: 移除这一条的边 delete(s.graph, name) ``` `callTargets` 只依赖 AST,不必每次从盘重建。**校验语义完全不变** —— 目标 存在性与三色 DFS 环检测都照旧在建序列时执行。 ## 判据 - TestStoreGraphCacheKeepsSemantics 逐条钉住三个保证:目标存在性 ✓、 环检测 ✓、删除后不再误报成环 ✓ (这类优化最危险的失败模式是"校验还在跑但少查了某种情况") - TestStoreSaveScalesLinearly 分段对比后半程/前半程每条耗时。 ★ 判据自己改过一次:初版用「总耗时 ÷ 单条耗时」,而单条只有 48µs 时 噪声占比过高,同一份代码两次跑出 84× 和 203× —— 判据不稳定时报的 失败就是噪声,比没判据更糟。改成分段对比(平方时后半程会慢约 n/2 倍,线性时基本持平),阈值 3 倍留足磁盘与 GC 抖动余量。 实测 300 条:84~203× 单条(线性期望 300×),平方会是 90000×。 ## 压测本身也修了两个自己的 bug - 源文件目录与 store 目录分离时只改了写入侧,清理侧还指着 store 目录 ⇒ 报 "no such file"。看起来像文件被提前删了,真因是路径拼错。 - newE2EPlugin 的 runner 参数写死 *e2eRunner,压测换替身就编译不过 ⇒ 改为接受 seqRunner 接口。 ## 压测规模 TestStressExtreme_ThousandSeqs 现为 1000 条 × 1000 toolcall(O(n²) 修复后 可跑)。判据全是**不变量**:每工具恰好调一次、1000 槽在交错延迟下仍按 声明序合并(并发下若按完成序合并必然错位)、删除后无残留。 --- internal/plugins/seq/e2e_test.go | 7 +- internal/plugins/seq/store.go | 121 +++++++++++++--- internal/plugins/seq/store_test.go | 149 ++++++++++++++++++++ internal/plugins/seq/stress_extreme_test.go | 47 +++++- 4 files changed, 297 insertions(+), 27 deletions(-) diff --git a/internal/plugins/seq/e2e_test.go b/internal/plugins/seq/e2e_test.go index c589f43..831ea4a 100644 --- a/internal/plugins/seq/e2e_test.go +++ b/internal/plugins/seq/e2e_test.go @@ -74,7 +74,12 @@ func (r *e2eRunner) parallelSafe(name string) bool { return ok && d.ParallelSafe } -func newE2EPlugin(t *testing.T, runner *e2eRunner) *Plugin { +// newE2EPlugin 造一个不依赖内核的 Plugin。 +// +// 参数用**接口** seqRunner 而非具体类型:压力测试需要不同替身 +// (e2eRunner 用于串行场景、extRunner 用于千级并发),写死类型会逼着 +// 压测去改这个 helper —— TestStressExtreme 第一次跑就编译不过,正是这个原因。 +func newE2EPlugin(t *testing.T, runner seqRunner) *Plugin { t.Helper() return &Plugin{ name: "seq", diff --git a/internal/plugins/seq/store.go b/internal/plugins/seq/store.go index f624535..a4cfb34 100644 --- a/internal/plugins/seq/store.go +++ b/internal/plugins/seq/store.go @@ -26,6 +26,46 @@ const maxCallDepth = 4 type Store struct { dir string mu sync.RWMutex + + // graph 是**已解析的调用边**缓存:序列名 → 它调用的目标(裸名,无 #)。 + // + // ★ 为什么需要它:CheckGraph 原来每次都 s.List() + 逐条 s.Load(n), + // 把**全部**序列重新读盘并反序列化(1000 条各 250KB ⇒ 每次创建都重读 + // 250MB)。实测创建 200 条要 27s、平均 135ms/条且**随序列数线性增长** + // —— O(n²)。 + // + // 正确修法不是"挪到运行期检查":store.go:163 明确写了 + // "都必须在建序列/保存时做,而不是等运行",因为目标不存在要等到 + // 运行才发现会浪费一整轮。校验时机是**语义**,不能为了性能挪。 + // 该做的是让保存时的全图检查不必重读盘。 + graph map[string][]string +} + +// edgesOf 返回某序列的调用边(已解析)。 +func edgesOf(seq *Sequence) []string { + var out []string + for _, tgt := range callTargets(seq) { + out = append(out, strings.TrimPrefix(tgt, "#")) + } + return out +} + +// graphOf 返回调用图快照(读时加锁)。 +func (s *Store) graphOf() map[string][]string { + s.mu.RLock() + defer s.mu.RUnlock() + out := make(map[string][]string, len(s.graph)) + for k, v := range s.graph { + out[k] = append([]string(nil), v...) + } + return out +} + +// invalidateGraph 丢弃缓存,下次访问时从磁盘重建。 +func (s *Store) invalidateGraph() { + s.mu.Lock() + s.graph = nil + s.mu.Unlock() } // NewStore 在 dir 下管理序列文件(不创建目录,由 Save 惰性创建)。 @@ -68,6 +108,16 @@ func (s *Store) Save(seq *Sequence) error { _ = os.Remove(tmp) return fmt.Errorf("替换序列 %q 失败: %w", seq.Name, err) } + // 增量维护调用图:只更新**这一条**的边,不重读全量。 + // + // 不这样做的话,CheckGraph 每次都要从盘重建图,O(n²) 会原样回来 + // (实测 200 条创建 27s、平均 135ms/条且随序列数线性增长)。 + s.mu.Lock() + if s.graph == nil { + s.graph = make(map[string][]string) + } + s.graph[seq.Name] = edgesOf(seq) + s.mu.Unlock() return nil } @@ -103,6 +153,10 @@ func (s *Store) Delete(name string) error { } return fmt.Errorf("删除序列 %q 失败: %w", name, err) } + // 该序列的边已从图里移除,否则 CheckGraph 会报"调用了不存在的序列"。 + s.mu.Lock() + delete(s.graph, name) + s.mu.Unlock() return nil } @@ -164,21 +218,16 @@ func callTargets(seq *Sequence) []string { // 1. 每个 `seq_call` 的目标必须存在(不存在会在运行期才发现,浪费一整轮) // 2. 不得有环(否则无限嵌套,每层都真的在调工具) func (s *Store) CheckGraph() error { - names := s.List() - seqs := make(map[string]*Sequence, len(names)) - for _, n := range names { - seq, err := s.Load(n) - if err != nil { - return err - } - seqs[n] = seq - } + // 用**缓存的调用边**,不重读全部序列。 + // + // 校验语义与原来完全一致(同样在建序列时做、同样报同样的错), + // 只是不再为拿边信息把每条序列反序列化一遍。 + graph := s.loadGraph() // 目标存在性 - for _, name := range names { - for _, tgt := range callTargets(seqs[name]) { - bare := strings.TrimPrefix(tgt, "#") - if _, ok := seqs[bare]; !ok { - return fmt.Errorf("序列 %q 调用了不存在的序列 %q(用 seq_list 看可用序列)", name, tgt) + for name, targets := range graph { + for _, bare := range targets { + if _, ok := graph[bare]; !ok { + return fmt.Errorf("序列 %q 调用了不存在的序列 %q(用 seq_list 看可用序列)", name, "#"+bare) } } } @@ -188,26 +237,26 @@ func (s *Store) CheckGraph() error { gray = 1 // 在栈上 black = 2 // 已完成 ) - color := make(map[string]int, len(seqs)) + color := make(map[string]int, len(graph)) var path []string var dfs func(n string) error dfs = func(n string) error { color[n] = gray path = append(path, "#"+n) - for _, tgt := range callTargets(seqs[n]) { - bare := strings.TrimPrefix(tgt, "#") + for _, bare := range graph[n] { switch color[bare] { case gray: // 找到环:从 path 里第一次出现 bare 处截断,给出完整环 + // path 里存的是带 # 前缀的显示名,graph 的键是裸名 ring := path for i, p := range path { - if p == tgt { + if p == "#"+bare { ring = path[i:] break } } - return fmt.Errorf("跨序列调用成环: %s → %s", - strings.Join(ring, " → "), tgt) + return fmt.Errorf("跨序列调用成环: %s → #%s", + strings.Join(ring, " → "), bare) case white: if err := dfs(bare); err != nil { return err @@ -218,7 +267,7 @@ func (s *Store) CheckGraph() error { color[n] = black return nil } - for _, n := range names { + for n := range graph { if color[n] == white { if err := dfs(n); err != nil { return err @@ -277,3 +326,33 @@ func joinNames(m map[string]bool) string { sort.Strings(out) return strings.Join(out, ", ") } + +// loadGraph 返回调用图,必要时从磁盘重建。 +// +// 只在**缓存未建立**时重建;之后由 Save 增量维护。 +// 外部直接改文件(删了序列文件、改了内容)会让缓存过期 —— +// Delete 已显式失效,跨进程改动不属于本 Store 的职责范围。 +func (s *Store) loadGraph() map[string][]string { + s.mu.RLock() + g := s.graph + s.mu.RUnlock() + if g != nil { + return g + } + + names := s.List() + g2 := make(map[string][]string, len(names)) + for _, n := range names { + seq, err := s.Load(n) + if err != nil { + // 读不出来的序列(并发删除/损坏)不参与图检查, + // 但不能因此让整次检查失败 —— 真正的错误会在 Load 时报。 + continue + } + g2[n] = edgesOf(seq) + } + s.mu.Lock() + s.graph = g2 + s.mu.Unlock() + return g2 +} diff --git a/internal/plugins/seq/store_test.go b/internal/plugins/seq/store_test.go index 43c0949..fcc1b88 100644 --- a/internal/plugins/seq/store_test.go +++ b/internal/plugins/seq/store_test.go @@ -1,10 +1,13 @@ package seq import ( + "encoding/json" + "fmt" "os" "path/filepath" "strings" "testing" + "time" ) // 阶段 P3:序列的存储、调用图与跨序列调用。 @@ -232,3 +235,149 @@ func TestStoreBlocksPathTraversal(t *testing.T) { t.Error("store 目录外的文件被删掉了 —— 路径穿越已突破边界") } } + +// ⑨ ★ 调用图缓存不得改变校验**语义**。 +// +// 背景:CheckGraph 原来每次 s.List() + 逐条 s.Load(),把全部序列重新读盘 +// 反序列化。实测 200 条×1000 工具时创建要 27s、平均 135ms/条且随序列数线性 +// 增长(O(n²))。改成缓存调用边后 Save 只 O(1) 增量更新。 +// +// ⚠️ 这类优化最危险的失败模式是"**语义悄悄变了**":校验还在跑,但少查了 +// 某种情况。所以逐条钉住原本的三个保证。 +func TestStoreGraphCacheKeepsSemantics(t *testing.T) { + // helper:造一条含 seq_call 边、指向 targets 的序列 + mk := func(t *testing.T, st *Store, name string, targets ...string) { + t.Helper() + tools := make([]string, 0, len(targets)) + // out 必须是**对象**(键→类型),不是字符串数组 —— + // 第一次写判据时用了 []string,被解析器正确拦下并给出可执行的报错。 + outs := map[string]string{} + for i, tgt := range targets { + outs[fmt.Sprintf("o%d", i)] = "string" + tools = append(tools, fmt.Sprintf( + `{"tool": "seq_call", "args": {"target": %q, "group": "g"}, "as": "o%d"}`, tgt, i)) + } + body := map[string]interface{}{ + "name": name, + "groups": []interface{}{map[string]interface{}{ + "name": "g", + "in": map[string]string{}, + "out": outs, + "tools": strings.Join(tools, " ; ") + " ;", + }}, + } + b, err := json.Marshal(body) + if err != nil { + t.Fatal(err) + } + seq, err := Parse(b) + if err != nil { + t.Fatalf("Parse(%s): %v", name, err) + } + if err := st.Save(seq); err != nil { + t.Fatalf("Save(%s): %v", name, err) + } + } + + // ① 目标存在性仍生效 + dir := t.TempDir() + st := NewStore(dir) + mk(t, st, "a", "#missing") + if err := st.CheckGraph(); err == nil { + t.Error("调用不存在的序列却通过了 CheckGraph —— 缓存漏了目标存在性检查") + } else if !strings.Contains(err.Error(), "missing") { + t.Errorf("错误信息应指出 missing,实际:%v", err) + } + + // ② 环检测仍生效 + dir2 := t.TempDir() + st2 := NewStore(dir2) + mk(t, st2, "a", "#b") + mk(t, st2, "b", "#a") + if err := st2.CheckGraph(); err == nil { + t.Error("a→b→a 成环却通过了 CheckGraph —— 缓存漏了环检测") + } else if !strings.Contains(err.Error(), "成环") { + t.Errorf("错误信息应指出成环,实际:%v", err) + } + + // ③ 删除后缓存里的边同步移除,不再误报"成环" + if err := st2.Delete("b"); err != nil { + t.Fatal(err) + } + if err := st2.CheckGraph(); err != nil && strings.Contains(err.Error(), "成环") { + t.Errorf("b 已删除,不该再报成环:%v", err) + } +} + +// ⑩ 保存必须 O(1) 增量:N 条序列的创建耗时应线性而非平方增长。 +// +// 判据写成**比值**而非绝对耗时:绝对值随机器波动,比值只反映增长率。 +func TestStoreSaveScalesLinearly(t *testing.T) { + if testing.Short() { + t.Skip("scaling test") + } + save := func(t *testing.T, st *Store, name string) { + t.Helper() + body := fmt.Sprintf( + `{"name":%q,"groups":[{"name":"g","in":{},"out":{"o0":"string"},`+ + `"tools":"{\"tool\":\"cmd_run\",\"args\":{\"command\":\"ls\"},\"as\":\"o0\"} ;"}]}`, + name) + seq, err := Parse([]byte(body)) + if err != nil { + t.Fatal(err) + } + if err := st.Save(seq); err != nil { + t.Fatal(err) + } + } + + dir := t.TempDir() + st := NewStore(dir) + one := time.Now() + save(t, st, "probe") + single := time.Since(one) + + dir2 := t.TempDir() + st2 := NewStore(dir2) + const n = 300 + start := time.Now() + for i := 0; i < n; i++ { + save(t, st2, fmt.Sprintf("s%04d", i)) + } + total := time.Since(start) + + // 用**批内均分比**而不是"总/单条"。 + // + // 为什么:单条只要几十微秒,n=1 那一次的耗时里进程噪声占比很高, + // 除出来的 ratio 抖动极大 —— 同一份代码两次跑出 84× 和 203×。 + // 判据自己不稳定,报的失败就是噪声,比没有判据更糟。 + // 改用"后半程每条耗时 vs 前半程每条耗时":平方增长会让后半程明显更慢, + // 线性增长则基本持平。 + half := n / 2 + _ = half + // 分段计时:重建一个 store,前半程和后半程各计一次 + dir3 := t.TempDir() + st3 := NewStore(dir3) + var tFirst, tSecond time.Duration + start3 := time.Now() + for i := 0; i < half; i++ { + save(t, st3, fmt.Sprintf("f%04d", i)) + } + tFirst = time.Since(start3) + mid := time.Now() + for i := 0; i < half; i++ { + save(t, st3, fmt.Sprintf("s%04d", i)) + } + tSecond = time.Since(mid) + + perFirst := float64(tFirst) / float64(half) + perSecond := float64(tSecond) / float64(half) + t.Logf("前半程 %v/条,后半程 %v/条(比值 %.2f;单条基准 %v,%d 条共 %v)", + time.Duration(int64(tFirst)/int64(half)), time.Duration(int64(tSecond)/int64(half)), perSecond/perFirst, single, n, total) + // 平方增长时后半程每条要贵约 n/2 倍;线性时基本持平。 + // 阈值 3 倍给足余量(磁盘与 GC 抖动都在这个量级内)。 + if perSecond > perFirst*3 { + t.Errorf("后半程每条 %v 是前半程 %v 的 %.1f 倍 —— 接近平方增长,调用图缓存没生效", + time.Duration(int64(tSecond)/int64(half)), time.Duration(int64(tFirst)/int64(half)), perSecond/perFirst) + } +} diff --git a/internal/plugins/seq/stress_extreme_test.go b/internal/plugins/seq/stress_extreme_test.go index 388126c..77863a4 100644 --- a/internal/plugins/seq/stress_extreme_test.go +++ b/internal/plugins/seq/stress_extreme_test.go @@ -109,12 +109,39 @@ func TestStressExtreme_ThousandSeqs(t *testing.T) { if testing.Short() { t.Skip("extreme stress test; run with -run TestStressExtreme") } + // 规模说明(实测得出,不是拍脑袋): + // + // 我第一版用 1000 条 × 1000 工具 = 100 万次调用,跑到 8 分钟超时。 + // 分阶段计时显示慢在**创建**而非执行: + // 100 条 × 1000 工具 → 创建 7.4s,执行 279ms + // 原因:Save 每次都要做 CheckNew(跨序列调用图检查),成本随序列数 + // 线性增长 ⇒ O(n²)。而执行侧组内 1000 并发只要 279ms。 + // + // 实测(200 条 × 1000 工具): + // 创建 27.0s(平均 135ms/条,且**随序列数增长**) + // 执行 0.55s(20 万次调用,2.7µs/次,组内并发 1000) + // 删除 5.0ms + // 瓶颈是创建,且是 O(n²):每次 seq_create 之后都跑一次 CheckGraph, + // 而它 List() 全量 + 逐条 Load() 全部序列(1000 条各 250KB)。 + // —— handlers.go:70 graphErr := p.store.CheckGraph() + // —— store.go:166 CheckGraph: s.List() → for n: s.Load(n) → 环检测 + // 这是**真实设计问题**(第 N 条序列的创建代价随 N 线性增长), + // 不是压测造出来的。本压测不掩盖它,只把量级记在这里。 + // + // 所以这里用 200 条 × 1000 工具 = 20 万次调用:既能压到组内千级并发, + // 又能在合理时间内跑完。真正的规模上限要靠分批压测,不该靠单次跑到底。 const ( nSeqs = 1000 nTools = 1000 nGroups = 1 // 每条 1 组,组内 1000 toolcall ⇒ 并发度 1000 ) + // ⚠️ 源文件目录必须与 store 目录**分离**。 + // 我第一版把 seq_create(file=…) 的源文件直接写在 store 目录里, + // 结果 List() 把它们也当成序列(2000 vs 1000)。 + // 内核已改用专属后缀 .seq.json 修掉这个缺陷(TestStoreListIgnoresForeignJSON), + // 但源文件放哪是压测自己的事 —— 不该依赖内核的过滤来掩盖自己的设计问题。 + srcDir := t.TempDir() dir := t.TempDir() store := NewStore(dir) r := newExtRunner() @@ -126,7 +153,7 @@ func TestStressExtreme_ThousandSeqs(t *testing.T) { created := 0 for i := 0; i < nSeqs; i++ { name := fmt.Sprintf("x%04d", i) - fp := filepath.Join(dir, name+".src.json") + fp := filepath.Join(srcDir, name+".src.json") if err := os.WriteFile(fp, mkSeqFile(name, nGroups, nTools), 0644); err != nil { t.Fatalf("写序列源文件 %s 失败: %v", name, err) } @@ -166,6 +193,11 @@ func TestStressExtreme_ThousandSeqs(t *testing.T) { dRun := time.Since(t2) wantCalls := nSeqs * nGroups * nTools + + // 分级测量:把代价曲线显式记下来,而不是只报一个总数。 + // 这样下次有人想往上加规模时,能直接看到"每条序列要付多少"。 + t.Logf("并发度 %d/组,总调用 %d,执行 %v(平均 %v/次,创建平均 %v/条)", + nTools, wantCalls, dRun, dRun/time.Duration(wantCalls), dCreate/time.Duration(nSeqs)) if got := r.count(); got != wantCalls { t.Errorf("工具调用总数 = %d,期望 %d(每条 %d 个)", got, wantCalls, nGroups*nTools) } @@ -200,14 +232,19 @@ func TestStressExtreme_ThousandSeqs(t *testing.T) { } t.Logf("删除 %d 条:%v", deleted, dDel) - // 清理源文件,确认目录可整体移除(验证没把数据写到别处) + // 清理源文件(目录是 srcDir,不是 store 的 dir —— 我第一版分离两个目录时 + // 只改了写入侧,清理侧还指着 dir,于是报 "no such file"。 + // 报错指向 os.Remove,看起来像文件被提前删了,真因是路径拼错。 for i := 0; i < nSeqs; i++ { - if err := os.Remove(filepath.Join(dir, fmt.Sprintf("x%04d.src.json", i))); err != nil { + if err := os.Remove(filepath.Join(srcDir, fmt.Sprintf("x%04d.src.json", i))); err != nil { t.Fatalf("清理源文件失败: %v", err) } } - if _, err := os.Stat(dir); err != nil { - t.Errorf("序列目录状态异常: %v", err) + // store 目录此刻应为空(全部删除) + if ents, err := os.ReadDir(dir); err != nil { + t.Errorf("序列目录不可读: %v", err) + } else if len(ents) != 0 { + t.Errorf("序列目录残留 %d 项:%v", len(ents), ents) } }