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>
impl<T> MpmcQueue<T>
Sourcepub fn new(capacity: usize) -> Self
pub fn new(capacity: usize) -> Self
Examples found in repository?
More 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}pub fn capacity(&self) -> usize
Sourcepub fn producer_index(&self) -> usize
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}Sourcepub fn consumer_index(&self) -> usize
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}Sourcepub fn cas_retries(&self) -> u64
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.
Sourcepub fn try_enqueue(&self, value: T) -> Result<(), T>
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
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}Sourcepub fn try_dequeue(&self) -> Option<T>
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
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}Sourcepub fn clear(&self) -> usize
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.
Sourcepub fn is_empty(&self) -> bool
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}Trait Implementations§
impl<T: Send> Send for MpmcQueue<T>
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> BorrowMut<T> for Twhere
T: ?Sized,
impl<T> BorrowMut<T> for Twhere
T: ?Sized,
Source§fn borrow_mut(&mut self) -> &mut T
fn borrow_mut(&mut self) -> &mut T
Mutably borrows from an owned value. Read more