mirror of
https://gitcode.com/JianFeeeee/HomeAgent.git
synced 2026-09-21 09:28:14 +00:00
proc: 性能基准 + 流式压测(Part 6.6 验收项)
此前只做了功能冒烟与内存快照,延迟与压测都没测。这两项是计划里 明确列出的验收条件,补上。 ## 基准结果(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 量级。 基准原名 BenchmarkStageLockRoundTrip 有误导性,已改为 BenchmarkStageLockArbitration,并在注释里写明不可对照的理由—— 否则日后有人拿 0.76µs 去比 19.4µs 会得出「优化了 25 倍」的错误结论。 **stage 往返 132µs 的成本构成**:共享段编解码只占 3.7µs,其余是 一次 stage 要走 3 次进程间往返(stage.invoke + 插件侧反向的 stage.lock / stage.unlock)。相对 LLM 往返 2-8 秒可忽略;要优化的方向是 把 lock/unlock 合入 stage.invoke 的请求/应答,省掉两次往返。 ## 流式压测:§4.3 标记「风险高」的那一项通过 原文担忧:「Bus.Publish 路径禁用任何锁/阻塞——流式输出逐 token 发布, 任何等待都会卡顿」。事件环是 Part 5 新加在这条路径上的,必须验。 ``` 5000 次 Publish + 每条睡 20µs 的慢消费者 实测 2.29ms,均摊 457 ns/token 同步语义理论下限 100ms 订阅者 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), 故「消费者卡住」不会连带拖慢内核主循环。 Ref: docs/zh/plugin-migration-plan.md Part 6.6、docs/zh/架构迁移评估.md §4.3
This commit is contained in:
222
internal/plugin/proc/bench_test.go
Normal file
222
internal/plugin/proc/bench_test.go
Normal file
@ -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()
|
||||
}
|
||||
}
|
||||
157
internal/plugin/proc/streaming_test.go
Normal file
157
internal/plugin/proc/streaming_test.go
Normal file
@ -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)
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user