Skip to main content

MpmcQueue

Struct MpmcQueue 

Source
pub struct MpmcQueue<T> { /* private fields */ }
Expand description

Bounded MPMC ring queue. Capacity is rounded up to the next power of two (minimum 2).

Implementations§

Source§

impl<T> MpmcQueue<T>

Source

pub fn new(capacity: usize) -> Self

Examples found in repository?
examples/perf_features.rs (line 317)
316        fn ring(n: usize) -> MpmcQueue<u64> {
317            let q = MpmcQueue::new(n);
318            for i in 0..n / 2 {
319                let _ = q.try_enqueue(i as u64);
320            }
321            q
322        }
More examples
Hide additional examples
examples/sample_app.rs (line 133)
126fn mpmc_sharded_match() {
127    use std::sync::atomic::{AtomicUsize, Ordering};
128
129    use subms_mpsc_queue::MpmcQueue;
130
131    println!("\n== mpmc: shard the match loop across several consumers ==");
132    let shards = 3usize;
133    let ring: Arc<MpmcQueue<u64>> = Arc::new(MpmcQueue::new(1_024));
134    let total = GATEWAYS * ORDERS_PER_GATEWAY;
135
136    let gateways: Vec<_> = (0..GATEWAYS)
137        .map(|g| {
138            let ring = Arc::clone(&ring);
139            thread::spawn(move || {
140                for seq in 0..ORDERS_PER_GATEWAY {
141                    let mut order = order_id(g, seq);
142                    while let Err(rejected) = ring.try_enqueue(order) {
143                        order = rejected;
144                        std::hint::spin_loop();
145                    }
146                }
147            })
148        })
149        .collect();
150
151    let matched = Arc::new(AtomicUsize::new(0));
152    let consumers: Vec<_> = (0..shards)
153        .map(|_| {
154            let ring = Arc::clone(&ring);
155            let matched = Arc::clone(&matched);
156            thread::spawn(move || {
157                let mut local = 0usize;
158                loop {
159                    if ring.try_dequeue().is_some() {
160                        local += 1;
161                        matched.fetch_add(1, Ordering::Relaxed);
162                    } else if matched.load(Ordering::Relaxed) >= total {
163                        break;
164                    } else {
165                        std::hint::spin_loop();
166                    }
167                }
168                local
169            })
170        })
171        .collect();
172
173    for h in gateways {
174        h.join().unwrap();
175    }
176    let drained: usize = consumers.into_iter().map(|c| c.join().unwrap()).sum();
177    // cas_retries() is the contention read-out, deliberately not printed: it is
178    // a property of how the OS scheduled these threads on this run.
179    println!(
180        "  {shards} shards drained {drained} orders, ring empty: {}",
181        ring.is_empty()
182    );
183    assert_eq!(
184        drained, total,
185        "shards together drain every order exactly once"
186    );
187    assert_eq!(
188        ring.producer_index(),
189        ring.consumer_index(),
190        "every claimed slot was consumed"
191    );
192}
Source

pub fn capacity(&self) -> usize

Source

pub fn producer_index(&self) -> usize

Monotonic count of slots ever claimed by producers.

Examples found in repository?
examples/sample_app.rs (line 188)
126fn mpmc_sharded_match() {
127    use std::sync::atomic::{AtomicUsize, Ordering};
128
129    use subms_mpsc_queue::MpmcQueue;
130
131    println!("\n== mpmc: shard the match loop across several consumers ==");
132    let shards = 3usize;
133    let ring: Arc<MpmcQueue<u64>> = Arc::new(MpmcQueue::new(1_024));
134    let total = GATEWAYS * ORDERS_PER_GATEWAY;
135
136    let gateways: Vec<_> = (0..GATEWAYS)
137        .map(|g| {
138            let ring = Arc::clone(&ring);
139            thread::spawn(move || {
140                for seq in 0..ORDERS_PER_GATEWAY {
141                    let mut order = order_id(g, seq);
142                    while let Err(rejected) = ring.try_enqueue(order) {
143                        order = rejected;
144                        std::hint::spin_loop();
145                    }
146                }
147            })
148        })
149        .collect();
150
151    let matched = Arc::new(AtomicUsize::new(0));
152    let consumers: Vec<_> = (0..shards)
153        .map(|_| {
154            let ring = Arc::clone(&ring);
155            let matched = Arc::clone(&matched);
156            thread::spawn(move || {
157                let mut local = 0usize;
158                loop {
159                    if ring.try_dequeue().is_some() {
160                        local += 1;
161                        matched.fetch_add(1, Ordering::Relaxed);
162                    } else if matched.load(Ordering::Relaxed) >= total {
163                        break;
164                    } else {
165                        std::hint::spin_loop();
166                    }
167                }
168                local
169            })
170        })
171        .collect();
172
173    for h in gateways {
174        h.join().unwrap();
175    }
176    let drained: usize = consumers.into_iter().map(|c| c.join().unwrap()).sum();
177    // cas_retries() is the contention read-out, deliberately not printed: it is
178    // a property of how the OS scheduled these threads on this run.
179    println!(
180        "  {shards} shards drained {drained} orders, ring empty: {}",
181        ring.is_empty()
182    );
183    assert_eq!(
184        drained, total,
185        "shards together drain every order exactly once"
186    );
187    assert_eq!(
188        ring.producer_index(),
189        ring.consumer_index(),
190        "every claimed slot was consumed"
191    );
192}
Source

