1use crate::mechanics::{DeckMechanicalControl, MotorMode};
9use crate::scratch_gate::ScratchPreset;
10use crate::timed_control::{PlayerControl, TimedPlayerControl};
11use std::cell::Cell;
12use std::error::Error;
13use std::fmt;
14use std::marker::PhantomData;
15use std::sync::atomic::{AtomicBool, AtomicU8, AtomicUsize, Ordering};
16use std::sync::Arc;
17
18pub trait SpscCodec<T: Copy, const BYTES: usize>: Send + Sync + 'static {
24 fn encode(value: T, destination: &mut [u8; BYTES]);
25 fn decode(source: &[u8; BYTES]) -> Option<T>;
26}
27
28#[derive(Debug, Clone, Copy, PartialEq, Eq)]
29pub enum SpscCreateError {
30 ZeroCapacity,
31 ZeroPayloadBytes,
32 CapacityTooLarge,
33}
34
35impl fmt::Display for SpscCreateError {
36 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
37 match self {
38 Self::ZeroCapacity => formatter.write_str("SPSC mailbox capacity must be positive"),
39 Self::ZeroPayloadBytes => {
40 formatter.write_str("SPSC mailbox payload size must be positive")
41 }
42 Self::CapacityTooLarge => formatter.write_str("SPSC mailbox capacity is too large"),
43 }
44 }
45}
46
47impl Error for SpscCreateError {}
48
49#[derive(Debug, Clone, Copy, PartialEq, Eq)]
51pub enum SpscPushError<T> {
52 Full(T),
53 Disconnected(T),
54}
55
56impl<T> SpscPushError<T> {
57 pub fn into_inner(self) -> T {
58 match self {
59 Self::Full(value) | Self::Disconnected(value) => value,
60 }
61 }
62}
63
64impl<T> fmt::Display for SpscPushError<T> {
65 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
66 match self {
67 Self::Full(_) => formatter.write_str("SPSC mailbox is full"),
68 Self::Disconnected(_) => formatter.write_str("SPSC mailbox consumer is disconnected"),
69 }
70 }
71}
72
73impl<T: fmt::Debug> Error for SpscPushError<T> {}
74
75#[derive(Debug, Clone, Copy, PartialEq, Eq)]
76pub enum SpscPopError {
77 Empty,
78 Disconnected,
79 InvalidEncoding,
80}
81
82impl fmt::Display for SpscPopError {
83 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
84 match self {
85 Self::Empty => formatter.write_str("SPSC mailbox is empty"),
86 Self::Disconnected => formatter.write_str("SPSC mailbox producer is disconnected"),
87 Self::InvalidEncoding => formatter.write_str("SPSC mailbox payload is invalid"),
88 }
89 }
90}
91
92impl Error for SpscPopError {}
93
94#[derive(Debug, Clone, Copy, PartialEq, Eq)]
95pub enum SpscCommitOutcome {
96 Value,
97 InvalidEncoding,
98}
99
100#[derive(Debug, Clone, Copy, PartialEq, Eq)]
101pub enum SpscCommitError {
102 NothingPeeked,
103}
104
105impl fmt::Display for SpscCommitError {
106 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
107 match self {
108 Self::NothingPeeked => formatter.write_str("SPSC mailbox has no peeked payload"),
109 }
110 }
111}
112
113impl Error for SpscCommitError {}
114
115#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
120pub struct SpscCounters {
121 pub accepted_pushes: usize,
122 pub rejected_full_pushes: usize,
123 pub rejected_disconnected_pushes: usize,
124 pub successful_pops: usize,
125 pub empty_pops: usize,
126 pub disconnected_pops: usize,
127 pub invalid_payloads: usize,
128}
129
130#[repr(align(64))]
131struct CacheLine<T>(T);
132
133#[repr(align(64))]
134struct Slot<const BYTES: usize> {
135 bytes: [AtomicU8; BYTES],
136}
137
138impl<const BYTES: usize> Slot<BYTES> {
139 fn new() -> Self {
140 Self {
141 bytes: std::array::from_fn(|_| AtomicU8::new(0)),
142 }
143 }
144
145 fn store(&self, source: &[u8; BYTES]) {
146 for (destination, source) in self.bytes.iter().zip(source.iter().copied()) {
147 destination.store(source, Ordering::Relaxed);
148 }
149 }
150
151 fn load(&self) -> [u8; BYTES] {
152 std::array::from_fn(|index| self.bytes[index].load(Ordering::Relaxed))
153 }
154}
155
156struct ProducerCounters {
157 accepted_pushes: AtomicUsize,
158 rejected_full_pushes: AtomicUsize,
159 rejected_disconnected_pushes: AtomicUsize,
160}
161
162struct ConsumerCounters {
163 successful_pops: AtomicUsize,
164 empty_pops: AtomicUsize,
165 disconnected_pops: AtomicUsize,
166 invalid_payloads: AtomicUsize,
167}
168
169struct Shared<T, Codec, const BYTES: usize> {
170 slots: Box<[Slot<BYTES>]>,
171 write_count: CacheLine<AtomicUsize>,
172 read_count: CacheLine<AtomicUsize>,
173 producer_alive: AtomicBool,
174 consumer_alive: AtomicBool,
175 producer_counters: CacheLine<ProducerCounters>,
176 consumer_counters: CacheLine<ConsumerCounters>,
177 marker: PhantomData<(T, Codec)>,
178}
179
180impl<T, Codec, const BYTES: usize> Shared<T, Codec, BYTES> {
181 fn capacity(&self) -> usize {
182 self.slots.len()
183 }
184
185 fn counters(&self) -> SpscCounters {
186 SpscCounters {
187 accepted_pushes: self
188 .producer_counters
189 .0
190 .accepted_pushes
191 .load(Ordering::Relaxed),
192 rejected_full_pushes: self
193 .producer_counters
194 .0
195 .rejected_full_pushes
196 .load(Ordering::Relaxed),
197 rejected_disconnected_pushes: self
198 .producer_counters
199 .0
200 .rejected_disconnected_pushes
201 .load(Ordering::Relaxed),
202 successful_pops: self
203 .consumer_counters
204 .0
205 .successful_pops
206 .load(Ordering::Relaxed),
207 empty_pops: self.consumer_counters.0.empty_pops.load(Ordering::Relaxed),
208 disconnected_pops: self
209 .consumer_counters
210 .0
211 .disconnected_pops
212 .load(Ordering::Relaxed),
213 invalid_payloads: self
214 .consumer_counters
215 .0
216 .invalid_payloads
217 .load(Ordering::Relaxed),
218 }
219 }
220
221 fn approximate_len(&self) -> usize {
222 let written = self.write_count.0.load(Ordering::Acquire);
223 let read = self.read_count.0.load(Ordering::Acquire);
224 written.wrapping_sub(read).min(self.capacity())
225 }
226}
227
228pub struct SpscProducer<T, Codec, const BYTES: usize> {
233 shared: Arc<Shared<T, Codec, BYTES>>,
234 write_count: usize,
235 next_slot: usize,
236 not_sync: PhantomData<Cell<()>>,
237}
238
239impl<T, Codec, const BYTES: usize> SpscProducer<T, Codec, BYTES>
240where
241 T: Copy + Send + Sync + 'static,
242 Codec: SpscCodec<T, BYTES>,
243{
244 pub fn try_push(&mut self, value: T) -> Result<(), SpscPushError<T>> {
245 if !self.shared.consumer_alive.load(Ordering::Acquire) {
246 increment(&self.shared.producer_counters.0.rejected_disconnected_pushes);
247 return Err(SpscPushError::Disconnected(value));
248 }
249
250 let read_count = self.shared.read_count.0.load(Ordering::Acquire);
251 if self.write_count.wrapping_sub(read_count) >= self.shared.capacity() {
252 increment(&self.shared.producer_counters.0.rejected_full_pushes);
253 return Err(SpscPushError::Full(value));
254 }
255
256 let mut encoded = [0; BYTES];
257 Codec::encode(value, &mut encoded);
258 self.shared.slots[self.next_slot].store(&encoded);
259
260 self.write_count = self.write_count.wrapping_add(1);
261 self.next_slot += 1;
262 if self.next_slot == self.shared.capacity() {
263 self.next_slot = 0;
264 }
265
266 self.shared
268 .write_count
269 .0
270 .store(self.write_count, Ordering::Release);
271 increment(&self.shared.producer_counters.0.accepted_pushes);
272 Ok(())
273 }
274
275 pub fn capacity(&self) -> usize {
276 self.shared.capacity()
277 }
278
279 pub fn approximate_len(&self) -> usize {
280 self.shared.approximate_len()
281 }
282
283 pub fn counters(&self) -> SpscCounters {
284 self.shared.counters()
285 }
286
287 pub fn is_consumer_connected(&self) -> bool {
288 self.shared.consumer_alive.load(Ordering::Acquire)
289 }
290}
291
292impl<T, Codec, const BYTES: usize> Drop for SpscProducer<T, Codec, BYTES> {
293 fn drop(&mut self) {
294 self.shared.producer_alive.store(false, Ordering::Release);
295 }
296}
297
298#[derive(Clone, Copy)]
299enum PeekedPayload<T> {
300 None,
301 Value(T),
302 InvalidEncoding,
303}
304
305pub struct SpscConsumer<T, Codec, const BYTES: usize> {
310 shared: Arc<Shared<T, Codec, BYTES>>,
311 read_count: usize,
312 next_slot: usize,
313 peeked: PeekedPayload<T>,
314 not_sync: PhantomData<Cell<()>>,
315}
316
317pub type SpscEndpoints<T, Codec, const BYTES: usize> =
318 (SpscProducer<T, Codec, BYTES>, SpscConsumer<T, Codec, BYTES>);
319
320impl<T, Codec, const BYTES: usize> SpscConsumer<T, Codec, BYTES>
321where
322 T: Copy + Send + Sync + 'static,
323 Codec: SpscCodec<T, BYTES>,
324{
325 pub fn try_peek(&mut self) -> Result<T, SpscPopError> {
330 match self.peeked {
331 PeekedPayload::Value(value) => return Ok(value),
332 PeekedPayload::InvalidEncoding => return Err(SpscPopError::InvalidEncoding),
333 PeekedPayload::None => {}
334 }
335
336 self.require_head()?;
337 let encoded = self.shared.slots[self.next_slot].load();
338 match Codec::decode(&encoded) {
339 Some(value) => {
340 self.peeked = PeekedPayload::Value(value);
341 Ok(value)
342 }
343 None => {
344 self.peeked = PeekedPayload::InvalidEncoding;
345 Err(SpscPopError::InvalidEncoding)
346 }
347 }
348 }
349
350 pub fn commit_peeked(&mut self) -> Result<SpscCommitOutcome, SpscCommitError> {
355 self.commit_cached().ok_or(SpscCommitError::NothingPeeked)
356 }
357
358 pub fn try_pop(&mut self) -> Result<T, SpscPopError> {
359 let value = match self.try_peek() {
360 Ok(value) => value,
361 Err(SpscPopError::Empty) => {
362 increment(&self.shared.consumer_counters.0.empty_pops);
363 return Err(SpscPopError::Empty);
364 }
365 Err(SpscPopError::Disconnected) => {
366 increment(&self.shared.consumer_counters.0.disconnected_pops);
367 return Err(SpscPopError::Disconnected);
368 }
369 Err(SpscPopError::InvalidEncoding) => {
370 let outcome = self
371 .commit_cached()
372 .expect("an invalid peek must reserve the queue head");
373 debug_assert_eq!(outcome, SpscCommitOutcome::InvalidEncoding);
374 return Err(SpscPopError::InvalidEncoding);
375 }
376 };
377
378 let outcome = self
379 .commit_cached()
380 .expect("a successful peek must reserve the queue head");
381 debug_assert_eq!(outcome, SpscCommitOutcome::Value);
382 Ok(value)
383 }
384
385 pub fn capacity(&self) -> usize {
386 self.shared.capacity()
387 }
388
389 pub fn approximate_len(&self) -> usize {
390 self.shared.approximate_len()
391 }
392
393 pub fn counters(&self) -> SpscCounters {
394 self.shared.counters()
395 }
396
397 pub fn is_producer_connected(&self) -> bool {
398 self.shared.producer_alive.load(Ordering::Acquire)
399 }
400
401 fn require_head(&self) -> Result<(), SpscPopError> {
402 let mut write_count = self.shared.write_count.0.load(Ordering::Acquire);
403 if self.read_count != write_count {
404 return Ok(());
405 }
406 if self.shared.producer_alive.load(Ordering::Acquire) {
407 return Err(SpscPopError::Empty);
408 }
409
410 write_count = self.shared.write_count.0.load(Ordering::Acquire);
412 if self.read_count == write_count {
413 Err(SpscPopError::Disconnected)
414 } else {
415 Ok(())
416 }
417 }
418
419 fn commit_cached(&mut self) -> Option<SpscCommitOutcome> {
420 let peeked = std::mem::replace(&mut self.peeked, PeekedPayload::None);
421 let outcome = match peeked {
422 PeekedPayload::None => return None,
423 PeekedPayload::Value(_) => SpscCommitOutcome::Value,
424 PeekedPayload::InvalidEncoding => SpscCommitOutcome::InvalidEncoding,
425 };
426
427 self.read_count = self.read_count.wrapping_add(1);
428 self.next_slot += 1;
429 if self.next_slot == self.shared.capacity() {
430 self.next_slot = 0;
431 }
432
433 self.shared
435 .read_count
436 .0
437 .store(self.read_count, Ordering::Release);
438 match outcome {
439 SpscCommitOutcome::Value => {
440 increment(&self.shared.consumer_counters.0.successful_pops);
441 }
442 SpscCommitOutcome::InvalidEncoding => {
443 increment(&self.shared.consumer_counters.0.invalid_payloads);
444 }
445 }
446 Some(outcome)
447 }
448}
449
450impl<T, Codec, const BYTES: usize> Drop for SpscConsumer<T, Codec, BYTES> {
451 fn drop(&mut self) {
452 self.shared.consumer_alive.store(false, Ordering::Release);
453 }
454}
455
456pub fn spsc_mailbox<T, Codec, const BYTES: usize>(
458 capacity: usize,
459) -> Result<SpscEndpoints<T, Codec, BYTES>, SpscCreateError>
460where
461 T: Copy + Send + Sync + 'static,
462 Codec: SpscCodec<T, BYTES>,
463{
464 if capacity == 0 {
465 return Err(SpscCreateError::ZeroCapacity);
466 }
467 if BYTES == 0 {
468 return Err(SpscCreateError::ZeroPayloadBytes);
469 }
470 if capacity > usize::MAX / 2 {
471 return Err(SpscCreateError::CapacityTooLarge);
472 }
473
474 let slots = std::iter::repeat_with(Slot::new)
475 .take(capacity)
476 .collect::<Vec<_>>()
477 .into_boxed_slice();
478 let shared = Arc::new(Shared {
479 slots,
480 write_count: CacheLine(AtomicUsize::new(0)),
481 read_count: CacheLine(AtomicUsize::new(0)),
482 producer_alive: AtomicBool::new(true),
483 consumer_alive: AtomicBool::new(true),
484 producer_counters: CacheLine(ProducerCounters {
485 accepted_pushes: AtomicUsize::new(0),
486 rejected_full_pushes: AtomicUsize::new(0),
487 rejected_disconnected_pushes: AtomicUsize::new(0),
488 }),
489 consumer_counters: CacheLine(ConsumerCounters {
490 successful_pops: AtomicUsize::new(0),
491 empty_pops: AtomicUsize::new(0),
492 disconnected_pops: AtomicUsize::new(0),
493 invalid_payloads: AtomicUsize::new(0),
494 }),
495 marker: PhantomData,
496 });
497
498 Ok((
499 SpscProducer {
500 shared: Arc::clone(&shared),
501 write_count: 0,
502 next_slot: 0,
503 not_sync: PhantomData,
504 },
505 SpscConsumer {
506 shared,
507 read_count: 0,
508 next_slot: 0,
509 peeked: PeekedPayload::None,
510 not_sync: PhantomData,
511 },
512 ))
513}
514
515fn increment(counter: &AtomicUsize) {
517 let next = counter.load(Ordering::Relaxed).wrapping_add(1);
518 counter.store(next, Ordering::Relaxed);
519}
520
521pub const TIMED_PLAYER_CONTROL_BYTES: usize = 78;
522
523#[derive(Debug, Clone, Copy, Default)]
524pub struct TimedPlayerControlCodec;
525
526impl SpscCodec<TimedPlayerControl, TIMED_PLAYER_CONTROL_BYTES> for TimedPlayerControlCodec {
527 fn encode(value: TimedPlayerControl, destination: &mut [u8; TIMED_PLAYER_CONTROL_BYTES]) {
528 let mut cursor = 0;
529 write_u64(destination, &mut cursor, value.absolute_frame);
530 write_u64(destination, &mut cursor, value.sequence);
531 destination[cursor] = match value.control.deck.motor_mode {
532 MotorMode::Off => 0,
533 MotorMode::Servo => 1,
534 MotorMode::Brake => 2,
535 };
536 cursor += 1;
537 write_f64(
538 destination,
539 &mut cursor,
540 value.control.deck.motor_target_angular_velocity_rad_s,
541 );
542 destination[cursor] = u8::from(value.control.deck.hand_contact);
543 cursor += 1;
544 match value.control.deck.hand_target_angle_rad {
545 None => {
546 destination[cursor] = 0;
547 cursor += 1;
548 write_f64(destination, &mut cursor, 0.0);
549 }
550 Some(angle) => {
551 destination[cursor] = 1;
552 cursor += 1;
553 write_f64(destination, &mut cursor, angle);
554 }
555 }
556 write_f64(
557 destination,
558 &mut cursor,
559 value.control.deck.hand_target_angular_velocity_rad_s,
560 );
561 write_f64(
562 destination,
563 &mut cursor,
564 value.control.deck.hand_normal_force_n,
565 );
566 write_f64(
567 destination,
568 &mut cursor,
569 value.control.deck.hand_contact_radius_m,
570 );
571 write_f64(
572 destination,
573 &mut cursor,
574 value.control.deck.stylus_torque_nm,
575 );
576 destination[cursor] = u8::from(value.control.stylus_lowered);
577 cursor += 1;
578 destination[cursor] = value.control.scratch_preset.id();
579 cursor += 1;
580 destination[cursor] = value.control.scratch_clicks;
581 cursor += 1;
582 write_f64(
583 destination,
584 &mut cursor,
585 value.control.manual_crossfader_gain,
586 );
587 debug_assert_eq!(cursor, TIMED_PLAYER_CONTROL_BYTES);
588 }
589
590 fn decode(source: &[u8; TIMED_PLAYER_CONTROL_BYTES]) -> Option<TimedPlayerControl> {
591 let mut cursor = 0;
592 let absolute_frame = read_u64(source, &mut cursor);
593 let sequence = read_u64(source, &mut cursor);
594 let motor_mode = match source[cursor] {
595 0 => MotorMode::Off,
596 1 => MotorMode::Servo,
597 2 => MotorMode::Brake,
598 _ => return None,
599 };
600 cursor += 1;
601 let motor_target_angular_velocity_rad_s = read_f64(source, &mut cursor);
602 let hand_contact = match source[cursor] {
603 0 => false,
604 1 => true,
605 _ => return None,
606 };
607 cursor += 1;
608 let hand_target_angle_rad = match source[cursor] {
609 0 => {
610 cursor += 1;
611 let _ = read_f64(source, &mut cursor);
612 None
613 }
614 1 => {
615 cursor += 1;
616 Some(read_f64(source, &mut cursor))
617 }
618 _ => return None,
619 };
620 let hand_target_angular_velocity_rad_s = read_f64(source, &mut cursor);
621 let hand_normal_force_n = read_f64(source, &mut cursor);
622 let hand_contact_radius_m = read_f64(source, &mut cursor);
623 let stylus_torque_nm = read_f64(source, &mut cursor);
624 let stylus_lowered = match source[cursor] {
625 0 => false,
626 1 => true,
627 _ => return None,
628 };
629 cursor += 1;
630 let scratch_preset = ScratchPreset::from_id(source[cursor])?;
631 cursor += 1;
632 let scratch_clicks = source[cursor];
633 cursor += 1;
634 let manual_crossfader_gain = read_f64(source, &mut cursor);
635 debug_assert_eq!(cursor, TIMED_PLAYER_CONTROL_BYTES);
636
637 Some(TimedPlayerControl::new(
638 absolute_frame,
639 sequence,
640 PlayerControl::new(
641 DeckMechanicalControl {
642 motor_mode,
643 motor_target_angular_velocity_rad_s,
644 hand_contact,
645 hand_target_angle_rad,
646 hand_target_angular_velocity_rad_s,
647 hand_normal_force_n,
648 hand_contact_radius_m,
649 stylus_torque_nm,
650 },
651 stylus_lowered,
652 )
653 .with_scratch(scratch_preset, scratch_clicks, manual_crossfader_gain),
654 ))
655 }
656}
657
658pub type TimedPlayerControlProducer =
659 SpscProducer<TimedPlayerControl, TimedPlayerControlCodec, TIMED_PLAYER_CONTROL_BYTES>;
660pub type TimedPlayerControlConsumer =
661 SpscConsumer<TimedPlayerControl, TimedPlayerControlCodec, TIMED_PLAYER_CONTROL_BYTES>;
662
663pub fn timed_player_control_mailbox(
664 capacity: usize,
665) -> Result<(TimedPlayerControlProducer, TimedPlayerControlConsumer), SpscCreateError> {
666 spsc_mailbox::<TimedPlayerControl, TimedPlayerControlCodec, TIMED_PLAYER_CONTROL_BYTES>(
667 capacity,
668 )
669}
670
671fn write_u64<const BYTES: usize>(destination: &mut [u8; BYTES], cursor: &mut usize, value: u64) {
672 let end = *cursor + 8;
673 destination[*cursor..end].copy_from_slice(&value.to_le_bytes());
674 *cursor = end;
675}
676
677fn read_u64<const BYTES: usize>(source: &[u8; BYTES], cursor: &mut usize) -> u64 {
678 let end = *cursor + 8;
679 let mut bytes = [0; 8];
680 bytes.copy_from_slice(&source[*cursor..end]);
681 *cursor = end;
682 u64::from_le_bytes(bytes)
683}
684
685fn write_f64<const BYTES: usize>(destination: &mut [u8; BYTES], cursor: &mut usize, value: f64) {
686 write_u64(destination, cursor, value.to_bits());
687}
688
689fn read_f64<const BYTES: usize>(source: &[u8; BYTES], cursor: &mut usize) -> f64 {
690 f64::from_bits(read_u64(source, cursor))
691}
692
693#[cfg(test)]
694mod tests {
695 use super::*;
696 use crate::timed_control::{ControlTimelinePushError, PlayerControlTimeline};
697 use std::thread;
698
699 #[derive(Debug, Clone, Copy, PartialEq, Eq)]
700 struct TestEvent {
701 sequence: u64,
702 value: u64,
703 }
704
705 struct TestCodec;
706
707 impl SpscCodec<TestEvent, 16> for TestCodec {
708 fn encode(value: TestEvent, destination: &mut [u8; 16]) {
709 destination[..8].copy_from_slice(&value.sequence.to_le_bytes());
710 destination[8..].copy_from_slice(&value.value.to_le_bytes());
711 }
712
713 fn decode(source: &[u8; 16]) -> Option<TestEvent> {
714 let mut sequence = [0; 8];
715 let mut value = [0; 8];
716 sequence.copy_from_slice(&source[..8]);
717 value.copy_from_slice(&source[8..]);
718 Some(TestEvent {
719 sequence: u64::from_le_bytes(sequence),
720 value: u64::from_le_bytes(value),
721 })
722 }
723 }
724
725 fn mailbox(
726 capacity: usize,
727 ) -> (
728 SpscProducer<TestEvent, TestCodec, 16>,
729 SpscConsumer<TestEvent, TestCodec, 16>,
730 ) {
731 spsc_mailbox(capacity).unwrap()
732 }
733
734 fn timed_event(absolute_frame: u64, sequence: u64) -> TimedPlayerControl {
735 TimedPlayerControl::new(
736 absolute_frame,
737 sequence,
738 PlayerControl::new(DeckMechanicalControl::default(), true),
739 )
740 }
741
742 #[test]
743 fn rejects_invalid_capacities() {
744 assert!(matches!(
745 spsc_mailbox::<TestEvent, TestCodec, 16>(0),
746 Err(SpscCreateError::ZeroCapacity)
747 ));
748
749 struct EmptyCodec;
750 impl SpscCodec<TestEvent, 0> for EmptyCodec {
751 fn encode(_value: TestEvent, _destination: &mut [u8; 0]) {}
752 fn decode(_source: &[u8; 0]) -> Option<TestEvent> {
753 None
754 }
755 }
756 assert!(matches!(
757 spsc_mailbox::<TestEvent, EmptyCodec, 0>(1),
758 Err(SpscCreateError::ZeroPayloadBytes)
759 ));
760 }
761
762 #[test]
763 fn rejects_full_push_without_overwriting_values() {
764 let (mut producer, mut consumer) = mailbox(2);
765 let first = TestEvent {
766 sequence: 1,
767 value: 11,
768 };
769 let second = TestEvent {
770 sequence: 2,
771 value: 22,
772 };
773 let rejected = TestEvent {
774 sequence: 3,
775 value: 33,
776 };
777
778 producer.try_push(first).unwrap();
779 producer.try_push(second).unwrap();
780 assert_eq!(
781 producer.try_push(rejected),
782 Err(SpscPushError::Full(rejected))
783 );
784 assert_eq!(consumer.try_pop(), Ok(first));
785 assert_eq!(consumer.try_pop(), Ok(second));
786 assert_eq!(consumer.try_pop(), Err(SpscPopError::Empty));
787
788 let counters = consumer.counters();
789 assert_eq!(counters.accepted_pushes, 2);
790 assert_eq!(counters.rejected_full_pushes, 1);
791 assert_eq!(counters.successful_pops, 2);
792 assert_eq!(counters.empty_pops, 1);
793 }
794
795 #[test]
796 fn repeated_peek_waits_for_explicit_commit() {
797 let (mut producer, mut consumer) = mailbox(2);
798 let event = TestEvent {
799 sequence: 1,
800 value: 11,
801 };
802 producer.try_push(event).unwrap();
803
804 assert_eq!(consumer.try_peek(), Ok(event));
805 assert_eq!(consumer.try_peek(), Ok(event));
806 assert_eq!(consumer.approximate_len(), 1);
807 assert_eq!(consumer.counters().successful_pops, 0);
808 assert_eq!(consumer.commit_peeked(), Ok(SpscCommitOutcome::Value));
809 assert_eq!(consumer.approximate_len(), 0);
810 assert_eq!(consumer.counters().successful_pops, 1);
811 assert_eq!(
812 consumer.commit_peeked(),
813 Err(SpscCommitError::NothingPeeked)
814 );
815 }
816
817 #[test]
818 fn full_timeline_backpressure_does_not_lose_mailbox_event() {
819 let initial = PlayerControl::default();
820 let mut timeline = PlayerControlTimeline::new(1, 0, initial).unwrap();
821 timeline.enqueue(timed_event(0, 1)).unwrap();
822 let incoming = timed_event(1, 2);
823 let (mut producer, mut consumer) = timed_player_control_mailbox(1).unwrap();
824 producer.try_push(incoming).unwrap();
825
826 let peeked = consumer.try_peek().unwrap();
827 assert_eq!(
828 timeline.enqueue(peeked),
829 Err(ControlTimelinePushError::Full { capacity: 1 })
830 );
831 assert_eq!(consumer.approximate_len(), 1);
832 assert_eq!(consumer.try_peek(), Ok(incoming));
833
834 timeline.visit_block(1, |_| {}).unwrap();
835 timeline.enqueue(incoming).unwrap();
836 assert_eq!(consumer.commit_peeked(), Ok(SpscCommitOutcome::Value));
837 assert_eq!(consumer.approximate_len(), 0);
838 assert_eq!(timeline.next_event(), Some(incoming));
839 }
840
841 #[test]
842 fn producer_cannot_overwrite_a_peeked_slot() {
843 let (mut producer, mut consumer) = mailbox(1);
844 let first = TestEvent {
845 sequence: 1,
846 value: 11,
847 };
848 let second = TestEvent {
849 sequence: 2,
850 value: 22,
851 };
852 producer.try_push(first).unwrap();
853 assert_eq!(consumer.try_peek(), Ok(first));
854
855 assert_eq!(producer.try_push(second), Err(SpscPushError::Full(second)));
856 assert_eq!(consumer.try_peek(), Ok(first));
857 assert_eq!(consumer.commit_peeked(), Ok(SpscCommitOutcome::Value));
858 producer.try_push(second).unwrap();
859 assert_eq!(consumer.try_pop(), Ok(second));
860 }
861
862 #[test]
863 fn try_pop_consumes_an_existing_peek() {
864 let (mut producer, mut consumer) = mailbox(1);
865 let event = TestEvent {
866 sequence: 4,
867 value: 44,
868 };
869 producer.try_push(event).unwrap();
870 assert_eq!(consumer.try_peek(), Ok(event));
871 assert_eq!(consumer.try_pop(), Ok(event));
872 assert_eq!(consumer.approximate_len(), 0);
873 assert_eq!(consumer.counters().successful_pops, 1);
874 }
875
876 #[test]
877 fn wraps_non_power_of_two_capacity_without_reordering() {
878 let (mut producer, mut consumer) = mailbox(3);
879 for sequence in 0..20_000 {
880 let event = TestEvent {
881 sequence,
882 value: !sequence,
883 };
884 producer.try_push(event).unwrap();
885 assert_eq!(consumer.try_pop(), Ok(event));
886 }
887 assert_eq!(producer.approximate_len(), 0);
888 }
889
890 #[test]
891 fn peek_and_commit_wrap_non_power_of_two_capacity() {
892 let (mut producer, mut consumer) = mailbox(3);
893 for sequence in 0..20_000 {
894 let event = TestEvent {
895 sequence,
896 value: sequence.rotate_right(9),
897 };
898 producer.try_push(event).unwrap();
899 assert_eq!(consumer.try_peek(), Ok(event));
900 assert_eq!(consumer.commit_peeked(), Ok(SpscCommitOutcome::Value));
901 }
902 assert_eq!(producer.approximate_len(), 0);
903 assert_eq!(consumer.counters().successful_pops, 20_000);
904 }
905
906 #[test]
907 fn distinguishes_empty_from_disconnected() {
908 let (producer, mut consumer) = mailbox(1);
909 assert_eq!(consumer.try_pop(), Err(SpscPopError::Empty));
910 drop(producer);
911 assert_eq!(consumer.try_pop(), Err(SpscPopError::Disconnected));
912 assert!(!consumer.is_producer_connected());
913 assert_eq!(consumer.counters().disconnected_pops, 1);
914 }
915
916 #[test]
917 fn peek_drains_final_value_before_disconnect() {
918 let (mut producer, mut consumer) = mailbox(1);
919 let event = TestEvent {
920 sequence: 5,
921 value: 55,
922 };
923 assert_eq!(consumer.try_peek(), Err(SpscPopError::Empty));
924 producer.try_push(event).unwrap();
925 drop(producer);
926
927 assert_eq!(consumer.try_peek(), Ok(event));
928 assert_eq!(consumer.try_peek(), Ok(event));
929 assert_eq!(consumer.commit_peeked(), Ok(SpscCommitOutcome::Value));
930 assert_eq!(consumer.try_peek(), Err(SpscPopError::Disconnected));
931 }
932
933 #[test]
934 fn returns_value_when_consumer_is_disconnected() {
935 let (mut producer, consumer) = mailbox(1);
936 drop(consumer);
937 let event = TestEvent {
938 sequence: 7,
939 value: 9,
940 };
941 assert_eq!(
942 producer.try_push(event),
943 Err(SpscPushError::Disconnected(event))
944 );
945 assert!(!producer.is_consumer_connected());
946 assert_eq!(producer.counters().rejected_disconnected_pushes, 1);
947 }
948
949 #[test]
950 fn transfers_values_between_threads_in_order() {
951 const EVENT_COUNT: u64 = 200_000;
952 let (mut producer, mut consumer) = mailbox(64);
953
954 let producer_thread = thread::spawn(move || {
955 for sequence in 0..EVENT_COUNT {
956 let mut pending = TestEvent {
957 sequence,
958 value: sequence.rotate_left(17),
959 };
960 loop {
961 match producer.try_push(pending) {
962 Ok(()) => break,
963 Err(SpscPushError::Full(value)) => {
964 pending = value;
965 std::hint::spin_loop();
966 }
967 Err(SpscPushError::Disconnected(_)) => {
968 panic!("consumer disconnected")
969 }
970 }
971 }
972 }
973 producer.counters()
974 });
975
976 let consumer_thread = thread::spawn(move || {
977 for expected_sequence in 0..EVENT_COUNT {
978 loop {
979 match consumer.try_pop() {
980 Ok(event) => {
981 assert_eq!(event.sequence, expected_sequence);
982 assert_eq!(event.value, expected_sequence.rotate_left(17));
983 break;
984 }
985 Err(SpscPopError::Empty) => std::hint::spin_loop(),
986 Err(error) => panic!("unexpected pop error: {error}"),
987 }
988 }
989 }
990 consumer.counters()
991 });
992
993 let producer_counters = producer_thread.join().unwrap();
994 let consumer_counters = consumer_thread.join().unwrap();
995 assert_eq!(producer_counters.accepted_pushes, EVENT_COUNT as usize);
996 assert_eq!(consumer_counters.successful_pops, EVENT_COUNT as usize);
997 }
998
999 #[test]
1000 fn timed_control_codec_preserves_every_field() {
1001 let (mut producer, mut consumer) = timed_player_control_mailbox(3).unwrap();
1002 let events = [
1003 TimedPlayerControl::new(
1004 44_100,
1005 8,
1006 PlayerControl::new(
1007 DeckMechanicalControl {
1008 motor_mode: MotorMode::Off,
1009 motor_target_angular_velocity_rad_s: -0.0,
1010 hand_contact: false,
1011 hand_target_angle_rad: None,
1012 hand_target_angular_velocity_rad_s: -12.25,
1013 hand_normal_force_n: 0.0,
1014 hand_contact_radius_m: 0.149,
1015 stylus_torque_nm: -0.000_012,
1016 },
1017 false,
1018 )
1019 .with_scratch(ScratchPreset::Baby, 1, 0.25),
1020 ),
1021 TimedPlayerControl::new(
1022 44_101,
1023 9,
1024 PlayerControl::new(
1025 DeckMechanicalControl {
1026 motor_mode: MotorMode::Servo,
1027 motor_target_angular_velocity_rad_s: 3.49,
1028 hand_contact: true,
1029 hand_target_angle_rad: Some(-1.75),
1030 hand_target_angular_velocity_rad_s: 48.0,
1031 hand_normal_force_n: 4.5,
1032 hand_contact_radius_m: 0.08,
1033 stylus_torque_nm: 0.000_1,
1034 },
1035 true,
1036 )
1037 .with_scratch(ScratchPreset::Crab, 8, 0.75),
1038 ),
1039 TimedPlayerControl::new(
1040 44_102,
1041 10,
1042 PlayerControl::new(
1043 DeckMechanicalControl {
1044 motor_mode: MotorMode::Brake,
1045 motor_target_angular_velocity_rad_s: -3.49,
1046 hand_contact: true,
1047 hand_target_angle_rad: Some(2.25),
1048 hand_target_angular_velocity_rad_s: -60.0,
1049 hand_normal_force_n: 5.0,
1050 hand_contact_radius_m: 0.12,
1051 stylus_torque_nm: 0.0,
1052 },
1053 false,
1054 )
1055 .with_scratch(ScratchPreset::Drum, 1, 0.0),
1056 ),
1057 ];
1058
1059 for event in events {
1060 producer.try_push(event).unwrap();
1061 }
1062 for expected in events {
1063 let actual = consumer.try_pop().unwrap();
1064 assert_eq!(actual.absolute_frame, expected.absolute_frame);
1065 assert_eq!(actual.sequence, expected.sequence);
1066 assert_eq!(
1067 actual.control.stylus_lowered,
1068 expected.control.stylus_lowered
1069 );
1070 assert_eq!(
1071 actual.control.deck.motor_mode,
1072 expected.control.deck.motor_mode
1073 );
1074 assert_eq!(
1075 actual.control.deck.hand_contact,
1076 expected.control.deck.hand_contact
1077 );
1078 assert_eq!(
1079 actual
1080 .control
1081 .deck
1082 .motor_target_angular_velocity_rad_s
1083 .to_bits(),
1084 expected
1085 .control
1086 .deck
1087 .motor_target_angular_velocity_rad_s
1088 .to_bits()
1089 );
1090 assert_eq!(
1091 actual.control.deck.hand_target_angle_rad.map(f64::to_bits),
1092 expected
1093 .control
1094 .deck
1095 .hand_target_angle_rad
1096 .map(f64::to_bits)
1097 );
1098 assert_eq!(
1099 actual
1100 .control
1101 .deck
1102 .hand_target_angular_velocity_rad_s
1103 .to_bits(),
1104 expected
1105 .control
1106 .deck
1107 .hand_target_angular_velocity_rad_s
1108 .to_bits()
1109 );
1110 assert_eq!(
1111 actual.control.deck.hand_normal_force_n.to_bits(),
1112 expected.control.deck.hand_normal_force_n.to_bits()
1113 );
1114 assert_eq!(
1115 actual.control.deck.hand_contact_radius_m.to_bits(),
1116 expected.control.deck.hand_contact_radius_m.to_bits()
1117 );
1118 assert_eq!(
1119 actual.control.deck.stylus_torque_nm.to_bits(),
1120 expected.control.deck.stylus_torque_nm.to_bits()
1121 );
1122 assert_eq!(
1123 actual.control.scratch_preset,
1124 expected.control.scratch_preset
1125 );
1126 assert_eq!(
1127 actual.control.scratch_clicks,
1128 expected.control.scratch_clicks
1129 );
1130 assert_eq!(
1131 actual.control.manual_crossfader_gain.to_bits(),
1132 expected.control.manual_crossfader_gain.to_bits()
1133 );
1134 }
1135 }
1136
1137 #[test]
1138 fn invalid_payload_is_consumed_and_counted() {
1139 struct RejectCodec;
1140 impl SpscCodec<TestEvent, 1> for RejectCodec {
1141 fn encode(_value: TestEvent, destination: &mut [u8; 1]) {
1142 destination[0] = 255;
1143 }
1144 fn decode(_source: &[u8; 1]) -> Option<TestEvent> {
1145 None
1146 }
1147 }
1148
1149 let (mut producer, mut consumer) = spsc_mailbox::<TestEvent, RejectCodec, 1>(1).unwrap();
1150 producer
1151 .try_push(TestEvent {
1152 sequence: 1,
1153 value: 2,
1154 })
1155 .unwrap();
1156 assert_eq!(consumer.try_pop(), Err(SpscPopError::InvalidEncoding));
1157 assert_eq!(consumer.approximate_len(), 0);
1158 assert_eq!(consumer.counters().invalid_payloads, 1);
1159 }
1160
1161 #[test]
1162 fn invalid_peek_remains_reserved_until_commit() {
1163 struct RejectCodec;
1164 impl SpscCodec<TestEvent, 1> for RejectCodec {
1165 fn encode(_value: TestEvent, destination: &mut [u8; 1]) {
1166 destination[0] = 255;
1167 }
1168 fn decode(_source: &[u8; 1]) -> Option<TestEvent> {
1169 None
1170 }
1171 }
1172
1173 let (mut producer, mut consumer) = spsc_mailbox::<TestEvent, RejectCodec, 1>(1).unwrap();
1174 let first = TestEvent {
1175 sequence: 1,
1176 value: 2,
1177 };
1178 let blocked = TestEvent {
1179 sequence: 2,
1180 value: 3,
1181 };
1182 producer.try_push(first).unwrap();
1183
1184 assert_eq!(consumer.try_peek(), Err(SpscPopError::InvalidEncoding));
1185 assert_eq!(consumer.try_peek(), Err(SpscPopError::InvalidEncoding));
1186 assert_eq!(consumer.counters().invalid_payloads, 0);
1187 assert_eq!(
1188 producer.try_push(blocked),
1189 Err(SpscPushError::Full(blocked))
1190 );
1191 assert_eq!(
1192 consumer.commit_peeked(),
1193 Ok(SpscCommitOutcome::InvalidEncoding)
1194 );
1195 assert_eq!(consumer.approximate_len(), 0);
1196 assert_eq!(consumer.counters().invalid_payloads, 1);
1197 }
1198
1199 #[test]
1200 fn endpoints_are_send() {
1201 fn assert_send<T: Send>() {}
1202 assert_send::<SpscProducer<TestEvent, TestCodec, 16>>();
1203 assert_send::<SpscConsumer<TestEvent, TestCodec, 16>>();
1204 }
1205}