spmc

package module
v0.1.2 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Jun 29, 2026 License: MIT Imports: 9 Imported by: 0

README

go-spmc-ring

Go Reference Latest Release Go Version Go Report Card codecov

A single-writer, multiple-reader ring buffer in Go, inspired by the LMAX Disruptor pattern. The codebase focuses on cache-line awareness, false-sharing avoidance, and lock-free reader lifecycle management: unlike the classic Disruptor, where the consumer graph is wired once before the first event, readers here join and leave a live stream at runtime without locks, allocations, or writer stalls.

On a Ryzen 5 9600X (Linux, Go 1.26): a single publish takes ~5.2 ns (≈193 M ops/s), batched publishes reach ~2.1 ns/item (≈476 M items/s), and eight concurrent readers cost the writer less than 7%. On an Apple M4 Pro (macOS) the same publish is ~7.3 ns and eight readers cost under 2%. Full measurements for both architectures are in docs/PERFORMANCE.md.

Why it's fast

  • Single-writer contract. With exactly one publishing goroutine, sequence bookkeeping needs no atomics. The entire hot path synchronizes through one writeCursor.Store per publish (or per batch).
  • Cache-line padding everywhere it matters. The writer's hot fields and every reader's cursor sit on their own cache lines, so cores never invalidate each other's lines by accident. The false-sharing benchmark family measures exactly this.
  • Lock-free reader lifecycle. Readers claim slots by CAS on a 2×uint64 bitmap: add and remove are wait-free for the writer, and reader goroutines are pooled and reused across registrations.
  • Cached gating. The writer caches the slowest reader's position and only rescans all cursors when that cache is exhausted or the reader bitmap changes, keeping the common case scan-free.
  • Batch primitives. Reserve/Commit exposes the ring's backing array directly and makes an entire batch visible with a single atomic store, amortizing the fixed publish overhead to ~2.1 ns/item.
  • Devirtualized waiting. Wait strategies are a uint8 switch rather than an interface, so the backpressure loop pays no dynamic-dispatch cost.

Install

go get github.com/pintomau/go-spmc-ring

Requires Go 1.24+. No third-party dependencies and no cgo.

Quick start

ctx, cancel := context.WithCancel(context.Background())
defer cancel()

rb, _ := spmc.NewRingBuffer[int](ctx, 1024) // capacity rounds up to a power of 2

// Readers live in a Stage. Gate the writer on the stage so it
// never laps the slowest reader.
s := rb.NewStage(nil) // nil upstream = gated by the writer cursor
rb.SetGatingStage(s)  // must happen before the first Publish

slotID, _ := 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++ {
					process(*rv.Get(seq))
				}
				cur.Store(w) // the reader owns its cursor: advance it or the writer stalls
				expected = w + 1
			} else {
				time.Sleep(50 * time.Microsecond) // idle backoff (see "Writing a reader")
			}
		}
	}
})

for i := range 1_000_000 {
	rb.Publish(i)
}

s.RemoveReader(slotID)
s.Shutdown()

Architecture

                 ┌─────────────────────────────────────────────────────┐
 Writer          │ RingBuffer[T]                                       │
 (exactly one    │   buffer []T          power-of-2, mask-indexed      │
  goroutine)     │   writeCursor         atomic, one Store per publish │
   │             │   nextSequence        writer-private, no atomics    │
   ├─ Publish ──►│   cachedSlowestReader writer-private gate cache     │
   ├─ PublishBatch                                                     │
   └─ Reserve/Commit                     (hot fields cache-line padded)│
                 └──────────────┬──────────────────────────────────────┘
                                │ writeCursor visible to readers
                                ▼
                 ┌─────────────────────────────────────────────────────┐
                 │ Stage 1 (BitmapReaderPool)                          │
                 │   2×uint64 bitmap → up to 128 reader slots          │
                 │   one padded cursor per reader (no false sharing)   │
                 │   reader goroutines pooled & reused                 │
                 └──────────────┬──────────────────────────────────────┘
                                │ Barrier() = min(stage-1 cursors)
                                ▼
                 ┌─────────────────────────────────────────────────────┐
                 │ Stage 2 …N (optional pipeline stages)               │
                 │   readers only see events stage N−1 has committed   │
                 └──────────────┬──────────────────────────────────────┘
                                │ min cursor of the leaf stage
                                ▼
                  writer backpressure (SetGatingStage):
                  Publish blocks when the buffer would lap the
                  slowest reader of the gating stage

The writer is single-threaded by contract, so its sequence bookkeeping needs no atomics. One writeCursor.Store per publish (or per batch) is the only synchronization on the hot path. Readers register into a 128-slot bitmap pool. Each reader's cursor sits on its own cache line, and the writer only rescans all cursors when the bitmap topology changes or its cached gate is exhausted.

Usage patterns

Writing a reader

A ReaderFunc runs on its own pooled goroutine and owns its cursor: advance cur after consuming or the writer will eventually stall waiting for you. The canonical loop is in the quick start above. Two details that matter in practice:

  • See Throughput vs latency below.
  • Batch your cursor stores. Consume everything up to LoadWriterBarrier() and store the cursor once per batch, not once per event.
Adding and removing readers at runtime

Readers are not a startup-time decision. While the writer publishes at full rate, you can attach a metrics tap or debug subscriber to a live stream, run a canary consumer, or drain and detach a reader before reconfiguring:

tapID, _ := s.AddReader(metricsTap) // writer keeps publishing throughout

// ... observe for a while ...

s.RemoveReader(tapID)

What makes this safe and cheap:

  • AddReader is lock-free. The slot is claimed by a CAS on the reader bitmap. The writer is never blocked and simply picks up the new gating cursor on its next scan.
  • A new reader starts at the writer's current position. It sees only events published after it joined. There is no replay of history.
  • RemoveReader takes effect immediately from the writer's perspective. The slot leaves the gating set before the reader goroutine finishes winding down, so a slow reader that gets removed can never stall the writer again.
  • Goroutines are pooled. A removed reader's goroutine idles with exponential backoff and self-terminates after about a minute, so add/remove churn reuses live goroutines instead of paying spawn costs.

The dynamic part is the readers, not the topology: stages and the gating barrier must still be wired before the first publish (see Pipeline stages).

Consuming in batches

ReadView offers three access paths, and all of them take absolute sequence numbers while masking is internal. Get(seq) returns a single event. GetSegments(start, end) returns up to two zero-copy slices spanning start..end inclusive. The second is non-nil only when the range wraps the ring end, mirroring the two-segment shape Reserve gives the writer. Iterate(start, end) yields each event in order via an iterator.

if w := rv.LoadWriterBarrier(); expected <= w {
	seg1, seg2 := rv.GetSegments(expected, w) // zero-copy, even when the range wraps
	for i := range seg1 {
		process(seg1[i])
	}
	for i := range seg2 {
		process(seg2[i])
	}
	cur.Store(w)
	expected = w + 1
}

Both segments point straight into the ring, so the batch must be fully consumed before the cursor is advanced: after cur.Store(w) the writer may overwrite those slots. Iterate is built on GetSegments and measures at parity with direct segment loops when the body inlines. Reach for GetSegments when you need the slices themselves, such as bulk copies or vector processing. All read paths are allocation-free, wrapping included (see Batch read paths).

Publishing in batches

Single-item Publish pays a fixed ~3.1 ns of gating and cursor-store overhead per call. Batching amortizes it (~2.1 ns/item at batch ≥ 10, see Batch scaling).

// Copy a prepared slice:
rb.PublishBatch(events)

// Fill slots in place (no intermediate slice):
rb.PublishBatchFunc(n, func(i int64, slot *Event) { slot.ID = i })

// Zero-overhead variant that writes directly into the ring:
seg1, seg2, claim := rb.Reserve(n) // seg2 non-nil only when the batch wraps the ring end
fill(seg1)
fill(seg2)
rb.Commit(claim) // one atomic store makes the whole batch visible

Reserve requires 0 < n < bufferSize (a full-buffer reservation can never be satisfied and panics), and every Reserve must be paired with exactly one Commit, in order.

Non-blocking publishing

Publish blocks when the ring is full. When dropping or deferring work beats waiting on a slow reader, use the try variants. They return false instead of waiting:

if !rb.TryPublish(event) {
    // ring full: drop the event, buffer upstream, or push back on the source
}

if rb.Remaining() < lowWater {
    // adapt before hitting the wall: shed load or switch to batching
}

TryReserve is the batch sibling: it returns ok = false when the ring lacks room for the whole batch. The try variants share the writer's contract, so call them only from the writer goroutine.

A slow reader can also opt out from its side: returning from a ReaderFunc removes the reader from the pool and ungates the writer. The reader picks its own exit point, so it is never mid-read when its slots are reclaimed. See the self-evicting reader example in the package docs.

