diff --git a/internal/plugin/proc/arena.go b/internal/plugin/proc/arena.go new file mode 100644 index 0000000..b63e8c0 --- /dev/null +++ b/internal/plugin/proc/arena.go @@ -0,0 +1,411 @@ +package proc + +// Exchange Arena:**内核独占管理**的跨进程共享内存槽池(§13.2 重设计)。 +// +// ── 所有权模型(架构约束)──────────────────────────────────────────── +// +// 共享内存由内核全权管理。插件需要使用共享内存时,通过 syscall 风格的 +// RPC 向内核申请,内核返回偏移与大小;使用完毕后插件再通知内核回收。 +// +// 因此分配器只存在于**内核进程内**,用一把普通 sync.Mutex 保护即可: +// 不需要任何跨进程原子操作,也不存在"共享游标被两个进程各自更新"这一 +// 类竞态——这正是前两版(bump 游标 / CAS 位图跨进程分配)失败的根本原因。 +// +// 共享内存是**内部实现**,不对插件开发者暴露:插件的公开 API 仍是 +// 普通字符串/Map(见 SDK 的 IOInjector / ToolHandler)。模板运行时在 +// 传输层按 payload 大小自动决定走内联 JSON 还是共享槽,开发者无感。 +// +// ── 槽布局 ──────────────────────────────────────────────────────── +// +// ┌──────────────────────────────────────────────┐ +// │ Arena Header 64B │ +// │ magic / version / slotCount / slotSize │ +// │ bitmapOff / slotsOff / reserved │ +// ├──────────────────────────────────────────────┤ +// │ Bitmap:ceil(slotCount/32) 个 uint32,bit=占用│ +// ├──────────────────────────────────────────────┤ +// │ Slot 0: [dataLen(4) | owner(4) | data(...)] │ +// │ ... │ +// └──────────────────────────────────────────────┘ +// +// ── 所有权与回收 ────────────────────────────────────────────────── +// +// 每个插件进程从内核领一个不透明 ownerID。槽头记录 owner。 +// - 内核 → 插件:内核直接 Alloc,把 ref 放进 RPC 参数(插件不分配) +// - 插件 → 内核:插件 RPC arena.alloc 申请,用完 RPC arena.free 归还 +// - 插件退出:内核 ReclaimOwner 回收其残留槽,避免崩溃泄漏耗尽池 +// +// 插件归还时校验 owner,防止一个插件释放另一个插件的槽。 + +import ( + "fmt" + "sync" + "sync/atomic" + "unsafe" +) + +const ( + arenaMagic uint32 = 0x41524E41 // "ARNA" + arenaVersion uint32 = 1 + + arenaHeaderSize = 64 + + // Arena Header 字段偏移(相对 arena 起始)。 + arOffMagic = 0 + arOffVersion = 4 + arOffSlotCount = 8 + arOffSlotSize = 12 + arOffBitmapOff = 16 + arOffSlotsOff = 20 + arOffReserved = 24 + + // 槽头:数据长度 + 所有者。数据紧随其后。 + slotHeaderSize = 8 + slotOffDataLen = 0 + slotOffOwner = 4 + + // 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 +} + +// arenaPublishSize 是统一区域里预留给槽池的字节数。 +var arenaPublishSize = arenaRequiredSize(arenaDefaultSlotCount, arenaDefaultSlotSize) + +// arenaRegion 是槽池视图。 +// +// region 始终是**完整**统一区域 mmap:SharedRef.Offset 是相对区域起始的 +// 绝对偏移,因此读写都直接落在 region 上,不需要再换算。 +// +// mu 只在**内核进程内**使用——分配器完全由内核持有(见文件头所有权模型)。 +type arenaRegion struct { + mu sync.Mutex + + 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) + } + 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)) + } + + bitmapBytes := ((slotCount+31)/32)*4 + 7 + bitmapBytes &^= 7 + slotsOff := (arenaHeaderSize + bitmapBytes + 7) &^ 7 + + ap := region[base : base+need] + putU32(ap[arOffMagic:], arenaMagic) + putU32(ap[arOffVersion:], arenaVersion) + putU32(ap[arOffSlotCount:], slotCount) + putU32(ap[arOffSlotSize:], slotSize) + putU32(ap[arOffBitmapOff:], arenaHeaderSize) + putU32(ap[arOffSlotsOff:], slotsOff) + + // 显式清零位图与槽头,保证复用已存在区域时状态干净。 + 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) + } + + return &arenaRegion{ + region: region, + base: base, + cap: need, + slots: slotCount, + slotSize: slotSize, + bmOff: arenaHeaderSize, + slotsOff: slotsOff, + }, nil +} + +// attachArena 从已初始化的区域解析槽池(诊断用;模板不再需要解析位图)。 +func attachArena(region []byte, base, size uint32) (*arenaRegion, error) { + if size < arenaHeaderSize { + return nil, fmt.Errorf("arena: 段太小(%d)", size) + } + if uint64(base)+uint64(arenaHeaderSize) > uint64(len(region)) { + return nil, fmt.Errorf("arena: 头部越界(base=%d total=%d)", base, len(region)) + } + ap := region[base : base+size] + + if got := getU32(ap[arOffMagic:]); got != arenaMagic { + return nil, fmt.Errorf("arena: 魔数不匹配(0x%x)", got) + } + 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) + } + 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 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, + slots: slotCount, + slotSize: slotSize, + bmOff: bmOff, + slotsOff: slotsOff, + }, nil +} + +// SlotCount 返回槽总数。 +func (a *arenaRegion) SlotCount() uint32 { return a.slots } + +// 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)。 +// +// 必须由内核调用:分配器由内核独占(见文件头所有权模型)。 +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()) + } + 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 + } + } + return SharedRef{}, fmt.Errorf("arena: 槽池已满(%d 槽全部占用)", a.slots) +} + +// Put 分配并写入 payload,返回共享引用(内核方向发送用)。 +func (a *arenaRegion) Put(owner uint32, payload []byte, gen uint64) (SharedRef, error) { + ref, err := a.Alloc(owner, len(payload)) + 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 读取槽数据。 +// +// 校验四件事,任一不满足都返回 error(而不是像早期版本那样返回 nil, +// 让调用方分不清"空数据"和"非法引用"): +// 1. generation 与当前区域一致(remap 后旧引用失效) +// 2. 槽号在范围内,且 Offset 与槽数据区自洽 +// 3. Length 不超过槽容量 +// 4. 槽当前处于占用状态(已被释放的槽不可再读) +func (a *arenaRegion) Read(ref SharedRef, gen uint64) ([]byte, error) { + if ref.IsZero() { + return nil, nil + } + 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 !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) + } + if !a.isBusyLocked(slot) { + return fmt.Errorf("arena: 槽 %d 未分配(重复释放?)", slot) + } + 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) + } + a.clearLocked(slot) + return nil +} + +// ReclaimOwner 释放 owner 名下所有槽,返回回收数量。 +// +// 用于插件进程退出:崩溃的插件无法归还自己申请的槽,若不管会把池慢慢 +// 耗尽,最终让所有走共享内存的调用退化成内联 RPC。 +func (a *arenaRegion) ReclaimOwner(owner uint32) int { + if owner == OwnerHost { + return 0 // 内核自己的槽由正常路径释放,不在此回收 + } + a.mu.Lock() + defer a.mu.Unlock() + + reclaimed := 0 + for slot := uint32(0); slot < a.slots; slot++ { + if !a.isBusyLocked(slot) { + continue + } + base := a.slotBase(slot) + if getU32(a.region[base+slotOffOwner:]) != owner { + continue + } + a.clearLocked(slot) + reclaimed++ + } + return reclaimed +} + +// Stats 返回 (已用槽数, 总槽数),供诊断与测试断言。 +func (a *arenaRegion) Stats() (used, total uint32) { + for slot := uint32(0); slot < a.slots; slot++ { + if a.IsBusy(slot) { + used++ + } + } + return used, a.slots +} + +// 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 + } + } +} diff --git a/internal/plugin/proc/arena_test.go b/internal/plugin/proc/arena_test.go new file mode 100644 index 0000000..f6d1c98 --- /dev/null +++ b/internal/plugin/proc/arena_test.go @@ -0,0 +1,263 @@ +package proc + +import ( + "bytes" + "sync" + "testing" +) + +// newTestArena 构造一个独立槽池(不经 Host),保持单测快速。 +func newTestArena(t *testing.T, slots, slotSize uint32) *arenaRegion { + t.Helper() + const base = 64 + need := arenaRequiredSize(slots, slotSize) + region := make([]byte, base+need) + a, err := initArena(region, base, need, slots, slotSize) + if err != nil { + t.Fatalf("initArena: %v", err) + } + return a +} + +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) + } + + 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) + } + got, err := a.Read(ref, 0) + if err != nil { + t.Fatalf("Read: %v", err) + } + if string(got) != "hello" { + t.Fatalf("Read=%q, want hello", got) + } + + if err := a.Free(OwnerHost, ref.Flags); 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("已释放槽的引用应读取失败") + } +} + +func TestArena_ConcurrentAllocUnique(t *testing.T) { + const slots = 32 + a := newTestArena(t, slots, 64) + + refs := make(chan SharedRef, slots) + var wg sync.WaitGroup + for i := 0; i < slots; i++ { + wg.Add(1) + go func() { + defer wg.Done() + ref, err := a.Alloc(OwnerHost, 8) + if err != nil { + t.Errorf("Alloc: %v", err) + return + } + refs <- ref + }() + } + wg.Wait() + close(refs) + + seen := make(map[uint32]bool, slots) + for ref := range refs { + if seen[ref.Flags] { + t.Fatalf("并发分配拿到重复槽 %d", ref.Flags) + } + seen[ref.Flags] = 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) + } +} + +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) + } + } + if _, err := a.Alloc(OwnerHost, 8); 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 应被拒绝") + } +} + +func TestArena_FreeEnforcesOwnership(t *testing.T) { + a := newTestArena(t, 4, 64) + const ownerA, ownerB = 7, 9 + + refA, err := a.Alloc(ownerA, 8) + if err != nil { + t.Fatal(err) + } + // 另一个插件不能释放 A 的槽 + if err := a.Free(ownerB, refA.Flags); err == nil { + t.Fatal("跨 owner 释放应被拒绝") + } + if used, _ := a.Stats(); used != 1 { + t.Fatalf("拒绝释放后 used=%d, want 1", used) + } + 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) + } +} + +func TestArena_ReclaimOwnerReleasesOnlyItsSlots(t *testing.T) { + a := newTestArena(t, 8, 64) + const ownerA, ownerB = 7, 9 + + for i := 0; i < 3; i++ { + if _, err := a.Alloc(ownerA, 8); err != nil { + t.Fatal(err) + } + } + for i := 0; i < 2; i++ { + if _, err := a.Alloc(ownerB, 8); 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) + } + // 内核自己的槽不参与回收 + if n := a.ReclaimOwner(OwnerHost); n != 0 { + t.Fatalf("ReclaimOwner(host)=%d, want 0", n) + } +} + +func TestArena_ReclaimHostIsNoop(t *testing.T) { + a := newTestArena(t, 4, 64) + if _, err := a.Alloc(OwnerHost, 8); err != nil { + t.Fatal(err) + } + if n := a.ReclaimOwner(OwnerHost); n != 0 { + t.Fatalf("内核槽不应被 ReclaimOwner 回收,实际回收 %d", n) + } + if used, _ := a.Stats(); used != 1 { + t.Fatalf("used=%d, want 1", used) + } +} + +func TestArena_ReadRejectsInvalidRefs(t *testing.T) { + a := newTestArena(t, 4, 256) + ref, err := a.Put(OwnerHost, []byte("payload"), 0) + if err != nil { + t.Fatal(err) + } + + // generation 过期 + stale := ref + stale.Generation = 99 + if _, err := a.Read(stale, 0); err == nil { + t.Fatal("generation 不匹配应被拒绝") + } + + // offset 与槽不自洽 + badOffset := ref + badOffset.Offset++ + if _, err := a.Read(badOffset, 0); err == nil { + t.Fatal("offset 与槽不自洽应被拒绝") + } + + // 槽号越界 + badSlot := ref + badSlot.Flags = 999 + if _, err := a.Read(badSlot, 0); err == nil { + t.Fatal("槽号越界应被拒绝") + } + + // Length 超容量 + tooLong := ref + tooLong.Length = uint32(a.SlotPayloadCap() + 1) + 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 字节 + ref, err := a.Put(OwnerHost, payload, 3) + if err != nil { + t.Fatalf("Put: %v", err) + } + got, err := a.Read(ref, 3) + if err != nil { + t.Fatalf("Read: %v", err) + } + if !bytes.Equal(got, payload) { + t.Fatalf("读回数据不一致:len(got)=%d len(want)=%d", len(got), len(payload)) + } +} + +func TestArena_AttachRejectsCorruptLayout(t *testing.T) { + a := newTestArena(t, 4, 256) + + // 魔数被改 + 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) + if _, err := attachArena(a.region, a.base, a.cap); err == nil { + t.Fatal("槽数组越界应被拒绝") + } +} + +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) + } + 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) + } + } +} diff --git a/internal/plugin/proc/capability.go b/internal/plugin/proc/capability.go index b9e1991..3c9cdb7 100644 --- a/internal/plugin/proc/capability.go +++ b/internal/plugin/proc/capability.go @@ -84,6 +84,9 @@ var methodCapability = map[string]Capability{ MethodInputRegister: CapCore, MethodStageLock: CapCore, MethodStageUnlock: CapCore, + // 共享槽池申请/归还:内部传输层能力,等同于核心基础能力。 + MethodArenaAlloc: CapCore, + MethodArenaFree: CapCore, // 自身配置读写与元信息属基础能力 MethodSettingsGet: CapCore, MethodSettingsSet: CapCore, diff --git a/internal/plugin/proc/corehandler.go b/internal/plugin/proc/corehandler.go index a222020..3ad36b8 100644 --- a/internal/plugin/proc/corehandler.go +++ b/internal/plugin/proc/corehandler.go @@ -28,13 +28,17 @@ type coreHandler struct { // ❗ 必须是"全部插件共享一个 Host"——每插件一段会退化成副本模型。 host *Host + // owner 是本插件在共享槽池里的身份。内核用它校验 arena.free 的归属, + // 并在插件退出时 ReclaimOwner 回收残留槽。 + owner uint32 + // locks 是 host.locks 的引用,供 stage.lock/unlock 路由。 locks *lockRegistry // invokeTool/invokeCleaner/invokeStageFn/invokeOutput 反向调用插件(内核 → 插件)。 // 由 Plugin 注入,注册回调时用它们构造 handler。 invokeTool func(name string, args map[string]interface{}) (interface{}, error) - invokeCleaner func(scope, name string, textRef SharedRef) (SharedRef, error) + invokeCleaner func(params CleanerInvokeParams) (CleanerInvokeResult, error) invokeStageFn func(ctx context.Context, stage string, seq uint64) error invokeOutput func(channel string, args map[string]interface{}) (interface{}, error) @@ -550,6 +554,25 @@ func (h *coreHandler) Handle(method string, params json.RawMessage) (interface{} // 当前设计:子进程 Stop 时由内核统一清理其订阅。 return nil, nil + // ---- 共享槽池(内部传输层,见 protocol.go 注释)---- + case MethodArenaAlloc: + var p ArenaAllocParams + if err := unmarshal(params, &p); err != nil { + return nil, err + } + ref, err := h.arenaAlloc(p.Size) + if err != nil { + return nil, err + } + return ArenaAllocResult{Ref: ref}, nil + + case MethodArenaFree: + var p ArenaFreeParams + if err := unmarshal(params, &p); err != nil { + return nil, err + } + return nil, h.arenaFree(p.Ref) + // ---- 多模态注入(C ABI 侧空实现)---- case MethodIOSetToolBlocks: // Part 4 扩展:二进制落 arena、Slice 描述符回传(§3.8)。 @@ -559,11 +582,17 @@ func (h *coreHandler) Handle(method string, params json.RawMessage) (interface{} return nil, fmt.Errorf("未知 method: %s", method) } -// resolveText 从 injectParams 中提取 text:优先使用 TextRef(SharedRef), +// resolveText 从 injectParams 中提取 text:优先使用 TextRef(共享槽), // 否则使用内联 Text。兼容新旧两种协议。 +// +// 子进程分配自己的槽并在收到应答后释放,因此这里读到的一定是 busy 槽。 func (h *coreHandler) resolveText(p injectParams) string { if !p.TextRef.IsZero() { - return string(h.host.unified.arenaRead(p.TextRef)) + data, err := h.host.Arena().Read(p.TextRef, h.host.Generation()) + if err != nil { + return "" + } + return string(data) } return p.Text } @@ -610,28 +639,98 @@ func (h *coreHandler) cleanerProxy(scope, name string, enabled bool) (func(strin return nil, fmt.Errorf("%s %s 声明 Cleaner,但清洗回调通道未就绪", scope, name) } return func(text string) string { - // Cleaner 输入和输出共享同一 arena。串行覆盖完整往返,保证结果读完前 - // 不被其他调用重置;返回 string 会复制结果,随后即可回收整个临时区。 - h.host.arenaMu.Lock() - defer h.host.arenaMu.Unlock() - defer h.host.unified.arenaReset() - - ref, err := h.host.unified.arenaWrite([]byte(text)) + out, err := h.invokeCleanerText(scope, name, text) if err != nil { return text } - resultRef, err := h.invokeCleaner(scope, name, ref) - if err != nil { - return text - } - result := h.host.unified.arenaRead(resultRef) - if result == nil && !resultRef.IsZero() { - return text - } - return string(result) + return out }, nil } +// invokeCleanerText 完成一次 Cleaner 往返:payload 走共享槽,超限时整条链路退回内联。 +// +// 槽的申请与归还全部由内核负责(**谁分配谁释放**): +// +// reqSlot ──写入 input──▶ 随 RPC 把 TextRef 发给插件 +// respSlot ──预分配────▶ 随 RPC 把 RespRef 发给插件,插件把结果写进来 +// 读取结果后两个槽一起归还 +// +// 插件侧不申请任何槽,因此不存在跨进程分配器的竞争。 +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 + } + + // 响应槽:预分配满容量。插件只写,不申请。 + var respRef SharedRef + if ref, err := arena.Alloc(OwnerHost, arena.SlotPayloadCap()); err == nil { + params.RespRef = ref + respRef = ref + } + + defer func() { + if !respRef.IsZero() { + _ = arena.Free(OwnerHost, respRef.Flags) + } + if !reqRef.IsZero() { + _ = arena.Free(OwnerHost, reqRef.Flags) + } + }() + + 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 + } + return res.Text, nil +} + +// ---- 共享槽池的 syscall 风格接口(插件 RPC)---- +// +// 共享内存是内部实现,不向插件开发者暴露;模板运行时在传输层调用它们, +// 公开 SDK 仍是普通字符串/Map。 + +// 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)) +} + +// arenaFree 归还本插件申请的槽。内核校验归属,防止释放他人的槽。 +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) +} + // toolRegister 注册插件工具,handler 反向调用插件执行(原 case 1)。 func (h *coreHandler) toolRegister(params json.RawMessage) (interface{}, error) { var p struct { diff --git a/internal/plugin/proc/host.go b/internal/plugin/proc/host.go index 78c7037..6c038db 100644 --- a/internal/plugin/proc/host.go +++ b/internal/plugin/proc/host.go @@ -5,6 +5,7 @@ import ( "log" "os" "sync" + "sync/atomic" pubsdk "gitcode.com/JianFeeeee/homeagent-sdk/sdk" ) @@ -17,9 +18,9 @@ import ( // 副本模型——两个插件各写各的段、各自回读,最后回读者覆盖前者, // lost update 原样复现(§8.4 实测 35.8~36.8%)。 // -// 统一区域(§13.1):单 memfd 包含 SuperBlock + StageContext + EvtRing。 -// 子进程经 fd 3 mmap 同一 memfd → 同一份物理页,区域头 SuperBlock -// 告知各 segment 的偏移与大小。 +// 统一区域(§13.1):单 memfd 包含 SuperBlock + StageContext + EvtRing +// + Exchange Arena。子进程经 fd 3 mmap 同一 memfd → 同一份物理页, +// 区域头 SuperBlock 告知各 segment 的偏移与大小。 // // fd 分配:fd 3 = 统一区域,fd 4 = eventfd。 type Host struct { @@ -29,6 +30,13 @@ type Host struct { seg *Segment // StageContext segment(位于 unified ctxData) shmSize int // 统一区域总大小 + // arena 是跨进程共享槽池(§13.2 重设计)。分配/释放走共享位图 CAS, + // 任意进程/goroutine 并发调用都安全,不需要任何锁。 + arena *arenaRegion + + // nextOwner 给每个插件进程分配一个不透明 owner ID,供 ReclaimOwner 使用。 + nextOwner atomic.Uint32 + // 事件通知(独立于共享段) evtfd *os.File // Unix:eventfd/pipe 读端(fd 4)。Windows 为 nil。 evtNotifyFd int // 通知句柄的平台无关标识 @@ -36,7 +44,6 @@ type Host struct { evtSubscriber EvtRingSubscriber locks *lockRegistry - arenaMu sync.Mutex // 动态 arena 跨进程调用的生命周期锁 stageMu sync.Mutex coordMu sync.Mutex coord *stageCoordinator @@ -53,8 +60,11 @@ type Host struct { // 三者共同点:全部插件看到同一份物理页,段内一律用相对偏移而非指针 // (实验 2 已验证各进程 mmap 到不同虚拟地址时偏移解引用仍正确)。 func NewHost() (*Host, error) { - // 统一区域大小:SuperBlock + StageContext segment + EvtRing segment - unifiedSize := superBlockSize + shmDefaultSize + evtTotalSize + unifiedArenaSize + // 统一区域大小:SuperBlock + StageContext + EvtRing + Exchange Arena + // + // +8 是 arena 基址 8 字节对齐的 padding 余量:arenaOff 向上取整可能 + // 吃掉最多 4 字节,预留 8 字节保证 arenaCap 不会小于 arenaPublishSize。 + unifiedSize := superBlockSize + shmDefaultSize + evtTotalSize + int(arenaPublishSize) + 8 memfd, data, err := allocShm(unifiedSize) if err != nil { return nil, err @@ -67,6 +77,17 @@ func NewHost() (*Host, error) { return nil, err } + // 初始化 Exchange Arena 槽池 + arena, err := initArena(data, ur.arenaOff, ur.arenaCap, arenaDefaultSlotCount, arenaDefaultSlotSize) + if err != nil { + freeShm(memfd, data) + return nil, fmt.Errorf("槽池初始化: %w", err) + } + if used, total := arena.Stats(); used != 0 || total != arenaDefaultSlotCount { + freeShm(memfd, data) + return nil, fmt.Errorf("槽池初始化异常:used=%d total=%d", used, total) + } + // 创建 StageContext segment(位于 SuperBlock 之后) seg, err := NewSegment(ur.ctxData()) if err != nil { @@ -102,6 +123,7 @@ func NewHost() (*Host, error) { unified: ur, seg: seg, shmSize: unifiedSize, + arena: arena, evtfd: evtfdReadFile(efd), evtNotifyFd: efd, evtRing: evtRing, @@ -302,6 +324,20 @@ func (c *stageCoordinator) finish(sc *pubsdk.StageContext, written bool) error { return nil } +// Arena 返回跨进程共享槽池。 +func (h *Host) Arena() *arenaRegion { return h.arena } + +// Generation 返回统一区域当前 generation(SharedRef 校验用)。 +func (h *Host) Generation() uint64 { return h.unified.generation() } + +// NextOwnerID 分配一个插件专用的槽 owner ID。 +// +// 从 1 开始(OwnerHost=0 保留给内核),单调递增,不会重复。 +func (h *Host) NextOwnerID() uint32 { return h.nextOwner.Add(1) } + +// ReclaimOwner 回收某个 owner 名下所有槽(插件退出时调用)。 +func (h *Host) ReclaimOwner(owner uint32) int { return h.arena.ReclaimOwner(owner) } + // ShmSize 返回共享段大小(供诊断/日志)。 func (h *Host) ShmSize() int { return h.shmSize } diff --git a/internal/plugin/proc/plugin.go b/internal/plugin/proc/plugin.go index 372e5db..e9921c4 100644 --- a/internal/plugin/proc/plugin.go +++ b/internal/plugin/proc/plugin.go @@ -43,6 +43,10 @@ type Plugin struct { // caps 是 manifest 声明的能力集(§3.8 权限梯度)。 caps *capabilitySet + // ownerID 是本插件在共享槽池里的身份(由 Host 分配)。 + // 内核用它校验 arena.free 的归属,并在插件退出时回收残留槽。 + ownerID uint32 + stopOnce sync.Once // stopping 标记「本次退出是内核主动发起的」,用于压掉 onCrash。 @@ -109,10 +113,14 @@ func (p *Plugin) Start(core CoreSDK) error { return fmt.Errorf("proc: %s 缺少共享段 Host", p.name) } + // 领一个槽池身份;插件退出时用它回收残留槽。 + p.ownerID = p.host.NextOwnerID() + p.handler = &coreHandler{ sdk: core, name: p.name, host: p.host, + owner: p.ownerID, locks: p.host.locks, evtRing: p.host.evtSubscriber, caps: p.caps, @@ -191,8 +199,15 @@ func (p *Plugin) Close() error { // 后者是"锁仲裁回内核"的自愈价值:持锁者死亡不会导致全局死锁, // 无需 robust pthread_mutex(实验 9)。 func (p *Plugin) handleExit(name string, err error) { - if p.host != nil && p.host.ForceReleaseLock(name) { - log.Printf("[proc] %s 退出,内核已释放其持有的 stage 锁", name) + if p.host != nil { + if p.host.ForceReleaseLock(name) { + log.Printf("[proc] %s 退出,内核已释放其持有的 stage 锁", name) + } + // 回收该插件未归还的共享槽:崩溃的插件无法自己归还, + // 不回收会让槽池慢慢耗尽,最终所有共享内存调用退化成内联 RPC。 + if n := p.host.ReclaimOwner(p.ownerID); n > 0 { + log.Printf("[proc] %s 退出,回收 %d 个残留共享槽", name, n) + } } // 内核主动停止(Stop/Close,含宽限期超时后的 Kill)不算崩溃: // 否则重载/禁用/卸载都会误触发自动重启。 @@ -222,23 +237,22 @@ func (p *Plugin) invokeTool(name string, args map[string]interface{}) (interface } // invokeCleaner 在插件进程内执行工具或通道注册时提供的 Cleaner 函数。 -func (p *Plugin) invokeCleaner(scope, name string, textRef SharedRef) (SharedRef, error) { +// +// 数据面参数(输入槽 / 预分配响应槽)由内核在 params 里给出, +// 插件只负责读输入、写结果,不做任何分配。 +func (p *Plugin) invokeCleaner(params CleanerInvokeParams) (CleanerInvokeResult, error) { if p.proc == nil { - return SharedRef{}, ErrProcessExited + return CleanerInvokeResult{}, ErrProcessExited } - raw, err := p.proc.Call(MethodCleanerInvoke, CleanerInvokeParams{ - Scope: scope, - Name: name, - TextRef: textRef, - }) + raw, err := p.proc.Call(MethodCleanerInvoke, params) if err != nil { - return SharedRef{}, err + return CleanerInvokeResult{}, err } var res CleanerInvokeResult if err := json.Unmarshal(raw, &res); err != nil { - return SharedRef{}, fmt.Errorf("proc: %s %s Cleaner %s 应答解析失败: %w", p.name, scope, name, err) + return CleanerInvokeResult{}, fmt.Errorf("proc: %s %s Cleaner %s 应答解析失败: %w", p.name, params.Scope, params.Name, err) } - return res.TextRef, nil + return res, nil } func (p *Plugin) invokeStage(ctx context.Context, stage string, seq uint64) error { diff --git a/internal/plugin/proc/plugin_test.go b/internal/plugin/proc/plugin_test.go index 193b8a8..a7390fb 100644 --- a/internal/plugin/proc/plugin_test.go +++ b/internal/plugin/proc/plugin_test.go @@ -26,6 +26,9 @@ type fakeCoreSDK struct { inputDefs map[string]pubsdk.ChannelDef settings map[string]interface{} autoStart bool + + // injected 记录经 InjectText 注入的文本(验证跨进程共享槽路径)。 + injected []string } func newFakeCore() *fakeCoreSDK { @@ -40,17 +43,27 @@ func newFakeCore() *fakeCoreSDK { } } -func (f *fakeCoreSDK) PluginName() string { return "fake" } -func (f *fakeCoreSDK) Settings() pubsdk.SettingsAPI { return nil } -func (f *fakeCoreSDK) Memory() pubsdk.MemoryAPI { return nil } -func (f *fakeCoreSDK) TextMemory() pubsdk.TextMemoryAPI { return nil } -func (f *fakeCoreSDK) DocMemory() pubsdk.DocMemoryAPI { return nil } -func (f *fakeCoreSDK) Knowledge() pubsdk.KnowledgeAPI { return nil } -func (f *fakeCoreSDK) LLM() pubsdk.LLMAPI { return nil } -func (f *fakeCoreSDK) Social() pubsdk.SocialAPI { return nil } -func (f *fakeCoreSDK) PluginMgr() pubsdk.PluginMgrAPI { return nil } -func (f *fakeCoreSDK) RegisterPluginAPI(name string) error { return nil } -func (f *fakeCoreSDK) InjectText(s, c, t string) {} +func (f *fakeCoreSDK) PluginName() string { return "fake" } +func (f *fakeCoreSDK) Settings() pubsdk.SettingsAPI { return nil } +func (f *fakeCoreSDK) Memory() pubsdk.MemoryAPI { return nil } +func (f *fakeCoreSDK) TextMemory() pubsdk.TextMemoryAPI { return nil } +func (f *fakeCoreSDK) DocMemory() pubsdk.DocMemoryAPI { return nil } +func (f *fakeCoreSDK) Knowledge() pubsdk.KnowledgeAPI { return nil } +func (f *fakeCoreSDK) LLM() pubsdk.LLMAPI { return nil } +func (f *fakeCoreSDK) Social() pubsdk.SocialAPI { return nil } +func (f *fakeCoreSDK) PluginMgr() pubsdk.PluginMgrAPI { return nil } +func (f *fakeCoreSDK) RegisterPluginAPI(name string) error { return nil } +func (f *fakeCoreSDK) InjectText(s, c, t string) { + f.mu.Lock() + f.injected = append(f.injected, t) + f.mu.Unlock() +} + +func (f *fakeCoreSDK) injectedTexts() []string { + f.mu.Lock() + defer f.mu.Unlock() + return append([]string(nil), f.injected...) +} func (f *fakeCoreSDK) InjectInterruptText(s, c, t string) {} func (f *fakeCoreSDK) InjectTextNoMemory(s, c, t string) {} func (f *fakeCoreSDK) InjectInputSync(s, c, t string) string { return "" } @@ -439,3 +452,51 @@ func TestPlugin_FiveProcessesConcurrentAppendNoLostUpdate(t *testing.T) { } } } + +// 跨进程共享槽池:插件通过 arena.alloc 申请、写入、随业务 RPC 回传、arena.free 归还。 +// +// 验证内核独占管理的所有权模型在真进程 + 真 RPC 下成立: +// - 内核能把插件申请的槽内容正确读回来(偏移/长度无误) +// - 插件归还后槽确实回到池里(无泄漏) +func TestPlugin_ArenaAllocFreeAcrossProcess(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_inject"] + core.mu.Unlock() + if !ok { + t.Fatal("插件应注册 demo_inject 工具") + } + + // 用明显超过内联阈值的 payload,确保真的走共享槽而非内联。 + payload := strings.Repeat("共享内存", 500) // 约 6KB + if _, err := h(map[string]interface{}{"text": payload}); err != nil { + t.Fatalf("调用 demo_inject: %v", err) + } + + got := core.injectedTexts() + if len(got) != 1 { + t.Fatalf("应注入 1 条文本,实际 %d 条", len(got)) + } + if got[0] != payload { + t.Fatalf("经共享槽读到的内容不一致:len(got)=%d len(want)=%d", len(got[0]), len(payload)) + } + + // 插件已归还槽:池必须回到全空,否则说明 arena.free 没生效。 + if used, total := host.Arena().Stats(); used != 0 { + t.Fatalf("插件归还后槽池应全空,实际 used=%d/%d", used, total) + } +} diff --git a/internal/plugin/proc/protocol.go b/internal/plugin/proc/protocol.go index e79941d..4bf30ef 100644 --- a/internal/plugin/proc/protocol.go +++ b/internal/plugin/proc/protocol.go @@ -57,6 +57,14 @@ const ( MethodStageInvoke = "stage.invoke" // MethodOutputInvoke 经插件输出通道发送。 MethodOutputInvoke = "output.invoke" + // MethodArenaAlloc 向内核申请一块共享内存(返回偏移与大小)。 + // MethodArenaFree 通知内核回收先前申请的共享内存。 + // + // 这两个是**内部传输层接口**,不由插件开发者直接使用:共享内存是 + // 内核的内部实现,公开 SDK 仍是普通字符串/Map,模板运行时按 + // payload 大小自动选择内联 JSON 还是共享槽。 + MethodArenaAlloc = "arena.alloc" + MethodArenaFree = "arena.free" // MethodHandshake 建链首帧:交换协议版本、SDK 版本、共享段规格。 MethodHandshake = "handshake" ) @@ -202,26 +210,40 @@ type StageInvokeResult struct { } // ToolInvokeParams / ToolInvokeResult:工具调用(原 go_invoke_tool)。 +// +// ArgsRef/ResultRef 为后续 §13.3 预留:payload 走 Exchange Arena, +// RPC 帧只带 16 字节描述符;当前仍以 Args/Result 内联为主。 type ToolInvokeParams struct { - Name string `json:"name"` - Args map[string]interface{} `json:"args,omitempty"` + Name string `json:"name"` + Args map[string]interface{} `json:"args,omitempty"` + ArgsRef SharedRef `json:"args_ref,omitempty"` } type ToolInvokeResult struct { - Result interface{} `json:"result,omitempty"` + Result interface{} `json:"result,omitempty"` + ResultRef SharedRef `json:"result_ref,omitempty"` } // CleanerInvokeParams / CleanerInvokeResult:跨进程计算层清洗。 // Scope 取 tool / input / output,Name 是工具名或通道名。 -// TextRef 指向共享内存 arena 中的实际文本数据。 +// +// 数据面走 Exchange Arena: +// - TextRef 指向内核写入的输入文本;payload 超过槽容量时改用内联 Text +// - RespRef 指向内核**预分配**的响应槽,插件把结果写进去后回 TextRef; +// 结果放不下时插件改用内联 Text 返回 +// +// 插件侧不做任何分配,因此不存在跨进程分配器的竞争。 type CleanerInvokeParams struct { Scope string `json:"scope"` Name string `json:"name"` - TextRef SharedRef `json:"text_ref"` + Text string `json:"text,omitempty"` + TextRef SharedRef `json:"text_ref,omitempty"` + RespRef SharedRef `json:"resp_ref,omitempty"` } type CleanerInvokeResult struct { - TextRef SharedRef `json:"text_ref"` + Text string `json:"text,omitempty"` + TextRef SharedRef `json:"text_ref,omitempty"` } const ( @@ -241,6 +263,23 @@ type OutputInvokeParams struct { Args map[string]interface{} `json:"args,omitempty"` } +// ArenaAllocParams / ArenaAllocResult:插件向内核申请共享内存。 +// +// 内核返回的 Ref.Length 是槽容量(可写上限),插件写入后自行把 Length +// 改成实际 payload 长度再随业务 RPC 回传。 +type ArenaAllocParams struct { + Size uint32 `json:"size"` +} + +type ArenaAllocResult struct { + Ref SharedRef `json:"ref"` +} + +// ArenaFreeParams:插件通知内核回收共享内存。只带描述符,槽号在 Flags 里。 +type ArenaFreeParams struct { + Ref SharedRef `json:"ref"` +} + // PluginInitParams:插件构造参数(原 case init_plugin)。 type PluginInitParams struct { Name string `json:"name"` diff --git a/internal/plugin/proc/testdata/stageplugin.go b/internal/plugin/proc/testdata/stageplugin.go index 145d6cd..f96f59d 100644 --- a/internal/plugin/proc/testdata/stageplugin.go +++ b/internal/plugin/proc/testdata/stageplugin.go @@ -38,13 +38,63 @@ var ( pendMu sync.Mutex pending = map[uint64]chan json.RawMessage{} - shm []byte - shmSize int - region []byte // 完整统一区域 mmap + shm []byte + shmSize int + region []byte // 完整统一区域 mmap arenaOff uint32 - arenaUsed uint32 ) +// ---- Exchange Arena 槽池(与内核 proc/arena.go 一致)---- + +const ( + arOffSlotCount = 8 + arOffSlotSize = 12 + slotHeaderSize = 8 +) + +// arenaSlotSize 返回单个槽的总字节数(含 8 字节槽头)。 +func arenaSlotSize() uint32 { + return binary.LittleEndian.Uint32(region[arenaOff+arOffSlotSize:]) +} + +// arenaPayloadCap 返回槽可承载的最大 payload。 +func arenaPayloadCap() int { return int(arenaSlotSize()) - slotHeaderSize } + +// ---- 共享槽池的插件侧接口(内核 RPC,内部实现)---- + +// SharedRef 是内核下发的共享内存描述符。 +// +// 这是**内部实现**:真实插件不直接触碰它,模板运行时在传输层自动使用。 +// 本测试插件手写 RPC,所以需要显式声明。 +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 } + +// arenaAlloc 向内核申请一块共享内存,内核返回偏移与大小。 +func arenaAlloc(size uint32) (SharedRef, error) { + raw := callKernel("arena.alloc", map[string]interface{}{"size": size}) + 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("内核返回空引用") + } + return r.Ref, nil +} + +// arenaFree 通知内核回收先前申请的共享内存。 +func arenaFree(ref SharedRef) { + callKernel("arena.free", map[string]interface{}{"ref": ref}) +} + func send(v interface{}) { b, _ := json.Marshal(v) writeMu.Lock() @@ -255,6 +305,11 @@ func main() { "def": map[string]interface{}{"name": "demo_upper", "description": "转大写"}, "has_cleaner": true, }) + callKernel("tool.register", map[string]interface{}{ + "name": "demo_inject", + "def": map[string]interface{}{"name": "demo_inject", "description": "共享内存注入"}, + "has_cleaner": false, + }) callKernel("stage.register", map[string]interface{}{ "stage": "after_toolcall", "scope": "global", @@ -275,24 +330,56 @@ func main() { os.Exit(0) case "tool.invoke": - var p struct { - Name string `json:"name"` - Args map[string]interface{} `json:"args"` - } - json.Unmarshal(req.Params, &p) - text, _ := p.Args["text"].(string) - send(response{ID: req.ID, Result: map[string]interface{}{ - "result": strings.ToUpper(text), - }}) + // 必须在独立 goroutine 里处理:demo_inject 会在处理过程中反向 + // 调用内核(arena.alloc / io.injectText / arena.free),而内核的 + // 应答只能由主读循环接收。若同步处理就会自锁死——这也是真实 + // 模板对每个内核请求都 `go handleKernelRequest` 的原因。 + go func(id uint64, raw json.RawMessage) { + var p struct { + Name string `json:"name"` + Args map[string]interface{} `json:"args"` + } + json.Unmarshal(raw, &p) + text, _ := p.Args["text"].(string) + + // demo_inject 走插件侧共享内存路径: + // 申请 → 写入 → 随业务 RPC 回传 → 归还。 + if p.Name == "demo_inject" { + ref, err := arenaAlloc(uint32(len(text))) + if err != nil { + send(response{ID: id, Error: err.Error()}) + return + } + copy(region[ref.Offset:ref.Offset+uint32(len(text))], text) + ref.Length = uint32(len(text)) + callKernel("io.injectText", map[string]interface{}{ + "source": "plugin", "channel": "demo", "text_ref": ref, + }) + arenaFree(ref) + send(response{ID: id, Result: map[string]interface{}{"result": "injected"}}) + return + } + + send(response{ID: id, Result: map[string]interface{}{ + "result": strings.ToUpper(text), + }}) + }(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"` } json.Unmarshal(req.Params, &p) valid := (p.Scope == "tool" && p.Name == "demo_upper") || @@ -302,23 +389,28 @@ func main() { send(response{ID: req.ID, Error: "未注册 Cleaner"}) continue } - // 从 arena 读文本 - input := string(region[p.TextRef.Offset : p.TextRef.Offset+p.TextRef.Length]) - output := p.Scope + "-cleaned:" + input - // 写回 arena - off := arenaUsed - if off == 0 { - off = 1 + + // 读输入:优先共享槽,否则内联。 + input := p.Text + if p.TextRef.Length > 0 { + input = string(region[p.TextRef.Offset : p.TextRef.Offset+p.TextRef.Length]) } - end := off + uint32(len(output)) - copy(region[arenaOff+end-uint32(len(output)):arenaOff+end], output) - arenaUsed = end - send(response{ID: req.ID, Result: map[string]interface{}{ - "text_ref": map[string]interface{}{ - "offset": arenaOff + off, - "length": uint32(len(output)), - }, - }}) + 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 + } + send(response{ID: req.ID, Result: map[string]interface{}{"text": output}}) case "stage.invoke": go func(id uint64) { diff --git a/internal/plugin/proc/toollane.go b/internal/plugin/proc/toollane.go deleted file mode 100644 index ca3428e..0000000 --- a/internal/plugin/proc/toollane.go +++ /dev/null @@ -1,261 +0,0 @@ -package proc - -// ToolCall lane 实现(§13.3)。 -// -// 采用操作系统函数调用栈帧模型:每个工具调用 = 一个固定帧,帧满背压,完成即回收。 -// 不是 arena 堆分配,是 ring buffer 上的 push/pop。 -// -// +--------------------------------------------+ -// | ToolCall Ring Header 32B | -// | magic / cap / frameSize / writeIdx | -// +--------------------------------------------+ -// | Frame 0: state(4) + reqID(8) + name(64) + | -// | inputRef(16) + outputRef(16) | -// | Frame 1: ... | -// | Frame N-1: ... | -// +--------------------------------------------+ -// | Dynamic Arena (input/output JSON data) | -// +--------------------------------------------+ -// -// 状态机: -// FREE -> CALLING (内核写 inputRef,eventfd 通知插件) -// -> READING (插件读 input,执行,写 outputRef,通知内核) -// -> FREE (内核读 output,帧回收) -// -// 背压:writeIdx 追上最后的 busy 帧时返回 ErrToolCallRingFull。 -// 零竞争:插件只操作自己 requestID 对应的帧。 - -import ( - "errors" - "fmt" - "sync/atomic" - "unsafe" -) - -const ( - toolRingMagic uint32 = 0x54524C47 // "TRLG" — TooL Ring - toolRingVersion uint32 = 1 - toolRingCap uint32 = 64 // 最大并发工具调用 - toolFrameNameLen uint32 = 64 // 工具名最大长度 - - // 帧偏移(相对帧起始) - toolFrameOffState = 0 // uint32: 帧状态 - toolFrameOffReqID = 8 // uint64: 请求 ID - toolFrameOffName = 16 // [16, 80): 工具名(64B) - toolFrameOffInput = 80 // SharedRef(16): input 描述符 - toolFrameOffOutput = 96 // SharedRef(16): output 描述符 - toolFrameSize = 112 // 单帧总大小 - - // 帧状态 - toolFrameFree uint32 = 0 - toolFrameReserved uint32 = 1 // 内核已占帧,正在写入元数据 - toolFrameCalling uint32 = 2 // 内核写完 input,插件待读 - toolFrameReading uint32 = 3 // 插件正在执行 - toolFrameReady uint32 = 4 // 插件写完 output,内核待读 -) - -// ToolRing 头偏移(相对区域起始) -const ( - trlOffMagic uint32 = 0 - trlOffVersion uint32 = 4 - trlOffCap uint32 = 8 - trlOffFrameSize uint32 = 12 - trlOffFrameBase uint32 = 16 -) - -var ( - ErrToolCallRingFull = errors.New("proc: tool call ring full") - ErrToolCallNotFound = errors.New("proc: tool call frame not found") -) - -// toolFrame 是单个帧的内存布局。 -type toolFrame struct { - state uint32 // atomic - reqID uint64 - name [toolFrameNameLen]byte - input SharedRef - output SharedRef -} - -// ToolCallRing 内核侧的 ToolCall ring 视图。 -type ToolCallRing struct { - data []byte - cap uint32 - frameSize uint32 - framesBase uint32 - writeIdx atomic.Uint64 // 下一个可写帧 index - gen *uint32 // 统一区域 generation(变化则 snapshot 失效) -} - -// InitToolRing 在统一区域的 ToolCall 子区域上初始化 ring header。 -func InitToolRing(data []byte) error { - if uint32(len(data)) < trlOffFrameBase+toolRingCap*toolFrameSize { - return fmt.Errorf("tool ring: 区域过小(%d 字节,至少需要 %d)", - len(data), trlOffFrameBase+toolRingCap*toolFrameSize) - } - - putU32(data[trlOffMagic:], toolRingMagic) - putU32(data[trlOffVersion:], toolRingVersion) - putU32(data[trlOffCap:], toolRingCap) - putU32(data[trlOffFrameSize:], toolFrameSize) - - // 初始化所有帧为 FREE - for i := uint32(0); i < toolRingCap; i++ { - frameOff := trlOffFrameBase + i*toolFrameSize - putU32(data[frameOff+toolFrameOffState:], toolFrameFree) - } - - return nil -} - -// AttachToolRing 挂载已初始化的 ToolCall ring(插件侧调用)。 -func AttachToolRing(data []byte) (*ToolCallRing, error) { - if uint32(len(data)) < trlOffFrameBase { - return nil, fmt.Errorf("tool ring: 区域过小") - } - magic := getU32(data[trlOffMagic:]) - if magic != toolRingMagic { - return nil, fmt.Errorf("tool ring: 魔数不匹配(0x%x)", magic) - } - cap := getU32(data[trlOffCap:]) - frameSize := getU32(data[trlOffFrameSize:]) - version := getU32(data[trlOffVersion:]) - if version != toolRingVersion { - return nil, fmt.Errorf("tool ring: 版本不匹配(%d,期望 %d)", version, toolRingVersion) - } - if cap == 0 || frameSize != toolFrameSize || uint64(trlOffFrameBase)+uint64(cap)*uint64(frameSize) > uint64(len(data)) { - return nil, fmt.Errorf("tool ring: 布局非法(cap=%d frameSize=%d total=%d)", cap, frameSize, len(data)) - } - return &ToolCallRing{ - data: data, - cap: cap, - frameSize: frameSize, - framesBase: trlOffFrameBase, - }, nil -} - -// Reserve 为一次工具调用预留帧(内核调用)。 -// -// 环形扫描:从 writeIdx 开始,绕环一圈找 FREE 帧。全部忙时背压。 -func (r *ToolCallRing) Reserve() (frameIdx uint32, err error) { - start := r.writeIdx.Add(1) - 1 - for i := uint64(0); i < uint64(r.cap); i++ { - idx := start + i - fi := uint32(idx % uint64(r.cap)) - if atomic.CompareAndSwapUint32(r.statePtr(fi), toolFrameFree, toolFrameReserved) { - return fi, nil - } - } - return 0, ErrToolCallRingFull -} - -// SetCalling 将帧状态设为 CALLING 并写入 requestID 和工具名(内核调用)。 -func (r *ToolCallRing) SetCalling(frameIdx uint32, reqID uint64, toolName string, inputRef SharedRef) { - off := r.frameOff(frameIdx) - // 写 requestID - putU64(r.data[off+toolFrameOffReqID:], reqID) - // 写工具名(截断到 64B) - nameBytes := []byte(toolName) - if len(nameBytes) > int(toolFrameNameLen) { - nameBytes = nameBytes[:toolFrameNameLen] - } - for i := range r.data[off+toolFrameOffName : off+toolFrameOffName+toolFrameNameLen] { - r.data[off+toolFrameOffName+uint32(i)] = 0 - } - copy(r.data[off+toolFrameOffName:], nameBytes) - // 写 input SharedRef - packSharedRef(r.data[off+toolFrameOffInput:], inputRef) - // 原子设置状态为 CALLING - atomic.StoreUint32(r.statePtr(frameIdx), toolFrameCalling) -} - -// FindFrameByReqID 查找对应 requestID 的帧(插件调用)。 -// 返回帧索引和帧内偏移数据,或 ErrToolCallNotFound。 -func (r *ToolCallRing) FindFrameByReqID(reqID uint64) (uint32, error) { - for i := uint32(0); i < r.cap; i++ { - off := r.frameOff(i) - state := atomic.LoadUint32(r.statePtr(i)) - if state == toolFrameCalling || state == toolFrameReading { - fid := getU64(r.data[off+toolFrameOffReqID:]) - if fid == reqID { - return i, nil - } - } - } - return 0, ErrToolCallNotFound -} - -// GetCallingFrame 读取帧的工具名和 input 描述符(插件侧,正在 CALLING 的帧)。 -func (r *ToolCallRing) GetCallingFrame(frameIdx uint32) (name string, inputRef SharedRef) { - off := r.frameOff(frameIdx) - nameBytes := r.data[off+toolFrameOffName : off+toolFrameOffName+toolFrameNameLen] - // 去尾零 - n := len(nameBytes) - for n > 0 && nameBytes[n-1] == 0 { - n-- - } - name = string(nameBytes[:n]) - inputRef = unpackSharedRef(r.data[off+toolFrameOffInput:]) - return -} - -// SetReading 将 CALLING 帧原子转为 READING(插件开始执行)。 -func (r *ToolCallRing) SetReading(frameIdx uint32) bool { - return atomic.CompareAndSwapUint32(r.statePtr(frameIdx), toolFrameCalling, toolFrameReading) -} - -// SetReady 将 READING 帧设为 READY 并写入 output 描述符(插件执行完毕)。 -func (r *ToolCallRing) SetReady(frameIdx uint32, outputRef SharedRef) bool { - if atomic.LoadUint32(r.statePtr(frameIdx)) != toolFrameReading { - return false - } - off := r.frameOff(frameIdx) - packSharedRef(r.data[off+toolFrameOffOutput:], outputRef) - atomic.StoreUint32(r.statePtr(frameIdx), toolFrameReady) - return true -} - -// GetReadyFrame 读取帧的 requestID 和 output 描述符(内核消费 READY 帧)。 -func (r *ToolCallRing) GetReadyFrame(frameIdx uint32) (reqID uint64, outputRef SharedRef) { - off := r.frameOff(frameIdx) - reqID = getU64(r.data[off+toolFrameOffReqID:]) - outputRef = unpackSharedRef(r.data[off+toolFrameOffOutput:]) - return -} - -// ReleaseFrame 回收帧为 FREE(内核消费完毕后调用)。 -func (r *ToolCallRing) ReleaseFrame(frameIdx uint32) { - off := r.frameOff(frameIdx) - // state 必须最后发布 FREE;否则并发 Reserve 可能在其余字段尚未清零时复用帧。 - for i := uint32(4); i < r.frameSize; i++ { - r.data[off+i] = 0 - } - atomic.StoreUint32(r.statePtr(frameIdx), toolFrameFree) -} - -func (r *ToolCallRing) frameOff(idx uint32) uint32 { - return r.framesBase + idx*r.frameSize -} - -func (r *ToolCallRing) statePtr(idx uint32) *uint32 { - off := r.framesBase + idx*r.frameSize + toolFrameOffState - return (*uint32)(unsafe.Pointer( - uintptr(unsafe.Pointer(&r.data[0])) + uintptr(off), - )) -} - -func packSharedRef(b []byte, ref SharedRef) { - putU32(b[0:], ref.Offset) - putU32(b[4:], ref.Length) - putU32(b[8:], ref.Generation) - putU32(b[12:], ref.Flags) -} - -func unpackSharedRef(b []byte) SharedRef { - return SharedRef{ - Offset: getU32(b[0:]), - Length: getU32(b[4:]), - Generation: getU32(b[8:]), - Flags: getU32(b[12:]), - } -} diff --git a/internal/plugin/proc/toollane_test.go b/internal/plugin/proc/toollane_test.go deleted file mode 100644 index c722513..0000000 --- a/internal/plugin/proc/toollane_test.go +++ /dev/null @@ -1,217 +0,0 @@ -package proc - -import ( - "sync" - "testing" - "unsafe" -) - -func TestToolCallRing_InitAndReserve(t *testing.T) { - arenaSize := 256 * 1024 - size := superBlockSize + int(shmDefaultSize) + int(evtTotalSize) + arenaSize + int(trlOffFrameBase+toolRingCap*toolFrameSize) - data := make([]byte, size) - putU32(data[0:], unifiedMagic) - putU32(data[4:], unifiedVersion) - putU32(data[16:], uint32(size)) - - ringData := data[size-int(trlOffFrameBase+toolRingCap*toolFrameSize) : size] - if err := InitToolRing(ringData); err != nil { - t.Fatalf("InitToolRing: %v", err) - } - - ring, err := AttachToolRing(ringData) - if err != nil { - t.Fatalf("AttachToolRing: %v", err) - } - if ring.cap != toolRingCap { - t.Fatalf("cap: got %d, want %d", ring.cap, toolRingCap) - } - - // Reserve 帧应成功 - idx, err := ring.Reserve() - if err != nil { - t.Fatalf("Reserve: %v", err) - } - if idx != 0 { - t.Fatalf("首帧 index: got %d, want 0", idx) - } - - // SetCalling -> GetCallingFrame -> SetReady -> GetReadyFrame -> ReleaseFrame - inputRef := SharedRef{Offset: 1024, Length: 100, Generation: 0, Flags: 1} - ring.SetCalling(idx, 42, "demo_upper", inputRef) - - name, gotInput := ring.GetCallingFrame(idx) - if name != "demo_upper" { - t.Fatalf("GetCallingFrame name: got %q, want %q", name, "demo_upper") - } - if gotInput.Offset != 1024 || gotInput.Length != 100 { - t.Fatalf("GetCallingFrame input: got %+v", gotInput) - } - - ring.SetReading(idx) - - outputRef := SharedRef{Offset: 2048, Length: 50, Generation: 0, Flags: 0} - ring.SetReady(idx, outputRef) - - reqID, gotOutput := ring.GetReadyFrame(idx) - if reqID != 42 { - t.Fatalf("GetReadyFrame reqID: got %d, want 42", reqID) - } - if gotOutput.Offset != 2048 || gotOutput.Length != 50 { - t.Fatalf("GetReadyFrame output: got %+v", gotOutput) - } - - ring.ReleaseFrame(idx) - - // 释放后应能重用 - idx2, err := ring.Reserve() - if err != nil { - t.Fatalf("Reserve after release: %v", err) - } - if idx2 != 1 { - t.Fatalf("第二帧 index: got %d, want 1", idx2) - } -} - -func TestToolCallRing_FullBackpressure(t *testing.T) { - size := int(toolRingCap)*int(toolFrameSize) + 64*1024 - data := make([]byte, size) - if err := InitToolRing(data); err != nil { - t.Fatalf("InitToolRing: %v", err) - } - ring, err := AttachToolRing(data) - if err != nil { - t.Fatalf("AttachToolRing: %v", err) - } - - // Reserve 所有帧 - for i := uint32(0); i < toolRingCap; i++ { - idx, err := ring.Reserve() - if err != nil { - t.Fatalf("Reserve %d: %v", i, err) - } - ring.SetCalling(idx, uint64(i+1), "tool", SharedRef{}) - ring.SetReading(idx) - // 不释放 - } - - // 再 Reserve 应失败(背压) - _, err = ring.Reserve() - if err != ErrToolCallRingFull { - t.Fatalf("满帧后 Reserve 应返回 ErrToolCallRingFull,实际: %v", err) - } - - // 释放一帧后应恢复 - ring.ReleaseFrame(0) - idx, err := ring.Reserve() - if err != nil { - t.Fatalf("释放后 Reserve: %v", err) - } - _ = idx -} - -func TestToolCallRing_FindFrameByReqID(t *testing.T) { - size := int(toolRingCap)*int(toolFrameSize) + 64*1024 - data := make([]byte, size) - if err := InitToolRing(data); err != nil { - t.Fatal(err) - } - ring, _ := AttachToolRing(data) - - // 没有帧时应找不到 - _, err := ring.FindFrameByReqID(999) - if err != ErrToolCallNotFound { - t.Fatalf("空 ring 应返回 ErrToolCallNotFound,实际: %v", err) - } - - // Reserve + SetCalling 后应能找到 - idx, _ := ring.Reserve() - ring.SetCalling(idx, 100, "tool_a", SharedRef{}) - ring.SetReading(idx) - - found, err := ring.FindFrameByReqID(100) - if err != nil { - t.Fatalf("FindFrameByReqID: %v", err) - } - if found != idx { - t.Fatalf("FindFrameByReqID: got %d, want %d", found, idx) - } - - // Release 后应找不到 - ring.ReleaseFrame(idx) - _, err = ring.FindFrameByReqID(100) - if err != ErrToolCallNotFound { - t.Fatalf("释放后应找不到,实际: %v", err) - } -} - -func TestSharedRef_PackUnpack(t *testing.T) { - b := make([]byte, sharedRefSize) - ref := SharedRef{Offset: 1024, Length: 256, Generation: 7, Flags: 3} - packSharedRef(b, ref) - got := unpackSharedRef(b) - if got.Offset != ref.Offset || got.Length != ref.Length || got.Generation != ref.Generation || got.Flags != ref.Flags { - t.Fatalf("pack/unpack: got %+v, want %+v", got, ref) - } -} - -func TestToolCallRing_ConcurrentReserveUnique(t *testing.T) { - data := make([]byte, int(trlOffFrameBase+toolRingCap*toolFrameSize)) - if err := InitToolRing(data); err != nil { - t.Fatal(err) - } - ring, err := AttachToolRing(data) - if err != nil { - t.Fatal(err) - } - - indices := make(chan uint32, toolRingCap) - var wg sync.WaitGroup - for i := uint32(0); i < toolRingCap; i++ { - wg.Add(1) - go func() { - defer wg.Done() - idx, reserveErr := ring.Reserve() - if reserveErr != nil { - t.Errorf("Reserve: %v", reserveErr) - return - } - indices <- idx - }() - } - wg.Wait() - close(indices) - - seen := make(map[uint32]bool, toolRingCap) - for idx := range indices { - if seen[idx] { - t.Fatalf("并发 Reserve 重复分配帧 %d", idx) - } - seen[idx] = true - } - if len(seen) != int(toolRingCap) { - t.Fatalf("唯一帧数=%d,期望 %d", len(seen), toolRingCap) - } -} - -func TestAttachToolRingRejectsInvalidLayout(t *testing.T) { - data := make([]byte, trlOffFrameBase) - putU32(data[trlOffMagic:], toolRingMagic) - putU32(data[trlOffVersion:], toolRingVersion) - putU32(data[trlOffCap:], toolRingCap) - putU32(data[trlOffFrameSize:], toolFrameSize) - if _, err := AttachToolRing(data); err == nil { - t.Fatal("越界布局应被拒绝") - } -} - -func TestToolFrameLayoutAligned(t *testing.T) { - // 确保帧布局字段偏移与内存布局一致(安全断言) - var frame toolFrame - off := func(ptr *uint32) uint32 { - return uint32(uintptr(unsafe.Pointer(ptr)) - uintptr(unsafe.Pointer(&frame))) - } - if off(&frame.state) != toolFrameOffState { - t.Errorf("state offset: got %d, want %d", off(&frame.state), toolFrameOffState) - } -} diff --git a/internal/plugin/proc/unified.go b/internal/plugin/proc/unified.go index 0f3bf63..f9855ef 100644 --- a/internal/plugin/proc/unified.go +++ b/internal/plugin/proc/unified.go @@ -50,9 +50,9 @@ const ( // [24,28) ctxSize(StageContext segment 字节数) // [28,32) evtOff(EvtRing segment 偏移) // [32,36) evtSize(EvtRing segment 字节数) -// [36,48) reserved(未来 lane 偏移/大小) -// [48,56) reserved -// [56,64) reserved(对齐到 8 字节) +// [36,40) arenaOff(Exchange Arena 起点,8 字节对齐) +// [40,44) arenaCap(Exchange Arena 字节数) +// [44,64) reserved const ( sbOffMagic = 0 sbOffVersion = 4 @@ -62,17 +62,15 @@ const ( sbOffCtxSize = 24 sbOffEvtOff = 28 sbOffEvtSize = 32 - sbOffArenaOff = 36 // 动态 arena 起始偏移 - sbOffArenaCap = 40 // 动态 arena 总容量 - sbOffArenaUsed = 44 // 动态 arena 已用字节(原子) + sbOffArenaOff = 36 // Exchange Arena 起点 + sbOffArenaCap = 40 // Exchange Arena 容量 + sbOffReserved1 = 44 sbOffReserved5 = 48 sbOffReserved6 = 52 sbOffReserved7 = 56 sbOffReserved8 = 60 ) -const unifiedArenaSize = 256 * 1024 // 默认动态 arena 256KB - // SharedRef 是跨进程共享内存描述符,替代内联 JSON 数据。 // // 所有数据交换(工具调用参数/结果、Cleaner、输入/输出通道消息) @@ -93,6 +91,9 @@ func (r SharedRef) IsZero() bool { } // Slice 从共享内存中按 SharedRef 切片。data 必须是完整的 mmap 区域。 +// +// 注意:这是**无校验**的原始切片,仅供已知安全的路径使用。 +// 跨进程引用一律走 arenaRegion.Read(校验 generation / 槽号 / 占用状态)。 func (r SharedRef) Slice(data []byte) []byte { end := uint64(r.Offset) + uint64(r.Length) if r.IsZero() || end > uint64(len(data)) { @@ -102,6 +103,9 @@ func (r SharedRef) Slice(data []byte) []byte { } // unifiedRegion 统一共享内存区域的内核侧视图。 +// +// arenaOff/arenaCap 指向 Exchange Arena 段;槽池句柄由 Host 持有 +// (见 arena.go),unifiedRegion 只负责布局与 generation。 type unifiedRegion struct { data []byte @@ -110,9 +114,8 @@ type unifiedRegion struct { evtOff uint32 evtSize uint32 - arenaOff uint32 // 动态 arena 起始偏移(相对 data) - arenaCap uint32 // 动态 arena 总容量 - arenaUsed *uint32 // 指向 SuperBlock 中的 arenaUsed 字段(原子 bump 游标) + arenaOff uint32 // Exchange Arena 起点(8 字节对齐) + arenaCap uint32 // Exchange Arena 容量 } // initUnifiedRegion 在 mmap 区域上初始化 SuperBlock + 两个 segment。 @@ -125,7 +128,8 @@ func initUnifiedRegion(data []byte, ctxTotal, evtTotal int) (*unifiedRegion, err ctxOff := uint32(superBlockSize) evtOff := ctxOff + uint32(ctxTotal) - arenaOff := evtOff + uint32(evtTotal) + // arena 基址必须 8 字节对齐:位图用 4 字节原子操作,槽头含 uint32。 + arenaOff := (evtOff + uint32(evtTotal) + 7) &^ 7 arenaCap := cap - arenaOff putU32(data[sbOffMagic:], unifiedMagic) @@ -138,7 +142,6 @@ func initUnifiedRegion(data []byte, ctxTotal, evtTotal int) (*unifiedRegion, err putU32(data[sbOffEvtSize:], uint32(evtTotal)) putU32(data[sbOffArenaOff:], arenaOff) putU32(data[sbOffArenaCap:], arenaCap) - putU32(data[sbOffArenaUsed:], 1) return &unifiedRegion{ data: data, @@ -148,9 +151,6 @@ func initUnifiedRegion(data []byte, ctxTotal, evtTotal int) (*unifiedRegion, err evtSize: uint32(evtTotal), arenaOff: arenaOff, arenaCap: arenaCap, - arenaUsed: (*uint32)(unsafe.Pointer( - uintptr(unsafe.Pointer(&data[0])) + uintptr(sbOffArenaUsed), - )), }, nil } @@ -197,9 +197,6 @@ func attachUnifiedRegion(data []byte) (*unifiedRegion, error) { evtSize: evtSize, arenaOff: arenaOff, arenaCap: arenaCap, - arenaUsed: (*uint32)(unsafe.Pointer( - uintptr(unsafe.Pointer(&data[0])) + uintptr(sbOffArenaUsed), - )), }, nil } @@ -253,56 +250,7 @@ func getU64(b []byte) uint64 { uint64(b[4])<<32 | uint64(b[5])<<40 | uint64(b[6])<<48 | uint64(b[7])<<56 } -// arenaAlloc 在动态 arena 中分配 n 字节,返回相对 arenaOff 的偏移。 -// append-only,offset 0 保留给“空”语义。 -func (r *unifiedRegion) arenaAlloc(n int) (uint32, error) { - if n <= 0 { - return 0, fmt.Errorf("arena: 非法分配长度 %d", n) - } - for { - used := atomic.LoadUint32(r.arenaUsed) - if used == 0 { - used = 1 - } - end := uint64(used) + uint64(n) - if end > uint64(r.arenaCap) { - return 0, fmt.Errorf("arena: 空间不足(需 %d,剩 %d)", n, uint64(r.arenaCap)-uint64(used)) - } - if atomic.CompareAndSwapUint32(r.arenaUsed, used, uint32(end)) { - return used, nil - } - } -} - -// arenaWrite 把 b 写入动态 arena 并返回 SharedRef。 -func (r *unifiedRegion) arenaWrite(b []byte) (SharedRef, error) { - if len(b) == 0 { - return SharedRef{}, nil - } - off, err := r.arenaAlloc(len(b)) - if err != nil { - return SharedRef{}, err - } - base := r.arenaOff + off - copy(r.data[base:base+uint32(len(b))], b) - gen := uint32(r.generation()) - return SharedRef{Offset: base, Length: uint32(len(b)), Generation: gen}, nil -} - -// arenaRead 按 SharedRef 读取数据。generation 不匹配或引用越出 arena 时拒绝。 -func (r *unifiedRegion) arenaRead(ref SharedRef) []byte { - if ref.IsZero() || ref.Generation != uint32(r.generation()) { - return nil - } - end := uint64(ref.Offset) + uint64(ref.Length) - arenaEnd := uint64(r.arenaOff) + uint64(r.arenaCap) - if uint64(ref.Offset) < uint64(r.arenaOff)+1 || end > arenaEnd { - return nil - } - return ref.Slice(r.data) -} - -// arenaReset 压实后重置游标(仅在无活跃 slot 时调用)。 -func (r *unifiedRegion) arenaReset() { - atomic.StoreUint32(r.arenaUsed, 1) -} +// arenaReset 已移除。 +// +// 旧实现把共享游标重置为 1,会覆盖其他进程/goroutine 仍在使用的分配。 +// 现设计用槽位图:释放是精确的 per-slot CAS,不存在"整体重置"语义。 diff --git a/plan.md b/plan.md index b2d1a15..48f7f2e 100644 --- a/plan.md +++ b/plan.md @@ -1147,41 +1147,56 @@ 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 共享槽池(内核独占所有权) -**目标**:ftruncate 扩展 + mremap/mmap 重映射 + generation 感知。 +**设计约束(用户明确)**: + +> 内核应当全权管理共享内存,插件需要共享内存要向内核申请,内核给插件返回偏移与大小,使用完成后插件通知内核回收。 +> 内核暴露类似 syscall 的 RPC 接口;共享内存是内部实现,不对插件开发者暴露。 + +**为什么不是跨进程分配器**(前两版都被推翻): + +- v1:SuperBlock 放 `arenaUsed` 游标,内核 CAS bump。但**插件模板里的 `arenaUsed` 是进程本地变量**,两个进程各自 bump,必然写到同一段内存;`arenaReset` 还会重置共享游标覆盖对方数据。 +- v2:把位图 CAS 下沉到插件模板。虽然正确,但把分配器实现细节泄漏进了插件运行时,且插件必须与内核保持位图布局同步。 +- v3(当前):分配器只存在于**内核进程内**,一把 `sync.Mutex` 即可。插件只通过 RPC 申请/归还,不做任何分配决策。 **实施**: -1. SuperBlock 加 generation + capacity 字段 -2. Host.Grow(newSize):ftruncate → 更新 capacity/gen → 通知子进程 -3. 子进程下次 stage 读 gen → remap -4. Host.Shrink():compact 后 free > 50% 且持续 5min → 缩容 -5. 无活跃 slot 时才允许缩容 +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**(原有正确路径),不静默截断。 **验证**: -- [ ] Grow 后子进程正确读写 -- [ ] Shrink 不裁掉活跃 slot -- [ ] git commit -m "feat(shm): grow/shrink with generation awareness" +- [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" -### 13.3 ToolCall lane +### 13.3 ToolCall 数据面走共享槽 -**目标**:工具调用参数/结果从 JSON RPC 改为 SharedRef。 +**目标**:工具调用的**参数/结果**走共享内存,RPC 只传描述符。 -**实施**: +**控制面 vs 数据面**(与初版设计的差异): -1. ToolCall Slot 状态机:FREE→WRITING→READY→READING→DONE→FREE -2. 内核写请求 → READY → eventfd 通知 -3. 插件读请求 → 执行 → 写结果 → DONE → eventfd 通知 -4. RPC tool.invoke 改为只传 {requestID, argsRef, resultRef} -5. has_cleaner 保持不变 +初版计划用 ring 状态机(FREE→WRITING→READY→READING→DONE)把 `tool.invoke` 也搬出 RPC。实际实现后发现这是错的方向: + +- RPC 已经提供请求 ID 关联、错误传递、ctx 取消、崩溃唤醒(`Process.CallContext`),ring 状态机是把这些重新实现一遍。 +- ring 版本写完从未接线,属死代码,已删除。 + +因此**控制面保留 RPC**,只把 payload 搬进共享槽: + +1. 小 payload 走内联 JSON(省两次 RPC)。 +2. 大 payload 走内核分配的槽,RPC 只带 `SharedRef`。 +3. 插件不需要分配:内核可预分配响应槽(如 Cleaner 路径)。 **验证**: -- [ ] TestPlugin_RegisteredToolInvokable 通过 -- [ ] Bench: ToolInvoke 延迟改善 -- [ ] git commit -m "feat(shm): toolcall lane with slot state machine" +- [x] `TestPlugin_RegisteredToolInvokable` 通过 +- [ ] Bench: ToolInvoke 延迟改善(待 `tool.invoke` 接入 SharedRef 后再测) ### 13.4 Cleaner 迁移至 SharedRef