mirror of
https://gitcode.com/JianFeeeee/HomeAgent.git
synced 2026-09-23 10:28:06 +00:00
事件环(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
231 lines
6.5 KiB
Go
231 lines
6.5 KiB
Go
package proc
|
||
|
||
import (
|
||
"context"
|
||
"encoding/json"
|
||
"fmt"
|
||
"log"
|
||
"os"
|
||
"sync"
|
||
|
||
pubsdk "gitcode.com/JianFeeeee/homeagent-sdk/sdk"
|
||
)
|
||
|
||
// Plugin 是 registry 可加载的子进程插件,与内置插件同构的启停接口。
|
||
//
|
||
// 生命周期:
|
||
//
|
||
// New() 创建(尚未 spawn)
|
||
// Start(core) spawn 子进程 → 握手(传共享段 fd)→ plugin.init → plugin.start
|
||
// (plugin.start 期间插件反向注册工具/阶段/通道)
|
||
// Stop() plugin.stop → 宽限期 → 必要时 Kill
|
||
// Close() 强制结束(registry 卸载/重载路径)
|
||
//
|
||
// **共享段不属于 Plugin**:它属于 Host,被全部子进程插件共享。
|
||
// 若每插件一段,「内核 ctx → 段 → 插件改 → 回读 ctx」在多插件下会退化成
|
||
// 副本模型,lost update 原样复现(§8.4)。
|
||
type Plugin struct {
|
||
name string
|
||
bin string
|
||
dir string
|
||
config map[string]interface{}
|
||
|
||
host *Host
|
||
proc *Process
|
||
handler *coreHandler
|
||
|
||
// env 追加到子进程环境变量(测试用;生产由 registry 按需设置)。
|
||
env []string
|
||
|
||
// onCrash 由 registry 注入,把进程退出喂给 plugin_health.recordCrash(§2.3)。
|
||
onCrash func(name string, err error)
|
||
|
||
stopOnce sync.Once
|
||
}
|
||
|
||
// New 创建子进程插件(不启动进程)。
|
||
//
|
||
// host 必须是全部子进程插件共用的实例(由 registry 创建一次)。
|
||
func New(name, bin, dir string, config map[string]interface{}, host *Host, onCrash func(string, error)) *Plugin {
|
||
return &Plugin{
|
||
name: name,
|
||
bin: bin,
|
||
dir: dir,
|
||
config: config,
|
||
host: host,
|
||
onCrash: onCrash,
|
||
}
|
||
}
|
||
|
||
// Name 实现 sdk.Plugin。
|
||
func (p *Plugin) Name() string { return p.name }
|
||
|
||
// Start 启动子进程并完成注册。
|
||
//
|
||
// core 是内核为该插件构建的能力面(internal/sdk.PluginSDK 天然满足 CoreSDK)。
|
||
func (p *Plugin) Start(core CoreSDK) error {
|
||
if p.host == nil {
|
||
return fmt.Errorf("proc: %s 缺少共享段 Host", p.name)
|
||
}
|
||
|
||
p.handler = &coreHandler{
|
||
sdk: core,
|
||
name: p.name,
|
||
host: p.host,
|
||
locks: p.host.locks,
|
||
evtRing: p.host.evtSubscriber,
|
||
}
|
||
// 反向调用闭包:注册回调时捕获,运行期经 RPC 打到插件进程。
|
||
p.handler.invokeTool = p.invokeTool
|
||
p.handler.invokeStageFn = p.invokeStage
|
||
p.handler.invokeOutput = p.invokeOutput
|
||
|
||
proc, err := Spawn(p.name, p.bin, Options{
|
||
Dir: p.dir,
|
||
Env: p.env,
|
||
// 子进程 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,
|
||
})
|
||
if err != nil {
|
||
return err
|
||
}
|
||
p.proc = proc
|
||
|
||
// plugin.init:构造插件实例
|
||
if _, err := proc.Call(MethodPluginInit, PluginInitParams{
|
||
Name: p.name,
|
||
Config: p.config,
|
||
}); err != nil {
|
||
proc.Kill()
|
||
return fmt.Errorf("proc: %s plugin.init 失败: %w", p.name, err)
|
||
}
|
||
|
||
// plugin.start:插件在此期间反向注册工具/阶段/通道
|
||
if _, err := proc.Call(MethodPluginStart, nil); err != nil {
|
||
proc.Kill()
|
||
return fmt.Errorf("proc: %s plugin.start 失败: %w", p.name, err)
|
||
}
|
||
return nil
|
||
}
|
||
|
||
// Stop 优雅停止(实现 sdk.Plugin)。
|
||
func (p *Plugin) Stop() error {
|
||
var err error
|
||
p.stopOnce.Do(func() {
|
||
if p.proc != nil {
|
||
err = p.proc.Stop()
|
||
}
|
||
})
|
||
return err
|
||
}
|
||
|
||
// Close 强制结束子进程。
|
||
//
|
||
// **这里是真 kill + wait**——对比 cabi 路径的 Close 只做 dlclose,
|
||
// 而 dlclose 对 Go c-shared 是 no-op(§1.1,热重载静默失效的根因)。
|
||
func (p *Plugin) Close() error {
|
||
var err error
|
||
p.stopOnce.Do(func() {
|
||
if p.proc != nil {
|
||
err = p.proc.Kill()
|
||
}
|
||
})
|
||
return err
|
||
}
|
||
|
||
// handleExit 在子进程退出时把信号喂给 plugin_health(§2.3 逻辑复用),
|
||
// 并释放该插件可能持有的 stage 锁。
|
||
//
|
||
// 后者是"锁仲裁回内核"的自愈价值:持锁者死亡不会导致全局死锁,
|
||
// 无需 robust pthread_mutex(实验 9)。
|
||
func (p *Plugin) handleExit(name string, err error) {
|
||
if p.host != nil && p.host.ForceReleaseLock(name) {
|
||
log.Printf("[proc] %s 退出,内核已释放其持有的 stage 锁", name)
|
||
}
|
||
if err != nil && p.onCrash != nil {
|
||
p.onCrash(name, err)
|
||
}
|
||
}
|
||
|
||
// ---- 内核 → 插件的反向调用 ----
|
||
|
||
func (p *Plugin) invokeTool(name string, args map[string]interface{}) (interface{}, error) {
|
||
if p.proc == nil {
|
||
return nil, ErrProcessExited
|
||
}
|
||
raw, err := p.proc.Call(MethodToolInvoke, ToolInvokeParams{Name: name, Args: args})
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
var res ToolInvokeResult
|
||
if err := json.Unmarshal(raw, &res); err != nil {
|
||
return nil, fmt.Errorf("proc: %s 工具 %s 应答解析失败: %w", p.name, name, err)
|
||
}
|
||
return res.Result, nil
|
||
}
|
||
|
||
func (p *Plugin) invokeStage(ctx context.Context, stage string, seq uint64) error {
|
||
if p.proc == nil {
|
||
return ErrProcessExited
|
||
}
|
||
raw, err := p.proc.CallContext(ctx, MethodStageInvoke, StageInvokeParams{
|
||
Stage: stage,
|
||
Seq: seq,
|
||
})
|
||
if err != nil {
|
||
return err
|
||
}
|
||
var res StageInvokeResult
|
||
if len(raw) > 0 {
|
||
if err := json.Unmarshal(raw, &res); err != nil {
|
||
return fmt.Errorf("proc: %s stage %s 应答解析失败: %w", p.name, stage, err)
|
||
}
|
||
}
|
||
if res.DirtyFields > 0 {
|
||
log.Printf("[proc] %s stage %s 改写了 %d 个字段", p.name, stage, res.DirtyFields)
|
||
}
|
||
return nil
|
||
}
|
||
|
||
// invokeOutput 经插件输出通道发送,**同步等待真实结果**。
|
||
//
|
||
// 这是 §9.4 的根治:C ABI 下 cgo 不可嵌套,只能异步 fire-and-forget,
|
||
// 导致 output_send 永远返回 {status:queued} + err=nil,模型永远以为发送成功
|
||
// (现网 7 天内 2 次消息实际发不出)。进程模型下 RPC 天然可等应答。
|
||
func (p *Plugin) invokeOutput(channel string, args map[string]interface{}) (interface{}, error) {
|
||
if p.proc == nil {
|
||
return nil, ErrProcessExited
|
||
}
|
||
raw, err := p.proc.Call(MethodOutputInvoke, OutputInvokeParams{
|
||
Channel: channel,
|
||
Args: args,
|
||
})
|
||
if err != nil {
|
||
return nil, err // 真实失败上报,模型可感知并重试
|
||
}
|
||
if len(raw) == 0 {
|
||
return map[string]interface{}{"status": "sent"}, nil
|
||
}
|
||
var res map[string]interface{}
|
||
if err := json.Unmarshal(raw, &res); err != nil {
|
||
return map[string]interface{}{"status": "sent"}, nil
|
||
}
|
||
if _, ok := res["status"]; !ok {
|
||
res["status"] = "sent"
|
||
}
|
||
return res, nil
|
||
}
|
||
|
||
// 编译期确认 Plugin 具备 registry 需要的启停形状。
|
||
var _ interface {
|
||
Name() string
|
||
Stop() error
|
||
Close() error
|
||
} = (*Plugin)(nil)
|
||
|
||
// 引用一下公开 SDK,确保本文件的类型假设与它同版本。
|
||
var _ = pubsdk.StageScopeGlobal
|