Repository navigation
Add SPMC broadcast ring buffer; harden both rings for wrap-safety and bounded progress (v0.2.0) - #6
Merged
Merged
Conversation
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
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)
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
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::broadcastbroadcast::new(capacity) -> (Producer, Consumer); additional consumers register byConsumer::clone()(no allocation, no producer-side registry; clones start at the producer's current position, no replay).overwrite_count.push_slice, and apop_slicewhose uncontended fast path copies the whole run against onewrite_possnapshot with a single trailing validation (~2 synchronized loads per call instead of 2 per sample).Hardening (both rings)
usizewrap (~25 h of continuous 48 kHz audio on 32-bit targets) permanently stalled the SPSC consumer, made broadcastavailable()report 0 for poppable data, and panicked anyoverflow-checks = truebuild.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 withinlog2(capacity)retries; if it wins even those races,popreturnsNonewith all skips counted, rather than spinning. Hard iteration bound: ~2·log2(cap) + 1.Syncremoved. SPSCProducer/Consumerandbroadcast::Producerare nowSend + !Sync(enforced via marker +compile_faildoctests). Sharing&Produceracross threads compiled before and silently lost samples by racingwrite_pos.Specs / guarantees
capacitycapacity − 1(strictlag < capacity)overwrites + popped == pushed, exactoverwrites + popped == pushed_since_cloneunsafe, no allocations/locks/IO in any hot path, zero runtime deps.--cfg loomswaps the atomics inside the library (single crate-internalsyncshim 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
Syncremoved from SPSCProducer/Consumer: code sharing them viaArc/&across threads no longer compiles. That pattern silently dropped samples; move the endpoint into one thread instead.loom-testscargo feature removed (it required an undocumentedRUSTFLAGSto even compile; gate is--cfg loomalone now). Anyone using--all-featuresis unaffected — that invocation is fixed by this change.popmay returnNoneunder extreme producer pressure (skips pre-counted); broadcastavailable()caps atcapacity − 1.Non-breaking for everyone else: the existing SPSC
push/pop/push_slice/pop_slice/available/overwrite_count/capacitysignatures and semantics are unchanged.