Skip to content

Add SPMC broadcast ring buffer; harden both rings for wrap-safety and bounded progress (v0.2.0) - #6

Merged
hasanzakeri merged 18 commits into
mainfrom
feat/broadcast-spmc
Aug 22, 2026
Merged

hasanzakeri merged 18 commits into
mainfrom
feat/broadcast-spmc

Conversation

@hasanzakeri

Copy link
Copy Markdown
Owner

What this delivers

A second ring flavour, rt_ring::broadcast — a single-producer, multi-consumer fan-out ring where every consumer owns an independent cursor and sees every sample it does not fall behind on — plus protocol-level hardening of both rings that came out of review.

New: rt_ring::broadcast

  • broadcast::new(capacity) -> (Producer, Consumer); additional consumers register by Consumer::clone() (no allocation, no producer-side registry; clones start at the producer's current position, no replay).
  • Producer is oblivious: push is 3 atomic ops, no CAS, no branches, independent of consumer count. Consumers detect lapping themselves and account skipped samples in a per-consumer overwrite_count.
  • Bulk ops: push_slice, and a pop_slice whose uncontended fast path copies the whole run against one write_pos snapshot with a single trailing validation (~2 synchronized loads per call instead of 2 per sample).

Hardening (both rings)

  • Wrap-safe position counters. All position arithmetic is wrapping. Before this, usize wrap (~25 h of continuous 48 kHz audio on 32-bit targets) permanently stalled the SPSC consumer, made broadcast available() report 0 for poppable data, and panicked any overflow-checks = true build.
  • Bounded, wait-free broadcast pop. Lap recovery doubles its headroom over the producer on each consecutive lap within a call, so a producer re-lapping the consumer is outrun within log2(capacity) retries; if it wins even those races, pop returns None with all skips counted, rather than spinning. Hard iteration bound: ~2·log2(cap) + 1.
  • Accidental Sync removed. SPSC Producer/Consumer and broadcast::Producer are now Send + !Sync (enforced via marker + compile_fail doctests). Sharing &Producer across threads compiled before and silently lost samples by racing write_pos.

Specs / guarantees

Property SPSC (root) broadcast
Topology 1 producer ↔ 1 consumer 1 producer → N consumers
Producer push lock-free (CAS only when full) wait-free, 3 atomic ops
Consumer pop lock-free (CAS retry only on producer progress) wait-free, bounded retries
Overwrite window capacity capacity − 1 (strict lag < capacity)
Accounting overwrites + popped == pushed, exact per-consumer, overwrites + popped == pushed_since_clone
Capacity next power of two next power of two, floor 2
Counters wrap-safe wrap-safe
  • No unsafe, no allocations/locks/IO in any hot path, zero runtime deps.
  • Memory orderings are model-checked: --cfg loom swaps the atomics inside the library (single crate-internal sync shim shared by both rings), so loom verifies the orderings themselves, not just schedules. Suite: unit / integration / proptest-vs-reference-model / loom / miri / two libFuzzer targets, all green; the broadcast fuzz target shares its reference model with the proptest (single-sourced).

Backward compatibility — breaking, requires 0.2.0

  1. Sync removed from SPSC Producer/Consumer: code sharing them via Arc/& across threads no longer compiles. That pattern silently dropped samples; move the endpoint into one thread instead.
  2. MSRV 1.65 → 1.87, now actually enforced in CI across all targets.
  3. loom-tests cargo feature removed (it required an undocumented RUSTFLAGS to even compile; gate is --cfg loom alone now). Anyone using --all-features is unaffected — that invocation is fixed by this change.
  4. New behavior, documented: broadcast pop may return None under extreme producer pressure (skips pre-counted); broadcast available() caps at capacity − 1.
  5. Cross-thread bench numbers are not comparable to 0.1.0 (harness now times only the transfer window, uniformly across the ringbuf/tokio comparisons).

Non-breaking for everyone else: the existing SPSC push/pop/push_slice/ pop_slice/available/overwrite_count/capacity signatures and semantics are unchanged.

hasanzakeri added 18 commits May 9, 2026 17:37
Adds rt_ring::broadcast as a sibling to the SPSC API. Every consumer sees
every sample; additional consumers register dynamically by cloning. The
producer is oblivious — push is 3 atomic ops with no CAS, independent of
how many consumers exist. Each consumer detects lapping itself via a
double-checked write_pos guard around the slot read.

Customer-facing trade-offs documented in module/README/CHANGELOG:

- Capacity floored to 2 (cap == 1 unsupported by the lapping protocol).
- At the lap boundary (lag == capacity), pop conservatively retries to
  avoid a producer/consumer race; this can over-count overwrites by 1
  under tight pacing at very small caps. The accounting invariant
  overwrites + popped == samples_pushed_since_clone is unaffected.
- Consumer is Send + !Sync (cursor is a Cell, single-thread access only).

Verification mirrors the SPSC posture: parallel basic / overwrite /
concurrent / proptest / loom test files plus a model-based fuzz target.
Adds three broadcast benches alongside the SPSC ones:
- broadcast push 1M (1 cons, single-thread)
- broadcast push+pop 1M (single-thread)
- broadcast cross-thread 1M (1 and 4 consumers)

Quick run on Apple M-series:
- broadcast push+pop is ~2x faster than SPSC push+pop in single thread
  (no CAS on either side).
- broadcast push (alone) is comparable to SPSC push (LSE atomics make
  SPSC's overflow CAS very cheap).
- broadcast cross-thread is ~2.7x slower than SPSC cross-thread, mainly
  from the consumer's two write_pos Acquire loads per pop. Cost scales
  roughly linearly with consumer count.

These match the design's documented trade-off: cheaper hot paths in
single-thread, slightly more wp cache traffic in cross-thread.
Under miri's deterministic scheduler with miri-reduced iteration counts
(count=100, cap=64), the consumer can keep up with the producer and never
get lapped — making the 'overwrites > 0' sanity check fail. Gate it on
!cfg!(miri); the accounting invariant assertion stays unconditional.
Adds two bench-only comparisons (no public API changes):

1. A sequencer-based broadcast prototype using packed
   AtomicU64 = (seq:u32 << 32) | (bits:u32). Single atomic load gets a
   consistent (seq, data) pair, sidestepping the producer/consumer race
   that two-atomic designs need a double-check around. Trade-offs vs the
   shipped lapping design: cap=1 works, no lap-boundary false positive,
   2x slot memory, 32-bit seq wraps after ~2^32 pushes.

2. A tokio::sync::broadcast comparison for context. tokio adds as a
   dev-dep with sync/rt features only.

Quick measurements (Apple M-series, 1M samples; full criterion):
- broadcast push+pop:       ~2.1 ms
- sequencer push 1M:        ~1.7 ms  (fastest single-thread)
- sequencer cross-thread:   ~5.4 ms  (~3-4x faster than lapping broadcast)
- broadcast cross-thread:   ~13-20 ms (high variance)
- tokio::broadcast:         ~24 ms

Bottom line: the sequencer's packed-atomic approach is meaningfully
faster AND fixes the cap=1 / boundary-false-positive limitations, at the
cost of a 32-bit seq wraparound. Not landing it as the default — keeping
the lapping design's unbounded lifetime guarantee — but the prototype
stays in benches as a reference and an option for callers that need it.
…rn the truly-poppable count from broadcast available()

Producer (SPSC + broadcast) and SPSC Consumer were accidentally Sync
via Arc<Shared>. Concurrent push/pop through Arc<Producer> /
Arc<Consumer> compiled cleanly and silently dropped samples by racing
write_pos (broadcast) or stepping on the same CAS retry (SPSC). A
PhantomData<Cell<()>> marker on each affected struct now turns that
into a compile error, and the broadcast::Consumer (already !Sync via
its Cell fields) gains a matching documentation guarantee. The intent
is verified by per-type compile_fail doctests that fail to fail the
moment any of these types regains Sync.

broadcast::Consumer::available() previously returned `lag.min(cap)`,
which advertised a value (cap) that no single call to pop() can ever
yield -- at lag == cap the lap-detection protocol triggers recovery
and only cap - 1 samples are actually drainable. The clamp now reads
`lag.min(cap - 1)`, so callers polling for "can I read N samples now"
get a truthful answer; a new test guards the boundary case.

Module docs and README are tightened to match the broadcast type's
real promise: every consumer sees every sample it does not fall
behind on, with a usable window of cap - 1; cap == 1 is documented
as silently rounded up to 2 rather than "unsupported".
…edules

The pre-existing loom tests used `std::sync::atomic` unconditionally,
which meant loom could only permute the thread scheduler at yield
points -- it had no visibility into the acquire/release pairings or
the Relaxed/Release sequencing that the protocols depend on. A
ring-buffer crate marketing "real-time, lock-free, verified" wants
loom looking inside the atomics, not just around them.

Both shared modules now select between std and loom atomics via a
`cfg(loom)` shim. `loom` itself moves from a flat dev-dependency to
a `[target.'cfg(loom)'.dependencies]` entry so it is only pulled in
when the cfg is set; otherwise the library compiles and ships
exactly as before. `Cargo.toml` declares `cfg(loom)` as expected via
`lints.rust = { unexpected_cfgs = ... }` so `-D warnings` does not
trip on the new gate.

The loom CI job now sets `RUSTFLAGS=--cfg loom` and runs both the
SPSC and broadcast loom targets (the previous `--test loom_spsc`
filter silently excluded broadcast). With the shim active, each loom
file takes ~4 s of real ordering exploration instead of the ~10 ms
scheduler-only pass it did before -- a useful signal on its own.
The previous MSRV of 1.65 was unenforced and silently broken: the
MSRV CI step ran `cargo check`, which never compiles tests or
dev-dependencies. Several test files used `u64::is_multiple_of`
(stabilized in 1.87), the resolved `tokio 1.52.x` dev-dep needs
Rust >= 1.70, and so on -- none of which the job would have caught.

`rust-version` is now 1.87 (the floor of everything actually used in
the repo), and the MSRV job runs `cargo test --no-run --all-targets`
so an MSRV regression in any compiled artifact -- library, tests,
benches, dev-deps -- blocks CI.
…pposed to measure

Two changes that make the broadcast numbers in the README actually
mean what the column headers claim:

bench_tokio_broadcast_cross_thread previously built a fresh
multi-thread Runtime inside `b.iter()`. Worker-thread spawn plus
buffer allocation cost several ms on its own and swamped the
broadcast itself, making the tokio comparison artificially worse.
The runtime is now constructed once at function entry and reused
across iterations, so the bench measures one `tokio::broadcast`
round-trip instead of a runtime cold-start every time.

bench_cross_thread (rt-ring SPSC) used a "two-phase exit" --
`count > 0 && available == 0`, then "try once more" -- that could
terminate before all 1M samples had been popped if the producer was
mid-push at the check, and could spin without yielding when the
consumer woke before any pushes were visible. Replaced with the
robust drain-until-N-1 pattern that bench_broadcast_cross_thread
already uses, with `spin_loop_hint` while empty. Counts are now
deterministic and the number reported is the cost of actually
moving 1M samples cross-thread.
…oducer signals done

clone_during_active_production_is_monotonic drains the clone in two
modes: the steady "Some(v) -> bump popped + monotonicity" branch,
and a post-done sweep that retries pop() once after the producer
signals completion. The sweep was written as
`if clone.pop().is_none() { break; }` -- if the retry returned
Some(v), the value was thrown away without bumping `popped` or
checking monotonicity. The end-of-test invariant
`overwrites + popped <= total` is a `<=`, so the silent drop never
tripped it; it just made the assertion strictly weaker than the
docstring promises.

The sweep now matches on the result and treats a `Some` exactly like
the steady branch -- monotonicity check + popped += 1 -- so every
sample the clone actually receives contributes to the bound the test
is supposed to enforce.

The miri gate added in the previous review pass for
slow_consumer_forces_overwrites also lives here as part of the same
file -- preserved.
…d capacity setup

Both rings carried private copies of the loom/std atomic import shim,
the 64-byte CachePadded wrapper, and the power-of-two capacity rounding
plus slot allocation. Centralize them in a crate-internal sync module so
the --cfg loom swap happens in exactly one place and the two rings
cannot drift apart on padding or capacity rules.
Positions are monotonically increasing usize counters, which wrap after
2^32 samples on 32-bit targets (~25 hours of continuous 48 kHz audio).
The arithmetic relied on unchecked +/- and two constructs that break
outright at the wrap:

- The SPSC emptiness test (rp >= wp) becomes permanently true once
  write_pos wraps below read_pos, stalling the consumer forever.
  Equality is the complete and wrap-safe test, since read_pos never
  passes write_pos.
- broadcast available() used saturating_sub while pop used plain
  subtraction, so after the wrap available() reported 0 for data that
  pop would return - a consumer gating its drain on available() > 0
  stalls permanently.
- Any build with overflow-checks enabled panicked in the hot path at
  the wrap.

All position math now uses wrapping_add/wrapping_sub, SPSC available()
loads read_pos first and clamps to capacity so it can neither underflow
nor over-report, and the wrap invariants are documented at the top of
each ring's core module.
The lap-recovery loop in pop was unbounded: recovery parked the cursor
one slot ahead of the producer's overwrite frontier, so a producer that
landed a single push between the consumer's cursor update and its next
write_pos read would re-trigger recovery. A producer pushing flat-out on
a faster core could repeat that win indefinitely and starve pop while
the ring was full of data.

Recovery now doubles its headroom over the producer on every consecutive
lap within a call: re-lapping requires the producer to push headroom
more samples within one loop iteration, so any realistic producer is
outrun within log2(capacity) retries. If the producer wins even those
races, pop gives up and returns None - with every skipped sample already
counted in overwrite_count - instead of spinning. The loop is hard
bounded at ~2*log2(capacity) + 1 iterations.

BREAKING CHANGE: pop returning None now means "nothing readable right
now": the ring is empty, or the producer is overwriting faster than
this consumer can safely read. The cases are distinguishable via
overwrite_count. Single-threaded and uncontended behavior is unchanged.
…hared

pop_slice paid two synchronized write_pos loads per element by looping
pop. When the consumer is not lapped, it now copies the whole run
against a single write_pos snapshot and validates once at the end: a
producer write can only have raced slot rc+i by reaching position
rc+i+cap, which the trailing check exposes for every copied slot at
once (the same conservative margin as pop's second lap check). Lapped
or contended calls fall back to the per-element path.

push_slice stays a loop over push on purpose: the consumers' torn-read
guard assumes at most one unpublished slot store is in flight, so
batching slot stores ahead of a single write_pos publish would silently
break it. That constraint is now documented in push.

Both slice ops live in Shared, mirroring the SPSC layering.
The loom test files were gated on a loom-tests cargo feature while the
loom crate itself was only a dependency under cfg(loom). Enabling the
feature without RUSTFLAGS="--cfg loom" - for example plain
`cargo test --all-features` - was therefore a guaranteed compile
error.

Gate the loom tests on cfg(loom) alone: one knob instead of two
half-switches. Without the cfg the test files compile to empty
binaries; with it, `cargo test --test loom_spsc --test broadcast_loom`
runs them as before. Makefile and CI updated.

BREAKING CHANGE: the loom-tests cargo feature no longer exists.
- tests/common/mod.rs now holds iteration_count and assert_monotonic,
  previously duplicated byte-for-byte across the SPSC and broadcast
  suites.
- The broadcast reference model (ref_pop) was duplicated between the
  proptest suite and the fuzz target; it now lives in
  tests/common/ref_model.rs, used by the proptest through the common
  module and textually included by the fuzz target.
- The fuzz target's op decoding is a plain high-bit test instead of an
  obscured two-bit match (same byte-to-operation mapping, so existing
  corpus entries keep their meaning), and the pop payload is no longer
  computed on push ops.
- The unreachable final-value check in the broadcast drain helper is
  gone.
- make fuzz runs the broadcast fuzz target in addition to push_pop; it
  was registered but never executed by any automation.
The cross-thread benches allocated the ring, built the barrier, and
spawned/joined threads inside the measured iteration, folding scheduler
noise into the reported throughput. Every cross-thread bench now uses
iter_custom and times only barrier-release to drain-complete, uniformly
across the rt-ring, ringbuf, and tokio comparisons, so the numbers stay
apples-to-apples. The tokio bench also gains the start synchronization
it previously lacked.

The drain-until-final-value consumer loop existed in four near-identical
copies, each carrying a provably unreachable break arm; one helper
replaces them.

Bench ids are unchanged; absolute numbers are not comparable to earlier
runs.
- Key Properties no longer claims unqualified overwrite-on-full and
  "real-time safe": broadcast's effective window is capacity - 1, and
  the progress properties are stated per ring (wait-free broadcast push
  and pop, CAS-retry-on-progress SPSC), along with the wrap-safety of
  the position counters.
- available() is documented as SPSC occupancy vs broadcast per-consumer
  poppable count capped at capacity - 1, and the broadcast doc no
  longer claims the returned samples can be popped "without triggering
  lap recovery", which was false at lag == capacity. It also notes the
  value is a racy snapshot under a concurrent producer.
- New README sections: installation with the MSRV, and the development
  commands including how the loom model-checking is wired.
Bump the version for the breaking changes in this cycle (Sync removal
on the SPSC endpoints, MSRV 1.87, loom-tests feature removal, the new
broadcast pop None semantics), retitle the changelog accordingly, and
refresh the crates.io keywords to cover the broadcast module.
@hasanzakeri
hasanzakeri merged commit 4f0967f into main Aug 22, 2026
6 checks passed
@hasanzakeri
hasanzakeri deleted the feat/broadcast-spmc branch August 22, 2026 20:16
hasanzakeri added a commit that referenced this pull request Oct 8, 2026
Add SPMC broadcast ring buffer; harden both rings for wrap-safety and bounded progress (v0.2.0)
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant