mirror of
https://gitcode.com/JianFeeeee/HomeAgent.git
synced 2026-10-03 15:53:56 +00:00
feat(shm): 输出通道 payload 走共享内存调用帧(§13.6)
§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),不是仅打勾的项。
This commit is contained in:
@ -62,6 +62,19 @@ func (p *e2ePlugin) Start(s *sdk.PluginSDK) error {
|
|||||||
ctx.FinalText = ctx.FinalText + "|stage-touched"
|
ctx.FinalText = ctx.FinalText + "|stage-touched"
|
||||||
return nil
|
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
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
@ -273,3 +286,54 @@ func (p *readerPlugin) Stop() error { return nil }
|
|||||||
t.Errorf("FinalText 改写被覆盖:实际 %q", sc.FinalText)
|
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)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@ -338,10 +338,39 @@ func (p *Plugin) invokeOutput(channel string, args map[string]interface{}) (inte
|
|||||||
if p.proc == nil {
|
if p.proc == nil {
|
||||||
return nil, ErrProcessExited
|
return nil, ErrProcessExited
|
||||||
}
|
}
|
||||||
raw, err := p.proc.Call(MethodOutputInvoke, OutputInvokeParams{
|
|
||||||
Channel: channel,
|
argJSON, err := json.Marshal(args)
|
||||||
Args: 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 {
|
if err != nil {
|
||||||
return nil, err // 真实失败上报,模型可感知并重试
|
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
|
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{}
|
var res map[string]interface{}
|
||||||
if err := json.Unmarshal(raw, &res); err == nil {
|
if err := json.Unmarshal(raw, &res); err == nil {
|
||||||
if _, ok := res["status"]; !ok {
|
if _, ok := res["status"]; !ok {
|
||||||
|
|||||||
@ -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 根治)。
|
// 输出通道**同步等真实结果**:失败必须上报(§9.4 根治)。
|
||||||
func TestPlugin_OutputChannelReportsRealFailure(t *testing.T) {
|
func TestPlugin_OutputChannelReportsRealFailure(t *testing.T) {
|
||||||
bin := buildTestPlugin(t, "stageplugin.go")
|
bin := buildTestPlugin(t, "stageplugin.go")
|
||||||
|
|||||||
@ -278,8 +278,22 @@ const (
|
|||||||
// 永远返回成功(§9.4,现网 2 次消息发不出而模型以为成功)。
|
// 永远返回成功(§9.4,现网 2 次消息发不出而模型以为成功)。
|
||||||
// 进程模型下 RPC 天然可等应答,该缺陷从根上消失。
|
// 进程模型下 RPC 天然可等应答,该缺陷从根上消失。
|
||||||
type OutputInvokeParams struct {
|
type OutputInvokeParams struct {
|
||||||
Channel string `json:"channel"`
|
Channel string `json:"channel"`
|
||||||
Args map[string]interface{} `json:"args,omitempty"`
|
// 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:插件向内核申请共享内存。
|
// ArenaAllocParams / ArenaAllocResult:插件向内核申请共享内存。
|
||||||
|
|||||||
46
internal/plugin/proc/testdata/stageplugin.go
vendored
46
internal/plugin/proc/testdata/stageplugin.go
vendored
@ -526,18 +526,40 @@ func main() {
|
|||||||
}(req.ID)
|
}(req.ID)
|
||||||
|
|
||||||
case "output.invoke":
|
case "output.invoke":
|
||||||
var p struct {
|
// §13.6:payload 走共享内存调用帧,与 tool.invoke 同一模型。
|
||||||
Channel string `json:"channel"`
|
// 必须在独立 goroutine 里处理:插件可能反向调用内核。
|
||||||
Args map[string]interface{} `json:"args"`
|
go func(id uint64, raw json.RawMessage) {
|
||||||
}
|
var p struct {
|
||||||
json.Unmarshal(req.Params, &p)
|
Channel string `json:"channel"`
|
||||||
payload, _ := p.Args["payload"].(string)
|
Args map[string]interface{} `json:"args"`
|
||||||
if payload == "fail" {
|
Frame SharedRef `json:"frame"`
|
||||||
// 模拟现网 qq 插件的真实失败:meta 缺 user_id
|
ArgsLen uint32 `json:"args_len"`
|
||||||
send(response{ID: req.ID, Error: "meta 中缺少 user_id 字段"})
|
}
|
||||||
continue
|
json.Unmarshal(raw, &p)
|
||||||
}
|
|
||||||
send(response{ID: req.ID, Result: map[string]interface{}{"status": "sent"}})
|
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:
|
default:
|
||||||
if req.ID != 0 {
|
if req.ID != 0 {
|
||||||
|
|||||||
32
plan.md
32
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
|
1. InputChannel Slot:source/channel/text/blocks 写入 arena
|
||||||
2. 插件消费后设 DONE
|
2. 插件消费后设 DONE
|
||||||
3. 同步注入仍走 RPC + SharedRef
|
3. 同步注入仍走 RPC + SharedRef ← 需要接线内核侧
|
||||||
4. 异步注入改写 arena + eventfd
|
4. 异步注入改写 arena + eventfd
|
||||||
|
|
||||||
**验证**:
|
**验证**:
|
||||||
|
|
||||||
- [ ] 真实 QQ 消息注入测试
|
- [ ] 真实 QQ 消息注入测试
|
||||||
- [ ] git commit -m "feat(shm): input channel lane"
|
- [ ] git commit -m "feat(shm): input channel lane" ← 未提交,内核侧待接线
|
||||||
|
|
||||||
### 13.6 OutputChannel 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 输出正常
|
- [x] `TestPlugin_OutputPayloadViaArena`:9000 字节 payload 经帧完整送达(插件
|
||||||
- [ ] git commit -m "feat(shm): output channel lane"
|
回报实收长度)+ 调用后 arena 归零
|
||||||
|
- [x] `TestE2E_RealTemplateOutputPayloadViaFrame`:同上但用**真实 SDK 模板**
|
||||||
|
编译的插件(生产插件走的就是模板,模板不读帧该改动就等于没做)
|
||||||
|
- [x] `TestPlugin_OutputChannelReportsRealFailure` 仍绿:同步等真实结果、
|
||||||
|
失败必须上报(§9.4)不被破坏
|
||||||
|
- [ ] QQ 输出正常(需生产部署后验证)
|
||||||
|
- [x] git commit -m "feat(shm): output channel lane"
|
||||||
|
|
||||||
### 13.7 RuntimeManager + 分组 worker
|
### 13.7 RuntimeManager + 分组 worker
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user