Skip to main content

record_player/
spsc.rs

1//! A bounded mailbox for one producer and one consumer.
2//!
3//! The mailbox uses only safe Rust atomic operations. Its storage does not
4//! allocate after construction. The built-in player-control codec also does not
5//! allocate. A full mailbox rejects the incoming value and keeps all accepted
6//! values in order.
7
8use 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
18/// Converts one value to and from a fixed-size atomic payload.
19///
20/// The producer calls `encode`. The consumer calls `decode`. Implementations
21/// must decode every payload that they encode. A real-time codec must not
22/// allocate in either function.
23pub 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/// A rejected push returns ownership of the incoming value.
50#[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/// Counters can wrap after `usize::MAX` operations.
116///
117/// A snapshot can observe counters from different instants. Use it for
118/// telemetry, not synchronization.
119#[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
228/// The only producer endpoint for a mailbox.
229///
230/// The endpoint is `Send` but not `Sync`. Its operations use a bounded number
231/// of lock-free atomic loads and stores.
232pub 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        // This release publishes every payload byte to the consumer.
267        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
305/// The only consumer endpoint for a mailbox.
306///
307/// The endpoint is `Send` but not `Sync`. Its operations use a bounded number
308/// of lock-free atomic loads and stores.
309pub 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    /// Returns the current value without permitting slot reuse.
326    ///
327    /// Repeated calls return the same value. Call `commit_peeked` only after the
328    /// downstream consumer accepts the value.
329    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    /// Consumes the value returned by `try_peek`.
351    ///
352    /// The release store permits the producer to reuse this slot. An invalid
353    /// payload can also be committed so later values remain accessible.
354    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        // The second load observes the producer's final publication.
411        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        // This release permits the producer to reuse the consumed slot.
434        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
456/// Creates one bounded mailbox and its two unique endpoints.
457pub 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
515/// Each counter has one endpoint writer.
516fn 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}