Files
HomeAgent/internal/plugin/proc/streaming_test.go
JianFeeeee 2572688c51 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
2026-09-02 21:42:18 +08:00

158 lines
4.4 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

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→8Publish 从 %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)
}
}