lib

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

commit f6542e571bb50e1f13c295dd6324743c52d5f94b
parent ac81b246ded9c0459eb50809fe948b9fd1b83acf
Author: triesap <tyson@radroots.org>
Date:   Mon,  3 Aug 2026 21:37:16 +0000

runtime: serialize account mutation through actor

- add one storage-backed runtime actor that owns the persistent application core
- route generation, import, selection, activation, sign-out, and removal commands
- replace direct FFI mutation access with bounded mailbox dispatch
- preserve read-only snapshot subscriptions and cover the full account lifecycle

Diffstat:
Mcrates/studio_application/src/actor.rs | 9++++++++-
Mcrates/studio_application/src/lib.rs | 2+-
Mcrates/studio_ffi/src/commands.rs | 206+++++++++++++++++++++++++++++++++++--------------------------------------------
Mcrates/studio_ffi/src/observer.rs | 47++++++++++++++++++++++++++---------------------
Mcrates/studio_storage/Cargo.toml | 1+
Mcrates/studio_storage/src/lib.rs | 2++
Acrates/studio_storage/src/runtime_actor.rs | 557+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
7 files changed, 686 insertions(+), 138 deletions(-)

diff --git a/crates/studio_application/src/actor.rs b/crates/studio_application/src/actor.rs @@ -372,11 +372,18 @@ impl<C, R> CommandEnvelope<C, R> { } } -#[derive(Clone)] pub struct ActorMailbox<C, R> { sender: mpsc::Sender<CommandEnvelope<C, R>>, } +impl<C, R> Clone for ActorMailbox<C, R> { + fn clone(&self) -> Self { + Self { + sender: self.sender.clone(), + } + } +} + impl<C, R> ActorMailbox<C, R> { #[must_use] pub fn bounded(capacity: NonZeroUsize) -> (Self, mpsc::Receiver<CommandEnvelope<C, R>>) { diff --git a/crates/studio_application/src/lib.rs b/crates/studio_application/src/lib.rs @@ -18,7 +18,7 @@ pub use accounts::{ InMemoryOperationJournal, }; pub use actor::{ - ActorMailbox, CommandContext, CommandReceipt, CommandRejection, CommandResult, + ActorMailbox, CommandContext, CommandEnvelope, CommandReceipt, CommandRejection, CommandResult, CommandSubmission, CommandTicket, LifecycleGate, RequestId, RuntimeCommandClass, RuntimeLifecycle, }; diff --git a/crates/studio_ffi/src/commands.rs b/crates/studio_ffi/src/commands.rs @@ -1,5 +1,6 @@ use std::collections::BTreeSet; use std::fmt::{self, Display, Formatter}; +use std::num::NonZeroUsize; use std::path::{Path, PathBuf}; use std::sync::atomic::AtomicBool; use std::sync::{Arc, Mutex, OnceLock}; @@ -7,11 +8,11 @@ use std::time::{Duration, SystemTime, UNIX_EPOCH}; use directories::ProjectDirs; use radroots_studio_application::{ - Clock, RelayRuntimeMode, RemovalConfirmationToken, SdkNostrClient, SecretStore, + Clock, RelayRuntimeMode, RemovalConfirmationToken, SdkNostrClient, relay_configuration_from_environment, }; use radroots_studio_domain::{PublicKey, SafeError, SecretKeyInput, UnixTimestamp}; -use radroots_studio_storage::{OsKeyringSecretStore, PersistentAppCore}; +use radroots_studio_storage::{OsKeyringSecretStore, RuntimeActorHandle}; use crate::{AccountDto, AppSnapshotDto}; @@ -19,6 +20,7 @@ const DATABASE_QUALIFIER: &str = "org"; const DATABASE_ORGANIZATION: &str = "radroots"; const DATABASE_APPLICATION: &str = "studio"; const DATABASE_FILENAME: &str = "studio.sqlite3"; +pub(crate) const ACTOR_MAILBOX_CAPACITY: usize = 64; #[derive(Debug, uniffi::Error)] pub enum StudioError { @@ -65,10 +67,7 @@ impl RemovalRequest { } pub(crate) struct RuntimeCore { - pub(crate) adapter: PersistentAppCore, - pub(crate) secrets: Arc<dyn SecretStore>, - pub(crate) clock: SystemClock, - pub(crate) nostr: SdkNostrClient, + pub(crate) actor: RuntimeActorHandle, pub(crate) observers: Mutex<BTreeSet<radroots_studio_application::ObserverHandle>>, pub(crate) closed: AtomicBool, } @@ -99,19 +98,17 @@ impl StudioAppCore { /// /// Returns a safe storage, recovery, or application-state error. pub async fn bootstrap(&self) -> Result<AppSnapshotDto, StudioError> { - let inner = Arc::clone(&self.inner); - blocking(move || { - inner - .adapter - .bootstrap(inner.secrets.as_ref(), &inner.clock) - .map(|snapshot| (&snapshot).into()) - }) - .await + self.inner + .actor + .bootstrap() + .await + .map(|snapshot| (&snapshot).into()) + .map_err(StudioError::from) } #[must_use] pub fn snapshot(&self) -> AppSnapshotDto { - (&self.inner.adapter.core().snapshot()).into() + (&self.inner.actor.snapshot()).into() } /// Generates and stores one local account with a one-time backup receipt. @@ -120,18 +117,16 @@ impl StudioAppCore { /// /// Returns a safe keyring, storage, or account error. pub async fn generate_account(&self) -> Result<GeneratedAccountDto, StudioError> { - let inner = Arc::clone(&self.inner); - blocking(move || { - let receipt = inner - .adapter - .generate_account(inner.secrets.as_ref(), &inner.clock)?; - Ok(GeneratedAccountDto { + self.inner + .actor + .generate_account() + .await + .map(|receipt| GeneratedAccountDto { account: receipt.account().into(), - snapshot: (&inner.adapter.core().snapshot()).into(), + snapshot: (&self.inner.actor.snapshot()).into(), nsec: receipt.generated_nsec().with_exposed_secret(str::to_owned), }) - }) - .await + .map_err(StudioError::from) } /// Imports one nsec or canonical secret-key hex value. @@ -144,14 +139,12 @@ impl StudioAppCore { secret_key: String, ) -> Result<AppSnapshotDto, StudioError> { let input = SecretKeyInput::parse(secret_key).map_err(StudioError::from)?; - let inner = Arc::clone(&self.inner); - blocking(move || { - inner - .adapter - .import_secret_key(input, inner.secrets.as_ref(), &inner.clock)?; - Ok((&inner.adapter.core().snapshot()).into()) - }) - .await + self.inner + .actor + .import_secret_key(input) + .await + .map(|_| (&self.inner.actor.snapshot()).into()) + .map_err(StudioError::from) } /// Selects one saved account without activating it. @@ -164,14 +157,12 @@ impl StudioAppCore { public_key_hex: String, ) -> Result<AppSnapshotDto, StudioError> { let public_key = parse_public_key(&public_key_hex)?; - let inner = Arc::clone(&self.inner); - blocking(move || { - inner - .adapter - .select_account(public_key) - .map(|snapshot| (&snapshot).into()) - }) - .await + self.inner + .actor + .select_account(public_key) + .await + .map(|snapshot| (&snapshot).into()) + .map_err(StudioError::from) } /// Activates one saved account after validating its credential. @@ -184,14 +175,12 @@ impl StudioAppCore { public_key_hex: String, ) -> Result<AppSnapshotDto, StudioError> { let public_key = parse_public_key(&public_key_hex)?; - let inner = Arc::clone(&self.inner); - blocking(move || { - inner - .adapter - .activate_account(public_key, inner.secrets.as_ref(), &inner.clock) - .map(|snapshot| (&snapshot).into()) - }) - .await + self.inner + .actor + .activate_account(public_key) + .await + .map(|snapshot| (&snapshot).into()) + .map_err(StudioError::from) } /// Signs out while retaining accounts and credentials. @@ -200,8 +189,12 @@ impl StudioAppCore { /// /// Returns a safe application-state error. pub async fn sign_out(&self) -> Result<AppSnapshotDto, StudioError> { - let inner = Arc::clone(&self.inner); - blocking(move || inner.adapter.sign_out().map(|snapshot| (&snapshot).into())).await + self.inner + .actor + .sign_out() + .await + .map(|snapshot| (&snapshot).into()) + .map_err(StudioError::from) } /// Refreshes the active Nostr profile from configured relays. @@ -210,19 +203,12 @@ impl StudioAppCore { /// /// Returns a safe storage or application-state error. pub async fn refresh_active_profile(&self) -> Result<AppSnapshotDto, StudioError> { - let inner = Arc::clone(&self.inner); - runtime() - .spawn(async move { - inner - .adapter - .core() - .refresh_active_profile(inner.adapter.database(), &inner.nostr, &inner.clock) - .await - .map(|snapshot| (&snapshot).into()) - .map_err(StudioError::from) - }) + self.inner + .actor + .refresh_active_profile() .await - .map_err(|_| runtime_unavailable())? + .map(|snapshot| (&snapshot).into()) + .map_err(StudioError::from) } /// Issues a revision-bound removal confirmation object. @@ -235,15 +221,17 @@ impl StudioAppCore { public_key_hex: String, ) -> Result<Arc<RemovalRequest>, StudioError> { let public_key = parse_public_key(&public_key_hex)?; - let inner = Arc::clone(&self.inner); - blocking(move || { - let token = inner.adapter.request_account_removal(public_key)?; - Ok(Arc::new(RemovalRequest { - public_key_hex, - token: Mutex::new(Some(token)), - })) - }) - .await + self.inner + .actor + .request_account_removal(public_key) + .await + .map(|token| { + Arc::new(RemovalRequest { + public_key_hex, + token: Mutex::new(Some(token)), + }) + }) + .map_err(StudioError::from) } /// Permanently removes the account represented by a one-time request. @@ -261,14 +249,12 @@ impl StudioAppCore { .unwrap_or_else(std::sync::PoisonError::into_inner) .take() .ok_or_else(confirmation_expired)?; - let inner = Arc::clone(&self.inner); - blocking(move || { - inner - .adapter - .confirm_account_removal(token, inner.secrets.as_ref(), &inner.clock) - .map(|snapshot| (&snapshot).into()) - }) - .await + self.inner + .actor + .confirm_account_removal(token) + .await + .map(|snapshot| (&snapshot).into()) + .map_err(StudioError::from) } } @@ -280,13 +266,18 @@ impl StudioAppCore { RelayRuntimeMode::Packaged }; let relays = relay_configuration_from_environment(mode)?; - let adapter = PersistentAppCore::open(path, relays)?; + let actor = RuntimeActorHandle::open( + path, + relays, + Arc::new(OsKeyringSecretStore), + Arc::new(SystemClock), + Arc::new(SdkNostrClient::new(Duration::from_secs(5))), + NonZeroUsize::new(ACTOR_MAILBOX_CAPACITY).expect("nonzero actor mailbox capacity"), + runtime().handle(), + )?; Ok(Arc::new(Self { inner: Arc::new(RuntimeCore { - adapter, - secrets: Arc::new(OsKeyringSecretStore), - clock: SystemClock, - nostr: SdkNostrClient::new(Duration::from_secs(5)), + actor, observers: Mutex::new(BTreeSet::new()), closed: AtomicBool::new(false), }), @@ -322,19 +313,7 @@ fn parse_public_key(value: &str) -> Result<PublicKey, StudioError> { PublicKey::from_hex(value).map_err(StudioError::from) } -async fn blocking<T, F>(operation: F) -> Result<T, StudioError> -where - T: Send + 'static, - F: FnOnce() -> Result<T, SafeError> + Send + 'static, -{ - runtime() - .spawn_blocking(operation) - .await - .map_err(|_| runtime_unavailable())? - .map_err(StudioError::from) -} - -fn runtime() -> &'static tokio::runtime::Runtime { +pub(crate) fn runtime() -> &'static tokio::runtime::Runtime { static RUNTIME: OnceLock<tokio::runtime::Runtime> = OnceLock::new(); RUNTIME.get_or_init(|| { tokio::runtime::Builder::new_multi_thread() @@ -352,13 +331,6 @@ fn path_unavailable() -> StudioError { } } -fn runtime_unavailable() -> StudioError { - StudioError::Failure { - code: "InvalidApplicationState".to_owned(), - safe_message: "The application runtime is unavailable.".to_owned(), - } -} - fn confirmation_expired() -> StudioError { StudioError::Failure { code: "InvalidApplicationState".to_owned(), @@ -368,28 +340,32 @@ fn confirmation_expired() -> StudioError { #[cfg(test)] mod tests { + use std::num::NonZeroUsize; use std::sync::Arc; - use radroots_studio_application::RelayConfiguration; - use radroots_studio_storage::PersistentAppCore; + use radroots_studio_application::{InMemorySecretStore, RelayConfiguration, SdkNostrClient}; + use radroots_studio_storage::RuntimeActorHandle; use radroots_studio_storage::{CREDENTIAL_SERVICE, CURRENT_SCHEMA_VERSION}; use super::{ - DATABASE_APPLICATION, DATABASE_FILENAME, DATABASE_ORGANIZATION, DATABASE_QUALIFIER, - RuntimeCore, StudioAppCore, SystemClock, + ACTOR_MAILBOX_CAPACITY, DATABASE_APPLICATION, DATABASE_FILENAME, DATABASE_ORGANIZATION, + DATABASE_QUALIFIER, RuntimeCore, StudioAppCore, SystemClock, runtime, }; fn in_memory_core() -> Arc<StudioAppCore> { + let actor = RuntimeActorHandle::in_memory( + RelayConfiguration::default(), + Arc::new(InMemorySecretStore::default()), + Arc::new(SystemClock), + Arc::new(SdkNostrClient::new(std::time::Duration::from_millis(10))), + NonZeroUsize::new(ACTOR_MAILBOX_CAPACITY).expect("capacity"), + runtime().handle(), + ) + .expect("in-memory actor"); Arc::new(StudioAppCore { inner: Arc::new(RuntimeCore { - adapter: PersistentAppCore::in_memory(RelayConfiguration::default()) - .expect("in-memory core"), - secrets: Arc::new(radroots_studio_application::InMemorySecretStore::default()), - clock: SystemClock, - nostr: radroots_studio_application::SdkNostrClient::new( - std::time::Duration::from_millis(10), - ), + actor, observers: std::sync::Mutex::new(std::collections::BTreeSet::new()), closed: std::sync::atomic::AtomicBool::new(false), }), diff --git a/crates/studio_ffi/src/observer.rs b/crates/studio_ffi/src/observer.rs @@ -38,7 +38,7 @@ impl ObserverSubscription { let (Some(core), Some(handle)) = (self.core.upgrade(), handle) else { return; }; - let _ = core.adapter.core().unsubscribe(handle); + let _ = core.actor.unsubscribe(handle); core.observers .lock() .unwrap_or_else(std::sync::PoisonError::into_inner) @@ -71,8 +71,7 @@ impl StudioAppCore { }); let handle = self .inner - .adapter - .core() + .actor .subscribe(bridge) .map_err(StudioError::from)?; self.inner @@ -98,9 +97,12 @@ impl StudioAppCore { .unwrap_or_else(std::sync::PoisonError::into_inner), ); for handle in handles { - let _ = self.inner.adapter.core().unsubscribe(handle); + let _ = self.inner.actor.unsubscribe(handle); } - let _ = self.inner.adapter.sign_out(); + let actor = self.inner.actor.clone(); + crate::commands::runtime().spawn(async move { + let _ = actor.sign_out().await; + }); } } @@ -113,17 +115,18 @@ fn closed_error() -> StudioError { #[cfg(test)] mod tests { + use std::num::NonZeroUsize; use std::sync::{Arc, Mutex}; use std::time::Duration; use nostr::{EventBuilder, Keys, Metadata}; use nostr_relay_builder::MockRelay; use nostr_sdk::Client; - use radroots_studio_application::RelayConfiguration; + use radroots_studio_application::{InMemorySecretStore, RelayConfiguration, SdkNostrClient}; use radroots_studio_domain::RelayUrl; - use radroots_studio_storage::PersistentAppCore; + use radroots_studio_storage::RuntimeActorHandle; - use crate::commands::{RuntimeCore, SystemClock}; + use crate::commands::{ACTOR_MAILBOX_CAPACITY, RuntimeCore, SystemClock, runtime}; use crate::{AppSnapshotDto, ProfileLoadStateDto, StudioAppCore, StudioObserver}; const SECRET_HEX: &str = "7e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7"; @@ -148,14 +151,18 @@ mod tests { } fn core_with_relays(relays: RelayConfiguration) -> Arc<StudioAppCore> { + let actor = RuntimeActorHandle::in_memory( + relays, + Arc::new(InMemorySecretStore::default()), + Arc::new(SystemClock), + Arc::new(SdkNostrClient::new(std::time::Duration::from_millis(10))), + NonZeroUsize::new(ACTOR_MAILBOX_CAPACITY).expect("capacity"), + runtime().handle(), + ) + .expect("actor"); Arc::new(StudioAppCore { inner: Arc::new(RuntimeCore { - adapter: PersistentAppCore::in_memory(relays).expect("core"), - secrets: Arc::new(radroots_studio_application::InMemorySecretStore::default()), - clock: SystemClock, - nostr: radroots_studio_application::SdkNostrClient::new( - std::time::Duration::from_millis(10), - ), + actor, observers: Mutex::new(std::collections::BTreeSet::new()), closed: std::sync::atomic::AtomicBool::new(false), }), @@ -165,7 +172,6 @@ mod tests { #[test] fn callbacks_allow_reentry_and_stop_after_subscription_close() { let core = core(); - core.inner.adapter.core().bootstrap().expect("bootstrap"); let observer = Arc::new(RecordingObserver::default()); *observer.core.lock().expect("core") = Some(Arc::clone(&core)); let subscription = core @@ -173,21 +179,20 @@ mod tests { .expect("subscribe"); assert_eq!(observer.snapshots.lock().expect("snapshots").len(), 1); - core.inner - .adapter - .core() - .bootstrap() + runtime() + .block_on(core.inner.actor.bootstrap()) .expect("idempotent bootstrap"); assert_eq!(observer.snapshots.lock().expect("snapshots").len(), 1); subscription.unsubscribe(); - core.inner.adapter.core().sign_out().expect("sign out"); + runtime() + .block_on(core.inner.actor.sign_out()) + .expect("sign out"); assert_eq!(observer.snapshots.lock().expect("snapshots").len(), 1); } #[test] fn core_close_deregisters_all_observers_and_rejects_new_subscriptions() { let core = core(); - core.inner.adapter.core().bootstrap().expect("bootstrap"); let observer = Arc::new(RecordingObserver::default()); let _subscription = core .subscribe(Box::new(ArcObserver(observer.clone()))) diff --git a/crates/studio_storage/Cargo.toml b/crates/studio_storage/Cargo.toml @@ -12,6 +12,7 @@ radroots-studio-domain = { path = "../domain" } keyring.workspace = true refinery.workspace = true rusqlite.workspace = true +tokio.workspace = true [dev-dependencies] nostr.workspace = true diff --git a/crates/studio_storage/src/lib.rs b/crates/studio_storage/src/lib.rs @@ -7,7 +7,9 @@ pub mod db; pub mod journal; pub mod os_keyring; pub mod profiles; +pub mod runtime_actor; pub use application_adapter::PersistentAppCore; pub use db::{CURRENT_SCHEMA_VERSION, Database}; pub use os_keyring::{CREDENTIAL_SERVICE, OsKeyringSecretStore}; +pub use runtime_actor::RuntimeActorHandle; diff --git a/crates/studio_storage/src/runtime_actor.rs b/crates/studio_storage/src/runtime_actor.rs @@ -0,0 +1,557 @@ +use std::num::NonZeroUsize; +use std::path::Path; +use std::sync::atomic::{AtomicU64, Ordering}; +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, +}; +use radroots_studio_domain::{PublicKey, SafeError, SafeErrorCode, SafeMessage, SecretKeyInput}; +use tokio::runtime::Handle; + +use crate::PersistentAppCore; + +const DEFAULT_COMMAND_TIMEOUT: Duration = Duration::from_secs(30); + +enum RuntimeCommand { + Snapshot, + GenerateAccount, + ImportSecretKey(SecretKeyInput), + SelectAccount(PublicKey), + ActivateAccount(PublicKey), + SignOut, + RefreshActiveProfile, + RequestAccountRemoval(PublicKey), + ConfirmAccountRemoval(RemovalConfirmationToken), +} + +enum RuntimeCommandValue { + Snapshot(Box<AppSnapshot>), + Generated(GenerateAccountReceipt), + Imported(ImportAccountReceipt), + RemovalRequest(RemovalConfirmationToken), +} + +impl RuntimeCommand { + const fn class(&self) -> RuntimeCommandClass { + match self { + Self::Snapshot => RuntimeCommandClass::Observe, + Self::GenerateAccount + | Self::ImportSecretKey(_) + | Self::ActivateAccount(_) + | Self::ConfirmAccountRemoval(_) => RuntimeCommandClass::UseCredential, + Self::SelectAccount(_) | Self::SignOut | Self::RequestAccountRemoval(_) => { + RuntimeCommandClass::MutateLocalState + } + Self::RefreshActiveProfile => RuntimeCommandClass::UseRelay, + } + } +} + +struct RuntimeActor { + adapter: Arc<PersistentAppCore>, + secrets: Arc<dyn SecretStore>, + clock: Arc<dyn Clock>, + nostr: Arc<dyn NostrClient>, + lifecycle: Arc<Mutex<LifecycleGate>>, + runtime: Handle, +} + +#[derive(Clone)] +pub struct RuntimeActorHandle { + mailbox: ActorMailbox<RuntimeCommand, RuntimeCommandValue>, + adapter: Arc<PersistentAppCore>, + lifecycle: Arc<Mutex<LifecycleGate>>, + next_request: Arc<AtomicU64>, +} + +impl RuntimeActorHandle { + /// Opens, migrates, recovers, and starts one actor-owned file-backed runtime. + /// + /// # Errors + /// + /// Returns a safe storage, recovery, or lifecycle error before the actor is + /// published when opening cannot reach ready state. + pub fn open( + path: &Path, + relay_configuration: RelayConfiguration, + secrets: Arc<dyn SecretStore>, + clock: Arc<dyn Clock>, + nostr: Arc<dyn NostrClient>, + capacity: NonZeroUsize, + runtime: &Handle, + ) -> Result<Self, SafeError> { + Self::start( + PersistentAppCore::open(path, relay_configuration)?, + secrets, + clock, + nostr, + capacity, + runtime, + ) + } + + /// Starts one isolated actor-owned in-memory runtime for tests. + /// + /// # Errors + /// + /// Returns a safe storage, recovery, or lifecycle error before publication. + pub fn in_memory( + relay_configuration: RelayConfiguration, + secrets: Arc<dyn SecretStore>, + clock: Arc<dyn Clock>, + nostr: Arc<dyn NostrClient>, + capacity: NonZeroUsize, + runtime: &Handle, + ) -> Result<Self, SafeError> { + Self::start( + PersistentAppCore::in_memory(relay_configuration)?, + secrets, + clock, + nostr, + capacity, + runtime, + ) + } + + fn start( + adapter: PersistentAppCore, + secrets: Arc<dyn SecretStore>, + clock: Arc<dyn Clock>, + nostr: Arc<dyn NostrClient>, + capacity: NonZeroUsize, + runtime: &Handle, + ) -> Result<Self, SafeError> { + let mut gate = LifecycleGate::opening(); + gate.begin_compatibility_check()?; + gate.compatibility_accepted()?; + gate.ownership_acquired()?; + gate.migration_complete()?; + adapter.bootstrap(secrets.as_ref(), clock.as_ref())?; + gate.recovery_complete()?; + + let adapter = Arc::new(adapter); + let lifecycle = Arc::new(Mutex::new(gate)); + let (mailbox, receiver) = ActorMailbox::bounded(capacity); + let actor = RuntimeActor { + adapter: Arc::clone(&adapter), + secrets, + clock, + nostr, + lifecycle: Arc::clone(&lifecycle), + runtime: runtime.clone(), + }; + runtime.spawn_blocking(move || actor.run(receiver)); + Ok(Self { + mailbox, + adapter, + lifecycle, + next_request: Arc::new(AtomicU64::new(1)), + }) + } + + #[must_use] + pub fn lifecycle(&self) -> RuntimeLifecycle { + self.lifecycle + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .lifecycle() + } + + #[must_use] + pub fn snapshot(&self) -> AppSnapshot { + self.adapter.core().snapshot() + } + + /// Registers a read-only snapshot observer. + /// + /// # Errors + /// + /// Returns a safe observer-registration error. + pub fn subscribe(&self, observer: Arc<dyn AppObserver>) -> Result<ObserverHandle, SafeError> { + self.adapter.core().subscribe(observer) + } + + #[must_use] + pub fn unsubscribe(&self, handle: ObserverHandle) -> bool { + self.adapter.core().unsubscribe(handle) + } + + /// Returns the ready snapshot through the actor command boundary. + /// + /// # Errors + /// + /// Returns a typed safe actor error. + pub async fn bootstrap(&self) -> Result<AppSnapshot, SafeError> { + Self::expect_snapshot(self.dispatch(RuntimeCommand::Snapshot, None).await?) + } + + /// Generates one account through the serialized actor boundary. + /// + /// # Errors + /// + /// Returns a safe account, storage, keyring, timeout, or actor error. + pub async fn generate_account(&self) -> Result<GenerateAccountReceipt, SafeError> { + match self.dispatch(RuntimeCommand::GenerateAccount, None).await? { + RuntimeCommandValue::Generated(receipt) => Ok(receipt), + _ => Err(invalid_actor_response()), + } + } + + /// Imports one account through the serialized actor boundary. + /// + /// # Errors + /// + /// Returns a safe account, storage, keyring, timeout, or actor error. + pub async fn import_secret_key( + &self, + input: SecretKeyInput, + ) -> Result<ImportAccountReceipt, SafeError> { + match self + .dispatch(RuntimeCommand::ImportSecretKey(input), None) + .await? + { + RuntimeCommandValue::Imported(receipt) => Ok(receipt), + _ => Err(invalid_actor_response()), + } + } + + /// Selects one account through the serialized actor boundary. + /// + /// # Errors + /// + /// Returns a safe account, storage, timeout, or actor error. + pub async fn select_account(&self, public_key: PublicKey) -> Result<AppSnapshot, SafeError> { + let value = self + .dispatch(RuntimeCommand::SelectAccount(public_key), None) + .await?; + Self::expect_snapshot(value) + } + + /// Activates one account through the serialized actor boundary. + /// + /// # Errors + /// + /// Returns a safe account, credential, storage, timeout, or actor error. + pub async fn activate_account(&self, public_key: PublicKey) -> Result<AppSnapshot, SafeError> { + let value = self + .dispatch(RuntimeCommand::ActivateAccount(public_key), None) + .await?; + Self::expect_snapshot(value) + } + + /// Signs out through the serialized actor boundary. + /// + /// # Errors + /// + /// Returns a safe timeout or actor error. + pub async fn sign_out(&self) -> Result<AppSnapshot, SafeError> { + let value = self.dispatch(RuntimeCommand::SignOut, None).await?; + Self::expect_snapshot(value) + } + + /// Refreshes the active profile through the serialized actor boundary. + /// + /// # Errors + /// + /// Returns a safe relay, storage, timeout, or actor error. + pub async fn refresh_active_profile(&self) -> Result<AppSnapshot, SafeError> { + let value = self + .dispatch(RuntimeCommand::RefreshActiveProfile, None) + .await?; + Self::expect_snapshot(value) + } + + /// Creates one removal request through the serialized actor boundary. + /// + /// # Errors + /// + /// Returns a safe account, timeout, or actor error. + pub async fn request_account_removal( + &self, + public_key: PublicKey, + ) -> Result<RemovalConfirmationToken, SafeError> { + match self + .dispatch(RuntimeCommand::RequestAccountRemoval(public_key), None) + .await? + { + RuntimeCommandValue::RemovalRequest(token) => Ok(token), + _ => Err(invalid_actor_response()), + } + } + + /// Confirms one removal through the serialized actor boundary. + /// + /// # Errors + /// + /// Returns a safe account, credential, storage, timeout, or actor error. + pub async fn confirm_account_removal( + &self, + token: RemovalConfirmationToken, + ) -> Result<AppSnapshot, SafeError> { + let value = self + .dispatch(RuntimeCommand::ConfirmAccountRemoval(token), None) + .await?; + Self::expect_snapshot(value) + } + + async fn dispatch( + &self, + command: RuntimeCommand, + expected_revision: Option<SnapshotRevision>, + ) -> Result<RuntimeCommandValue, 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)?; + let context = CommandContext::new( + request_id, + expected_revision, + Instant::now() + DEFAULT_COMMAND_TIMEOUT, + ); + let receipt = match self.mailbox.submit(context, command) { + CommandSubmission::Accepted(ticket) => ticket.receipt().await, + CommandSubmission::Rejected(receipt) => receipt, + }; + match receipt.into_result() { + CommandResult::Completed(value) => Ok(value), + CommandResult::Conflicted { .. } => Err(command_conflicted()), + CommandResult::Rejected(_) => Err(command_rejected()), + CommandResult::TimedOut => Err(command_timed_out()), + CommandResult::Closed => Err(runtime_closed()), + CommandResult::Failed(error) => Err(error), + } + } + + fn expect_snapshot(value: RuntimeCommandValue) -> Result<AppSnapshot, SafeError> { + match value { + RuntimeCommandValue::Snapshot(snapshot) => Ok(*snapshot), + _ => Err(invalid_actor_response()), + } + } +} + +impl RuntimeActor { + fn run( + self, + mut receiver: tokio::sync::mpsc::Receiver< + radroots_studio_application::CommandEnvelope<RuntimeCommand, RuntimeCommandValue>, + >, + ) { + while let Some(envelope) = receiver.blocking_recv() { + let (context, command, reply) = envelope.into_parts(); + let result = self.execute(context, command); + let _ = reply.send(CommandReceipt::new(context.request_id(), result)); + } + } + + fn execute( + &self, + context: CommandContext, + command: RuntimeCommand, + ) -> CommandResult<RuntimeCommandValue> { + if context.is_expired(Instant::now()) { + return CommandResult::TimedOut; + } + let lifecycle = self + .lifecycle + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .to_owned(); + if matches!(lifecycle.lifecycle(), RuntimeLifecycle::Closed) { + return CommandResult::Closed; + } + if !lifecycle.allows(command.class()) { + return 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 }; + } + let result = match command { + RuntimeCommand::Snapshot => Ok(RuntimeCommandValue::Snapshot(Box::new( + self.adapter.core().snapshot(), + ))), + RuntimeCommand::GenerateAccount => self + .adapter + .generate_account(self.secrets.as_ref(), self.clock.as_ref()) + .map(RuntimeCommandValue::Generated), + RuntimeCommand::ImportSecretKey(input) => self + .adapter + .import_secret_key(input, self.secrets.as_ref(), self.clock.as_ref()) + .map(RuntimeCommandValue::Imported), + RuntimeCommand::SelectAccount(public_key) => self + .adapter + .select_account(public_key) + .map(Box::new) + .map(RuntimeCommandValue::Snapshot), + RuntimeCommand::ActivateAccount(public_key) => self + .adapter + .activate_account(public_key, self.secrets.as_ref(), self.clock.as_ref()) + .map(Box::new) + .map(RuntimeCommandValue::Snapshot), + RuntimeCommand::SignOut => self + .adapter + .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::RequestAccountRemoval(public_key) => self + .adapter + .request_account_removal(public_key) + .map(RuntimeCommandValue::RemovalRequest), + RuntimeCommand::ConfirmAccountRemoval(token) => self + .adapter + .confirm_account_removal(token, self.secrets.as_ref(), self.clock.as_ref()) + .map(Box::new) + .map(RuntimeCommandValue::Snapshot), + }; + result.map_or_else(CommandResult::Failed, CommandResult::Completed) + } +} + +const fn request_space_exhausted() -> SafeError { + SafeError::new( + SafeErrorCode::InvalidApplicationState, + SafeMessage::new("The runtime request identifier space is exhausted."), + ) +} + +const fn command_conflicted() -> SafeError { + SafeError::new( + SafeErrorCode::InvalidApplicationState, + SafeMessage::new("The command conflicts with newer application state."), + ) +} + +const fn command_rejected() -> SafeError { + SafeError::new( + SafeErrorCode::InvalidApplicationState, + SafeMessage::new("The runtime is busy. Try again."), + ) +} + +const fn command_timed_out() -> SafeError { + SafeError::new( + SafeErrorCode::InvalidApplicationState, + SafeMessage::new("The runtime command timed out."), + ) +} + +const fn runtime_closed() -> SafeError { + SafeError::new( + SafeErrorCode::InvalidApplicationState, + SafeMessage::new("The application runtime is closed."), + ) +} + +const fn command_unavailable() -> SafeError { + SafeError::new( + SafeErrorCode::InvalidApplicationState, + SafeMessage::new("The command is unavailable in the current runtime state."), + ) +} + +const fn invalid_actor_response() -> SafeError { + SafeError::new( + SafeErrorCode::InvalidApplicationState, + SafeMessage::new("The runtime returned an invalid command response."), + ) +} + +#[cfg(test)] +mod tests { + use std::num::NonZeroUsize; + use std::sync::Arc; + + use radroots_studio_application::{ + BoxFuture, Clock, InMemorySecretStore, NostrClient, RelayConfiguration, RuntimeLifecycle, + SecretStore, SessionState, + }; + use radroots_studio_domain::{ + Kind0ProfileCandidate, PublicKey, RelayUrl, SafeError, SecretKeyInput, UnixTimestamp, + }; + + use super::RuntimeActorHandle; + + struct FixedClock; + + impl Clock for FixedClock { + fn now(&self) -> UnixTimestamp { + UnixTimestamp::from_seconds(50).expect("time") + } + } + + struct OfflineNostr; + + impl NostrClient for OfflineNostr { + fn fetch_profile<'a>( + &'a self, + _public_key: PublicKey, + _relays: &'a [RelayUrl], + ) -> BoxFuture<'a, Result<Option<Kind0ProfileCandidate>, SafeError>> { + Box::pin(async { Ok(None) }) + } + } + + fn actor() -> (RuntimeActorHandle, Arc<InMemorySecretStore>) { + let secrets = Arc::new(InMemorySecretStore::default()); + let secret_port: Arc<dyn SecretStore> = secrets.clone(); + let actor = RuntimeActorHandle::in_memory( + RelayConfiguration::default(), + secret_port, + Arc::new(FixedClock), + Arc::new(OfflineNostr), + NonZeroUsize::new(8).expect("capacity"), + &tokio::runtime::Handle::current(), + ) + .expect("actor"); + (actor, secrets) + } + + #[tokio::test(flavor = "multi_thread")] + async fn account_mutations_run_serially_through_one_ready_actor() { + let (actor, secrets) = actor(); + assert_eq!(actor.lifecycle(), RuntimeLifecycle::Ready); + + let imported = actor + .import_secret_key( + SecretKeyInput::parse( + "7e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7".to_owned(), + ) + .expect("input"), + ) + .await + .expect("import"); + let public_key = imported.account().public_key(); + let activated = actor.activate_account(public_key).await.expect("activate"); + assert_eq!(activated.session(), SessionState::Active); + assert!(secrets.contains(public_key).expect("credential")); + + let signed_out = actor.sign_out().await.expect("sign out"); + assert_eq!(signed_out.session(), SessionState::SignedOut); + let removal = actor + .request_account_removal(public_key) + .await + .expect("removal request"); + let removed = actor + .confirm_account_removal(removal) + .await + .expect("remove"); + assert!(removed.accounts().is_empty()); + assert!(!secrets.contains(public_key).expect("credential removed")); + } +}