pub fn consumer_index(&self) -> usize

Monotonic count of slots ever claimed by consumers.

Examples found in repository?
examples/sample_app.rs (line 189)
126fn mpmc_sharded_match() {
127    use std::sync::atomic::{AtomicUsize, Ordering};
128
129    use subms_mpsc_queue::MpmcQueue;
130
131    println!("\n== mpmc: shard the match loop across several consumers ==");
132    let shards = 3usize;
133    let ring: Arc<MpmcQueue<u64>> = Arc::new(MpmcQueue::new(1_024));
134    let total = GATEWAYS * ORDERS_PER_GATEWAY;
135
136    let gateways: Vec<_> = (0..GATEWAYS)
137        .map(|g| {
138            let ring = Arc::clone(&ring);
139            thread::spawn(move || {
140                for seq in 0..ORDERS_PER_GATEWAY {
141                    let mut order = order_id(g, seq);
142                    while let Err(rejected) = ring.try_enqueue(order) {
143                        order = rejected;
144                        std::hint::spin_loop();
145                    }
146                }
147            })
148        })
149        .collect();
150
151    let matched = Arc::new(AtomicUsize::new(0));
152    let consumers: Vec<_> = (0..shards)
153        .map(|_| {
154            let ring = Arc::clone(&ring);
155            let matched = Arc::clone(&matched);
156            thread::spawn(move || {
157                let mut local = 0usize;
158                loop {
159                    if ring.try_dequeue().is_some() {
160                        local += 1;
161                        matched.fetch_add(1, Ordering::Relaxed);
162                    } else if matched.load(Ordering::Relaxed) >= total {
163                        break;
164                    } else {
165                        std::hint::spin_loop();
166                    }
167                }
168                local
169            })
170        })
171        .collect();
172
173    for h in gateways {
174        h.join().unwrap();
175    }
176    let drained: usize = consumers.into_iter().map(|c| c.join().unwrap()).sum();
177    // cas_retries() is the contention read-out, deliberately not printed: it is
178    // a property of how the OS scheduled these threads on this run.
179    println!(
180        "  {shards} shards drained {drained} orders, ring empty: {}",
181        ring.is_empty()
182    );
183    assert_eq!(
184        drained, total,
185        "shards together drain every order exactly once"
186    );
187    assert_eq!(
188        ring.producer_index(),
189        ring.consumer_index(),
190        "every claimed slot was consumed"
191    );
192}
Source

pub fn cas_retries(&self) -> u64

Total CAS retries (both producers losing tail-CAS and consumers losing head-CAS). Useful for diagnosing contention; ignored by the hot path otherwise.

Source

pub fn try_enqueue(&self, value: T) -> Result<(), T>

Multi-producer enqueue. Returns Err(value) if the ring is full.

