mirror of
https://gitcode.com/JianFeeeee/HomeAgent.git
synced 2026-09-21 17:38:10 +00:00
三件事,都是 plan.md §13 里未打勾的项:
1. §13.3「Bench: ToolInvoke 延迟对比」补测(bench_test.go)
改造后原有的 BenchmarkToolInvoke 仍走 legacy 内联 Args,已不代表生产
路径。拆成 inline / frame 两个子基准并各测两个尺寸(-count=3 取中位):
payload inline(RPC 报文) frame(共享内存帧) 差
16 B 26.4 µs/op 48.0 µs/op +21.6 µs
32 KiB 695.6 µs/op 391.7 µs/op −303.9 µs
取舍随尺寸翻转:小 payload 多付 ~22µs 帧固定开销,大 payload 省掉整份
JSON 编解码与管道拷贝、快 44%。因为线上工具结果动辄几十 KB 且 22µs
相对 LLM 往返 2-8 秒可忽略,所以统一走帧而不按大小分叉(§13.3 第 4 条)
是对的。
2. 删除 §13.1 遗留死代码(evtring.go)
allocEvtRing 自统一共享区域之后从未被调用——旧版它单独建 memfd + eventfd,
现在只有一个 memfd,事件环只是区内的一个 segment(NewEvtRing 操作区内
切片)。留着会让人误以为事件环还有独立段。
3. 修正误导性注释(plugin.go)
仍写着旧 3-fd 布局「3=StageContext, 4=事件环, 5=通知」,与实现
(procExtraFilesForShm 只返回 memfd+evtfd)和 shmpass_unix.go 的权威
注释互相矛盾。
plan.md:§13.1/13.3/13.4/13.8 的验证项按实测打勾,并注明证据(提交号 /
测试名 / 实测数字),不留无依据的勾。
305 lines
8.2 KiB
Go
305 lines
8.2 KiB
Go
package proc
|
||
|
||
import (
|
||
"encoding/binary"
|
||
"encoding/json"
|
||
"fmt"
|
||
"sync"
|
||
"sync/atomic"
|
||
"time"
|
||
|
||
pubsdk "gitcode.com/JianFeeeee/homeagent-sdk/sdk"
|
||
)
|
||
|
||
// ---- 事件类型编码(编译时确定,与 pubsdk.EventType 一一对应)----
|
||
|
||
var evtTypeNames = [evtTypeMax]string{
|
||
"raw_input",
|
||
"agent_output",
|
||
"agent_llm_chain",
|
||
"tool_call",
|
||
"reasoning",
|
||
"stage",
|
||
"system",
|
||
"reasoning_delta",
|
||
"content_delta",
|
||
"skill_detected",
|
||
}
|
||
|
||
var evtTypeIndex = map[string]uint32{
|
||
"raw_input": evtTypeRawInput,
|
||
"agent_output": evtTypeAgentOutput,
|
||
"agent_llm_chain": evtTypeAgentLLMChain,
|
||
"tool_call": evtTypeToolCall,
|
||
"reasoning": evtTypeReasoning,
|
||
"stage": evtTypeStage,
|
||
"system": evtTypeSystem,
|
||
"reasoning_delta": evtTypeReasoningDelta,
|
||
"content_delta": evtTypeContentDelta,
|
||
"skill_detected": evtTypeSkillDetected,
|
||
}
|
||
|
||
func encodeEvtType(t pubsdk.EventType) uint32 {
|
||
if idx, ok := evtTypeIndex[string(t)]; ok {
|
||
return idx
|
||
}
|
||
return 0xFFFFFFFF // 未知类型:子进程 typeMask 用 0 匹配全部,此值不影响
|
||
}
|
||
|
||
func decodeEvtType(idx uint32) pubsdk.EventType {
|
||
if int(idx) < len(evtTypeNames) {
|
||
return pubsdk.EventType(evtTypeNames[idx])
|
||
}
|
||
return ""
|
||
}
|
||
|
||
func evtTypeMask(types ...pubsdk.EventType) uint32 {
|
||
var mask uint32
|
||
for _, t := range types {
|
||
if idx, ok := evtTypeIndex[string(t)]; ok {
|
||
mask |= 1 << idx
|
||
}
|
||
}
|
||
return mask
|
||
}
|
||
|
||
// ---- 事件环共享段布局(§3.6)----
|
||
|
||
const (
|
||
evtRingMagic uint32 = 0x48455654 // "HEVT"
|
||
evtRingVersion uint32 = 1
|
||
evtRingCap uint32 = 8192 // 2^13,满足流式场景突发(实验 4)
|
||
evtRingSlotLen uint32 = 32 // seq(8)+type(4)+off(4)+len(4)+pad(12)
|
||
|
||
evtOffMagic uint32 = 0
|
||
evtOffVersion uint32 = 4
|
||
evtOffWriteSeq uint32 = 8
|
||
evtOffCap uint32 = 16
|
||
evtOffSlots uint32 = 20
|
||
|
||
evtTypeRawInput uint32 = 0
|
||
evtTypeAgentOutput uint32 = 1
|
||
evtTypeAgentLLMChain uint32 = 2
|
||
evtTypeToolCall uint32 = 3
|
||
evtTypeReasoning uint32 = 4
|
||
evtTypeStage uint32 = 5
|
||
evtTypeSystem uint32 = 6
|
||
evtTypeReasoningDelta uint32 = 7
|
||
evtTypeContentDelta uint32 = 8
|
||
evtTypeSkillDetected uint32 = 9
|
||
evtTypeMax uint32 = 10
|
||
|
||
evtHeaderSize = 20
|
||
evtArenaCap = 64 * 1024
|
||
evtTotalSize = int(evtHeaderSize + evtRingCap*evtRingSlotLen + evtArenaCap)
|
||
)
|
||
|
||
// ---- 内核侧:EvtRing ----
|
||
|
||
type EvtRing struct {
|
||
data []byte
|
||
writeSeq atomic.Uint64
|
||
cap uint32
|
||
slotsBase uint32
|
||
arenaBase uint32
|
||
arenaCap uint32
|
||
arenaUsed atomic.Uint32
|
||
mu sync.Mutex
|
||
}
|
||
|
||
func NewEvtRing(data []byte) (*EvtRing, error) {
|
||
if uint32(len(data)) < evtHeaderSize+evtRingCap*evtRingSlotLen+evtArenaCap {
|
||
return nil, fmt.Errorf("事件环段太小:需要 %d,实际 %d", evtTotalSize, len(data))
|
||
}
|
||
if got := binary.LittleEndian.Uint32(data[evtOffMagic:]); got != evtRingMagic {
|
||
return nil, fmt.Errorf("事件环魔数不匹配(0x%x)", got)
|
||
}
|
||
return &EvtRing{
|
||
data: data,
|
||
cap: evtRingCap,
|
||
slotsBase: evtOffSlots,
|
||
arenaBase: evtOffSlots + evtRingCap*evtRingSlotLen,
|
||
arenaCap: evtArenaCap,
|
||
}, nil
|
||
}
|
||
|
||
func (r *EvtRing) Init() {
|
||
binary.LittleEndian.PutUint32(r.data[evtOffMagic:], evtRingMagic)
|
||
binary.LittleEndian.PutUint32(r.data[evtOffVersion:], evtRingVersion)
|
||
binary.LittleEndian.PutUint32(r.data[evtOffCap:], r.cap)
|
||
r.writeSeq.Store(0)
|
||
}
|
||
|
||
// WritePush post-and-forget,**绝不阻塞**(§3.6 约束 B)。
|
||
func (r *EvtRing) WritePush(evtType pubsdk.EventType, payload []byte) {
|
||
seq := r.writeSeq.Add(1) - 1
|
||
var off uint32
|
||
r.mu.Lock()
|
||
used := r.arenaUsed.Load()
|
||
if used+uint32(len(payload)) <= r.arenaCap {
|
||
off = r.arenaBase + used
|
||
r.arenaUsed.Store(used + uint32(len(payload)))
|
||
copy(r.data[off:], payload)
|
||
}
|
||
r.mu.Unlock()
|
||
idx := seq % uint64(r.cap)
|
||
slotOff := r.slotsBase + uint32(idx)*evtRingSlotLen
|
||
binary.LittleEndian.PutUint64(r.data[slotOff:], seq)
|
||
binary.LittleEndian.PutUint32(r.data[slotOff+8:], encodeEvtType(evtType))
|
||
binary.LittleEndian.PutUint32(r.data[slotOff+12:], off)
|
||
binary.LittleEndian.PutUint32(r.data[slotOff+16:], uint32(len(payload)))
|
||
binary.LittleEndian.PutUint64(r.data[evtOffWriteSeq:], seq+1)
|
||
}
|
||
|
||
// ---- 子进程侧:EvtConsumer ----
|
||
|
||
type EvtConsumer struct {
|
||
ringData []byte
|
||
evtfd evtfdReader
|
||
handler func(*pubsdk.Event) error
|
||
readSeq uint64
|
||
typeMask uint32
|
||
mu sync.Mutex
|
||
running bool
|
||
stop chan struct{}
|
||
stopOnce sync.Once
|
||
done chan struct{}
|
||
}
|
||
|
||
type evtfdReader interface {
|
||
Read(b []byte) (int, error)
|
||
}
|
||
|
||
// evtfdFder 是可取出原始 fd 的通知句柄(*os.File 满足)。
|
||
//
|
||
// 有它才能用 poll(2) 加超时等待可读。**为什么不能用 SetReadDeadline**:
|
||
// eventfd/pipe 经 os.NewFile 包装后不会注册进 Go netpoller(os.NewFile 对
|
||
// 非 open 得到的 fd 一律按非 pollable 处理),Read 退化成阻塞 syscall,
|
||
// SetReadDeadline 返回错误且不生效——Run 会永久卡在 syscall.Read,
|
||
// 此时调用方若已 munmap 区域(host.Close),恢复后的 drainEvents 就是
|
||
// 读已解除映射的内存:SIGSEGV,recover 捕不到。
|
||
type evtfdFder interface {
|
||
Fd() uintptr
|
||
}
|
||
|
||
func NewEvtConsumer(ringData []byte, evtfd evtfdReader, mask uint32, handler func(*pubsdk.Event) error) *EvtConsumer {
|
||
return &EvtConsumer{
|
||
ringData: ringData,
|
||
evtfd: evtfd,
|
||
handler: handler,
|
||
typeMask: mask,
|
||
stop: make(chan struct{}),
|
||
done: make(chan struct{}),
|
||
}
|
||
}
|
||
|
||
// readWakeInterval 是 poll 超时间隔:保证 Run 至少这么频繁地检查 stop。
|
||
const readWakeInterval = 100 * time.Millisecond
|
||
|
||
func (c *EvtConsumer) Run() {
|
||
c.mu.Lock()
|
||
if c.running {
|
||
c.mu.Unlock()
|
||
return
|
||
}
|
||
c.running = true
|
||
c.mu.Unlock()
|
||
defer func() {
|
||
c.mu.Lock()
|
||
c.running = false
|
||
c.mu.Unlock()
|
||
close(c.done)
|
||
}()
|
||
|
||
buf := make([]byte, 8)
|
||
for {
|
||
select {
|
||
case <-c.stop:
|
||
return
|
||
default:
|
||
}
|
||
|
||
// 显式 poll(2) 加超时:只有它能打破阻塞 Read,让 Stop 真正生效。
|
||
if f, ok := c.evtfd.(evtfdFder); ok {
|
||
ready, err := pollEvtfd(int(f.Fd()), int(readWakeInterval/time.Millisecond))
|
||
if err != nil {
|
||
return // 通知句柄已失效,退出以免空转
|
||
}
|
||
if !ready {
|
||
continue // 超时:回到顶部检查 stop
|
||
}
|
||
}
|
||
|
||
// 阻塞等待内核通知(走 netpoller,只 park goroutine)
|
||
if _, err := c.evtfd.Read(buf); err != nil {
|
||
select {
|
||
case <-c.stop:
|
||
return
|
||
default:
|
||
}
|
||
continue
|
||
}
|
||
c.drainEvents()
|
||
}
|
||
}
|
||
|
||
func (c *EvtConsumer) drainEvents() {
|
||
writeSeq := binary.LittleEndian.Uint64(c.ringData[evtOffWriteSeq:])
|
||
cap := uint64(evtRingCap)
|
||
for c.readSeq < writeSeq {
|
||
// 每处理一条就检查一次 stop:handler 可能很慢,
|
||
// 不加这个检查的话 Stop 要等整轮 drain 完才生效。
|
||
select {
|
||
case <-c.stop:
|
||
return
|
||
default:
|
||
}
|
||
if writeSeq-c.readSeq > cap {
|
||
c.readSeq = writeSeq - cap
|
||
}
|
||
idx := c.readSeq % cap
|
||
slotOff := evtOffSlots + uint32(idx)*evtRingSlotLen
|
||
seq := binary.LittleEndian.Uint64(c.ringData[slotOff:])
|
||
etype := binary.LittleEndian.Uint32(c.ringData[slotOff+8:])
|
||
off := binary.LittleEndian.Uint32(c.ringData[slotOff+12:])
|
||
slen := binary.LittleEndian.Uint32(c.ringData[slotOff+16:])
|
||
if seq != c.readSeq {
|
||
// slot 已被新事件覆盖——逐个扫太慢(溢出场景 readSeq=0 要跳 100+ 步),
|
||
// 直接跳到 writeSeq 附近找下一个可读 slot。
|
||
// 简化:溢出后直接跳到 writeSeq - cap(最旧的可读事件)。
|
||
if writeSeq > cap {
|
||
c.readSeq = writeSeq - cap
|
||
} else {
|
||
c.readSeq = writeSeq
|
||
}
|
||
continue
|
||
}
|
||
// 位掩码过滤
|
||
if c.typeMask != 0 && (1<<etype)&c.typeMask == 0 {
|
||
c.readSeq++
|
||
continue
|
||
}
|
||
if off > 0 && slen > 0 && uint64(off)+uint64(slen) <= uint64(len(c.ringData)) {
|
||
payload := make([]byte, slen)
|
||
copy(payload, c.ringData[off:off+slen])
|
||
var evt pubsdk.Event
|
||
if err := json.Unmarshal(payload, &evt); err == nil {
|
||
c.handler(&evt)
|
||
}
|
||
}
|
||
c.readSeq++
|
||
}
|
||
}
|
||
|
||
// Stop 请求消费者退出。
|
||
//
|
||
// 非阻塞:Stop 返回**不代表** Run 已退出(最多 readWakeInterval 后退出)。
|
||
// 若要在 Stop 之后释放 ringData(host.Close 会 munmap 整个区域),
|
||
// 必须先 Stop() 再 Wait()。
|
||
func (c *EvtConsumer) Stop() {
|
||
c.stopOnce.Do(func() { close(c.stop) })
|
||
}
|
||
|
||
// Wait 阻塞至 Run 退出。返回后 drainEvents 保证不会再访问 ringData。
|
||
func (c *EvtConsumer) Wait() { <-c.done }
|