field_ios

In-the-field app for Radroots on iOS
git clone https://radroots.dev/git/field_ios.git
Log | Files | Refs | README | LICENSE

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 }