app

Local-first trade for farms and co-ops
git clone https://radroots.dev/git/app.git
Log | Files | Refs | README | LICENSE

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 }