diff --git a/internal/plugin/proc/arena.go b/internal/plugin/proc/arena.go index b63e8c0..8587db9 100644 --- a/internal/plugin/proc/arena.go +++ b/internal/plugin/proc/arena.go @@ -1,6 +1,6 @@ package proc -// Exchange Arena:**内核独占管理**的跨进程共享内存槽池(§13.2 重设计)。 +// Exchange Arena:**内核独占管理**的跨进程共享内存块分配器(§13.2 重设计)。 // // ── 所有权模型(架构约束)──────────────────────────────────────────── // @@ -9,84 +9,102 @@ package proc // // 因此分配器只存在于**内核进程内**,用一把普通 sync.Mutex 保护即可: // 不需要任何跨进程原子操作,也不存在"共享游标被两个进程各自更新"这一 -// 类竞态——这正是前两版(bump 游标 / CAS 位图跨进程分配)失败的根本原因。 +// 类竞态——这正是前几版(bump 游标 / CAS 位图跨进程分配)失败的根本原因。 // // 共享内存是**内部实现**,不对插件开发者暴露:插件的公开 API 仍是 // 普通字符串/Map(见 SDK 的 IOInjector / ToolHandler)。模板运行时在 -// 传输层按 payload 大小自动决定走内联 JSON 还是共享槽,开发者无感。 +// 传输层完成全部搬运,开发者无感。 // -// ── 槽布局 ──────────────────────────────────────────────────────── +// ── 为什么是变长块而不是定长槽 ────────────────────────────────────── +// +// 定长槽唯一的理由是"跨进程无法安全地做变长分配"。分配器收回内核后这个 +// 约束消失,于是可以采用真正的变长块分配(first-fit + 邻块合并): +// +// - 内核可以按需**标定**每块大小(funccall 帧模型)而不是一律给固定槽 +// - 大 payload(文件内容、长工具输出)不再受 16KB 槽容量限制 +// - 块用尽时插件才请求扩容,而不是事先把池铺满 +// +// ── 布局 ────────────────────────────────────────────────────────── // // ┌──────────────────────────────────────────────┐ // │ Arena Header 64B │ -// │ magic / version / slotCount / slotSize │ -// │ bitmapOff / slotsOff / reserved │ +// │ magic / version / capacity / blockBase │ // ├──────────────────────────────────────────────┤ -// │ Bitmap:ceil(slotCount/32) 个 uint32,bit=占用│ -// ├──────────────────────────────────────────────┤ -// │ Slot 0: [dataLen(4) | owner(4) | data(...)] │ -// │ ... │ +// │ Block 0: [size | state | owner | prevSize] │ +// │ [data ...] │ +// │ Block 1: ... │ +// │ Block N-1: ...(最后一块的 size 到 arena 末尾)│ // └──────────────────────────────────────────────┘ // -// ── 所有权与回收 ────────────────────────────────────────────────── +// 块头 16B,size 含头并按 8 字节对齐,因此所有数据起点都是 8 字节对齐的。 +// prevSize 让 Free 能 O(1) 找到前驱做向后合并(经典分离/合并做法)。 // -// 每个插件进程从内核领一个不透明 ownerID。槽头记录 owner。 -// - 内核 → 插件:内核直接 Alloc,把 ref 放进 RPC 参数(插件不分配) -// - 插件 → 内核:插件 RPC arena.alloc 申请,用完 RPC arena.free 归还 -// - 插件退出:内核 ReclaimOwner 回收其残留槽,避免崩溃泄漏耗尽池 +// ── 安全边界 ────────────────────────────────────────────────────── // -// 插件归还时校验 owner,防止一个插件释放另一个插件的槽。 +// 插件归还/读取时提交的是 SharedRef(offset 由插件转述)。内核必须把它 +// 当成不可信输入:Free 会重新走一遍块链确认 offset 确实是一个已分配块的 +// 数据起点、owner 匹配,才改动分配器状态;否则伪造的 offset 会直接破坏 +// 块链。Read 同理。 +// +// ── 惰性物理内存 ────────────────────────────────────────────────── +// +// arena 很大(MB 级)但底层是 memfd:**未触碰的页不占物理内存**。所以 +// 把容量开大不会带来常驻内存开销,只是虚拟地址空间。 import ( "fmt" "sync" - "sync/atomic" - "unsafe" ) const ( arenaMagic uint32 = 0x41524E41 // "ARNA" - arenaVersion uint32 = 1 + arenaVersion uint32 = 2 // v2:定长槽 → 变长块 arenaHeaderSize = 64 // Arena Header 字段偏移(相对 arena 起始)。 arOffMagic = 0 arOffVersion = 4 - arOffSlotCount = 8 - arOffSlotSize = 12 - arOffBitmapOff = 16 - arOffSlotsOff = 20 - arOffReserved = 24 + arOffCapacity = 8 // arena 可用字节数 + arOffBlockBase = 12 // 首个块相对 arena 起始的偏移 + arOffBlockUsed = 16 // 已分配的数据字节数(诊断) + arOffReserved = 20 - // 槽头:数据长度 + 所有者。数据紧随其后。 - slotHeaderSize = 8 - slotOffDataLen = 0 - slotOffOwner = 4 + // 块头 16B:size 含头,按 8 字节对齐。 + blockHeaderSize = 16 + blockOffSize = 0 // uint32 总大小(含头) + blockOffState = 4 // uint32 0=free 1=used + blockOffOwner = 8 // uint32 所有者 + blockOffPrevSize = 12 // uint32 前一块总大小(0=无前块) - // OwnerHost 标记由内核分配的槽。插件用 Host.NextOwnerID 分配的值。 + blockFree uint32 = 0 + blockUsed uint32 = 1 + + // minBlockSize 是拆分后允许的最小块(含头)。低于它就不再拆, + // 避免产生一堆无法再利用的碎片。 + minBlockSize = 32 + + // SharedRef.Flags 语义位。 + sharedRefFlagJSON = 1 << 0 // 载荷是 JSON(工具参数/结果) + sharedRefFlagExpand = 1 << 1 // 引用指向插件申请的扩容块(内核需单独归还) + + // OwnerHost 标记由内核分配的块。插件用 Host.NextOwnerID 分配的值。 OwnerHost uint32 = 0 - - // 槽池默认规格:32 槽 × 16KB ≈ 512KB。 - arenaDefaultSlotCount = 32 - arenaDefaultSlotSize = 16 * 1024 ) -// arenaRequiredSize 返回给定规格的槽池所需字节数(含头、位图对齐与槽数组)。 -func arenaRequiredSize(slotCount, slotSize uint32) uint32 { - bitmapBytes := ((slotCount+31)/32)*4 + 7 - bitmapBytes &^= 7 // 位图按 8 字节对齐 - slotsOff := (arenaHeaderSize + bitmapBytes + 7) &^ 7 - return slotsOff + slotCount*slotSize -} +// arenaDefaultCapacity 是统一区域里预留给 arena 的字节数。 +// +// 取 4MB:memfd 惰性分配,未触碰的页不占物理内存,所以开大无成本; +// 但它让单个工具调用可以承载 MB 级 payload,不必再退回内联 RPC。 +const arenaDefaultCapacity = 4 * 1024 * 1024 -// arenaPublishSize 是统一区域里预留给槽池的字节数。 -var arenaPublishSize = arenaRequiredSize(arenaDefaultSlotCount, arenaDefaultSlotSize) +// arenaPublishSize 是统一区域里预留给 arena 的字节数。 +var arenaPublishSize = uint32(arenaDefaultCapacity) -// arenaRegion 是槽池视图。 +// arenaRegion 是块分配器视图。 // // region 始终是**完整**统一区域 mmap:SharedRef.Offset 是相对区域起始的 -// 绝对偏移,因此读写都直接落在 region 上,不需要再换算。 +// 绝对偏移,因此读写都直接落在 region 上。 // // mu 只在**内核进程内**使用——分配器完全由内核持有(见文件头所有权模型)。 type arenaRegion struct { @@ -95,60 +113,50 @@ type arenaRegion struct { region []byte base uint32 // arena 在 region 内的起始偏移 cap uint32 // arena 可用字节数 - - slots uint32 - slotSize uint32 - bmOff uint32 // 位图在 arena 内的偏移 - slotsOff uint32 // 槽数组在 arena 内的偏移 } -// initArena 在统一区域的 arena 段上初始化槽池并返回句柄。 -func initArena(region []byte, base, size uint32, slotCount, slotSize uint32) (*arenaRegion, error) { - if slotCount == 0 || slotSize <= slotHeaderSize { - return nil, fmt.Errorf("arena: 非法规格(slots=%d slotSize=%d)", slotCount, slotSize) +// arenaMinSize 返回块分配器能工作的最小字节数。 +func arenaMinSize() uint32 { + return arenaHeaderSize + blockHeaderSize + minBlockSize +} + +// align8 把 n 向上取整到 8 的倍数(块大小与数据起点都必须 8 字节对齐)。 +func align8(n uint32) uint32 { return (n + 7) &^ 7 } + +// initArena 在统一区域的 arena 段上初始化块分配器。 +func initArena(region []byte, base, size uint32) (*arenaRegion, error) { + if size < arenaMinSize() { + return nil, fmt.Errorf("arena: 段太小(%d 字节,至少需要 %d)", size, arenaMinSize()) } - need := arenaRequiredSize(slotCount, slotSize) - if need > size { - return nil, fmt.Errorf("arena: 段太小(需要 %d,实际 %d)", need, size) - } - if uint64(base)+uint64(need) > uint64(len(region)) { - return nil, fmt.Errorf("arena: 越出映射(base=%d need=%d total=%d)", base, need, len(region)) + if uint64(base)+uint64(size) > uint64(len(region)) { + return nil, fmt.Errorf("arena: 越出映射(base=%d size=%d total=%d)", base, size, len(region)) } - bitmapBytes := ((slotCount+31)/32)*4 + 7 - bitmapBytes &^= 7 - slotsOff := (arenaHeaderSize + bitmapBytes + 7) &^ 7 + a := &arenaRegion{region: region, base: base, cap: size} - ap := region[base : base+need] + blockBase := align8(arenaHeaderSize) + if blockBase+blockHeaderSize+minBlockSize > size { + return nil, fmt.Errorf("arena: 头部过大(blockBase=%d size=%d)", blockBase, size) + } + + ap := region[base : base+size] putU32(ap[arOffMagic:], arenaMagic) putU32(ap[arOffVersion:], arenaVersion) - putU32(ap[arOffSlotCount:], slotCount) - putU32(ap[arOffSlotSize:], slotSize) - putU32(ap[arOffBitmapOff:], arenaHeaderSize) - putU32(ap[arOffSlotsOff:], slotsOff) + putU32(ap[arOffCapacity:], size) + putU32(ap[arOffBlockBase:], blockBase) + putU32(ap[arOffBlockUsed:], 0) - // 显式清零位图与槽头,保证复用已存在区域时状态干净。 - for i := uint32(0); i < bitmapBytes; i++ { - ap[arenaHeaderSize+i] = 0 - } - for i := uint32(0); i < slotCount; i++ { - sb := slotsOff + i*slotSize - putU32(ap[sb+slotOffDataLen:], 0) - putU32(ap[sb+slotOffOwner:], 0) - } + // 初始状态:一整块 free 覆盖剩余空间。 + first := blockBase + putU32(ap[first+blockOffSize:], size-first) + putU32(ap[first+blockOffState:], blockFree) + putU32(ap[first+blockOffOwner:], OwnerHost) + putU32(ap[first+blockOffPrevSize:], 0) - return &arenaRegion{ - region: region, - base: base, - cap: need, - slots: slotCount, - slotSize: slotSize, - bmOff: arenaHeaderSize, - slotsOff: slotsOff, - }, nil + return a, nil } -// attachArena 从已初始化的区域解析槽池(诊断用;模板不再需要解析位图)。 +// attachArena 从已初始化的区域解析分配器(诊断用)。 func attachArena(region []byte, base, size uint32) (*arenaRegion, error) { if size < arenaHeaderSize { return nil, fmt.Errorf("arena: 段太小(%d)", size) @@ -164,122 +172,211 @@ func attachArena(region []byte, base, size uint32) (*arenaRegion, error) { if got := getU32(ap[arOffVersion:]); got != arenaVersion { return nil, fmt.Errorf("arena: 版本不匹配(%d,期望 %d)", got, arenaVersion) } - slotCount := getU32(ap[arOffSlotCount:]) - slotSize := getU32(ap[arOffSlotSize:]) - bmOff := getU32(ap[arOffBitmapOff:]) - slotsOff := getU32(ap[arOffSlotsOff:]) - - if slotCount == 0 || slotSize <= slotHeaderSize { - return nil, fmt.Errorf("arena: 非法规格(slots=%d slotSize=%d)", slotCount, slotSize) + capacity := getU32(ap[arOffCapacity:]) + blockBase := getU32(ap[arOffBlockBase:]) + if capacity != size { + return nil, fmt.Errorf("arena: capacity 不匹配(header=%d mapped=%d)", capacity, size) } - if uint64(slotsOff)+uint64(slotCount)*uint64(slotSize) > uint64(size) { - return nil, fmt.Errorf("arena: 槽数组越出段(slotsOff=%d slots=%d slotSize=%d size=%d)", - slotsOff, slotCount, slotSize, size) + if blockBase+blockHeaderSize > size { + return nil, fmt.Errorf("arena: blockBase 越界(%d,size=%d)", blockBase, size) } - if uint64(bmOff)+uint64(((slotCount+31)/32)*4) > uint64(size) { - return nil, fmt.Errorf("arena: 位图越出段(bmOff=%d slots=%d)", bmOff, slotCount) + return &arenaRegion{region: region, base: base, cap: size}, nil +} + +// MaxPayload 返回单次分配可承载的最大 payload 字节数(上界)。 +func (a *arenaRegion) MaxPayload() int { + return int(a.cap) - arenaHeaderSize - blockHeaderSize +} + +// Capacity 返回 arena 总字节数。 +func (a *arenaRegion) Capacity() uint32 { return a.cap } + +// ---- 块头访问(相对 region 的绝对偏移)---- + +func (a *arenaRegion) blockBase() uint32 { + return a.base + getU32(a.region[a.base+arOffBlockBase:]) +} + +func (a *arenaRegion) blockEnd() uint32 { return a.base + a.cap } + +func (a *arenaRegion) blockSize(off uint32) uint32 { + return getU32(a.region[off+blockOffSize:]) +} + +func (a *arenaRegion) setBlockSize(off, v uint32) { + putU32(a.region[off+blockOffSize:], v) +} + +func (a *arenaRegion) blockState(off uint32) uint32 { + return getU32(a.region[off+blockOffState:]) +} + +func (a *arenaRegion) setBlockState(off, v uint32) { + putU32(a.region[off+blockOffState:], v) +} + +func (a *arenaRegion) blockOwner(off uint32) uint32 { + return getU32(a.region[off+blockOffOwner:]) +} + +func (a *arenaRegion) setBlockOwner(off, v uint32) { + putU32(a.region[off+blockOffOwner:], v) +} + +func (a *arenaRegion) blockPrevSize(off uint32) uint32 { + return getU32(a.region[off+blockOffPrevSize:]) +} + +func (a *arenaRegion) setBlockPrevSize(off, v uint32) { + putU32(a.region[off+blockOffPrevSize:], v) +} + +// blockDataOff 返回块数据区起点(SharedRef.Offset 即此值)。 +func (a *arenaRegion) blockDataOff(off uint32) uint32 { return off + blockHeaderSize } + +// nextBlock 返回下一块偏移;到末尾返回 blockEnd()。 +func (a *arenaRegion) nextBlock(off uint32) uint32 { + sz := a.blockSize(off) + if sz < blockHeaderSize { + return a.blockEnd() // 块链损坏:交给 walk 的步数上限兜底 } - - return &arenaRegion{ - region: region, - base: base, - cap: size, - slots: slotCount, - slotSize: slotSize, - bmOff: bmOff, - slotsOff: slotsOff, - }, nil + next := off + sz + if next > a.blockEnd() { + return a.blockEnd() + } + return next } -// SlotCount 返回槽总数。 -func (a *arenaRegion) SlotCount() uint32 { return a.slots } +// DataOffsetOf 返回块偏移对应的数据起点(诊断/测试用)。 +func (a *arenaRegion) DataOffsetOf(off uint32) uint32 { return a.blockDataOff(off) } -// SlotPayloadCap 返回单个槽可承载的最大 payload 字节数。 -func (a *arenaRegion) SlotPayloadCap() int { return int(a.slotSize) - slotHeaderSize } - -func (a *arenaRegion) bitmapWord(i uint32) *uint32 { - off := a.base + a.bmOff + i*4 - return (*uint32)(unsafe.Pointer(&a.region[off])) -} - -func (a *arenaRegion) slotBase(slot uint32) uint32 { - return a.base + a.slotsOff + slot*a.slotSize -} - -// dataOffset 返回槽数据区的绝对偏移(SharedRef.Offset 即此值)。 -func (a *arenaRegion) dataOffset(slot uint32) uint32 { - return a.slotBase(slot) + slotHeaderSize -} - -// Alloc 预留一个槽并写入 owner,返回共享引用(Flags 即槽号,Length 为 n)。 +// Alloc 分配一块至少 n 字节的块,owner 写入块头。 // -// 必须由内核调用:分配器由内核独占(见文件头所有权模型)。 -func (a *arenaRegion) Alloc(owner uint32, n int) (SharedRef, error) { - if n < 0 || n > a.SlotPayloadCap() { - return SharedRef{}, fmt.Errorf("arena: payload %d 超出槽容量 %d", n, a.SlotPayloadCap()) +// first-fit:从首个块开始找第一个容量足够的空闲块。分配器由内核独占, +// 因此这里的线性扫描不需要任何跨进程同步。 +func (a *arenaRegion) Alloc(owner uint32, n int, gen uint64) (SharedRef, error) { + if n < 0 { + return SharedRef{}, fmt.Errorf("arena: 非法长度 %d", n) } + need := align8(uint32(n) + blockHeaderSize) + if need > a.cap-blockHeaderSize { + return SharedRef{}, fmt.Errorf("arena: 请求 %d 字节超出 arena 容量 %d", n, a.MaxPayload()) + } + a.mu.Lock() defer a.mu.Unlock() - for slot := uint32(0); slot < a.slots; slot++ { - if a.testAndSetLocked(slot) { - base := a.slotBase(slot) - putU32(a.region[base+slotOffDataLen:], uint32(n)) - putU32(a.region[base+slotOffOwner:], owner) - return SharedRef{ - Offset: a.dataOffset(slot), - Length: uint32(n), - Flags: slot, - }, nil + for off := a.blockBase(); off < a.blockEnd(); { + sz := a.blockSize(off) + if sz < blockHeaderSize { + return SharedRef{}, fmt.Errorf("arena: 块链损坏(off=%d size=%d)", off, sz) } + if a.blockState(off) == blockFree && sz >= need { + a.carveLocked(off, need) + a.setBlockState(off, blockUsed) + a.setBlockOwner(off, owner) + a.accountUsed(int32(align8(uint32(n)))) + return a.refOf(off, n, gen), nil + } + off = a.nextBlock(off) } - return SharedRef{}, fmt.Errorf("arena: 槽池已满(%d 槽全部占用)", a.slots) + return SharedRef{}, fmt.Errorf("arena: 空间不足(请求 %d 字节,容量 %d)", n, a.MaxPayload()) } -// Put 分配并写入 payload,返回共享引用(内核方向发送用)。 +// carveLocked 把 off 处的空闲块劈成「已用 need + 剩余空闲」,并把剩余块 +// 的头与后继块的 prevSize 维护好。 +func (a *arenaRegion) carveLocked(off, need uint32) { + sz := a.blockSize(off) + if sz-need < minBlockSize { + return // 剩余太小,整块给出去,不拆 + } + rest := off + need + restSize := sz - need + a.setBlockSize(rest, restSize) + a.setBlockState(rest, blockFree) + a.setBlockOwner(rest, OwnerHost) + a.setBlockPrevSize(rest, need) + a.setBlockSize(off, need) + + // rest 的后继块现在以 rest 为前驱 + if nxt := rest + restSize; nxt < a.blockEnd() { + a.setBlockPrevSize(nxt, restSize) + } +} + +// refOf 由块偏移构造引用。 +func (a *arenaRegion) refOf(off uint32, n int, gen uint64) SharedRef { + return SharedRef{ + Offset: a.blockDataOff(off), + Length: uint32(n), + Generation: uint32(gen), + } +} + +// Put 分配并写入 payload。 func (a *arenaRegion) Put(owner uint32, payload []byte, gen uint64) (SharedRef, error) { - ref, err := a.Alloc(owner, len(payload)) + ref, err := a.Alloc(owner, len(payload), gen) if err != nil { return SharedRef{}, err } copy(a.region[ref.Offset:ref.Offset+ref.Length], payload) - ref.Generation = uint32(gen) return ref, nil } -// MakeRef 为已分配槽构造引用(Length 为 n)。 -func (a *arenaRegion) MakeRef(slot uint32, n int, gen uint64) SharedRef { - return SharedRef{ - Offset: a.dataOffset(slot), - Length: uint32(n), - Generation: uint32(gen), - Flags: slot, - } -} - -// SlotOf 从 SharedRef 解出槽号,并校验它与 Offset 自洽。 -func (a *arenaRegion) SlotOf(ref SharedRef) (uint32, bool) { - slot := ref.Flags - if slot >= a.slots { - return 0, false - } - if uint32(ref.Offset) != a.dataOffset(slot) { - return 0, false - } - if int(ref.Length) > a.SlotPayloadCap() { - return 0, false - } - return slot, true -} - -// Read 按 SharedRef 读取槽数据。 +// findBlockByData 在块链上找到数据起点等于 dataOff 的块。 // -// 校验四件事,任一不满足都返回 error(而不是像早期版本那样返回 nil, -// 让调用方分不清"空数据"和"非法引用"): -// 1. generation 与当前区域一致(remap 后旧引用失效) -// 2. 槽号在范围内,且 Offset 与槽数据区自洽 -// 3. Length 不超过槽容量 -// 4. 槽当前处于占用状态(已被释放的槽不可再读) +// 这是对**插件提交的不可信 offset** 的校验入口:只有真的走完块链确认 +// 该 offset 是一个块的数据起点,才允许后续改动分配器状态。 +func (a *arenaRegion) findBlockByData(dataOff uint32) (uint32, bool) { + steps := 0 + limit := int(a.cap/minBlockSize) + 4 + for off := a.blockBase(); off < a.blockEnd(); { + if steps++; steps > limit { + return 0, false // 链损坏,防死循环 + } + sz := a.blockSize(off) + if sz < blockHeaderSize || off+sz > a.blockEnd() { + return 0, false + } + if a.blockDataOff(off) == dataOff { + return off, true + } + off += sz + } + return 0, false +} + +// findBlockContaining 找到数据区包含 off 的块。 +// +// 与 findBlockByData 的区别:后者要求 off 恰好是块数据起点(用于 Free), +// 前者允许 off 落在块数据区的任意位置(用于 Read:调用帧内的结果区就 +// 位于帧块中间,而不是块起点)。 +func (a *arenaRegion) findBlockContaining(off uint32) (uint32, bool) { + steps := 0 + limit := int(a.cap/minBlockSize) + 4 + for bo := a.blockBase(); bo < a.blockEnd(); { + if steps++; steps > limit { + return 0, false + } + sz := a.blockSize(bo) + if sz < blockHeaderSize || bo+sz > a.blockEnd() { + return 0, false + } + if off >= a.blockDataOff(bo) && off < bo+sz { + return bo, true + } + bo += sz + } + return 0, false +} + +// Read 按 SharedRef 读取数据。 +// +// 校验:generation 与当前区域一致、offset 落在某个**已分配**块的数据区内、 +// 引用不跨越块边界。任一不满足都返回 error(而不是 nil——早期版本返回 +// nil 让调用方分不清“空数据”和“非法引用”)。 +// +// 允许 offset 不必是块起点:调用帧的结果区就在帧块内部。 func (a *arenaRegion) Read(ref SharedRef, gen uint64) ([]byte, error) { if ref.IsZero() { return nil, nil @@ -287,125 +384,155 @@ func (a *arenaRegion) Read(ref SharedRef, gen uint64) ([]byte, error) { if ref.Generation != uint32(gen) { return nil, fmt.Errorf("arena: 引用 generation 过期(ref=%d now=%d)", ref.Generation, gen) } - slot, ok := a.SlotOf(ref) - if !ok { - return nil, fmt.Errorf("arena: 非法引用(offset=%d len=%d slot=%d)", ref.Offset, ref.Length, ref.Flags) + if uint64(ref.Offset)+uint64(ref.Length) > uint64(len(a.region)) { + return nil, fmt.Errorf("arena: 引用越界(offset=%d len=%d total=%d)", + ref.Offset, ref.Length, len(a.region)) } - if !a.IsBusy(slot) { - return nil, fmt.Errorf("arena: 槽 %d 已被释放,引用失效", slot) - } - start := a.dataOffset(slot) - return a.region[start : start+ref.Length], nil -} -// Free 释放 owner 名下的槽。 -// -// 校验 owner 是关键安全边界:一个插件不能释放另一个插件(或内核)的槽。 -func (a *arenaRegion) Free(owner uint32, slot uint32) error { a.mu.Lock() defer a.mu.Unlock() - if slot >= a.slots { - return fmt.Errorf("arena: 槽号越界 %d", slot) + bo, ok := a.findBlockContaining(ref.Offset) + if !ok { + return nil, fmt.Errorf("arena: 非法引用(offset=%d 不在任何块的数据区)", ref.Offset) } - if !a.isBusyLocked(slot) { - return fmt.Errorf("arena: 槽 %d 未分配(重复释放?)", slot) + if a.blockState(bo) != blockUsed { + return nil, fmt.Errorf("arena: 块(offset=%d)已释放,引用失效", ref.Offset) } - base := a.slotBase(slot) - got := getU32(a.region[base+slotOffOwner:]) - if got != owner { - return fmt.Errorf("arena: 槽 %d 不属于调用者(owner=%d caller=%d)", slot, got, owner) + if uint64(ref.Offset)+uint64(ref.Length) > uint64(bo+a.blockSize(bo)) { + return nil, fmt.Errorf("arena: 引用跨越块边界(offset=%d len=%d blockEnd=%d)", + ref.Offset, ref.Length, bo+a.blockSize(bo)) } - a.clearLocked(slot) + return a.region[ref.Offset : ref.Offset+ref.Length], nil +} + +// Free 归还 owner 名下的块,并与相邻空闲块合并。 +// +// 会走块链校验 offset 与 owner:伪造引用不能改动分配器状态。 +func (a *arenaRegion) Free(owner uint32, ref SharedRef) error { + if ref.IsZero() { + return nil + } + a.mu.Lock() + defer a.mu.Unlock() + + off, ok := a.findBlockByData(ref.Offset) + if !ok { + return fmt.Errorf("arena: 非法引用(offset=%d 不是块数据起点)", ref.Offset) + } + if a.blockState(off) != blockUsed { + return fmt.Errorf("arena: 块(offset=%d)未分配(重复归还?)", ref.Offset) + } + if got := a.blockOwner(off); got != owner { + return fmt.Errorf("arena: 块(offset=%d)不属于调用者(owner=%d caller=%d)", ref.Offset, got, owner) + } + + a.setBlockState(off, blockFree) + a.setBlockOwner(off, OwnerHost) + a.accountUsed(-int32(align8(ref.Length))) + a.coalesceLocked(off) return nil } -// ReclaimOwner 释放 owner 名下所有槽,返回回收数量。 +// coalesceLocked 向后、向前合并相邻空闲块,并修正后继块的 prevSize。 +func (a *arenaRegion) coalesceLocked(off uint32) { + // 向后合并 + for { + nxt := a.nextBlock(off) + if nxt >= a.blockEnd() || a.blockState(nxt) != blockFree { + break + } + a.setBlockSize(off, a.blockSize(off)+a.blockSize(nxt)) + } + // 向前合并 + if ps := a.blockPrevSize(off); ps > 0 && off > a.blockBase() { + prev := off - ps + if prev >= a.blockBase() && a.blockState(prev) == blockFree { + a.setBlockSize(prev, a.blockSize(prev)+a.blockSize(off)) + off = prev + } + } + // 后继块的 prevSize 现在应等于合并后本块的大小 + if nxt := a.nextBlock(off); nxt < a.blockEnd() { + a.setBlockPrevSize(nxt, a.blockSize(off)) + } +} + +// ReclaimOwner 归还 owner 名下所有块,返回回收块数。 // -// 用于插件进程退出:崩溃的插件无法归还自己申请的槽,若不管会把池慢慢 -// 耗尽,最终让所有走共享内存的调用退化成内联 RPC。 +// 用于插件进程退出:崩溃的插件无法归还自己申请的块,若不管会把 arena +// 慢慢耗尽,最终让所有走共享内存的调用退化成内联 RPC。 func (a *arenaRegion) ReclaimOwner(owner uint32) int { if owner == OwnerHost { - return 0 // 内核自己的槽由正常路径释放,不在此回收 + return 0 // 内核自己的块由正常路径归还,不在此回收 } a.mu.Lock() defer a.mu.Unlock() - reclaimed := 0 - for slot := uint32(0); slot < a.slots; slot++ { - if !a.isBusyLocked(slot) { - continue + // 先收集偏移再逐个归还:归还过程中会做合并,边遍历边改块链不可靠。 + var targets []uint32 + var reclaimedBytes uint32 + limit := int(a.cap/minBlockSize) + 4 + steps := 0 + for off := a.blockBase(); off < a.blockEnd(); { + if steps++; steps > limit { + break } - base := a.slotBase(slot) - if getU32(a.region[base+slotOffOwner:]) != owner { - continue + sz := a.blockSize(off) + if sz < blockHeaderSize || off+sz > a.blockEnd() { + break } - a.clearLocked(slot) - reclaimed++ + if a.blockState(off) == blockUsed && a.blockOwner(off) == owner { + targets = append(targets, off) + } + off += sz } - return reclaimed + + for _, off := range targets { + if a.blockState(off) != blockUsed { + continue // 已被前面的合并吸收 + } + if a.blockOwner(off) != owner { + continue + } + reclaimedBytes += a.blockSize(off) - blockHeaderSize + a.setBlockState(off, blockFree) + a.setBlockOwner(off, OwnerHost) + a.coalesceLocked(off) + } + if reclaimedBytes > 0 { + a.accountUsed(-int32(align8(reclaimedBytes))) + } + return len(targets) } -// Stats 返回 (已用槽数, 总槽数),供诊断与测试断言。 +// Stats 返回 (已用数据字节, 总字节)。已用字节按 8 字节对齐记账。 func (a *arenaRegion) Stats() (used, total uint32) { - for slot := uint32(0); slot < a.slots; slot++ { - if a.IsBusy(slot) { - used++ - } + return getU32(a.region[a.base+arOffBlockUsed:]), a.cap +} + +// accountUsed 累加/扣减已用字节(相对量,正数表示分配)。 +func (a *arenaRegion) accountUsed(delta int32) { + off := a.base + arOffBlockUsed + cur := int32(getU32(a.region[off:])) + next := cur + delta + if next < 0 { + next = 0 } - return used, a.slots + putU32(a.region[off:], uint32(next)) } -// OwnerOf 返回槽当前的 owner(诊断用,槽空闲时返回 OwnerHost)。 -func (a *arenaRegion) OwnerOf(slot uint32) uint32 { - if slot >= a.slots { - return OwnerHost - } - return getU32(a.region[a.slotBase(slot)+slotOffOwner:]) -} - -// IsBusy 报告槽当前是否被占用。 -func (a *arenaRegion) IsBusy(slot uint32) bool { - if slot >= a.slots { - return false - } - return atomic.LoadUint32(a.bitmapWord(slot/32))&(1<<(slot%32)) != 0 -} - -// ---- 位图内部操作(调用方必须持有 a.mu)---- - -func (a *arenaRegion) isBusyLocked(slot uint32) bool { - return atomic.LoadUint32(a.bitmapWord(slot/32))&(1<<(slot%32)) != 0 -} - -// testAndSetLocked 尝试占用 slot;已占用返回 false。 -func (a *arenaRegion) testAndSetLocked(slot uint32) bool { - w := a.bitmapWord(slot / 32) - bit := uint32(1) << (slot % 32) - for { - cur := atomic.LoadUint32(w) - if cur&bit != 0 { - return false - } - if atomic.CompareAndSwapUint32(w, cur, cur|bit) { - return true - } - } -} - -func (a *arenaRegion) clearLocked(slot uint32) { - w := a.bitmapWord(slot / 32) - bit := uint32(1) << (slot % 32) - for { - cur := atomic.LoadUint32(w) - if cur&bit == 0 { - return - } - if atomic.CompareAndSwapUint32(w, cur, cur&^bit) { - base := a.slotBase(slot) - putU32(a.region[base+slotOffDataLen:], 0) - putU32(a.region[base+slotOffOwner:], 0) - return +// BlockCount 返回块链上的块数(诊断/测试用)。 +func (a *arenaRegion) BlockCount() int { + a.mu.Lock() + defer a.mu.Unlock() + n, limit := 0, int(a.cap/minBlockSize)+4 + for off := a.blockBase(); off < a.blockEnd() && n < limit; n++ { + sz := a.blockSize(off) + if sz < blockHeaderSize { + break } + off += sz } + return n } diff --git a/internal/plugin/proc/arena_test.go b/internal/plugin/proc/arena_test.go index f6d1c98..0827006 100644 --- a/internal/plugin/proc/arena_test.go +++ b/internal/plugin/proc/arena_test.go @@ -6,13 +6,12 @@ import ( "testing" ) -// newTestArena 构造一个独立槽池(不经 Host),保持单测快速。 -func newTestArena(t *testing.T, slots, slotSize uint32) *arenaRegion { +// newTestArena 构造一个独立块分配器(不经 Host),保持单测快速。 +func newTestArena(t *testing.T, size uint32) *arenaRegion { t.Helper() - const base = 64 - need := arenaRequiredSize(slots, slotSize) - region := make([]byte, base+need) - a, err := initArena(region, base, need, slots, slotSize) + const base = 64 // 8 字节对齐 + region := make([]byte, base+size) + a, err := initArena(region, base, size) if err != nil { t.Fatalf("initArena: %v", err) } @@ -20,17 +19,17 @@ func newTestArena(t *testing.T, slots, slotSize uint32) *arenaRegion { } func TestArena_AllocFreeRoundTrip(t *testing.T) { - a := newTestArena(t, 4, 64) - if got, total := a.Stats(); got != 0 || total != 4 { - t.Fatalf("初始 Stats: got (%d,%d), want (0,4)", got, total) + a := newTestArena(t, 8*1024) + if used, total := a.Stats(); used != 0 || total != 8*1024 { + t.Fatalf("初始 Stats: got (%d,%d), want (0,%d)", used, total, 8*1024) } ref, err := a.Put(OwnerHost, []byte("hello"), 0) if err != nil { t.Fatalf("Put: %v", err) } - if used, _ := a.Stats(); used != 1 { - t.Fatalf("Put 后 used=%d, want 1", used) + if used, _ := a.Stats(); used == 0 { + t.Fatal("Put 后 used 应大于 0") } got, err := a.Read(ref, 0) if err != nil { @@ -40,29 +39,29 @@ func TestArena_AllocFreeRoundTrip(t *testing.T) { t.Fatalf("Read=%q, want hello", got) } - if err := a.Free(OwnerHost, ref.Flags); err != nil { + if err := a.Free(OwnerHost, ref); err != nil { t.Fatalf("Free: %v", err) } if used, _ := a.Stats(); used != 0 { t.Fatalf("Free 后 used=%d, want 0", used) } - // 归还后再读必须失败(槽已回收) + // 归还后再读必须失败(块已回收) if _, err := a.Read(ref, 0); err == nil { - t.Fatal("已释放槽的引用应读取失败") + t.Fatal("已释放块的引用应读取失败") } } func TestArena_ConcurrentAllocUnique(t *testing.T) { - const slots = 32 - a := newTestArena(t, slots, 64) + const workers = 64 + a := newTestArena(t, 256*1024) - refs := make(chan SharedRef, slots) + refs := make(chan SharedRef, workers) var wg sync.WaitGroup - for i := 0; i < slots; i++ { + for i := 0; i < workers; i++ { wg.Add(1) go func() { defer wg.Done() - ref, err := a.Alloc(OwnerHost, 8) + ref, err := a.Alloc(OwnerHost, 1024, 0) if err != nil { t.Errorf("Alloc: %v", err) return @@ -73,105 +72,121 @@ func TestArena_ConcurrentAllocUnique(t *testing.T) { wg.Wait() close(refs) - seen := make(map[uint32]bool, slots) + seen := make(map[uint32]bool, workers) for ref := range refs { - if seen[ref.Flags] { - t.Fatalf("并发分配拿到重复槽 %d", ref.Flags) + if seen[ref.Offset] { + t.Fatalf("并发分配拿到重复 offset=%d", ref.Offset) } - seen[ref.Flags] = true + seen[ref.Offset] = true } - if len(seen) != slots { - t.Fatalf("唯一槽数=%d, want %d", len(seen), slots) - } - if used, _ := a.Stats(); used != slots { - t.Fatalf("used=%d, want %d", used, slots) + if len(seen) != workers { + t.Fatalf("唯一块数=%d, want %d", len(seen), workers) } } func TestArena_ExhaustionReturnsError(t *testing.T) { - a := newTestArena(t, 2, 64) - for i := 0; i < 2; i++ { - if _, err := a.Alloc(OwnerHost, 8); err != nil { - t.Fatalf("第 %d 次 Alloc: %v", i, err) - } + a := newTestArena(t, 4*1024) + // 第一次分配吃掉几乎整块 + if _, err := a.Alloc(OwnerHost, 3*1024, 0); err != nil { + t.Fatalf("首次 Alloc: %v", err) } - if _, err := a.Alloc(OwnerHost, 8); err == nil { - t.Fatal("槽池耗尽时 Alloc 应返回错误") + // 再申请一大块必然失败 + if _, err := a.Alloc(OwnerHost, 3*1024, 0); err == nil { + t.Fatal("空间不足时 Alloc 应返回错误") } } func TestArena_RejectsOversizePayload(t *testing.T) { - a := newTestArena(t, 2, 64) - if _, err := a.Alloc(OwnerHost, a.SlotPayloadCap()+1); err == nil { - t.Fatal("超出槽容量的 payload 应被拒绝") + a := newTestArena(t, 4*1024) + if _, err := a.Alloc(OwnerHost, a.MaxPayload()+1, 0); err == nil { + t.Fatal("超出 arena 容量的请求应被拒绝") } } func TestArena_FreeEnforcesOwnership(t *testing.T) { - a := newTestArena(t, 4, 64) + a := newTestArena(t, 8*1024) const ownerA, ownerB = 7, 9 - refA, err := a.Alloc(ownerA, 8) + refA, err := a.Alloc(ownerA, 128, 0) if err != nil { t.Fatal(err) } - // 另一个插件不能释放 A 的槽 - if err := a.Free(ownerB, refA.Flags); err == nil { - t.Fatal("跨 owner 释放应被拒绝") + // 另一个插件不能释放 A 的块 + if err := a.Free(ownerB, refA); err == nil { + t.Fatal("跨 owner 归还应被拒绝") } - if used, _ := a.Stats(); used != 1 { - t.Fatalf("拒绝释放后 used=%d, want 1", used) + if used, _ := a.Stats(); used == 0 { + t.Fatal("拒绝归还后块应仍然占用") } - if a.OwnerOf(refA.Flags) != ownerA { - t.Fatalf("槽 %d 的 owner 应为 %d", refA.Flags, ownerA) - } - // 正确的 owner 可以释放 - if err := a.Free(ownerA, refA.Flags); err != nil { - t.Fatalf("同 owner 释放: %v", err) + // 正确的 owner 可以归还 + if err := a.Free(ownerA, refA); err != nil { + t.Fatalf("同 owner 归还: %v", err) } } -func TestArena_ReclaimOwnerReleasesOnlyItsSlots(t *testing.T) { - a := newTestArena(t, 8, 64) +func TestArena_FreeRejectsForgedOffset(t *testing.T) { + a := newTestArena(t, 8*1024) + ref, err := a.Alloc(OwnerHost, 256, 0) + if err != nil { + t.Fatal(err) + } + + // 未对齐到块数据起点的 offset 必须被拒绝——否则会直接破坏块链。 + forged := ref + forged.Offset += 8 + if err := a.Free(OwnerHost, forged); err == nil { + t.Fatal("伪造 offset 应被拒绝") + } + // 合法引用仍能正常归还(分配器状态未被破坏) + if err := a.Free(OwnerHost, ref); err != nil { + t.Fatalf("合法引用不应受影响: %v", err) + } + if used, _ := a.Stats(); used != 0 { + t.Fatalf("归还后 used=%d, want 0", used) + } +} + +func TestArena_ReclaimOwnerReleasesOnlyItsBlocks(t *testing.T) { + a := newTestArena(t, 64*1024) const ownerA, ownerB = 7, 9 for i := 0; i < 3; i++ { - if _, err := a.Alloc(ownerA, 8); err != nil { + if _, err := a.Alloc(ownerA, 1024, 0); err != nil { t.Fatal(err) } } for i := 0; i < 2; i++ { - if _, err := a.Alloc(ownerB, 8); err != nil { + if _, err := a.Alloc(ownerB, 1024, 0); err != nil { t.Fatal(err) } } if n := a.ReclaimOwner(ownerA); n != 3 { t.Fatalf("ReclaimOwner(A)=%d, want 3", n) } - if used, _ := a.Stats(); used != 2 { - t.Fatalf("回收 A 后 used=%d, want 2(B 的槽必须保留)", used) + // B 的两块必须保留:仍能读回 + if n := a.ReclaimOwner(ownerB); n != 2 { + t.Fatalf("ReclaimOwner(B)=%d, want 2(A 的回收不能误伤 B)", n) } - // 内核自己的槽不参与回收 - if n := a.ReclaimOwner(OwnerHost); n != 0 { - t.Fatalf("ReclaimOwner(host)=%d, want 0", n) + if used, _ := a.Stats(); used != 0 { + t.Fatalf("全部回收后 used=%d, want 0", used) } } func TestArena_ReclaimHostIsNoop(t *testing.T) { - a := newTestArena(t, 4, 64) - if _, err := a.Alloc(OwnerHost, 8); err != nil { + a := newTestArena(t, 4*1024) + if _, err := a.Alloc(OwnerHost, 128, 0); err != nil { t.Fatal(err) } if n := a.ReclaimOwner(OwnerHost); n != 0 { - t.Fatalf("内核槽不应被 ReclaimOwner 回收,实际回收 %d", n) + t.Fatalf("内核块不应被 ReclaimOwner 回收,实际回收 %d", n) } - if used, _ := a.Stats(); used != 1 { - t.Fatalf("used=%d, want 1", used) + if used, _ := a.Stats(); used == 0 { + t.Fatal("used 应仍大于 0") } } func TestArena_ReadRejectsInvalidRefs(t *testing.T) { - a := newTestArena(t, 4, 256) + a := newTestArena(t, 8*1024) ref, err := a.Put(OwnerHost, []byte("payload"), 0) if err != nil { t.Fatal(err) @@ -184,31 +199,24 @@ func TestArena_ReadRejectsInvalidRefs(t *testing.T) { t.Fatal("generation 不匹配应被拒绝") } - // offset 与槽不自洽 + // offset 不是块数据起点 badOffset := ref - badOffset.Offset++ + badOffset.Offset += 8 if _, err := a.Read(badOffset, 0); err == nil { - t.Fatal("offset 与槽不自洽应被拒绝") + t.Fatal("offset 不是块数据起点应被拒绝") } - // 槽号越界 - badSlot := ref - badSlot.Flags = 999 - if _, err := a.Read(badSlot, 0); err == nil { - t.Fatal("槽号越界应被拒绝") - } - - // Length 超容量 + // Length 超出块容量 tooLong := ref - tooLong.Length = uint32(a.SlotPayloadCap() + 1) + tooLong.Length = 1 << 20 if _, err := a.Read(tooLong, 0); err == nil { t.Fatal("Length 超容量应被拒绝") } } func TestArena_PutReadLargePayload(t *testing.T) { - a := newTestArena(t, 4, 4096) - payload := bytes.Repeat([]byte("abcdefgh"), 400) // 3200 字节 + a := newTestArena(t, 512*1024) + payload := bytes.Repeat([]byte("abcdefgh"), 8192) // 64KB ref, err := a.Put(OwnerHost, payload, 3) if err != nil { t.Fatalf("Put: %v", err) @@ -220,44 +228,80 @@ func TestArena_PutReadLargePayload(t *testing.T) { if !bytes.Equal(got, payload) { t.Fatalf("读回数据不一致:len(got)=%d len(want)=%d", len(got), len(payload)) } + // 大 payload 不能被 16KB 定长槽时代的容量假设卡住 + if len(got) <= 16*1024 { + t.Fatalf("本用例应验证超过旧 16KB 槽容量的 payload,实际 %d", len(got)) + } +} + +func TestArena_FreeCoalescesAdjacentBlocks(t *testing.T) { + a := newTestArena(t, 64*1024) + + refs := make([]SharedRef, 0, 8) + for i := 0; i < 8; i++ { + r, err := a.Alloc(OwnerHost, 4*1024, 0) + if err != nil { + t.Fatalf("第 %d 次 Alloc: %v", i, err) + } + refs = append(refs, r) + } + for _, r := range refs { + if err := a.Free(OwnerHost, r); err != nil { + t.Fatalf("Free: %v", err) + } + } + + // 全部归还并合并后,应能再分配一个接近整块的大块。 + big, err := a.Alloc(OwnerHost, 48*1024, 0) + if err != nil { + t.Fatalf("合并后应能分配大块: %v", err) + } + if err := a.Free(OwnerHost, big); err != nil { + t.Fatal(err) + } + if n := a.BlockCount(); n != 1 { + t.Fatalf("全部归还并合并后应只剩 1 块,实际 %d", n) + } } func TestArena_AttachRejectsCorruptLayout(t *testing.T) { - a := newTestArena(t, 4, 256) + a := newTestArena(t, 4*1024) - // 魔数被改 putU32(a.region[a.base+arOffMagic:], 0xDEADBEEF) if _, err := attachArena(a.region, a.base, a.cap); err == nil { t.Fatal("魔数错误应被拒绝") } putU32(a.region[a.base+arOffMagic:], arenaMagic) - // 版本被改 putU32(a.region[a.base+arOffVersion:], 99) if _, err := attachArena(a.region, a.base, a.cap); err == nil { t.Fatal("版本不匹配应被拒绝") } putU32(a.region[a.base+arOffVersion:], arenaVersion) - // 槽数组越界 - putU32(a.region[a.base+arOffSlotCount:], 1<<20) + putU32(a.region[a.base+arOffCapacity:], 1<<30) if _, err := attachArena(a.region, a.base, a.cap); err == nil { - t.Fatal("槽数组越界应被拒绝") + t.Fatal("capacity 不匹配应被拒绝") } } -func TestArenaRequiredSizeAlignment(t *testing.T) { - // 位图按 8 字节对齐、槽数组 8 字节对齐,保证 4 字节原子操作不跨页/不越界。 - for _, tc := range []struct{ slots, slotSize uint32 }{ - {1, 16}, {8, 64}, {32, 16384}, {33, 1024}, {64, 512}, {100, 256}, - } { - need := arenaRequiredSize(tc.slots, tc.slotSize) - if need%8 != 0 { - t.Errorf("slots=%d slotSize=%d: arenaRequiredSize=%d 不是 8 的倍数", tc.slots, tc.slotSize, need) +func TestArena_InitRejectsTooSmall(t *testing.T) { + const base = 64 + region := make([]byte, base+arenaMinSize()-8) + if _, err := initArena(region, base, arenaMinSize()-8); err == nil { + t.Fatal("过小的段应被拒绝") + } +} + +func TestArena_DataOffsetsAligned(t *testing.T) { + a := newTestArena(t, 64*1024) + for i, size := range []uint32{1, 7, 8, 9, 100, 4096} { + ref, err := a.Alloc(OwnerHost, int(size), 0) + if err != nil { + t.Fatalf("第 %d 次 Alloc(size=%d): %v", i, size, err) } - region := make([]byte, 64+need) - if _, err := initArena(region, 64, need, tc.slots, tc.slotSize); err != nil { - t.Errorf("slots=%d slotSize=%d: initArena 失败: %v", tc.slots, tc.slotSize, err) + if ref.Offset%8 != 0 { + t.Fatalf("size=%d 的数据起点 %d 未 8 字节对齐", size, ref.Offset) } } } diff --git a/internal/plugin/proc/corehandler.go b/internal/plugin/proc/corehandler.go index 3ad36b8..88bb78e 100644 --- a/internal/plugin/proc/corehandler.go +++ b/internal/plugin/proc/corehandler.go @@ -647,88 +647,70 @@ func (h *coreHandler) cleanerProxy(scope, name string, enabled bool) (func(strin }, nil } -// invokeCleanerText 完成一次 Cleaner 往返:payload 走共享槽,超限时整条链路退回内联。 +// invokeCleanerText 完成一次 Cleaner 往返,使用与工具调用相同的 funccall 帧模型。 // -// 槽的申请与归还全部由内核负责(**谁分配谁释放**): -// -// reqSlot ──写入 input──▶ 随 RPC 把 TextRef 发给插件 -// respSlot ──预分配────▶ 随 RPC 把 RespRef 发给插件,插件把结果写进来 -// 读取结果后两个槽一起归还 -// -// 插件侧不申请任何槽,因此不存在跨进程分配器的竞争。 +// 内核(caller)标定帧:输入段 + 结果预算段,插件在帧内写结果; +// 只有结果超出预算时插件才向内核申请扩容块(插件只申请,回收由内核做)。 func (h *coreHandler) invokeCleanerText(scope, name, text string) (string, error) { arena := h.host.Arena() gen := h.host.Generation() - params := CleanerInvokeParams{Scope: scope, Name: name} - - // 请求槽:放得下就走共享内存,否则内联。 - // 注意“放不下”不是错误,只是退化成原有正确路径。 - var reqRef SharedRef - if len(text) <= arena.SlotPayloadCap() { - if ref, err := arena.Put(OwnerHost, []byte(text), gen); err == nil { - params.TextRef = ref - reqRef = ref - } - } - if reqRef.IsZero() { - params.Text = text + frame, err := arena.Alloc(OwnerHost, len(text)+cleanerResultBudget, gen) + if err != nil { + return "", fmt.Errorf("%s %s Cleaner 分配调用帧失败: %w", scope, name, err) } + defer func() { _ = arena.Free(OwnerHost, frame) }() - // 响应槽:预分配满容量。插件只写,不申请。 - var respRef SharedRef - if ref, err := arena.Alloc(OwnerHost, arena.SlotPayloadCap()); err == nil { - params.RespRef = ref - respRef = ref + area, err := arena.Read(frame, gen) + if err != nil { + return "", err } + copy(area[:len(text)], text) - defer func() { - if !respRef.IsZero() { - _ = arena.Free(OwnerHost, respRef.Flags) - } - if !reqRef.IsZero() { - _ = arena.Free(OwnerHost, reqRef.Flags) - } - }() + params := CleanerInvokeParams{ + Scope: scope, + Name: name, + Frame: frame, + InputLen: uint32(len(text)), + } res, err := h.invokeCleaner(params) if err != nil { return "", err } - // 结果优先取共享槽;插件放不下时会内联返回。 - if !res.TextRef.IsZero() { - data, err := arena.Read(res.TextRef, gen) - if err != nil { - return "", err - } - return string(data), nil + // 插件申请了扩容块:内核负责归还(插件只会申请,回收由内核做)。 + if res.TextRef.IsZero() { + return "", fmt.Errorf("%s %s Cleaner 未返回结果引用", scope, name) } - return res.Text, nil + if res.TextRef.Flags&sharedRefFlagExpand != 0 { + defer func() { _ = arena.Free(OwnerHost, res.TextRef) }() + } + data, err := arena.Read(res.TextRef, gen) + if err != nil { + return "", err + } + return string(data), nil } -// ---- 共享槽池的 syscall 风格接口(插件 RPC)---- +// ---- 共享内存的 syscall 风格接口(插件 RPC)---- // // 共享内存是内部实现,不向插件开发者暴露;模板运行时在传输层调用它们, // 公开 SDK 仍是普通字符串/Map。 -// arenaAlloc 给本插件分配一块共享槽并返回描述符。 +// arenaAlloc 给本插件分配一块共享内存。 +// +// 用途:结果超出内核标定帧的预算时,插件据此申请扩容块。 func (h *coreHandler) arenaAlloc(size uint32) (SharedRef, error) { - if int(size) > h.host.Arena().SlotPayloadCap() { - return SharedRef{}, fmt.Errorf("arena.alloc: 请求 %d 字节超出槽容量 %d", - size, h.host.Arena().SlotPayloadCap()) - } - return h.host.Arena().Alloc(h.owner, int(size)) + return h.host.Arena().Alloc(h.owner, int(size), h.host.Generation()) } -// arenaFree 归还本插件申请的槽。内核校验归属,防止释放他人的槽。 +// arenaFree 归还本插件申请的块。 +// +// 内核会走块链校验 offset 确实是某个已分配块的数据起点、owner 匹配, +// 伪造引用不能改动分配器状态。 func (h *coreHandler) arenaFree(ref SharedRef) error { - slot, ok := h.host.Arena().SlotOf(ref) - if !ok { - return fmt.Errorf("arena.free: 非法引用(offset=%d len=%d slot=%d)", - ref.Offset, ref.Length, ref.Flags) - } - return h.host.Arena().Free(h.owner, slot) + return h.host.Arena().Free(h.owner, ref) } // toolRegister 注册插件工具,handler 反向调用插件执行(原 case 1)。 diff --git a/internal/plugin/proc/host.go b/internal/plugin/proc/host.go index 6c038db..2ccae7a 100644 --- a/internal/plugin/proc/host.go +++ b/internal/plugin/proc/host.go @@ -30,8 +30,11 @@ type Host struct { seg *Segment // StageContext segment(位于 unified ctxData) shmSize int // 统一区域总大小 - // arena 是跨进程共享槽池(§13.2 重设计)。分配/释放走共享位图 CAS, - // 任意进程/goroutine 并发调用都安全,不需要任何锁。 + // arena 是跨进程共享内存分配器(§13.2 重设计),由内核独占管理: + // 插件通过 RPC 申请/归还,不做任何分配决策。 + // + // 分配与回收都在内核进程内进行,一把 mu 即可保证安全, + // 不存在跨进程分配器那种“共享游标被两个进程各自更新”的竞态。 arena *arenaRegion // nextOwner 给每个插件进程分配一个不透明 owner ID,供 ReclaimOwner 使用。 @@ -77,15 +80,15 @@ func NewHost() (*Host, error) { return nil, err } - // 初始化 Exchange Arena 槽池 - arena, err := initArena(data, ur.arenaOff, ur.arenaCap, arenaDefaultSlotCount, arenaDefaultSlotSize) + // 初始化 Exchange Arena(内核独占的变长块分配器) + arena, err := initArena(data, ur.arenaOff, ur.arenaCap) if err != nil { freeShm(memfd, data) - return nil, fmt.Errorf("槽池初始化: %w", err) + return nil, fmt.Errorf("共享内存分配器初始化: %w", err) } - if used, total := arena.Stats(); used != 0 || total != arenaDefaultSlotCount { + if used, total := arena.Stats(); used != 0 || total != ur.arenaCap { freeShm(memfd, data) - return nil, fmt.Errorf("槽池初始化异常:used=%d total=%d", used, total) + return nil, fmt.Errorf("分配器初始化异常:used=%d total=%d", used, total) } // 创建 StageContext segment(位于 SuperBlock 之后) diff --git a/internal/plugin/proc/plugin.go b/internal/plugin/proc/plugin.go index e9921c4..a79531f 100644 --- a/internal/plugin/proc/plugin.go +++ b/internal/plugin/proc/plugin.go @@ -221,11 +221,46 @@ func (p *Plugin) handleExit(name string, err error) { // ---- 内核 → 插件的反向调用 ---- +// invokeTool 在插件进程内执行工具。 +// +// 按 **funccall 模型**:内核是 caller,为每次调用**标定一块内存帧** +// (参数段 + 结果预算段)交给插件(callee)。参数永远写在帧里,不再有 +// “小 payload 走内联”的按大小分支。 +// +// 结果超出预算时插件才向内核申请扩容块,并在引用上打 +// sharedRefFlagExpand,内核据此单独归还。 +// +// 控制面仍是 RPC(请求 ID 关联、ctx 取消、崩溃唤醒都由 Process 承载)。 func (p *Plugin) invokeTool(name string, args map[string]interface{}) (interface{}, error) { if p.proc == nil { return nil, ErrProcessExited } - raw, err := p.proc.Call(MethodToolInvoke, ToolInvokeParams{Name: name, Args: args}) + arena := p.host.Arena() + gen := p.host.Generation() + + argJSON, err := json.Marshal(args) + if err != nil { + return nil, fmt.Errorf("proc: %s 工具 %s 参数序列化失败: %w", p.name, name, err) + } + + // 内核标定调用帧:参数段 + 结果预算段。 + frame, err := arena.Alloc(OwnerHost, len(argJSON)+toolResultBudget, gen) + if err != nil { + return nil, fmt.Errorf("proc: %s 工具 %s 分配调用帧失败: %w", p.name, name, 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, name, err) + } + copy(area[:len(argJSON)], argJSON) + + raw, err := p.proc.Call(MethodToolInvoke, ToolInvokeParams{ + Name: name, + Frame: frame, + ArgsLen: uint32(len(argJSON)), + }) if err != nil { return nil, err } @@ -233,7 +268,22 @@ func (p *Plugin) invokeTool(name string, args map[string]interface{}) (interface if err := json.Unmarshal(raw, &res); err != nil { return nil, fmt.Errorf("proc: %s 工具 %s 应答解析失败: %w", p.name, name, err) } - return res.Result, nil + if res.ResultRef.IsZero() { + return nil, fmt.Errorf("proc: %s 工具 %s 未返回结果引用", p.name, name) + } + // 插件申请了扩容块:内核负责归还。 + if res.ResultRef.Flags&sharedRefFlagExpand != 0 { + defer func() { _ = arena.Free(OwnerHost, res.ResultRef) }() + } + data, err := arena.Read(res.ResultRef, gen) + if err != nil { + return nil, fmt.Errorf("proc: %s 工具 %s 读取共享结果失败: %w", p.name, name, err) + } + var out interface{} + if err := json.Unmarshal(data, &out); err != nil { + return nil, fmt.Errorf("proc: %s 工具 %s 解析共享结果失败: %w", p.name, name, err) + } + return out, nil } // invokeCleaner 在插件进程内执行工具或通道注册时提供的 Cleaner 函数。 diff --git a/internal/plugin/proc/plugin_test.go b/internal/plugin/proc/plugin_test.go index a7390fb..941de0a 100644 --- a/internal/plugin/proc/plugin_test.go +++ b/internal/plugin/proc/plugin_test.go @@ -500,3 +500,70 @@ func TestPlugin_ArenaAllocFreeAcrossProcess(t *testing.T) { t.Fatalf("插件归还后槽池应全空,实际 used=%d/%d", used, total) } } + +// 工具调用的参数/结果走共享槽(§13.3)。 +// +// 控制面仍是 RPC(请求 ID 关联、ctx 取消、崩溃唤醒都由它承载), +// 只有 payload 走共享内存: +// - 参数超过阈值时内核写入槽,把 ArgsRef 发给插件 +// - 结果放得下时插件写入内核预分配的响应槽,回 ResultRef +// - 超限/池满时退回内联 JSON,不能影响功能 +// +// demo_big 会在结果里报出它到底从哪里读到参数,因此本测试验证的是 +// “真的走了共享内存”,而不只是“返回值对”。 +func TestPlugin_ToolInvokeArgsResultViaArena(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.tools["demo_big"] + core.mu.Unlock() + if !ok { + t.Fatal("插件应注册 demo_big 工具") + } + + t.Run("大 payload 走共享槽", func(t *testing.T) { + // 4 字节 × 3 × 1000 = 12000 字节,明显超过内联阈值且能放进槽 + big := strings.Repeat("共享内存", 1000) + res, err := h(map[string]interface{}{"text": big}) + if err != nil { + t.Fatalf("调用 demo_big: %v", err) + } + s, _ := res.(string) + if !strings.HasPrefix(s, "shared:") { + t.Fatalf("大参数应经共享槽传递,实际结果前缀不对(len=%d, head=%.40q)", len(s), s) + } + if want := "shared:" + strings.ToUpper(big); s != want { + t.Fatalf("经共享槽往返的内容不一致:got len=%d want len=%d", len(s), len(want)) + } + if used, total := host.Arena().Stats(); used != 0 { + t.Fatalf("调用结束后槽池应全空,实际 used=%d/%d", used, total) + } + }) + + t.Run("小 payload 同样走调用帧", func(t *testing.T) { + // 设计上不再有“小 payload 走内联”的按大小分支:内核总是标定调用帧。 + res, err := h(map[string]interface{}{"text": "abc"}) + if err != nil { + t.Fatalf("调用 demo_big: %v", err) + } + if res != "shared:ABC" { + t.Fatalf("小参数也应走内核标定的调用帧,实际 %v", res) + } + if used, total := host.Arena().Stats(); used != 0 { + t.Fatalf("调用结束后应全部归还,实际 used=%d/%d", used, total) + } + }) +} diff --git a/internal/plugin/proc/process_test.go b/internal/plugin/proc/process_test.go index d931656..14592e2 100644 --- a/internal/plugin/proc/process_test.go +++ b/internal/plugin/proc/process_test.go @@ -341,6 +341,7 @@ func TestProcess_SpawnRequiresHandler(t *testing.T) { // - 在途调用挂到自己的超时; // - OnExit 不触发 → 崩溃计数、工具摘除、自动重启全都不发生; // - 进程表里插件已是僵尸,注册表里却一切正常。 +// // 生产上 browser 拉 chromium、editdoc 拉 python 正是这个形状。 // 现在由专职 waitLoop 直接 wait4(2) 判定,不再依赖 fd 生命周期。 func TestProcess_ExitDetectedDespiteInheritedStdout(t *testing.T) { diff --git a/internal/plugin/proc/protocol.go b/internal/plugin/proc/protocol.go index 4bf30ef..f6cba77 100644 --- a/internal/plugin/proc/protocol.go +++ b/internal/plugin/proc/protocol.go @@ -209,43 +209,62 @@ type StageInvokeResult struct { Seq uint64 `json:"seq"` } -// ToolInvokeParams / ToolInvokeResult:工具调用(原 go_invoke_tool)。 +// ToolInvokeParams:工具调用(原 go_invoke_tool)。 // -// ArgsRef/ResultRef 为后续 §13.3 预留:payload 走 Exchange Arena, -// RPC 帧只带 16 字节描述符;当前仍以 Args/Result 内联为主。 +// 按 **funccall 模型**,内核(caller)为每次调用标定一块内存帧交给插件 +// (callee),插件在这块内存里工作: +// +// [0, ArgsLen) 参数 JSON +// [ArgsLen, Frame.Length) 结果区(内核预留的预算) +// +// 结果放得下就写在帧内;**不够用时插件才向内核申请扩容**(arena.alloc), +// 并在返回引用上打 sharedRefFlagExpand,内核据此单独归还扩容块。 +// +// 参数永远在共享内存里,不存在“小 payload 走内联”的按大小分支。 +// +// Args 仅剩给**直连 RPC 的调用方**(process/bench 测试不建 Host,拿不到 +// 共享内存);内核的 invokeTool 始终走 Frame。 type ToolInvokeParams struct { Name string `json:"name"` - Args map[string]interface{} `json:"args,omitempty"` - ArgsRef SharedRef `json:"args_ref,omitempty"` + Frame SharedRef `json:"frame,omitempty"` + ArgsLen uint32 `json:"args_len,omitempty"` + Args map[string]interface{} `json:"args,omitempty"` // 仅直连 RPC 调用方使用 } type ToolInvokeResult struct { - Result interface{} `json:"result,omitempty"` ResultRef SharedRef `json:"result_ref,omitempty"` + Result interface{} `json:"result,omitempty"` // 仅直连 RPC 调用方使用 } // CleanerInvokeParams / CleanerInvokeResult:跨进程计算层清洗。 // Scope 取 tool / input / output,Name 是工具名或通道名。 // -// 数据面走 Exchange Arena: -// - TextRef 指向内核写入的输入文本;payload 超过槽容量时改用内联 Text -// - RespRef 指向内核**预分配**的响应槽,插件把结果写进去后回 TextRef; -// 结果放不下时插件改用内联 Text 返回 +// 与工具调用用**同一个 funccall 帧模型**: // -// 插件侧不做任何分配,因此不存在跨进程分配器的竞争。 +// [0, InputLen) 输入文本 +// [InputLen, Frame.Length) 结果区(内核预留的预算) +// +// 结果放不下时插件申请扩容块,并在 TextRef 上打 sharedRefFlagExpand。 type CleanerInvokeParams struct { - Scope string `json:"scope"` - Name string `json:"name"` - Text string `json:"text,omitempty"` - TextRef SharedRef `json:"text_ref,omitempty"` - RespRef SharedRef `json:"resp_ref,omitempty"` + Scope string `json:"scope"` + Name string `json:"name"` + Frame SharedRef `json:"frame,omitempty"` + InputLen uint32 `json:"input_len,omitempty"` } type CleanerInvokeResult struct { - Text string `json:"text,omitempty"` TextRef SharedRef `json:"text_ref,omitempty"` } +// 调用帧的结果预算。 +// +// 内核按“参数长度 + 预算”标定帧;结果超出预算不是失败,插件会申请扩容块。 +// 预算取 64KB:覆盖绝大多数工具结果,使常态调用完全免于第二次分配。 +const ( + toolResultBudget = 64 * 1024 + cleanerResultBudget = 64 * 1024 +) + const ( CleanerScopeTool = "tool" CleanerScopeInput = "input" diff --git a/internal/plugin/proc/testdata/stageplugin.go b/internal/plugin/proc/testdata/stageplugin.go index f96f59d..8960e50 100644 --- a/internal/plugin/proc/testdata/stageplugin.go +++ b/internal/plugin/proc/testdata/stageplugin.go @@ -57,8 +57,58 @@ func arenaSlotSize() uint32 { return binary.LittleEndian.Uint32(region[arenaOff+arOffSlotSize:]) } -// arenaPayloadCap 返回槽可承载的最大 payload。 -func arenaPayloadCap() int { return int(arenaSlotSize()) - slotHeaderSize } +// ---- 调用帧(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 +} // ---- 共享槽池的插件侧接口(内核 RPC,内部实现)---- @@ -75,6 +125,19 @@ type SharedRef struct { 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 // 引用指向插件申请的扩容块 +) + // arenaAlloc 向内核申请一块共享内存,内核返回偏移与大小。 func arenaAlloc(size uint32) (SharedRef, error) { raw := callKernel("arena.alloc", map[string]interface{}{"size": size}) @@ -310,6 +373,11 @@ func main() { "def": map[string]interface{}{"name": "demo_inject", "description": "共享内存注入"}, "has_cleaner": false, }) + callKernel("tool.register", map[string]interface{}{ + "name": "demo_big", + "def": map[string]interface{}{"name": "demo_big", "description": "报告参数量及结果回传路径"}, + "has_cleaner": false, + }) callKernel("stage.register", map[string]interface{}{ "stage": "after_toolcall", "scope": "global", @@ -336,11 +404,29 @@ func main() { // 模板对每个内核请求都 `go handleKernelRequest` 的原因。 go func(id uint64, raw json.RawMessage) { var p struct { - Name string `json:"name"` - Args map[string]interface{} `json:"args"` + Name string `json:"name"` + Args map[string]interface{} `json:"args"` + Frame SharedRef `json:"frame"` + ArgsLen uint32 `json:"args_len"` } json.Unmarshal(raw, &p) - text, _ := p.Args["text"].(string) + + // 参数:内核标定帧的前段。只有直连 RPC 的调用方(无帧)才走 + // 内联 Args——生产路径永远走帧。 + args := p.Args + argsFrom := "inline" + if !p.Frame.IsZero() { + argsFrom = "shared" + 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 + } + } + text, _ := args["text"].(string) // demo_inject 走插件侧共享内存路径: // 申请 → 写入 → 随业务 RPC 回传 → 归还。 @@ -356,30 +442,51 @@ func main() { "source": "plugin", "channel": "demo", "text_ref": ref, }) arenaFree(ref) - send(response{ID: id, Result: map[string]interface{}{"result": "injected"}}) + + blob, _ := json.Marshal("injected") + outRef, err := frameOutput(p.Frame, p.ArgsLen, blob, true) + if err != nil { + outRef, err = arenaPut(blob, true) + if err != nil { + send(response{ID: id, Error: "结果扩容失败: " + err.Error()}) + return + } + } + send(response{ID: id, Result: map[string]interface{}{"result_ref": outRef}}) return } - send(response{ID: id, Result: map[string]interface{}{ - "result": strings.ToUpper(text), - }}) + var out interface{} + if p.Name == "demo_big" { + // 把参数来路编进结果,让测试能验证真的走了共享内存而非只看返回值。 + out = argsFrom + ":" + strings.ToUpper(text) + } else { + out = strings.ToUpper(text) + } + + blob, err := json.Marshal(out) + if err != nil { + send(response{ID: id, Error: "序列化结果失败: " + err.Error()}) + return + } + // 结果优先写进内核标定的帧;放不下才申请扩容块。 + ref, err := frameOutput(p.Frame, p.ArgsLen, blob, true) + if err != nil { + ref, err = arenaPut(blob, true) + if err != nil { + send(response{ID: id, Error: "结果扩容失败: " + err.Error()}) + return + } + } + send(response{ID: id, Result: map[string]interface{}{"result_ref": ref}}) }(req.ID, req.Params) case "cleaner.invoke": var p struct { - Scope string `json:"scope"` - Name string `json:"name"` - Text string `json:"text"` - TextRef struct { - Offset uint32 `json:"offset"` - Length uint32 `json:"length"` - Flags uint32 `json:"flags"` - } `json:"text_ref"` - RespRef struct { - Offset uint32 `json:"offset"` - Length uint32 `json:"length"` - Flags uint32 `json:"flags"` - } `json:"resp_ref"` + Scope string `json:"scope"` + Name string `json:"name"` + Frame SharedRef `json:"frame"` + InputLen uint32 `json:"input_len"` } json.Unmarshal(req.Params, &p) valid := (p.Scope == "tool" && p.Name == "demo_upper") || @@ -390,27 +497,19 @@ func main() { continue } - // 读输入:优先共享槽,否则内联。 - input := p.Text - if p.TextRef.Length > 0 { - input = string(region[p.TextRef.Offset : p.TextRef.Offset+p.TextRef.Length]) - } + input := string(frameInput(p.Frame, p.InputLen)) output := p.Scope + "-cleaned:" + input - // 写结果:有预分配槽且放得下就写槽,否则内联。 - // 插件不分配任何槽(“谁分配谁释放”全部在内核侧)。 - if p.RespRef.Offset > 0 && len(output) <= arenaPayloadCap() { - copy(region[p.RespRef.Offset:], output) - send(response{ID: req.ID, Result: map[string]interface{}{ - "text_ref": map[string]interface{}{ - "offset": p.RespRef.Offset, - "length": uint32(len(output)), - "flags": p.RespRef.Flags, - }, - }}) - continue + // 结果优先写进内核标定的帧;放不下才申请扩容块。 + ref, err := frameOutput(p.Frame, p.InputLen, []byte(output), false) + if err != nil { + ref, err = arenaPut([]byte(output), false) + if err != nil { + send(response{ID: req.ID, Error: "结果扩容失败: " + err.Error()}) + continue + } } - send(response{ID: req.ID, Result: map[string]interface{}{"text": output}}) + send(response{ID: req.ID, Result: map[string]interface{}{"text_ref": ref}}) case "stage.invoke": go func(id uint64) { diff --git a/plan.md b/plan.md index 48f7f2e..0b2a833 100644 --- a/plan.md +++ b/plan.md @@ -1096,8 +1096,8 @@ SDK 仓 `v1.0.0` / `v1.1.0`。main 的版本路牌现为 `1.2.0`(尚无 tag) | 步骤 | 内容 | 依赖 | 产出 | |---|---|---|---| | 13.1 | 统一共享内存布局 | 无 | 单一 memfd,StageContext+EvtRing 成为区内 segment | -| 13.2 | 共享内存扩缩容 | 13.1 | ftruncate/remap,generation 感知 | -| 13.3 | ToolCall lane | 13.1 | 跨进程工具调用走 SharedRef,RPC 退化为控制信号 | +| 13.2 | 共享内存分配器 | 13.1 | 内核独占的变长块分配器,插件经 RPC 申请/归还 | +| 13.3 | 工具调用调用帧 | 13.1/13.2 | payload 始终走共享内存,内核标定帧,插件按需扩容 | | 13.4 | Cleaner 迁移至 SharedRef | 13.1 | cleaner.invoke 参数/结果走 SharedRef | | 13.5 | InputChannel lane | 13.1 | 输入通道消息走共享内存 | | 13.6 | OutputChannel lane | 13.1 | 输出通道消息走共享内存 | @@ -1147,56 +1147,69 @@ SDK 仓 `v1.0.0` / `v1.1.0`。main 的版本路牌现为 `1.2.0`(尚无 tag) - [ ] fd 数从 3 降到 2 - [ ] git commit -m "feat(shm): unified shared memory region" -### 13.2 共享槽池(内核独占所有权) +### 13.2 内核独占的共享内存分配器 **设计约束(用户明确)**: > 内核应当全权管理共享内存,插件需要共享内存要向内核申请,内核给插件返回偏移与大小,使用完成后插件通知内核回收。 > 内核暴露类似 syscall 的 RPC 接口;共享内存是内部实现,不对插件开发者暴露。 -**为什么不是跨进程分配器**(前两版都被推翻): +**为什么不是跨进程分配器**(前几版都被推翻): - v1:SuperBlock 放 `arenaUsed` 游标,内核 CAS bump。但**插件模板里的 `arenaUsed` 是进程本地变量**,两个进程各自 bump,必然写到同一段内存;`arenaReset` 还会重置共享游标覆盖对方数据。 - v2:把位图 CAS 下沉到插件模板。虽然正确,但把分配器实现细节泄漏进了插件运行时,且插件必须与内核保持位图布局同步。 -- v3(当前):分配器只存在于**内核进程内**,一把 `sync.Mutex` 即可。插件只通过 RPC 申请/归还,不做任何分配决策。 +- v3:定长槽 + 共享位图 CAS。正确,但定长槽唯一的理由是“跨进程没法安全做变长分配”。 +- v4(当前):分配器收回内核进程后那个约束消失,改成**变长块分配器**(first-fit + 邻块合并)。内核可以按需**标定**每块大小,大 payload 不再受固定槽容量限制。 **实施**: -1. `arena.go`:定长槽 + 位图,`Alloc/Put/Read/Free/ReclaimOwner`,内核独占。 -2. 槽头记录 `owner`;`Free` 校验归属,插件不能释放他人(或内核)的槽。 -3. 协议新增 `arena.alloc` / `arena.free`(`CapCore`,属基础能力)。 -4. `Plugin` 在 `Start` 领取 ownerID;`handleExit` 调 `ReclaimOwner` 回收残留槽,防崩溃把池耗尽。 -5. 插件侧不再写位图:`proc_main.go.tmpl` 删掉本地 bump 分配器,改为 `arenaAlloc`/`arenaFree` 走 RPC;SDK 公开 API 仍是普通字符串/Map,开发者无感。 -6. payload 超槽容量时**退回内联 RPC**(原有正确路径),不静默截断。 +1. `arena.go`:块头 16B(`size/state/owner/prevSize`),`Alloc/Put/Read/Free/ReclaimOwner`,内核独占一把 `sync.Mutex`。`prevSize` 让 `Free` 能 O(1) 找到前驱做向后合并。 +2. 块头记录 `owner`;`Free` 校验归属 + 走块链确认 offset 是已分配块的数据起点,**伪造引用不能改动分配器状态**。 +3. `Read` 允许块内偏移(调用帧的结果区就在帧块中间),但要求不跨越块边界。 +4. 协议新增 `arena.alloc` / `arena.free`(`CapCore`,属基础能力)。 +5. `Plugin` 在 `Start` 领取 ownerID;`handleExit` 调 `ReclaimOwner` 回收残留块,防崩溃把 arena 耗尽。 +6. 插件侧不写任何分配器状态:模板只通过 RPC 申请/归还;SDK 公开 API 仍是普通字符串/Map,开发者无感。 +7. arena 容量 4MB,但底层是 memfd:**未触碰的页不占物理内存**,所以开大无成本。 **验证**: -- [x] `TestArena_*`:分配/归还/归属校验/回收/并发唯一/耗尽/超限/非法引用/布局校验 -- [x] `TestPlugin_ArenaAllocFreeAcrossProcess`:真进程申请→写入→随业务 RPC 回传→归还,内核读回内容一致且池归零 -- [ ] **Grow/Shrink 未实现**:跨进程 remap 很危险(对端可能正在读,munmap 会 SIGSEGV)。当前槽池规格固定,用尽后退回内联 RPC。generation 字段已就位但只在 resize 时才会 bump。 -- [x] git commit -m "refactor(shm): kernel-owned exchange arena" +- [x] `TestArena_*`:分配/归还/归属校验/回收/并发唯一/耗尽/超限/非法引用/伪造 offset/相邻合并/对齐/布局校验 +- [x] `TestPlugin_ArenaAllocFreeAcrossProcess`:真进程申请→写入→随业务 RPC 回传→归还,内核读回内容一致且 arena 归零 +- [ ] **Grow/Shrink 未实现**:跨进程 remap 会让正在读的对端 SIGSEGV。当前容量固定,用尽时调用失败(不再退回内联)。 +- [x] git commit -m "refactor(shm): kernel-owned variable-size arena" -### 13.3 ToolCall 数据面走共享槽 +### 13.3 工具调用走 funccall 调用帧 -**目标**:工具调用的**参数/结果**走共享内存,RPC 只传描述符。 +**目标**:工具调用的 payload **始终**在共享内存里;内核作为 caller 标定内存块交给插件。 -**控制面 vs 数据面**(与初版设计的差异): +**模型(用户明确)**: -初版计划用 ring 状态机(FREE→WRITING→READY→READING→DONE)把 `tool.invoke` 也搬出 RPC。实际实现后发现这是错的方向: +> 工具调用的 payload 应当始终在共享内存中。因为工具调用是内核发出的,按 funccall 方式,内存块应当由内核标定后交给子进程。当内核提前给的不够用时,插件侧才请求扩容。 -- RPC 已经提供请求 ID 关联、错误传递、ctx 取消、崩溃唤醒(`Process.CallContext`),ring 状态机是把这些重新实现一遍。 -- ring 版本写完从未接线,属死代码,已删除。 +调用帧布局: -因此**控制面保留 RPC**,只把 payload 搬进共享槽: +```text +[0, ArgsLen) 参数 JSON +[ArgsLen, Frame.Length) 结果区(内核预留的预算) +``` -1. 小 payload 走内联 JSON(省两次 RPC)。 -2. 大 payload 走内核分配的槽,RPC 只带 `SharedRef`。 -3. 插件不需要分配:内核可预分配响应槽(如 Cleaner 路径)。 +**为什么不用 ring 状态机**:初版计划用 FREE→WRITING→READY→READING→DONE 的 ring 把 `tool.invoke` 整个搬出 RPC。写完发现是错的方向——RPC 已经提供请求 ID 关联、错误传递、ctx 取消、崩溃唤醒(`Process.CallContext`),ring 只是把它们重新实现一遍,且从未接线(死代码,已删)。 + +**实施**: + +1. 内核 `invokeTool`:序列化参数 → `Alloc(len(args)+toolResultBudget)` → 参数写帧前段 → 发 `ToolInvokeParams{Name, Frame, ArgsLen}`。 +2. 插件:从帧读参数;结果优先写帧的结果区。 +3. 结果超出预算 → 插件 `arena.alloc` 扩容块,引用上打 `sharedRefFlagExpand`;内核据此单独归还。 +4. **没有按大小切换内联的分支**:小 payload 同样走帧。 +5. Cleaner 复用同一帧模型(`Frame` + `InputLen`)。 +6. `ToolInvokeParams.Args` / `ToolInvokeResult.Result` 仅剩给**直连 RPC 的测试**(process/bench 不建 Host,拿不到共享内存);生产路径永远走帧。 **验证**: -- [x] `TestPlugin_RegisteredToolInvokable` 通过 -- [ ] Bench: ToolInvoke 延迟改善(待 `tool.invoke` 接入 SharedRef 后再测) +- [x] `TestPlugin_ToolInvokeArgsResultViaArena`:大/小 payload 都经帧往返,结果内容一致且 arena 归零 +- [x] `TestE2E_RealTemplatePluginFullLifecycle`:真实 SDK 模板编译的插件跑通 +- [x] `TestProcTemplate_ToolInvokeUsesSharedRef`:模板必须处理 frame/args_len/result_ref(防漂移) +- [ ] Bench: ToolInvoke 延迟对比(待补) ### 13.4 Cleaner 迁移至 SharedRef