Files
HomeAgent/internal/plugin/proc/plugin.go
JianFeeeee d027c964e2 proc: Windows 共享内存 + 事件通知适配(Part 6.2 内核侧)
补齐内核侧的 Windows 创建端,与 6.1 的插件侧打开端配对。三平台
(linux/darwin/windows)现在都能构建 internal/plugin/proc。

## Windows 走命名内核对象(无 fd 继承语义)

os/exec 的 ExtraFiles 在 Windows 实现里不被支持,故:
- shmalloc_windows.go:CreateFileMappingW(INVALID_HANDLE_VALUE + 命名 →
  系统页文件支撑的匿名段,不落盘)+ MapViewOfFile
- evtfd_windows.go:CreateEventW 命名 Event 对象 + SetEvent 通知
- shmpass_windows.go:把段名/对象名经环境变量注入子进程
  (HOMEAGENT_SHM_STAGE / HOMEAGENT_SHM_EVTRING / HOMEAGENT_EVT_EVENT)

名字带 PID + 递增序号:多个 homed 实例并存时不能撞名。

Event 与 eventfd 的语义差异:Event 是二元信号,多次 SetEvent 只对应一次
唤醒,不累积。不影响正确性——消费者被唤醒后按 readSeq 追 writeSeq 批量
drain,丢的是"唤醒次数"不是"事件";事件环本身就允许溢出丢弃并让消费者
知道丢了(dropped 计数),通知面从来不是可靠投递语义。

## 传递机制抽象为 shmpass_*.go

Plugin.Start 不再直接构造 ExtraFiles 列表,改为问 Host 要:
  Env:        p.host.procEnvForShm()        // Windows 返回段名,Unix 返回 nil
  ExtraFiles: p.host.procExtraFilesForShm() // Unix 返回 fd 列表,Windows 返回 nil

平台差异被收敛到这一对函数,Plugin/coreHandler/stage 全部平台无关。

## macOS pipe 生命周期修正

原实现只返回读端 fd,写端 *os.File 无人持有 → 可能被 GC 回收 →
读端收到 EOF 而非阻塞 → 消费循环变忙转。改为 pipePair 表同时持有两端,
evtfdClose 一并关闭。

## E2E 测试跟进模板拆分

模板从单文件拆成三个(主体 + unix/windows 挂载),测试需要一并落盘,
否则编译报 attachStageShm undefined。procRuntimeTemplates 表必须与
SDK 仓 proc_runtime.go 的 procRuntimeFiles 一致。

验证:三平台 go build ./internal/plugin/... 通过(gojieba 的 cgo 依赖
导致 internal/memory 在非 linux 失败,与本次无关);
go test -race ./internal/plugin/... 全绿,含 2 项真实模板 E2E。

Ref: docs/zh/架构迁移评估.md §9.2、docs/zh/plugin-migration-plan.md Part 6
2026-09-02 19:07:14 +08:00

233 lines
6.7 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

package proc
import (
"context"
"encoding/json"
"fmt"
"log"
"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,
// 共享段的传递机制按平台不同shmpass_*.go
// Unix 经 ExtraFiles 传继承 fd 3=StageContext, 4=事件环, 5=通知);
// Windows 无 fd 继承语义,改用命名内核对象,名字经环境变量传入。
Env: append(p.env, p.host.procEnvForShm()...),
ExtraFiles: p.host.procExtraFilesForShm(),
ShmSize: p.host.shmSize,
EvtRingSize: evtTotalSize,
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