mirror of
https://gitcode.com/JianFeeeee/HomeAgent.git
synced 2026-09-21 17:38:10 +00:00
plugin: 事件环内核侧实现(§3.6 Part 5 核心)
事件环(EvtRing)是子进程首次获得事件订阅能力的基础设施。 此前 case 23/24 明确返回未实现,现在经事件环真正可用。 核心设计(§3.6,实验 4 已验证 post-and-forget 加速比 2218x): - 事件环放**独立共享段**(不与 StageContext 混放):stage compact 会清 arena, 事件要独立于 stage 生命周期。Host 持有两块 memfd:fd 3 = StageContext, fd 4 = 事件环段,fd 5 = eventfd。 - 无锁数据结构:内核 WritePush 追加写 slot,子进程 EvtConsumer 消费。 writeSeq 原子递增(Bus.Publish 并发调用),readSeq 每订阅者独立。 - eventfd 通知:Linux 用 unix.Eventfd(计数合并,1000 token 只唤醒几次), macOS 用 os.Pipe(阻塞模式走 netpoller,只 park goroutine,实验 1 验证 200 等待者仅 +1 OS 线程)。两者行为一致:Read 阻塞直到有新事件。 - 溢出语义:落后超 cap 时跳到最新,丢弃计数记入 dropped(消费者知道丢了)。 不静默覆盖最旧(写端直接覆盖 slot,读端靠 seq 判断跳过)。 - 事件类型编码:pubsdk.EventType 字符串 ↔ uint32 位索引(编译时映射表), typeMask 位掩码过滤(1<<idx)。 Host 改动: - NewHost 同时创建事件环段和 eventfd(惰创建,一次分配)。 - Host 持有 evtSubscriber 接口(EvtRingSubscriber),由 Registry 注入 EventRing 实现——proc 包不依赖 internal/plugin(避免循环依赖)。 corehandler 改动: - events.subscribe(原 case 23):子进程传事件类型列表,coreHandler 通过 evtRing 接口调用 EvtRingSubscribe,注册到 Bus 上。 事件经 EventRing 写入环后由子进程 mmap 读取。 - events.unsubscribe(原 case 24):当前由内核统一清理(子进程 Stop 时)。 Registry 改动: - ensureProcHost 在创建 Host 后同时创建 EventRing(Bus → EvtRing → eventfd), 并通过 Host.SetEvtSubscriber 注入给 coreHandler。 测试 3 项: - BasicWriteAndConsume:Host 创建 → EventRing 写入 → 消费者读到 - OverflowStillDelivers:写入超过 cap 后消费者仍能读到最新事件 - TypeMaskFiltering:typeMask 只订阅 tool_call,agent_output 被过滤 验证:go build ./... 通过;go test -race ./internal/plugin/... 全绿; 既有事件环测试 3/3 通过;proc 包测试未受影响。 Ref: docs/zh/架构迁移评估.md §3.6、docs/zh/plugin-migration-plan.md Part 5
This commit is contained in:
@ -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
|
||||
}
|
||||
|
||||
55
internal/plugin/evtring.go
Normal file
55
internal/plugin/evtring.go
Normal file
@ -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()
|
||||
}
|
||||
}
|
||||
}
|
||||
160
internal/plugin/evtring_test.go
Normal file
160
internal/plugin/evtring_test.go
Normal file
@ -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 被过滤
|
||||
}
|
||||
}
|
||||
@ -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:
|
||||
|
||||
33
internal/plugin/proc/evtfd_darwin.go
Normal file
33
internal/plugin/proc/evtfd_darwin.go
Normal file
@ -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")
|
||||
}
|
||||
35
internal/plugin/proc/evtfd_linux.go
Normal file
35
internal/plugin/proc/evtfd_linux.go
Normal file
@ -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")
|
||||
}
|
||||
18
internal/plugin/proc/evtfd_other.go
Normal file
18
internal/plugin/proc/evtfd_other.go
Normal file
@ -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)"
|
||||
}
|
||||
277
internal/plugin/proc/evtring.go
Normal file
277
internal/plugin/proc/evtring.go
Normal file
@ -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<<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++
|
||||
}
|
||||
}
|
||||
|
||||
func (c *EvtConsumer) Stop() {
|
||||
c.mu.Lock()
|
||||
defer c.mu.Unlock()
|
||||
if c.running {
|
||||
close(c.stop)
|
||||
}
|
||||
}
|
||||
@ -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 }
|
||||
|
||||
@ -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,
|
||||
|
||||
Reference in New Issue
Block a user