Skip to main content

Crate rapidfire

Crate rapidfire 

Source
Expand description

§rapidfire

An ultra low-latency, lock-free, dependency-free async MPMC channel for Rust.

Designed for high-throughput, latency-critical pipelines (trading bots, WebSocket fan-in, market-data multiplexers) where message loss is unacceptable and every nanosecond on the hot path counts.

§Architecture

  • Own lock-free block queue, zero dependencies. Values live in 63-slot blocks chained through prev/next pointers. Producers claim a slot on the tail index (a wait-free fetch_add while uncontended, a CAS once contention is observed, because with several producers CAS serialises the claims and keeps adjacent-slot writes from fighting over one cache line); consumers check the head slot’s state and claim it with a CAS only once the value is there, so they never touch the producers’ cache line. See queue.rs for the full protocol and its proofs.
  • Lap-tagged slot states, never-freed blocks. A slot’s state word carries the block’s lap number, so recycled slot states need no reset and a stale state can never be mistaken for a fresh one. Blocks cycle through a spare slot and a small pool instead of being freed, which is what makes it sound to look at a slot before owning it. Memory stays at the channel’s high-water mark until it is dropped.
  • 128-byte cache-line isolation. Producer state, consumer state, the spare block, each block’s producer header and consumer header (read marks) and the rarely written wake-up flags all live on separate lines. On the hot path the only lines that cross cores are the slot lines carrying the values.
  • Zero-cost sleep and wake, no SeqCst on the hot path. While messages flow no lock is taken and no waker is cloned. A producer only does a Relaxed load of one read-mostly flag right after its claim; the party that goes to sleep pays instead, with one RMW on the other side’s index that every later claim synchronises with (release-sequence argument, model-checked with loom).
  • Cancellation safe. Dropping a Recv/Send future (e.g. in tokio::select!) removes its waker; no counter leaks.
  • Drop-in API. Matches async-channel: unbounded, bounded, send, recv, try_send, try_recv, close, counts and lengths.

§Example

use rapidfire::unbounded;

let (tx, rx) = unbounded::<u64>();
tx.send(42).await.unwrap();
assert_eq!(rx.recv().await.unwrap(), 42);

Modules§

mpsc
Channels with many senders and exactly one receiver.

Structs§

Receiver
The receiving half of a channel.
Recv
Future returned by Receiver::recv.
RecvError
An error returned when attempting to receive from a closed and empty channel.
Send
Future returned by Sender::send.
SendError
An error returned when attempting to send a message on a closed channel.
Sender
The sending half of a channel.

Enums§

TryRecvError
An error returned by Receiver::try_recv.
TrySendError
An error returned by Sender::try_send.

Functions§

bounded
Creates a bounded channel with the specified capacity.
unbounded
Creates an unbounded channel.