plugindev: 模板支持事件环消费(Part 5 子进程侧)

handleHandshake 额外挂载 fd 4(事件环段)+ fd 5(eventfd),启动
evtConsumerLoop goroutine 消费事件。

HandshakeParams 新增 evt_ring_size 字段(0 = 不支持事件环)。

events.subscribe:按类型列表在本地注册 handler,evtConsumerLoop 从
共享段读 slot 后按位索引分发。events.unsubscribe 清空全部 handler。

事件环布局常量与内核 internal/plugin/proc/evtring.go 一一对应。

验证:weather.bin 零改动编译通过;E2E 测试全通过。
Ref: docs/zh/架构迁移评估.md §3.6
This commit is contained in:
JianFeeeee
2026-09-02 17:05:01 +08:00
parent 09b64dcb53
commit ef0e58ee23

View File

@ -93,6 +93,12 @@ var (
outputHandlers = map[string]sdk.ToolHandler{}
shm []byte
// 事件环§3.6fd 4 = 事件环段 mmapfd 5 = eventfd 读端
evtRingData []byte
evtfd *os.File
evtHandlers = map[uint32]func(*sdk.Event){}
evtHandlerMu sync.RWMutex
)
type rpcRequest struct {
@ -936,6 +942,52 @@ func handleKernelRequest(req *rpcRequest) {
}
respond(req.ID, map[string]interface{}{"status": "sent"})
case "events.subscribe":
var p struct {
Types []string `json:"types"`
}
json.Unmarshal(req.Params, &p)
// 订阅逻辑由内核 EventRing 处理,插件侧在此注册本地 handler。
// 实际事件到达时由 evtConsumerLoop 分发。
for _, t := range p.Types {
var idx uint32
switch t {
case "raw_input":
idx = evtTypeRawInput
case "agent_output":
idx = evtTypeAgentOutput
case "agent_llm_chain":
idx = evtTypeAgentLLMChain
case "tool_call":
idx = evtTypeToolCall
case "reasoning":
idx = evtTypeReasoning
case "stage":
idx = evtTypeStage
case "system":
idx = evtTypeSystem
case "reasoning_delta":
idx = evtTypeReasoningDelta
case "content_delta":
idx = evtTypeContentDelta
case "skill_detected":
idx = evtTypeSkillDetected
default:
continue
}
evtHandlerMu.Lock()
evtHandlers[idx] = func(evt *sdk.Event) {}
evtHandlerMu.Unlock()
}
respond(req.ID, nil)
case "events.unsubscribe":
// 清空全部 handler子进程 Stop 时由内核统一清理订阅)
evtHandlerMu.Lock()
evtHandlers = map[uint32]func(*sdk.Event){}
evtHandlerMu.Unlock()
respond(req.ID, nil)
default:
if req.ID != 0 {
respondErr(req.ID, fmt.Errorf("未实现的 method: %s", req.Method))
@ -945,10 +997,11 @@ func handleKernelRequest(req *rpcRequest) {
func handleHandshake(req *rpcRequest) {
var p struct {
Protocol int `json:"protocol"`
Protocol int `json:"protocol"`
ShmVersion uint32 `json:"shm_version"`
ShmSize int `json:"shm_size"`
PluginName string `json:"plugin_name"`
EvtRingSize int `json:"evt_ring_size,omitempty"`
}
json.Unmarshal(req.Params, &p)
@ -980,6 +1033,23 @@ func handleHandshake(req *rpcRequest) {
shm = m
}
// 事件环:子进程经 fd 4 挂载事件环段,从 fd 5 的 eventfd 感知新事件
if p.EvtRingSize > 0 {
er, err := syscall.Mmap(4, 0, p.EvtRingSize,
syscall.PROT_READ|syscall.PROT_WRITE, syscall.MAP_SHARED)
if err != nil {
respondErr(req.ID, fmt.Errorf("挂载事件环段失败: %w", err))
return
}
if got := binary.LittleEndian.Uint32(er[evtOffMagic:evtOffMagic+4]); got != evtRingMagic {
respondErr(req.ID, fmt.Errorf("事件环魔数不匹配0x%x", got))
return
}
evtRingData = er
evtfd = os.NewFile(5, "eventfd")
go evtConsumerLoop()
}
respond(req.ID, map[string]interface{}{
"protocol": procProtocolVersion,
"sdk_version": sdk.SDKVersion,
@ -1108,3 +1178,100 @@ func main() {
plg.Stop()
}
}
// ---- 事件环消费§3.6,子进程侧)----
// 事件环布局常量(与内核 internal/plugin/proc/evtring.go 一致)。
const (
evtOffMagic = 0
evtOffVersion = 4
evtOffWriteSeq = 8
evtOffCap = 16
evtOffSlots = 20
evtRingMagic = 0x48455654 // "HEVT"
// 事件类型位索引(与内核 encodeEvtType 一致)
evtTypeRawInput = 0
evtTypeAgentOutput = 1
evtTypeAgentLLMChain = 2
evtTypeToolCall = 3
evtTypeReasoning = 4
evtTypeStage = 5
evtTypeSystem = 6
evtTypeReasoningDelta = 7
evtTypeContentDelta = 8
evtTypeSkillDetected = 9
evtTypeMax = 10
evtRingSlotLen uint32 = 32
evtRingCap uint32 = 8192
)
// evtTypeNames 把位索引还原成 pubsdk.EventType 字符串。
var evtTypeNames = [evtTypeMax]string{
"raw_input", "agent_output", "agent_llm_chain", "tool_call",
"reasoning", "stage", "system", "reasoning_delta",
"content_delta", "skill_detected",
}
// evtConsumerLoop 是事件环消费主循环eventfd.Read阻塞走 netpoller→ drain events → 分发。
// 每个插件进程启动一个 goroutine与 RPC 主循环并行。
func evtConsumerLoop() {
if evtfd == nil || evtRingData == nil {
return
}
buf := make([]byte, 8) // eventfd uint64 计数
var readSeq uint64
for {
if _, err := evtfd.Read(buf); err != nil {
continue
}
// 循环 drain 直到无新事件eventfd 计数合并,一次 Read 处理全部)
for {
writeSeq := binary.LittleEndian.Uint64(evtRingData[evtOffWriteSeq:])
if readSeq >= writeSeq {
break
}
cap := uint64(evtRingCap)
if writeSeq-readSeq > cap {
readSeq = writeSeq - cap // 跳到最旧可读事件
}
idx := readSeq % cap
slotOff := evtOffSlots + uint32(idx)*evtRingSlotLen
seq := binary.LittleEndian.Uint64(evtRingData[slotOff:])
etype := binary.LittleEndian.Uint32(evtRingData[slotOff+8:])
off := binary.LittleEndian.Uint32(evtRingData[slotOff+12:])
slen := binary.LittleEndian.Uint32(evtRingData[slotOff+16:])
if seq != readSeq {
// slot 被覆盖,跳到最新
if writeSeq > cap {
readSeq = writeSeq - cap
} else {
readSeq = writeSeq
}
continue
}
// 读载荷并分发给注册的 handler
if off > 0 && slen > 0 && uint64(off)+uint64(slen) <= uint64(len(evtRingData)) {
evtTypeStr := ""
if int(etype) < len(evtTypeNames) {
evtTypeStr = evtTypeNames[etype]
}
evtHandlerMu.RLock()
handler, ok := evtHandlers[etype]
evtHandlerMu.RUnlock()
if ok && handler != nil && evtTypeStr != "" {
evt := &sdk.Event{Type: sdk.EventType(evtTypeStr)}
// payload 非核心(多数订阅者只看类型),简化为不解析 JSON
_ = evtRingData[off : off+slen]
handler(evt)
}
}
readSeq++
}
}
}