diff --git a/docs/zh/experiments/plugin-arch/19-migration-verify/README.md b/docs/zh/experiments/plugin-arch/19-migration-verify/README.md index 9b69747..bfa29c6 100644 --- a/docs/zh/experiments/plugin-arch/19-migration-verify/README.md +++ b/docs/zh/experiments/plugin-arch/19-migration-verify/README.md @@ -81,3 +81,57 @@ RSS 随二进制体积线性增长,故绝对数字不可比。可比的是结 这些测试用**真实 example 产物**而非 testdata 假插件,且 manifest 刻意写 `"entry":"plugin.so"`——验证「业务代码零改动」这一承诺在完整内核装配下成立。 未重编时 skip 而非 fail,CI 不强制先跑重编脚本。 + +## 压测与延迟(Part 6.6 验收) + +基准与压测在代码里而非独立脚本: +`internal/plugin/proc/bench_test.go` + `streaming_test.go`。 + +```bash +go test -run '^$' -bench . ./internal/plugin/proc/ +go test -run 'TestStreaming_' -v ./internal/plugin/proc/ +``` + +### 实测(2026-09-02,AMD Ryzen 7 7840HS) + +| 项目 | 实测 | 基线 | 判断 | +|---|---|---|---| +| 工具调用 RPC 往返 | 24.1 µs | 实验 11: 19.6 µs | 同量级 | +| 锁仲裁(内核侧) | 0.76 µs | — | 见下注 | +| 事件环写入 | 95 ns | — | 亚微秒 | +| 事件环并发写入 | 83 ns | — | 无锁竞争恶化 | +| 完整 stage 往返 | 132 µs | — | 含 3 次进程间往返 | +| 共享段编解码 | 3.7 µs | — | 占 stage 的 2.8% | + +**锁仲裁 0.76µs 不可与实验 3 的 19.40µs 对照**——两者测的不是同一个东西: +实验 3 测插件经 RPC 请求锁的完整跨进程往返,本基准只测内核侧 +`lockRegistry.acquire/release`。真实成本仍在 20µs 量级(那部分是 RPC 往返)。 +基准原名 `BenchmarkStageLockRoundTrip` 有误导性,已改为 +`BenchmarkStageLockArbitration`。 + +**stage 往返 132µs 的成本构成**:共享段编解码只占 3.7µs(2.8%), +其余是**一次 stage 要走 3 次进程间往返**——`stage.invoke` 加上插件侧反向的 +`stage.lock` / `stage.unlock`。相对 LLM 往返 2-8 秒可忽略;若日后要优化, +方向是把 lock/unlock 合入 `stage.invoke` 的请求/应答,省掉两次往返。 + +### 流式压测(§4.3 标记「风险高」的那一项) + +原文的担忧:「`Bus.Publish` 路径禁用任何锁/阻塞——流式输出逐 token 发布, +任何等待都会卡顿」。 + +``` +5000 次 Publish + 每条睡 20µs 的慢消费者 + 实测 2.29ms,均摊 457 ns/token + 同步语义理论下限 100ms(5000 × 20µs) + +订阅者 1 个:1.547ms(515 ns/次) +订阅者 8 个:1.518ms(506 ns/次) ← 几乎不变,无线性恶化 + +环溢出(无消费者写 30000 次,cap=8192):均摊 35 ns/次 ← 仍 O(1) +``` + +2.29ms 与实验 4 的数字完全一致(那次也是 2.29ms / 0.46µs per token)—— +post-and-forget 在实现中成立。 + +最后一项的意义:消费者完全停摆时写端覆盖最旧 slot,这条路径仍是 O(1), +故「消费者卡住」不会连带拖慢内核主循环。 diff --git a/internal/plugin/proc/bench_test.go b/internal/plugin/proc/bench_test.go new file mode 100644 index 0000000..24f4bb2 --- /dev/null +++ b/internal/plugin/proc/bench_test.go @@ -0,0 +1,222 @@ +package proc + +import ( + "os" + "os/exec" + "path/filepath" + "testing" + + pubsdk "gitcode.com/JianFeeeee/homeagent-sdk/sdk" +) + +// 子进程架构的性能基准(Part 6.6 验收项)。 +// +// 对照基线来自 docs/zh/experiments/plugin-arch: +// +// 实验 3 锁仲裁 RPC 往返 19.40 µs/次 +// 实验 4 post-and-forget 5.07s → 2.29ms(5000 token + 20µs 慢消费者) +// 实验 11 工具调用 RPC p50 19.6 µs +// +// 这些基准回答的是「进程边界的代价是否可忽略」——相对 stage handler 的实际 +// 工作量(LLM 往返 2-8 秒),微秒级往返不构成问题;但若退化到毫秒级, +// 高频工具调用就会被感知。 +// +// 实测结果(2026-09-02,AMD Ryzen 7 7840HS): +// +// ToolInvoke 24.1 µs/op ← 对照实验 11 的 19.6µs,同量级 +// StageLockArbitration 0.76 µs/op ← 仅内核侧仲裁,不跨进程 +// EvtRingWritePush 95 ns/op +// EvtRingWritePushConcurrent 83 ns/op ← 并发不恶化 +// StageInvokeSharedMemory 132 µs/op ← 含 3 次进程间往返 +// SegmentWriteAllReadInto 3.7 µs/op ← 占 stage 的 2.8% + +// buildBenchPlugin 编译 testdata 里的测试插件(benchmark 版)。 +func buildBenchPlugin(b *testing.B, srcName string) string { + b.Helper() + src := filepath.Join("testdata", srcName) + if _, err := os.Stat(src); err != nil { + b.Skipf("测试插件源码缺失 %s: %v", src, err) + } + bin := filepath.Join(b.TempDir(), "benchplugin") + cmd := exec.Command("go", "build", "-o", bin, src) + cmd.Env = append(os.Environ(), "CGO_ENABLED=0") + if out, err := cmd.CombinedOutput(); err != nil { + b.Fatalf("编译 %s: %v\n%s", srcName, err, out) + } + return bin +} + +// BenchmarkToolInvoke 测量内核 → 插件的工具调用往返。 +// +// 链路:Call 写 stdin → 插件读循环 → handler → 写 stdout → +// 内核 readLoop → pending channel 唤醒。对照实验 11 的 19.6µs。 +func BenchmarkToolInvoke(b *testing.B) { + bin := buildBenchPlugin(b, "echoplugin.go") + + p, err := Spawn("echo", bin, Options{Handler: noopHandler}) + if err != nil { + b.Fatalf("Spawn: %v", err) + } + defer p.Kill() + + args := map[string]interface{}{"text": "benchmark"} + + b.ResetTimer() + for i := 0; i < b.N; i++ { + if _, err := p.Call(MethodToolInvoke, ToolInvokeParams{ + Name: "echo_tool", + Args: args, + }); err != nil { + b.Fatalf("第 %d 次调用失败: %v", i, err) + } + } +} + +// BenchmarkStageLockArbitration 测量**内核侧锁仲裁本身**的成本。 +// +// ⚠️ 不要拿这个数字对照实验 3 的 19.40µs——两者测的不是同一个东西: +// - 实验 3:插件经 RPC 请求锁的**完整跨进程往返** +// - 本基准:仅 lockRegistry.acquire/release,不跨进程 +// +// 真实成本仍在 20µs 量级(那部分是 RPC 往返,见 BenchmarkToolInvoke)。 +// 本基准的用途是确认仲裁逻辑自身不是瓶颈:若它也到了微秒级, +// 说明 sync.Mutex 之外又引入了什么开销。 +func BenchmarkStageLockArbitration(b *testing.B) { + lock := newStageLock() + r := &lockRegistry{} + r.bind(lock) + + b.ResetTimer() + for i := 0; i < b.N; i++ { + if err := r.acquire("bench"); err != nil { + b.Fatalf("acquire: %v", err) + } + if err := r.release("bench"); err != nil { + b.Fatalf("release: %v", err) + } + } +} + +// BenchmarkEvtRingWritePush 测量事件环写入(Bus.Publish 路径)。 +// +// 这是 §4.3 标记「风险高」的那一项:流式输出逐 token 发布, +// Publish 路径上任何阻塞都会直接卡顿。 +func BenchmarkEvtRingWritePush(b *testing.B) { + host, err := NewHost() + if err != nil { + b.Fatalf("NewHost: %v", err) + } + defer host.Close() + + ring := host.EvtRing() + payload := []byte(`{"type":"content_delta","payload":{"text":"token"}}`) + + b.ResetTimer() + for i := 0; i < b.N; i++ { + ring.WritePush(pubsdk.EventContentDelta, payload) + } +} + +// BenchmarkEvtRingWritePushConcurrent 并发写入。 +// +// 内核有多条路径并发发布事件(主循环、工具调用、流式增量), +// writeSeq 是 atomic 而 arena 分配有锁——确认锁不是瓶颈。 +func BenchmarkEvtRingWritePushConcurrent(b *testing.B) { + host, err := NewHost() + if err != nil { + b.Fatalf("NewHost: %v", err) + } + defer host.Close() + + ring := host.EvtRing() + payload := []byte(`{"type":"content_delta","payload":{"text":"tok"}}`) + + b.ResetTimer() + b.RunParallel(func(pb *testing.PB) { + for pb.Next() { + ring.WritePush(pubsdk.EventContentDelta, payload) + } + }) +} + +// BenchmarkStageInvokeSharedMemory 测量完整 stage 往返: +// 写共享段 → RPC → 插件读改写 → 回读 → 压实。 +// +// 这是迁移引入的最重路径,每次 stage 都要走一遍。 +// +// 实测 ~132µs,比单次 RPC(~24µs)高 5 倍,因为**一次 stage 要走 3 次 +// 进程间往返**:stage.invoke + 插件侧反向的 stage.lock / stage.unlock。 +// 共享段编解码只占 3.7µs(2.8%)——成本在往返次数而非数据搬运。 +// +// 相对 LLM 往返 2-8 秒可忽略。若日后要优化,方向是把 lock/unlock +// 合入 stage.invoke 的请求/应答,省掉两次往返。 +func BenchmarkStageInvokeSharedMemory(b *testing.B) { + bin := buildBenchPlugin(b, "stageplugin.go") + + host, err := NewHost() + if err != nil { + b.Fatalf("NewHost: %v", err) + } + defer host.Close() + + core := newFakeCore() + p := New("sanitizer", bin, b.TempDir(), nil, host, nil) + if err := p.Start(core); err != nil { + b.Fatalf("Start: %v", err) + } + defer p.Close() + + handlers := core.stageHandlers(pubsdk.StageAfterToolcall) + if len(handlers) != 1 { + b.Fatalf("应注册 1 个 handler,实际 %d", len(handlers)) + } + handler := handlers[0] + + b.ResetTimer() + for i := 0; i < b.N; i++ { + sc := &pubsdk.StageContext{ + Phase: pubsdk.StageAfterToolcall, + ToolResults: []pubsdk.ToolResult{ + {CallID: "c1", Name: "t", Result: "结果:\x1b[31m脏\x1b[0m"}, + }, + } + if err := handler(sc); err != nil { + b.Fatalf("第 %d 次 stage 失败: %v", i, err) + } + } +} + +// BenchmarkSegmentWriteAllReadInto 只测共享段编解码(不含 RPC)。 +// +// 用于拆分 stage 往返的成本构成:编解码 vs 进程间通信。 +func BenchmarkSegmentWriteAllReadInto(b *testing.B) { + host, err := NewHost() + if err != nil { + b.Fatalf("NewHost: %v", err) + } + defer host.Close() + seg := host.Segment() + + sc := &pubsdk.StageContext{ + Phase: pubsdk.StageAfterToolcall, + RawMessage: "用户输入的一段话", + UserID: "u1", + LLMText: "模型输出的文本", + FinalText: "最终文本", + ToolResults: []pubsdk.ToolResult{ + {CallID: "c1", Name: "tool_a", Result: "结果 A"}, + {CallID: "c2", Name: "tool_b", Result: "结果 B"}, + }, + } + + b.ResetTimer() + for i := 0; i < b.N; i++ { + if err := seg.WriteAll(sc); err != nil { + b.Fatalf("WriteAll: %v", err) + } + if err := seg.ReadInto(sc); err != nil { + b.Fatalf("ReadInto: %v", err) + } + seg.Compact() + } +} diff --git a/internal/plugin/proc/streaming_test.go b/internal/plugin/proc/streaming_test.go new file mode 100644 index 0000000..0418f94 --- /dev/null +++ b/internal/plugin/proc/streaming_test.go @@ -0,0 +1,157 @@ +package proc + +import ( + "sync" + "sync/atomic" + "testing" + "time" + + pubsdk "gitcode.com/JianFeeeee/homeagent-sdk/sdk" +) + +// 流式输出压测(§4.3 标记「风险高」的那一项)。 +// +// 担忧的原文:「Bus.Publish 路径禁用任何锁/阻塞——流式输出逐 token 发布, +// 任何等待都会卡顿」。实验 4 的数据:同步 Publish + 一个 20µs 慢订阅者, +// 5000 token 耗时 5.07s;改为写环 + post 后 2.29ms(加速比 2218x)。 +// +// 这里验证事件环侧的 post-and-forget 性质在实现中成立。 + +// 慢消费者不拖慢 Publish。 +// +// 判据:若 Publish 等消费者,5000 × 20µs = 100ms 是理论下限。 +// post-and-forget 应远低于此。 +func TestStreaming_SlowConsumerDoesNotBlockPublish(t *testing.T) { + host, err := NewHost() + if err != nil { + t.Fatalf("NewHost: %v", err) + } + defer host.Close() + + ring := host.EvtRing() + + var consumed atomic.Int64 + consumer := NewEvtConsumer(host.EvtData(), host.EvtfdReadFile(), 0, + func(evt *pubsdk.Event) error { + time.Sleep(20 * time.Microsecond) // 刻意的慢订阅者 + consumed.Add(1) + return nil + }) + go consumer.Run() + defer consumer.Stop() + + const tokens = 5000 + payload := []byte(`{"type":"content_delta","payload":{"text":"t"}}`) + + start := time.Now() + for i := 0; i < tokens; i++ { + ring.WritePush(pubsdk.EventContentDelta, payload) + EvtfdNotify(host.EvtNotifyFd()) + } + elapsed := time.Since(start) + perToken := elapsed / tokens + + t.Logf("%d 次 Publish 耗时 %v,均摊 %v/token(消费者每条睡 20µs)", + tokens, elapsed, perToken) + t.Logf("同步语义下的理论下限:%v", tokens*20*time.Microsecond) + + if elapsed > 100*time.Millisecond { + t.Errorf("Publish 疑似被慢消费者阻塞:耗时 %v ≥ 同步下限 100ms", elapsed) + } + if perToken > 20*time.Microsecond { + t.Errorf("均摊 %v/token ≥ 消费者处理时间 20µs,说明存在等待", perToken) + } +} + +// 订阅者增多不使 Publish 线性恶化。 +// +// §4.3 的具体要求:「长回复下 Publish 单次耗时不随订阅者数线性恶化」。 +func TestStreaming_PublishLatencyFlatAcrossSubscribers(t *testing.T) { + host, err := NewHost() + if err != nil { + t.Fatalf("NewHost: %v", err) + } + defer host.Close() + + ring := host.EvtRing() + payload := []byte(`{"type":"content_delta","payload":{"text":"t"}}`) + const rounds = 3000 + + measure := func(consumers int) time.Duration { + var wg sync.WaitGroup + active := make([]*EvtConsumer, 0, consumers) + for i := 0; i < consumers; i++ { + c := NewEvtConsumer(host.EvtData(), host.EvtfdReadFile(), 0, + func(evt *pubsdk.Event) error { + time.Sleep(10 * time.Microsecond) + return nil + }) + active = append(active, c) + wg.Add(1) + go func(cc *EvtConsumer) { + defer wg.Done() + cc.Run() + }(c) + } + defer func() { + for _, c := range active { + c.Stop() + } + }() + + // 让消费者先就位 + time.Sleep(10 * time.Millisecond) + + start := time.Now() + for i := 0; i < rounds; i++ { + ring.WritePush(pubsdk.EventContentDelta, payload) + EvtfdNotify(host.EvtNotifyFd()) + } + return time.Since(start) + } + + d1 := measure(1) + d8 := measure(8) + + t.Logf("1 个消费者:%v(均摊 %v/次)", d1, d1/rounds) + t.Logf("8 个消费者:%v(均摊 %v/次)", d8, d8/rounds) + + // 线性恶化的判据:8 倍订阅者不应接近 8 倍耗时。 + // 阈值取 4 倍——测量噪声与调度抖动都会影响。 + if d8 > d1*4 { + t.Errorf("订阅者 1→8,Publish 从 %v 涨到 %v(>4 倍),疑似线性恶化", d1, d8) + } +} + +// 事件环溢出时 Publish 不退化。 +// +// 消费者完全停摆时写端会覆盖最旧 slot。这条路径必须仍是 O(1), +// 否则「消费者卡住」会连带拖慢内核主循环。 +func TestStreaming_PublishStaysFastWhenRingOverflows(t *testing.T) { + host, err := NewHost() + if err != nil { + t.Fatalf("NewHost: %v", err) + } + defer host.Close() + + ring := host.EvtRing() + payload := []byte(`{"type":"content_delta","payload":{"text":"t"}}`) + + // 无消费者:环必然溢出(cap=8192) + const rounds = 30000 + + start := time.Now() + for i := 0; i < rounds; i++ { + ring.WritePush(pubsdk.EventContentDelta, payload) + } + elapsed := time.Since(start) + perPush := elapsed / rounds + + t.Logf("无消费者写入 %d 次(环 cap=%d,必然溢出):%v,均摊 %v/次", + rounds, evtRingCap, elapsed, perPush) + + // 溢出路径仍应是亚微秒级 + if perPush > 5*time.Microsecond { + t.Errorf("溢出时均摊 %v/次,超出预期(应亚微秒级)", perPush) + } +}