Examples found in repository?
examples/perf_features.rs (line 319)
195fn main() -> io::Result<()> {
196    let path = PathBuf::from(env!("CARGO_MANIFEST_DIR"))
197        .join("..")
198        .join(".subms")
199        .join("features")
200        .join("rust.json");
201    let existing = std::fs::read_to_string(&path).unwrap_or_default();
202    let mut manifest = SubMsFeatureManifest::load_str("rust", &existing);
203    // Stamp the box these numbers came from. The bench runs wherever it is
204    // invoked, so an unstamped manifest is indistinguishable from a fleet
205    // capture; the renderer will not publish one it cannot attribute.
206    let (source, instance) = SubMsP99Source::from_env();
207    manifest.set_p99_source(source, instance.as_deref());
208
209    // Optional, off by default: pin this thread to one core for the whole run.
210    // On a heterogeneous laptop the scheduler moves the bench between core
211    // clusters and every measurement lands in one of two clock states 1.31x
212    // apart - a spread three times wider than the deltas being classified, and
213    // large enough on its own to flip a feature between auxiliary and hot-path.
214    // Pinned, the same sweep repeats to within 1%. Left OFF by default because a
215    // fleet box isolates cores outside the process, and pinning from in here
216    // would override that placement with a core the orchestrator did not choose.
217    #[cfg(feature = "affinity")]
218    if let Some(core) = std::env::var("SUBMS_PIN").ok().and_then(|v| v.parse().ok()) {
219        let _ = subms_mpsc_queue::set_affinity(&[core]);
220    }
221
222    // Burn before the first measurement, not just before each one. Every
223    // `batched` call warms itself, but the FIRST measurement in the process pays
224    // a ramp the per-measurement warm sits inside rather than absorbs, and the
225    // sweep runs smallest-first: without this the base curve read 71800 / 45300 /
226    // 46500 ns, a 1.6x fall with size that is the process settling, not the
227    // queue.
228    {
229        let mut q = filled(CANON);
230        let start = std::time::Instant::now();
231        while (start.elapsed().as_nanos() as u64) < BURN_NANOS {
232            for i in 0..ITEMS_PER_SAMPLE {
233                q.push(i as u64);
234                black_box(q.try_pop());
235            }
236        }
237    }
238
239    // The baseline: the base queue's push + try_pop round trip. Swept as well as
240    // sampled, because whether queue depth moves the BASE op is the context
241    // every feature curve is read against.
242    let base_sweep = sweep("base/push+pop", |n| {
243        let mut q = filled(n);
244        batched(ITEMS_PER_SAMPLE, |i| {
245            q.push(i as u64);
246            black_box(q.try_pop());
247        })
248    });
249    let base_p50 = base_sweep
250        .iter()
251        .find(|(n, _)| *n == CANON)
252        .map_or(0, |(_, v)| *v);
253    eprintln!("base push+pop p50 per {ITEMS_PER_SAMPLE}-item sample: {base_p50}ns");
254
255    // ---------- bounded: fixed-capacity ring, backpressure on enqueue ----------
256    #[cfg(feature = "bounded")]
257    {
258        use subms_mpsc_queue::BoundedMpscQueue;
259        // Half full at every sweep point. Filled to a FIXED element count
260        // instead, the big rings would sit 98% empty and the enqueue would be
261        // measuring the fill fraction rather than the footprint.
262        fn ring(n: usize) -> BoundedMpscQueue<u64> {
263            let q = BoundedMpscQueue::new(n);
264            for i in 0..n / 2 {
265                let _ = q.try_enqueue(i as u64);
266            }
267            q
268        }
269        let sw = sweep("bounded/enqueue+dequeue", |n| {
270            let mut q = ring(n);
271            batched(ITEMS_PER_SAMPLE, |i| {
272                let _ = q.try_enqueue(i as u64);
273                black_box(q.try_dequeue());
274            })
275        });
276        let (cat, reason) = classify_feature(&sw, Some(base_p50), None);
277
278        let mut q = ring(CANON);
279        let mut p99 = BTreeMap::new();
280        p99.insert(
281            "enqueue".to_string(),
282            single(|st, i| {
283                st.time(|| {
284                    let _ = q.try_enqueue(i as u64);
285                });
286                black_box(q.try_dequeue());
287            }),
288        );
289        p99.insert(
290            "dequeue".to_string(),
291            single(|st, i| {
292                let _ = q.try_enqueue(i as u64);
293                st.time(|| black_box(q.try_dequeue()));
294            }),
295        );
296        // The reject path, which is the reason the feature exists. The ring is
297        // filled to capacity once, OUTSIDE the timed region; every timed call
298        // then takes the full branch and hands the value back to the caller.
299        let full: BoundedMpscQueue<u64> = BoundedMpscQueue::new(CANON);
300        while full.try_enqueue(0).is_ok() {}
301        p99.insert(
302            "enqueue_full".to_string(),
303            single(|st, i| {
304                st.time(|| {
305                    let _ = full.try_enqueue(i as u64);
306                });
307            }),
308        );
309        manifest.set_feature("bounded", cat, &p99, &reason);
310    }
311
312    // ---------- mpmc: bounded ring, sequence CAS on both ends ----------
313    #[cfg(feature = "mpmc")]
314    {
315        use subms_mpsc_queue::MpmcQueue;
316        fn ring(n: usize) -> MpmcQueue<u64> {
317            let q = MpmcQueue::new(n);
318            for i in 0..n / 2 {
319                let _ = q.try_enqueue(i as u64);
320            }
321            q
322        }
323        // Uncontended, so every CAS succeeds first try. That is the figure the
324        // category is about: what the multi-consumer claim costs a queue that is
325        // NOT contended, which is the state a well-sized pipeline runs in.
326        let sw = sweep("mpmc/enqueue+dequeue", |n| {
327            let q = ring(n);
328            batched(ITEMS_PER_SAMPLE, |i| {
329                let _ = q.try_enqueue(i as u64);
330                black_box(q.try_dequeue());
331            })
332        });
333        let (cat, reason) = classify_feature(&sw, Some(base_p50), None);
334
335        let q = ring(CANON);
336        let mut p99 = BTreeMap::new();
337        p99.insert(
338            "enqueue".to_string(),
339            single(|st, i| {
340                st.time(|| {
341                    let _ = q.try_enqueue(i as u64);
342                });
343                black_box(q.try_dequeue());
344            }),
345        );
346        p99.insert(
347            "dequeue".to_string(),
348            single(|st, i| {
349                let _ = q.try_enqueue(i as u64);
350                st.time(|| black_box(q.try_dequeue()));
351            }),
352        );
353        manifest.set_feature("mpmc", cat, &p99, &reason);
354    }
355
356    // ---------- batch: drain up to BATCH items behind one acquire fence ----------
357    #[cfg(feature = "batch")]
358    {
359        use subms_mpsc_queue::BatchMpscQueue;
360        fn filled_batch(n: usize) -> BatchMpscQueue<u64> {
361            let q = BatchMpscQueue::new();
362            for i in 0..n {
363                q.push(i as u64);
364            }
365            q
366        }
367        // A sample moves ITEMS_PER_SAMPLE items either way; only the call width
368        // differs. That is why the reps count is divided rather than the batch
369        // grown - growing it would sweep the batch size, and the number would
370        // stop being comparable to the base round trip.
371        let sw = sweep("batch/push+dequeue_batch", |n| {
372            let mut q = filled_batch(n);
373            let mut buf: Vec<Option<u64>> = (0..BATCH).map(|_| None).collect();
374            batched(ITEMS_PER_SAMPLE / BATCH, |i| {
375                for j in 0..BATCH {
376                    q.push((i + j) as u64);
377                }
378                black_box(q.try_dequeue_batch(&mut buf));
379            })
380        });
381        let (cat, reason) = classify_feature(&sw, Some(base_p50), None);
382
383        let mut q = filled_batch(CANON);
384        let mut buf: Vec<Option<u64>> = (0..BATCH).map(|_| None).collect();
385        let mut p99 = BTreeMap::new();
386        // The refill is outside the timed region: timing it would put a BATCH of
387        // pushes inside the drain's number and the stage would stop being a
388        // drain figure at all.
389        p99.insert(
390            "dequeue_batch".to_string(),
391            single(|st, i| {
392                st.time(|| black_box(q.try_dequeue_batch(&mut buf)));
393                for j in 0..BATCH {
394                    q.push((i + j) as u64);
395                }
396            }),
397        );
398        p99.insert(
399            "enqueue".to_string(),
400            single(|st, i| {
401                st.time(|| q.push(i as u64));
402                let _ = q.try_dequeue_batch(&mut buf[..1]);
403            }),
404        );
405        // The producer mirror: BATCH items published behind one head swap. The
406        // drain that puts the queue back is outside the timed region for the
407        // same reason the refill is above.
408        p99.insert(
409            "enqueue_batch".to_string(),
410            single(|st, i| {
411                let base = i as u64;
412                st.time(|| black_box(q.push_batch(base..base + BATCH as u64)));
413                let _ = q.try_dequeue_batch(&mut buf);
414            }),
415        );
416        manifest.set_feature("batch", cat, &p99, &reason);
417    }
418
419    // ---------- metrics: relaxed atomic counters around each op ----------
420    #[cfg(feature = "metrics")]
421    {
422        use subms_mpsc_queue::MetricsMpscQueue;
423        fn filled_metrics(n: usize) -> MetricsMpscQueue<u64> {
424            let q = MetricsMpscQueue::new();
425            for i in 0..n {
426                q.push(i as u64);
427            }
428            q
429        }
430        let sw = sweep("metrics/push+pop", |n| {
431            let mut q = filled_metrics(n);
432            batched(ITEMS_PER_SAMPLE, |i| {
433                q.push(i as u64);
434                black_box(q.try_pop());
435            })
436        });
437        let (cat, reason) = classify_feature(&sw, Some(base_p50), None);
438
439        let mut q = filled_metrics(CANON);
440        let mut p99 = BTreeMap::new();
441        p99.insert(
442            "enqueue".to_string(),
443            single(|st, i| {
444                st.time(|| q.push(i as u64));
445                black_box(q.try_pop());
446            }),
447        );
448        p99.insert(
449            "dequeue".to_string(),
450            single(|st, i| {
451                q.push(i as u64);
452                st.time(|| black_box(q.try_pop()));
453            }),
454        );
455        p99.insert(
456            "snapshot".to_string(),
457            single(|st, _| {
458                st.time(|| black_box(q.snapshot()));
459            }),
460        );
461        manifest.set_feature("metrics", cat, &p99, &reason);
462    }
463
464    // ---------- affinity: pin the calling thread, once, at startup ----------
465    // Runs LAST because measuring it pins THIS process to core 0, and every
466    // number taken afterwards would be a number taken on one core.
467    #[cfg(feature = "affinity")]
468    {
469        use subms_mpsc_queue::set_affinity;
470        // Swept over the same axis to show what it is: a call that touches no
471        // queue state and cannot move with queue size. PINNED auxiliary rather
472        // than left to the base-delta test, which would see a syscall costing
473        // more than an enqueue and call it hot-path. It is not on the hot path at
474        // any price - `set_affinity` is called once per thread at startup and
475        // appears in neither `push` nor `try_pop`. The two ports are not even
476        // measuring the same thing: Rust issues a real `SetThreadAffinityMask` /
477        // `sched_setaffinity`, while the Java sibling validates its argument and
478        // returns UNSUPPORTED because the stock JDK has no pinning API. The
479        // per-call figure, not the sample, is the interpretable one and it is in
480        // `p99ByStage`.
481        let sw = sweep("affinity/set_affinity", |_| {
482            batched(ITEMS_PER_SAMPLE, |_| {
483                let _ = set_affinity(&[0]);
484            })
485        });
486        let (cat, reason) = classify_feature(
487            &sw,
488            Some(base_p50),
489            Some(subms::SubMsFeatureCategory::Auxiliary),
490        );
491
492        let mut p99 = BTreeMap::new();
493        p99.insert(
494            "set_affinity".to_string(),
495            single(|st, _| {
496                st.time(|| {
497                    let _ = set_affinity(&[0]);
498                });
499            }),
500        );
501        manifest.set_feature("affinity", cat, &p99, &reason);
502
503        let cores: Vec<usize> = (0..std::thread::available_parallelism()
504            .map_or(1, std::num::NonZeroUsize::get))
505            .collect();
506        let _ = set_affinity(&cores);
507    }
508
509    std::fs::create_dir_all(path.parent().unwrap())?;
510    std::fs::write(&path, manifest.to_json())?;
511    io::stdout().write_all(manifest.to_json().as_bytes())?;
512    Ok(())
513}
More examples
Hide additional examples
examples/sample_app.rs (line 142)
126fn mpmc_sharded_match() {
127    use std::sync::atomic::{AtomicUsize, Ordering};
128
129    use subms_mpsc_queue::MpmcQueue;
130
131    println!("\n== mpmc: shard the match loop across several consumers ==");
132    let shards = 3usize;
133    let ring: Arc<MpmcQueue<u64>> = Arc::new(MpmcQueue::new(1_024));
134    let total = GATEWAYS * ORDERS_PER_GATEWAY;
135
136    let gateways: Vec<_> = (0..GATEWAYS)
137        .map(|g| {
138            let ring = Arc::clone(&ring);
139            thread::spawn(move || {
140                for seq in 0..ORDERS_PER_GATEWAY {
141                    let mut order = order_id(g, seq);
142                    while let Err(rejected) = ring.try_enqueue(order) {
143                        order = rejected;
144                        std::hint::spin_loop();
145                    }
146                }
147            })
148        })
149        .collect();
150
151    let matched = Arc::new(AtomicUsize::new(0));
152    let consumers: Vec<_> = (0..shards)
153        .map(|_| {
154            let ring = Arc::clone(&ring);
155            let matched = Arc::clone(&matched);
156            thread::spawn(move || {
157                let mut local = 0usize;
158                loop {
159                    if ring.try_dequeue().is_some() {
160                        local += 1;
161                        matched.fetch_add(1, Ordering::Relaxed);
162                    } else if matched.load(Ordering::Relaxed) >= total {
163                        break;
164                    } else {
165                        std::hint::spin_loop();
166                    }
167                }
168                local
169            })
170        })
171        .collect();
172
173    for h in gateways {
174        h.join().unwrap();
175    }
176    let drained: usize = consumers.into_iter().map(|c| c.join().unwrap()).sum();
177    // cas_retries() is the contention read-out, deliberately not printed: it is
178    // a property of how the OS scheduled these threads on this run.
179    println!(
180        "  {shards} shards drained {drained} orders, ring empty: {}",
181        ring.is_empty()
182    );
183    assert_eq!(
184        drained, total,
185        "shards together drain every order exactly once"
186    );
187    assert_eq!(
188        ring.producer_index(),
189        ring.consumer_index(),
190        "every claimed slot was consumed"
191    );
192}
Source

