111 lines
4.2 KiB
Go
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()
|
||
|
|
}
|