Files
petrbalvin af4ee19703
Release / gates (push) Successful in 4m38s
Test / test (push) Successful in 5m16s
Release / release (push) Successful in 35s
feat: initial release
Assisted-by: GLM 5.3 Flash
2026-09-03 10:00:00 +02:00

111 lines
4.2 KiB
Go

// Copyright (c) 2026 Petr Balvín <opensource@petrbalvin.org> (https://petrbalvin.org)
// SPDX-License-Identifier: MIT
package spmd
import "sync"
// The routed frame pool recycles the payload buffers of the frames the
// hub routes on, the frames whose destination is another rank. Its
// safety rests on one ownership chain and nothing else: the pump takes
// a buffer when it reads a routed frame, the frame hands the buffer to
// the destination link's outbox channel, and the drain, the chain's
// single consumer, writes the frame and returns the buffer at once,
// on a landed write and on a benign skip at a departing peer alike.
// No frame travels a pending stash with the pool's mark on, no
// collective ever sees a pooled buffer, and no count, atom or second
// return point exists: a channel handoff is the ownership transfer.
// A buffer the chain drops on the way, a mislabelled frame or a world
// that ends mid-route, is left to the collector, which makes the loss
// a pool miss and never a double return. This single owner is why the
// first pool attempt's failure mode, a payload retired while a decode
// still read it, cannot arise here.
//
// Two bounds keep the pool from pinning memory. A buffer above the
// retention ceiling is never kept, so one enormous frame cannot pin
// the heap across rounds, and the pool holds at most a fixed number
// of buffers, so a flood of mid-sized ones cannot grow without end.
// Both refusals are pool misses: the buffer goes to the collector,
// which is exactly the behaviour the package had before the pool.
const (
// framePoolCeiling is the largest buffer the pool retains, one
// mebibyte. A payload beyond it is a pool miss by rule.
framePoolCeiling = int64(1 << 20)
// framePoolBound is the largest number of buffers the pool holds
// at once; a return that would exceed it is dropped.
framePoolBound = 32
)
// framePoolEnabled switches the pool at package level. It exists for
// the A/B measurement, which flips it between alternating
// sub-benchmarks of one binary; the library itself always runs pooled.
var framePoolEnabled = true
// framePool is the free list behind the pool: payload buffers waiting
// for the next routed frame whose length fits.
type framePool struct {
mu sync.Mutex
free [][]byte
}
// routedFrames is the process's one pool. The worlds of one process
// share it, which is what makes its bounds process-wide facts.
var routedFrames framePool
// takeFrameBuffer answers the buffer a routed frame's payload reads
// into, or nil when the frame allocates as usual, outside the pool's
// chain: when the pool is switched off, when the frame stays at the
// hub, or when the length the header announced is beyond the ceiling.
// A qualified frame rides the chain whatever the free list holds: the
// buffer comes from the pool when one fits, and a fresh one joins the
// chain in its place, so the drain has a buffer to return and the
// pool warms with the traffic it serves. The link read calls this
// with the header already parsed, before the payload moves.
func takeFrameBuffer(m message, length int64) []byte {
if !framePoolEnabled || m.dest == 0 || length <= 0 || length > framePoolCeiling {
return nil
}
if buf := routedFrames.take(length); buf != nil {
return buf
}
return make([]byte, length)
}
// take answers a buffer of exactly length bytes with at least that
// capacity, the smallest retained buffer that fits, or nil when the
// pool holds none.
func (p *framePool) take(length int64) []byte {
if length <= 0 || length > framePoolCeiling {
return nil
}
p.mu.Lock()
defer p.mu.Unlock()
best := -1
for i, buf := range p.free {
if int64(cap(buf)) >= length && (best < 0 || cap(buf) < cap(p.free[best])) {
best = i
}
}
if best < 0 {
return nil
}
buf := p.free[best]
p.free = append(p.free[:best], p.free[best+1:]...)
return buf[:length]
}
// retire returns one buffer to the pool under the two bounds: a
// buffer whose capacity exceeds the ceiling, and a return that would
// push the pool past its count bound, are dropped to the collector.
func (p *framePool) retire(data []byte) {
if !framePoolEnabled || cap(data) == 0 || int64(cap(data)) > framePoolCeiling {
return
}
p.mu.Lock()
if len(p.free) < framePoolBound {
p.free = append(p.free, data[:cap(data)])
}
p.mu.Unlock()
}