pub fn try_dequeue(&self) -> Option<T>

Multi-consumer dequeue. Returns None if the ring is empty.

Examples found in repository?
examples/sample_app.rs (line 159)
126fn mpmc_sharded_match() {
127    use std::sync::atomic::{AtomicUsize, Ordering};
128
129    use subms_mpsc_queue::MpmcQueue;
130
131    println!("\n== mpmc: shard the match loop across several consumers ==");
132    let shards = 3usize;
133    let ring: Arc<MpmcQueue<u64>> = Arc::new(MpmcQueue::new(1_024));
134    let total = GATEWAYS * ORDERS_PER_GATEWAY;
135
136    let gateways: Vec<_> = (0..GATEWAYS)
137        .map(|g| {
138            let ring = Arc::clone(&ring);
139            thread::spawn(move || {
140                for seq in 0..ORDERS_PER_GATEWAY {
141                    let mut order = order_id(g, seq);
142                    while let Err(rejected) = ring.try_enqueue(order) {
143                        order = rejected;
144                        std::hint::spin_loop();
145                    }
146                }
147            })
148        })
149        .collect();
150
151    let matched = Arc::new(AtomicUsize::new(0));
152    let consumers: Vec<_> = (0..shards)
153        .map(|_| {
154            let ring = Arc::clone(&ring);
155            let matched = Arc::clone(&matched);
156            thread::spawn(move || {
157                let mut local = 0usize;
158                loop {
159                    if ring.try_dequeue().is_some() {
160                        local += 1;
161                        matched.fetch_add(1, Ordering::Relaxed);
162                    } else if matched.load(Ordering::Relaxed) >= total {
163                        break;
164                    } else {
165                        std::hint::spin_loop();
166                    }
167                }
168                local
169            })
170        })
171        .collect();
172
173    for h in gateways {
174        h.join().unwrap();
175    }
176    let drained: usize = consumers.into_iter().map(|c| c.join().unwrap()).sum();
177    // cas_retries() is the contention read-out, deliberately not printed: it is
178    // a property of how the OS scheduled these threads on this run.
179    println!(
180        "  {shards} shards drained {drained} orders, ring empty: {}",
181        ring.is_empty()
182    );
183    assert_eq!(
184        drained, total,
185        "shards together drain every order exactly once"
186    );
187    assert_eq!(
188        ring.producer_index(),
189        ring.consumer_index(),
190        "every claimed slot was consumed"
191    );
192}
More examples
Hide additional examples
examples/perf_features.rs (line 330)
195fn main() -> io::Result<()> {
196    let path = PathBuf::from(env!("CARGO_MANIFEST_DIR"))
197        .join("..")
198        .join(".subms")
199        .join("features")
200        .join("rust.json");
201    let existing = std::fs::read_to_string(&path).unwrap_or_default();
202    let mut manifest = SubMsFeatureManifest::load_str("rust", &existing);
203    // Stamp the box these numbers came from. The bench runs wherever it is
204    // invoked, so an unstamped manifest is indistinguishable from a fleet
205    // capture; the renderer will not publish one it cannot attribute.
206    let (source, instance) = SubMsP99Source::from_env();
207    manifest.set_p99_source(source, instance.as_deref());
208
209    // Optional, off by default: pin this thread to one core for the whole run.
210    // On a heterogeneous laptop the scheduler moves the bench between core
211    // clusters and every measurement lands in one of two clock states 1.31x
212    // apart - a spread three times wider than the deltas being classified, and
213    // large enough on its own to flip a feature between auxiliary and hot-path.
214    // Pinned, the same sweep repeats to within 1%. Left OFF by default because a
215    // fleet box isolates cores outside the process, and pinning from in here
216    // would override that placement with a core the orchestrator did not choose.
217    #[cfg(feature = "affinity")]
218    if let Some(core) = std::env::var("SUBMS_PIN").ok().and_then(|v| v.parse().ok()) {
219        let _ = subms_mpsc_queue::set_affinity(&[core]);
220    }
221
222    // Burn before the first measurement, not just before each one. Every
223    // `batched` call warms itself, but the FIRST measurement in the process pays
224    // a ramp the per-measurement warm sits inside rather than absorbs, and the
225    // sweep runs smallest-first: without this the base curve read 71800 / 45300 /
226    // 46500 ns, a 1.6x fall with size that is the process settling, not the
227    // queue.
228    {
229        let mut q = filled(CANON);
230        let start = std::time::Instant::now();
231        while (start.elapsed().as_nanos() as u64) < BURN_NANOS {
232            for i in 0..ITEMS_PER_SAMPLE {
233                q.push(i as u64);
234                black_box(q.try_pop());
235            }
236        }
237    }
238
239    // The baseline: the base queue's push + try_pop round trip. Swept as well as
240    // sampled, because whether queue depth moves the BASE op is the context
241    // every feature curve is read against.
242    let base_sweep = sweep("base/push+pop", |n| {
243        let mut q = filled(n);
244        batched(ITEMS_PER_SAMPLE, |i| {
245            q.push(i as u64);
246            black_box(q.try_pop());
247        })
248    });
249    let base_p50 = base_sweep
250        .iter()
251        .find(|(n, _)| *n == CANON)
252        .map_or(0, |(_, v)| *v);
253    eprintln!("base push+pop p50 per {ITEMS_PER_SAMPLE}-item sample: {base_p50}ns");
254
255    // ---------- bounded: fixed-capacity ring, backpressure on enqueue ----------
256    #[cfg(feature = "bounded")]
257    {
258        use subms_mpsc_queue::BoundedMpscQueue;
259        // Half full at every sweep point. Filled to a FIXED element count
260        // instead, the big rings would sit 98% empty and the enqueue would be
261        // measuring the fill fraction rather than the footprint.
262        fn ring(n: usize) -> BoundedMpscQueue<u64> {
263            let q = BoundedMpscQueue::new(n);
264            for i in 0..n / 2 {
265                let _ = q.try_enqueue(i as u64);
266            }
267            q
268        }
269        let sw = sweep("bounded/enqueue+dequeue", |n| {
270            let mut q = ring(n);
271            batched(ITEMS_PER_SAMPLE, |i| {
272                let _ = q.try_enqueue(i as u64);
273                black_box(q.try_dequeue());
274            })
275        });
276        let (cat, reason) = classify_feature(&sw, Some(base_p50), None);
277
278        let mut q = ring(CANON);
279        let mut p99 = BTreeMap::new();
280        p99.insert(
281            "enqueue".to_string(),
282            single(|st, i| {
283                st.time(|| {
284                    let _ = q.try_enqueue(i as u64);
285                });
286                black_box(q.try_dequeue());
287            }),
288        );
289        p99.insert(
290            "dequeue".to_string(),
291            single(|st, i| {
292                let _ = q.try_enqueue(i as u64);
293                st.time(|| black_box(q.try_dequeue()));
294            }),
295        );
296        // The reject path, which is the reason the feature exists. The ring is
297        // filled to capacity once, OUTSIDE the timed region; every timed call
298        // then takes the full branch and hands the value back to the caller.
299        let full: BoundedMpscQueue<u64> = BoundedMpscQueue::new(CANON);
300        while full.try_enqueue(0).is_ok() {}
301        p99.insert(
302            "enqueue_full".to_string(),
303            single(|st, i| {
304                st.time(|| {
305                    let _ = full.try_enqueue(i as u64);
306                });
307            }),
308        );
309        manifest.set_feature("bounded", cat, &p99, &reason);
310    }
311
312    // ---------- mpmc: bounded ring, sequence CAS on both ends ----------
313    #[cfg(feature = "mpmc")]
314    {
315        use subms_mpsc_queue::MpmcQueue;
316        fn ring(n: usize) -> MpmcQueue<u64> {
317            let q = MpmcQueue::new(n);
318            for i in 0..n / 2 {
319                let _ = q.try_enqueue(i as u64);
320            }
321            q
322        }
323        // Uncontended, so every CAS succeeds first try. That is the figure the
324        // category is about: what the multi-consumer claim costs a queue that is
325        // NOT contended, which is the state a well-sized pipeline runs in.
326        let sw = sweep("mpmc/enqueue+dequeue", |n| {
327            let q = ring(n);
328            batched(ITEMS_PER_SAMPLE, |i| {
329                let _ = q.try_enqueue(i as u64);
330                black_box(q.try_dequeue());
331            })
332        });
333        let (cat, reason) = classify_feature(&sw, Some(base_p50), None);
334
335        let q = ring(CANON);
336        let mut p99 = BTreeMap::new();
337        p99.insert(
338            "enqueue".to_string(),
339            single(|st, i| {
340                st.time(|| {
341                    let _ = q.try_enqueue(i as u64);
342                });
343                black_box(q.try_dequeue());
344            }),
345        );
346        p99.insert(
347            "dequeue".to_string(),
348            single(|st, i| {
349                let _ = q.try_enqueue(i as u64);
350                st.time(|| black_box(q.try_dequeue()));
351            }),
352        );
353        manifest.set_feature("mpmc", cat, &p99, &reason);
354    }
355
356    // ---------- batch: drain up to BATCH items behind one acquire fence ----------
357    #[cfg(feature = "batch")]
358    {
359        use subms_mpsc_queue::BatchMpscQueue;
360        fn filled_batch(n: usize) -> BatchMpscQueue<u64> {
361            let q = BatchMpscQueue::new();
362            for i in 0..n {
363                q.push(i as u64);
364            }
365            q
366        }
367        // A sample moves ITEMS_PER_SAMPLE items either way; only the call width
368        // differs. That is why the reps count is divided rather than the batch
369        // grown - growing it would sweep the batch size, and the number would
370        // stop being comparable to the base round trip.
371        let sw = sweep("batch/push+dequeue_batch", |n| {
372            let mut q = filled_batch(n);
373            let mut buf: Vec<Option<u64>> = (0..BATCH).map(|_| None).collect();
374            batched(ITEMS_PER_SAMPLE / BATCH, |i| {
375                for j in 0..BATCH {
376                    q.push((i + j) as u64);
377                }
378                black_box(q.try_dequeue_batch(&mut buf));
379            })
380        });
381        let (cat, reason) = classify_feature(&sw, Some(base_p50), None);
382
383        let mut q = filled_batch(CANON);
384        let mut buf: Vec<Option<u64>> = (0..BATCH).map(|_| None).collect();
385        let mut p99 = BTreeMap::new();
386        // The refill is outside the timed region: timing it would put a BATCH of
387        // pushes inside the drain's number and the stage would stop being a
388        // drain figure at all.
389        p99.insert(
390            "dequeue_batch".to_string(),
391            single(|st, i| {
392                st.time(|| black_box(q.try_dequeue_batch(&mut buf)));
393                for j in 0..BATCH {
394                    q.push((i + j) as u64);
395                }
396            }),
397        );
398        p99.insert(
399            "enqueue".to_string(),
400            single(|st, i| {
401                st.time(|| q.push(i as u64));
402                let _ = q.try_dequeue_batch(&mut buf[..1]);
403            }),
404        );
405        // The producer mirror: BATCH items published behind one head swap. The
406        // drain that puts the queue back is outside the timed region for the
407        // same reason the refill is above.
408        p99.insert(
409            "enqueue_batch".to_string(),
410            single(|st, i| {
411                let base = i as u64;
412                st.time(|| black_box(q.push_batch(base..base + BATCH as u64)));
413                let _ = q.try_dequeue_batch(&mut buf);
414            }),
415        );
416        manifest.set_feature("batch", cat, &p99, &reason);
417    }
418
419    // ---------- metrics: relaxed atomic counters around each op ----------
420    #[cfg(feature = "metrics")]
421    {
422        use subms_mpsc_queue::MetricsMpscQueue;
423        fn filled_metrics(n: usize) -> MetricsMpscQueue<u64> {
424            let q = MetricsMpscQueue::new();
425            for i in 0..n {
426                q.push(i as u64);
427            }
428            q
429        }
430        let sw = sweep("metrics/push+pop", |n| {
431            let mut q = filled_metrics(n);
432            batched(ITEMS_PER_SAMPLE, |i| {
433                q.push(i as u64);
434                black_box(q.try_pop());
435            })
436        });
437        let (cat, reason) = classify_feature(&sw, Some(base_p50), None);
438
439        let mut q = filled_metrics(CANON);
440        let mut p99 = BTreeMap::new();
441        p99.insert(
442            "enqueue".to_string(),
443            single(|st, i| {
444                st.time(|| q.push(i as u64));
445                black_box(q.try_pop());
446            }),
447        );
448        p99.insert(
449            "dequeue".to_string(),
450            single(|st, i| {
451                q.push(i as u64);
452                st.time(|| black_box(q.try_pop()));
453            }),
454        );
455        p99.insert(
456            "snapshot".to_string(),
457            single(|st, _| {
458                st.time(|| black_box(q.snapshot()));
459            }),
460        );
461        manifest.set_feature("metrics", cat, &p99, &reason);
462    }
463
464    // ---------- affinity: pin the calling thread, once, at startup ----------
465    // Runs LAST because measuring it pins THIS process to core 0, and every
466    // number taken afterwards would be a number taken on one core.
467    #[cfg(feature = "affinity")]
468    {
469        use subms_mpsc_queue::set_affinity;
470        // Swept over the same axis to show what it is: a call that touches no
471        // queue state and cannot move with queue size. PINNED auxiliary rather
472        // than left to the base-delta test, which would see a syscall costing
473        // more than an enqueue and call it hot-path. It is not on the hot path at
474        // any price - `set_affinity` is called once per thread at startup and
475        // appears in neither `push` nor `try_pop`. The two ports are not even
476        // measuring the same thing: Rust issues a real `SetThreadAffinityMask` /
477        // `sched_setaffinity`, while the Java sibling validates its argument and
478        // returns UNSUPPORTED because the stock JDK has no pinning API. The
479        // per-call figure, not the sample, is the interpretable one and it is in
480        // `p99ByStage`.
481        let sw = sweep("affinity/set_affinity", |_| {
482            batched(ITEMS_PER_SAMPLE, |_| {
483                let _ = set_affinity(&[0]);
484            })
485        });
486        let (cat, reason) = classify_feature(
487            &sw,
488            Some(base_p50),
489            Some(subms::SubMsFeatureCategory::Auxiliary),
490        );
491
492        let mut p99 = BTreeMap::new();
493        p99.insert(
494            "set_affinity".to_string(),
495            single(|st, _| {
496                st.time(|| {
497                    let _ = set_affinity(&[0]);
498                });
499            }),
500        );
501        manifest.set_feature("affinity", cat, &p99, &reason);
502
503        let cores: Vec<usize> = (0..std::thread::available_parallelism()
504            .map_or(1, std::num::NonZeroUsize::get))
505            .collect();
506        let _ = set_affinity(&cores);
507    }
508
509    std::fs::create_dir_all(path.parent().unwrap())?;
510    std::fs::write(&path, manifest.to_json())?;
511    io::stdout().write_all(manifest.to_json().as_bytes())?;
512    Ok(())
513}
Source

