package main // 子进程插件入口(由 plugindev 自动生成,请勿手工编辑)。 // // 与旧 C ABI bridge(z_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/unified.go 一致)---- const ( unifiedMagic = 0x554D5352 // "UMSR" — Unified Memory Shared Region unifiedVersion = 1 superBlockSize = 64 sbOffMagic = 0 sbOffVersion = 4 sbOffGeneration = 8 sbOffCapacity = 16 sbOffCtxOff = 20 sbOffCtxSize = 24 sbOffEvtOff = 28 sbOffEvtSize = 32 sbOffArenaOff = 36 sbOffArenaCap = 40 sbOffArenaUsed = 44 ) // ---- StageContext 段内部布局(与内核 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 ) // SharedRef 跨进程共享内存描述符。 type SharedRef struct { Offset uint32 `json:"offset"` Length uint32 `json:"length"` Generation uint32 `json:"generation"` Flags uint32 `json:"flags"` } func (r SharedRef) IsZero() bool { return r.Offset == 0 && r.Length == 0 } func (r SharedRef) Slice(data []byte) []byte { if r.IsZero() || int(r.Offset)+int(r.Length) > len(data) { return nil } return data[r.Offset : r.Offset+r.Length] } // arena 辅助(插件侧简化版:bump 分配) var arenaOff uint32 var arenaUsed uint32 func arenaWrite(b []byte) (SharedRef, error) { if len(b) == 0 { return SharedRef{}, nil } off := arenaUsed if off == 0 { off = 1 } end := off + uint32(len(b)) if int(end) > len(region)-int(arenaOff) { return SharedRef{}, fmt.Errorf("arena 空间不足") } copy(region[arenaOff+off:arenaOff+end], b) arenaUsed = end return SharedRef{Offset: arenaOff + off, Length: uint32(len(b))}, nil } // ---- 全局状态 ---- 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{} toolCleaners = map[string]func(string) string{} inputCleaners = map[string]func(string) string{} stageHandlers = map[string]sdk.StageHandler{} outputHandlers = map[string]sdk.ToolHandler{} outputCleaners = map[string]func(string) string{} shm []byte region []byte // 统一区域完整 mmap(用于 SharedRef 读写) // 事件环(§3.6):fd 4 = 事件环段 mmap,fd 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 if def.Cleaner != nil { toolCleaners[toolName] = def.Cleaner } else { delete(toolCleaners, toolName) } handlerMu.Unlock() return callCoreVoid("tool.register", map[string]interface{}{ "name": toolName, "def": def, "has_cleaner": def.Cleaner != nil, }) }, 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 if def.Cleaner != nil { outputCleaners[chName] = def.Cleaner } else { delete(outputCleaners, chName) } handlerMu.Unlock() return callCoreVoid("output.register", map[string]interface{}{ "name": chName, "caps": caps, "desc": desc, "def": map[string]interface{}{"NoMemory": def.NoMemory}, "has_cleaner": def.Cleaner != nil, }) }, ) 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 { handlerMu.Lock() if def.Cleaner != nil { inputCleaners[chName] = def.Cleaner } else { delete(inputCleaners, chName) } handlerMu.Unlock() return callCoreVoid("input.register", map[string]interface{}{ "name": chName, "def": map[string]interface{}{"NoMemory": def.NoMemory}, "has_cleaner": def.Cleaner != nil, }) }) 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) } } // 带媒体的注入:插件主动发起一轮带图/音频的对话。 // 与 SetToolBlocks 的区别是媒体在**本轮**就到模型手上,而不是等下一条 tool message。 func (procIO) InjectInputMedia(s, c, t string, blocks []sdk.ContentBlock) { callCoreVoid("io.injectMedia", map[string]interface{}{ "source": s, "channel": c, "text": t, "blocks": blocks, }) } func (procIO) InjectInputMediaSync(s, c, t string, blocks []sdk.ContentBlock) string { raw, err := callCore("io.injectMediaSync", map[string]interface{}{ "source": s, "channel": c, "text": t, "blocks": blocks, }) if err != nil { return "" } var r struct { Reply string `json:"reply"` } json.Unmarshal(raw, &r) return r.Reply } func (procIO) InjectInterruptMedia(s, c, t string, blocks []sdk.ContentBlock) { callCoreVoid("io.injectInterruptMedia", map[string]interface{}{ "source": s, "channel": c, "text": t, "blocks": blocks, }) } 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}) } // InsertWithMedia 写入文档并关联媒体。 // // 内核会把 `[mime <短digest>] <描述>` 标记补进 Content 并挂上引用,回传的 // doc 带着补好的 Content/ID/MediaDigests——回写进 d 让调用方能拿到这些。 func (procDocMemory) InsertWithMedia(d *sdk.Doc, atts []sdk.MediaAttachment) error { raw, err := callCore("doc.insertWithMedia", map[string]interface{}{ "doc": d, "attachments": atts, }) if err != nil { return err } var r struct { Doc *sdk.Doc `json:"doc"` } if json.Unmarshal(raw, &r) == nil && r.Doc != nil { *d = *r.Doc } return nil } 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 "cleaner.invoke": var p struct { Scope string `json:"scope"` Name string `json:"name"` TextRef SharedRef `json:"text_ref"` } if err := json.Unmarshal(req.Params, &p); err != nil { respondErr(req.ID, fmt.Errorf("解析 Cleaner 参数: %w", err)) return } handlerMu.RLock() var cleaner func(string) string switch p.Scope { case "tool": cleaner = toolCleaners[p.Name] case "input": cleaner = inputCleaners[p.Name] case "output": cleaner = outputCleaners[p.Name] } handlerMu.RUnlock() if cleaner == nil { respondErr(req.ID, fmt.Errorf("%s %s 未注册 Cleaner", p.Scope, p.Name)) return } input := string(p.TextRef.Slice(region)) output := cleaner(input) outRef, err := arenaWrite([]byte(output)) if err != nil { respondErr(req.ID, fmt.Errorf("arena 写入失败: %w", err)) return } respond(req.ID, map[string]interface{}{"text_ref": outRef}) 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 } // 挂载统一共享内存区域。 // fd 3 (Unix) / 命名对象 (Windows) 传给插件子进程,包含 SuperBlock + // StageContext + EvtRing 两段。SuperBlock 记录各段的偏移与大小。 if p.ShmSize > 0 { m, err := attachUnifiedShm(p.ShmSize) if err != nil { respondErr(req.ID, fmt.Errorf("挂载统一共享区域失败: %w", err)) return } if got := binary.LittleEndian.Uint32(m[sbOffMagic:]); got != unifiedMagic { respondErr(req.ID, fmt.Errorf("统一区域魔数不匹配(0x%x,期望 0x%x)", got, unifiedMagic)) return } ctxOff := binary.LittleEndian.Uint32(m[sbOffCtxOff:]) ctxSize := binary.LittleEndian.Uint32(m[sbOffCtxSize:]) evtOff := binary.LittleEndian.Uint32(m[sbOffEvtOff:]) evtSize := binary.LittleEndian.Uint32(m[sbOffEvtSize:]) // shm 指向 StageContext 段,后续代码用 shm[off...] 访问该段内部字段 shm = m[ctxOff : ctxOff+ctxSize] // region 保存完整 mmap 区域,供 SharedRef 读写 arena region = m arenaOff = binary.LittleEndian.Uint32(m[sbOffArenaOff:]) arenaUsed = 1 // 挂载事件环段 + 打开通知句柄(§13.1:EvtRing 在统一区域内) if p.EvtRingSize > 0 && evtSize > 0 { er := m[evtOff : evtOff+evtSize] 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() { // 日志走 stderr:stdout 是 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 } // 每个请求独立 goroutine:handler 内可能反向调用内核, // 在读循环里同步处理会死锁(等应答但没人读)。 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 抽象事件通知等待。 // // Unix:eventfd(Linux)/ pipe(macOS)的读端,阻塞 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 }