Pipeline stages

Stages chain reader groups so stage N only consumes events all of stage N−1 has committed, which is the LMAX barrier-group pattern. Example: journal every event before business logic is allowed to see it.

s1 := rb.NewStage(nil)          // gated by the writer cursor
s2 := rb.NewStage(s1.Barrier()) // gated by s1's slowest reader
rb.SetGatingStage(s2)           // writer waits for the leaf stage

s1.AddReader(journalFn)
s2.AddReader(businessFn)

// Each stage shuts down independently:
defer s1.Shutdown()
defer s2.Shutdown()

Constraints:

  • SetGatingStage must be called before the first publish, because the gating field is not synchronized.
  • Wire downstream stages with Barrier() (concurrent-safe), never with Load() (writer-only cache path).
  • Depth is effectively free: overhead tracks total reader count, not stage count (see Stage / pipeline scaling).
Wait strategies

The writer's backpressure behavior when the buffer is full is a WaitStrategy (a uint8 switch, not an interface, to keep the wait loop devirtualized):

Strategy Behavior
WaitStrategySpin busy-spin
WaitStrategyYield runtime.Gosched() (default)
WaitStrategySleep time.Sleep
WaitStrategyHybrid spin N times, then sleep
rb, _ := spmc.NewRingBuffer[Event](ctx, 1<<16,
	spmc.WithWaitStrategy[Event](spmc.WaitStrategySpin))

Which one to pick, and how the same choice plays out differently for writers and readers, is covered in Throughput vs latency below.

Things to know
  • One writer, always. All publish methods assume a single publishing goroutine. There is no multi-producer mode.
  • Capacity rounds up to the next power of two.
  • With no readers registered, the writer free-runs and never blocks.
  • Readers are limited to 128 per stage (the 2×64-bit bitmap). AddReader returns an error when the stage is full.

Performance

Headline numbers, arithmetic mean of 10 runs, on two machines: x86-64 (Ryzen 5 9600X, 6C/12T, Linux) and arm64 (Apple M4 Pro, 14C/14T, macOS), both Go 1.26.2:

Measurement x86-64 (Ryzen) arm64 (M4 Pro)
Single Publish 5.2 ns/op (193 M ops/s) 7.3 ns/op (138 M ops/s)
Batch publish, size ≥ 10 2.1–2.55 ns/item (all APIs converge) 1.76 ns/item floor (bulk PublishBatch, per-item cost is fill-pattern dependent)
8 concurrent readers +6.2% writer cost vs. 1 reader +1.5% vs. 1 reader
128 readers (capacity limit) 3.0× single-reader cost 2.4× single-reader cost
Pipeline depth (2–3 stages, ≤4 readers each) within noise of single-stage within noise of single-stage
Best end-to-end p99 (burst workload) 13.2µs (Yield wait, batch poll) 9.0µs (Hybrid wait, batch poll)

Absolute numbers capture the whole platform (CPU, scheduler, OS timer), not just the CPU. See the per-section commentary for where the two architectures diverge in shape (notably the Reserve vs PublishBatch ordering, the multi-reader cliff position, and the FixedRate Spin vs Batch poll recommendation). The full data, including multi-reader scaling, the LockOSThread study, batch-size sweeps, and the complete latency matrix with HDR percentiles, is in docs/PERFORMANCE.md.

Throughput vs latency: choosing spin, yield, and sleep

Three independent decisions hide behind "wait strategy", and the right answer is different for each. Two mechanisms drive all of it:

  1. time.Sleep has a wakeup floor. Asking Linux for a 1µs sleep takes 15–50µs in practice (Go runtime timer plus OS scheduler wakeup). Any sleep on an event-delivery path donates that floor to your latency percentiles.
  2. runtime.Gosched does not park. It yields the processor but re-enters the run queue immediately, so an idle loop built on it busy-spins through the scheduler, churning run queues and stealing cycles from the writer. Measured cost: −157% writer throughput with 8 idle-spinning readers.
Decision point Best choice Why
Writer backpressure (WithWaitStrategy) Yield (default), or Spin when latency-critical the buffer-full path is rare when the ring is sized right, so politeness is cheap. spin only buys the last microseconds
Reader polling while events flow drain to the barrier each pass, never sleep between polls every sleep adds the 15–50µs wakeup floor to every event that arrives during it
Reader idle after catching up time.Sleep(50µs), never Gosched sleep parks the goroutine and frees the scheduler. Gosched creates a scheduler storm

sleeping is the right way to idle (it parks the goroutine, protecting throughput) and the wrong way to poll (it pays the wakeup floor per event, destroying latency). Workload-specific recommendations, including how the picture changes at p99.9 and beyond, are in Tail-aware recommendations.

Other findings worth knowing
  • Batching pays for itself by size 10. The ~3.1 ns fixed cost per publish (gating check plus cursor store) amortizes away. On x86 the per-item cost is then flat at the payload-write floor. On arm64 it dips into a mid-size hump before reaching its lowest at large batches (see Batch scaling).
  • Reader scaling is sub-linear up to the hardware thread count, then degrades as the writer's full cursor scan starts to dominate.
  • Publishing is O(1) in reader count. A channel-based broadcast is O(readers). Fanning the same event out through one buffered channel per reader costs 6x the ring at a single reader and ~288x at 128, because every channel needs its own send while the ring's readers all watch one cursor (see Baseline: channel fan-out).
  • Batching closes most of the per-op gap, but not the allocation gap. A chan []T fan-out cuts per-item cost 17x (266 to 15 ns/item from batch size 1 to 1000), yet every batch still allocates ~64 B per item, while the ring publishes batches in place with zero allocations (see Batched insert and read).
  • runtime.LockOSThread does not pay. Go's scheduler beats manual OS-thread pinning for this workload at almost every reader count. The idle strategy matters far more than pinning.
  • Batch-draining readers get more valuable at higher percentiles. At p99.9+ they win in almost every measured scenario.

SIMD segment fills (experimental)

A GOEXPERIMENT=simd experiment (amd64, Go 1.26+, not part of the shipped library) found one pattern that pays: filling Reserve segments with vector-generated payloads runs about 3.5× a scalar PublishBatchFunc at batch ≥ 64 while the ring is cache-resident. Plain copies gain nothing (memmove is already vectorized) and the win disappears at DRAM bandwidth. Full tables, alignment findings, and the reader-side analysis: SIMD segment fills. Recipe in example_simd_test.go.

How it's tested

  • Property-based simulation (pgregory.net/rapid): randomized writers, batch sizes, and reader add/remove churn are checked against ordering and visibility invariants. The full suite runs in minutes locally. A nightly CI job runs each property at a check budget sized to its cost, from hundreds of cases for the heaviest property to tens of thousands for the cheapest.
  • Deterministic replay. When a simulation fails, rapid saves a failfile under testdata/rapid/ and the next test run replays it automatically. A specific run can be reproduced with RAPID_SEED=<seed> go test -run Simulation ..
  • Multi-platform CI. Build and short tests across Linux, macOS, and Windows on Go 1.24–1.26, plus golangci-lint and CodeQL.
  • Coordinated-omission-aware latency harness. internal/cmd/latency stamps events with their intended dispatch time and records end-to-end, writer-stall, and reader-lag percentiles in HDR histograms, so writer stalls can't hide queueing delay.

License

MIT, see LICENSE.

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

Examples

Constants

View Source
const CacheLineSize = 64

Variables

This section is empty.

Functions

This section is empty.

Types

type BackoffParams

type BackoffParams struct {
	InitialBackoff    time.Duration
	MaxBackoff        time.Duration
	TerminateAfter    time.Duration
	BackoffMultiplier float64
}

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]) Get

func (r *ReadView[T]) Get(seq int64) *T

func (*ReadView[T]) GetSegments

func (r *ReadView[T]) GetSegments(start, end int64) (seg1, seg2 []T)

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

func (r *ReadView[T]) Iterate(start, end int64) iter.Seq[*T]

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

func (r *ReadView[T]) LoadWriterBarrier() int64

type ReaderBarrier

type ReaderBarrier interface {
	Load() int64
}

type ReaderFunc

type ReaderFunc[T any] func(ctx context.Context, readView ReadView[T], readerCursor *atomic.Int64)

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]) AddReader

func (s *Stage[T]) AddReader(fn ReaderFunc[T]) (int, error)

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

func (s *Stage[T]) Load() int64

Load implements WriterBarrier using the cached scan path. Called only by the single writer goroutine via SetGatingStage.

func (*Stage[T]) RemoveReader

func (s *Stage[T]) RemoveReader(slotId int) error

func (*Stage[T]) Shutdown

func (s *Stage[T]) Shutdown()

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
}

Directories

Path Synopsis
internal
cmd/latency command

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL