Skip to main content

acex_sim/
bus.rs

1// region: Imports
2
3use crate::clock::{Clock, Duration, Instant, SimClock};
4use crate::fault::FaultConfig;
5use crate::io::NodeAddress;
6use crate::rng::{Rng, Xorshift64};
7use acex_core::Vec;
8
9// endregion: Imports
10
11// region: Envelope
12
13/// A message in-flight on the simulation bus.
14#[derive(Debug, Clone)]
15#[cfg_attr(all(feature = "defmt", not(feature = "alloc")), derive(defmt::Format))]
16pub struct Envelope<const MAX_DATA: usize> {
17    pub src: NodeAddress,
18    pub dst: NodeAddress,
19    pub data: Vec<u8, MAX_DATA>,
20    /// Earliest time this message may be delivered.
21    pub deliver_at: Instant,
22}
23
24// endregion: Envelope
25
26// region: SimBus
27
28/// The simulation message bus.
29///
30/// Connects all nodes in the simulation. Drives time forward, delivers messages, and injects
31/// faults according to [`FaultConfig`].
32///
33/// `N` - max message payload bytes
34/// `Q` - max messages in-flight simultaneously
35#[derive(Debug)]
36#[cfg_attr(all(feature = "defmt", not(feature = "alloc")), derive(defmt::Format))]
37pub struct SimBus<const MAX_DATA: usize, const MAX_QUEUED: usize> {
38    clock: SimClock,
39    rng: Xorshift64,
40    faults: FaultConfig,
41    queue: Vec<Envelope<MAX_DATA>, MAX_QUEUED>,
42}
43
44impl<const MAX_DATA: usize, const MAX_QUEUED: usize> SimBus<MAX_DATA, MAX_QUEUED> {
45    pub fn new(seed: u64, faults: FaultConfig) -> Self {
46        Self {
47            clock: SimClock::new(),
48            rng: Xorshift64::new(seed),
49            faults,
50            queue: Vec::new(),
51        }
52    }
53
54    /// Returns the current simulation time.
55    pub fn now(&self) -> Instant {
56        self.clock.now()
57    }
58
59    /// Advances simulation time by `duration` and returns all messages that are due for delivery
60    /// at or before the new time.
61    pub fn tick(&mut self, duration: Duration) -> Vec<Envelope<MAX_DATA>, MAX_QUEUED> {
62        self.clock.advance(duration);
63        let now = self.clock.now();
64
65        let mut delivered = Vec::new();
66        let mut remaining = Vec::new();
67
68        for envelope in self.queue.drain(..) {
69            if envelope.deliver_at <= now {
70                let _ = delivered.push(envelope);
71            } else {
72                let _ = remaining.push(envelope);
73            }
74        }
75
76        if delivered.len() > 1 {
77            for i in 0..delivered.len() - 1 {
78                if self
79                    .rng
80                    .chance(self.faults.message_reorder.0, self.faults.message_reorder.1)
81                {
82                    delivered.swap(i, i + 1);
83                }
84            }
85        }
86
87        self.queue = remaining;
88        delivered
89    }
90
91    /// Enqueues a message from `src` to `dst` with fault injection applied.
92    ///
93    /// Returns `true` if the message was enqueued, `false` if it was dropped by fault injection or
94    /// the queue is full
95    pub fn send(&mut self, src: NodeAddress, dst: NodeAddress, data: &[u8]) -> bool {
96        let r = self
97            .rng
98            .chance(self.faults.message_loss.0, self.faults.message_loss.1);
99
100        if r {
101            return false;
102        }
103
104        if self
105            .rng
106            .chance(self.faults.timeout.0, self.faults.timeout.1)
107        {
108            return false;
109        }
110
111        let mut payload = Vec::new();
112
113        for &byte in data {
114            if self
115                .rng
116                .chance(self.faults.corruption.0, self.faults.corruption.1)
117            {
118                let _ = payload.push(byte ^ self.rng.next_u8());
119            } else {
120                let _ = payload.push(byte);
121            }
122        }
123
124        let deliver_at = if self
125            .rng
126            .chance(self.faults.message_delay.0, self.faults.message_delay.1)
127        {
128            let delay_us = self.rng.next_u64() % self.faults.max_delay.as_micros().max(1);
129            self.clock.now() + Duration::from_micros(delay_us)
130        } else {
131            self.clock.now()
132        };
133
134        let envelope = Envelope {
135            src,
136            dst,
137            data: payload,
138            deliver_at,
139        };
140
141        if self.queue.len() >= MAX_QUEUED {
142            return false;
143        }
144
145        let _ = self.queue.push(envelope);
146
147        true
148    }
149
150    /// Returns a reference to the current fault config
151    pub fn faults(&self) -> &FaultConfig {
152        &self.faults
153    }
154
155    /// Returns the fault config - allows escalating fault severity during a simulation run.
156    pub fn set_faults(&mut self, faults: FaultConfig) {
157        self.faults = faults;
158    }
159
160    /// Returns the RNG seed-derived next value - useful for injecting spontaneous NRCs at the node
161    /// level.
162    pub fn next_u8(&mut self) -> u8 {
163        self.rng.next_u8()
164    }
165
166    pub fn chance(&mut self, numerator: u32, denominator: u32) -> bool {
167        self.rng.chance(numerator, denominator)
168    }
169}
170
171// endregion: SimBus