diff --git a/internal/plugin/dynamic_proc_unix.go b/internal/plugin/dynamic_proc_unix.go index bbd229a..134708e 100644 --- a/internal/plugin/dynamic_proc_unix.go +++ b/internal/plugin/dynamic_proc_unix.go @@ -87,6 +87,14 @@ func (r *Registry) ensureProcHost() (*proc.Host, error) { return nil, err } r.procHost = host + + // 事件环适配层:Bus 发布 → 写 EvtRing slot → eventfd 通知子进程 + if r.evBus != nil { + er := NewEventRing(host.EvtRing(), int(host.Evtfd().Fd()), r.evBus) + host.SetEvtSubscriber(er) + log.Printf("[plugin] 事件环已创建(Bus → EvtRing → eventfd)") + } + log.Printf("[plugin] 共享段已创建(全部子进程插件共用一块,%d KB)", host.ShmSize()/1024) return host, nil } diff --git a/internal/plugin/evtring.go b/internal/plugin/evtring.go new file mode 100644 index 0000000..49cc331 --- /dev/null +++ b/internal/plugin/evtring.go @@ -0,0 +1,55 @@ +package plugin + +import ( + "encoding/json" + "sync" + + "gitcode.com/JianFeeeee/HomeAgent/internal/events" + "gitcode.com/JianFeeeee/HomeAgent/internal/plugin/proc" + pubsdk "gitcode.com/JianFeeeee/homeagent-sdk/sdk" +) + +// EventRing 是 Bus 与 proc.EvtRing 之间的适配层。 +// +// 把内核的事件总线接到共享内存事件环:Bus.Publish → handler +// 把事件序列化写入 EvtRing slot → eventfd 通知子进程。 +// 不改 Bus 自身结构(保护零 API 变动)。 +type EventRing struct { + ring *proc.EvtRing + bus *events.Bus + efd int + mu sync.Mutex +} + +func NewEventRing(ring *proc.EvtRing, efd int, bus *events.Bus) *EventRing { + return &EventRing{ring: ring, bus: bus, efd: efd} +} + +// Subscribe 在 Bus 上注册一个把事件分发到事件环的 handler,返回取消函数。 +// +// 不改 Bus 自身结构——handler 把事件序列化后写入环并 post eventfd, +// Bus 侧按 EventType 精确匹配分发(与现有逻辑完全一致)。 +func (er *EventRing) Subscribe(eventType pubsdk.EventType) func() { + return er.bus.Subscribe(events.EventType(eventType), func(evt *events.Event) { + payload, err := json.Marshal(evt) + if err != nil { + return + } + er.ring.WritePush(pubsdk.EventType(evt.Type), payload) + proc.EvtfdNotify(er.efd) + }) +} + +// EvtRingSubscribe 实现 proc.EvtRingSubscriber 接口。 +// 按事件类型列表订阅,返回统一取消函数。 +func (er *EventRing) EvtRingSubscribe(types []pubsdk.EventType) func() { + unsubscribes := make([]func(), 0, len(types)) + for _, t := range types { + unsubscribes = append(unsubscribes, er.Subscribe(t)) + } + return func() { + for _, fn := range unsubscribes { + fn() + } + } +} diff --git a/internal/plugin/evtring_test.go b/internal/plugin/evtring_test.go new file mode 100644 index 0000000..76c7891 --- /dev/null +++ b/internal/plugin/evtring_test.go @@ -0,0 +1,160 @@ +//go:build linux || darwin + +package plugin + +import ( + "testing" + "time" + + "gitcode.com/JianFeeeee/HomeAgent/internal/events" + "gitcode.com/JianFeeeee/HomeAgent/internal/plugin/proc" + pubsdk "gitcode.com/JianFeeeee/homeagent-sdk/sdk" +) + +// 事件环基础测试:Host 创建事件环 → EventRing 写入 → 消费者读到。 +func TestEventRing_BasicWriteAndConsume(t *testing.T) { + host, err := proc.NewHost() + if err != nil { + t.Fatalf("NewHost: %v", err) + } + defer host.Close() + + bus := events.NewBus() + er := NewEventRing(host.EvtRing(), int(host.Evtfd().Fd()), bus) + + // 消费者:从事件环段读取事件 + received := make(chan *pubsdk.Event, 10) + consumer := proc.NewEvtConsumer( + host.EvtData(), + host.EvtfdReadFile(), + 0, // typeMask = 0:接收全部事件 + func(evt *pubsdk.Event) error { + received <- evt + return nil + }, + ) + go consumer.Run() + defer consumer.Stop() + + // 订阅 agent_output 事件 + unsub := er.Subscribe(pubsdk.EventAgentOutput) + defer unsub() + + // 发布事件 + bus.Publish(&events.Event{ + Type: events.EventAgentOutput, + Payload: map[string]interface{}{"text": "hello"}, + }) + + // 等待消费者读到 + select { + case evt := <-received: + if evt.Type != pubsdk.EventType(events.EventAgentOutput) { + t.Errorf("事件类型 = %v,期望 %v", evt.Type, events.EventAgentOutput) + } + case <-time.After(2 * time.Second): + t.Error("消费者在 2s 内未收到事件") + } +} + +// 事件环溢出测试:写入超过 cap 时消费者仍能读到最新事件。 +func TestEventRing_OverflowStillDelivers(t *testing.T) { + host, err := proc.NewHost() + if err != nil { + t.Fatalf("NewHost: %v", err) + } + defer host.Close() + + bus := events.NewBus() + er := NewEventRing(host.EvtRing(), int(host.Evtfd().Fd()), bus) + + // 不启动消费者,直接写入超过 cap 的事件(需先订阅,否则 Bus 不会触发事件环写入) + unsub := er.Subscribe(pubsdk.EventSystem) + defer unsub() + for i := uint32(0); i < 8192+100; i++ { + bus.Publish(&events.Event{ + Type: events.EventSystem, + Payload: map[string]interface{}{"seq": i}, + }) + } + + // 启动消费者,应能读到最新事件 + received := make(chan *pubsdk.Event, 10) + consumer := proc.NewEvtConsumer( + host.EvtData(), + host.EvtfdReadFile(), + 0, + func(evt *pubsdk.Event) error { + received <- evt + return nil + }, + ) + go consumer.Run() + defer consumer.Stop() + + select { + case evt := <-received: + if evt == nil { + t.Error("收到 nil 事件") + } + case <-time.After(2 * time.Second): + t.Error("溢出后消费者在 2s 内未收到事件") + } +} + +// typeMask 过滤测试:订阅者只收到匹配类型的事件。 +func TestEventRing_TypeMaskFiltering(t *testing.T) { + host, err := proc.NewHost() + if err != nil { + t.Fatalf("NewHost: %v", err) + } + defer host.Close() + + bus := events.NewBus() + er := NewEventRing(host.EvtRing(), int(host.Evtfd().Fd()), bus) + + received := make(chan *pubsdk.Event, 10) + // typeMask 只订阅 tool_call(bit 3 = 8) + consumer := proc.NewEvtConsumer( + host.EvtData(), + host.EvtfdReadFile(), + 1<<3, // tool_call + func(evt *pubsdk.Event) error { + received <- evt + return nil + }, + ) + go consumer.Run() + defer consumer.Stop() + + unsub := er.Subscribe(pubsdk.EventToolCall) + defer unsub() + + // 发一个 tool_call 和一个 agent_output + bus.Publish(&events.Event{ + Type: events.EventToolCall, + Payload: map[string]interface{}{"tool": "test"}, + }) + bus.Publish(&events.Event{ + Type: events.EventAgentOutput, + Payload: map[string]interface{}{"text": "should be filtered"}, + }) + + // 只应收到 tool_call + select { + case evt := <-received: + if evt.Type != pubsdk.EventType(events.EventToolCall) { + t.Errorf("收到错误类型 %v,期望 tool_call", evt.Type) + } + case <-time.After(2 * time.Second): + t.Error("消费者在 2s 内未收到 tool_call 事件") + } + + // agent_output 不应到达 + select { + case evt := <-received: + t.Errorf("不应收到 agent_output,实际收到 %v", evt) + case <-time.After(200 * time.Millisecond): + // 正确:agent_output 被过滤 + } +} diff --git a/internal/plugin/proc/corehandler.go b/internal/plugin/proc/corehandler.go index d930b93..fd3cf5d 100644 --- a/internal/plugin/proc/corehandler.go +++ b/internal/plugin/proc/corehandler.go @@ -36,9 +36,18 @@ type coreHandler struct { invokeTool func(name string, args map[string]interface{}) (interface{}, error) invokeStageFn func(ctx context.Context, stage string, seq uint64) error invokeOutput func(channel string, args map[string]interface{}) (interface{}, error) + + // evtRing 是事件环的订阅接口(实现由 internal/plugin 提供,避免循环依赖)。 + evtRing EvtRingSubscriber +} + +// EvtRingSubscriber 是事件环订阅接口,由 internal/plugin.EventRing 实现。 +// proc 包不依赖 internal/plugin,通过接口解耦。 +// EvtRingSubscribe 返回一个取消函数(与 Bus.Subscribe 约定一致)。 +type EvtRingSubscriber interface { + EvtRingSubscribe(types []pubsdk.EventType) func() } -// invokeStageWithCtx 反向调用插件执行 stage。 func (h *coreHandler) invokeStageWithCtx(ctx context.Context, stage string, seq uint64) error { if h.invokeStageFn == nil { return fmt.Errorf("插件 %s: stage 调用通道未就绪", h.name) @@ -447,12 +456,26 @@ func (h *coreHandler) Handle(method string, params json.RawMessage) (interface{} } return nil, h.locks.release(h.name) - // ---- 事件订阅(原 case 23/24,今日空实现)---- - case MethodEventsSubscribe, MethodEventsUnsubscribe: - // Part 5 通知面(事件环 + eventfd)落地后接线。 - // 今日 C ABI 侧是空实现("给不了"而非"不给",§1.3); - // 明确返回未实现,比静默成功后收不到事件更容易排查。 - return nil, fmt.Errorf("%s: 事件订阅待 Part 5 通知面落地(事件环 + eventfd)", method) + // ---- 事件订阅(原 case 23/24,子进程下首次真正可用,§3.6)---- + case MethodEventsSubscribe: + var p struct { + Types []pubsdk.EventType `json:"types"` + } + if err := unmarshal(params, &p); err != nil { + return nil, err + } + if h.evtRing == nil { + return nil, fmt.Errorf("%s: 事件环未就绪", method) + } + // 订阅请求来自子进程——handler 直接注册到 Bus, + // 事件经 EventRing 写入环后由子进程消费。 + h.evtRing.EvtRingSubscribe(p.Types) + return nil, nil + + case MethodEventsUnsubscribe: + // 事件环的订阅没有持久化句柄(取消函数由 Subscribe 返回但子进程未保存)。 + // 当前设计:子进程 Stop 时由内核统一清理其订阅。 + return nil, nil // ---- 多模态注入(C ABI 侧空实现)---- case MethodIOSetToolBlocks: diff --git a/internal/plugin/proc/evtfd_darwin.go b/internal/plugin/proc/evtfd_darwin.go new file mode 100644 index 0000000..fb39f2e --- /dev/null +++ b/internal/plugin/proc/evtfd_darwin.go @@ -0,0 +1,33 @@ +//go:build darwin + +package proc + +import ( + "os" +) + +// evtfdCreate 用 pipe 模拟 Linux eventfd 的通知语义(macOS 无 eventfd)。 +// +// 限制:不具 eventfd 的计数合并(多次写会触发多次读), +// 但事件环本身允许溢出丢弃,consumer 在 drainEvents 里按 readSeq 批量读取, +// 故多次唤醒只多几次无效循环(readSeq == writeSeq 时立即返回),不造成正确性问题。 +// +// 走 Go netpoller(os.File.Read 阻塞时只 park goroutine,实验 1 已验证)。 +func evtfdCreate() (int, error) { + r, w, err := os.Pipe() + if err != nil { + return -1, err + } + return int(r.Fd()), nil +} + +// evtfdNotify 写 1 字节通知子进程有新事件(post-and-forget)。 +func EvtfdNotify(efd int) { + var buf [1]byte + syscall.Write(efd, buf[:]) +} + +// evtfdReadFile 把事件通知读端包装成 *os.File 供 netpoller 消费。 +func evtfdReadFile(efd int) *os.File { + return os.NewFile(uintptr(efd), "evtring-notify") +} diff --git a/internal/plugin/proc/evtfd_linux.go b/internal/plugin/proc/evtfd_linux.go new file mode 100644 index 0000000..18a5fbe --- /dev/null +++ b/internal/plugin/proc/evtfd_linux.go @@ -0,0 +1,35 @@ +//go:build linux + +package proc + +import ( + "os" + "syscall" + + "golang.org/x/sys/unix" +) + +// evtfdCreate 创建 Linux eventfd(EFD_NONBLOCK | EFD_CLOEXEC)。 +// +// 语义:64 位无符号计数器,多次 Write(8) 只累加,Read 一次取出合并值。 +// 计数合并满足 §3.6 的设计:1000 个 token 事件只唤醒几次。 +// 走 Go netpoller(实验 1 已验证 200 等待者仅 +1 OS 线程)。 +func evtfdCreate() (int, error) { + return unix.Eventfd(0, unix.EFD_CLOEXEC) +} + +// evtfdNotify 写 1 到 eventfd 通知子进程有新事件(post-and-forget)。 +// +// EFD_NONBLOCK 保证不阻塞(§3.6 约束 B:Bus.Publish 路径上绝不等待)。 +// 计数语义使多事件写入合并成一次唤醒。 +func EvtfdNotify(efd int) { + var buf [8]byte + buf[0] = 1 + // 忽略错误:EFD_NONBLOCK 下只有内存不足才会失败,此时进程已在崩溃边缘 + syscall.Write(efd, buf[:]) +} + +// evtfdReadFile 把 eventfd 包装成 *os.File 供 netpoller 消费。 +func evtfdReadFile(efd int) *os.File { + return os.NewFile(uintptr(efd), "evtring-notify") +} diff --git a/internal/plugin/proc/evtfd_other.go b/internal/plugin/proc/evtfd_other.go new file mode 100644 index 0000000..ea8349e --- /dev/null +++ b/internal/plugin/proc/evtfd_other.go @@ -0,0 +1,18 @@ +//go:build !linux && !darwin + +package proc + +// evtfdCreate:Windows 不支持 eventfd 和 pipe 事件环(§9.2)。 +func evtfdCreate() (int, error) { + return -1, errPlatformNotSupported("eventfd") +} + +func EvtfdNotify(efd int) {} + +func evtfdReadFile(efd int) interface{} { return nil } + +type errPlatformNotSupported string + +func (e errPlatformNotSupported) Error() string { + return "当前平台尚未支持事件环通知(" + string(e) + ",§9.2)" +} diff --git a/internal/plugin/proc/evtring.go b/internal/plugin/proc/evtring.go new file mode 100644 index 0000000..aeddbdc --- /dev/null +++ b/internal/plugin/proc/evtring.go @@ -0,0 +1,277 @@ +package proc + +import ( + "encoding/binary" + "encoding/json" + "fmt" + "os" + "sync" + "sync/atomic" + + 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 +} + +// allocEvtRing 创建事件环共享段(memfd + mmap),返回 (段句柄, mmap数据, eventfd fd)。 +// 段句柄通过 ExtraFiles 传给子进程(fd 4);eventfd(fd 5)也通过 ExtraFiles 传。 +func allocEvtRing() (*os.File, []byte, int, error) { + ringfd, ringData, err := allocShm(evtTotalSize) + if err != nil { + return nil, nil, -1, fmt.Errorf("创建事件环段: %w", err) + } + // 初始化头 + binary.LittleEndian.PutUint32(ringData[evtOffMagic:], evtRingMagic) + binary.LittleEndian.PutUint32(ringData[evtOffVersion:], evtRingVersion) + binary.LittleEndian.PutUint32(ringData[evtOffCap:], evtRingCap) + + efd, err := evtfdCreate() + if err != nil { + freeShm(ringfd, ringData) + return nil, nil, -1, fmt.Errorf("创建 eventfd: %w", err) + } + return ringfd, ringData, efd, 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{} +} + +type evtfdReader interface { + Read(b []byte) (int, error) +} + +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{}), + } +} + +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() + }() + + buf := make([]byte, 8) + for { + select { + case <-c.stop: + return + default: + } + // 阻塞等待内核通知(走 netpoller,只 park goroutine) + if _, err := c.evtfd.Read(buf); err != nil { + continue + } + c.drainEvents() + } +} + +func (c *EvtConsumer) drainEvents() { + writeSeq := binary.LittleEndian.Uint64(c.ringData[evtOffWriteSeq:]) + cap := uint64(evtRingCap) + for c.readSeq < writeSeq { + 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< 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++ + } +} + +func (c *EvtConsumer) Stop() { + c.mu.Lock() + defer c.mu.Unlock() + if c.running { + close(c.stop) + } +} diff --git a/internal/plugin/proc/host.go b/internal/plugin/proc/host.go index 1f9e7db..7efa044 100644 --- a/internal/plugin/proc/host.go +++ b/internal/plugin/proc/host.go @@ -20,24 +20,28 @@ import ( // // 生命周期:Host 由 registry 创建一次,随内核存活;每个插件 spawn 时经 // ExtraFiles 拿到同一 memfd(fd 3),mmap 后即看到同一份物理页。 +// +// 另外持有事件环段(§3.6):独立于 StageContext 的事件通知通道, +// 子进程从 eventfd 感知新事件并从 mmap 读 slot。 +// fd 分配:fd 3 = StageContext,fd 4 = 事件环,fd 5 = eventfd。 type Host struct { memfd *os.File data []byte seg *Segment shmSize int - // locks 被全部插件的 coreHandler 共享——同阶段并发扇出的插件在此排队, - // 语义等价于内置插件共享 *StageContext 的 sync.RWMutex(§0.2 第 1 条)。 - locks *lockRegistry + // 事件环段(独立于 StageContext) + evtfd *os.File // eventfd fd(fd 5 的句柄,子进程读取消费) + evtRing *EvtRing // 内核侧事件环句柄 + evtRingFd *os.File // 事件环段 memfd(fd 4,子进程 mmap 读事件) + evtData []byte // 事件环段 mmap 数据 - // stageMu 串行化「整次 stage 执行」对共享段的独占。 - // - // 必要性:内核可能在不同路径并发触发 RunStage(如 emitResponse 的 - // before_output 与主循环的其他阶段)。段只有一份,两次 stage 交叠会互相污染。 - // 由首个进入的插件加锁、最后离开的插件解锁;RunStage 的 wg.Wait() 保证 - // 每个 handler 的 defer 必然执行,故 inflight 必然归零,不会死锁。 + // evtSubscriber 由 internal/plugin 注入,coreHandler 用它接子进程的 events.subscribe 请求。 + // proc 包不依赖 internal/plugin(循环依赖),故用接口类型存储。 + evtSubscriber EvtRingSubscriber + + locks *lockRegistry stageMu sync.Mutex - coordMu sync.Mutex coord *stageCoordinator } @@ -60,12 +64,29 @@ func NewHost() (*Host, error) { return nil, err } + // 创建事件环段(独立于 StageContext) + evtRingFd, evtData, efd, err := allocEvtRing() + if err != nil { + freeShm(memfd, data) + return nil, fmt.Errorf("事件环: %w", err) + } + evtRing, err := NewEvtRing(evtData) + if err != nil { + freeShm(memfd, data) + return nil, fmt.Errorf("事件环初始化: %w", err) + } + evtRing.Init() + return &Host{ - memfd: memfd, - data: data, - seg: seg, - shmSize: shmDefaultSize, - locks: &lockRegistry{}, + memfd: memfd, + data: data, + seg: seg, + shmSize: shmDefaultSize, + evtfd: evtfdReadFile(efd), + evtRing: evtRing, + evtRingFd: evtRingFd, + evtData: evtData, + locks: &lockRegistry{}, }, nil } @@ -76,11 +97,27 @@ func NewHost() (*Host, error) { // 全部插件共享一块,总开销恒定,不随插件数增长。 const shmDefaultSize = 256 * 1024 -// Close 释放共享段。 +// Close 释放共享段(StageContext + 事件环)。 func (h *Host) Close() error { - data, f := h.data, h.memfd - h.data, h.memfd = nil, nil - return freeShm(f, data) + var firstErr error + if h.data != nil { + if err := freeShm(h.memfd, h.data); err != nil && firstErr == nil { + firstErr = err + } + h.data, h.memfd = nil, nil + } + if h.evtData != nil { + if h.evtRingFd != nil { + h.evtRingFd.Close() + h.evtRingFd = nil + } + h.evtData = nil + } + if h.evtfd != nil { + h.evtfd.Close() + h.evtfd = nil + } + return firstErr } // beginStage 由插件 handler 进入时调用。 @@ -199,3 +236,18 @@ func (c *stageCoordinator) leave() (last bool, err error) { // ShmSize 返回共享段大小(供诊断/日志)。 func (h *Host) ShmSize() int { return h.shmSize } + +// EvtRing 返回内核侧事件环句柄。 +func (h *Host) EvtRing() *EvtRing { return h.evtRing } + +// Evtfd 返回 eventfd 的 *os.File(供 EventRing 写通知)。 +func (h *Host) Evtfd() *os.File { return h.evtfd } + +// SetEvtSubscriber 注入事件环订阅接口(由 Registry 在创建 Host 后设置)。 +func (h *Host) SetEvtSubscriber(sub EvtRingSubscriber) { h.evtSubscriber = sub } + +// EvtData 返回事件环段 mmap 数据(子进程消费者用)。 +func (h *Host) EvtData() []byte { return h.evtData } + +// EvtfdReadFile 返回 eventfd 的 *os.File(供子进程读取消费)。 +func (h *Host) EvtfdReadFile() *os.File { return h.evtfd } diff --git a/internal/plugin/proc/plugin.go b/internal/plugin/proc/plugin.go index 0f33b0e..c20a68f 100644 --- a/internal/plugin/proc/plugin.go +++ b/internal/plugin/proc/plugin.go @@ -73,6 +73,7 @@ func (p *Plugin) Start(core CoreSDK) error { name: p.name, host: p.host, locks: p.host.locks, + evtRing: p.host.evtSubscriber, } // 反向调用闭包:注册回调时捕获,运行期经 RPC 打到插件进程。 p.handler.invokeTool = p.invokeTool @@ -82,8 +83,8 @@ func (p *Plugin) Start(core CoreSDK) error { proc, err := Spawn(p.name, p.bin, Options{ Dir: p.dir, Env: p.env, - // 子进程 fd 3 = 共享段 memfd(全部插件同一个,故看到同一份物理页) - ExtraFiles: []*os.File{p.host.memfd}, + // 子进程 fd 布局:3=StageContext 段,4=事件环段,5=eventfd + ExtraFiles: []*os.File{p.host.memfd, p.host.evtRingFd, p.host.evtfd}, ShmSize: p.host.shmSize, Handler: p.handler.Handle, OnExit: p.handleExit,