Files
homeagent-sdk/tools/plugindev/templates/proc_main.go.tmpl
JianFeeeee 9f844123fe plugindev: entry 语义收敛 + 删 C ABI 工具链 + Windows 共享内存适配(Part 6.1)
## entry 不再是通道开关 —— 外部插件零改动的关键

17 个存量插件的 plg.json 都写着 "entry": "plugin.so"。若把 entry 当通道
开关,迁移就得改 17 个文件,而「外部插件零改动」是本次迁移的硬约束。

改法:Go 插件一律产出 plugin.bin,不看 entry 值。isProcEntry 删除,
resolveBuild 去掉 proc 参数。entry 现在只剩区分 Lua(main.lua)一个用途。

实测:weather 的 plg.json 一行不改(仍写 plugin.so),plugindev build
直接产出三平台 plugin.bin。

## Windows 不再是能力退化的第三套实现(§9.2 的正解)

C ABI 时代 Windows 是独立的第三套 ABI:dynamic_dll_windows.go 的 stage
只下发 3 个字段(raw_message/user_id/phase)且完全没有写回,sanitizer
这类改写型插件在 Windows 上静默失效,且无任何运行时警告。

现在 Windows 与 Unix 共用同一份 RPC 逻辑与同一份共享段布局。平台差异
收敛到三个挂载函数:
- Unix(linux/darwin/freebsd):内核经 ExtraFiles 传继承 fd(3=StageContext
  段,4=事件环段,5=eventfd/pipe)
- Windows:没有 fd 继承语义(os/exec 的 ExtraFiles 在 Windows 不支持),
  改用命名内核对象——父进程 CreateFileMapping/CreateEvent 建带名字的对象,
  子进程 OpenFileMappingW/OpenEventW 按同名打开。名字经环境变量传入而非
  硬编码:多个 homed 实例并存时不能撞名。

Windows 绑定用 syscall.NewLazyDLL 而非 golang.org/x/sys/windows:
OpenFileMappingW/OpenEventW 未被标准库 syscall 导出,而引入 x/sys 会给
**每个插件的 go.mod** 加一个新依赖,违反「插件仅依赖公开 SDK」。
LazyDLL 属标准库,零新增依赖。

新增 evtWaiter 接口抽象等待语义:eventfd 是计数器(多事件合并成一次
唤醒),Windows Event 是二元信号。不影响正确性——消费者被唤醒后按
readSeq 追 writeSeq 批量 drain,一次唤醒能处理累积的全部事件。

模板拆成三个文件:
  proc_main.go.tmpl          平台无关(RPC + 共享段布局 + stage + 事件环消费)
  proc_shm_unix.go.tmpl      继承 fd 挂载
  proc_shm_windows.go.tmpl   命名对象挂载

## 删除 C ABI 工具链

templates.go 1296 → 516 行:
- tmplBridge(Windows DLL bridge)      -265 行
- tmplLinuxBridge(Linux c-shared)     -457 行
- tmplPluginInitC(C 入口)              -57 行
另删 generateBridge / detectWindowsCC(MinGW 探测)/ tmplCABIHeader /
InitData.CABIVersion+CABIHeader。

交叉编译不再需要目标平台 C 工具链——这是 -buildmode=c-shared 退场的
连带收益(§3.1)。

## 测试

15 项全过,新增 4 项守护迁移不变量:
- AllPlatformsProduceBin:6 个 GOOS/GOARCH 组合统一产出 plugin.bin
- LuaIsSeparatePath:Lua 仍走解释器路径
- UnsupportedOSErrors:不支持平台明确报错,不静默产出错误产物
- NoCABIResiduals:代码中不得再出现 c-shared / CGO_ENABLED=1 /
  detectWindowsCC / tmplLinuxBridge / tmplPluginInitC(注释除外)
- IgnoresEntryForGoPlugins:isProcEntry 必须已删除

验证:go build/vet/test 全通过;三平台交叉编译产出 plugin.bin;
git diff sdk/ 为空(接口冻结)。

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

1292 lines
34 KiB
Cheetah
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 main
// 子进程插件入口(由 plugindev 自动生成,请勿手工编辑)。
//
// 与旧 C ABI bridgez_bridge_gen.go的关键差异
// - **零 cgo**:没有 //export、没有 C.CString/C.free、不需要 -buildmode=c-shared
// - 51 个整数 method id 换成可读 method 名(内核侧 internal/plugin/proc/protocol.go
// - StageContext 走共享内存fd 3 传入的 memfd插件在同一份状态上读改写
// 消除副本模型的 lost update实测 35.8~36.8% → 0
// - **插件业务代码零改动**:仍是 NewPluginFactory + sdk.PluginSDK
//
// 设计依据docs/zh/架构迁移评估.md 第三章
import (
"bufio"
"encoding/binary"
"encoding/json"
"fmt"
"log"
"os"
"sync"
sdk "gitcode.com/JianFeeeee/homeagent-sdk/sdk"
)
// ---- 协议常量(须与内核 internal/plugin/proc/protocol.go 一致)----
const procProtocolVersion = 1
// ---- 共享段布局(须与内核 internal/plugin/proc/shm.go 一致)----
const (
shmStageFieldCount = 18
shmSliceSize = 8
shmOffMagic = 0
shmOffVersion = 4
shmOffArenaBase = 8
shmOffArenaCap = 12
shmOffArenaUsed = 16
shmOffCtxBase = 20
shmOffSeq = 24
shmMagic = 0x48415348
shmVersion = 1
)
// 字段索引(顺序须与内核 stageField 枚举一致)
const (
fRawMessage = iota
fUserID
fGroupID
fLLMText
fReasoningContent
fFinalText
fResponse
fPhase
fContextMsgs
fToolCalls
fToolResults
fMemory
fTokenUsage
fErrors
fExtraMediaBlocks
fExtraMediaType
fExtraInputSource
fExtraOutputChannel
)
const (
flagNoMemory = 0
flagResponseSet = 1
)
// ---- 全局状态 ----
var (
stdoutW = bufio.NewWriter(os.Stdout)
writeMu sync.Mutex
nextID uint64
pendMu sync.Mutex
pending = map[uint64]chan rpcResponse{}
plg sdk.Plugin
pluginSDK *sdk.PluginSDK
pluginName string
handlerMu sync.RWMutex
toolHandlers = map[string]sdk.ToolHandler{}
stageHandlers = map[string]sdk.StageHandler{}
outputHandlers = map[string]sdk.ToolHandler{}
shm []byte
// 事件环§3.6fd 4 = 事件环段 mmapfd 5 = eventfd 读端
evtRingData []byte
evtNotifier evtWaiter
evtHandlers = map[uint32]func(*sdk.Event){}
evtHandlerMu sync.RWMutex
)
type rpcRequest struct {
ID uint64 `json:"id,omitempty"`
Method string `json:"method"`
Params json.RawMessage `json:"params,omitempty"`
}
type rpcResponse struct {
ID uint64 `json:"id"`
Result json.RawMessage `json:"result,omitempty"`
Error string `json:"error,omitempty"`
}
func writeFrame(v interface{}) {
b, err := json.Marshal(v)
if err != nil {
log.Printf("序列化帧失败: %v", err)
return
}
writeMu.Lock()
stdoutW.Write(b)
stdoutW.WriteByte('\n')
stdoutW.Flush()
writeMu.Unlock()
}
func respond(id uint64, result interface{}) {
resp := rpcResponse{ID: id}
if result != nil {
if b, err := json.Marshal(result); err == nil {
resp.Result = b
}
}
writeFrame(&resp)
}
func respondErr(id uint64, err error) {
writeFrame(&rpcResponse{ID: id, Error: err.Error()})
}
// callCore 反向调用内核(对应旧 bridge 的 callVoid/callString
func callCore(method string, params interface{}) (json.RawMessage, error) {
pendMu.Lock()
nextID++
id := nextID
ch := make(chan rpcResponse, 1)
pending[id] = ch
pendMu.Unlock()
var raw json.RawMessage
if params != nil {
b, err := json.Marshal(params)
if err != nil {
return nil, err
}
raw = b
}
writeFrame(&rpcRequest{ID: id, Method: method, Params: raw})
resp := <-ch
if resp.Error != "" {
return nil, fmt.Errorf("%s", resp.Error)
}
return resp.Result, nil
}
func callCoreVoid(method string, params interface{}) error {
_, err := callCore(method, params)
return err
}
// ---- 共享段访问插件作者永远不接触这些§3.4----
func shmU32(off int) uint32 { return binary.LittleEndian.Uint32(shm[off:]) }
func shmArenaBase() uint32 { return shmU32(shmOffArenaBase) }
func shmArenaCap() uint32 { return shmU32(shmOffArenaCap) }
func shmCtxBase() uint32 { return shmU32(shmOffCtxBase) }
func shmDescOff(field int) uint32 {
return shmCtxBase() + uint32(field*shmSliceSize)
}
func shmGetDesc(field int) (off, ln uint32) {
o := shmDescOff(field)
return binary.LittleEndian.Uint32(shm[o:]), binary.LittleEndian.Uint32(shm[o+4:])
}
func shmSetDesc(field int, off, ln uint32) {
o := shmDescOff(field)
binary.LittleEndian.PutUint32(shm[o:], off)
binary.LittleEndian.PutUint32(shm[o+4:], ln)
}
func shmFlagsOff() uint32 {
return shmCtxBase() + uint32(shmStageFieldCount*shmSliceSize)
}
func shmGetFlag(bit int) bool { return shm[shmFlagsOff()+uint32(bit)] != 0 }
func shmSetFlag(bit int, v bool) {
b := byte(0)
if v {
b = 1
}
shm[shmFlagsOff()+uint32(bit)] = b
}
func shmRead(field int) []byte {
off, ln := shmGetDesc(field)
if off == 0 && ln == 0 {
return nil
}
if ln == 0 {
return []byte{}
}
base := shmArenaBase()
return shm[base+off : base+off+ln]
}
// shmWrite 在 arena 上 append-only 分配并更新描述符。
// arena 用尽显式报错,不静默截断(与内核侧同一约定)。
func shmWrite(field int, data []byte) error {
if len(data) == 0 {
shmSetDesc(field, 1, 0)
return nil
}
used := shmU32(shmOffArenaUsed)
if used == 0 {
used = 1
}
end := used + uint32(len(data))
if end > shmArenaCap() {
return fmt.Errorf("共享段 arena 空间不足:需要 %d 字节,容量 %d已用 %d",
len(data), shmArenaCap(), used)
}
base := shmArenaBase()
copy(shm[base+used:], data)
binary.LittleEndian.PutUint32(shm[shmOffArenaUsed:], end)
shmSetDesc(field, used, uint32(len(data)))
return nil
}
func shmBumpSeq() {
v := binary.LittleEndian.Uint64(shm[shmOffSeq:])
binary.LittleEndian.PutUint64(shm[shmOffSeq:], v+1)
}
// readStageContext 从共享段构造插件侧原生 StageContext。
// 全 16 个字段可见——C ABI 下只有 10 个§8.3)。
func readStageContext() (*sdk.StageContext, error) {
sc := &sdk.StageContext{}
sc.RawMessage = string(shmRead(fRawMessage))
sc.UserID = string(shmRead(fUserID))
sc.GroupID = string(shmRead(fGroupID))
sc.LLMText = string(shmRead(fLLMText))
sc.ReasoningContent = string(shmRead(fReasoningContent))
sc.FinalText = string(shmRead(fFinalText))
sc.Phase = sdk.Stage(string(shmRead(fPhase)))
sc.NoMemory = shmGetFlag(flagNoMemory)
if shmGetFlag(flagResponseSet) {
r := string(shmRead(fResponse))
sc.Response = &r
}
unmarshalField := func(field int, out interface{}) error {
b := shmRead(field)
if len(b) == 0 {
return nil
}
return json.Unmarshal(b, out)
}
if err := unmarshalField(fContextMsgs, &sc.ContextMsgs); err != nil {
return nil, err
}
if err := unmarshalField(fToolCalls, &sc.ToolCalls); err != nil {
return nil, err
}
if err := unmarshalField(fToolResults, &sc.ToolResults); err != nil {
return nil, err
}
if err := unmarshalField(fMemory, &sc.Memory); err != nil {
return nil, err
}
if err := unmarshalField(fTokenUsage, &sc.TokenUsage); err != nil {
return nil, err
}
if err := unmarshalField(fErrors, &sc.Errors); err != nil {
return nil, err
}
extra := map[string]interface{}{}
for _, pair := range []struct {
field int
key string
}{
{fExtraMediaBlocks, "media_blocks"},
{fExtraMediaType, "media_type"},
{fExtraInputSource, "input_source"},
{fExtraOutputChannel, "output_channel"},
} {
var v interface{}
if err := unmarshalField(pair.field, &v); err != nil {
return nil, err
}
if v != nil {
extra[pair.key] = v
}
}
if len(extra) > 0 {
sc.Extra = extra
}
return sc, nil
}
// stageSnapshot 是 handler 运行前的序列化快照,用于计算脏字段。
//
// ❗ 必须存序列化后的字符串handler 原地改切片元素
// sc.ToolResults[0].Result = x直接持有的 Go 值快照会跟着变,
// 脏字段计算失效——这个坑在修 C ABI 侧的 11.3 时已经踩过一次。
type stageSnapshot struct {
strs map[int]string
jsons map[int]string
response string
responseSet bool
noMemory bool
}
func takeStageSnapshot(sc *sdk.StageContext) *stageSnapshot {
sn := &stageSnapshot{strs: map[int]string{}, jsons: map[int]string{}}
sn.strs[fRawMessage] = sc.RawMessage
sn.strs[fUserID] = sc.UserID
sn.strs[fGroupID] = sc.GroupID
sn.strs[fLLMText] = sc.LLMText
sn.strs[fReasoningContent] = sc.ReasoningContent
sn.strs[fFinalText] = sc.FinalText
sn.strs[fPhase] = string(sc.Phase)
marshal := func(v interface{}, n int) string {
if n == 0 {
return ""
}
b, err := json.Marshal(v)
if err != nil {
return ""
}
return string(b)
}
sn.jsons[fContextMsgs] = marshal(sc.ContextMsgs, len(sc.ContextMsgs))
sn.jsons[fToolCalls] = marshal(sc.ToolCalls, len(sc.ToolCalls))
sn.jsons[fToolResults] = marshal(sc.ToolResults, len(sc.ToolResults))
sn.jsons[fMemory] = marshal(sc.Memory, len(sc.Memory))
sn.jsons[fTokenUsage] = marshal(sc.TokenUsage, len(sc.TokenUsage))
sn.jsons[fErrors] = marshal(sc.Errors, len(sc.Errors))
if sc.Response != nil {
sn.response = *sc.Response
sn.responseSet = true
}
sn.noMemory = sc.NoMemory
return sn
}
// writeStageDirty 只把变更字段写回共享段,返回写回字段数。
//
// **这是消除 lost update 的核心**:只读插件的脏字段集为空 → 零写入 →
// 不可能覆盖其他插件的改写(对照 C ABI 副本模型实测 35.8~36.8% 丢失)。
func writeStageDirty(sc *sdk.StageContext, base *stageSnapshot) (int, error) {
now := takeStageSnapshot(sc)
changed := 0
for field, cur := range now.strs {
if base.strs[field] != cur {
if err := shmWrite(field, []byte(cur)); err != nil {
return changed, err
}
changed++
}
}
for field, cur := range now.jsons {
if base.jsons[field] == cur {
continue
}
if err := shmWrite(field, []byte(cur)); err != nil {
return changed, err
}
changed++
}
if base.responseSet != now.responseSet || base.response != now.response {
if now.responseSet {
if err := shmWrite(fResponse, []byte(now.response)); err != nil {
return changed, err
}
shmSetFlag(flagResponseSet, true)
changed++
}
// Response 置回 nil 不清空内核已设的值:短路语义不应被撑销
}
if base.noMemory != now.noMemory {
shmSetFlag(flagNoMemory, now.noMemory)
changed++
}
if changed > 0 {
shmBumpSeq()
}
return changed, nil
}
// ---- SDK 装配:全部 API 经 RPC 打回内核51 个 method 的插件侧一半)----
func buildPluginSDK(name string) *sdk.PluginSDK {
base := sdk.New(name, procSettings{},
func(toolName string, def sdk.ToolDef, handler sdk.ToolHandler) error {
handlerMu.Lock()
toolHandlers[toolName] = handler
handlerMu.Unlock()
return callCoreVoid("tool.register", map[string]interface{}{
"name": toolName, "def": def,
})
},
func(stage sdk.Stage, handler sdk.StageHandler) {
handlerMu.Lock()
stageHandlers[string(stage)] = handler
handlerMu.Unlock()
if err := callCoreVoid("stage.register", map[string]interface{}{
"stage": string(stage), "scope": "global",
}); err != nil {
log.Printf("注册阶段 %s 失败: %v", stage, err)
}
},
func(apiName string) error {
return callCoreVoid("api.register", map[string]interface{}{"name": apiName})
},
func(chName string, caps int, desc string, def sdk.ChannelDef, handler sdk.ToolHandler) error {
handlerMu.Lock()
outputHandlers[chName] = handler
handlerMu.Unlock()
return callCoreVoid("output.register", map[string]interface{}{
"name": chName, "caps": caps, "desc": desc,
"def": map[string]interface{}{"NoMemory": def.NoMemory},
})
},
)
base.SetIOInjector(procIO{})
base.SetMemoryAPI(procMemory{})
base.SetDocMemoryAPI(procDocMemory{})
base.SetKnowledgeAPI(procKnowledge{})
base.SetLLMAPI(procLLM{})
base.SetSocialAPI(procSocial{})
base.SetTextMemoryAPI(procTextMemory{})
base.SetPluginMgrAPI(procPluginMgr{})
base.SetInputChannelRegistrar(func(chName string, def sdk.ChannelDef) error {
return callCoreVoid("input.register", map[string]interface{}{
"name": chName,
"def": map[string]interface{}{"NoMemory": def.NoMemory},
})
})
return base
}
type procIO struct{}
func (procIO) InjectText(s, c, t string) {
callCoreVoid("io.injectText", map[string]string{"source": s, "channel": c, "text": t})
}
func (procIO) InjectInterruptText(s, c, t string) {
callCoreVoid("io.injectInterrupt", map[string]string{"source": s, "channel": c, "text": t})
}
func (procIO) InjectTextNoMemory(s, c, t string) {
callCoreVoid("io.injectTextNoMem", map[string]string{"source": s, "channel": c, "text": t})
}
func (procIO) InjectInputSync(s, c, t string) string {
raw, err := callCore("io.injectInputSync", map[string]string{"source": s, "channel": c, "text": t})
if err != nil {
return ""
}
var r struct {
Reply string `json:"reply"`
}
json.Unmarshal(raw, &r)
return r.Reply
}
func (procIO) SetToolBlocks(blocks []sdk.ContentBlock) {
if err := callCoreVoid("io.setToolBlocks", map[string]interface{}{"blocks": blocks}); err != nil {
log.Printf("SetToolBlocks: %v", err)
}
}
type procMemory struct{}
func (procMemory) Recall(q []string, d int) ([]sdk.Entity, []sdk.Relation, error) {
raw, err := callCore("memory.recall", map[string]interface{}{"query": q, "depth": d})
if err != nil {
return nil, nil, err
}
var r struct {
Entities []sdk.Entity `json:"entities"`
Relations []sdk.Relation `json:"relations"`
}
if err := json.Unmarshal(raw, &r); err != nil {
return nil, nil, err
}
return r.Entities, r.Relations, nil
}
func (procMemory) Commit(t []sdk.Triple) error {
return callCoreVoid("memory.commit", map[string]interface{}{"triples": t})
}
func (procMemory) Introspect() (map[string]interface{}, error) {
raw, err := callCore("memory.introspect", nil)
if err != nil {
return nil, err
}
var m map[string]interface{}
json.Unmarshal(raw, &m)
return m, nil
}
func (procMemory) MergeEntities(s, t string) (int, error) {
raw, err := callCore("memory.merge", map[string]string{"source": s, "target": t})
if err != nil {
return 0, err
}
var r struct {
Merged int `json:"merged"`
}
json.Unmarshal(raw, &r)
return r.Merged, nil
}
func (procMemory) Purge(c map[string]string, mode string) (int, error) {
raw, err := callCore("memory.purge", map[string]interface{}{"criteria": c, "mode": mode})
if err != nil {
return 0, err
}
var r struct {
Purged int `json:"purged"`
}
json.Unmarshal(raw, &r)
return r.Purged, nil
}
type procDocMemory struct{}
func (procDocMemory) Query(text string, topK int) []*sdk.Doc {
raw, err := callCore("doc.query", map[string]interface{}{"text": text, "top_k": topK})
if err != nil {
return nil
}
var r struct {
Docs []*sdk.Doc `json:"docs"`
}
json.Unmarshal(raw, &r)
return r.Docs
}
func (procDocMemory) Insert(d *sdk.Doc) error {
return callCoreVoid("doc.insert", map[string]interface{}{"doc": d})
}
func (procDocMemory) Remove(id string) {
callCoreVoid("doc.remove", map[string]string{"id": id})
}
func (procDocMemory) Stats() map[string]interface{} {
raw, err := callCore("doc.stats", nil)
if err != nil {
return nil
}
var m map[string]interface{}
json.Unmarshal(raw, &m)
return m
}
type procKnowledge struct{}
func (procKnowledge) Search(q string, topK int) ([]*sdk.Knowledge, error) {
raw, err := callCore("knowledge.search", map[string]interface{}{"query": q, "top_k": topK})
if err != nil {
return nil, err
}
var r struct {
Results []*sdk.Knowledge `json:"results"`
}
json.Unmarshal(raw, &r)
return r.Results, nil
}
func (procKnowledge) Add(name, content string) error {
return callCoreVoid("knowledge.add", map[string]string{"name": name, "content": content})
}
func (procKnowledge) List() ([]string, error) {
raw, err := callCore("knowledge.list", nil)
if err != nil {
return nil, err
}
var r struct {
Names []string `json:"names"`
}
json.Unmarshal(raw, &r)
return r.Names, nil
}
type procTextMemory struct{}
func (procTextMemory) Append(evt sdk.TextEvent) error {
return callCoreVoid("textmemory.append", map[string]interface{}{"event": evt})
}
type procLLM struct{}
func (procLLM) ListSources() []string {
raw, err := callCore("llm.listSources", nil)
if err != nil {
return nil
}
var r struct {
Sources []string `json:"sources"`
}
json.Unmarshal(raw, &r)
return r.Sources
}
func (procLLM) SetSource(name string) error {
return callCoreVoid("llm.setSource", map[string]string{"name": name})
}
func (procLLM) CurrentSource() string {
raw, err := callCore("llm.currentSource", nil)
if err != nil {
return ""
}
var r struct {
Source string `json:"source"`
}
json.Unmarshal(raw, &r)
return r.Source
}
type procSocial struct{}
func (procSocial) GetPerson(name string) (*sdk.PersonProfile, error) {
raw, err := callCore("social.getPerson", map[string]string{"name": name})
if err != nil {
return nil, err
}
var r struct {
Person *sdk.PersonProfile `json:"person"`
}
json.Unmarshal(raw, &r)
return r.Person, nil
}
func (procSocial) GetTrait(name, trait string) (string, bool) {
raw, err := callCore("social.getTrait", map[string]string{"name": name, "trait": trait})
if err != nil {
return "", false
}
var r struct {
Value string `json:"value"`
Found bool `json:"found"`
}
json.Unmarshal(raw, &r)
return r.Value, r.Found
}
func (procSocial) GetRelations(name string) ([]sdk.SocialRelation, error) {
raw, err := callCore("social.getRelations", map[string]string{"name": name})
if err != nil {
return nil, err
}
var r struct {
Relations []sdk.SocialRelation `json:"relations"`
}
json.Unmarshal(raw, &r)
return r.Relations, nil
}
func (procSocial) GetNetwork(name string, depth int) ([]*sdk.PersonProfile, error) {
raw, err := callCore("social.getNetwork", map[string]interface{}{"name": name, "depth": depth})
if err != nil {
return nil, err
}
var r struct {
Network []*sdk.PersonProfile `json:"network"`
}
json.Unmarshal(raw, &r)
return r.Network, nil
}
func (procSocial) ListPersons() ([]string, error) {
raw, err := callCore("social.listPersons", nil)
if err != nil {
return nil, err
}
var r struct {
Persons []string `json:"persons"`
}
json.Unmarshal(raw, &r)
return r.Persons, nil
}
type procPluginMgr struct{}
func (procPluginMgr) ReloadOne(name string) error {
return callCoreVoid("plugin.reloadOne", map[string]string{"name": name})
}
func (procPluginMgr) ListLoadedPlugins() []string {
raw, err := callCore("plugin.listLoaded", nil)
if err != nil {
return nil
}
var r struct {
Plugins []string `json:"plugins"`
}
json.Unmarshal(raw, &r)
return r.Plugins
}
func (procPluginMgr) IsPluginDisabled(name string) bool {
raw, err := callCore("plugin.isDisabled", map[string]string{"name": name})
if err != nil {
return false
}
var r struct {
Disabled bool `json:"disabled"`
}
json.Unmarshal(raw, &r)
return r.Disabled
}
type procSettings struct{}
func (procSettings) Get(key string) (interface{}, error) {
return settingsValue("settings.get", map[string]string{"key": key})
}
func (procSettings) Set(key string, v interface{}) error {
return callCoreVoid("settings.set", map[string]interface{}{"key": key, "value": v})
}
func (procSettings) List(prefix string) ([]string, error) {
return settingsKeys("settings.list", map[string]string{"prefix": prefix})
}
func (procSettings) GetCore(key string) (interface{}, error) {
return settingsValue("settings.getCore", map[string]string{"key": key})
}
func (procSettings) SetCore(key string, v interface{}) error {
return callCoreVoid("settings.setCore", map[string]interface{}{"key": key, "value": v})
}
func (procSettings) ListCore(prefix string) ([]string, error) {
return settingsKeys("settings.listCore", map[string]string{"prefix": prefix})
}
func (procSettings) DataDir() string {
raw, err := callCore("settings.dataDir", nil)
if err != nil {
return ""
}
var r struct {
Dir string `json:"dir"`
}
json.Unmarshal(raw, &r)
return r.Dir
}
func (procSettings) GetPlugin(plugin, key string) (interface{}, error) {
return settingsValue("settings.getPlugin", map[string]string{"plugin": plugin, "key": key})
}
func (procSettings) SetPlugin(plugin, key string, v interface{}) error {
return callCoreVoid("settings.setPlugin", map[string]interface{}{
"plugin": plugin, "key": key, "value": v,
})
}
func (procSettings) ListPlugin(plugin, prefix string) ([]string, error) {
return settingsKeys("settings.listPlugin", map[string]string{"plugin": plugin, "prefix": prefix})
}
func (procSettings) RegisterDef(def sdk.ConfigDef) {
callCoreVoid("settings.registerDef", map[string]interface{}{"def": def})
}
func (procSettings) Defs(prefix string) []*sdk.ConfigDef {
raw, err := callCore("settings.defs", map[string]string{"prefix": prefix})
if err != nil {
return nil
}
var r struct {
Defs []*sdk.ConfigDef `json:"defs"`
}
json.Unmarshal(raw, &r)
return r.Defs
}
func (procSettings) Dump() map[string]interface{} {
raw, err := callCore("settings.dump", nil)
if err != nil {
return nil
}
var m map[string]interface{}
json.Unmarshal(raw, &m)
return m
}
func (procSettings) Plugins() []string {
raw, err := callCore("settings.plugins", nil)
if err != nil {
return nil
}
var r struct {
Plugins []string `json:"plugins"`
}
json.Unmarshal(raw, &r)
return r.Plugins
}
func settingsValue(method string, params interface{}) (interface{}, error) {
raw, err := callCore(method, params)
if err != nil {
return nil, err
}
var r struct {
Value interface{} `json:"value"`
}
if err := json.Unmarshal(raw, &r); err != nil {
return nil, err
}
return r.Value, nil
}
func settingsKeys(method string, params interface{}) ([]string, error) {
raw, err := callCore(method, params)
if err != nil {
return nil, err
}
var r struct {
Keys []string `json:"keys"`
}
if err := json.Unmarshal(raw, &r); err != nil {
return nil, err
}
return r.Keys, nil
}
// ---- 内核 → 插件的调用处理 ----
func handleKernelRequest(req *rpcRequest) {
defer func() {
if r := recover(); r != nil {
// handler panic 只影响本次调用,不带崩进程;
// 真崩溃时进程退出,内核经 EOF 感知并按 recordCrash 处理。
if req.ID != 0 {
respondErr(req.ID, fmt.Errorf("插件 handler panic: %v", r))
}
log.Printf("handler panic (%s): %v", req.Method, r)
}
}()
switch req.Method {
case "handshake":
handleHandshake(req)
case "plugin.init":
var p struct {
Name string `json:"name"`
Config map[string]interface{} `json:"config"`
}
json.Unmarshal(req.Params, &p)
if p.Name != "" {
pluginName = p.Name
}
instance, err := NewPluginFactory(pluginName, p.Config)
if err != nil {
respondErr(req.ID, err)
return
}
plg = instance
respond(req.ID, nil)
case "plugin.start":
if plg == nil {
respondErr(req.ID, fmt.Errorf("plugin.start 前未 init"))
return
}
pluginSDK = buildPluginSDK(pluginName)
if err := plg.Start(pluginSDK); err != nil {
respondErr(req.ID, err)
return
}
// 上报 AutoRestart公开 SDK 的 SetAutoRestart 是纯 setter无回调 hook
// 插件在 Start() 里调它只改进程内副本。C ABI 路径下内核在 Start 返回后
// 直接读 plgSDK.AutoRestart();子进程隔着进程边界读不到,故在此显式上报。
// **不改公开 SDK 接口**(接口冻结约束)。
if err := callCoreVoid("lifecycle.autoRestart", map[string]interface{}{
"enabled": pluginSDK.AutoRestart(),
}); err != nil {
log.Printf("上报 autoRestart 失败: %v", err)
}
respond(req.ID, nil)
case "plugin.stop":
if pluginSDK != nil {
pluginSDK.RunStopHandlers()
}
if plg != nil {
if err := plg.Stop(); err != nil {
log.Printf("Stop: %v", err)
}
}
respond(req.ID, nil)
stdoutW.Flush()
os.Exit(0)
case "tool.invoke":
var p struct {
Name string `json:"name"`
Args map[string]interface{} `json:"args"`
}
json.Unmarshal(req.Params, &p)
handlerMu.RLock()
h, ok := toolHandlers[p.Name]
handlerMu.RUnlock()
if !ok {
respondErr(req.ID, fmt.Errorf("未注册的工具: %s", p.Name))
return
}
res, err := h(p.Args)
if err != nil {
respondErr(req.ID, err)
return
}
respond(req.ID, map[string]interface{}{"result": res})
case "stage.invoke":
handleStageInvoke(req)
case "output.invoke":
var p struct {
Channel string `json:"channel"`
Args map[string]interface{} `json:"args"`
}
json.Unmarshal(req.Params, &p)
handlerMu.RLock()
h, ok := outputHandlers[p.Channel]
handlerMu.RUnlock()
if !ok {
respondErr(req.ID, fmt.Errorf("未注册的输出通道: %s", p.Channel))
return
}
// 同步返回真实结果——内核据此告知模型成功/失败不再假成功§9.4
res, err := h(p.Args)
if err != nil {
respondErr(req.ID, err)
return
}
if m, ok := res.(map[string]interface{}); ok {
respond(req.ID, m)
return
}
respond(req.ID, map[string]interface{}{"status": "sent"})
case "events.subscribe":
var p struct {
Types []string `json:"types"`
}
json.Unmarshal(req.Params, &p)
// 订阅逻辑由内核 EventRing 处理,插件侧在此注册本地 handler。
// 实际事件到达时由 evtConsumerLoop 分发。
for _, t := range p.Types {
var idx uint32
switch t {
case "raw_input":
idx = evtTypeRawInput
case "agent_output":
idx = evtTypeAgentOutput
case "agent_llm_chain":
idx = evtTypeAgentLLMChain
case "tool_call":
idx = evtTypeToolCall
case "reasoning":
idx = evtTypeReasoning
case "stage":
idx = evtTypeStage
case "system":
idx = evtTypeSystem
case "reasoning_delta":
idx = evtTypeReasoningDelta
case "content_delta":
idx = evtTypeContentDelta
case "skill_detected":
idx = evtTypeSkillDetected
default:
continue
}
evtHandlerMu.Lock()
evtHandlers[idx] = func(evt *sdk.Event) {}
evtHandlerMu.Unlock()
}
respond(req.ID, nil)
case "events.unsubscribe":
// 清空全部 handler子进程 Stop 时由内核统一清理订阅)
evtHandlerMu.Lock()
evtHandlers = map[uint32]func(*sdk.Event){}
evtHandlerMu.Unlock()
respond(req.ID, nil)
default:
if req.ID != 0 {
respondErr(req.ID, fmt.Errorf("未实现的 method: %s", req.Method))
}
}
}
func handleHandshake(req *rpcRequest) {
var p struct {
Protocol int `json:"protocol"`
ShmVersion uint32 `json:"shm_version"`
ShmSize int `json:"shm_size"`
PluginName string `json:"plugin_name"`
EvtRingSize int `json:"evt_ring_size,omitempty"`
}
json.Unmarshal(req.Params, &p)
if p.Protocol != procProtocolVersion {
respondErr(req.ID, fmt.Errorf("协议版本不匹配(内核 %d插件 %d——请用配套 plugindev 重编",
p.Protocol, procProtocolVersion))
return
}
if p.ShmVersion != shmVersion {
respondErr(req.ID, fmt.Errorf("共享段版本不匹配(内核 %d插件 %d", p.ShmVersion, shmVersion))
return
}
if p.PluginName != "" {
pluginName = p.PluginName
}
// 挂载 StageContext 共享段。
// 传递机制按平台不同Unix 用继承的 fdWindows 用命名段),
// 由 z_proc_shm_*.go 承担——本文件保持平台无关。
if p.ShmSize > 0 {
m, err := attachStageShm(p.ShmSize)
if err != nil {
respondErr(req.ID, fmt.Errorf("挂载共享段失败: %w", err))
return
}
if got := binary.LittleEndian.Uint32(m[shmOffMagic:]); got != shmMagic {
respondErr(req.ID, fmt.Errorf("共享段魔数不匹配0x%x", got))
return
}
shm = m
}
// 挂载事件环段 + 打开通知句柄§3.6
if p.EvtRingSize > 0 {
er, err := attachEvtRingShm(p.EvtRingSize)
if err != nil {
respondErr(req.ID, fmt.Errorf("挂载事件环段失败: %w", err))
return
}
if got := binary.LittleEndian.Uint32(er[evtOffMagic : evtOffMagic+4]); got != evtRingMagic {
respondErr(req.ID, fmt.Errorf("事件环魔数不匹配0x%x", got))
return
}
notifier, err := openEvtNotifier()
if err != nil {
respondErr(req.ID, fmt.Errorf("打开事件通知句柄失败: %w", err))
return
}
evtRingData = er
evtNotifier = notifier
go evtConsumerLoop()
}
respond(req.ID, map[string]interface{}{
"protocol": procProtocolVersion,
"sdk_version": sdk.SDKVersion,
"plugin_name": pluginName,
"pid": os.Getpid(),
})
}
// handleStageInvoke 执行阶段处理器:拿锁 → 读共享段 → handler → 只写脏字段 → 放锁。
//
// 插件作者的 handler 与 .so 时代完全一致(仍是 func(ctx *sdk.StageContext) error
// 共享内存与锁的复杂度全部由本模板承担§3.4)。
func handleStageInvoke(req *rpcRequest) {
var p struct {
Stage string `json:"stage"`
Seq uint64 `json:"seq"`
}
json.Unmarshal(req.Params, &p)
handlerMu.RLock()
h, ok := stageHandlers[p.Stage]
handlerMu.RUnlock()
if !ok {
respond(req.ID, map[string]interface{}{"dirty_fields": 0})
return
}
if shm == nil {
respondErr(req.ID, fmt.Errorf("共享段未挂载"))
return
}
// 跨进程写锁内核仲裁§3.7),持锁进程崩溃由内核代为释放
if err := callCoreVoid("stage.lock", nil); err != nil {
respondErr(req.ID, fmt.Errorf("申请 stage 锁: %w", err))
return
}
unlocked := false
unlock := func() {
if !unlocked {
unlocked = true
if err := callCoreVoid("stage.unlock", nil); err != nil {
log.Printf("释放 stage 锁: %v", err)
}
}
}
defer unlock()
sc, err := readStageContext()
if err != nil {
respondErr(req.ID, fmt.Errorf("读共享段: %w", err))
return
}
snap := takeStageSnapshot(sc)
if err := h(sc); err != nil {
respondErr(req.ID, err)
return
}
dirty, err := writeStageDirty(sc, snap)
if err != nil {
respondErr(req.ID, fmt.Errorf("写回共享段: %w", err))
return
}
unlock()
respond(req.ID, map[string]interface{}{"dirty_fields": dirty, "seq": p.Seq})
}
// ---- 主循环 ----
func main() {
// 日志走 stderrstdout 是 RPC 通道,写日志会破坏帧
log.SetOutput(os.Stderr)
log.SetPrefix("[plugin] ")
in := bufio.NewScanner(bufio.NewReader(os.Stdin))
// 单帧上限 1MB控制面帧本应很小大 payload 走共享段
in.Buffer(make([]byte, 0, 64*1024), 1024*1024)
for in.Scan() {
line := make([]byte, len(in.Bytes()))
copy(line, in.Bytes())
var probe struct {
ID uint64 `json:"id"`
Method string `json:"method"`
}
if err := json.Unmarshal(line, &probe); err != nil {
log.Printf("非法 JSON 帧: %v", err)
continue
}
// method 为空 = 内核对我们反向调用的应答
if probe.Method == "" {
var resp rpcResponse
if err := json.Unmarshal(line, &resp); err != nil {
continue
}
pendMu.Lock()
ch, ok := pending[resp.ID]
delete(pending, resp.ID)
pendMu.Unlock()
if ok {
ch <- resp
}
continue
}
var req rpcRequest
if err := json.Unmarshal(line, &req); err != nil {
continue
}
// 每个请求独立 goroutinehandler 内可能反向调用内核,
// 在读循环里同步处理会死锁(等应答但没人读)。
go handleKernelRequest(&req)
}
if err := in.Err(); err != nil {
log.Printf("读 stdin 出错: %v", err)
}
// stdin 关闭 = 内核结束了我们
if pluginSDK != nil {
pluginSDK.RunStopHandlers()
}
if plg != nil {
plg.Stop()
}
}
// ---- 事件环消费§3.6,子进程侧)----
// 事件环布局常量(与内核 internal/plugin/proc/evtring.go 一致)。
const (
evtOffMagic = 0
evtOffVersion = 4
evtOffWriteSeq = 8
evtOffCap = 16
evtOffSlots = 20
evtRingMagic = 0x48455654 // "HEVT"
// 事件类型位索引(与内核 encodeEvtType 一致)
evtTypeRawInput = 0
evtTypeAgentOutput = 1
evtTypeAgentLLMChain = 2
evtTypeToolCall = 3
evtTypeReasoning = 4
evtTypeStage = 5
evtTypeSystem = 6
evtTypeReasoningDelta = 7
evtTypeContentDelta = 8
evtTypeSkillDetected = 9
evtTypeMax = 10
evtRingSlotLen uint32 = 32
evtRingCap uint32 = 8192
)
// evtTypeNames 把位索引还原成 pubsdk.EventType 字符串。
var evtTypeNames = [evtTypeMax]string{
"raw_input", "agent_output", "agent_llm_chain", "tool_call",
"reasoning", "stage", "system", "reasoning_delta",
"content_delta", "skill_detected",
}
// evtConsumerLoop 是事件环消费主循环eventfd.Read阻塞走 netpoller→ drain events → 分发。
// 每个插件进程启动一个 goroutine与 RPC 主循环并行。
func evtConsumerLoop() {
if evtNotifier == nil || evtRingData == nil {
return
}
buf := make([]byte, 8) // eventfd uint64 计数
var readSeq uint64
for {
if err := evtNotifier.Wait(buf); err != nil {
continue
}
// 循环 drain 直到无新事件eventfd 计数合并,一次 Read 处理全部)
for {
writeSeq := binary.LittleEndian.Uint64(evtRingData[evtOffWriteSeq:])
if readSeq >= writeSeq {
break
}
cap := uint64(evtRingCap)
if writeSeq-readSeq > cap {
readSeq = writeSeq - cap // 跳到最旧可读事件
}
idx := readSeq % cap
slotOff := evtOffSlots + uint32(idx)*evtRingSlotLen
seq := binary.LittleEndian.Uint64(evtRingData[slotOff:])
etype := binary.LittleEndian.Uint32(evtRingData[slotOff+8:])
off := binary.LittleEndian.Uint32(evtRingData[slotOff+12:])
slen := binary.LittleEndian.Uint32(evtRingData[slotOff+16:])
if seq != readSeq {
// slot 被覆盖,跳到最新
if writeSeq > cap {
readSeq = writeSeq - cap
} else {
readSeq = writeSeq
}
continue
}
// 读载荷并分发给注册的 handler
if off > 0 && slen > 0 && uint64(off)+uint64(slen) <= uint64(len(evtRingData)) {
evtTypeStr := ""
if int(etype) < len(evtTypeNames) {
evtTypeStr = evtTypeNames[etype]
}
evtHandlerMu.RLock()
handler, ok := evtHandlers[etype]
evtHandlerMu.RUnlock()
if ok && handler != nil && evtTypeStr != "" {
evt := &sdk.Event{Type: sdk.EventType(evtTypeStr)}
// payload 非核心(多数订阅者只看类型),简化为不解析 JSON
_ = evtRingData[off : off+slen]
handler(evt)
}
}
readSeq++
}
}
}
// evtWaiter 抽象事件通知等待。
//
// UnixeventfdLinux/ pipemacOS的读端阻塞 Read 走 netpoller。
// Windows命名 Event 对象WaitForSingleObject。
// 平台实现在 z_proc_shm_unix.go / z_proc_shm_windows.go。
type evtWaiter interface {
// Wait 阻塞直到有新事件buf 供实现复用Unix 读 8 字节计数)。
Wait(buf []byte) error
}