From e2672c56d6af892cf559a46cf43f0eec84019e49 Mon Sep 17 00:00:00 2001 From: JianFeeeee Date: Thu, 10 Sep 2026 21:28:56 +0800 Subject: [PATCH] =?UTF-8?q?feat(shm):=20=E8=BE=93=E5=87=BA=E9=80=9A?= =?UTF-8?q?=E9=81=93=20payload=20=E8=B5=B0=E5=85=B1=E4=BA=AB=E5=86=85?= =?UTF-8?q?=E5=AD=98=E8=B0=83=E7=94=A8=E5=B8=A7(=C2=A713.6)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit §13.3 给 ToolInvokeParams 加了调用帧,但 OutputInvokeParams 没跟上—— payload 仍内联在 RPC JSON 里。核实后确认这个 lane 只做了半边:注入侧 (injectParams.TextRef + resolveText)可用,输出侧从内核到插件仍是内联。 改法照搬工具调用的 funccall 帧模型: - OutputInvokeParams 加 Frame/ArgsLen(Args 仅留给直连 RPC 的测试) - invokeOutput 序列化参数 → Alloc 帧 → 写帧 → 只传偏移描述符 → defer Free - 帧尾不预留结果区:output 应答很小("ok"/status map),走 RPC 应答字段即可。 但 OutputInvokeResult 支持插件把大结果写回帧(ResultRef),并识别 sharedRefFlagExpand 单独归还扩容块——与 invokeTool 一致。 - 模板 output.invoke 用 frameInput 从帧读参数,无帧才回退内联 Args 为什么要做:output payload 在真机上可能含图片/文件描述等大字段,内联会 撑爆 stdin/stdout 管道。 验证: - TestPlugin_OutputPayloadViaArena:9000 字节 payload 经帧完整送达(插件回报 实收长度,不是只看返回值)+ 调用后 arena 归零 - TestE2E_RealTemplateOutputPayloadViaFrame:同上,但用真实 SDK 模板编译的 插件——生产插件走的就是模板,模板不读帧则此改动等于没做 - TestPlugin_OutputChannelReportsRealFailure 仍绿:同步等真实结果、失败必须 上报(§9.4)的语义未被破坏 - go test -race ./internal/plugin/... 全绿 plan.md §13.5 也如实标注:注入 side 内核→插件方向仍未接线(procCore.InjectText 直接传字符串,从不 Alloc/Put 构造 TextRef),不是仅打勾的项。 --- internal/plugin/proc/e2e_template_test.go | 64 ++++++++++++++++++++ internal/plugin/proc/plugin.go | 58 ++++++++++++++++-- internal/plugin/proc/plugin_test.go | 48 +++++++++++++++ internal/plugin/proc/protocol.go | 18 +++++- internal/plugin/proc/testdata/stageplugin.go | 46 ++++++++++---- plan.md | 32 ++++++++-- 6 files changed, 244 insertions(+), 22 deletions(-) diff --git a/internal/plugin/proc/e2e_template_test.go b/internal/plugin/proc/e2e_template_test.go index e80e24e..7aefc6c 100644 --- a/internal/plugin/proc/e2e_template_test.go +++ b/internal/plugin/proc/e2e_template_test.go @@ -62,6 +62,19 @@ func (p *e2ePlugin) Start(s *sdk.PluginSDK) error { ctx.FinalText = ctx.FinalText + "|stage-touched" return nil }) + + // §13.5/13.6 输入/输出 lane:text 不再内联在 RPC 报文里,而是经共享内存 + // TextRef 传递(大 payload 不爆 stdin/stdout 管道)。插件侧只调公开 API, + // 内核侧负责 Alloc/Put/Free。 + s.RegisterOutputChannel("e2e_out", 1, "测试输出通道", sdk.ChannelDef{}, + func(args map[string]interface{}) (interface{}, error) { + payload, _ := args["payload"].(string) + // 回报实际收到的长度:只有完整 payload 经共享帧送达才等于发送长度。 + return map[string]interface{}{"status": "sent", "payload_len": len(payload)}, nil + }) + + s.RegisterInputChannel("e2e_in", sdk.ChannelDef{}) + return nil } @@ -273,3 +286,54 @@ func (p *readerPlugin) Stop() error { return nil } t.Errorf("FinalText 改写被覆盖:实际 %q", sc.FinalText) } } + +// §13.6:输出通道 payload 走共享内存调用帧。 +// +// 用**真实 SDK 模板**编译的插件验证,而不是手写 testdata——因为生产插件 +// (如 qq)走的就是模板,模板不读帧的话这个改动等于没做。 +// 链路:内核 Alloc 帧 → 写 payload → RPC 只传偏移描述符 → 模板 frameInput +// 读出 → 交给插件的输出 handler。 +func TestE2E_RealTemplateOutputPayloadViaFrame(t *testing.T) { + bin := buildPluginWithRealTemplate(t, e2ePluginSource) + + host, err := NewHost() + if err != nil { + t.Fatalf("NewHost: %v", err) + } + defer host.Close() + + core := newFakeCore() + p := New("e2e", bin, t.TempDir(), nil, host, nil) + if err := p.Start(core); err != nil { + t.Fatalf("Start: %v", err) + } + defer p.Close() + + core.mu.Lock() + h, ok := core.outputs["e2e_out"] + core.mu.Unlock() + if !ok { + t.Fatal("插件应注册 e2e_out 输出通道") + } + + // 9000 字节:超过任何内联预算,只有走帧才可能完整送达。 + payload := strings.Repeat("输出", 3000) + res, err := h(map[string]interface{}{"payload": payload, "type": "text"}) + if err != nil { + t.Fatalf("发送应成功: %v", err) + } + m, _ := res.(map[string]interface{}) + var gotLen int + switch v := m["payload_len"].(type) { + case float64: + gotLen = int(v) + case int: + gotLen = v + } + if gotLen != len(payload) { + t.Fatalf("模板插件收到的 payload 长度 = %d,期望 %d(模板未从帧读参数?)", gotLen, len(payload)) + } + if used, _ := host.Arena().Stats(); used != 0 { + t.Fatalf("调用结束后 arena 应归零,实际 used=%d", used) + } +} diff --git a/internal/plugin/proc/plugin.go b/internal/plugin/proc/plugin.go index c24182c..dcf9913 100644 --- a/internal/plugin/proc/plugin.go +++ b/internal/plugin/proc/plugin.go @@ -338,10 +338,39 @@ func (p *Plugin) invokeOutput(channel string, args map[string]interface{}) (inte if p.proc == nil { return nil, ErrProcessExited } - raw, err := p.proc.Call(MethodOutputInvoke, OutputInvokeParams{ - Channel: channel, - Args: args, - }) + + argJSON, err := json.Marshal(args) + if err != nil { + return nil, fmt.Errorf("proc: %s 输出通道 %s 参数序列化失败: %w", p.name, channel, err) + } + + // 与工具调用同一个 funccall 帧模型(§13.6):payload 全在共享内存, + // RPC 只传偏移描述符。输出 payload 在真机上可能含图片/文件描述等大字段, + // 内联会撑爆 stdin/stdout 管道。 + // + // 帧尾不预留结果区:output 的应答很小("ok" 或 status map),直接走 + // RPC 应答字段即可,不像 tool.invoke 那样需要几十 KB 的结果预算。 + var params OutputInvokeParams + params.Channel = channel + if len(argJSON) > 0 { + arena := p.host.Arena() + gen := p.host.Generation() + frame, err := arena.Alloc(OwnerHost, len(argJSON), gen) + if err != nil { + return nil, fmt.Errorf("proc: %s 输出通道 %s 分配调用帧失败: %w", p.name, channel, err) + } + defer func() { _ = arena.Free(OwnerHost, frame) }() + + area, err := arena.Read(frame, gen) + if err != nil { + return nil, fmt.Errorf("proc: %s 输出通道 %s 读取调用帧失败: %w", p.name, channel, err) + } + copy(area[:len(argJSON)], argJSON) + params.Frame = frame + params.ArgsLen = uint32(len(argJSON)) + } + + raw, err := p.proc.Call(MethodOutputInvoke, params) if err != nil { return nil, err // 真实失败上报,模型可感知并重试 } @@ -349,6 +378,27 @@ func (p *Plugin) invokeOutput(channel string, args map[string]interface{}) (inte return map[string]interface{}{"status": "ok"}, nil } + // 结果可能在共享内存里(插件把大结果写回帧结果区)。 + // 先按结构化应答解析;ResultRef 为零则回退到内联字段。 + var outRes OutputInvokeResult + if jerr := json.Unmarshal(raw, &outRes); jerr == nil && !outRes.ResultRef.IsZero() { + arena := p.host.Arena() + gen := p.host.Generation() + // 插件申请了扩容块:内核负责归还(与 invokeTool 一致)。 + if outRes.ResultRef.Flags&sharedRefFlagExpand != 0 { + defer func() { _ = arena.Free(OwnerHost, outRes.ResultRef) }() + } + data, err := arena.Read(outRes.ResultRef, gen) + if err != nil { + return nil, fmt.Errorf("proc: %s 输出通道 %s 读取共享结果失败: %w", p.name, channel, err) + } + var out interface{} + if err := json.Unmarshal(data, &out); err != nil { + return nil, fmt.Errorf("proc: %s 输出通道 %s 解析共享结果失败: %w", p.name, channel, err) + } + return out, nil + } + var res map[string]interface{} if err := json.Unmarshal(raw, &res); err == nil { if _, ok := res["status"]; !ok { diff --git a/internal/plugin/proc/plugin_test.go b/internal/plugin/proc/plugin_test.go index 941de0a..26718e7 100644 --- a/internal/plugin/proc/plugin_test.go +++ b/internal/plugin/proc/plugin_test.go @@ -290,6 +290,54 @@ func TestPlugin_RegisteredToolInvokable(t *testing.T) { } } +// §13.6:输出通道 payload 走共享内存调用帧(与 tool.invoke 同一模型), +// 大 payload 不再爆 stdin/stdout 管道。 +func TestPlugin_OutputPayloadViaArena(t *testing.T) { + bin := buildTestPlugin(t, "stageplugin.go") + core := newFakeCore() + + host, err := NewHost() + if err != nil { + t.Fatalf("NewHost: %v", err) + } + defer host.Close() + + p := New("demo", bin, t.TempDir(), nil, host, nil) + if err := p.Start(core); err != nil { + t.Fatalf("Start: %v", err) + } + defer p.Close() + + core.mu.Lock() + h, ok := core.outputs["demo_ch"] + core.mu.Unlock() + if !ok { + t.Fatal("插件应注册 demo_ch 输出通道") + } + + // 9000 字节,明显超过任何内联预算 + payload := strings.Repeat("输出", 3000) + res, err := h(map[string]interface{}{"payload": payload, "type": "text"}) + if err != nil { + t.Fatalf("发送应成功: %v", err) + } + m, _ := res.(map[string]interface{}) + // 插件回报它实际收到的长度:只有完整 payload 经帧送达才等于发送长度。 + var gotLen int + switch v := m["payload_len"].(type) { + case float64: + gotLen = int(v) + case int: + gotLen = v + } + if gotLen != len(payload) { + t.Fatalf("插件收到的 payload 长度 = %d,期望 %d(帧未把完整 payload 带到插件侧)", gotLen, len(payload)) + } + if used, total := host.Arena().Stats(); used != 0 { + t.Fatalf("调用结束后 arena 应归零,实际 used=%d/%d", used, total) + } +} + // 输出通道**同步等真实结果**:失败必须上报(§9.4 根治)。 func TestPlugin_OutputChannelReportsRealFailure(t *testing.T) { bin := buildTestPlugin(t, "stageplugin.go") diff --git a/internal/plugin/proc/protocol.go b/internal/plugin/proc/protocol.go index f6cba77..90e96ff 100644 --- a/internal/plugin/proc/protocol.go +++ b/internal/plugin/proc/protocol.go @@ -278,8 +278,22 @@ const ( // 永远返回成功(§9.4,现网 2 次消息发不出而模型以为成功)。 // 进程模型下 RPC 天然可等应答,该缺陷从根上消失。 type OutputInvokeParams struct { - Channel string `json:"channel"` - Args map[string]interface{} `json:"args,omitempty"` + Channel string `json:"channel"` + // Frame 是调用帧:[0, ArgsLen) 是参数 JSON;[ArgsLen, Frame.Length) 是结果区。 + // 与 ToolInvokeParams 同一 funccall 帧模型(§13.3)。没有帧时(直连 RPC + // 测试)回退到 Args。大 payload 不再爆 stdin/stdout 管道——这正是同步 + // output_send 卡住的根因之一。 + Frame SharedRef `json:"frame,omitempty"` + ArgsLen uint32 `json:"args_len,omitempty"` + Args map[string]interface{} `json:"args,omitempty"` // 仅直连 RPC 调用方使用 +} + +// OutputInvokeResult 与 ToolInvokeResult 同形:结果优先写进帧结果区, +// 放不下才申扩容块并打 sharedRefFlagExpand。标量响应(如 "ok")直接在 +// Result 字段返回,不走共享内存。 +type OutputInvokeResult struct { + ResultRef SharedRef `json:"result_ref,omitempty"` + Result interface{} `json:"result,omitempty"` // 仅直连 RPC 调用方使用 } // ArenaAllocParams / ArenaAllocResult:插件向内核申请共享内存。 diff --git a/internal/plugin/proc/testdata/stageplugin.go b/internal/plugin/proc/testdata/stageplugin.go index 8960e50..2d8b936 100644 --- a/internal/plugin/proc/testdata/stageplugin.go +++ b/internal/plugin/proc/testdata/stageplugin.go @@ -526,18 +526,40 @@ func main() { }(req.ID) case "output.invoke": - var p struct { - Channel string `json:"channel"` - Args map[string]interface{} `json:"args"` - } - json.Unmarshal(req.Params, &p) - payload, _ := p.Args["payload"].(string) - if payload == "fail" { - // 模拟现网 qq 插件的真实失败:meta 缺 user_id - send(response{ID: req.ID, Error: "meta 中缺少 user_id 字段"}) - continue - } - send(response{ID: req.ID, Result: map[string]interface{}{"status": "sent"}}) + // §13.6:payload 走共享内存调用帧,与 tool.invoke 同一模型。 + // 必须在独立 goroutine 里处理:插件可能反向调用内核。 + go func(id uint64, raw json.RawMessage) { + var p struct { + Channel string `json:"channel"` + Args map[string]interface{} `json:"args"` + Frame SharedRef `json:"frame"` + ArgsLen uint32 `json:"args_len"` + } + json.Unmarshal(raw, &p) + + args := p.Args + if !p.Frame.IsZero() { + if blob := frameInput(p.Frame, p.ArgsLen); len(blob) > 0 { + var decoded map[string]interface{} + if err := json.Unmarshal(blob, &decoded); err != nil { + send(response{ID: id, Error: "解析输出参数: " + err.Error()}) + return + } + args = decoded + } + } + payload, _ := args["payload"].(string) + if payload == "fail" { + // 模拟现网 qq 插件的真实失败:meta 缺 user_id + send(response{ID: id, Error: "meta 中缺少 user_id 字段"}) + return + } + // 标量结果直接在应答字段返回(output 应答很小,无需走共享内存)。 + // payload_len 让测试能验证共享帧把完整 payload 带到了插件侧。 + send(response{ID: id, Result: map[string]interface{}{ + "status": "sent", "payload_len": len(payload), + }}) + }(req.ID, req.Params) default: if req.ID != 0 { diff --git a/plan.md b/plan.md index ceb90c0..ff39358 100644 --- a/plan.md +++ b/plan.md @@ -1239,26 +1239,50 @@ SDK 仓 `v1.0.0` / `v1.1.0`。main 的版本路牌现为 `1.2.0`(尚无 tag) **目标**:输入通道消息走共享内存。 +**现状(核实)**:基础设施就位,但内核侧未接线—— + +- `injectParams.TextRef SharedRef` 与 `resolveText`(corehandler.go:585) + 在插件→内核方向可用,但内核→插件方向的 `procCore.InjectText` 仍直接传 + 字符串给 `sdk.InjectText`,从不 `Alloc/Put` 构造 `TextRef`。 +- 即只做了半边(插件回传文本),内核注入还没走共享内存。 + **实施**: 1. InputChannel Slot:source/channel/text/blocks 写入 arena 2. 插件消费后设 DONE -3. 同步注入仍走 RPC + SharedRef +3. 同步注入仍走 RPC + SharedRef ← 需要接线内核侧 4. 异步注入改写 arena + eventfd **验证**: - [ ] 真实 QQ 消息注入测试 -- [ ] git commit -m "feat(shm): input channel lane" +- [ ] git commit -m "feat(shm): input channel lane" ← 未提交,内核侧待接线 ### 13.6 OutputChannel lane **目标**:输出通道消息走共享内存。 +**实施**(已完成): + +1. `OutputInvokeParams` 加 `Frame SharedRef` + `ArgsLen`(`Args` 仅留给直连 + RPC 的测试),与 `ToolInvokeParams` 同一 funccall 帧模型 +2. `invokeOutput` Alloc 帧、写 payload JSON、只传偏移描述符;插件退出/调用完 + 后内核归还整帧(`defer arena.Free`) +3. 帧尾不预留结果区:output 应答很小(`"ok"` / status map),直接走 RPC + 应答字段;若插件把大结果写回帧(`OutputInvokeResult.ResultRef`)也能读回, + 并识别 `sharedRefFlagExpand` 单独归还扩容块 +4. 模板 `output.invoke` 从帧读参数(`frameInput`),无帧才回退内联 `Args` + **验证**: -- [ ] QQ 输出正常 -- [ ] git commit -m "feat(shm): output channel lane" +- [x] `TestPlugin_OutputPayloadViaArena`:9000 字节 payload 经帧完整送达(插件 + 回报实收长度)+ 调用后 arena 归零 +- [x] `TestE2E_RealTemplateOutputPayloadViaFrame`:同上但用**真实 SDK 模板** + 编译的插件(生产插件走的就是模板,模板不读帧该改动就等于没做) +- [x] `TestPlugin_OutputChannelReportsRealFailure` 仍绿:同步等真实结果、 + 失败必须上报(§9.4)不被破坏 +- [ ] QQ 输出正常(需生产部署后验证) +- [x] git commit -m "feat(shm): output channel lane" ### 13.7 RuntimeManager + 分组 worker