diff --git a/sdk/plugin.go b/sdk/plugin.go index ebd234d..11aac9c 100644 --- a/sdk/plugin.go +++ b/sdk/plugin.go @@ -100,12 +100,13 @@ type ToolResult struct { // ToolDef describes a tool that the plugin exposes. type ToolDef struct { - Name string `json:"name"` - Plugin string `json:"plugin,omitempty"` - Description string `json:"description"` - Parameters map[string]interface{} `json:"parameters"` - NoMemory bool `json:"no_memory,omitempty"` // 此工具输出不参与记忆计算,但原文保留 - Cleaner func(string) string `json:"-"` // 计算层过滤函数,不改原文;仅在向量化/jieba/蒸馏时调用 + Name string `json:"name"` + Plugin string `json:"plugin,omitempty"` + Description string `json:"description"` + Parameters map[string]interface{} `json:"parameters"` + NoMemory bool `json:"no_memory,omitempty"` // 此工具输出不参与记忆计算,但原文保留 + Cleaner func(string) string `json:"-"` // 计算层过滤函数,不改原文;仅在向量化/jieba/蒸馏时调用 + ContextPolicy string `json:"context_policy,omitempty"` // 工具上下文策略:"none"(默认) / "prune" } // IOInjector provides methods for injecting input and interrupts into the agent pipeline. diff --git a/tools/plugindev/proc_runtime_test.go b/tools/plugindev/proc_runtime_test.go index 3681612..c14ec79 100644 --- a/tools/plugindev/proc_runtime_test.go +++ b/tools/plugindev/proc_runtime_test.go @@ -133,13 +133,13 @@ func TestProcTemplate_CoversAllCoreMethods(t *testing.T) { } } -// 模板必须处理内核发来的全部 7 个调用(原 C ABI 的 7 个 //export)。 +// 模板必须处理内核发来的全部调用(含无法 JSON 序列化的 Cleaner 回调)。 func TestProcTemplate_HandlesAllKernelCalls(t *testing.T) { src := loadProcTemplate(t) for _, m := range []string{ "handshake", "plugin.init", "plugin.start", "plugin.stop", - "tool.invoke", "stage.invoke", "output.invoke", + "tool.invoke", "cleaner.invoke", "stage.invoke", "output.invoke", } { if !strings.Contains(src, `case "`+m+`"`) { t.Errorf("模板未处理内核调用 %q", m) @@ -147,6 +147,31 @@ func TestProcTemplate_HandlesAllKernelCalls(t *testing.T) { } } +// 模板必须通过 arena.alloc / arena.free 向内核申请与归还共享内存。 +// +// 共享内存是内核独占管理的**内部实现**:插件不能自己维护分配游标。 +// 历史上两版跨进程分配器(bump 游标 / 模板内位图 CAS)都因为把可变 +// 分配状态放在共享内存里而出竞态,所以这里做回归保护。 +func TestProcTemplate_UsesKernelArenaRPC(t *testing.T) { + src := loadProcTemplate(t) + + for _, m := range []string{`"arena.alloc"`, `"arena.free"`} { + if !strings.Contains(src, m) { + t.Errorf("模板缺少内核共享内存 RPC %s(插件必须向内核申请/归还)", m) + } + } + + // 禁止插件侧再出现本地分配器符号。 + // + // 只查代码不查注释:注释里会解释“为什么不再这么做”。 + code := stripComments(t, src) + for _, forbidden := range []string{"arenaUsed", "arenaWrite"} { + if strings.Contains(code, forbidden) { + t.Errorf("模板不应再出现插件侧分配器 %q(共享内存由内核独占管理)", forbidden) + } + } +} + // 共享段布局常量必须与内核 internal/plugin/proc/shm.go 一致。 // // 字段索引错位是最危险的漂移:插件会读到相邻字段的数据, @@ -281,10 +306,11 @@ func TestProcTemplate_DispatchesRequestsConcurrently(t *testing.T) { } } -// 协议与共享段版本不匹配必须拒绝,不得半兼容运行。 +// 协议与共享内存区域版本/魔数不匹配必须拒绝,不得半兼容运行。 func TestProcTemplate_RejectsVersionMismatch(t *testing.T) { src := loadProcTemplate(t) - for _, want := range []string{"协议版本不匹配", "共享段版本不匹配", "共享段魔数不匹配"} { + // §13.1 起共享段合并为单一「统一区域」,魔数校验文案随之更新。 + for _, want := range []string{"协议版本不匹配", "共享段版本不匹配", "统一区域魔数不匹配"} { if !strings.Contains(src, want) { t.Errorf("握手应校验并拒绝 %q", want) } diff --git a/tools/plugindev/templates/proc_main.go.tmpl b/tools/plugindev/templates/proc_main.go.tmpl index b2b0a42..0622130 100644 --- a/tools/plugindev/templates/proc_main.go.tmpl +++ b/tools/plugindev/templates/proc_main.go.tmpl @@ -93,6 +93,10 @@ const ( ) // SharedRef 跨进程共享内存描述符。 +// +// ⚠️ 这是**内部实现细节**:插件开发者永远看不到它。公开 SDK 只暴露普通 +// 字符串与 Map;模板运行时在传输层按 payload 大小自动选择内联 JSON 还是 +// 共享槽。直接使用 SharedRef 属于运行时内部行为,不是插件 API。 type SharedRef struct { Offset uint32 `json:"offset"` Length uint32 `json:"length"` @@ -108,25 +112,78 @@ func (r SharedRef) Slice(data []byte) []byte { return data[r.Offset : r.Offset+r.Length] } -// arena 辅助(插件侧简化版:bump 分配) -var arenaOff uint32 -var arenaUsed uint32 +// region 是内核传入的统一共享区域 mmap(handshake 时设置)。 +var region []byte -func arenaWrite(b []byte) (SharedRef, error) { - if len(b) == 0 { - return SharedRef{}, nil +// arenaAlloc 向内核申请一块共享内存,内核返回偏移与大小(Length 为槽容量)。 +// +// 分配器由内核独占管理(见内核 proc/arena.go):插件只申请与归还, +// 不做任何分配决策,因此不存在跨进程分配器的竞争。 +func arenaAlloc(size uint32) (SharedRef, error) { + raw, err := callCore("arena.alloc", map[string]interface{}{"size": size}) + if err != nil { + return SharedRef{}, err } - off := arenaUsed - if off == 0 { - off = 1 + var r struct { + Ref SharedRef `json:"ref"` } - end := off + uint32(len(b)) - if int(end) > len(region)-int(arenaOff) { - return SharedRef{}, fmt.Errorf("arena 空间不足") + if err := json.Unmarshal(raw, &r); err != nil { + return SharedRef{}, err } - copy(region[arenaOff+off:arenaOff+end], b) - arenaUsed = end - return SharedRef{Offset: arenaOff + off, Length: uint32(len(b))}, nil + if r.Ref.IsZero() { + return SharedRef{}, fmt.Errorf("arena.alloc: 内核返回空引用") + } + return r.Ref, nil +} + +// arenaFree 通知内核回收先前申请的共享内存。 +func arenaFree(ref SharedRef) { + if ref.IsZero() { + return + } + callCoreVoid("arena.free", map[string]interface{}{"ref": ref}) +} + +// inlinePayloadLimit 是走内联 JSON 的上限。 +// +// 小 payload 走内联省两次 RPC(申请 + 归还);大 payload 走共享内存, +// 避免把长文本塞进 NDJSON 帧。这是纯传输层优化,插件开发者无感。 +const inlinePayloadLimit = 512 + +// putInArena 把 payload 写入内核分配的共享槽,返回可随业务 RPC 回传的引用。 +// +// 任一步失败都返回 ok=false,让调用方退回内联:共享内存只是优化, +// 池满或超限绝不能影响功能。 +func putInArena(payload string) (SharedRef, bool) { + if len(payload) <= inlinePayloadLimit || len(region) == 0 { + return SharedRef{}, false + } + ref, err := arenaAlloc(uint32(len(payload))) + if err != nil { + return SharedRef{}, false + } + if len(payload) > int(ref.Length) || int(ref.Offset)+len(payload) > len(region) { + arenaFree(ref) + return SharedRef{}, false + } + copy(region[ref.Offset:ref.Offset+uint32(len(payload))], payload) + ref.Length = uint32(len(payload)) + return ref, true +} + +// callWithText 按 payload 大小自动选择共享槽或内联,发起一次带文本的业务 RPC。 +// +// 共享内存对插件开发者完全透明:SDK 层只看得到 string。 +func callWithText(method, source, channel, text string) (json.RawMessage, error) { + if ref, ok := putInArena(text); ok { + defer arenaFree(ref) + return callCore(method, map[string]interface{}{ + "source": source, "channel": channel, "text_ref": ref, + }) + } + return callCore(method, map[string]string{ + "source": source, "channel": channel, "text": text, + }) } // ---- 全局状态 ---- @@ -152,12 +209,13 @@ var ( outputCleaners = map[string]func(string) string{} shm []byte - region []byte // 统一区域完整 mmap(用于 SharedRef 读写) + + // region 见文件头部 SharedRef 注释(handshake 时设置)。 // 事件环(§3.6):fd 4 = 事件环段 mmap,fd 5 = eventfd 读端 - evtRingData []byte - evtNotifier evtWaiter - evtHandlers = map[uint32]func(*sdk.Event){} + evtRingData []byte + evtNotifier evtWaiter + evtHandlers = map[uint32]func(*sdk.Event){} evtHandlerMu sync.RWMutex ) @@ -511,7 +569,7 @@ func buildPluginSDK(name string) *sdk.PluginSDK { handlerMu.Unlock() return callCoreVoid("output.register", map[string]interface{}{ "name": chName, "caps": caps, "desc": desc, - "def": map[string]interface{}{"NoMemory": def.NoMemory}, + "def": map[string]interface{}{"NoMemory": def.NoMemory}, "has_cleaner": def.Cleaner != nil, }) }, @@ -534,8 +592,8 @@ func buildPluginSDK(name string) *sdk.PluginSDK { } handlerMu.Unlock() return callCoreVoid("input.register", map[string]interface{}{ - "name": chName, - "def": map[string]interface{}{"NoMemory": def.NoMemory}, + "name": chName, + "def": map[string]interface{}{"NoMemory": def.NoMemory}, "has_cleaner": def.Cleaner != nil, }) }) @@ -545,16 +603,17 @@ func buildPluginSDK(name string) *sdk.PluginSDK { type procIO struct{} func (procIO) InjectText(s, c, t string) { - callCoreVoid("io.injectText", map[string]string{"source": s, "channel": c, "text": t}) + // 忽略错误:注入是 fire-and-forget,与内联路径语义一致 + _, _ = callWithText("io.injectText", s, c, t) } func (procIO) InjectInterruptText(s, c, t string) { - callCoreVoid("io.injectInterrupt", map[string]string{"source": s, "channel": c, "text": t}) + _, _ = callWithText("io.injectInterrupt", s, c, t) } func (procIO) InjectTextNoMemory(s, c, t string) { - callCoreVoid("io.injectTextNoMem", map[string]string{"source": s, "channel": c, "text": t}) + _, _ = callWithText("io.injectTextNoMem", s, c, t) } func (procIO) InjectInputSync(s, c, t string) string { - raw, err := callCore("io.injectInputSync", map[string]string{"source": s, "channel": c, "text": t}) + raw, err := callWithText("io.injectInputSync", s, c, t) if err != nil { return "" } @@ -1043,9 +1102,11 @@ func handleKernelRequest(req *rpcRequest) { case "cleaner.invoke": var p struct { - Scope string `json:"scope"` - Name string `json:"name"` + Scope string `json:"scope"` + Name string `json:"name"` + Text string `json:"text"` TextRef SharedRef `json:"text_ref"` + RespRef SharedRef `json:"resp_ref"` } if err := json.Unmarshal(req.Params, &p); err != nil { respondErr(req.ID, fmt.Errorf("解析 Cleaner 参数: %w", err)) @@ -1066,14 +1127,23 @@ func handleKernelRequest(req *rpcRequest) { respondErr(req.ID, fmt.Errorf("%s %s 未注册 Cleaner", p.Scope, p.Name)) return } - input := string(p.TextRef.Slice(region)) + // 读输入:优先内核写入的共享槽,否则内联。 + input := p.Text + if !p.TextRef.IsZero() { + 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)) + + // 写结果:内核预分配了响应槽且放得下就写槽,否则内联。 + // 插件不申请任何槽(“谁分配谁归还”全部在内核侧)。 + if !p.RespRef.IsZero() && len(output) <= int(p.RespRef.Length) { + copy(region[p.RespRef.Offset:p.RespRef.Offset+uint32(len(output))], output) + outRef := p.RespRef + outRef.Length = uint32(len(output)) + respond(req.ID, map[string]interface{}{"text_ref": outRef}) return } - respond(req.ID, map[string]interface{}{"text_ref": outRef}) + respond(req.ID, map[string]interface{}{"text": output}) case "stage.invoke": handleStageInvoke(req) @@ -1159,10 +1229,10 @@ func handleKernelRequest(req *rpcRequest) { 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"` + 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) @@ -1198,10 +1268,9 @@ func handleHandshake(req *rpcRequest) { evtSize := binary.LittleEndian.Uint32(m[sbOffEvtSize:]) // shm 指向 StageContext 段,后续代码用 shm[off...] 访问该段内部字段 shm = m[ctxOff : ctxOff+ctxSize] - // region 保存完整 mmap 区域,供 SharedRef 读写 arena + // region 保存完整 mmap 区域。arena 的偏移与大小由内核在 + // arena.alloc 的应答里下发,插件侧不再自己解析槽池布局。 region = m - arenaOff = binary.LittleEndian.Uint32(m[sbOffArenaOff:]) - arenaUsed = 1 // 挂载事件环段 + 打开通知句柄(§13.1:EvtRing 在统一区域内) if p.EvtRingSize > 0 && evtSize > 0 { er := m[evtOff : evtOff+evtSize]