diff --git a/internal/plugin/evtring_test.go b/internal/plugin/evtring_test.go index 76c7891..5980748 100644 --- a/internal/plugin/evtring_test.go +++ b/internal/plugin/evtring_test.go @@ -34,6 +34,9 @@ func TestEventRing_BasicWriteAndConsume(t *testing.T) { }, ) go consumer.Run() + // LIFO:先 Stop(打断阻塞的 Read)再 Wait(等 Run 退出), + // 两者都必须在 host.Close(munmap 整个区域)之前完成。 + defer consumer.Wait() defer consumer.Stop() // 订阅 agent_output 事件 @@ -85,11 +88,18 @@ func TestEventRing_OverflowStillDelivers(t *testing.T) { host.EvtfdReadFile(), 0, func(evt *pubsdk.Event) error { - received <- evt + // 非阻塞投递:本用例写入了 8292 条事件,若这里阻塞在 + // channel 上,drainEvents 会卡在 handler 里,Stop 就无法 + // 让 Run 退出。 + select { + case received <- evt: + default: + } return nil }, ) go consumer.Run() + defer consumer.Wait() defer consumer.Stop() select { @@ -125,6 +135,7 @@ func TestEventRing_TypeMaskFiltering(t *testing.T) { }, ) go consumer.Run() + defer consumer.Wait() defer consumer.Stop() unsub := er.Subscribe(pubsdk.EventToolCall) diff --git a/internal/plugin/proc/evtpoll_other.go b/internal/plugin/proc/evtpoll_other.go new file mode 100644 index 0000000..24949e3 --- /dev/null +++ b/internal/plugin/proc/evtpoll_other.go @@ -0,0 +1,12 @@ +//go:build !unix + +package proc + +// pollEvtfd 在非 Unix 平台不可用。 +// +// Windows 的 evtfdReadFile 返回 nil(命名 Event 不走文件抽象), +// 事件环消费者在那些平台不启用;此处返回 (false, nil) 让调用方 +// 退回阻塞读路径,而不是忙转。 +func pollEvtfd(fd int, timeoutMs int) (bool, error) { + return false, nil +} diff --git a/internal/plugin/proc/evtpoll_unix.go b/internal/plugin/proc/evtpoll_unix.go new file mode 100644 index 0000000..5f2361d --- /dev/null +++ b/internal/plugin/proc/evtpoll_unix.go @@ -0,0 +1,34 @@ +//go:build unix + +package proc + +import ( + "errors" + + "golang.org/x/sys/unix" +) + +// pollEvtfd 等待 fd 可读,最多 timeoutMs 毫秒。超时返回 (false, nil)。 +// +// 为什么不用 os.File.SetReadDeadline:eventfd/pipe 经 os.NewFile 包装后 +// **不会**注册进 Go netpoller(os.NewFile 对非 open 得到的 fd 一律按非 +// pollable 处理),Read 退化成阻塞 syscall,SetReadDeadline 返回错误且 +// 不生效。实测表现是 Run 永久卡在 syscall.Read,Stop 无法打断。 +// +// 用 poll(2) 显式加超时,消费循环才能周期性回到 stop 检查。 +func pollEvtfd(fd int, timeoutMs int) (bool, error) { + if fd < 0 { + return false, errors.New("evtfd: 非法 fd") + } + fds := []unix.PollFd{{Fd: int32(fd), Events: unix.POLLIN}} + for { + n, err := unix.Poll(fds, timeoutMs) + if err == unix.EINTR { + continue // 被信号打断:重试,超时预算不变 + } + if err != nil { + return false, err + } + return n > 0, nil + } +} diff --git a/internal/plugin/proc/evtring.go b/internal/plugin/proc/evtring.go index 604164d..42e52b6 100644 --- a/internal/plugin/proc/evtring.go +++ b/internal/plugin/proc/evtring.go @@ -7,6 +7,7 @@ import ( "os" "sync" "sync/atomic" + "time" pubsdk "gitcode.com/JianFeeeee/homeagent-sdk/sdk" ) @@ -182,12 +183,26 @@ type EvtConsumer struct { 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, @@ -195,9 +210,13 @@ func NewEvtConsumer(ringData []byte, evtfd evtfdReader, mask uint32, handler fun 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 { @@ -210,6 +229,7 @@ func (c *EvtConsumer) Run() { c.mu.Lock() c.running = false c.mu.Unlock() + close(c.done) }() buf := make([]byte, 8) @@ -219,8 +239,25 @@ func (c *EvtConsumer) Run() { 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() @@ -231,6 +268,13 @@ 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 } @@ -268,10 +312,14 @@ func (c *EvtConsumer) drainEvents() { } } +// Stop 请求消费者退出。 +// +// 非阻塞:Stop 返回**不代表** Run 已退出(最多 readWakeInterval 后退出)。 +// 若要在 Stop 之后释放 ringData(host.Close 会 munmap 整个区域), +// 必须先 Stop() 再 Wait()。 func (c *EvtConsumer) Stop() { - c.mu.Lock() - defer c.mu.Unlock() - if c.running { - close(c.stop) - } + c.stopOnce.Do(func() { close(c.stop) }) } + +// Wait 阻塞至 Run 退出。返回后 drainEvents 保证不会再访问 ringData。 +func (c *EvtConsumer) Wait() { <-c.done }