Files
homeagent-sdk/tools/plugindev/templates/proc_main.go.tmpl
JianFeeeee fc236120e3 feat(shm): 统一共享内存区域(§13.1) — 子进程侧模板适配
- proc_shm_unix.go.tmpl: fd 3 = 统一区域,fd 4 = eventfd
- proc_main.go.tmpl: 握手解析 SuperBlock,从 ctxOff/evtOff 定位两段
- 新增 unified region 常量(magic/version/offset) + StageContext 内部布局常量
2026-09-10 11:46:12 +08:00

1406 lines
37 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 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
)
// ---- 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
)
// ---- 全局状态 ----
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
// 事件环(§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"`
Text string `json:"text"`
}
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
}
respond(req.ID, map[string]interface{}{"text": cleaner(p.Text)})
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]
// 挂载事件环段 + 打开通知句柄(§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
}