From ef0e58ee2365360e431c49070ad904f851cb7fa9 Mon Sep 17 00:00:00 2001 From: JianFeeeee Date: Wed, 2 Sep 2026 17:05:01 +0800 Subject: [PATCH] =?UTF-8?q?plugindev:=20=E6=A8=A1=E6=9D=BF=E6=94=AF?= =?UTF-8?q?=E6=8C=81=E4=BA=8B=E4=BB=B6=E7=8E=AF=E6=B6=88=E8=B4=B9=EF=BC=88?= =?UTF-8?q?Part=205=20=E5=AD=90=E8=BF=9B=E7=A8=8B=E4=BE=A7=EF=BC=89?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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 --- tools/plugindev/templates/proc_main.go.tmpl | 169 +++++++++++++++++++- 1 file changed, 168 insertions(+), 1 deletion(-) diff --git a/tools/plugindev/templates/proc_main.go.tmpl b/tools/plugindev/templates/proc_main.go.tmpl index d743ed0..1198253 100644 --- a/tools/plugindev/templates/proc_main.go.tmpl +++ b/tools/plugindev/templates/proc_main.go.tmpl @@ -93,6 +93,12 @@ var ( outputHandlers = map[string]sdk.ToolHandler{} shm []byte + + // 事件环(§3.6):fd 4 = 事件环段 mmap,fd 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++ + } + } +}