pub fn clear(&self) -> usize

Drop everything currently readable and return the count. Any consumer may call it, and other consumers keep draining alongside, so the count is this caller’s share rather than the queue’s total.

Source

pub fn len(&self) -> usize

Approximate length.

Source

pub fn is_empty(&self) -> bool

Examples found in repository?
examples/sample_app.rs (line 181)
126fn mpmc_sharded_match() {
127    use std::sync::atomic::{AtomicUsize, Ordering};
128
129    use subms_mpsc_queue::MpmcQueue;
130
131    println!("\n== mpmc: shard the match loop across several consumers ==");
132    let shards = 3usize;
133    let ring: Arc<MpmcQueue<u64>> = Arc::new(MpmcQueue::new(1_024));
134    let total = GATEWAYS * ORDERS_PER_GATEWAY;
135
136    let gateways: Vec<_> = (0..GATEWAYS)
137        .map(|g| {
138            let ring = Arc::clone(&ring);
139            thread::spawn(move || {
140                for seq in 0..ORDERS_PER_GATEWAY {
141                    let mut order = order_id(g, seq);
142                    while let Err(rejected) = ring.try_enqueue(order) {
143                        order = rejected;
144                        std::hint::spin_loop();
145                    }
146                }
147            })
148        })
149        .collect();
150
151    let matched = Arc::new(AtomicUsize::new(0));
152    let consumers: Vec<_> = (0..shards)
153        .map(|_| {
154            let ring = Arc::clone(&ring);
155            let matched = Arc::clone(&matched);
156            thread::spawn(move || {
157                let mut local = 0usize;
158                loop {
159                    if ring.try_dequeue().is_some() {
160                        local += 1;
161                        matched.fetch_add(1, Ordering::Relaxed);
162                    } else if matched.load(Ordering::Relaxed) >= total {
163                        break;
164                    } else {
165                        std::hint::spin_loop();
166                    }
167                }
168                local
169            })
170        })
171        .collect();
172
173    for h in gateways {
174        h.join().unwrap();
175    }
176    let drained: usize = consumers.into_iter().map(|c| c.join().unwrap()).sum();
177    // cas_retries() is the contention read-out, deliberately not printed: it is
178    // a property of how the OS scheduled these threads on this run.
179    println!(
180        "  {shards} shards drained {drained} orders, ring empty: {}",
181        ring.is_empty()
182    );
183    assert_eq!(
184        drained, total,
185        "shards together drain every order exactly once"
186    );
187    assert_eq!(
188        ring.producer_index(),
189        ring.consumer_index(),
190        "every claimed slot was consumed"
191    );
192}
Source

