app

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

commit 721f75a4354de820412d0001cf84fde26f6d682d
parent afc8650e20ed785e19836203873c29ae27f5b36e
Author: triesap <tyson@radroots.org>
Date:   Mon,  3 Aug 2026 21:43:27 +0000

runtime: supervise correlated profile tasks

- split profile refresh into actor-owned prepare and completion phases
- correlate tasks by request, account, revision, and session generation
- cancel pending work when activation, sign-out, or removal replaces a session
- prove cancelled refreshes cannot publish after the foreground session changes

Diffstat:
Mcore/Cargo.toml | 2+-
Mcore/crates/application/src/actor.rs | 75++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-
Mcore/crates/application/src/lib.rs | 3++-
Mcore/crates/application/src/profile_refresh.rs | 98++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++---------------
Mcore/crates/storage/src/runtime_actor.rs | 309+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++--------
5 files changed, 436 insertions(+), 51 deletions(-)

diff --git a/core/Cargo.toml b/core/Cargo.toml @@ -34,7 +34,7 @@ rusqlite = { version = "=0.39.0", features = ["bundled"] } secrecy = "=0.10.3" url = "=2.5.8" zeroize = "=1.9.0" -tokio = { version = "=1.47.1", features = ["rt-multi-thread", "sync"] } +tokio = { version = "=1.47.1", features = ["macros", "rt-multi-thread", "sync"] } uniffi = "=0.32.0" [patch.crates-io] diff --git a/core/crates/application/src/actor.rs b/core/crates/application/src/actor.rs @@ -1,11 +1,84 @@ use std::num::{NonZeroU64, NonZeroUsize}; use std::time::Instant; -use radroots_studio_domain::SafeError; +use radroots_studio_domain::{PublicKey, SafeError}; use tokio::sync::{mpsc, oneshot}; use crate::SnapshotRevision; +#[derive(Clone, Copy, Debug, Default, Eq, Hash, Ord, PartialEq, PartialOrd)] +pub struct SessionGeneration(u64); + +impl SessionGeneration { + #[must_use] + pub const fn initial() -> Self { + Self(0) + } + + #[must_use] + pub const fn from_value(value: u64) -> Self { + Self(value) + } + + #[must_use] + pub const fn value(self) -> u64 { + self.0 + } + + #[must_use] + pub const fn next(self) -> Option<Self> { + match self.0.checked_add(1) { + Some(value) => Some(Self(value)), + None => None, + } + } +} + +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub struct TaskCorrelation { + request_id: RequestId, + account: PublicKey, + expected_revision: SnapshotRevision, + session_generation: SessionGeneration, +} + +impl TaskCorrelation { + #[must_use] + pub const fn new( + request_id: RequestId, + account: PublicKey, + expected_revision: SnapshotRevision, + session_generation: SessionGeneration, + ) -> Self { + Self { + request_id, + account, + expected_revision, + session_generation, + } + } + + #[must_use] + pub const fn request_id(self) -> RequestId { + self.request_id + } + + #[must_use] + pub const fn account(self) -> PublicKey { + self.account + } + + #[must_use] + pub const fn expected_revision(self) -> SnapshotRevision { + self.expected_revision + } + + #[must_use] + pub const fn session_generation(self) -> SessionGeneration { + self.session_generation + } +} + #[derive(Clone, Copy, Debug, Eq, PartialEq)] pub enum RuntimeLifecycle { Opening, diff --git a/core/crates/application/src/lib.rs b/core/crates/application/src/lib.rs @@ -20,7 +20,7 @@ pub use accounts::{ pub use actor::{ ActorMailbox, CommandContext, CommandEnvelope, CommandReceipt, CommandRejection, CommandResult, CommandSubmission, CommandTicket, LifecycleGate, RequestId, RuntimeCommandClass, - RuntimeLifecycle, + RuntimeLifecycle, SessionGeneration, TaskCorrelation, }; pub use app_core::{AppCore, AppObserver, ObserverHandle, RemovalConfirmationToken}; pub use config::{ @@ -33,6 +33,7 @@ pub use ports::{ OperationDiagnostic, OperationId, OperationJournal, PendingAccountOperation, ProfileRefreshStatus, ProfileRepository, }; +pub use profile_refresh::ProfileRefreshPlan; pub use secrets::{ FailureSecretStore, InMemorySecretStore, SecretStore, SecretStoreCall, SecretStoreOperation, }; diff --git a/core/crates/application/src/profile_refresh.rs b/core/crates/application/src/profile_refresh.rs @@ -1,11 +1,41 @@ -use radroots_studio_domain::{PublicKey, SafeError, SafeErrorCode}; +use radroots_studio_domain::{PublicKey, RelayUrl, SafeError, SafeErrorCode}; use crate::{ ActiveAccountSnapshot, AppCore, AppSnapshot, CachedProfile, Clock, NostrClient, ProfileLoadState, ProfileRefreshStatus, ProfileRepository, RelayConnectionState, - StateTransition, + SnapshotRevision, StateTransition, }; +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct ProfileRefreshPlan { + public_key: PublicKey, + active_account: ActiveAccountSnapshot, + relays: Vec<RelayUrl>, + expected_revision: SnapshotRevision, +} + +impl ProfileRefreshPlan { + #[must_use] + pub const fn public_key(&self) -> PublicKey { + self.public_key + } + + #[must_use] + pub const fn active_account(&self) -> &ActiveAccountSnapshot { + &self.active_account + } + + #[must_use] + pub fn relays(&self) -> &[RelayUrl] { + &self.relays + } + + #[must_use] + pub const fn expected_revision(&self) -> SnapshotRevision { + self.expected_revision + } +} + impl AppCore { /// Manually refreshes the active account's Nostr kind-0 profile. /// @@ -40,11 +70,24 @@ impl AppCore { client: &(impl NostrClient + ?Sized), clock: &(impl Clock + ?Sized), ) -> Result<AppSnapshot, SafeError> { - let Some(active) = self.snapshot().active_account().cloned() else { + let Some(plan) = self.begin_profile_refresh()? else { return Ok(self.snapshot()); }; + let result = client.fetch_profile(plan.public_key(), plan.relays()).await; + self.complete_profile_refresh(&plan, result, profiles, clock) + } + + /// Begins a refresh on the actor and returns the immutable network plan. + /// + /// # Errors + /// + /// Returns a safe state error when the loading transition is invalid. + pub fn begin_profile_refresh(&self) -> Result<Option<ProfileRefreshPlan>, SafeError> { + let Some(active) = self.snapshot().active_account().cloned() else { + return Ok(None); + }; let public_key = active.account().public_key(); - self.apply_transition(StateTransition::UpdateActiveAccount { + let loading = self.apply_transition(StateTransition::UpdateActiveAccount { expected: public_key, active_account: Box::new(ActiveAccountSnapshot::new( active.account().clone(), @@ -54,11 +97,30 @@ impl AppCore { )), problem: None, })?; + Ok(Some(ProfileRefreshPlan { + public_key, + active_account: active, + relays: loading.relay_configuration().relays().to_vec(), + expected_revision: loading.revision(), + })) + } - let result = client - .fetch_profile(public_key, self.snapshot().relay_configuration().relays()) - .await; - if !is_current_active(self, public_key) { + /// Applies a correlated refresh result on the actor. + /// + /// # Errors + /// + /// Returns a safe storage or application-state error. Stale results are + /// discarded without persistence or publication. + pub fn complete_profile_refresh( + &self, + plan: &ProfileRefreshPlan, + result: Result<Option<radroots_studio_domain::Kind0ProfileCandidate>, SafeError>, + profiles: &(impl ProfileRepository + ?Sized), + clock: &(impl Clock + ?Sized), + ) -> Result<AppSnapshot, SafeError> { + if !is_current_active(self, plan.public_key()) + || self.snapshot().revision() != plan.expected_revision() + { return Ok(self.snapshot()); } @@ -71,9 +133,9 @@ impl AppCore { ); profiles.save_profile(&cached)?; self.apply_transition(StateTransition::UpdateActiveAccount { - expected: public_key, + expected: plan.public_key(), active_account: Box::new(ActiveAccountSnapshot::new( - active.account().clone(), + plan.active_account().account().clone(), RelayConnectionState::Connected, ProfileLoadState::Fresh, Some(candidate.metadata().clone()), @@ -82,29 +144,29 @@ impl AppCore { }) } Ok(None) => self.apply_transition(StateTransition::UpdateActiveAccount { - expected: public_key, + expected: plan.public_key(), active_account: Box::new(ActiveAccountSnapshot::new( - active.account().clone(), + plan.active_account().account().clone(), RelayConnectionState::Connected, - if active.profile().is_some() { + if plan.active_account().profile().is_some() { ProfileLoadState::Cached } else { ProfileLoadState::Empty }, - active.profile().cloned(), + plan.active_account().profile().cloned(), )), problem: None, }), Err(error) => { let status = refresh_status(error); - profiles.record_refresh_status(public_key, clock.now(), status)?; + profiles.record_refresh_status(plan.public_key(), clock.now(), status)?; self.apply_transition(StateTransition::UpdateActiveAccount { - expected: public_key, + expected: plan.public_key(), active_account: Box::new(ActiveAccountSnapshot::new( - active.account().clone(), + plan.active_account().account().clone(), RelayConnectionState::Degraded, ProfileLoadState::Error(error), - active.profile().cloned(), + plan.active_account().profile().cloned(), )), problem: Some(error), }) diff --git a/core/crates/storage/src/runtime_actor.rs b/core/crates/storage/src/runtime_actor.rs @@ -1,3 +1,4 @@ +use std::collections::BTreeMap; use std::num::NonZeroUsize; use std::path::Path; use std::sync::atomic::{AtomicU64, Ordering}; @@ -5,17 +6,22 @@ use std::sync::{Arc, Mutex}; use std::time::{Duration, Instant}; use radroots_studio_application::{ - ActorMailbox, AppObserver, AppSnapshot, Clock, CommandContext, CommandReceipt, CommandResult, - CommandSubmission, GenerateAccountReceipt, ImportAccountReceipt, LifecycleGate, NostrClient, - ObserverHandle, RelayConfiguration, RemovalConfirmationToken, RequestId, RuntimeCommandClass, - RuntimeLifecycle, SecretStore, SnapshotRevision, + ActorMailbox, AppObserver, AppSnapshot, Clock, CommandContext, CommandEnvelope, CommandReceipt, + CommandResult, CommandSubmission, GenerateAccountReceipt, ImportAccountReceipt, LifecycleGate, + NostrClient, ObserverHandle, ProfileRefreshPlan, RelayConfiguration, RemovalConfirmationToken, + RequestId, RuntimeCommandClass, RuntimeLifecycle, SecretStore, SessionGeneration, + SnapshotRevision, TaskCorrelation, +}; +use radroots_studio_domain::{ + Kind0ProfileCandidate, PublicKey, SafeError, SafeErrorCode, SafeMessage, SecretKeyInput, }; -use radroots_studio_domain::{PublicKey, SafeError, SafeErrorCode, SafeMessage, SecretKeyInput}; use tokio::runtime::Handle; +use tokio::sync::{mpsc, oneshot}; use crate::PersistentAppCore; const DEFAULT_COMMAND_TIMEOUT: Duration = Duration::from_secs(30); +const DEFAULT_TASK_CAPACITY: usize = 64; enum RuntimeCommand { Snapshot, @@ -59,6 +65,21 @@ struct RuntimeActor { nostr: Arc<dyn NostrClient>, lifecycle: Arc<Mutex<LifecycleGate>>, runtime: Handle, + session_generation: SessionGeneration, + published_session_generation: Arc<AtomicU64>, + profile_tasks: BTreeMap<RequestId, PendingProfileTask>, +} + +struct PendingProfileTask { + correlation: TaskCorrelation, + plan: ProfileRefreshPlan, + reply: oneshot::Sender<CommandReceipt<RuntimeCommandValue>>, + handle: tokio::task::JoinHandle<()>, +} + +struct ProfileCompletion { + request_id: RequestId, + result: Result<Option<Kind0ProfileCandidate>, SafeError>, } #[derive(Clone)] @@ -67,6 +88,7 @@ pub struct RuntimeActorHandle { adapter: Arc<PersistentAppCore>, lifecycle: Arc<Mutex<LifecycleGate>>, next_request: Arc<AtomicU64>, + session_generation: Arc<AtomicU64>, } impl RuntimeActorHandle { @@ -137,6 +159,7 @@ impl RuntimeActorHandle { let adapter = Arc::new(adapter); let lifecycle = Arc::new(Mutex::new(gate)); let (mailbox, receiver) = ActorMailbox::bounded(capacity); + let session_generation = Arc::new(AtomicU64::new(SessionGeneration::initial().value())); let actor = RuntimeActor { adapter: Arc::clone(&adapter), secrets, @@ -144,13 +167,17 @@ impl RuntimeActorHandle { nostr, lifecycle: Arc::clone(&lifecycle), runtime: runtime.clone(), + session_generation: SessionGeneration::initial(), + published_session_generation: Arc::clone(&session_generation), + profile_tasks: BTreeMap::new(), }; - runtime.spawn_blocking(move || actor.run(receiver)); + drop(runtime.spawn(actor.run(receiver))); Ok(Self { mailbox, adapter, lifecycle, next_request: Arc::new(AtomicU64::new(1)), + session_generation, }) } @@ -163,6 +190,11 @@ impl RuntimeActorHandle { } #[must_use] + pub fn session_generation(&self) -> SessionGeneration { + SessionGeneration::from_value(self.session_generation.load(Ordering::Acquire)) + } + + #[must_use] pub fn snapshot(&self) -> AppSnapshot { self.adapter.core().snapshot() } @@ -334,26 +366,63 @@ impl RuntimeActorHandle { } impl RuntimeActor { - fn run( - self, - mut receiver: tokio::sync::mpsc::Receiver< - radroots_studio_application::CommandEnvelope<RuntimeCommand, RuntimeCommandValue>, - >, + async fn run( + mut self, + mut receiver: mpsc::Receiver<CommandEnvelope<RuntimeCommand, RuntimeCommandValue>>, + ) { + let (completion_sender, mut completions) = mpsc::channel(DEFAULT_TASK_CAPACITY); + loop { + tokio::select! { + envelope = receiver.recv() => { + let Some(envelope) = envelope else { + break; + }; + self.handle_command(envelope, &completion_sender); + } + completion = completions.recv(), if !self.profile_tasks.is_empty() => { + if let Some(completion) = completion { + self.complete_profile_task(completion); + } + } + } + } + self.cancel_profile_tasks(None); + } + + fn handle_command( + &mut self, + envelope: CommandEnvelope<RuntimeCommand, RuntimeCommandValue>, + completion_sender: &mpsc::Sender<ProfileCompletion>, ) { - while let Some(envelope) = receiver.blocking_recv() { - let (context, command, reply) = envelope.into_parts(); - let result = self.execute(context, command); + 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; + } + if matches!(command, RuntimeCommand::RefreshActiveProfile) { + self.start_profile_task(context, reply, completion_sender.clone()); + return; } + let changes_session = matches!( + command, + RuntimeCommand::ActivateAccount(_) + | RuntimeCommand::SignOut + | RuntimeCommand::ConfirmAccountRemoval(_) + ); + let result = self.execute_sync(command); + if changes_session && matches!(result, CommandResult::Completed(_)) { + self.advance_session_generation(); + } + let _ = reply.send(CommandReceipt::new(context.request_id(), result)); } - fn execute( + fn preflight( &self, context: CommandContext, - command: RuntimeCommand, - ) -> CommandResult<RuntimeCommandValue> { + command: &RuntimeCommand, + ) -> Option<CommandResult<RuntimeCommandValue>> { if context.is_expired(Instant::now()) { - return CommandResult::TimedOut; + return Some(CommandResult::TimedOut); } let lifecycle = self .lifecycle @@ -361,18 +430,22 @@ impl RuntimeActor { .unwrap_or_else(std::sync::PoisonError::into_inner) .to_owned(); if matches!(lifecycle.lifecycle(), RuntimeLifecycle::Closed) { - return CommandResult::Closed; + return Some(CommandResult::Closed); } if !lifecycle.allows(command.class()) { - return CommandResult::Failed(command_unavailable()); + return Some(CommandResult::Failed(command_unavailable())); } let current_revision = self.adapter.core().snapshot().revision(); if context .expected_revision() .is_some_and(|expected| expected != current_revision) { - return CommandResult::Conflicted { current_revision }; + return Some(CommandResult::Conflicted { current_revision }); } + None + } + + fn execute_sync(&self, command: RuntimeCommand) -> CommandResult<RuntimeCommandValue> { let result = match command { RuntimeCommand::Snapshot => Ok(RuntimeCommandValue::Snapshot(Box::new( self.adapter.core().snapshot(), @@ -400,15 +473,7 @@ impl RuntimeActor { .sign_out() .map(Box::new) .map(RuntimeCommandValue::Snapshot), - RuntimeCommand::RefreshActiveProfile => self - .runtime - .block_on(self.adapter.core().refresh_active_profile( - self.adapter.database(), - self.nostr.as_ref(), - self.clock.as_ref(), - )) - .map(Box::new) - .map(RuntimeCommandValue::Snapshot), + RuntimeCommand::RefreshActiveProfile => Err(invalid_actor_response()), RuntimeCommand::RequestAccountRemoval(public_key) => self .adapter .request_account_removal(public_key) @@ -421,6 +486,118 @@ impl RuntimeActor { }; result.map_or_else(CommandResult::Failed, CommandResult::Completed) } + + fn start_profile_task( + &mut self, + context: CommandContext, + reply: oneshot::Sender<CommandReceipt<RuntimeCommandValue>>, + completion_sender: mpsc::Sender<ProfileCompletion>, + ) { + let plan = match self.adapter.core().begin_profile_refresh() { + Ok(Some(plan)) => plan, + Ok(None) => { + let _ = reply.send(CommandReceipt::new( + context.request_id(), + CommandResult::Completed(RuntimeCommandValue::Snapshot(Box::new( + self.adapter.core().snapshot(), + ))), + )); + return; + } + Err(error) => { + let _ = reply.send(CommandReceipt::new( + context.request_id(), + CommandResult::Failed(error), + )); + return; + } + }; + let correlation = TaskCorrelation::new( + context.request_id(), + plan.public_key(), + plan.expected_revision(), + self.session_generation, + ); + let client = Arc::clone(&self.nostr); + let relays = plan.relays().to_vec(); + let request_id = context.request_id(); + let handle = self.runtime.spawn(async move { + let result = client.fetch_profile(correlation.account(), &relays).await; + let _ = completion_sender + .send(ProfileCompletion { request_id, result }) + .await; + }); + let previous = self.profile_tasks.insert( + request_id, + PendingProfileTask { + correlation, + plan, + reply, + handle, + }, + ); + debug_assert!(previous.is_none(), "request identifiers are unique"); + } + + fn complete_profile_task(&mut self, completion: ProfileCompletion) { + let Some(task) = self.profile_tasks.remove(&completion.request_id) else { + return; + }; + let current = self.adapter.core().snapshot(); + let correlated = task.correlation.session_generation() == self.session_generation + && task.correlation.expected_revision() == current.revision() + && current + .active_account() + .is_some_and(|active| active.account().public_key() == task.correlation.account()); + let result = if correlated { + self.adapter + .core() + .complete_profile_refresh( + &task.plan, + completion.result, + self.adapter.database(), + self.clock.as_ref(), + ) + .map(Box::new) + .map(RuntimeCommandValue::Snapshot) + .map_or_else(CommandResult::Failed, CommandResult::Completed) + } else { + CommandResult::Completed(RuntimeCommandValue::Snapshot(Box::new(current))) + }; + let _ = task + .reply + .send(CommandReceipt::new(task.correlation.request_id(), result)); + } + + fn advance_session_generation(&mut self) { + let Some(next) = self.session_generation.next() else { + self.lifecycle + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .fail(request_space_exhausted()); + self.cancel_profile_tasks(None); + return; + }; + self.session_generation = next; + self.published_session_generation + .store(next.value(), Ordering::Release); + let snapshot = self.adapter.core().snapshot(); + self.cancel_profile_tasks(Some(&snapshot)); + } + + fn cancel_profile_tasks(&mut self, snapshot: Option<&AppSnapshot>) { + let tasks = std::mem::take(&mut self.profile_tasks); + for (_, task) in tasks { + task.handle.abort(); + let receipt_result = snapshot.map_or(CommandResult::Closed, |snapshot| { + CommandResult::Completed(RuntimeCommandValue::Snapshot(Box::new(snapshot.clone()))) + }); + let _ = task.reply.send(CommandReceipt::new( + task.correlation.request_id(), + receipt_result, + )); + } + } } const fn request_space_exhausted() -> SafeError { @@ -507,6 +684,35 @@ mod tests { } } + struct BlockingNostr { + started: tokio::sync::Semaphore, + release: tokio::sync::Semaphore, + } + + impl BlockingNostr { + fn new() -> Self { + Self { + started: tokio::sync::Semaphore::new(0), + release: tokio::sync::Semaphore::new(0), + } + } + } + + impl NostrClient for BlockingNostr { + fn fetch_profile<'a>( + &'a self, + _public_key: PublicKey, + _relays: &'a [RelayUrl], + ) -> BoxFuture<'a, Result<Option<Kind0ProfileCandidate>, SafeError>> { + Box::pin(async move { + self.started.add_permits(1); + let permit = self.release.acquire().await.expect("release"); + permit.forget(); + Ok(None) + }) + } + } + fn actor() -> (RuntimeActorHandle, Arc<InMemorySecretStore>) { let secrets = Arc::new(InMemorySecretStore::default()); let secret_port: Arc<dyn SecretStore> = secrets.clone(); @@ -554,4 +760,47 @@ mod tests { assert!(removed.accounts().is_empty()); assert!(!secrets.contains(public_key).expect("credential removed")); } + + #[tokio::test(flavor = "multi_thread")] + async fn session_generation_cancels_correlated_profile_work_on_sign_out() { + 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( + SecretKeyInput::parse( + "7e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7".to_owned(), + ) + .expect("input"), + ) + .await + .expect("import"); + actor + .activate_account(imported.account().public_key()) + .await + .expect("activate"); + assert_eq!(actor.session_generation().value(), 1); + + 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(); + let signed_out = actor.sign_out().await.expect("sign out"); + let cancelled = refresh + .await + .expect("refresh task") + .expect("safe cancellation"); + + assert_eq!(actor.session_generation().value(), 2); + assert_eq!(signed_out.session(), SessionState::SignedOut); + assert_eq!(cancelled.session(), SessionState::SignedOut); + assert!(cancelled.active_account().is_none()); + } }