dora_node_api/event_stream/
extensions.rs1use std::collections::VecDeque;
16use std::sync::{LazyLock, Mutex};
17
18const MAX_PENDING: usize = 4096;
23
24type Dropped = (String, String);
25
26static DROPPED: LazyLock<Mutex<VecDeque<Dropped>>> = LazyLock::new(|| Mutex::new(VecDeque::new()));
27
28pub(crate) fn push_dropped(namespace: String, key: String) {
30 let mut queue = DROPPED.lock().unwrap_or_else(|e| e.into_inner());
31 if queue.len() >= MAX_PENDING {
32 queue.pop_front();
33 }
34 queue.push_back((namespace, key));
35}
36
37pub fn drain_dropped_keys(namespace: &str) -> Vec<String> {
42 let mut queue = DROPPED.lock().unwrap_or_else(|e| e.into_inner());
43 let mut taken = Vec::new();
44 queue.retain(|(ns, key)| {
45 if ns == namespace {
46 taken.push(key.clone());
47 false
48 } else {
49 true
50 }
51 });
52 taken
53}
54
55#[cfg(test)]
56mod tests {
57 use super::*;
58
59 static TEST_LOCK: Mutex<()> = Mutex::new(());
63
64 fn guard() -> std::sync::MutexGuard<'static, ()> {
65 let g = TEST_LOCK.lock().unwrap_or_else(|e| e.into_inner());
66 DROPPED.lock().unwrap_or_else(|e| e.into_inner()).clear();
67 g
68 }
69
70 #[test]
71 fn drain_returns_only_the_requested_namespace() {
72 let _g = guard();
73 push_dropped("pool".into(), "a".into());
74 push_dropped("other".into(), "b".into());
75 push_dropped("pool".into(), "c".into());
76
77 assert_eq!(drain_dropped_keys("pool"), vec!["a", "c"]);
78 assert_eq!(drain_dropped_keys("other"), vec!["b"]);
80 }
81
82 #[test]
83 fn drain_is_exhaustive() {
84 let _g = guard();
85 push_dropped("pool".into(), "a".into());
86 assert_eq!(drain_dropped_keys("pool"), vec!["a"]);
87 assert!(drain_dropped_keys("pool").is_empty());
88 }
89
90 #[test]
91 fn queue_is_bounded_and_drops_oldest() {
92 let _g = guard();
93 for i in 0..MAX_PENDING + 10 {
94 push_dropped("pool".into(), i.to_string());
95 }
96 let drained = drain_dropped_keys("pool");
97 assert_eq!(drained.len(), MAX_PENDING);
98 assert_eq!(drained.first().map(String::as_str), Some("10"));
100 }
101}