Files
homeagent-sdk/tools/plugindev/templates/proc_main.go.tmpl
JianFeeeee ce5bff9275 feat(sdk): 多模态贯通插件边界——媒体字段、媒体注入接口与并发修复
记忆系统在核心 1.1.0 支持了二进制多媒体节点,但那条链路只对**内核自己**开放:
插件把 Triple / Doc 交进来,媒体一律无处安放,且**不报错**。本版补上公开接口
侧缺失的表达能力。

## 一、类型与接口(全部新增,无签名变更)

- `Triple` += `SentenceText`、`MediaDigests`
- `Doc` += `MediaDigests`、`Attachments`;新增 `MediaAttachment`
- `TextEvent` += `Attachments`
- `DocMemoryAPI` += `InsertWithMedia`
- `IOInjector` += `InjectInputMedia` / `InjectInputMediaSync` / `InjectInterruptMedia`
- `PluginSDK` 补上一直缺失的 `SetToolBlocks` 包装(接口里有、便捷方法里没有,
  插件只能自己去拿 injector)

`MediaAttachment` 一个类型服务两个方向:给 `Data`+`MIME` 是新内容(内核按字节
去重),只给 `Digest` 是引用已有内容。读路径**只回元数据不回字节**——一次检索
可能命中几十份媒体,把字节全塞回来会撑爆跨进程消息。

媒体注入为什么不能搭 `SetToolBlocks` 的车:那个方法只在工具处理函数内部可用,
且媒体要等**下一条** tool message 才到模型手上。插件主动发起一轮带媒体的对话、
以及中断注入,需要各自的签名,且媒体在**本轮**就随消息发出。

`Triple.MediaDigests` 非空而 `SentenceText` 为空时,内核会用媒体标记本身充当句子
——媒体引用挂在句子上,没有句子就无处挂起。插件只需填 digest,标记由内核拼:
要求调用方知道格式,等于让一个拼写错误静默切断引用绑定而全链路无人报错。

## 二、修掉两处并发竞态

`sdk/stress_test.go` 的 `-race` 实测报 11 处 DATA RACE,收敛到两个字段:

1. **`PluginSDK` 的 API 字段无锁**。写方是内核(加载/重载插件时依次注入
   injector、memory、doc、llm…),读方是插件在 `Start()` 里起的后台 goroutine
   ——轮询、监听、定时器都要拿 injector 往管道注消息。生产表现是插件重载瞬间
   偶发崩溃:读到半个接口值就 nil 解引用。
2. **`autoRestart` 标志无锁**。`SetAutoRestart` 的文档用法本身就是「外部连接建好
   后再决定能否自动重启」,而连接建立通常在后台 goroutine;内核 registry 在另一个
   goroutine 读 `AutoRestart()` 决定崩溃后重启策略。这对读写天然跨 goroutine。

加 `apiMu sync.RWMutex`。关键约定写进注释:**只在持锁期间取字段值,取完立刻
释放再调用**。持锁调用会把 `InjectInputSync`(阻塞到 agent 回复,可达数分钟)
与 `SetIOInjector` 串到一起,让插件重载卡死。

## 三、压测(sdk/stress_test.go,13 例)

SDK 是被多个 goroutine 同时使用的共享对象,单线程单测全绿不代表并发路径成立。
断言的是不变量而非吞吐:

- 媒体注入高并发不丢不串——每次调用带唯一 tag,逐条校验文本与图片 URL 配对。
  「不串」是重点:若实现里出现任何共享中间状态(把 blocks 暂存到字段再读出),
  高并发下会出现 A 的文本配 B 的图,而两者单独看都「成功」了;
- injector 热替换(含替换成 nil,即内核卸载 API 的真实状态);
- stop / onRemove handler 恰好一次——契约是「执行后清空,幂等」,执行两次的后果
  从重复写文件到 close 已关闭 channel 直接 panic;
- `StageContext` 并发读改写无 lost update(媒体链路让 Extra 成为新热点,
  而 map 并发写在 Go 里是直接 fatal,recover 接不住);
- `OwnTools` scope 不跨插件泄漏;
- 媒体类型 JSON 往返字节级一致(9 种长度,含 0/1/2/3 与 base64 分组边界)
  ——`[]byte` 在 JSON 里是 base64,往返不一致意味着图片静默损坏,
  要到 CAS 校验 digest 时才发现,那时已无从追查;
- `omitempty` 真的生效(读路径不能出现 `"data"` 键);
- nil 依赖全部静默降级不 panic。

## 四、工具链同步

- `proc_main.go.tmpl`:`procIO` 三个媒体方法、`procDocMemory.InsertWithMedia`。
  模板不跟上的后果是**每个外部插件都编不过**(接口未实现),是硬失败;
- `proc_runtime_test.go`:方法清单补 `io.injectMedia*` 与 `doc.insertWithMedia`。
  漏接线时插件调 `InjectInputMedia` 会静默无效果——模板不发这个 RPC,内核也就
  收不到,两边都不报错;
- `yaegi/mocksdk`:与公开 SDK 对齐。它此前漂移严重且**没有任何代码对着它编译**,
  所以漂移不会被编译器抓到:`Triple` 用的是 `Predicate`,而公开 SDK 一直叫
  `Relation` —— 插件在 yaegi 调试期写 `Relation:` 报未知字段,写 `Predicate:` 则
  编成 plugin.bin 时报错,两边都不对。
- README 中英双语补媒体接口文档与用法示例。

## 兼容性

存量插件不需要改一行也不需要重编:新增方法由**插件调用、内核实现**,不调就不
受影响。17 个 example 插件源码零改动通过类型检查;用 SDK 0.9.2 编的旧 plugin.bin
在新内核上直接建链通过(握手校验的是 ProtocolVersion=1,不是 SDK 版本)。

媒体接口需要核心 1.1.1+(更早的核心没有对应 RPC,调用返回 unknown method)。
`CoreVersion` 保持 1.0.0:它是「SDK 能在其上运行」的下限,媒体是可选能力。
2026-09-06 10:03:36 +08:00

1340 lines
35 KiB
Cheetah
Raw Permalink 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/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.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
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)
}
}
// 带媒体的注入:插件主动发起一轮带图/音频的对话。
// 与 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 "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 用继承的 fd,Windows 用命名段),
// 由 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() {
// 日志走 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
}