pub fn is_full(&self) -> bool

Best-effort fullness. Stale the instant any consumer drains a slot.

Trait Implementations§

Source§

impl<T> Drop for MpmcQueue<T>

Source§

fn drop(&mut self)

Executes the destructor for this type. Read more
Source§

fn pin_drop(self: Pin<&mut Self>)

🔬This is a nightly-only experimental API. (pin_ergonomics)
Execute the destructor for this type, but different to Drop::drop, it requires self to be pinned. Read more
Source§

impl<T: Send> Send for MpmcQueue<T>

Source§

impl<T: Send> Sync for MpmcQueue<T>

Auto Trait Implementations§

§

impl<T> !Freeze for MpmcQueue<T>

§

impl<T> !RefUnwindSafe for MpmcQueue<T>

§

impl<T> Unpin for MpmcQueue<T>

§

impl<T> UnsafeUnpin for MpmcQueue<T>

§

impl<T> UnwindSafe for MpmcQueue<T>
where T: UnwindSafe,

Blanket Implementations§

Source§

impl<T> Any for T
where T: 'static + ?Sized,

Source§

fn type_id(&self) -> TypeId

Gets the TypeId of self. Read more
Source§

impl<T> Borrow<T> for T
where T: ?Sized,

Source§

fn borrow(&self) -> &T

Immutably borrows from an owned value. Read more
Source§

impl<T> BorrowMut<T> for T
where T: ?Sized,

Source§

fn borrow_mut(&mut self) -> &mut T

Mutably borrows from an owned value. Read more
Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

Source§

impl<T, U> Into<U> for T
where U: From<T>,

Source§

fn into(self) -> U

Calls U::from(self).

That is, this conversion is whatever the implementation of From<T> for U chooses to do.

Source§

impl<T, U> TryFrom<U> for T
where U: Into<T>,

Source§

type Error = Infallible

The type returned in the event of a conversion error.
Source§

fn try_from(value: U) -> Result<T, <T as TryFrom<U>>::Error>

Performs the conversion.
Source§

impl<T, U> TryInto<U> for T
where U: TryFrom<T>,

Source§

type Error = <U as TryFrom<T>>::Error

The type returned in the event of a conversion error.
Source§

fn try_into(self) -> Result<U, <U as TryFrom<T>>::Error>

Performs the conversion.