subscription_queue.rs (8422B)
1 //! One bounded observer queue with an atomic out-of-band resnapshot signal. 2 //! 3 //! Queue contents, loss and closure share the same lock and wait predicate. 4 //! A producer cannot strand a final gap between the consumer's empty check and 5 //! sleep. No callback or asynchronous work runs while the lock is held. 6 7 use std::collections::VecDeque; 8 use std::sync::{Arc, Condvar, Mutex}; 9 10 use crate::{FfiRuntimeChangeKind, FfiRuntimeChangeRecord}; 11 12 pub(crate) const CHANGE_BUFFER_CAPACITY: usize = 16; 13 14 pub(crate) struct SubscriptionQueue { 15 state: Mutex<State>, 16 ready: Condvar, 17 gap: FfiRuntimeChangeRecord, 18 } 19 20 struct State { 21 changes: VecDeque<FfiRuntimeChangeRecord>, 22 dirty: bool, 23 closed: bool, 24 cancelled: bool, 25 #[cfg(test)] 26 waiting: bool, 27 } 28 29 impl SubscriptionQueue { 30 pub(crate) fn new(initial: FfiRuntimeChangeRecord) -> Arc<Self> { 31 let gap = initial.clone().requiring_resnapshot(); 32 let mut changes = VecDeque::with_capacity(CHANGE_BUFFER_CAPACITY); 33 changes.push_back(initial); 34 Arc::new(Self { 35 state: Mutex::new(State { 36 changes, 37 dirty: false, 38 closed: false, 39 cancelled: false, 40 #[cfg(test)] 41 waiting: false, 42 }), 43 ready: Condvar::new(), 44 gap, 45 }) 46 } 47 48 /// Returns false only after this queue stops admitting producer work. 49 pub(crate) fn send(&self, change: FfiRuntimeChangeRecord) -> bool { 50 let mut state = self 51 .state 52 .lock() 53 .unwrap_or_else(std::sync::PoisonError::into_inner); 54 if state.closed { 55 return false; 56 } 57 if state.changes.len() < CHANGE_BUFFER_CAPACITY { 58 state.changes.push_back(change); 59 } else { 60 state.dirty = true; 61 } 62 self.ready.notify_one(); 63 true 64 } 65 66 pub(crate) fn receive(&self) -> Option<FfiRuntimeChangeRecord> { 67 let mut state = self 68 .state 69 .lock() 70 .unwrap_or_else(std::sync::PoisonError::into_inner); 71 loop { 72 if state.cancelled { 73 return None; 74 } 75 // Preserve first epoch admission even if publication beats the worker. 76 if state 77 .changes 78 .front() 79 .is_some_and(|value| value.kind == FfiRuntimeChangeKind::Initial) 80 { 81 return state.changes.pop_front(); 82 } 83 if state.dirty { 84 state.dirty = false; 85 return Some(self.gap.clone()); 86 } 87 if let Some(change) = state.changes.pop_front() { 88 return Some(change); 89 } 90 if state.closed { 91 return None; 92 } 93 #[cfg(test)] 94 { 95 state.waiting = true; 96 } 97 state = self 98 .ready 99 .wait(state) 100 .unwrap_or_else(std::sync::PoisonError::into_inner); 101 #[cfg(test)] 102 { 103 state.waiting = false; 104 } 105 } 106 } 107 108 /// Drop stale ordinary callbacks but retain pending loss and final lifecycle. 109 pub(crate) fn close(&self, lifecycle: FfiRuntimeChangeRecord) { 110 let mut state = self 111 .state 112 .lock() 113 .unwrap_or_else(std::sync::PoisonError::into_inner); 114 if !state.closed { 115 state.changes.clear(); 116 state.changes.push_back(lifecycle); 117 state.closed = true; 118 self.ready.notify_one(); 119 } 120 } 121 122 /// An admitted callback may finish; no queued callback remains owned afterward. 123 pub(crate) fn cancel(&self) { 124 let mut state = self 125 .state 126 .lock() 127 .unwrap_or_else(std::sync::PoisonError::into_inner); 128 state.changes.clear(); 129 state.dirty = false; 130 state.closed = true; 131 state.cancelled = true; 132 self.ready.notify_one(); 133 } 134 } 135 136 #[cfg(test)] 137 mod tests { 138 use super::*; 139 use crate::{ 140 FfiInvalidationRevision, FfiRuntimeChangeDelivery, FfiRuntimeChangeScope, 141 RUNTIME_CHANGE_SCHEMA_VERSION, 142 }; 143 use std::time::Duration; 144 145 fn change(kind: FfiRuntimeChangeKind, revision: u64) -> FfiRuntimeChangeRecord { 146 FfiRuntimeChangeRecord { 147 schema_version: RUNTIME_CHANGE_SCHEMA_VERSION, 148 scope: FfiRuntimeChangeScope { 149 public_key: "a".repeat(64), 150 source_generation: "b".repeat(64), 151 context: None, 152 }, 153 epoch: "c".repeat(32), 154 revision: FfiInvalidationRevision::Current { value: revision }, 155 delivery: FfiRuntimeChangeDelivery::Change, 156 kind, 157 entity_id: Some(revision.to_string()), 158 } 159 } 160 161 #[test] 162 fn exact_capacity_and_one_more_retain_a_gap_without_a_later_event() { 163 let initial = change(FfiRuntimeChangeKind::Initial, 0); 164 let queue = SubscriptionQueue::new(initial.clone()); 165 assert_eq!(queue.receive(), Some(initial.clone())); 166 for revision in 1..=CHANGE_BUFFER_CAPACITY { 167 assert!(queue.send(change(FfiRuntimeChangeKind::Drafts, revision as u64))); 168 } 169 { 170 let state = queue.state.lock().unwrap(); 171 assert_eq!(state.changes.len(), CHANGE_BUFFER_CAPACITY); 172 assert!(!state.dirty); 173 } 174 assert!(queue.send(change(FfiRuntimeChangeKind::Media, 99))); 175 assert_eq!( 176 queue.state.lock().unwrap().changes.len(), 177 CHANGE_BUFFER_CAPACITY 178 ); 179 let gap = queue.receive().unwrap(); 180 assert_eq!(gap, initial.requiring_resnapshot()); 181 for revision in 1..=CHANGE_BUFFER_CAPACITY { 182 assert_eq!( 183 queue.receive(), 184 Some(change(FfiRuntimeChangeKind::Drafts, revision as u64)) 185 ); 186 } 187 queue.cancel(); 188 assert_eq!(queue.receive(), None); 189 assert!(!queue.send(change(FfiRuntimeChangeKind::Media, 100))); 190 } 191 192 #[test] 193 fn close_preserves_final_gap_and_lifecycle_but_cancel_discards_queued_work() { 194 for cancel in [false, true] { 195 let initial = change(FfiRuntimeChangeKind::Initial, 0); 196 let queue = SubscriptionQueue::new(initial.clone()); 197 for revision in 1..=CHANGE_BUFFER_CAPACITY { 198 queue.send(change(FfiRuntimeChangeKind::Today, revision as u64)); 199 } 200 // A raced publisher cannot replace the first epoch notification. 201 assert_eq!(queue.receive(), Some(initial.clone())); 202 let lifecycle = change(FfiRuntimeChangeKind::Lifecycle, 1); 203 queue.close(lifecycle.clone()); 204 if cancel { 205 queue.cancel(); 206 } else { 207 assert_eq!(queue.receive(), Some(initial.requiring_resnapshot())); 208 assert_eq!(queue.receive(), Some(lifecycle)); 209 } 210 assert_eq!(queue.receive(), None); 211 assert!(queue.state.lock().unwrap().changes.is_empty()); 212 } 213 } 214 215 #[test] 216 fn idle_receiver_wakes_on_close_or_cancellation() { 217 for cancel in [false, true] { 218 let queue = SubscriptionQueue::new(change(FfiRuntimeChangeKind::Initial, 0)); 219 queue.receive().unwrap(); 220 let receiver = Arc::clone(&queue); 221 let (sender, result) = std::sync::mpsc::channel(); 222 let worker = std::thread::spawn(move || sender.send(receiver.receive()).unwrap()); 223 let deadline = std::time::Instant::now() + Duration::from_secs(5); 224 while !queue.state.lock().unwrap().waiting && std::time::Instant::now() < deadline { 225 std::thread::yield_now(); 226 } 227 assert!(queue.state.lock().unwrap().waiting); 228 let lifecycle = change(FfiRuntimeChangeKind::Lifecycle, 1); 229 if cancel { 230 queue.cancel(); 231 } else { 232 queue.close(lifecycle.clone()); 233 } 234 assert_eq!( 235 result.recv_timeout(Duration::from_secs(5)).unwrap(), 236 if cancel { None } else { Some(lifecycle) } 237 ); 238 worker.join().unwrap(); 239 assert_eq!(queue.receive(), None); 240 } 241 } 242 }