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
16pub struct SequencedSender<T> {
18 next_seq: u64,
19 state: Arc<State<T>>,
20}
21
22pub 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 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 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}