lib

Core libraries for Radroots
git clone https://radroots.dev/git/lib.git
Log | Files | Refs | README

commit f8fe9c96528cbbb2ed29a0ed619f19eb3c6abced
parent 53cc1d73cdb51ae268ecc7224d398c2ebf31d569
Author: triesap <tyson@radroots.org>
Date:   Mon,  3 Aug 2026 21:57:18 +0000

runtime: make shutdown terminal and bounded

- terminate actor admission after a successful close receipt
- cancel supervised work and close every ordered change receiver
- prevent expired queued shutdown requests from taking effect later
- cover task cancellation, publication closure, and shutdown deadlines

Diffstat:
Mcrates/studio_application/src/change_stream.rs | 30+++++++++++++++++++++++++++++-
Mcrates/studio_storage/src/runtime_actor.rs | 125+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++----
2 files changed, 148 insertions(+), 7 deletions(-)

diff --git a/crates/studio_application/src/change_stream.rs b/crates/studio_application/src/change_stream.rs @@ -63,6 +63,7 @@ pub struct OrderedSnapshotChanges { latest: AppSnapshot, next_subscription: u64, subscribers: BTreeMap<ChangeSubscriptionId, mpsc::Sender<SnapshotChange>>, + closed: bool, } impl OrderedSnapshotChanges { @@ -72,6 +73,7 @@ impl OrderedSnapshotChanges { latest: initial_snapshot, next_subscription: 1, subscribers: BTreeMap::new(), + closed: false, } } @@ -89,6 +91,9 @@ impl OrderedSnapshotChanges { &mut self, capacity: NonZeroUsize, ) -> Option<(ChangeSubscriptionId, SnapshotChangeReceiver)> { + if self.closed { + return None; + } let id = ChangeSubscriptionId(NonZeroU64::new(self.next_subscription)?); self.next_subscription = self.next_subscription.checked_add(1)?; let (sender, receiver) = mpsc::channel(capacity.get()); @@ -108,7 +113,7 @@ impl OrderedSnapshotChanges { } pub fn publish(&mut self, snapshot: AppSnapshot) { - if snapshot.revision() <= self.latest.revision() { + if self.closed || snapshot.revision() <= self.latest.revision() { return; } let change = SnapshotChange { @@ -122,6 +127,11 @@ impl OrderedSnapshotChanges { Err(mpsc::error::TrySendError::Closed(_)) => false, }); } + + pub fn close(&mut self) { + self.closed = true; + self.subscribers.clear(); + } } #[cfg(test)] @@ -188,6 +198,24 @@ mod tests { assert!(recovered.recovers_gap_after(first.revision())); } + #[tokio::test] + async fn close_terminates_consumers_and_rejects_later_subscriptions() { + let mut changes = OrderedSnapshotChanges::new(snapshot(0)); + let (_, mut receiver) = changes + .subscribe(NonZeroUsize::new(1).expect("capacity")) + .expect("subscription"); + receiver.receive().await.expect("initial"); + + changes.close(); + changes.publish(snapshot(1)); + assert!(receiver.receive().await.is_none()); + assert!( + changes + .subscribe(NonZeroUsize::new(1).expect("capacity")) + .is_none() + ); + } + fn revision(value: u64) -> SnapshotRevision { SnapshotRevision::from_value(value) } diff --git a/crates/studio_storage/src/runtime_actor.rs b/crates/studio_storage/src/runtime_actor.rs @@ -366,7 +366,27 @@ impl RuntimeActorHandle { /// /// Returns a safe timeout or actor error. Repeated calls return closed. pub async fn close(&self) -> Result<(), SafeError> { - match self.dispatch(RuntimeCommand::Close, None).await? { + self.close_with_timeout(DEFAULT_COMMAND_TIMEOUT).await + } + + /// Closes the runtime within the supplied command deadline. + /// + /// # Errors + /// + /// Returns a safe timeout or actor error. An expired queued close cannot + /// later change runtime state. + pub async fn close_with_timeout(&self, timeout: Duration) -> Result<(), SafeError> { + let raw_request = self.next_request.fetch_add(1, Ordering::Relaxed); + let request_id = RequestId::new(raw_request).ok_or_else(request_space_exhausted)?; + match self + .dispatch_with_deadline( + RuntimeCommand::Close, + None, + request_id, + Instant::now() + timeout, + ) + .await? + { RuntimeCommandValue::Closed => Ok(()), _ => Err(invalid_actor_response()), } @@ -491,7 +511,9 @@ impl RuntimeActor { let Some(envelope) = envelope else { break; }; - self.handle_command(envelope, &completion_sender); + if !self.handle_command(envelope, &completion_sender) { + break; + } } completion = completions.recv(), if !self.profile_tasks.is_empty() => { if let Some(completion) = completion { @@ -507,20 +529,21 @@ impl RuntimeActor { &mut self, envelope: CommandEnvelope<RuntimeCommand, RuntimeCommandValue>, completion_sender: &mpsc::Sender<ProfileCompletion>, - ) { + ) -> bool { let (context, command, reply) = envelope.into_parts(); if let Some(result) = self.preflight(context, &command) { let _ = reply.send(CommandReceipt::new(context.request_id(), result)); - return; + return true; } if matches!(command, RuntimeCommand::RefreshActiveProfile) { self.start_profile_task(context, reply, completion_sender.clone()); - return; + return true; } if matches!(command, RuntimeCommand::Close) { let result = self.close_actor(); + let closed = matches!(result, CommandResult::Completed(_)); let _ = reply.send(CommandReceipt::new(context.request_id(), result)); - return; + return !closed; } let changes_session = matches!( command, @@ -536,6 +559,7 @@ impl RuntimeActor { self.changes.publish(self.adapter.core().snapshot()); } let _ = reply.send(CommandReceipt::new(context.request_id(), result)); + true } fn preflight( @@ -685,6 +709,7 @@ impl RuntimeActor { match transition { Ok(()) => { self.cancel_profile_tasks(None); + self.changes.close(); CommandResult::Completed(RuntimeCommandValue::Closed) } Err(error) => CommandResult::Failed(error), @@ -1175,6 +1200,94 @@ mod tests { ); } + #[tokio::test(flavor = "multi_thread", worker_threads = 4)] + async fn expired_queued_shutdown_does_not_close_runtime_later() { + let secrets = Arc::new(BlockingSecretStore::new()); + let actor = RuntimeActorHandle::in_memory( + RelayConfiguration::default(), + secrets.clone(), + Arc::new(FixedClock), + Arc::new(OfflineNostr), + NonZeroUsize::new(1).expect("capacity"), + &tokio::runtime::Handle::current(), + ) + .expect("actor"); + let import_actor = actor.clone(); + let import = tokio::spawn(async move { + import_actor + .import_secret_key(secret( + "7e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7", + )) + .await + }); + secrets.wait_until_put_started().await; + + let timeout = actor + .close_with_timeout(Duration::from_millis(10)) + .await + .expect_err("queued shutdown must expire"); + assert_eq!(timeout.message().as_str(), "The runtime command timed out."); + secrets.release(); + import.await.expect("import task").expect("import"); + assert_eq!(actor.lifecycle(), RuntimeLifecycle::Ready); + assert_eq!( + actor + .bootstrap() + .await + .expect("still open") + .accounts() + .len(), + 1 + ); + actor.close().await.expect("later close"); + } + + #[tokio::test(flavor = "multi_thread", worker_threads = 4)] + async fn shutdown_cancels_in_flight_work_and_terminates_publication() { + let client = Arc::new(BlockingNostr::new()); + let actor = RuntimeActorHandle::in_memory( + RelayConfiguration::new(vec![RelayUrl::parse("ws://localhost:8080").expect("relay")]), + Arc::new(InMemorySecretStore::default()), + Arc::new(FixedClock), + client.clone(), + NonZeroUsize::new(8).expect("capacity"), + &tokio::runtime::Handle::current(), + ) + .expect("actor"); + let imported = actor + .import_secret_key(secret( + "7e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7", + )) + .await + .expect("import"); + actor + .activate_account(imported.account().public_key()) + .await + .expect("activate"); + let mut changes = actor + .subscribe_changes(NonZeroUsize::new(4).expect("capacity")) + .await + .expect("subscribe"); + changes.receive().await.expect("initial"); + + let refresh_actor = actor.clone(); + let refresh = tokio::spawn(async move { refresh_actor.refresh_active_profile().await }); + let started = client.started.acquire().await.expect("refresh started"); + started.forget(); + actor.close().await.expect("close"); + + let cancelled = refresh + .await + .expect("refresh task") + .expect_err("refresh closes"); + assert_eq!( + cancelled.message().as_str(), + "The application runtime is closed." + ); + assert!(changes.receive().await.is_none()); + assert_eq!(actor.lifecycle(), RuntimeLifecycle::Closed); + } + fn secret(value: &str) -> SecretKeyInput { SecretKeyInput::parse(value.to_owned()).expect("valid test secret") }