Skip to main content

sequenced_broadcast/
lib.rs

1use std::{
2    collections::VecDeque,
3    sync::{
4        Arc,
5        atomic::{AtomicBool, Ordering},
6    },
7};
8
9use arc_metrics::{IntCounter, IntGauge};
10use tokio::sync::{Notify, RwLock, broadcast};
11
12pub struct SequencedBroadcast<T> {
13    state: Arc<State<T>>,
14}
15
16/// Sends sequence-numbered messages into a [`SequencedBroadcast`].
17pub struct SequencedSender<T> {
18    next_seq: u64,
19    state: Arc<State<T>>,
20}
21
22/// Receives sequence-numbered messages from replay history and live broadcast.
23///
24/// A receiver may observe the same sequence in both its catch-up replay and the
25/// live broadcast channel. This type treats those repeated sequences as
26/// duplicates and skips them before returning the next expected message.
27pub struct SequencedReceiver<T> {
28    state: Arc<State<T>>,
29    next_seq: u64,
30    replay: VecDeque<SequencedItem<T>>,
31    live_rx: broadcast::Receiver<SequencedItem<T>>,
32    terminal: Option<SequencedRecvError>,
33    active: bool,
34}
35
36#[derive(Default, Debug)]
37pub struct SequencedBroadcastMetrics {
38    pub oldest_sequence: IntGauge,
39    pub next_sequence: IntGauge,
40    pub new_client_drop_count: IntCounter,
41    pub new_client_accept_count: IntCounter,
42    pub active_subs_gauge: IntGauge,
43    pub disconnect_count: IntCounter,
44    pub duplicate_skip_count: IntCounter,
45    pub lagged_receiver_count: IntCounter,
46}
47
48#[derive(Debug, Clone)]
49pub struct SequencedBroadcastSettings {
50    pub history_capacity: usize,
51    pub broadcast_capacity: usize,
52}
53
54#[derive(Debug, Clone, PartialEq, Eq)]
55pub enum SettingsError {
56    ZeroHistoryCapacity,
57    ZeroBroadcastCapacity,
58}
59
60#[derive(Debug, Clone, PartialEq, Eq)]
61pub enum SubscribeError {
62    SequenceTooFarAhead { seq: u64, max: u64 },
63    SequenceTooFarBehind { seq: u64, min: u64 },
64    Closed,
65}
66
67#[derive(Debug, PartialEq, Eq)]
68pub enum SequencedSenderError<T> {
69    InvalidSequence { expected: u64, got: u64, item: T },
70    Closed(T),
71}
72
73#[derive(Debug, Clone, PartialEq, Eq)]
74pub enum SequencedRecvError {
75    Closed,
76    Lagged {
77        expected: u64,
78        got: u64,
79        skipped: u64,
80    },
81}
82
83#[derive(Debug, Clone, PartialEq, Eq)]
84pub enum SequencedTryRecvError {
85    Empty,
86    Closed,
87    Lagged {
88        expected: u64,
89        got: u64,
90        skipped: u64,
91    },
92}
93
94struct State<T> {
95    live_tx: broadcast::Sender<SequencedItem<T>>,
96    history: RwLock<History<T>>,
97    closed: AtomicBool,
98    close_notify: Notify,
99    metrics: Arc<SequencedBroadcastMetrics>,
100}
101
102#[derive(Debug, Clone)]
103struct SequencedItem<T> {
104    seq: u64,
105    item: T,
106}
107
108struct History<T> {
109    oldest_seq: u64,
110    next_seq: u64,
111    entries: VecDeque<SequencedItem<T>>,
112    capacity: usize,
113}
114
115impl Default for SequencedBroadcastSettings {
116    fn default() -> Self {
117        SequencedBroadcastSettings {
118            history_capacity: 16 * 1024,
119            broadcast_capacity: 16 * 1024,
120        }
121    }
122}
123
124impl<T> SequencedBroadcast<T>
125where
126    T: Send + Clone + 'static,
127{
128    pub fn new(
129        next_seq: u64,
130        settings: SequencedBroadcastSettings,
131    ) -> Result<(Self, SequencedSender<T>), SettingsError> {
132        if settings.history_capacity == 0 {
133            return Err(SettingsError::ZeroHistoryCapacity);
134        }
135
136        if settings.broadcast_capacity == 0 {
137            return Err(SettingsError::ZeroBroadcastCapacity);
138        }
139
140        let (live_tx, _) = broadcast::channel(settings.broadcast_capacity);
141        let metrics = Arc::new(SequencedBroadcastMetrics {
142            oldest_sequence: {
143                let i = IntGauge::default();
144                i.set(next_seq);
145                i
146            },
147            next_sequence: {
148                let i = IntGauge::default();
149                i.set(next_seq);
150                i
151            },
152            ..Default::default()
153        });
154
155        let state = Arc::new(State {
156            live_tx,
157            history: RwLock::new(History {
158                oldest_seq: next_seq,
159                next_seq,
160                entries: VecDeque::with_capacity(settings.history_capacity),
161                capacity: settings.history_capacity,
162            }),
163            closed: AtomicBool::new(false),
164            close_notify: Notify::new(),
165            metrics,
166        });
167
168        Ok((
169            Self {
170                state: state.clone(),
171            },
172            SequencedSender { next_seq, state },
173        ))
174    }
175
176    pub async fn subscribe_from(
177        &self,
178        next_sequence: u64,
179    ) -> Result<SequencedReceiver<T>, SubscribeError> {
180        // Subscribe to live messages before copying history. That ordering can
181        // duplicate messages across replay and live delivery, but it prevents a
182        // missed sequence between the history snapshot and live subscription.
183        let live_rx = self.state.live_tx.subscribe();
184        let history = self.state.history.read().await;
185
186        if next_sequence < history.oldest_seq {
187            self.state.metrics.new_client_drop_count.inc();
188            return Err(SubscribeError::SequenceTooFarBehind {
189                seq: next_sequence,
190                min: history.oldest_seq,
191            });
192        }
193
194        if history.next_seq < next_sequence {
195            self.state.metrics.new_client_drop_count.inc();
196            return Err(SubscribeError::SequenceTooFarAhead {
197                seq: next_sequence,
198                max: history.next_seq,
199            });
200        }
201
202        let replay = history
203            .entries
204            .iter()
205            .filter(|entry| next_sequence <= entry.seq)
206            .cloned()
207            .collect();
208
209        drop(history);
210
211        self.state.metrics.new_client_accept_count.inc();
212        self.state.metrics.active_subs_gauge.inc();
213
214        Ok(SequencedReceiver {
215            state: self.state.clone(),
216            next_seq: next_sequence,
217            replay,
218            live_rx,
219            terminal: None,
220            active: true,
221        })
222    }
223
224    pub fn metrics_ref(&self) -> &SequencedBroadcastMetrics {
225        &self.state.metrics
226    }
227
228    pub fn metrics(&self) -> Arc<SequencedBroadcastMetrics> {
229        self.state.metrics.clone()
230    }
231
232    pub fn is_closed(&self) -> bool {
233        self.state.closed.load(Ordering::Acquire)
234    }
235
236    pub async fn closed(&self) {
237        loop {
238            let notified = self.state.close_notify.notified();
239            if self.is_closed() {
240                return;
241            }
242
243            notified.await;
244        }
245    }
246}
247
248impl<T> SequencedSender<T> {
249    pub fn seq(&self) -> u64 {
250        self.next_seq
251    }
252
253    pub fn is_closed(&self) -> bool {
254        self.state.closed.load(Ordering::Acquire)
255    }
256
257    pub async fn closed(&self) {
258        loop {
259            let notified = self.state.close_notify.notified();
260            if self.is_closed() {
261                return;
262            }
263
264            notified.await;
265        }
266    }
267
268    pub fn close(&mut self) {
269        if !self.state.closed.swap(true, Ordering::AcqRel) {
270            self.state.close_notify.notify_waiters();
271        }
272    }
273}
274
275impl<T> SequencedSender<T>
276where
277    T: Send + Clone + 'static,
278{
279    pub async fn send(&mut self, item: T) -> Result<u64, SequencedSenderError<T>> {
280        self.send_at(self.next_seq, item).await
281    }
282
283    pub async fn send_at(&mut self, seq: u64, item: T) -> Result<u64, SequencedSenderError<T>> {
284        if self.is_closed() {
285            return Err(SequencedSenderError::Closed(item));
286        }
287
288        if seq != self.next_seq {
289            return Err(SequencedSenderError::InvalidSequence {
290                expected: self.next_seq,
291                got: seq,
292                item,
293            });
294        }
295
296        let mut history = self.state.history.write().await;
297        if self.is_closed() {
298            return Err(SequencedSenderError::Closed(item));
299        }
300
301        /* note, do capacity check before push to ensure we don't allocate more memory */
302        if history.capacity <= history.entries.len() {
303            history.entries.pop_front();
304            history.oldest_seq += 1;
305        }
306
307        let message = SequencedItem { seq, item };
308        history.entries.push_back(message.clone());
309        history.next_seq += 1;
310
311        self.state.metrics.oldest_sequence.set(history.oldest_seq);
312        self.state.metrics.next_sequence.set(history.next_seq);
313
314        drop(history);
315
316        let _ = self.state.live_tx.send(message);
317        self.next_seq += 1;
318
319        Ok(seq)
320    }
321}
322
323impl<T> Drop for SequencedSender<T> {
324    fn drop(&mut self) {
325        self.close();
326    }
327}
328
329impl<T> SequencedReceiver<T>
330where
331    T: Send + Clone + 'static,
332{
333    pub async fn recv(&mut self) -> Result<(u64, T), SequencedRecvError> {
334        if let Some(error) = &self.terminal {
335            return Err(error.clone());
336        }
337
338        loop {
339            if let Some(item) = self.pop_replay()? {
340                return Ok(item);
341            }
342
343            if self.state.closed.load(Ordering::Acquire) {
344                match self.live_rx.try_recv() {
345                    Ok(message) => match self.handle_message(message) {
346                        Ok(Some(item)) => return Ok(item),
347                        Ok(None) => continue,
348                        Err(error) => return Err(error),
349                    },
350                    Err(broadcast::error::TryRecvError::Empty)
351                    | Err(broadcast::error::TryRecvError::Closed) => {
352                        return Err(self.terminate(SequencedRecvError::Closed));
353                    }
354                    Err(broadcast::error::TryRecvError::Lagged(skipped)) => {
355                        return Err(self.terminate_lagged(skipped));
356                    }
357                }
358            }
359
360            tokio::select! {
361                result = self.live_rx.recv() => {
362                    match result {
363                        Ok(message) => match self.handle_message(message) {
364                            Ok(Some(item)) => return Ok(item),
365                            Ok(None) => continue,
366                            Err(error) => return Err(error),
367                        },
368                        Err(broadcast::error::RecvError::Closed) => {
369                            return Err(self.terminate(SequencedRecvError::Closed));
370                        }
371                        Err(broadcast::error::RecvError::Lagged(skipped)) => {
372                            return Err(self.terminate_lagged(skipped));
373                        }
374                    }
375                }
376                _ = self.state.close_notify.notified() => {
377                    continue;
378                }
379            }
380        }
381    }
382
383    pub fn try_recv(&mut self) -> Result<(u64, T), SequencedTryRecvError> {
384        if let Some(error) = &self.terminal {
385            return Err(error.clone().into());
386        }
387
388        loop {
389            if let Some(item) = self.pop_replay()? {
390                return Ok(item);
391            }
392
393            match self.live_rx.try_recv() {
394                Ok(message) => match self
395                    .handle_message(message)
396                    .map_err(SequencedTryRecvError::from)?
397                {
398                    Some(item) => return Ok(item),
399                    None => continue,
400                },
401                Err(broadcast::error::TryRecvError::Empty) => {
402                    if self.state.closed.load(Ordering::Acquire) {
403                        return Err(SequencedTryRecvError::from(
404                            self.terminate(SequencedRecvError::Closed),
405                        ));
406                    }
407
408                    return Err(SequencedTryRecvError::Empty);
409                }
410                Err(broadcast::error::TryRecvError::Closed) => {
411                    return Err(SequencedTryRecvError::from(
412                        self.terminate(SequencedRecvError::Closed),
413                    ));
414                }
415                Err(broadcast::error::TryRecvError::Lagged(skipped)) => {
416                    return Err(SequencedTryRecvError::from(self.terminate_lagged(skipped)));
417                }
418            }
419        }
420    }
421
422    pub fn next_seq(&self) -> u64 {
423        self.next_seq
424    }
425
426    pub fn is_closed(&self) -> bool {
427        self.terminal.is_some() || self.state.closed.load(Ordering::Acquire)
428    }
429
430    fn pop_replay(&mut self) -> Result<Option<(u64, T)>, SequencedRecvError> {
431        while let Some(message) = self.replay.pop_front() {
432            match self.handle_message(message)? {
433                Some(item) => return Ok(Some(item)),
434                None => continue,
435            }
436        }
437
438        Ok(None)
439    }
440
441    fn handle_message(
442        &mut self,
443        message: SequencedItem<T>,
444    ) -> Result<Option<(u64, T)>, SequencedRecvError> {
445        if message.seq < self.next_seq {
446            self.state.metrics.duplicate_skip_count.inc();
447            return Ok(None);
448        }
449
450        if self.next_seq < message.seq {
451            let expected = self.next_seq;
452            let skipped = message.seq - expected;
453            return Err(self.terminate(SequencedRecvError::Lagged {
454                expected,
455                got: message.seq,
456                skipped,
457            }));
458        }
459
460        self.next_seq = message.seq + 1;
461        Ok(Some((message.seq, message.item)))
462    }
463
464    fn terminate_lagged(&mut self, skipped: u64) -> SequencedRecvError {
465        let expected = self.next_seq;
466        self.terminate(SequencedRecvError::Lagged {
467            expected,
468            got: expected.saturating_add(skipped),
469            skipped,
470        })
471    }
472
473    fn terminate(&mut self, error: SequencedRecvError) -> SequencedRecvError {
474        if self.terminal.is_none() {
475            if self.active {
476                self.active = false;
477                self.state.metrics.active_subs_gauge.dec();
478                self.state.metrics.disconnect_count.inc();
479            }
480
481            if matches!(error, SequencedRecvError::Lagged { .. }) {
482                self.state.metrics.lagged_receiver_count.inc();
483            }
484
485            self.terminal = Some(error.clone());
486        }
487
488        error
489    }
490}
491
492impl<T> Drop for SequencedReceiver<T> {
493    fn drop(&mut self) {
494        if self.active {
495            self.active = false;
496            self.state.metrics.active_subs_gauge.dec();
497        }
498    }
499}
500
501impl From<SequencedRecvError> for SequencedTryRecvError {
502    fn from(value: SequencedRecvError) -> Self {
503        match value {
504            SequencedRecvError::Closed => SequencedTryRecvError::Closed,
505            SequencedRecvError::Lagged {
506                expected,
507                got,
508                skipped,
509            } => SequencedTryRecvError::Lagged {
510                expected,
511                got,
512                skipped,
513            },
514        }
515    }
516}
517
518#[cfg(test)]
519mod test {
520    use super::*;
521    use tokio::{
522        task::JoinHandle,
523        time::{Duration, Instant, sleep, timeout},
524    };
525
526    fn settings(history_capacity: usize, broadcast_capacity: usize) -> SequencedBroadcastSettings {
527        SequencedBroadcastSettings {
528            history_capacity,
529            broadcast_capacity,
530        }
531    }
532
533    #[tokio::test]
534    async fn basic_live_delivery() {
535        let (subs, mut tx) = SequencedBroadcast::new(10, SequencedBroadcastSettings::default())
536            .expect("valid settings");
537        let mut rx = subs.subscribe_from(10).await.unwrap();
538
539        assert_eq!(tx.send("a").await.unwrap(), 10);
540        assert_eq!(tx.send("b").await.unwrap(), 11);
541        assert_eq!(tx.send("c").await.unwrap(), 12);
542
543        assert_eq!(rx.recv().await.unwrap(), (10, "a"));
544        assert_eq!(rx.recv().await.unwrap(), (11, "b"));
545        assert_eq!(rx.recv().await.unwrap(), (12, "c"));
546    }
547
548    #[tokio::test]
549    async fn history_catchup_delivery() {
550        let (subs, mut tx) = SequencedBroadcast::new(0, SequencedBroadcastSettings::default())
551            .expect("valid settings");
552
553        tx.send("a").await.unwrap();
554        tx.send("b").await.unwrap();
555
556        let mut rx = subs.subscribe_from(0).await.unwrap();
557        assert_eq!(rx.recv().await.unwrap(), (0, "a"));
558        assert_eq!(rx.recv().await.unwrap(), (1, "b"));
559
560        tx.send("c").await.unwrap();
561        assert_eq!(rx.recv().await.unwrap(), (2, "c"));
562    }
563
564    #[tokio::test]
565    async fn subscribe_from_middle_of_history() {
566        let (subs, mut tx) = SequencedBroadcast::new(0, SequencedBroadcastSettings::default())
567            .expect("valid settings");
568
569        tx.send("a").await.unwrap();
570        tx.send("b").await.unwrap();
571        tx.send("c").await.unwrap();
572
573        let mut rx = subs.subscribe_from(1).await.unwrap();
574        assert_eq!(rx.recv().await.unwrap(), (1, "b"));
575        assert_eq!(rx.recv().await.unwrap(), (2, "c"));
576        assert_eq!(rx.try_recv(), Err(SequencedTryRecvError::Empty));
577    }
578
579    #[tokio::test]
580    async fn reject_too_far_behind() {
581        let (subs, mut tx) = SequencedBroadcast::new(0, settings(2, 16)).expect("valid settings");
582
583        tx.send("a").await.unwrap();
584        tx.send("b").await.unwrap();
585        tx.send("c").await.unwrap();
586
587        let error = match subs.subscribe_from(0).await {
588            Ok(_) => panic!("expected subscribe error"),
589            Err(error) => error,
590        };
591        assert_eq!(
592            error,
593            SubscribeError::SequenceTooFarBehind { seq: 0, min: 1 }
594        );
595        assert_eq!(subs.metrics_ref().new_client_drop_count.load(), 1);
596    }
597
598    #[tokio::test]
599    async fn reject_too_far_ahead() {
600        let (subs, mut tx) = SequencedBroadcast::new(0, SequencedBroadcastSettings::default())
601            .expect("valid settings");
602
603        tx.send("a").await.unwrap();
604        tx.send("b").await.unwrap();
605        tx.send("c").await.unwrap();
606
607        let error = match subs.subscribe_from(4).await {
608            Ok(_) => panic!("expected subscribe error"),
609            Err(error) => error,
610        };
611        assert_eq!(
612            error,
613            SubscribeError::SequenceTooFarAhead { seq: 4, max: 3 }
614        );
615        assert_eq!(subs.metrics_ref().new_client_drop_count.load(), 1);
616    }
617
618    #[tokio::test]
619    async fn duplicate_live_messages_are_ignored() {
620        let (subs, mut tx) = SequencedBroadcast::new(0, SequencedBroadcastSettings::default())
621            .expect("valid settings");
622        let mut rx = subs.subscribe_from(0).await.unwrap();
623
624        rx.replay.push_back(SequencedItem { seq: 0, item: "a" });
625
626        tx.send("a").await.unwrap();
627        tx.send("b").await.unwrap();
628
629        assert_eq!(rx.recv().await.unwrap(), (0, "a"));
630        assert_eq!(rx.recv().await.unwrap(), (1, "b"));
631        assert_eq!(subs.metrics_ref().duplicate_skip_count.load(), 1);
632    }
633
634    #[tokio::test]
635    async fn broadcast_lag_returns_error() {
636        let (subs, mut tx) = SequencedBroadcast::new(0, settings(16, 2)).expect("valid settings");
637        let mut rx = subs.subscribe_from(0).await.unwrap();
638
639        for i in 0..8 {
640            tx.send(i).await.unwrap();
641        }
642
643        assert!(matches!(
644            rx.recv().await,
645            Err(SequencedRecvError::Lagged { .. })
646        ));
647        assert_eq!(subs.metrics_ref().active_subs_gauge.load(), 0);
648        assert_eq!(subs.metrics_ref().disconnect_count.load(), 1);
649        assert_eq!(subs.metrics_ref().lagged_receiver_count.load(), 1);
650    }
651
652    #[tokio::test]
653    async fn gap_returns_lagged_error() {
654        let (subs, _tx) = SequencedBroadcast::new(5, SequencedBroadcastSettings::default())
655            .expect("valid settings");
656        let mut rx = subs.subscribe_from(5).await.unwrap();
657
658        let _ = subs.state.live_tx.send(SequencedItem {
659            seq: 7,
660            item: "gap",
661        });
662
663        assert_eq!(
664            rx.recv().await.unwrap_err(),
665            SequencedRecvError::Lagged {
666                expected: 5,
667                got: 7,
668                skipped: 2,
669            }
670        );
671    }
672
673    #[tokio::test]
674    async fn send_succeeds_without_receivers() {
675        let (subs, mut tx) = SequencedBroadcast::new(0, SequencedBroadcastSettings::default())
676            .expect("valid settings");
677
678        assert_eq!(tx.send("a").await.unwrap(), 0);
679        assert_eq!(tx.send("b").await.unwrap(), 1);
680
681        let mut rx = subs.subscribe_from(0).await.unwrap();
682        assert_eq!(rx.recv().await.unwrap(), (0, "a"));
683        assert_eq!(rx.recv().await.unwrap(), (1, "b"));
684    }
685
686    #[tokio::test]
687    async fn sender_close_closes_receivers_after_replay() {
688        let (subs, mut tx) = SequencedBroadcast::new(0, SequencedBroadcastSettings::default())
689            .expect("valid settings");
690
691        tx.send("a").await.unwrap();
692        let mut rx = subs.subscribe_from(0).await.unwrap();
693        tx.close();
694
695        assert_eq!(rx.recv().await.unwrap(), (0, "a"));
696        assert_eq!(rx.recv().await.unwrap_err(), SequencedRecvError::Closed);
697    }
698
699    #[tokio::test]
700    async fn subscribe_after_close_can_replay_history() {
701        let (subs, mut tx) = SequencedBroadcast::new(0, SequencedBroadcastSettings::default())
702            .expect("valid settings");
703
704        tx.send("a").await.unwrap();
705        tx.send("b").await.unwrap();
706        tx.close();
707
708        let mut rx = subs.subscribe_from(0).await.unwrap();
709        assert_eq!(rx.recv().await.unwrap(), (0, "a"));
710        assert_eq!(rx.recv().await.unwrap(), (1, "b"));
711        assert_eq!(rx.recv().await.unwrap_err(), SequencedRecvError::Closed);
712    }
713
714    #[tokio::test]
715    async fn send_after_close_returns_item() {
716        let (_subs, mut tx) = SequencedBroadcast::new(0, SequencedBroadcastSettings::default())
717            .expect("valid settings");
718
719        tx.close();
720
721        assert_eq!(
722            tx.send("a").await.unwrap_err(),
723            SequencedSenderError::Closed("a")
724        );
725    }
726
727    #[tokio::test]
728    async fn send_at_validates_sequence() {
729        let (_subs, mut tx) = SequencedBroadcast::new(10, SequencedBroadcastSettings::default())
730            .expect("valid settings");
731
732        assert_eq!(
733            tx.send_at(11, "a").await.unwrap_err(),
734            SequencedSenderError::InvalidSequence {
735                expected: 10,
736                got: 11,
737                item: "a",
738            }
739        );
740        assert_eq!(tx.seq(), 10);
741    }
742
743    #[tokio::test]
744    async fn try_recv_empty() {
745        let (subs, _tx) =
746            SequencedBroadcast::<&'static str>::new(0, SequencedBroadcastSettings::default())
747                .expect("valid settings");
748        let mut rx = subs.subscribe_from(0).await.unwrap();
749
750        assert_eq!(rx.try_recv(), Err(SequencedTryRecvError::Empty));
751    }
752
753    #[tokio::test]
754    async fn active_metrics_decrement_on_drop() {
755        let (subs, _tx) =
756            SequencedBroadcast::<&'static str>::new(0, SequencedBroadcastSettings::default())
757                .expect("valid settings");
758        let rx_1 = subs.subscribe_from(0).await.unwrap();
759        let _rx_2 = subs.subscribe_from(0).await.unwrap();
760
761        assert_eq!(subs.metrics_ref().active_subs_gauge.load(), 2);
762        drop(rx_1);
763        assert_eq!(subs.metrics_ref().active_subs_gauge.load(), 1);
764    }
765
766    #[tokio::test]
767    async fn closed_waits_for_sender_close() {
768        let (subs, mut tx) =
769            SequencedBroadcast::<&'static str>::new(0, SequencedBroadcastSettings::default())
770                .expect("valid settings");
771
772        assert!(
773            timeout(Duration::from_millis(10), subs.closed())
774                .await
775                .is_err()
776        );
777        tx.close();
778        timeout(Duration::from_millis(10), subs.closed())
779            .await
780            .expect("closed should resolve");
781    }
782
783    #[tokio::test]
784    async fn settings_validation() {
785        let err = match SequencedBroadcast::<()>::new(
786            0,
787            SequencedBroadcastSettings {
788                history_capacity: 0,
789                broadcast_capacity: 1,
790            },
791        ) {
792            Ok(_) => panic!("expected settings error"),
793            Err(error) => error,
794        };
795        assert_eq!(err, SettingsError::ZeroHistoryCapacity);
796
797        let err = match SequencedBroadcast::<()>::new(
798            0,
799            SequencedBroadcastSettings {
800                history_capacity: 1,
801                broadcast_capacity: 0,
802            },
803        ) {
804            Ok(_) => panic!("expected settings error"),
805            Err(error) => error,
806        };
807        assert_eq!(err, SettingsError::ZeroBroadcastCapacity);
808    }
809
810    #[derive(Clone)]
811    struct FuzzRng {
812        state: u64,
813    }
814
815    impl FuzzRng {
816        fn new(seed: u64) -> Self {
817            Self { state: seed }
818        }
819
820        fn next(&mut self) -> u64 {
821            self.state = self
822                .state
823                .wrapping_mul(6_364_136_223_846_793_005)
824                .wrapping_add(1_442_695_040_888_963_407);
825            self.state
826        }
827
828        fn range(&mut self, upper: usize) -> usize {
829            if upper == 0 {
830                return 0;
831            }
832
833            (self.next() as usize) % upper
834        }
835
836        fn one_in(&mut self, denominator: usize) -> bool {
837            self.range(denominator) == 0
838        }
839    }
840
841    async fn run_fuzz_receiver(mut rx: SequencedReceiver<u64>, mut rng: FuzzRng) -> u64 {
842        let mut received = 0;
843
844        loop {
845            if rng.one_in(3) {
846                let expected = rx.next_seq();
847                match rx.try_recv() {
848                    Ok((seq, item)) => {
849                        assert_eq!(seq, expected);
850                        assert_eq!(item, seq);
851                        received += 1;
852                    }
853                    Err(SequencedTryRecvError::Empty) => {
854                        sleep(Duration::from_micros((rng.range(500) + 1) as u64)).await;
855                    }
856                    Err(SequencedTryRecvError::Closed)
857                    | Err(SequencedTryRecvError::Lagged { .. }) => break,
858                }
859            } else {
860                let expected = rx.next_seq();
861                match timeout(Duration::from_millis(50), rx.recv()).await {
862                    Ok(Ok((seq, item))) => {
863                        assert_eq!(seq, expected);
864                        assert_eq!(item, seq);
865                        received += 1;
866                    }
867                    Ok(Err(SequencedRecvError::Closed))
868                    | Ok(Err(SequencedRecvError::Lagged { .. })) => break,
869                    Err(_) => {}
870                }
871            }
872
873            if received != 0 && rng.one_in(128) {
874                break;
875            }
876
877            if rng.one_in(8) {
878                sleep(Duration::from_micros((rng.range(1_000) + 1) as u64)).await;
879            }
880        }
881
882        received
883    }
884
885    async fn join_finished(tasks: &mut Vec<JoinHandle<u64>>) -> u64 {
886        let mut received = 0;
887
888        while let Some(pos) = tasks.iter().position(JoinHandle::is_finished) {
889            received += tasks
890                .swap_remove(pos)
891                .await
892                .expect("fuzz receiver task panicked");
893        }
894
895        received
896    }
897
898    #[tokio::test(flavor = "multi_thread", worker_threads = 4)]
899    #[ignore = "30 second fuzzy stability test; run with `cargo test fuzzy_stability_30_seconds -- --ignored`"]
900    async fn fuzzy_stability_30_seconds() {
901        let deadline = Instant::now() + Duration::from_secs(30);
902        let (subs, mut tx) = SequencedBroadcast::new(0, settings(512, 64)).expect("valid settings");
903        let sender_deadline = deadline;
904
905        let sender = tokio::spawn(async move {
906            let mut rng = FuzzRng::new(0xa511_ce5d_f00d);
907            let mut sent = 0;
908
909            while Instant::now() < sender_deadline {
910                let seq = tx.seq();
911                let sent_seq = tx
912                    .send(seq)
913                    .await
914                    .expect("send should succeed before close");
915                assert_eq!(sent_seq, seq);
916                sent += 1;
917
918                if rng.one_in(4) {
919                    tokio::task::yield_now().await;
920                }
921
922                if rng.one_in(16) {
923                    sleep(Duration::from_micros((rng.range(250) + 1) as u64)).await;
924                }
925            }
926
927            tx.close();
928            sent
929        });
930
931        let mut rng = FuzzRng::new(0x5eed_5eed_cafe);
932        let mut tasks: Vec<JoinHandle<u64>> = Vec::new();
933        let mut total_received = 0;
934
935        while Instant::now() < deadline {
936            total_received += join_finished(&mut tasks).await;
937
938            if tasks.len() < 128 {
939                let history = subs.state.history.read().await;
940                let oldest = history.oldest_seq;
941                let next = history.next_seq;
942                drop(history);
943
944                let valid_span = next.saturating_sub(oldest) + 1;
945                let seq = oldest + rng.range(valid_span as usize) as u64;
946
947                match subs.subscribe_from(seq).await {
948                    Ok(rx) => {
949                        let seed = rng.next();
950                        tasks.push(tokio::spawn(run_fuzz_receiver(rx, FuzzRng::new(seed))));
951                    }
952                    Err(SubscribeError::SequenceTooFarBehind { .. })
953                    | Err(SubscribeError::SequenceTooFarAhead { .. }) => {}
954                    Err(SubscribeError::Closed) => break,
955                }
956            }
957
958            if rng.one_in(4) {
959                let history = subs.state.history.read().await;
960                let too_old = history.oldest_seq.saturating_sub(1);
961                let too_new = history.next_seq.saturating_add(1_000);
962                let oldest = history.oldest_seq;
963                drop(history);
964
965                if oldest != 0 {
966                    assert!(matches!(
967                        subs.subscribe_from(too_old).await,
968                        Err(SubscribeError::SequenceTooFarBehind { .. })
969                    ));
970                }
971
972                assert!(matches!(
973                    subs.subscribe_from(too_new).await,
974                    Err(SubscribeError::SequenceTooFarAhead { .. })
975                ));
976            }
977
978            sleep(Duration::from_millis((rng.range(5) + 1) as u64)).await;
979        }
980
981        let sent = sender.await.expect("fuzz sender task panicked");
982
983        for task in tasks {
984            total_received += task.await.expect("fuzz receiver task panicked");
985        }
986
987        assert!(sent > 1_000);
988        assert!(total_received > 0);
989        assert_eq!(subs.metrics_ref().next_sequence.load(), sent);
990        assert_eq!(subs.metrics_ref().active_subs_gauge.load(), 0);
991    }
992}