mirror of
https://gitcode.com/JianFeeeee/homeagent-sdk.git
synced 2026-09-20 00:48:12 +00:00
内核的输入调度器区分两类别:中断输入(可抢占)与排队输入(可被任何中断打断)。
中断的级别是“这项工作有多不能等”的声明,由插件在注入时给出:
p.sdk.InjectInterruptTextOpts(src, ch, text, sdk.InjectOptions{
NoMemory: true,
Priority: sdk.PriorityL2, // L1 完全可等 / L2 一般提醒 / L3 需及时
})
- 新增 `InjectOptions.Priority string` 与 `PriorityL1/L2/L3` 常量(纯追加)。
- 空/非法值一律降级为 L1(默认级)——拼写错误不会被静默当成别的级别。
- **L4 由内核独占**(panic 中断、内核事件中断 selfip),插件声明 L4 会被内核
夹到 L3,远端常量的取值域里也不提供 L4。
- 排队注入(InjectText*/InjectInputSync*)没有级别:它们本就是“不需及时处理”
的那一类,可被任何中断打断;传了 Priority 也不会生效。
- 贯通链路:sdk.InjectOptions -> proc RPC 参数(priority)-> 内核 payload;
tools/hmapdev 模板同步透传(三个注入的 6 个 Opts 变体共用 applyInjectOpts)。
- example/qq 显式声明 L1:QQ 消息既不是时钟那样的实时工作,也不是紧急工作。
兼容性:零值等价于旧行为(L1),既有插件无需改动。
1772 lines
50 KiB
Cheetah
1772 lines
50 KiB
Cheetah
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 一致)----
|
||
|
||
// procProtocolVersion 必须与内核的 proc.ProtocolVersion 完全一致。
|
||
//
|
||
// v2:内核→插件的 payload 改用调用帧(tool/cleaner/output),媒体块改走
|
||
// blocks_ref。v1 插件只读内联 args,遇上 v2 内核会拿到空参数;反过来 v2
|
||
// 插件发 blocks_ref,v1 内核也会静默忽略。两边错配都不报错、只是静默失效,
|
||
// 所以靠这个常量在握手上显式拦下。
|
||
const procProtocolVersion = 2
|
||
|
||
// ---- 统一共享内存区域布局(与内核 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
|
||
sbOffArenaOff = 36
|
||
sbOffArenaCap = 40
|
||
sbOffArenaUsed = 44
|
||
)
|
||
|
||
// ---- 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
|
||
)
|
||
|
||
// SharedRef 跨进程共享内存描述符。
|
||
//
|
||
// ⚠️ 这是**内部实现细节**:插件开发者永远看不到它。公开 SDK 只暴露普通
|
||
// 字符串与 Map;模板运行时在传输层按 payload 大小自动选择内联 JSON 还是
|
||
// 共享槽。直接使用 SharedRef 属于运行时内部行为,不是插件 API。
|
||
type SharedRef struct {
|
||
Offset uint32 `json:"offset"`
|
||
Length uint32 `json:"length"`
|
||
Generation uint32 `json:"generation"`
|
||
Flags uint32 `json:"flags"`
|
||
}
|
||
|
||
func (r SharedRef) IsZero() bool { return r.Offset == 0 && r.Length == 0 }
|
||
func (r SharedRef) Slice(data []byte) []byte {
|
||
if r.IsZero() || int(r.Offset)+int(r.Length) > len(data) {
|
||
return nil
|
||
}
|
||
return data[r.Offset : r.Offset+r.Length]
|
||
}
|
||
|
||
// SharedRef.Flags 语义位(须与内核 internal/plugin/proc/arena.go 一致)。
|
||
const (
|
||
sharedRefFlagJSON = 1 << 0 // 载荷是 JSON
|
||
sharedRefFlagExpand = 1 << 1 // 引用指向插件申请的扩容块
|
||
)
|
||
|
||
// region 是内核传入的统一共享区域 mmap(handshake 时设置)。
|
||
var region []byte
|
||
|
||
// 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
|
||
}
|
||
var r struct {
|
||
Ref SharedRef `json:"ref"`
|
||
}
|
||
if err := json.Unmarshal(raw, &r); err != nil {
|
||
return SharedRef{}, err
|
||
}
|
||
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})
|
||
}
|
||
|
||
// ---- 调用帧(funccall 模型)辅助 ----
|
||
//
|
||
// 工具调用/清洗由内核发起:内核标定一块内存帧交给插件,插件在帧内工作,
|
||
// 只有结果超出内核预留的预算时才向内核申请扩容块。
|
||
|
||
// frameInput 返回帧内的输入段(内核写入的参数/输入文本)。
|
||
func frameInput(frame SharedRef, inputLen uint32) []byte {
|
||
if frame.IsZero() || int(inputLen) > len(frame.Slice(region)) {
|
||
return nil
|
||
}
|
||
return frame.Slice(region)[:inputLen]
|
||
}
|
||
|
||
// frameOutput 尝试把 payload 写进帧的结果区(帧内 [inputLen, frame.Length))。
|
||
// 放不下时返回错误,由调用方决定是否申请扩容块。
|
||
func frameOutput(frame SharedRef, inputLen uint32, payload []byte, jsonFlag bool) (SharedRef, error) {
|
||
if frame.IsZero() {
|
||
return SharedRef{}, fmt.Errorf("无调用帧")
|
||
}
|
||
area := frame.Slice(region)
|
||
start := int(inputLen)
|
||
if start > len(area) || len(payload) > len(area)-start {
|
||
return SharedRef{}, fmt.Errorf("帧内空间不足(需 %d,剩 %d)", len(payload), len(area)-start)
|
||
}
|
||
copy(region[frame.Offset+uint32(start):], payload)
|
||
ref := SharedRef{
|
||
Offset: frame.Offset + uint32(start),
|
||
Length: uint32(len(payload)),
|
||
Generation: frame.Generation,
|
||
}
|
||
if jsonFlag {
|
||
ref.Flags |= sharedRefFlagJSON
|
||
}
|
||
return ref, nil
|
||
}
|
||
|
||
// arenaPut 申请一块扩容块并写入 payload,引用上打 sharedRefFlagExpand
|
||
// 告知内核该块需单独归还(插件只申请,回收由内核做)。
|
||
func arenaPut(payload []byte, jsonFlag bool) (SharedRef, error) {
|
||
ref, err := arenaAlloc(uint32(len(payload)))
|
||
if err != nil {
|
||
return SharedRef{}, err
|
||
}
|
||
if len(payload) > int(ref.Length) {
|
||
arenaFree(ref)
|
||
return SharedRef{}, fmt.Errorf("扩容块容量不足(需 %d,得 %d)", len(payload), ref.Length)
|
||
}
|
||
copy(region[ref.Offset:ref.Offset+uint32(len(payload))], payload)
|
||
ref.Length = uint32(len(payload))
|
||
ref.Flags |= sharedRefFlagExpand
|
||
if jsonFlag {
|
||
ref.Flags |= sharedRefFlagJSON
|
||
}
|
||
return ref, nil
|
||
}
|
||
|
||
// 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
|
||
}
|
||
|
||
// putValueInArena 把任意值 JSON 序列化后放进共享槽,太小或 arena 不可用时
|
||
// 返回 false(调用方退到内联)。
|
||
func putValueInArena(v interface{}) (SharedRef, bool) {
|
||
blob, err := json.Marshal(v)
|
||
if err != nil {
|
||
return SharedRef{}, false
|
||
}
|
||
return putInArena(string(blob))
|
||
}
|
||
|
||
// callWithText 按 payload 大小自动选择共享槽或内联,发起一次带文本的业务 RPC。
|
||
//
|
||
// 共享内存对插件开发者完全透明:SDK 层只看得到 string。
|
||
func callWithText(method, source, channel, text string) (json.RawMessage, error) {
|
||
return callWithTextOpts(method, source, channel, text, sdk.InjectOptions{})
|
||
}
|
||
|
||
// callWithTextOpts 是 callWithText 的带标志位版本。
|
||
//
|
||
// 只在标志位非零时才写入参数:零值(记入记忆 + 不裁剪)与旧参数形态完全一致,
|
||
// 便于内核侧做兼容与灰度。
|
||
func callWithTextOpts(method, source, channel, text string, opts sdk.InjectOptions) (json.RawMessage, error) {
|
||
if ref, ok := putInArena(text); ok {
|
||
defer arenaFree(ref)
|
||
args := map[string]interface{}{
|
||
"source": source, "channel": channel, "text_ref": ref,
|
||
}
|
||
applyInjectOpts(args, opts)
|
||
return callCore(method, args)
|
||
}
|
||
args := map[string]interface{}{
|
||
"source": source, "channel": channel, "text": text,
|
||
}
|
||
applyInjectOpts(args, opts)
|
||
return callCore(method, args)
|
||
}
|
||
|
||
// applyInjectOpts 把 InjectOptions 摊进注入参数字典(仅非零值)。
|
||
func applyInjectOpts(args map[string]interface{}, opts sdk.InjectOptions) {
|
||
if opts.NoMemory {
|
||
args["no_memory"] = true
|
||
}
|
||
if opts.ContextPolicy != "" {
|
||
args["context_policy"] = opts.ContextPolicy
|
||
}
|
||
if opts.CleanerName != "" {
|
||
args["cleaner_name"] = opts.CleanerName
|
||
}
|
||
// priority 只对中断注入有意义(排队注入没有级别)。
|
||
if opts.Priority != "" {
|
||
args["priority"] = opts.Priority
|
||
}
|
||
}
|
||
|
||
// ---- 全局状态 ----
|
||
|
||
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
|
||
|
||
// region 见文件头部 SharedRef 注释(handshake 时设置)。
|
||
|
||
// 事件环(§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,
|
||
// 整个结构体:手写字段白名单会把新增字段静默丢掉
|
||
// (ChannelDef.Cleaner 已标 json:"-",可以整体 marshal)。
|
||
"def": def,
|
||
"has_cleaner": def.Cleaner != nil,
|
||
})
|
||
})
|
||
return base
|
||
}
|
||
|
||
type procIO struct{}
|
||
|
||
// 下面六个三参数方法是 *Opts 变体的零值糖:记入记忆 + 不裁剪。
|
||
|
||
func (procIO) InjectText(s, c, t string) {
|
||
procIO{}.InjectTextOpts(s, c, t, sdk.InjectOptions{})
|
||
}
|
||
func (procIO) InjectInterruptText(s, c, t string) {
|
||
procIO{}.InjectInterruptTextOpts(s, c, t, sdk.InjectOptions{})
|
||
}
|
||
func (procIO) InjectTextNoMemory(s, c, t string) {
|
||
procIO{}.InjectTextOpts(s, c, t, sdk.InjectOptions{NoMemory: true})
|
||
}
|
||
func (procIO) InjectInputSync(s, c, t string) string {
|
||
return procIO{}.InjectInputSyncOpts(s, c, t, sdk.InjectOptions{})
|
||
}
|
||
|
||
// 以下为带标志位的注入:opts 决定这次注入是否进记忆、是否据此裁剪上下文。
|
||
|
||
func (procIO) InjectTextOpts(s, c, t string, opts sdk.InjectOptions) {
|
||
// 忽略错误:注入是 fire-and-forget,与内联路径语义一致
|
||
_, _ = callWithTextOpts("io.injectText", s, c, t, opts)
|
||
}
|
||
func (procIO) InjectInterruptTextOpts(s, c, t string, opts sdk.InjectOptions) {
|
||
_, _ = callWithTextOpts("io.injectInterrupt", s, c, t, opts)
|
||
}
|
||
func (procIO) InjectInputSyncOpts(s, c, t string, opts sdk.InjectOptions) string {
|
||
raw, err := callWithTextOpts("io.injectInputSync", s, c, t, opts)
|
||
if err != nil {
|
||
return ""
|
||
}
|
||
var r struct {
|
||
Reply string `json:"reply"`
|
||
}
|
||
json.Unmarshal(raw, &r)
|
||
return r.Reply
|
||
}
|
||
func (procIO) SetToolBlocks(blocks []sdk.ContentBlock) {
|
||
// 媒体块经共享内存(blocks_ref):本地生成的图/音频是 base64 data URL,
|
||
// 一张图可达数 MB;内联时整份 base64 还要在 RPC 报文里再编码/再拷贝一遍。
|
||
// 更重要的是内容本体落在共享段里,插件回调才能就地改写。
|
||
// 小 payload(如纯文本块)仍走内联,省一次 RPC。
|
||
args := map[string]interface{}{}
|
||
if ref, ok := putValueInArena(blocks); ok {
|
||
defer arenaFree(ref)
|
||
args["blocks_ref"] = ref
|
||
} else {
|
||
args["blocks"] = blocks
|
||
}
|
||
if err := callCoreVoid("io.setToolBlocks", args); err != nil {
|
||
log.Printf("SetToolBlocks: %v", err)
|
||
}
|
||
}
|
||
|
||
// 带媒体的注入:插件主动发起一轮带图/音频的对话。
|
||
// 与 SetToolBlocks 的区别是媒体在**本轮**就到模型手上,而不是等下一条 tool message。
|
||
func (procIO) InjectInputMedia(s, c, t string, blocks []sdk.ContentBlock) {
|
||
procIO{}.InjectInputMediaOpts(s, c, t, blocks, sdk.InjectOptions{})
|
||
}
|
||
|
||
func (procIO) InjectInputMediaSync(s, c, t string, blocks []sdk.ContentBlock) string {
|
||
return procIO{}.InjectInputMediaSyncOpts(s, c, t, blocks, sdk.InjectOptions{})
|
||
}
|
||
|
||
func (procIO) InjectInterruptMedia(s, c, t string, blocks []sdk.ContentBlock) {
|
||
procIO{}.InjectInterruptMediaOpts(s, c, t, blocks, sdk.InjectOptions{})
|
||
}
|
||
|
||
func (procIO) InjectInputMediaOpts(s, c, t string, blocks []sdk.ContentBlock, opts sdk.InjectOptions) {
|
||
args, free := mediaArgsOwned(s, c, t, blocks)
|
||
defer free()
|
||
applyInjectOpts(args, opts)
|
||
callCoreVoid("io.injectMedia", args)
|
||
}
|
||
|
||
func (procIO) InjectInputMediaSyncOpts(s, c, t string, blocks []sdk.ContentBlock, opts sdk.InjectOptions) string {
|
||
args, free := mediaArgsOwned(s, c, t, blocks)
|
||
defer free()
|
||
applyInjectOpts(args, opts)
|
||
raw, err := callCore("io.injectMediaSync", args)
|
||
if err != nil {
|
||
return ""
|
||
}
|
||
var r struct {
|
||
Reply string `json:"reply"`
|
||
}
|
||
json.Unmarshal(raw, &r)
|
||
return r.Reply
|
||
}
|
||
|
||
func (procIO) InjectInterruptMediaOpts(s, c, t string, blocks []sdk.ContentBlock, opts sdk.InjectOptions) {
|
||
args, free := mediaArgsOwned(s, c, t, blocks)
|
||
defer free()
|
||
applyInjectOpts(args, opts)
|
||
callCoreVoid("io.injectInterruptMedia", args)
|
||
}
|
||
|
||
// mediaArgsOwned 构造媒体注入参数,并返回释放函数。
|
||
//
|
||
// 为什么要返回释放函数而不是自己 defer:调用方可能是需要等应答的同步调用
|
||
// (injectMediaSync),槽在应答到达前不能被回收,否则内核读到的是已释放的内存。
|
||
func mediaArgsOwned(s, c, t string, blocks []sdk.ContentBlock) (map[string]interface{}, func()) {
|
||
args := map[string]interface{}{"source": s, "channel": c, "text": t}
|
||
ref, ok := putValueInArena(blocks)
|
||
if !ok {
|
||
args["blocks"] = blocks
|
||
return args, func() {}
|
||
}
|
||
args["blocks_ref"] = ref
|
||
return args, func() { arenaFree(ref) }
|
||
}
|
||
|
||
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", docInsertArgs(d))
|
||
}
|
||
|
||
// docInsertArgs 构造 doc.insert 参数:正文优先走共享内存。
|
||
//
|
||
// doc_content 可达几十 KB~数 MB,内联时整份要在 RPC 报文里再编码再拷贝一遍;
|
||
// 且内容本体落在共享段里,插件回调才能就地改写。小文档仍内联。
|
||
func docInsertArgs(d *sdk.Doc) map[string]interface{} {
|
||
if ref, ok := putValueInArena(d); ok {
|
||
defer arenaFree(ref)
|
||
return map[string]interface{}{"doc_ref": ref}
|
||
}
|
||
return map[string]interface{}{"doc": d}
|
||
}
|
||
|
||
// InsertWithMedia 写入文档并关联媒体。
|
||
//
|
||
// 内核会把 `[mime <短digest>] <描述>` 标记补进 Content 并挂上引用,回传的
|
||
// doc 带着补好的 Content/ID/MediaDigests——回写进 d 让调用方能拿到这些。
|
||
func (procDocMemory) InsertWithMedia(d *sdk.Doc, atts []sdk.MediaAttachment) error {
|
||
// 文档正文与附件(含媒体 data URL)都优先走共享内存。
|
||
args := map[string]interface{}{}
|
||
var frees []func()
|
||
defer func() {
|
||
for _, f := range frees {
|
||
f()
|
||
}
|
||
}()
|
||
if ref, ok := putValueInArena(d); ok {
|
||
frees = append(frees, func() { arenaFree(ref) })
|
||
args["doc_ref"] = ref
|
||
} else {
|
||
args["doc"] = d
|
||
}
|
||
if ref, ok := putValueInArena(atts); ok {
|
||
frees = append(frees, func() { arenaFree(ref) })
|
||
args["attachments_ref"] = ref
|
||
} else {
|
||
args["attachments"] = atts
|
||
}
|
||
|
||
raw, err := callCore("doc.insertWithMedia", args)
|
||
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 {
|
||
// 知识正文可达数十 KB,优先走共享内存(内容是 JSON 字符串)。
|
||
if ref, ok := putValueInArena(content); ok {
|
||
defer arenaFree(ref)
|
||
return callCoreVoid("knowledge.add", map[string]interface{}{
|
||
"name": name, "content_ref": ref,
|
||
})
|
||
}
|
||
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"`
|
||
Frame SharedRef `json:"frame"`
|
||
ArgsLen uint32 `json:"args_len"`
|
||
}
|
||
json.Unmarshal(req.Params, &p)
|
||
|
||
// 参数:内核标定帧的前段。只有直连 RPC 的调用方(无帧)才走
|
||
// 内联 Args——生产路径永远走帧。
|
||
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 {
|
||
respondErr(req.ID, fmt.Errorf("解析共享参数: %w", err))
|
||
return
|
||
}
|
||
args = decoded
|
||
}
|
||
}
|
||
|
||
handlerMu.RLock()
|
||
h, ok := toolHandlers[p.Name]
|
||
handlerMu.RUnlock()
|
||
if !ok {
|
||
respondErr(req.ID, fmt.Errorf("未注册的工具: %s", p.Name))
|
||
return
|
||
}
|
||
res, err := h(args)
|
||
if err != nil {
|
||
respondErr(req.ID, err)
|
||
return
|
||
}
|
||
|
||
blob, err := json.Marshal(res)
|
||
if err != nil {
|
||
respondErr(req.ID, fmt.Errorf("序列化结果: %w", err))
|
||
return
|
||
}
|
||
// 结果优先写进内核标定的帧;放不下才申请扩容块。
|
||
outRef, err := frameOutput(p.Frame, p.ArgsLen, blob, true)
|
||
if err != nil {
|
||
outRef, err = arenaPut(blob, true)
|
||
if err != nil {
|
||
respondErr(req.ID, fmt.Errorf("结果扩容失败: %w", err))
|
||
return
|
||
}
|
||
}
|
||
respond(req.ID, map[string]interface{}{"result_ref": outRef})
|
||
|
||
case "cleaner.invoke":
|
||
var p struct {
|
||
Scope string `json:"scope"`
|
||
Name string `json:"name"`
|
||
Frame SharedRef `json:"frame"`
|
||
InputLen uint32 `json:"input_len"`
|
||
}
|
||
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
|
||
}
|
||
|
||
input := string(frameInput(p.Frame, p.InputLen))
|
||
output := cleaner(input)
|
||
|
||
// 结果优先写进内核标定的帧;放不下才申请扩容块。
|
||
outRef, err := frameOutput(p.Frame, p.InputLen, []byte(output), false)
|
||
if err != nil {
|
||
outRef, err = arenaPut([]byte(output), false)
|
||
if err != nil {
|
||
respondErr(req.ID, fmt.Errorf("结果扩容失败: %w", err))
|
||
return
|
||
}
|
||
}
|
||
respond(req.ID, map[string]interface{}{"text_ref": outRef})
|
||
|
||
case "stage.invoke":
|
||
handleStageInvoke(req)
|
||
|
||
case "output.invoke":
|
||
var p struct {
|
||
Channel string `json:"channel"`
|
||
Args map[string]interface{} `json:"args"`
|
||
Frame SharedRef `json:"frame"`
|
||
ArgsLen uint32 `json:"args_len"`
|
||
}
|
||
json.Unmarshal(req.Params, &p)
|
||
// 参数:内核标定帧的前段(§13.6)。只有直连 RPC 的调用方(无帧)
|
||
// 才走内联 Args——生产路径永远走帧,大 payload 不再爆管道。
|
||
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 {
|
||
respondErr(req.ID, fmt.Errorf("解析共享输出参数: %w", err))
|
||
return
|
||
}
|
||
args = decoded
|
||
}
|
||
}
|
||
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(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]
|
||
// region 保存完整 mmap 区域。arena 的偏移与大小由内核在
|
||
// arena.alloc 的应答里下发,插件侧不再自己解析槽池布局。
|
||
region = m
|
||
// 挂载事件环段 + 打开通知句柄(§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
|
||
}
|