observer.rs (55248B)
1 use std::num::NonZeroUsize; 2 use std::panic::{AssertUnwindSafe, catch_unwind}; 3 use std::sync::atomic::Ordering; 4 use std::sync::{Arc, Mutex, Weak}; 5 6 use harvestcircle_application::ChangeSubscriptionId; 7 8 use crate::commands::RuntimeCore; 9 use crate::{AppSnapshotDto, HarvestCircleAppCore, HarvestCircleError}; 10 11 const OBSERVER_CHANGE_CAPACITY: NonZeroUsize = NonZeroUsize::MIN.saturating_add(63); 12 pub(crate) const MAX_OBSERVERS: usize = 32; 13 14 pub(crate) struct ObserverTask { 15 handle: tokio::sync::Mutex<Option<tokio::task::JoinHandle<()>>>, 16 stop: Mutex<Option<tokio::sync::oneshot::Sender<()>>>, 17 _admission: Arc<tokio::sync::OwnedSemaphorePermit>, 18 } 19 20 struct ObserverResources { 21 // Field drop order keeps callback destruction inside its admission reservation. 22 observer: Box<dyn HarvestCircleChangeObserver>, 23 admission: Arc<tokio::sync::OwnedSemaphorePermit>, 24 } 25 26 #[derive(Clone, Debug, Eq, PartialEq)] 27 #[cfg_attr(not(coverage_nightly), derive(uniffi::Record))] 28 pub struct SnapshotChangeDto { 29 pub snapshot: AppSnapshotDto, 30 pub previous_revision: Option<u64>, 31 } 32 33 #[derive(Clone, Copy, Debug, Eq, PartialEq)] 34 #[cfg_attr(not(coverage_nightly), derive(uniffi::Record))] 35 pub struct ShutdownReceiptDto { 36 pub final_revision: u64, 37 pub closed: bool, 38 } 39 40 #[cfg_attr(not(coverage_nightly), uniffi::export(callback_interface))] 41 pub trait HarvestCircleChangeObserver: Send + Sync { 42 fn on_change(&self, change: SnapshotChangeDto); 43 } 44 45 #[cfg_attr(not(coverage_nightly), derive(uniffi::Object))] 46 pub struct ObserverSubscription { 47 core: Weak<RuntimeCore>, 48 id: Mutex<Option<ChangeSubscriptionId>>, 49 } 50 51 impl Drop for ObserverSubscription { 52 fn drop(&mut self) { 53 let Ok(retained_id) = self.id.get_mut() else { 54 return; 55 }; 56 let (Some(core), Some(id)) = (self.core.upgrade(), *retained_id) else { 57 return; 58 }; 59 let task = { 60 let Ok(observers) = core.observers.lock() else { 61 return; 62 }; 63 let Ok(retired) = core.retired_observers.lock() else { 64 return; 65 }; 66 observers.get(&id).or_else(|| retired.get(&id)).cloned() 67 }; 68 if let Some(task) = task 69 && let Ok(mut stop) = task.stop.lock() 70 { 71 drop(stop.take()); 72 } 73 } 74 } 75 76 #[cfg_attr(not(coverage_nightly), uniffi::export)] 77 impl ObserverSubscription { 78 pub async fn unsubscribe(&self) { 79 let id = { 80 let Ok(retained_id) = self.id.lock() else { 81 return; 82 }; 83 *retained_id 84 }; 85 let (Some(core), Some(id)) = (self.core.upgrade(), id) else { 86 return; 87 }; 88 let _close = core.close_gate.lock().await; 89 if finish_observers(&core, Some(id)).await.is_err() { 90 return; 91 } 92 let _ = core.actor.unsubscribe_changes(id).await; 93 if let Ok(mut retained_id) = self.id.lock() { 94 *retained_id = None; 95 } 96 } 97 } 98 99 async fn finish_observers( 100 core: &RuntimeCore, 101 selected: Option<ChangeSubscriptionId>, 102 ) -> Result<(), HarvestCircleError> { 103 let tasks = { 104 let observers = core 105 .observers 106 .lock() 107 .map_err(|_| crate::commands::internal_state_unavailable())?; 108 let retired = core 109 .retired_observers 110 .lock() 111 .map_err(|_| crate::commands::internal_state_unavailable())?; 112 observers 113 .iter() 114 .chain(retired.iter()) 115 .filter(|(id, _)| selected.is_none_or(|selected| selected == **id)) 116 .map(|(id, task)| (*id, Arc::clone(task))) 117 .collect::<Vec<_>>() 118 }; 119 for (_, task) in &tasks { 120 if let Some(task) = task.handle.lock().await.as_ref() { 121 task.abort(); 122 } 123 } 124 for (id, task) in tasks { 125 { 126 let mut retained = task.handle.lock().await; 127 if let Some(task) = retained.as_mut() { 128 // The registry retains the exact handle if this awaiting future is cancelled. 129 let _ = task.await; 130 } 131 drop(retained.take()); 132 } 133 // Admission stays reserved until the join, including callback destruction, finishes. 134 let mut observers = core 135 .observers 136 .lock() 137 .map_err(|_| crate::commands::internal_state_unavailable())?; 138 let mut retired = core 139 .retired_observers 140 .lock() 141 .map_err(|_| crate::commands::internal_state_unavailable())?; 142 observers.remove(&id); 143 retired.remove(&id); 144 } 145 Ok(()) 146 } 147 148 fn retire_observer(core: &RuntimeCore, id: ChangeSubscriptionId) { 149 let Ok(mut observers) = core.observers.lock() else { 150 return; 151 }; 152 let Ok(mut retired) = core.retired_observers.lock() else { 153 return; 154 }; 155 if let Some(task) = observers.remove(&id) { 156 retired.insert(id, task); 157 } 158 } 159 160 async fn forward_observer( 161 resources: ObserverResources, 162 runtime_core: Weak<RuntimeCore>, 163 mut subscription: harvestcircle_runtime::RuntimeChangeSubscription, 164 mut stopped: tokio::sync::oneshot::Receiver<()>, 165 ) { 166 let id = subscription.id(); 167 loop { 168 let change = tokio::select! { 169 biased; 170 _ = &mut stopped => break, 171 change = subscription.receive() => change, 172 }; 173 let Some(change) = change else { 174 break; 175 }; 176 let Some(runtime_core) = runtime_core.upgrade() else { 177 break; 178 }; 179 let delivery = SnapshotChangeDto { 180 snapshot: AppSnapshotDto::from_runtime( 181 change.snapshot(), 182 runtime_core.effective_lifecycle(), 183 ), 184 previous_revision: change 185 .previous_revision() 186 .map(harvestcircle_application::SnapshotRevision::value), 187 }; 188 if catch_unwind(AssertUnwindSafe(|| resources.observer.on_change(delivery))).is_err() { 189 break; 190 } 191 } 192 if let Some(runtime_core) = runtime_core.upgrade() { 193 let _ = runtime_core.actor.unsubscribe_changes(id).await; 194 retire_observer(&runtime_core, id); 195 } 196 // Consume the whole bundle so async capture cannot separate the callback and reservation. 197 drop(resources); 198 } 199 200 #[cfg_attr(not(coverage_nightly), uniffi::export)] 201 impl HarvestCircleAppCore { 202 /// Subscribes to ordered revision changes including predecessor metadata. 203 /// 204 /// # Errors 205 /// 206 /// Returns a safe observer or lifecycle error. 207 pub async fn subscribe_changes_v2( 208 &self, 209 observer: Box<dyn HarvestCircleChangeObserver>, 210 ) -> Result<Arc<ObserverSubscription>, HarvestCircleError> { 211 let resources = { 212 let _observers = self 213 .inner 214 .observers 215 .lock() 216 .map_err(|_| observer_registration_error())?; 217 let mut retired = self 218 .inner 219 .retired_observers 220 .lock() 221 .map_err(|_| observer_registration_error())?; 222 retired.retain(|_, task| match task.handle.try_lock() { 223 Ok(retained) => retained.as_ref().is_some_and(|task| !task.is_finished()), 224 Err(_) => true, 225 }); 226 if !self.inner.is_open() { 227 return Err(closed_error()); 228 } 229 let admission = Arc::clone(&self.inner.observer_admission) 230 .try_acquire_owned() 231 .map_err(|_| observer_registration_error())?; 232 ObserverResources { 233 observer, 234 admission: Arc::new(admission), 235 } 236 }; 237 let subscription = self 238 .inner 239 .actor 240 .subscribe_changes(OBSERVER_CHANGE_CAPACITY) 241 .await 242 .map_err(HarvestCircleError::from)?; 243 let id = subscription.id(); 244 let runtime_core = Arc::downgrade(&self.inner); 245 let (stop, stopped) = tokio::sync::oneshot::channel(); 246 let admitted = { 247 let mut observers = self 248 .inner 249 .observers 250 .lock() 251 .map_err(|_| observer_registration_error())?; 252 let _retired = self 253 .inner 254 .retired_observers 255 .lock() 256 .map_err(|_| observer_registration_error())?; 257 if !self.inner.is_open() { 258 false 259 } else { 260 let admission = Arc::clone(&resources.admission); 261 let task = self.inner.runtime.spawn(forward_observer( 262 resources, 263 runtime_core, 264 subscription, 265 stopped, 266 )); 267 observers.insert( 268 id, 269 Arc::new(ObserverTask { 270 handle: tokio::sync::Mutex::new(Some(task)), 271 stop: Mutex::new(Some(stop)), 272 _admission: admission, 273 }), 274 ); 275 true 276 } 277 }; 278 if !admitted { 279 self.inner 280 .actor 281 .unsubscribe_changes(id) 282 .await 283 .map_err(HarvestCircleError::from)?; 284 return Err(observer_registration_error()); 285 } 286 Ok(Arc::new(ObserverSubscription { 287 core: Arc::downgrade(&self.inner), 288 id: Mutex::new(Some(id)), 289 })) 290 } 291 292 /// Stops admission and waits for observer, actor, keyring, and runtime shutdown. 293 /// 294 /// Once shutdown begins, dropping or cancelling the calling future does 295 /// not reopen admission. A later call resumes the same close sequence and 296 /// successful calls are idempotent. 297 /// 298 /// # Errors 299 /// 300 /// Returns a safe closed or timeout error when shutdown cannot complete. 301 pub async fn shutdown_v2(&self) -> Result<ShutdownReceiptDto, HarvestCircleError> { 302 { 303 let _observers = self 304 .inner 305 .observers 306 .lock() 307 .map_err(|_| crate::commands::internal_state_unavailable())?; 308 let _retired = self 309 .inner 310 .retired_observers 311 .lock() 312 .map_err(|_| crate::commands::internal_state_unavailable())?; 313 let _ = 314 self.inner 315 .close_state 316 .compare_exchange(0, 1, Ordering::AcqRel, Ordering::Acquire); 317 } 318 let _close = self.inner.close_gate.lock().await; 319 if self.inner.close_state.load(Ordering::Acquire) == 2 { 320 return Ok(ShutdownReceiptDto { 321 final_revision: self.inner.actor.snapshot().revision().value(), 322 closed: true, 323 }); 324 } 325 finish_observers(&self.inner, None).await?; 326 self.inner 327 .actor 328 .close() 329 .await 330 .map_err(HarvestCircleError::from)?; 331 match ( 332 self.inner.keyring.as_ref(), 333 self.inner.host_runtime.as_ref(), 334 ) { 335 (Some(keyring), Some(runtime)) => { 336 let keyring = Arc::clone(keyring); 337 runtime 338 .run(async move { keyring.close().await }) 339 .await 340 .map_err(|()| closed_error())? 341 .map_err(HarvestCircleError::from)?; 342 } 343 (Some(_), None) => return Err(closed_error()), 344 (None, _) => {} 345 } 346 if let Some(runtime) = self.inner.host_runtime.as_ref() { 347 runtime.shutdown().await.map_err(|()| closed_error())?; 348 } 349 self.inner.close_state.store(2, Ordering::Release); 350 Ok(ShutdownReceiptDto { 351 final_revision: self.inner.actor.snapshot().revision().value(), 352 closed: true, 353 }) 354 } 355 } 356 357 fn closed_error() -> HarvestCircleError { 358 crate::commands::runtime_closed_error() 359 } 360 361 fn observer_registration_error() -> HarvestCircleError { 362 HarvestCircleError::Failure { 363 code: crate::WireErrorCode::ObserverRegistrationFailed, 364 category: crate::WireErrorCategory::Lifecycle, 365 retryable: true, 366 recovery_action: crate::WireRecoveryAction::Retry, 367 correlation_id: None, 368 safe_message: "The change observer could not be registered.".to_owned(), 369 } 370 } 371 372 #[cfg(test)] 373 #[cfg_attr(coverage_nightly, coverage(off))] 374 mod tests { 375 use std::future::Future; 376 use std::sync::atomic::{AtomicBool, Ordering}; 377 use std::sync::{Arc, Mutex}; 378 use std::task::Poll; 379 use std::time::Duration; 380 381 use harvestcircle_application::{ 382 DurableRequestId, RelayAccess, RelayConfiguration, RelayEndpoint, RelayUrlPolicy, 383 }; 384 use harvestcircle_domain::SecretKeyInput; 385 use nostr::{EventBuilder, Keys, Metadata}; 386 use nostr_relay_builder::MockRelay; 387 use nostr_sdk::Client; 388 389 use crate::commands::{RuntimeCore, test_actor, test_actor_with_nostr_timeout}; 390 use crate::{ 391 AppSnapshotDto, HarvestCircleAppCore, HarvestCircleChangeObserver, ProfileLoadStateDto, 392 SnapshotChangeDto, 393 }; 394 395 const SECRET_HEX: &str = "7e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7"; 396 const TEST_RELAY_TIMEOUT: Duration = Duration::from_secs(2); 397 #[derive(Default)] 398 struct RecordingObserver { 399 snapshots: Mutex<Vec<AppSnapshotDto>>, 400 core: Mutex<Option<Arc<HarvestCircleAppCore>>>, 401 } 402 403 struct PanickingObserver; 404 405 #[derive(Default)] 406 struct ObserverDropState { 407 finished: AtomicBool, 408 released: AtomicBool, 409 } 410 411 struct GatedDropObserver { 412 entered: Option<tokio::sync::oneshot::Sender<()>>, 413 release: Mutex<std::sync::mpsc::Receiver<()>>, 414 state: Arc<ObserverDropState>, 415 } 416 417 impl HarvestCircleChangeObserver for GatedDropObserver { 418 fn on_change(&self, _change: SnapshotChangeDto) {} 419 } 420 421 impl Drop for GatedDropObserver { 422 fn drop(&mut self) { 423 if let Some(entered) = self.entered.take() { 424 let _ = entered.send(()); 425 } 426 let released = self 427 .release 428 .get_mut() 429 .is_ok_and(|release| release.recv_timeout(OBSERVER_DELIVERY_TIMEOUT).is_ok()); 430 self.state.released.store(released, Ordering::Release); 431 self.state.finished.store(true, Ordering::Release); 432 } 433 } 434 435 struct GatedObserver { 436 changes: Mutex<Vec<SnapshotChangeDto>>, 437 entered: Mutex<Option<tokio::sync::oneshot::Sender<()>>>, 438 release: Mutex<std::sync::mpsc::Receiver<()>>, 439 delivered: tokio::sync::Notify, 440 } 441 442 impl HarvestCircleChangeObserver for Arc<GatedObserver> { 443 fn on_change(&self, change: SnapshotChangeDto) { 444 self.changes.lock().expect("changes").push(change); 445 if let Some(entered) = self.entered.lock().expect("initial callback gate").take() { 446 entered.send(()).expect("initial callback entered"); 447 self.release 448 .lock() 449 .expect("callback release gate") 450 .recv_timeout(OBSERVER_DELIVERY_TIMEOUT) 451 .expect("release initial callback"); 452 } 453 self.delivered.notify_one(); 454 } 455 } 456 457 impl HarvestCircleChangeObserver for PanickingObserver { 458 fn on_change(&self, _change: SnapshotChangeDto) { 459 panic!("injected host callback failure"); 460 } 461 } 462 463 impl HarvestCircleChangeObserver for RecordingObserver { 464 fn on_change(&self, change: SnapshotChangeDto) { 465 let snapshot = change.snapshot; 466 if let Some(core) = self.core.lock().expect("core").as_ref() { 467 assert_eq!(core.snapshot().revision, snapshot.revision); 468 } 469 self.snapshots.lock().expect("snapshots").push(snapshot); 470 } 471 } 472 473 async fn core() -> Arc<HarvestCircleAppCore> { 474 core_with_relays(RelayConfiguration::default()).await 475 } 476 477 async fn core_with_relays(relays: RelayConfiguration) -> Arc<HarvestCircleAppCore> { 478 let (actor, directory) = test_actor(relays).await; 479 core_with_actor(actor, directory) 480 } 481 482 async fn core_with_live_relays(relays: RelayConfiguration) -> Arc<HarvestCircleAppCore> { 483 let (actor, directory) = test_actor_with_nostr_timeout(relays, TEST_RELAY_TIMEOUT).await; 484 core_with_actor(actor, directory) 485 } 486 487 fn core_with_actor( 488 actor: harvestcircle_runtime::RuntimeActorHandle, 489 directory: Arc<tempfile::TempDir>, 490 ) -> Arc<HarvestCircleAppCore> { 491 Arc::new(HarvestCircleAppCore { 492 inner: Arc::new(RuntimeCore { 493 actor, 494 runtime: tokio::runtime::Handle::current(), 495 host_runtime: None, 496 keyring: None, 497 observers: Mutex::new(std::collections::BTreeMap::new()), 498 retired_observers: Mutex::new(std::collections::BTreeMap::new()), 499 observer_admission: Arc::new(tokio::sync::Semaphore::new(super::MAX_OBSERVERS)), 500 close_state: std::sync::atomic::AtomicU8::new(0), 501 close_gate: tokio::sync::Mutex::new(()), 502 _test_directory: Some(directory), 503 }), 504 }) 505 } 506 507 async fn core_with_host_runtime( 508 host_runtime: Arc<crate::host_runtime::HostRuntime>, 509 ) -> Arc<HarvestCircleAppCore> { 510 let (actor, directory) = test_actor(RelayConfiguration::default()).await; 511 Arc::new(HarvestCircleAppCore { 512 inner: Arc::new(RuntimeCore { 513 actor, 514 runtime: tokio::runtime::Handle::current(), 515 host_runtime: Some(host_runtime), 516 keyring: None, 517 observers: Mutex::new(std::collections::BTreeMap::new()), 518 retired_observers: Mutex::new(std::collections::BTreeMap::new()), 519 observer_admission: Arc::new(tokio::sync::Semaphore::new(super::MAX_OBSERVERS)), 520 close_state: std::sync::atomic::AtomicU8::new(0), 521 close_gate: tokio::sync::Mutex::new(()), 522 _test_directory: Some(directory), 523 }), 524 }) 525 } 526 527 fn test_runtime() -> tokio::runtime::Runtime { 528 tokio::runtime::Runtime::new().expect("test runtime") 529 } 530 531 #[test] 532 fn callbacks_allow_reentry_and_stop_after_subscription_close() { 533 test_runtime().block_on(async { 534 let core = core().await; 535 let observer = Arc::new(RecordingObserver::default()); 536 *observer.core.lock().expect("core") = Some(Arc::clone(&core)); 537 let subscription = core 538 .subscribe_changes_v2(Box::new(ArcObserver(observer.clone()))) 539 .await 540 .expect("subscribe"); 541 wait_for_snapshot_count(&observer, 1).await; 542 core.inner 543 .actor 544 .bootstrap() 545 .await 546 .expect("idempotent bootstrap"); 547 assert_eq!(observer.snapshots.lock().expect("snapshots").len(), 1); 548 subscription.unsubscribe().await; 549 subscription.unsubscribe().await; 550 core.inner.actor.sign_out().await.expect("sign out"); 551 assert_eq!(observer.snapshots.lock().expect("snapshots").len(), 1); 552 }); 553 } 554 555 #[tokio::test(flavor = "multi_thread", worker_threads = 2)] 556 async fn cancelled_unsubscribe_retry_waits_for_the_retained_callback_task() { 557 let core = core().await; 558 let (observer, entered, release) = gated_observer(); 559 let subscription = core 560 .subscribe_changes_v2(Box::new(Arc::clone(&observer))) 561 .await 562 .expect("subscribe gated callback"); 563 tokio::time::timeout(OBSERVER_DELIVERY_TIMEOUT, entered) 564 .await 565 .expect("callback entry deadline") 566 .expect("callback entered"); 567 let task = observer_abort_handle(&core, &subscription); 568 let first_was_pending = { 569 let mut first = Box::pin(subscription.unsubscribe()); 570 std::future::poll_fn(|context| Poll::Ready(first.as_mut().poll(context))) 571 .await 572 .is_pending() 573 }; 574 let mut retry = Box::pin(subscription.unsubscribe()); 575 let retry_was_pending = 576 std::future::poll_fn(|context| Poll::Ready(retry.as_mut().poll(context))) 577 .await 578 .is_pending(); 579 let finished_before_release = task.is_finished(); 580 581 release 582 .send(()) 583 .expect("release callback before assertions"); 584 if retry_was_pending { 585 tokio::time::timeout(OBSERVER_DELIVERY_TIMEOUT, retry) 586 .await 587 .expect("unsubscribe retry completion"); 588 } 589 wait_for_aborted_task(&task).await; 590 core.shutdown_v2().await.expect("cleanup runtime"); 591 592 assert!( 593 first_was_pending, 594 "first unsubscribe must wait for callback exit" 595 ); 596 assert!(!finished_before_release, "callback was explicitly gated"); 597 assert!( 598 retry_was_pending, 599 "retry must retain ownership of callback cleanup" 600 ); 601 assert!(task.is_finished()); 602 assert_eq!(observer.changes.lock().expect("changes").len(), 1); 603 subscription.unsubscribe().await; 604 } 605 606 #[tokio::test(flavor = "multi_thread", worker_threads = 2)] 607 async fn cancelled_shutdown_retry_finishes_observers_before_terminal_close() { 608 let core = core().await; 609 let (observer, entered, release) = gated_observer(); 610 let subscription = core 611 .subscribe_changes_v2(Box::new(Arc::clone(&observer))) 612 .await 613 .expect("subscribe gated callback"); 614 tokio::time::timeout(OBSERVER_DELIVERY_TIMEOUT, entered) 615 .await 616 .expect("callback entry deadline") 617 .expect("callback entered"); 618 let task = observer_abort_handle(&core, &subscription); 619 let first_was_pending = { 620 let mut first = Box::pin(core.shutdown_v2()); 621 std::future::poll_fn(|context| Poll::Ready(first.as_mut().poll(context))) 622 .await 623 .is_pending() 624 }; 625 let mut retry = Box::pin(core.shutdown_v2()); 626 let initial_retry = 627 std::future::poll_fn(|context| Poll::Ready(retry.as_mut().poll(context))).await; 628 // This command orders actor work after any close submitted by the retry. 629 let _ = core.inner.actor.bootstrap().await; 630 let retried = if initial_retry.is_pending() { 631 std::future::poll_fn(|context| Poll::Ready(retry.as_mut().poll(context))).await 632 } else { 633 initial_retry 634 }; 635 let closed_before_callback_exit = retried.is_ready(); 636 let finished_before_release = task.is_finished(); 637 638 release 639 .send(()) 640 .expect("release callback before assertions"); 641 let receipt = match retried { 642 Poll::Ready(receipt) => receipt.expect("retry close receipt"), 643 Poll::Pending => tokio::time::timeout(OBSERVER_DELIVERY_TIMEOUT, retry) 644 .await 645 .expect("resumed shutdown completion") 646 .expect("resumed shutdown receipt"), 647 }; 648 wait_for_aborted_task(&task).await; 649 subscription.unsubscribe().await; 650 assert_eq!(core.shutdown_v2().await.expect("repeated close"), receipt); 651 652 assert!( 653 first_was_pending, 654 "first shutdown must wait for callback exit" 655 ); 656 assert!(!finished_before_release, "callback was explicitly gated"); 657 assert!( 658 !closed_before_callback_exit, 659 "terminal close must retain observer cleanup ownership" 660 ); 661 assert!(receipt.closed); 662 assert!(task.is_finished()); 663 assert!(core.inner.observers.lock().expect("observers").is_empty()); 664 assert_eq!(observer.changes.lock().expect("changes").len(), 1); 665 } 666 667 #[tokio::test] 668 async fn abandoned_observer_handles_release_bounded_registration_admission() { 669 let core = core().await; 670 let observer = Arc::new(RecordingObserver::default()); 671 for _ in 0..super::MAX_OBSERVERS { 672 let subscription = core 673 .subscribe_changes_v2(Box::new(ArcObserver(Arc::clone(&observer)))) 674 .await 675 .expect("admit observer before abandonment"); 676 drop(subscription); 677 } 678 core.inner 679 .actor 680 .bootstrap() 681 .await 682 .expect("actor ordering barrier"); 683 let cleanup_completed = tokio::time::timeout(OBSERVER_DELIVERY_TIMEOUT, async { 684 while !core.inner.observers.lock().expect("observers").is_empty() { 685 tokio::task::yield_now().await; 686 } 687 }) 688 .await 689 .is_ok(); 690 let registered_after_abandonment = core.inner.observers.lock().expect("observers").len(); 691 let replacement = core 692 .subscribe_changes_v2(Box::new(ArcObserver(observer))) 693 .await; 694 let replacement_was_admitted = replacement.is_ok(); 695 if let Ok(replacement) = replacement { 696 replacement.unsubscribe().await; 697 } 698 core.shutdown_v2() 699 .await 700 .expect("cleanup abandoned observers"); 701 702 assert!( 703 cleanup_completed, 704 "bounded observer task cleanup must finish after abandonment" 705 ); 706 assert_eq!( 707 registered_after_abandonment, 0, 708 "abandoned handles must release admission" 709 ); 710 assert!( 711 replacement_was_admitted, 712 "abandonment must not exhaust the fixed observer limit" 713 ); 714 } 715 716 #[tokio::test] 717 async fn cancelled_native_registration_before_and_after_actor_reply_never_delivers() { 718 let core = core().await; 719 for cancel_after_actor_reply in [false, true] { 720 let observer = Arc::new(RecordingObserver::default()); 721 let mut registration = 722 Box::pin(core.subscribe_changes_v2(Box::new(ArcObserver(Arc::clone(&observer))))); 723 let first_poll = 724 std::future::poll_fn(|context| Poll::Ready(registration.as_mut().poll(context))) 725 .await; 726 assert!( 727 first_poll.is_pending(), 728 "registration queued on the current-thread actor" 729 ); 730 if cancel_after_actor_reply { 731 core.inner 732 .actor 733 .bootstrap() 734 .await 735 .expect("actor installed registration before cancellation"); 736 } 737 drop(registration); 738 core.inner 739 .actor 740 .bootstrap() 741 .await 742 .expect("actor cancellation ordering barrier"); 743 assert!(core.inner.observers.lock().expect("observers").is_empty()); 744 assert!(observer.snapshots.lock().expect("snapshots").is_empty()); 745 } 746 core.shutdown_v2().await.expect("cleanup runtime"); 747 } 748 749 #[tokio::test] 750 async fn pending_native_registrations_share_the_fixed_observer_admission_limit() { 751 let core = core().await; 752 let observer = Arc::new(RecordingObserver::default()); 753 let mut pending = Vec::with_capacity(super::MAX_OBSERVERS); 754 for _ in 0..super::MAX_OBSERVERS { 755 let mut registration = 756 Box::pin(core.subscribe_changes_v2(Box::new(ArcObserver(Arc::clone(&observer))))); 757 let first_poll = 758 std::future::poll_fn(|context| Poll::Ready(registration.as_mut().poll(context))) 759 .await; 760 assert!( 761 first_poll.is_pending(), 762 "admitted registration awaits the actor" 763 ); 764 core.inner 765 .actor 766 .bootstrap() 767 .await 768 .expect("actor installed the unpolled reply"); 769 pending.push(registration); 770 } 771 assert!(core.inner.observers.lock().expect("observers").is_empty()); 772 let mut excess = 773 Box::pin(core.subscribe_changes_v2(Box::new(ArcObserver(Arc::clone(&observer))))); 774 let refused = 775 std::future::poll_fn(|context| Poll::Ready(excess.as_mut().poll(context))).await; 776 let refused_before_actor_allocation = matches!( 777 refused, 778 Poll::Ready(Err(crate::HarvestCircleError::Failure { 779 code: crate::WireErrorCode::ObserverRegistrationFailed, 780 .. 781 })) 782 ); 783 drop(excess); 784 drop(pending); 785 core.inner 786 .actor 787 .bootstrap() 788 .await 789 .expect("cancelled reply cleanup barrier"); 790 let replacement = core 791 .subscribe_changes_v2(Box::new(ArcObserver(Arc::clone(&observer)))) 792 .await 793 .expect("replacement after pending cancellation"); 794 replacement.unsubscribe().await; 795 core.shutdown_v2().await.expect("cleanup runtime"); 796 797 assert!( 798 refused_before_actor_allocation, 799 "pending replies must consume the same fixed observer slots" 800 ); 801 assert_eq!( 802 core.inner.observer_admission.available_permits(), 803 super::MAX_OBSERVERS 804 ); 805 } 806 807 #[tokio::test] 808 async fn cancelled_pending_registration_retains_admission_until_callback_destruction_finishes() 809 { 810 let core = core().await; 811 let state = Arc::new(ObserverDropState::default()); 812 let (entered_sender, entered_receiver) = tokio::sync::oneshot::channel(); 813 let (release_sender, release_receiver) = std::sync::mpsc::sync_channel(1); 814 let pending_core = Arc::clone(&core); 815 let drop_state = Arc::clone(&state); 816 let mut gated_registration = Box::pin(async move { 817 pending_core 818 .subscribe_changes_v2(Box::new(GatedDropObserver { 819 entered: Some(entered_sender), 820 release: Mutex::new(release_receiver), 821 state: drop_state, 822 })) 823 .await 824 }); 825 let first_poll = 826 std::future::poll_fn(|context| Poll::Ready(gated_registration.as_mut().poll(context))) 827 .await; 828 assert!( 829 first_poll.is_pending(), 830 "gated registration awaits the actor" 831 ); 832 core.inner 833 .actor 834 .bootstrap() 835 .await 836 .expect("gated actor reply ready but unpolled"); 837 838 let observer = Arc::new(RecordingObserver::default()); 839 let mut pending = Vec::with_capacity(super::MAX_OBSERVERS - 1); 840 for _ in 1..super::MAX_OBSERVERS { 841 let mut registration = 842 Box::pin(core.subscribe_changes_v2(Box::new(ArcObserver(Arc::clone(&observer))))); 843 let first_poll = 844 std::future::poll_fn(|context| Poll::Ready(registration.as_mut().poll(context))) 845 .await; 846 assert!( 847 first_poll.is_pending(), 848 "peer registration awaits the actor" 849 ); 850 core.inner 851 .actor 852 .bootstrap() 853 .await 854 .expect("peer actor reply ready but unpolled"); 855 pending.push(registration); 856 } 857 assert_eq!(core.inner.observer_admission.available_permits(), 0); 858 859 let dropping = std::thread::spawn(move || drop(gated_registration)); 860 let entered = tokio::time::timeout(OBSERVER_DELIVERY_TIMEOUT, entered_receiver) 861 .await 862 .is_ok_and(|entered| entered.is_ok()); 863 let destructor_was_gated = !state.finished.load(Ordering::Acquire); 864 let mut replacement = 865 Box::pin(core.subscribe_changes_v2(Box::new(ArcObserver(Arc::clone(&observer))))); 866 let replacement_poll = 867 std::future::poll_fn(|context| Poll::Ready(replacement.as_mut().poll(context))).await; 868 let refused_during_destruction = matches!( 869 replacement_poll, 870 Poll::Ready(Err(crate::HarvestCircleError::Failure { 871 code: crate::WireErrorCode::ObserverRegistrationFailed, 872 .. 873 })) 874 ); 875 876 let released = release_sender.send(()).is_ok(); 877 let joined = dropping.join().is_ok(); 878 drop(replacement); 879 drop(pending); 880 core.inner 881 .actor 882 .bootstrap() 883 .await 884 .expect("cancelled registration cleanup barrier"); 885 let resumed = core 886 .subscribe_changes_v2(Box::new(ArcObserver(observer))) 887 .await 888 .expect("registration after callback destruction"); 889 resumed.unsubscribe().await; 890 core.shutdown_v2() 891 .await 892 .expect("cleanup runtime before assertions"); 893 894 assert!(entered, "callback destructor entry was observed"); 895 assert!( 896 destructor_was_gated, 897 "replacement was probed during callback destruction" 898 ); 899 assert!( 900 released && joined, 901 "the destructor gate was released and its owned thread joined" 902 ); 903 assert!(state.released.load(Ordering::Acquire)); 904 assert!(state.finished.load(Ordering::Acquire)); 905 assert!( 906 refused_during_destruction, 907 "callback destruction must retain its pending admission slot" 908 ); 909 assert_eq!( 910 core.inner.observer_admission.available_permits(), 911 super::MAX_OBSERVERS 912 ); 913 } 914 915 #[tokio::test] 916 async fn ready_native_registration_reply_cannot_install_after_terminal_close() { 917 let core = core().await; 918 let observer = Arc::new(RecordingObserver::default()); 919 let mut registration = 920 Box::pin(core.subscribe_changes_v2(Box::new(ArcObserver(Arc::clone(&observer))))); 921 let first_poll = 922 std::future::poll_fn(|context| Poll::Ready(registration.as_mut().poll(context))).await; 923 assert!( 924 first_poll.is_pending(), 925 "registration queued on the current-thread actor" 926 ); 927 core.inner 928 .actor 929 .bootstrap() 930 .await 931 .expect("actor reply ready before close"); 932 933 let closed = core 934 .shutdown_v2() 935 .await 936 .expect("close before host registration resumes"); 937 assert!( 938 registration.await.is_err(), 939 "closed registration must not install a callback task" 940 ); 941 942 assert!(closed.closed); 943 assert!(observer.snapshots.lock().expect("snapshots").is_empty()); 944 assert!(core.inner.observers.lock().expect("observers").is_empty()); 945 assert!( 946 core.inner 947 .retired_observers 948 .lock() 949 .expect("retired observers") 950 .is_empty() 951 ); 952 assert_eq!( 953 core.inner.observer_admission.available_permits(), 954 super::MAX_OBSERVERS 955 ); 956 } 957 958 #[tokio::test(flavor = "multi_thread", worker_threads = 2)] 959 async fn full_native_callback_queue_terminates_before_shutdown_returns() { 960 let core = core().await; 961 let mut public_keys = Vec::new(); 962 for request_id in [ 963 "01890f3e-7b1c-7000-8000-000000000052", 964 "01890f3e-7b1c-7000-8000-000000000053", 965 ] { 966 let imported = core 967 .inner 968 .actor 969 .import_secret_key( 970 DurableRequestId::parse(request_id).expect("test import request"), 971 core.inner.actor.snapshot().revision(), 972 SecretKeyInput::parse(Keys::generate().secret_key().to_secret_hex()) 973 .expect("ephemeral identity input"), 974 OBSERVER_DELIVERY_TIMEOUT, 975 ) 976 .await 977 .expect("isolated memory-store import"); 978 public_keys.push(imported.identity().public_key()); 979 } 980 core.inner 981 .actor 982 .select_identity(public_keys[1]) 983 .await 984 .expect("initial selection"); 985 let initial_revision = core.snapshot().revision; 986 let (observer, entered, release) = gated_observer(); 987 let subscription = core 988 .subscribe_changes_v2(Box::new(Arc::clone(&observer))) 989 .await 990 .expect("gated observer registration"); 991 tokio::time::timeout(OBSERVER_DELIVERY_TIMEOUT, entered) 992 .await 993 .expect("callback entry deadline") 994 .expect("callback entered"); 995 let task = observer_abort_handle(&core, &subscription); 996 let publications = super::OBSERVER_CHANGE_CAPACITY.get() + 2; 997 for offset in 0..publications { 998 core.inner 999 .actor 1000 .select_identity(public_keys[offset % 2]) 1001 .await 1002 .expect("revision-changing selection"); 1003 } 1004 let final_revision = core.snapshot().revision; 1005 let mut closing = Box::pin(core.shutdown_v2()); 1006 let first_close_poll = 1007 std::future::poll_fn(|context| Poll::Ready(closing.as_mut().poll(context))).await; 1008 let closing_was_pending = first_close_poll.is_pending(); 1009 release 1010 .send(()) 1011 .expect("release saturated callback before assertions"); 1012 let receipt = match first_close_poll { 1013 Poll::Ready(receipt) => receipt.expect("shutdown receipt"), 1014 Poll::Pending => tokio::time::timeout(OBSERVER_DELIVERY_TIMEOUT, closing) 1015 .await 1016 .expect("bounded full-queue close") 1017 .expect("shutdown receipt"), 1018 }; 1019 let callback_count_at_close = observer.changes.lock().expect("changes").len(); 1020 subscription.unsubscribe().await; 1021 subscription.unsubscribe().await; 1022 let repeated = core.shutdown_v2().await.expect("repeated close"); 1023 let refused = core 1024 .subscribe_changes_v2(Box::new(PanickingObserver)) 1025 .await 1026 .is_err(); 1027 1028 assert!( 1029 closing_was_pending, 1030 "close must wait for the running callback" 1031 ); 1032 assert_eq!( 1033 final_revision, 1034 initial_revision + u64::try_from(publications).expect("publication count") 1035 ); 1036 assert_eq!(receipt.final_revision, final_revision); 1037 assert!(receipt.closed); 1038 assert_eq!(repeated, receipt); 1039 assert!(task.is_finished()); 1040 assert!(refused); 1041 assert!(core.inner.observers.lock().expect("observers").is_empty()); 1042 assert_eq!( 1043 observer.changes.lock().expect("changes").len(), 1044 callback_count_at_close 1045 ); 1046 } 1047 1048 fn gated_observer() -> ( 1049 Arc<GatedObserver>, 1050 tokio::sync::oneshot::Receiver<()>, 1051 std::sync::mpsc::SyncSender<()>, 1052 ) { 1053 let (entered_sender, entered_receiver) = tokio::sync::oneshot::channel(); 1054 let (release_sender, release_receiver) = std::sync::mpsc::sync_channel(1); 1055 ( 1056 Arc::new(GatedObserver { 1057 changes: Mutex::new(Vec::new()), 1058 entered: Mutex::new(Some(entered_sender)), 1059 release: Mutex::new(release_receiver), 1060 delivered: tokio::sync::Notify::new(), 1061 }), 1062 entered_receiver, 1063 release_sender, 1064 ) 1065 } 1066 1067 fn observer_abort_handle( 1068 core: &HarvestCircleAppCore, 1069 subscription: &crate::ObserverSubscription, 1070 ) -> tokio::task::AbortHandle { 1071 let id = subscription.id.lock().expect("id").expect("registered id"); 1072 let task = core 1073 .inner 1074 .observers 1075 .lock() 1076 .expect("observers") 1077 .get(&id) 1078 .expect("registered observer") 1079 .clone(); 1080 let retained = task.handle.try_lock().expect("observer join not started"); 1081 retained.as_ref().expect("observer task").abort_handle() 1082 } 1083 1084 async fn wait_for_aborted_task(task: &tokio::task::AbortHandle) { 1085 tokio::time::timeout(OBSERVER_DELIVERY_TIMEOUT, async { 1086 while !task.is_finished() { 1087 tokio::task::yield_now().await; 1088 } 1089 }) 1090 .await 1091 .expect("released callback task finishes"); 1092 } 1093 1094 #[test] 1095 fn slow_callback_recovers_final_tail_after_actor_queue_saturation() { 1096 let runtime = tokio::runtime::Builder::new_multi_thread() 1097 .worker_threads(2) 1098 .enable_all() 1099 .build() 1100 .expect("two-worker callback test runtime"); 1101 runtime.block_on(async { 1102 let core = core().await; 1103 let mut public_keys = Vec::with_capacity(2); 1104 for request_id in [ 1105 "01890f3e-7b1c-7000-8000-000000000050", 1106 "01890f3e-7b1c-7000-8000-000000000051", 1107 ] { 1108 let imported = core 1109 .inner 1110 .actor 1111 .import_secret_key( 1112 DurableRequestId::parse(request_id).expect("test import request"), 1113 core.inner.actor.snapshot().revision(), 1114 SecretKeyInput::parse(Keys::generate().secret_key().to_secret_hex()) 1115 .expect("ephemeral identity input"), 1116 OBSERVER_DELIVERY_TIMEOUT, 1117 ) 1118 .await 1119 .expect("import into isolated memory secret store"); 1120 public_keys.push(imported.identity().public_key()); 1121 } 1122 assert_ne!(public_keys[0], public_keys[1]); 1123 core.inner 1124 .actor 1125 .select_identity(public_keys[1]) 1126 .await 1127 .expect("select second identity before observation"); 1128 let initial_revision = core.snapshot().revision; 1129 let (entered_sender, entered_receiver) = tokio::sync::oneshot::channel(); 1130 let (release_sender, release_receiver) = std::sync::mpsc::sync_channel(1); 1131 let observer = Arc::new(GatedObserver { 1132 changes: Mutex::new(Vec::new()), 1133 entered: Mutex::new(Some(entered_sender)), 1134 release: Mutex::new(release_receiver), 1135 delivered: tokio::sync::Notify::new(), 1136 }); 1137 let subscription = core 1138 .subscribe_changes_v2(Box::new(Arc::clone(&observer))) 1139 .await 1140 .expect("subscribe slow callback"); 1141 tokio::time::timeout(OBSERVER_DELIVERY_TIMEOUT, entered_receiver) 1142 .await 1143 .expect("initial callback deadline") 1144 .expect("initial callback entered"); 1145 1146 let queued_changes = super::OBSERVER_CHANGE_CAPACITY.get(); 1147 let publications = queued_changes + 2; 1148 let final_revision = tokio::time::timeout(OBSERVER_DELIVERY_TIMEOUT, async { 1149 let mut final_revision = initial_revision; 1150 for offset in 0..publications { 1151 final_revision = core 1152 .inner 1153 .actor 1154 .select_identity(public_keys[offset % public_keys.len()]) 1155 .await 1156 .expect("alternate selected identity") 1157 .revision() 1158 .value(); 1159 } 1160 final_revision 1161 }) 1162 .await 1163 .expect("bounded actor publications"); 1164 assert_eq!( 1165 final_revision, 1166 initial_revision + u64::try_from(publications).expect("publication count") 1167 ); 1168 assert_eq!(observer.changes.lock().expect("changes").len(), 1); 1169 release_sender.send(()).expect("release slow callback"); 1170 1171 tokio::time::timeout(OBSERVER_DELIVERY_TIMEOUT, async { 1172 loop { 1173 let delivered = observer.delivered.notified(); 1174 if observer 1175 .changes 1176 .lock() 1177 .expect("changes") 1178 .last() 1179 .is_some_and(|change| change.snapshot.revision == final_revision) 1180 { 1181 break; 1182 } 1183 delivered.await; 1184 } 1185 }) 1186 .await 1187 .expect("final callback without a later publication"); 1188 1189 let changes = observer.changes.lock().expect("changes").clone(); 1190 let mut expected_revisions = vec![initial_revision]; 1191 expected_revisions.extend( 1192 (1..=queued_changes) 1193 .map(|offset| initial_revision + u64::try_from(offset).expect("queue offset")), 1194 ); 1195 expected_revisions.push(final_revision); 1196 assert_eq!( 1197 changes 1198 .iter() 1199 .map(|change| change.snapshot.revision) 1200 .collect::<Vec<_>>(), 1201 expected_revisions 1202 ); 1203 assert_eq!(changes[0].previous_revision, None); 1204 for change in &changes[1..] { 1205 assert_eq!(change.previous_revision, Some(change.snapshot.revision - 1)); 1206 } 1207 subscription.unsubscribe().await; 1208 core.shutdown_v2().await.expect("shutdown"); 1209 }); 1210 } 1211 1212 #[test] 1213 fn core_close_deregisters_all_observers_and_rejects_new_subscriptions() { 1214 test_runtime().block_on(async { 1215 let core = core().await; 1216 let observer = Arc::new(RecordingObserver::default()); 1217 let subscription = core 1218 .subscribe_changes_v2(Box::new(ArcObserver(observer.clone()))) 1219 .await 1220 .expect("subscribe"); 1221 let _active_subscription = core 1222 .subscribe_changes_v2(Box::new(ArcObserver(observer.clone()))) 1223 .await 1224 .expect("second subscription"); 1225 let id = subscription 1226 .id 1227 .lock() 1228 .expect("subscription id") 1229 .expect("active subscription id"); 1230 let task = core 1231 .inner 1232 .observers 1233 .lock() 1234 .expect("observers") 1235 .get(&id) 1236 .expect("registered observer") 1237 .clone(); 1238 task.handle 1239 .lock() 1240 .await 1241 .as_ref() 1242 .expect("observer task") 1243 .abort(); 1244 1245 let first = core.shutdown_v2().await.expect("shutdown"); 1246 let repeated = core.shutdown_v2().await.expect("repeated shutdown"); 1247 assert_eq!(repeated, first); 1248 1249 assert!( 1250 core.subscribe_changes_v2(Box::new(ArcObserver(observer))) 1251 .await 1252 .is_err() 1253 ); 1254 assert!(core.inner.observers.lock().expect("observers").is_empty()); 1255 }); 1256 } 1257 1258 #[tokio::test] 1259 async fn cancelled_host_close_remains_non_admitting_and_resumes() { 1260 let gated = crate::host_runtime::HostRuntime::new_completion_gated_for_test() 1261 .expect("host runtime"); 1262 let core = core_with_host_runtime(gated.runtime).await; 1263 let closing_core = Arc::clone(&core); 1264 let closing = tokio::spawn(async move { closing_core.shutdown_v2().await }); 1265 tokio::task::spawn_blocking(move || gated.entered.recv()) 1266 .await 1267 .expect("entered join") 1268 .expect("close reached host completion gate"); 1269 closing.abort(); 1270 assert!(closing.await.is_err()); 1271 assert!( 1272 core.subscribe_changes_v2(Box::new(PanickingObserver)) 1273 .await 1274 .is_err() 1275 ); 1276 gated.release.send(()).expect("release host close"); 1277 assert!(core.shutdown_v2().await.expect("resumed close").closed); 1278 } 1279 1280 #[test] 1281 fn subscription_unsubscribe_tolerates_a_dropped_runtime_core() { 1282 test_runtime().block_on(async { 1283 let core = core().await; 1284 let observer = Arc::new(RecordingObserver::default()); 1285 let subscription = core 1286 .subscribe_changes_v2(Box::new(ArcObserver(observer.clone()))) 1287 .await 1288 .expect("subscribe"); 1289 wait_for_snapshot_count(&observer, 1).await; 1290 1291 drop(core); 1292 subscription.unsubscribe().await; 1293 }); 1294 } 1295 1296 #[test] 1297 fn observer_registration_is_bounded_and_callback_panics_are_contained() { 1298 test_runtime().block_on(async { 1299 let core = core().await; 1300 let panic_subscription = core 1301 .subscribe_changes_v2(Box::new(PanickingObserver)) 1302 .await 1303 .expect("panic observer registration"); 1304 wait_for_observer_count(&core, 0).await; 1305 panic_subscription.unsubscribe().await; 1306 1307 let observer = Arc::new(RecordingObserver::default()); 1308 let mut subscriptions = Vec::new(); 1309 for _ in 0..super::MAX_OBSERVERS { 1310 subscriptions.push( 1311 core.subscribe_changes_v2(Box::new(ArcObserver(observer.clone()))) 1312 .await 1313 .expect("bounded observer registration"), 1314 ); 1315 } 1316 assert!( 1317 core.subscribe_changes_v2(Box::new(ArcObserver(observer))) 1318 .await 1319 .is_err() 1320 ); 1321 for subscription in subscriptions { 1322 subscription.unsubscribe().await; 1323 } 1324 assert!(core.inner.observers.lock().expect("observers").is_empty()); 1325 }); 1326 } 1327 1328 #[tokio::test] 1329 async fn ffi_callback_receives_async_profile_refresh_and_stops_after_unsubscribe() { 1330 let local_relay = MockRelay::run().await.expect("local relay"); 1331 let relay_url = local_relay.url().await; 1332 let publisher = Client::new(Keys::parse(SECRET_HEX).expect("known key")); 1333 publisher 1334 .add_relay(relay_url.clone()) 1335 .await 1336 .expect("publisher relay"); 1337 publisher.connect().await; 1338 publisher.wait_for_connection(Duration::from_secs(2)).await; 1339 publisher 1340 .send_event_builder(EventBuilder::metadata( 1341 &Metadata::new().display_name("FFI Profile"), 1342 )) 1343 .await 1344 .expect("publish profile"); 1345 1346 let core = core_with_live_relays( 1347 RelayConfiguration::new(vec![ 1348 RelayEndpoint::new( 1349 relay_url.as_str(), 1350 RelayUrlPolicy::Local, 1351 RelayAccess::ReadWrite, 1352 ) 1353 .expect("relay endpoint"), 1354 ]) 1355 .expect("relay configuration"), 1356 ) 1357 .await; 1358 core.bootstrap().await.expect("bootstrap"); 1359 let observer = Arc::new(RecordingObserver::default()); 1360 *observer.core.lock().expect("core") = Some(Arc::clone(&core)); 1361 let subscription = core 1362 .subscribe_changes_v2(Box::new(ArcObserver(observer.clone()))) 1363 .await 1364 .expect("subscribe"); 1365 let imported = core 1366 .import_identity( 1367 crate::RequestContextDto { 1368 request_id: "01890f3e-7b1c-7000-8000-000000000049".to_owned(), 1369 expected_revision: core.snapshot().revision, 1370 deadline_millis: 5_000, 1371 }, 1372 SECRET_HEX.as_bytes().to_vec(), 1373 ) 1374 .await 1375 .expect("import") 1376 .snapshot; 1377 let public_key = imported.selected_public_key_hex.expect("selection"); 1378 core.activate_identity(public_key).await.expect("activate"); 1379 core.refresh_active_profile().await.expect("refresh"); 1380 1381 wait_for_fresh_profile(&observer).await; 1382 let snapshots = observer.snapshots.lock().expect("snapshots").clone(); 1383 assert!(snapshots.iter().any(|snapshot| { 1384 snapshot.active_identity.as_ref().is_some_and(|active| { 1385 active.profile_state == ProfileLoadStateDto::Fresh 1386 && active 1387 .profile 1388 .as_ref() 1389 .and_then(|profile| profile.display_name.as_deref()) 1390 == Some("FFI Profile") 1391 }) 1392 })); 1393 subscription.unsubscribe().await; 1394 let count = observer.snapshots.lock().expect("snapshots").len(); 1395 core.sign_out().await.expect("sign out"); 1396 assert_eq!(observer.snapshots.lock().expect("snapshots").len(), count); 1397 1398 core.shutdown_v2().await.expect("shutdown"); 1399 publisher.shutdown().await; 1400 local_relay.shutdown(); 1401 } 1402 1403 struct ArcObserver(Arc<RecordingObserver>); 1404 1405 impl HarvestCircleChangeObserver for ArcObserver { 1406 fn on_change(&self, change: SnapshotChangeDto) { 1407 self.0.on_change(change); 1408 } 1409 } 1410 1411 const OBSERVER_DELIVERY_TIMEOUT: Duration = Duration::from_secs(5); 1412 1413 async fn wait_for_snapshot_count(observer: &RecordingObserver, minimum: usize) { 1414 tokio::time::timeout(OBSERVER_DELIVERY_TIMEOUT, async { 1415 while observer.snapshots.lock().expect("snapshots").len() < minimum { 1416 tokio::task::yield_now().await; 1417 } 1418 }) 1419 .await 1420 .expect("snapshot delivery"); 1421 } 1422 1423 async fn wait_for_observer_count(core: &HarvestCircleAppCore, expected: usize) { 1424 tokio::time::timeout(OBSERVER_DELIVERY_TIMEOUT, async { 1425 while core.inner.observers.lock().expect("observers").len() != expected { 1426 tokio::task::yield_now().await; 1427 } 1428 }) 1429 .await 1430 .expect("observer deregistration"); 1431 } 1432 1433 async fn wait_for_fresh_profile(observer: &RecordingObserver) { 1434 tokio::time::timeout(OBSERVER_DELIVERY_TIMEOUT, async { 1435 loop { 1436 let fresh = observer 1437 .snapshots 1438 .lock() 1439 .expect("snapshots") 1440 .iter() 1441 .any(|snapshot| { 1442 snapshot.active_identity.as_ref().is_some_and(|active| { 1443 active.profile_state == ProfileLoadStateDto::Fresh 1444 }) 1445 }); 1446 if fresh { 1447 break; 1448 } 1449 tokio::task::yield_now().await; 1450 } 1451 }) 1452 .await 1453 .expect("fresh profile delivery"); 1454 } 1455 }