subscription_contract.rs (7665B)
1 use std::sync::mpsc::{Receiver, Sender, channel}; 2 use std::time::Duration; 3 4 use tera_ffi::{FfiRuntimeChangeKind, FfiRuntimeChangeRecord, TeraRuntimeObserver}; 5 6 mod support; 7 8 struct Observer(Sender<FfiRuntimeChangeRecord>); 9 10 impl TeraRuntimeObserver for Observer { 11 fn on_change(&self, change: FfiRuntimeChangeRecord) { 12 let _ = self.0.send(change); 13 } 14 } 15 16 fn observer() -> ( 17 Box<dyn TeraRuntimeObserver>, 18 Receiver<FfiRuntimeChangeRecord>, 19 ) { 20 let (sender, receiver) = channel(); 21 (Box::new(Observer(sender)), receiver) 22 } 23 24 fn receive(receiver: &Receiver<FfiRuntimeChangeRecord>) -> FfiRuntimeChangeRecord { 25 receiver 26 .recv_timeout(Duration::from_secs(1)) 27 .expect("bounded observer delivery") 28 } 29 30 #[tokio::test] 31 async fn subscriptions_are_independent_bounded_handles_and_stop_individually() { 32 let (_root, runtime) = support::runtime().await; 33 let settings = runtime.phase1_settings().await.unwrap(); 34 let (first_observer, first_receiver) = observer(); 35 let (second_observer, second_receiver) = observer(); 36 let first = runtime 37 .subscribe_changes(first_observer) 38 .expect("first subscription"); 39 let second = runtime 40 .subscribe_changes(second_observer) 41 .expect("second subscription"); 42 43 let first_initial = receive(&first_receiver); 44 let second_initial = receive(&second_receiver); 45 assert_eq!(first_initial.kind, FfiRuntimeChangeKind::Initial); 46 assert_eq!(second_initial, first_initial); 47 assert_eq!(first_initial.scope.public_key, support::PUBLIC_KEY); 48 assert_eq!(first_initial.scope.source_generation, support::GENERATION); 49 assert_eq!(first_initial.scope.context, None); 50 assert_eq!(first_initial.epoch.len(), 32); 51 assert_eq!(runtime.phase1_settings().await.unwrap(), settings); 52 first.unsubscribe(); 53 assert!(!first.is_active()); 54 assert!(second.is_active()); 55 56 runtime 57 .configure_public_relays(vec!["wss://write.example".to_owned()]) 58 .await 59 .expect("relay configuration"); 60 assert_eq!(receive(&second_receiver).kind, FfiRuntimeChangeKind::Relay); 61 assert!( 62 first_receiver 63 .recv_timeout(Duration::from_millis(50)) 64 .is_err() 65 ); 66 67 runtime.shutdown().await.expect("shutdown"); 68 assert_eq!( 69 receive(&second_receiver).kind, 70 FfiRuntimeChangeKind::Lifecycle 71 ); 72 assert!(!second.is_active()); 73 } 74 75 #[tokio::test] 76 async fn today_hints_retain_the_exact_query_context_and_domain_revision() { 77 use tera_ffi::{FfiInvalidationRevision, FfiLocalNetworkRecord, FfiTodayProjectionUpdate}; 78 let (root, runtime) = support::runtime().await; 79 let (first_observer, receiver) = observer(); 80 let handle = runtime.subscribe_changes(first_observer).unwrap(); 81 let initial = receive(&receiver); 82 for (index, relay) in ["wss://first.example", "wss://second.example"] 83 .into_iter() 84 .enumerate() 85 { 86 let context = FfiLocalNetworkRecord { 87 schema_version: 1, 88 id: "default".into(), 89 label: "Local network".into(), 90 relay_urls: vec![relay.into()], 91 locality: None, 92 followed_authors: vec![], 93 generation: 1, 94 }; 95 runtime 96 .phase1_refresh_today( 97 context.clone(), 98 1_800_000_000, 99 FfiTodayProjectionUpdate::Incremental, 100 ) 101 .await 102 .unwrap(); 103 let hint = receive(&receiver); 104 assert_eq!(hint.kind, FfiRuntimeChangeKind::Today); 105 assert_eq!(hint.epoch, initial.epoch); 106 assert_eq!(hint.scope.public_key, support::PUBLIC_KEY); 107 assert_eq!(hint.scope.source_generation, support::GENERATION); 108 assert_eq!(hint.scope.context, Some(context)); 109 assert_eq!( 110 hint.revision, 111 FfiInvalidationRevision::Current { 112 value: index as u64 + 1 113 } 114 ); 115 } 116 handle.unsubscribe(); 117 let settings = runtime.phase1_settings().await.unwrap(); 118 runtime.shutdown().await.unwrap(); 119 let reopened = tera_ffi::TeraRuntime::new( 120 root.path().to_string_lossy().into_owned(), 121 support::PUBLIC_KEY.into(), 122 support::GENERATION.into(), 123 1_800_000_000_000, 124 tera_ffi::ProtectedDataAvailability::Available, 125 ) 126 .await 127 .unwrap(); 128 assert_eq!(reopened.phase1_settings().await.unwrap(), settings); 129 let (next_observer, next_receiver) = observer(); 130 let next = reopened.subscribe_changes(next_observer).unwrap(); 131 let next_initial = receive(&next_receiver); 132 assert_eq!(next_initial.scope, initial.scope); 133 assert_ne!(next_initial.epoch, initial.epoch); 134 next.unsubscribe(); 135 reopened.shutdown().await.unwrap(); 136 } 137 138 struct PausedObserver { 139 sender: Sender<FfiRuntimeChangeRecord>, 140 release: std::sync::Arc<(std::sync::Mutex<bool>, std::sync::Condvar)>, 141 } 142 143 impl TeraRuntimeObserver for PausedObserver { 144 fn on_change(&self, change: FfiRuntimeChangeRecord) { 145 let initial = change.kind == FfiRuntimeChangeKind::Initial 146 && change.delivery == tera_ffi::FfiRuntimeChangeDelivery::Change; 147 let _ = self.sender.send(change); 148 if initial { 149 let (released, wake) = &*self.release; 150 drop( 151 wake.wait_while(released.lock().unwrap(), |value| !*value) 152 .unwrap(), 153 ); 154 } 155 } 156 } 157 158 struct ReleaseOnDrop(std::sync::Arc<(std::sync::Mutex<bool>, std::sync::Condvar)>); 159 160 impl Drop for ReleaseOnDrop { 161 fn drop(&mut self) { 162 *self.0.0.lock().unwrap() = true; 163 self.0.1.notify_all(); 164 } 165 } 166 167 #[tokio::test] 168 async fn slow_observer_receives_final_gap_and_can_query_all_committed_settings_without_more_events() 169 { 170 use tera_ffi::{FfiReplaceSettingsRecord, FfiRuntimeChangeDelivery}; 171 let (_root, runtime) = support::runtime().await; 172 let (sender, receiver) = channel(); 173 let release = std::sync::Arc::new((std::sync::Mutex::new(false), std::sync::Condvar::new())); 174 let release_on_drop = ReleaseOnDrop(std::sync::Arc::clone(&release)); 175 let handle = runtime 176 .subscribe_changes(Box::new(PausedObserver { sender, release })) 177 .unwrap(); 178 let initial = receive(&receiver); 179 let mut settings = runtime.phase1_settings().await.unwrap(); 180 let before = settings.revision; 181 // The production FFI queue admits sixteen pending records; the next commit 182 // must leave a gap even when this caller performs no subsequent mutation. 183 for _ in 0..17 { 184 let mut media_network = settings.media_network; 185 media_network.allow_cellular_downloads = !media_network.allow_cellular_downloads; 186 settings = runtime 187 .phase1_replace_settings(FfiReplaceSettingsRecord { 188 schema_version: settings.schema_version, 189 expected_revision: settings.revision, 190 relays: settings.relays, 191 blossom: settings.blossom, 192 media_network, 193 local_storage: settings.local_storage, 194 }) 195 .await 196 .unwrap() 197 .settings; 198 } 199 assert_eq!(settings.revision, before + 17); 200 drop(release_on_drop); 201 let gap = receive(&receiver); 202 assert_eq!(gap.delivery, FfiRuntimeChangeDelivery::ResnapshotRequired); 203 assert_eq!(gap.scope, initial.scope); 204 assert_eq!(gap.epoch, initial.epoch); 205 assert_eq!(gap.kind, FfiRuntimeChangeKind::Initial); 206 assert_eq!(gap.entity_id, None); 207 assert_eq!(runtime.phase1_settings().await.unwrap(), settings); 208 handle.unsubscribe(); 209 runtime.shutdown().await.unwrap(); 210 }