Documentation
¶
Overview ¶
Package spmc implements a cache-line-aware, single-writer multiple-reader ring buffer inspired by the LMAX Disruptor pattern.
Key features:
- Lock-free reader lifecycle management via CAS bitmap
- Configurable wait strategies (Spin, Yield, Sleep, Hybrid)
- Batch publish primitives (Reserve/Commit, PublishBatch)
- Non-blocking variants (TryPublish, TryReserve) and a Remaining capacity signal
- Pipeline staging (Stage, SetGatingStage)
- False-sharing elimination through cache-line padding
Quick start:
rbuf, _ := spmc.NewRingBuffer[int](ctx, 1024)
stage := rbuf.NewStage(nil) // gated by the writer cursor
rbuf.SetGatingStage(stage) // writer waits for this stage's slowest reader
slotID, _ := stage.AddReader(func(ctx context.Context, rv spmc.ReadView[int], cur *atomic.Int64) {
// read loop; see ExampleRingBuffer
})
rbuf.Publish(42)
Index ¶
- Constants
- type BackoffParams
- type BitmapReaderPool
- type MinimumBarrier
- type ReadView
- type ReaderBarrier
- type ReaderFunc
- type RingBuffer
- func (r *RingBuffer[T]) Commit(claim int64)
- func (r *RingBuffer[T]) NewStage(upstream WriterBarrier) *Stage[T]
- func (r *RingBuffer[T]) Publish(payload T)
- func (r *RingBuffer[T]) PublishBatch(payloads []T)
- func (r *RingBuffer[T]) PublishBatchFunc(n int64, f func(i int64, slot *T))
- func (r *RingBuffer[T]) PublishFunc(f func(*T))
- func (r *RingBuffer[T]) Remaining() int64
- func (r *RingBuffer[T]) Reserve(n int64) (seg1, seg2 []T, claim int64)
- func (r *RingBuffer[T]) SetGatingStage(leaf *Stage[T])
- func (r *RingBuffer[T]) TryPublish(payload T) bool
- func (r *RingBuffer[T]) TryPublishFunc(f func(*T)) bool
- func (r *RingBuffer[T]) TryReserve(n int64) (seg1, seg2 []T, claim int64, ok bool)
- type RingBufferOption
- func WithBackoffParams[T any](params BackoffParams) RingBufferOption[T]
- func WithDefaultBackoffParams[T any]() RingBufferOption[T]
- func WithDefaultWaitStrategy[T any]() RingBufferOption[T]
- func WithDefaultWaitStrategyParams[T any]() RingBufferOption[T]
- func WithWaitStrategy[T any](strategy WaitStrategy) RingBufferOption[T]
- func WithWaitStrategyParams[T any](params WaitStrategyParams) RingBufferOption[T]
- type Stage
- type WaitStrategy
- type WaitStrategyParams
- type WriterBarrier
Examples ¶
Constants ¶
const CacheLineSize = 64
Variables ¶
This section is empty.
Functions ¶
This section is empty.
Types ¶
type BackoffParams ¶
type BitmapReaderPool ¶
type BitmapReaderPool[T any] struct { // contains filtered or unexported fields }
BitmapReaderPool supports up to 128 dynamic readers with zero false sharing
func NewBitmapReaderPool ¶
func NewBitmapReaderPool[T any](ctx context.Context, writerCursor WriterBarrier) *BitmapReaderPool[T]
func (*BitmapReaderPool[T]) AddReader ¶
func (b *BitmapReaderPool[T]) AddReader(fn ReaderFunc[T]) (int, error)
func (*BitmapReaderPool[T]) Load ¶
func (b *BitmapReaderPool[T]) Load() int64
Load returns minimum cursor across all active readers
Called by writer in reserve() loop (Hot Path)
- Optimized with "Gating" strategy to avoid scanning all readers on every call.
- Fast path: Only cache line 0 (bitmap + cache in same line)
- Slow path: cache line 0 + cursor cache lines (2-129)
- Never touches control cache lines (130-257)
- Zero false sharing with reader updates
func (*BitmapReaderPool[T]) RemoveReader ¶
func (b *BitmapReaderPool[T]) RemoveReader(slotId int) error
func (*BitmapReaderPool[T]) Shutdown ¶
func (b *BitmapReaderPool[T]) Shutdown()
Shutdown stops all internal goroutines and waits for them to exit.
type MinimumBarrier ¶
type MinimumBarrier []ReaderBarrier
func (MinimumBarrier) Load ¶
func (m MinimumBarrier) Load() int64
Load implements a branch-free min comparison to find the slowest reader's position
Example ¶
Let's say we have 3 readers at different positions:
Reader 0: sequence = 100 Reader 1: sequence = 95 (slowest) Reader 2: sequence = 102
Iteration 1: i=0 minimum = 100 (initialized from m[0])
Iteration 2: i=1
seq = 95 diff = minimum - seq = 100 - 95 = 5
## Binary representation (int64):
diff = 0000...0101 (positive number)
## Arithmetic right shift by 63:
mask = diff >> 63 = 0000...0000 = 0
## Update minimum:
minimum = seq + (diff & mask) = 95 + (5 & 0) = 95 + 0 = 95 ✓ (updated to smaller cursor)
Iteration 3: i=2
seq = 102 diff = minimum - seq = 95 - 102 = -7
## Binary representation (int64 two's complement):
diff = 1111...1001 (negative number, sign bit = 1)
## Arithmetic right shift by 63:
mask = diff >> 63 = 1111...1111 = -1 (all bits set)
## Update minimum:
minimum = seq + (diff & mask) = 102 + (-7 & -1) = 102 + (-7) = 95 ✓ (kept smaller cursor)
type ReadView ¶
type ReadView[T any] struct { // contains filtered or unexported fields }
func (*ReadView[T]) GetSegments ¶
GetSegments returns zero-copy views of the backing buffer spanning sequences start..end inclusive. seg2 is non-nil only when the range wraps the ring end. Sequences are absolute, as accepted by Get; masking is internal. The caller must consume both segments before advancing its cursor and must not retain them afterward.
Contract (not validated for max performance): start <= end, end-start < bufferSize, and end at or below the barrier the reader observed. The gating protocol guarantees any range inside [cursor+1, LoadWriterBarrier()] satisfies all three.
func (*ReadView[T]) Iterate ¶
Iterate yields pointers to sequences start..end inclusive, in order. Sequences are absolute, as accepted by Get and GetSegments. Implemented over GetSegments; when the compiler inlines the loop body it measures at parity with looping the segments directly (docs/PERFORMANCE.md, batch read paths). Prefer GetSegments when the slices themselves are needed, for bulk copies or vector processing.
func (*ReadView[T]) LoadWriterBarrier ¶
type ReaderBarrier ¶
type ReaderBarrier interface {
Load() int64
}
type ReaderFunc ¶
type RingBuffer ¶
type RingBuffer[T any] struct { // contains filtered or unexported fields }
Example ¶
This example demonstrates a basic publish-and-read cycle using a pipeline stage. A single reader goroutine processes values published by the main goroutine.
package main
import (
"context"
"fmt"
"sync/atomic"
"github.com/pintomau/go-spmc-ring"
)
func main() {
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
rb, err := spmc.NewRingBuffer[int](ctx, 64)
if err != nil {
fmt.Println("error creating ring buffer:", err)
return
}
// Create a stage and set it as the gating barrier.
s := rb.NewStage(nil)
rb.SetGatingStage(s)
// Channel to collect read values so the example can print them.
read := make(chan int, 10)
// Register a reader that collects published values.
slotID, err := s.AddReader(func(ctx context.Context, rv spmc.ReadView[int], cur *atomic.Int64) {
expected := cur.Load() + 1
for {
select {
case <-ctx.Done():
return
default:
if w := rv.LoadWriterBarrier(); expected <= w {
for seq := expected; seq <= w; seq++ {
read <- *rv.Get(seq)
}
cur.Store(w)
expected = w + 1
}
}
}
})
if err != nil {
fmt.Println("error adding reader:", err)
return
}
// Publish some values.
for _, v := range []int{10, 20, 30} {
rb.Publish(v)
}
// Collect the values the reader picked up.
for i := 0; i < 3; i++ {
fmt.Println("Read:", <-read)
}
// Clean up.
_ = s.RemoveReader(slotID)
s.Shutdown()
}
Output: Read: 10 Read: 20 Read: 30
Example (Batch) ¶
This example demonstrates batch I/O: Reserve/Commit on the write side and GetSegments on the read side. Reserve allocates contiguous slots (up to two segments when wrapping). After filling them, a single Commit makes the whole batch visible atomically. The reader then uses GetSegments to obtain zero-copy views of the available range and processes all events at once.
package main
import (
"context"
"fmt"
"sync/atomic"
"github.com/pintomau/go-spmc-ring"
)
func main() {
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
rb, err := spmc.NewRingBuffer[int](ctx, 8)
if err != nil {
fmt.Println("error creating ring buffer:", err)
return
}
s := rb.NewStage(nil)
rb.SetGatingStage(s)
stall := make(chan struct{}, 1)
// Channel to collect batch results from the reader.
batchSums := make(chan int, 10)
// Register a reader that processes events in batches using GetSegments.
slotID, err := s.AddReader(func(ctx context.Context, rv spmc.ReadView[int], cur *atomic.Int64) {
<-stall
expected := cur.Load() + 1
rv.Get(expected)
expected++
cur.Store(expected)
<-stall
for {
select {
case <-ctx.Done():
return
default:
if w := rv.LoadWriterBarrier(); expected <= w {
// GetSegments returns zero-copy views of the available
// range. seg2 is non-nil only when the range wraps the
// ring end.
seg1, seg2 := rv.GetSegments(expected, w)
// Process all events in the batch.
sum := 0
for i := range seg1 {
sum += seg1[i]
}
for i := range seg2 {
sum += seg2[i]
}
batchSums <- sum
cur.Store(w)
expected = w + 1
}
}
}
})
if err != nil {
fmt.Println("error adding reader:", err)
return
}
// force a wrapping on reserve
rb.Publish(1000)
stall <- struct{}{}
// Batch write: Reserve allocates 7 slots, returning up to two segments
seg1, seg2, claim := rb.Reserve(7)
// Fill both segments. seg1 contains slots up to the ring end, seg2
val := 100
for i := range seg1 {
seg1[i] = val
val++
}
for i := range seg2 {
seg2[i] = val
val++
}
// Commit makes the entire batch visible atomically.
rb.Commit(claim)
stall <- struct{}{}
// Collect the batch sum from the reader.
fmt.Println("Batch sum:", <-batchSums)
// Clean up.
_ = s.RemoveReader(slotID)
s.Shutdown()
close(stall)
}
Output: Batch sum: 721
Example (SelfEvictingReader) ¶
This example shows a reader that evicts itself when it falls too far behind, instead of stalling the writer forever. Returning from a ReaderFunc is a first-class way to leave the pool: the slot is deactivated and freed, and the writer is ungated. The reader picks its own exit point, so it is never mid-read when the writer reclaims its slots.
package main
import (
"context"
"fmt"
"sync/atomic"
"github.com/pintomau/go-spmc-ring"
)
func main() {
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
rb, err := spmc.NewRingBuffer[int](ctx, 8)
if err != nil {
fmt.Println("error creating ring buffer:", err)
return
}
s := rb.NewStage(nil)
rb.SetGatingStage(s)
evicted := make(chan struct{})
_, err = s.AddReader(func(ctx context.Context, rv spmc.ReadView[int], cur *atomic.Int64) {
defer close(evicted)
const maxLag = 4
for {
select {
case <-ctx.Done():
return
default:
if rv.LoadWriterBarrier()-cur.Load() > maxLag {
return // self-evict: slot is freed, writer is ungated
}
}
}
})
if err != nil {
fmt.Println("error adding reader:", err)
return
}
// 100 events exceed the capacity-8 ring many times over. The blocking
// Publish calls can only complete because the laggard removes itself.
for i := 1; i <= 100; i++ {
rb.Publish(i)
}
<-evicted
fmt.Println("published 100 events past a stalled reader")
s.Shutdown()
}
Output: published 100 events past a stalled reader
func NewRingBuffer ¶
func NewRingBuffer[T any](ctx context.Context, capacity int64, opts ...RingBufferOption[T]) (*RingBuffer[T], error)
func (*RingBuffer[T]) Commit ¶
func (r *RingBuffer[T]) Commit(claim int64)
Commit publishes a batch previously returned by Reserve. It must be called exactly once per Reserve, with the claim returned by that Reserve, and in the same order as the Reserve calls. A single atomic Store makes the entire batch visible to readers.
func (*RingBuffer[T]) NewStage ¶
func (r *RingBuffer[T]) NewStage(upstream WriterBarrier) *Stage[T]
NewStage creates a pipeline stage gated by upstream. Pass nil to gate on the writer cursor. Wire a two-stage pipeline:
s1 := rb.NewStage(nil) s2 := rb.NewStage(s1.Barrier()) rb.SetGatingStage(s2)
func (*RingBuffer[T]) Publish ¶
func (r *RingBuffer[T]) Publish(payload T)
func (*RingBuffer[T]) PublishBatch ¶
func (r *RingBuffer[T]) PublishBatch(payloads []T)
PublishBatch copies payloads into len(payloads) consecutive slots and publishes them atomically. Empty payloads is a no-op.
func (*RingBuffer[T]) PublishBatchFunc ¶
func (r *RingBuffer[T]) PublishBatchFunc(n int64, f func(i int64, slot *T))
PublishBatchFunc reserves n slots and invokes f for each in publish order, passing a pointer into the backing buffer. f must not retain the pointer past its call. Wrap-around is handled internally; f sees indices 0..n-1.
func (*RingBuffer[T]) PublishFunc ¶
func (r *RingBuffer[T]) PublishFunc(f func(*T))
func (*RingBuffer[T]) Remaining ¶
func (r *RingBuffer[T]) Remaining() int64
Remaining reports how many slots the writer can publish right now without blocking. It refreshes the gating barrier, so it reports actual capacity rather than the writer's cached view; the result is 0 exactly when TryPublish would return false. The maximum value is bufferSize-1 (strict < gating, consistent with Reserve's n < bufferSize rule). It does not account for an uncommitted Reserve claim, so call it between Commit and the next Reserve. Must only be called from the writer goroutine.
func (*RingBuffer[T]) Reserve ¶
func (r *RingBuffer[T]) Reserve(n int64) (seg1, seg2 []T, claim int64)
Reserve blocks until n contiguous sequence numbers are available and returns up to two slices into the backing buffer spanning those slots. seg2 is non-nil only when the reservation wraps the end of the ring. The caller must fill every slot in both segments and then pass the returned claim to Commit exactly once, in publish order. n must satisfy 0 < n < bufferSize.
Each segment is contiguous in memory, so callers may fill them with bulk techniques such as the GOEXPERIMENT=simd package's vector stores. No alignment beyond the element type's natural alignment is guaranteed.
func (*RingBuffer[T]) SetGatingStage ¶
func (r *RingBuffer[T]) SetGatingStage(leaf *Stage[T])
SetGatingStage points the writer's backpressure at the leaf stage. Only accepts *Stage[T] (not a bare WriterBarrier) to guarantee the cached Load() path is used.
func (*RingBuffer[T]) TryPublish ¶
func (r *RingBuffer[T]) TryPublish(payload T) bool
TryPublish attempts to publish without blocking. It returns false when the ring is full, after refreshing the gating barrier once. Like Publish, it must only be called from the writer goroutine.
func (*RingBuffer[T]) TryPublishFunc ¶
func (r *RingBuffer[T]) TryPublishFunc(f func(*T)) bool
TryPublishFunc is the non-blocking sibling of PublishFunc. It returns false when the ring is full, in which case f is not called. Like PublishFunc, it must only be called from the writer goroutine.
func (*RingBuffer[T]) TryReserve ¶
func (r *RingBuffer[T]) TryReserve(n int64) (seg1, seg2 []T, claim int64, ok bool)
TryReserve is the non-blocking sibling of Reserve. When n contiguous slots are free it behaves exactly like Reserve and returns ok = true; the caller must then pass claim to Commit exactly once, in publish order. When the ring lacks room for the whole batch it returns ok = false (after refreshing the gating barrier once) and no claim is made. Panics if n is out of range, matching Reserve. Must only be called from the writer goroutine.
type RingBufferOption ¶
type RingBufferOption[T any] func(*RingBuffer[T])
RingBufferOption configures a RingBuffer
func WithBackoffParams ¶
func WithBackoffParams[T any](params BackoffParams) RingBufferOption[T]
func WithDefaultBackoffParams ¶
func WithDefaultBackoffParams[T any]() RingBufferOption[T]
func WithDefaultWaitStrategy ¶
func WithDefaultWaitStrategy[T any]() RingBufferOption[T]
func WithDefaultWaitStrategyParams ¶
func WithDefaultWaitStrategyParams[T any]() RingBufferOption[T]
func WithWaitStrategy ¶
func WithWaitStrategy[T any](strategy WaitStrategy) RingBufferOption[T]
WithWaitStrategy sets the wait strategy for writer backpressure
func WithWaitStrategyParams ¶
func WithWaitStrategyParams[T any](params WaitStrategyParams) RingBufferOption[T]
WithWaitStrategyParams sets additional parameters for the wait strategy
type Stage ¶
type Stage[T any] struct { // contains filtered or unexported fields }
Stage is a named pipeline stage. Each stage owns a BitmapReaderPool and exposes:
- Barrier(): a concurrent-safe WriterBarrier for wiring a downstream stage
- Load(): the cached WriterBarrier the writer uses for gating (single-threaded)
func (*Stage[T]) Barrier ¶
func (s *Stage[T]) Barrier() WriterBarrier
Barrier returns a concurrent-safe, stateless WriterBarrier over this stage's minimum cursor. Pass this to rb.NewStage() as the upstream argument when wiring a downstream stage. Never pass to SetGatingStage; that path requires the cached Load() for performance.
func (*Stage[T]) Load ¶
Load implements WriterBarrier using the cached scan path. Called only by the single writer goroutine via SetGatingStage.
func (*Stage[T]) RemoveReader ¶
type WaitStrategy ¶
type WaitStrategy uint8
WaitStrategy defines the wait strategy for writer backpressure. Using a config enum + switch is usually faster than interface/generic dispatch in tight spin loops due to Go's interface/generic dictionary overhead. See docs/PERFORMANCE.md for the latency matrix.
const ( // WaitStrategySpin performs pure busy-spinning. // Lowest latency, highest CPU usage. Use for ultra-low-latency scenarios // where backpressure is rare and brief. WaitStrategySpin WaitStrategy = iota // WaitStrategyYield calls runtime.Gosched() to yield to other goroutines. // Conservative default. For latency-sensitive workloads with 1-4 readers, // WaitStrategySpin is faster. For bursty workloads on arm64, consider // WaitStrategyHybrid. WaitStrategyYield // WaitStrategySleep sleeps for a configurable duration. // Lowest CPU usage, higher latency. Good for throughput-oriented scenarios. WaitStrategySleep // WaitStrategyHybrid starts with yielding and backs off to sleeping // under sustained contention. Best for variable workloads. WaitStrategyHybrid )
func (WaitStrategy) String ¶
func (w WaitStrategy) String() string
String returns the name of the wait strategy
type WaitStrategyParams ¶
type WaitStrategyParams struct {
// SleepDuration for WaitStrategySleep (default: 1μs)
SleepDuration time.Duration // 8 bytes
// HybridSpinCount before backing off to sleep (default: 100)
HybridSpinCount uint64 // 8 bytes
// HybridSleepDuration after spin count exceeded (default: 1μs)
HybridSleepDuration time.Duration // 8 bytes
}
WaitStrategyParams holds optional parameters for wait strategies
type WriterBarrier ¶
type WriterBarrier interface {
Load() int64
}