app

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

commit 97bb016bcf29555bb890a2e3f1effb98890f5502
parent a4d7deebec3e2ce2c1daa455de6d79857839aed0
Author: triesap <tyson@radroots.org>
Date:   Sun,  9 Aug 2026 17:43:18 +0000

application: own the Studio application crate

- copy the exact locked Studio application source into HarvestCircle core
- register application policy as a local workspace member
- bind application code to the locally owned domain crate
- refresh the lockfile after the source-boundary change

Diffstat:
Mcore/Cargo.lock | 19++++++++++++++-----
Mcore/Cargo.toml | 3++-
Acore/crates/studio_application/Cargo.toml | 20++++++++++++++++++++
Acore/crates/studio_application/src/accounts.rs | 1822+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acore/crates/studio_application/src/actor.rs | 787+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acore/crates/studio_application/src/app_core.rs | 358+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acore/crates/studio_application/src/change_stream.rs | 239+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acore/crates/studio_application/src/config.rs | 107+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acore/crates/studio_application/src/custody.rs | 288+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acore/crates/studio_application/src/lib.rs | 58++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acore/crates/studio_application/src/ports.rs | 887+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acore/crates/studio_application/src/profile_refresh.rs | 659+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acore/crates/studio_application/src/recovery.rs | 788+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acore/crates/studio_application/src/secrets.rs | 308+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acore/crates/studio_application/src/session.rs | 268+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acore/crates/studio_application/src/snapshot.rs | 479+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acore/crates/studio_application/src/state_machine.rs | 616+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acore/crates/studio_application/src/test_support.rs | 58++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acore/crates/studio_application/tests/redaction.rs | 48++++++++++++++++++++++++++++++++++++++++++++++++
19 files changed, 7806 insertions(+), 6 deletions(-)

diff --git a/core/Cargo.lock b/core/Cargo.lock @@ -1943,6 +1943,15 @@ dependencies = [ [[package]] name = "radroots_studio_application" version = "0.1.0-alpha" +dependencies = [ + "radroots_studio_domain 0.1.0-alpha", + "secrecy", + "tokio", +] + +[[package]] +name = "radroots_studio_application" +version = "0.1.0-alpha" source = "git+https://github.com/radrootslabs/lib?rev=09065a610d95e57acdc895a14c07580fa099e7c3#09065a610d95e57acdc895a14c07580fa099e7c3" dependencies = [ "radroots_studio_domain 0.1.0-alpha (git+https://github.com/radrootslabs/lib?rev=09065a610d95e57acdc895a14c07580fa099e7c3)", @@ -1980,7 +1989,7 @@ source = "git+https://github.com/radrootslabs/lib?rev=09065a610d95e57acdc895a14c dependencies = [ "directories", "quote", - "radroots_studio_application", + "radroots_studio_application 0.1.0-alpha (git+https://github.com/radrootslabs/lib?rev=09065a610d95e57acdc895a14c07580fa099e7c3)", "radroots_studio_domain 0.1.0-alpha (git+https://github.com/radrootslabs/lib?rev=09065a610d95e57acdc895a14c07580fa099e7c3)", "radroots_studio_nostr", "radroots_studio_runtime", @@ -1999,7 +2008,7 @@ dependencies = [ "nostr 0.44.1", "nostr-sdk 0.44.0", "radroots_identity", - "radroots_studio_application", + "radroots_studio_application 0.1.0-alpha (git+https://github.com/radrootslabs/lib?rev=09065a610d95e57acdc895a14c07580fa099e7c3)", "radroots_studio_domain 0.1.0-alpha (git+https://github.com/radrootslabs/lib?rev=09065a610d95e57acdc895a14c07580fa099e7c3)", "radroots_transport", "radroots_transport_nostr", @@ -2018,7 +2027,7 @@ name = "radroots_studio_runtime" version = "0.1.0-alpha" source = "git+https://github.com/radrootslabs/lib?rev=09065a610d95e57acdc895a14c07580fa099e7c3#09065a610d95e57acdc895a14c07580fa099e7c3" dependencies = [ - "radroots_studio_application", + "radroots_studio_application 0.1.0-alpha (git+https://github.com/radrootslabs/lib?rev=09065a610d95e57acdc895a14c07580fa099e7c3)", "radroots_studio_domain 0.1.0-alpha (git+https://github.com/radrootslabs/lib?rev=09065a610d95e57acdc895a14c07580fa099e7c3)", "radroots_studio_nostr", "radroots_studio_storage", @@ -2030,7 +2039,7 @@ dependencies = [ name = "radroots_studio_source_lock" version = "0.1.0-alpha" dependencies = [ - "radroots_studio_application", + "radroots_studio_application 0.1.0-alpha", "radroots_studio_domain 0.1.0-alpha", "radroots_studio_ffi", "radroots_studio_nostr", @@ -2048,7 +2057,7 @@ dependencies = [ "getrandom 0.2.17", "hmac", "keyring", - "radroots_studio_application", + "radroots_studio_application 0.1.0-alpha (git+https://github.com/radrootslabs/lib?rev=09065a610d95e57acdc895a14c07580fa099e7c3)", "radroots_studio_domain 0.1.0-alpha (git+https://github.com/radrootslabs/lib?rev=09065a610d95e57acdc895a14c07580fa099e7c3)", "refinery", "rusqlite", diff --git a/core/Cargo.toml b/core/Cargo.toml @@ -1,6 +1,7 @@ [workspace] members = [ "crates/source_lock", + "crates/studio_application", "crates/studio_domain", "crates/studio_preferences", ] @@ -23,7 +24,7 @@ all = "deny" pedantic = "deny" [workspace.dependencies] -radroots_studio_application = { git = "https://github.com/radrootslabs/lib", rev = "09065a610d95e57acdc895a14c07580fa099e7c3", version = "=0.1.0-alpha" } +radroots_studio_application = { path = "crates/studio_application", version = "=0.1.0-alpha" } radroots_studio_domain = { path = "crates/studio_domain", version = "=0.1.0-alpha" } radroots_studio_ffi = { git = "https://github.com/radrootslabs/lib", rev = "09065a610d95e57acdc895a14c07580fa099e7c3", version = "=0.1.0-alpha" } radroots_studio_nostr = { git = "https://github.com/radrootslabs/lib", rev = "09065a610d95e57acdc895a14c07580fa099e7c3", version = "=0.1.0-alpha" } diff --git a/core/crates/studio_application/Cargo.toml b/core/crates/studio_application/Cargo.toml @@ -0,0 +1,20 @@ +[package] +name = "radroots_studio_application" +description = "Private application policy and ports for Radroots Studio" +version = "0.1.0-alpha" +edition.workspace = true +authors.workspace = true +rust-version.workspace = true +license = "GPL-3.0-only" +repository.workspace = true +homepage.workspace = true +publish = false +include = ["src/**", "tests/**", "Cargo.toml"] + +[dependencies] +radroots_studio_domain.workspace = true +secrecy = "=0.10.3" +tokio = { version = "=1.47.1", features = ["macros", "rt-multi-thread", "sync", "time"] } + +[lints] +workspace = true diff --git a/core/crates/studio_application/src/accounts.rs b/core/crates/studio_application/src/accounts.rs @@ -0,0 +1,1822 @@ +use std::sync::{Mutex, MutexGuard}; + +use crate::{ + AccountOperationKind, AccountOperationPhase, AccountRepository, AppCore, AppStateRepository, + Clock, DurableOperationKind, DurableOperationPhase, DurableOperationRepository, + DurableOperationStart, DurableRequestId, DurableTerminalOutcome, OperationDiagnostic, + OperationId, OperationJournal, OperationPriorState, PendingAccountOperation, + RemovalConfirmationToken, SecretStore, StagedGeneratedKey, StateTransition, +}; +use radroots_studio_domain::{ + AccountCreatedAt, AccountIdentity, AccountSummary, BindingAvailability, LocalSignerBinding, + Nsec, PublicKey, SafeError, SafeErrorCode, SafeMessage, SecretKeyInput, +}; + +pub struct GenerateAccountReceipt { + account: AccountSummary, + generated_nsec: Nsec, +} + +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct ImportAccountReceipt { + account: AccountSummary, +} + +impl ImportAccountReceipt { + #[must_use] + pub const fn account(&self) -> &AccountSummary { + &self.account + } +} + +impl GenerateAccountReceipt { + #[must_use] + pub const fn account(&self) -> &AccountSummary { + &self.account + } + + #[must_use] + pub const fn generated_nsec(&self) -> &Nsec { + &self.generated_nsec + } +} + +impl AppCore { + /// Commits a staged generated key only after its recovery acknowledgement. + /// + /// # Errors + /// + /// Returns a safe conflict, keyring, persistence, or recovery error. + #[allow(clippy::too_many_arguments)] + pub fn commit_staged_generated_key( + &self, + request_id: &DurableRequestId, + staged: StagedGeneratedKey, + accounts: &(impl AccountRepository + ?Sized), + app_state: &(impl AppStateRepository + ?Sized), + secrets: &(impl SecretStore + ?Sized), + operations: &(impl DurableOperationRepository + ?Sized), + clock: &(impl Clock + ?Sized), + ) -> Result<ImportAccountReceipt, SafeError> { + let expected_revision = staged.expected_revision(); + self.require_revision(expected_revision)?; + let (account, secret) = staged.into_commit_parts(); + self.persist_account_durable( + request_id, + DurableOperationKind::Create, + expected_revision, + &account, + secret, + None, + accounts, + app_state, + secrets, + operations, + clock, + )?; + Ok(ImportAccountReceipt { account }) + } + + /// Generates and commits one account under a durable caller request. + /// + /// # Errors + /// + /// Returns a safe conflict, keyring, persistence, or state error. Staged recovery transport + /// replaces this transitional generated-secret receipt in the custody phase. + #[allow(clippy::too_many_arguments)] + pub fn generate_account_durable( + &self, + request_id: &DurableRequestId, + expected_revision: u64, + accounts: &(impl AccountRepository + ?Sized), + app_state: &(impl AppStateRepository + ?Sized), + secrets: &(impl SecretStore + ?Sized), + operations: &(impl DurableOperationRepository + ?Sized), + clock: &(impl Clock + ?Sized), + ) -> Result<GenerateAccountReceipt, SafeError> { + self.require_revision(expected_revision)?; + let generated = self.key_material().generate()?; + let (public_key, npub, secret, nsec) = generated.into_parts(); + let account = AccountSummary::new( + AccountIdentity::verify(public_key, npub.as_str().to_owned())?, + LocalSignerBinding::new(public_key, BindingAvailability::Available), + None, + AccountCreatedAt::new(clock.now()), + None, + )?; + self.persist_account_durable( + request_id, + DurableOperationKind::Create, + expected_revision, + &account, + secret, + None, + accounts, + app_state, + secrets, + operations, + clock, + )?; + Ok(GenerateAccountReceipt { + account, + generated_nsec: nsec, + }) + } + + /// Imports or explicitly repairs one local account under a durable caller request. + /// + /// # Errors + /// + /// Returns a safe conflict, validation, keyring, persistence, or state error. + #[allow(clippy::too_many_arguments)] + pub fn import_secret_key_durable( + &self, + request_id: &DurableRequestId, + expected_revision: u64, + input: SecretKeyInput, + accounts: &(impl AccountRepository + ?Sized), + app_state: &(impl AppStateRepository + ?Sized), + secrets: &(impl SecretStore + ?Sized), + operations: &(impl DurableOperationRepository + ?Sized), + clock: &(impl Clock + ?Sized), + ) -> Result<ImportAccountReceipt, SafeError> { + if let Some(existing) = operations.load_durable_operation(request_id)? { + return if existing + .terminal() + .is_some_and(|receipt| receipt.outcome() == DurableTerminalOutcome::Completed) + { + accounts + .find_account(existing.account())? + .map(|account| ImportAccountReceipt { account }) + .ok_or_else(recovery_required) + } else { + Err(recovery_required()) + }; + } + self.require_revision(expected_revision)?; + let imported = self.key_material().import(input)?; + let (public_key, npub, secret) = imported.into_parts(); + let previous = accounts.find_account(public_key)?; + if let Some(existing) = &previous + && (existing.signer().availability() != BindingAvailability::CredentialMissing + || secrets.contains(public_key)?) + { + return Err(account_exists()); + } + if previous.is_none() && secrets.contains(public_key)? { + return Err(account_exists()); + } + let account = if let Some(existing) = &previous { + existing.with_binding_availability(BindingAvailability::Available) + } else { + AccountSummary::new( + AccountIdentity::verify(public_key, npub.as_str().to_owned())?, + LocalSignerBinding::new(public_key, BindingAvailability::Available), + None, + AccountCreatedAt::new(clock.now()), + None, + )? + }; + let kind = if previous.is_some() { + DurableOperationKind::Repair + } else { + DurableOperationKind::Import + }; + self.persist_account_durable( + request_id, + kind, + expected_revision, + &account, + secret, + previous.as_ref(), + accounts, + app_state, + secrets, + operations, + clock, + )?; + Ok(ImportAccountReceipt { account }) + } + + fn require_revision(&self, expected_revision: u64) -> Result<(), SafeError> { + if self.snapshot().revision().value() != expected_revision { + return Err(operation_conflict()); + } + Ok(()) + } + + #[allow(clippy::too_many_arguments)] + fn persist_account_durable( + &self, + request_id: &DurableRequestId, + kind: DurableOperationKind, + expected_revision: u64, + account: &AccountSummary, + secret: SecretKeyInput, + previous: Option<&AccountSummary>, + accounts: &(impl AccountRepository + ?Sized), + app_state: &(impl AppStateRepository + ?Sized), + secrets: &(impl SecretStore + ?Sized), + operations: &(impl DurableOperationRepository + ?Sized), + clock: &(impl Clock + ?Sized), + ) -> Result<(), SafeError> { + let prior = OperationPriorState::new( + app_state.load_selected_account()?, + previous.map(|account| account.signer().availability()), + ); + match operations.begin_durable_operation( + request_id, + kind, + account.public_key(), + Some(expected_revision), + prior, + clock.now(), + )? { + DurableOperationStart::Started(_) => {} + DurableOperationStart::Existing(operation) => { + return if operation + .terminal() + .is_some_and(|receipt| receipt.outcome() == DurableTerminalOutcome::Completed) + { + Ok(()) + } else { + Err(recovery_required()) + }; + } + } + secrets.put(account.public_key(), secret)?; + operations.advance_durable_operation( + request_id, + DurableOperationPhase::IntentRecorded, + DurableOperationPhase::CredentialWritten, + clock.now(), + None, + )?; + previous.map_or_else( + || accounts.insert_account(account), + |_| accounts.update_account(account), + )?; + operations.advance_durable_operation( + request_id, + DurableOperationPhase::CredentialWritten, + DurableOperationPhase::MetadataCommitted, + clock.now(), + None, + )?; + app_state.save_selected_account(Some(account.public_key()))?; + operations.advance_durable_operation( + request_id, + DurableOperationPhase::MetadataCommitted, + DurableOperationPhase::SelectionCommitted, + clock.now(), + None, + )?; + let snapshot = self.apply_transition(StateTransition::ReplaceRegistry { + accounts: accounts.list_accounts()?, + selected: Some(account.public_key()), + })?; + operations.finalize_durable_operation( + request_id, + DurableOperationPhase::SelectionCommitted, + DurableTerminalOutcome::Completed, + Some(snapshot.revision().value()), + clock.now(), + )?; + Ok(()) + } + + /// Issues a single-use confirmation bound to the target and current revision. + /// + /// # Errors + /// + /// Returns a safe account or application-state error. + pub fn request_account_removal( + &self, + public_key: PublicKey, + clock: &(impl Clock + ?Sized), + ) -> Result<RemovalConfirmationToken, SafeError> { + self.issue_removal_token(public_key, clock.now()) + } + + pub fn cancel_account_removal(&self, token: RemovalConfirmationToken) -> bool { + self.cancel_removal_token(token) + } + + /// Permanently removes a confirmed account and selects a deterministic fallback. + /// + /// # Errors + /// + /// Returns a safe confirmation, credential, persistence, recovery, or state error. + pub fn confirm_account_removal( + &self, + token: RemovalConfirmationToken, + accounts: &(impl AccountRepository + ?Sized), + app_state: &(impl AppStateRepository + ?Sized), + secrets: &(impl SecretStore + ?Sized), + journal: &(impl OperationJournal + ?Sized), + clock: &(impl Clock + ?Sized), + ) -> Result<crate::AppSnapshot, SafeError> { + let public_key = self.consume_removal_token(token, clock.now())?; + let registry = accounts.list_accounts()?; + let index = registry + .iter() + .position(|account| account.public_key() == public_key) + .ok_or_else(account_not_found)?; + let selected = if self.snapshot().selected_account() == Some(public_key) { + registry + .get(index + 1) + .or_else(|| index.checked_sub(1).and_then(|before| registry.get(before))) + .map(AccountSummary::public_key) + } else { + self.snapshot().selected_account() + }; + let operation = + journal.begin_operation(AccountOperationKind::Remove, public_key, clock.now())?; + let was_active = self + .snapshot() + .active_account() + .is_some_and(|active| active.account().public_key() == public_key); + if was_active { + self.sign_out()?; + } + let account = &registry[index]; + match secrets.delete(public_key) { + Ok(()) => {} + Err(error) + if error.code() == SafeErrorCode::CredentialMissing + && account.signer().availability() + == BindingAvailability::CredentialMissing => {} + Err(error) => return Err(error), + } + journal.update_operation( + operation, + AccountOperationPhase::CredentialDeleted, + clock.now(), + None, + )?; + accounts.remove_account(public_key)?; + app_state.save_selected_account(selected)?; + journal.update_operation( + operation, + AccountOperationPhase::MetadataDeleted, + clock.now(), + None, + )?; + journal.finalize_operation(operation)?; + self.apply_transition(StateTransition::ReplaceRegistryPreservingSession { + accounts: accounts.list_accounts()?, + selected, + }) + } + + /// Confirms and executes an expiring removal plan as a durable request. + /// + /// # Errors + /// + /// Returns a safe expiry, conflict, credential, persistence, or recovery error. + #[allow(clippy::too_many_arguments)] + pub fn confirm_account_removal_durable( + &self, + request_id: &DurableRequestId, + token: RemovalConfirmationToken, + accounts: &(impl AccountRepository + ?Sized), + app_state: &(impl AppStateRepository + ?Sized), + secrets: &(impl SecretStore + ?Sized), + operations: &(impl DurableOperationRepository + ?Sized), + clock: &(impl Clock + ?Sized), + ) -> Result<crate::AppSnapshot, SafeError> { + let expected_revision = token.revision().value(); + let public_key = self.consume_removal_token(token, clock.now())?; + self.require_revision(expected_revision)?; + let registry = accounts.list_accounts()?; + let index = registry + .iter() + .position(|account| account.public_key() == public_key) + .ok_or_else(account_not_found)?; + let selected = if self.snapshot().selected_account() == Some(public_key) { + registry + .get(index + 1) + .or_else(|| index.checked_sub(1).and_then(|before| registry.get(before))) + .map(AccountSummary::public_key) + } else { + self.snapshot().selected_account() + }; + let account = &registry[index]; + match operations.begin_durable_operation( + request_id, + DurableOperationKind::Remove, + public_key, + Some(expected_revision), + OperationPriorState::new(selected, Some(account.signer().availability())), + clock.now(), + )? { + DurableOperationStart::Started(_) => {} + DurableOperationStart::Existing(operation) => { + return if operation + .terminal() + .is_some_and(|receipt| receipt.outcome() == DurableTerminalOutcome::Completed) + { + Ok(self.snapshot()) + } else { + Err(recovery_required()) + }; + } + } + if self + .snapshot() + .active_account() + .is_some_and(|active| active.account().public_key() == public_key) + { + self.sign_out()?; + } + match secrets.delete(public_key) { + Ok(()) => {} + Err(error) + if error.code() == SafeErrorCode::CredentialMissing + && account.signer().availability() + == BindingAvailability::CredentialMissing => {} + Err(error) => return Err(error), + } + operations.advance_durable_operation( + request_id, + DurableOperationPhase::IntentRecorded, + DurableOperationPhase::CredentialDeleted, + clock.now(), + None, + )?; + accounts.remove_account(public_key)?; + operations.advance_durable_operation( + request_id, + DurableOperationPhase::CredentialDeleted, + DurableOperationPhase::MetadataDeleted, + clock.now(), + None, + )?; + app_state.save_selected_account(selected)?; + operations.advance_durable_operation( + request_id, + DurableOperationPhase::MetadataDeleted, + DurableOperationPhase::SelectionCommitted, + clock.now(), + None, + )?; + let snapshot = + self.apply_transition(StateTransition::ReplaceRegistryPreservingSession { + accounts: accounts.list_accounts()?, + selected, + })?; + operations.finalize_durable_operation( + request_id, + DurableOperationPhase::SelectionCommitted, + DurableTerminalOutcome::Completed, + Some(snapshot.revision().value()), + clock.now(), + )?; + Ok(snapshot) + } + + /// Persists and publishes a saved account selection without activating it. + /// + /// # Errors + /// + /// Returns a safe account, persistence, or application-state error. + pub fn select_account( + &self, + public_key: PublicKey, + accounts: &(impl AccountRepository + ?Sized), + app_state: &(impl AppStateRepository + ?Sized), + ) -> Result<crate::AppSnapshot, SafeError> { + if accounts.find_account(public_key)?.is_none() { + return Err(account_not_found()); + } + app_state.save_selected_account(Some(public_key))?; + self.apply_transition(StateTransition::Select(public_key)) + } + + /// Generates, stores, and selects one local Nostr account without activating it. + /// + /// # Errors + /// + /// Returns a safe key, credential, persistence, or application-state error. + pub fn generate_account( + &self, + accounts: &(impl AccountRepository + ?Sized), + app_state: &(impl AppStateRepository + ?Sized), + secrets: &(impl SecretStore + ?Sized), + journal: &(impl OperationJournal + ?Sized), + clock: &(impl Clock + ?Sized), + ) -> Result<GenerateAccountReceipt, SafeError> { + let generated = self.key_material().generate()?; + let (public_key, npub, secret, nsec) = generated.into_parts(); + let account = AccountSummary::new( + AccountIdentity::verify(public_key, npub.as_str().to_owned())?, + LocalSignerBinding::new(public_key, BindingAvailability::Available), + None, + AccountCreatedAt::new(clock.now()), + None, + )?; + Self::persist_account_transaction( + AccountOperationKind::Add, + &account, + secret, + None, + accounts, + app_state, + secrets, + journal, + clock, + )?; + let registry = accounts.list_accounts()?; + self.apply_transition(StateTransition::ReplaceRegistry { + accounts: registry, + selected: Some(public_key), + })?; + Ok(GenerateAccountReceipt { + account, + generated_nsec: nsec, + }) + } + + /// Imports, stores, and selects one local Nostr account without activating it. + /// + /// # Errors + /// + /// Returns a safe key, credential, persistence, or application-state error. + pub fn import_secret_key( + &self, + input: SecretKeyInput, + accounts: &(impl AccountRepository + ?Sized), + app_state: &(impl AppStateRepository + ?Sized), + secrets: &(impl SecretStore + ?Sized), + journal: &(impl OperationJournal + ?Sized), + clock: &(impl Clock + ?Sized), + ) -> Result<ImportAccountReceipt, SafeError> { + let imported = self.key_material().import(input)?; + let (public_key, npub, secret) = imported.into_parts(); + if let Some(existing) = accounts.find_account(public_key)? { + if existing.signer().availability() != BindingAvailability::CredentialMissing + || secrets.contains(public_key)? + { + return Err(account_exists()); + } + let repaired = existing.with_binding_availability(BindingAvailability::Available); + Self::persist_account_transaction( + AccountOperationKind::Import, + &repaired, + secret, + Some(&existing), + accounts, + app_state, + secrets, + journal, + clock, + )?; + self.apply_transition(StateTransition::ReplaceRegistry { + accounts: accounts.list_accounts()?, + selected: Some(public_key), + })?; + return Ok(ImportAccountReceipt { account: repaired }); + } + if secrets.contains(public_key)? { + return Err(account_exists()); + } + let account = AccountSummary::new( + AccountIdentity::verify(public_key, npub.as_str().to_owned())?, + LocalSignerBinding::new(public_key, BindingAvailability::Available), + None, + AccountCreatedAt::new(clock.now()), + None, + )?; + Self::persist_account_transaction( + AccountOperationKind::Import, + &account, + secret, + None, + accounts, + app_state, + secrets, + journal, + clock, + )?; + self.apply_transition(StateTransition::ReplaceRegistry { + accounts: accounts.list_accounts()?, + selected: Some(public_key), + })?; + Ok(ImportAccountReceipt { account }) + } + + #[allow(clippy::too_many_arguments)] + fn persist_account_transaction( + kind: AccountOperationKind, + account: &AccountSummary, + secret: SecretKeyInput, + previous: Option<&AccountSummary>, + accounts: &(impl AccountRepository + ?Sized), + app_state: &(impl AppStateRepository + ?Sized), + secrets: &(impl SecretStore + ?Sized), + journal: &(impl OperationJournal + ?Sized), + clock: &(impl Clock + ?Sized), + ) -> Result<(), SafeError> { + let public_key = account.public_key(); + let previous_selection = app_state.load_selected_account()?; + let operation = journal.begin_operation(kind, public_key, clock.now())?; + if let Err(error) = secrets.put(public_key, secret) { + let _ = journal.finalize_operation(operation); + return Err(error); + } + if let Err(error) = journal.update_operation( + operation, + AccountOperationPhase::CredentialWritten, + clock.now(), + None, + ) { + return compensate_account_write( + operation, + public_key, + error, + None, + previous_selection, + accounts, + app_state, + secrets, + journal, + clock, + ); + } + let metadata_result = previous.map_or_else( + || accounts.insert_account(account), + |_| accounts.update_account(account), + ); + if let Err(error) = metadata_result { + return compensate_account_write( + operation, + public_key, + error, + previous, + previous_selection, + accounts, + app_state, + secrets, + journal, + clock, + ); + } + if let Err(error) = app_state.save_selected_account(Some(public_key)) { + return compensate_account_write( + operation, + public_key, + error, + previous, + previous_selection, + accounts, + app_state, + secrets, + journal, + clock, + ); + } + journal.update_operation( + operation, + AccountOperationPhase::MetadataCommitted, + clock.now(), + None, + )?; + journal.finalize_operation(operation) + } +} + +#[allow(clippy::too_many_arguments)] +fn compensate_account_write( + operation: OperationId, + public_key: PublicKey, + original_error: SafeError, + previous: Option<&AccountSummary>, + previous_selection: Option<PublicKey>, + accounts: &(impl AccountRepository + ?Sized), + app_state: &(impl AppStateRepository + ?Sized), + secrets: &(impl SecretStore + ?Sized), + journal: &(impl OperationJournal + ?Sized), + clock: &(impl Clock + ?Sized), +) -> Result<(), SafeError> { + let metadata_rollback = if let Some(previous) = previous { + accounts.update_account(previous) + } else { + accounts.remove_account(public_key) + }; + let selection_rollback = app_state.save_selected_account(previous_selection); + let credential_rollback = secrets.delete(public_key); + if metadata_rollback.is_err() || selection_rollback.is_err() || credential_rollback.is_err() { + let _ = journal.update_operation( + operation, + AccountOperationPhase::CompensationPending, + clock.now(), + Some(OperationDiagnostic::CompensationFailed), + ); + return Err(recovery_required()); + } + let _ = journal.finalize_operation(operation); + Err(original_error) +} + +#[derive(Default)] +pub struct InMemoryOperationJournal { + state: Mutex<InMemoryJournalState>, +} + +#[derive(Default)] +struct InMemoryJournalState { + next_id: u64, + pending: Vec<PendingAccountOperation>, +} + +impl OperationJournal for InMemoryOperationJournal { + fn begin_operation( + &self, + kind: AccountOperationKind, + subject: PublicKey, + updated_at: radroots_studio_domain::UnixTimestamp, + ) -> Result<OperationId, SafeError> { + let mut state = self + .state + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + state.next_id = state.next_id.checked_add(1).ok_or_else(recovery_required)?; + let id = OperationId::from_raw(state.next_id); + state.pending.push(PendingAccountOperation::new( + id, + kind, + subject, + AccountOperationPhase::IntentRecorded, + updated_at, + None, + )); + Ok(id) + } + + fn update_operation( + &self, + id: OperationId, + phase: AccountOperationPhase, + updated_at: radroots_studio_domain::UnixTimestamp, + diagnostic: Option<OperationDiagnostic>, + ) -> Result<(), SafeError> { + let mut state = self + .state + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + let operation = state + .pending + .iter_mut() + .find(|operation| operation.id() == id) + .ok_or_else(recovery_required)?; + *operation = PendingAccountOperation::new( + id, + operation.kind(), + operation.subject(), + phase, + updated_at, + diagnostic, + ); + Ok(()) + } + + fn list_pending_operations(&self) -> Result<Vec<PendingAccountOperation>, SafeError> { + Ok(self + .state + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .pending + .clone()) + } + + fn finalize_operation(&self, id: OperationId) -> Result<(), SafeError> { + self.state + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .pending + .retain(|operation| operation.id() != id); + Ok(()) + } +} + +#[derive(Default)] +pub struct InMemoryAccountRepository { + state: Mutex<InMemoryAccountState>, +} + +#[derive(Default)] +struct InMemoryAccountState { + accounts: Vec<AccountSummary>, + selected: Option<PublicKey>, +} + +impl InMemoryAccountRepository { + fn state(&self) -> MutexGuard<'_, InMemoryAccountState> { + self.state + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + } +} + +impl AccountRepository for InMemoryAccountRepository { + fn list_accounts(&self) -> Result<Vec<AccountSummary>, SafeError> { + Ok(self.state().accounts.clone()) + } + + fn find_account(&self, public_key: PublicKey) -> Result<Option<AccountSummary>, SafeError> { + Ok(self + .state() + .accounts + .iter() + .find(|account| account.public_key() == public_key) + .cloned()) + } + + fn insert_account(&self, account: &AccountSummary) -> Result<(), SafeError> { + let mut state = self.state(); + if state + .accounts + .iter() + .any(|saved| saved.public_key() == account.public_key()) + { + return Err(account_exists()); + } + state.accounts.push(account.clone()); + state + .accounts + .sort_by_key(|saved| (saved.created_at().timestamp(), saved.public_key())); + Ok(()) + } + + fn update_account(&self, account: &AccountSummary) -> Result<(), SafeError> { + let mut state = self.state(); + let saved = state + .accounts + .iter_mut() + .find(|saved| saved.public_key() == account.public_key()) + .ok_or_else(account_not_found)?; + *saved = account.clone(); + Ok(()) + } + + fn remove_account(&self, public_key: PublicKey) -> Result<(), SafeError> { + let mut state = self.state(); + state + .accounts + .retain(|account| account.public_key() != public_key); + if state.selected == Some(public_key) { + state.selected = None; + } + Ok(()) + } +} + +impl AppStateRepository for InMemoryAccountRepository { + fn load_selected_account(&self) -> Result<Option<PublicKey>, SafeError> { + Ok(self.state().selected) + } + + fn save_selected_account(&self, public_key: Option<PublicKey>) -> Result<(), SafeError> { + let mut state = self.state(); + if public_key.is_some_and(|key| { + !state + .accounts + .iter() + .any(|account| account.public_key() == key) + }) { + return Err(account_not_found()); + } + state.selected = public_key; + Ok(()) + } +} + +const fn account_exists() -> SafeError { + SafeError::new( + SafeErrorCode::AccountAlreadyExists, + SafeMessage::new("The Nostr account is already saved."), + ) +} + +const fn account_not_found() -> SafeError { + SafeError::new( + SafeErrorCode::AccountNotFound, + SafeMessage::new("The account was not found."), + ) +} + +const fn recovery_required() -> SafeError { + SafeError::new( + SafeErrorCode::PendingOperationRecoveryRequired, + SafeMessage::new("Account recovery is required before this operation can continue."), + ) +} + +const fn operation_conflict() -> SafeError { + SafeError::new( + SafeErrorCode::InvalidApplicationState, + SafeMessage::new("The account operation conflicts with the current application state."), + ) +} + +#[cfg(test)] +mod tests { + use std::sync::atomic::{AtomicBool, Ordering}; + + use radroots_studio_domain::{ + AccountCreatedAt, AccountIdentity, AccountSummary, BindingAvailability, LocalSignerBinding, + PublicKey, SafeError, SafeErrorCode, SafeMessage, SecretKeyInput, UnixTimestamp, + }; + + use super::InMemoryAccountRepository; + use crate::{ + AccountOperationPhase, AccountRepository, AppCore, AppStateRepository, Clock, + DurableOperationKind, DurableOperationPhase, FailureSecretStore, InMemoryOperationJournal, + InMemorySecretStore, OperationJournal, ProfileRefreshStatus, ProfileRepository, + RelayConfiguration, SecretStore, SecretStoreOperation, SessionState, StateTransition, + recovery::tests::{TestDurableRepository, operation as durable_operation}, + }; + + struct FixedClock; + + impl Clock for FixedClock { + fn now(&self) -> UnixTimestamp { + UnixTimestamp::from_seconds(10).expect("time") + } + } + + struct LateClock; + + impl Clock for LateClock { + fn now(&self) -> UnixTimestamp { + UnixTimestamp::from_seconds(311).expect("time") + } + } + + struct EmptyProfiles; + + impl ProfileRepository for EmptyProfiles { + fn load_profile( + &self, + _public_key: PublicKey, + ) -> Result<Option<crate::CachedProfile>, SafeError> { + Ok(None) + } + + fn save_profile(&self, _profile: &crate::CachedProfile) -> Result<(), SafeError> { + Ok(()) + } + + fn record_refresh_status( + &self, + _public_key: PublicKey, + _refreshed_at: UnixTimestamp, + _status: ProfileRefreshStatus, + ) -> Result<(), SafeError> { + Ok(()) + } + + fn remove_profile(&self, _public_key: PublicKey) -> Result<(), SafeError> { + Ok(()) + } + } + + #[derive(Default)] + struct FailingUpdateJournal(InMemoryOperationJournal); + + impl OperationJournal for FailingUpdateJournal { + fn begin_operation( + &self, + kind: crate::AccountOperationKind, + subject: PublicKey, + updated_at: UnixTimestamp, + ) -> Result<crate::OperationId, SafeError> { + self.0.begin_operation(kind, subject, updated_at) + } + + fn update_operation( + &self, + _id: crate::OperationId, + _phase: AccountOperationPhase, + _updated_at: UnixTimestamp, + _diagnostic: Option<crate::OperationDiagnostic>, + ) -> Result<(), SafeError> { + Err(SafeError::new( + SafeErrorCode::StorageUnavailable, + SafeMessage::new("The test journal is unavailable."), + )) + } + + fn list_pending_operations( + &self, + ) -> Result<Vec<crate::PendingAccountOperation>, SafeError> { + self.0.list_pending_operations() + } + + fn finalize_operation(&self, id: crate::OperationId) -> Result<(), SafeError> { + self.0.finalize_operation(id) + } + } + + #[derive(Default)] + struct FailingInsertRepository { + inner: InMemoryAccountRepository, + } + + #[derive(Default)] + struct FailingSelectionRepository { + inner: InMemoryAccountRepository, + fail_next_selection: AtomicBool, + } + + impl AccountRepository for FailingSelectionRepository { + fn list_accounts(&self) -> Result<Vec<AccountSummary>, SafeError> { + self.inner.list_accounts() + } + + fn find_account(&self, public_key: PublicKey) -> Result<Option<AccountSummary>, SafeError> { + self.inner.find_account(public_key) + } + + fn insert_account(&self, account: &AccountSummary) -> Result<(), SafeError> { + self.inner.insert_account(account) + } + + fn update_account(&self, account: &AccountSummary) -> Result<(), SafeError> { + self.inner.update_account(account) + } + + fn remove_account(&self, public_key: PublicKey) -> Result<(), SafeError> { + self.inner.remove_account(public_key) + } + } + + impl AppStateRepository for FailingSelectionRepository { + fn load_selected_account(&self) -> Result<Option<PublicKey>, SafeError> { + self.inner.load_selected_account() + } + + fn save_selected_account(&self, public_key: Option<PublicKey>) -> Result<(), SafeError> { + if self.fail_next_selection.swap(false, Ordering::SeqCst) { + return Err(SafeError::new( + SafeErrorCode::StorageUnavailable, + SafeMessage::new("The test selection repository is unavailable."), + )); + } + self.inner.save_selected_account(public_key) + } + } + + impl AccountRepository for FailingInsertRepository { + fn list_accounts(&self) -> Result<Vec<AccountSummary>, SafeError> { + self.inner.list_accounts() + } + + fn find_account(&self, public_key: PublicKey) -> Result<Option<AccountSummary>, SafeError> { + self.inner.find_account(public_key) + } + + fn insert_account(&self, _account: &AccountSummary) -> Result<(), SafeError> { + Err(SafeError::new( + SafeErrorCode::StorageUnavailable, + SafeMessage::new("The test account repository is unavailable."), + )) + } + + fn update_account(&self, account: &AccountSummary) -> Result<(), SafeError> { + self.inner.update_account(account) + } + + fn remove_account(&self, public_key: PublicKey) -> Result<(), SafeError> { + self.inner.remove_account(public_key) + } + } + + impl AppStateRepository for FailingInsertRepository { + fn load_selected_account(&self) -> Result<Option<PublicKey>, SafeError> { + self.inner.load_selected_account() + } + + fn save_selected_account(&self, public_key: Option<PublicKey>) -> Result<(), SafeError> { + self.inner.save_selected_account(public_key) + } + } + + #[test] + fn generate_account_stores_selects_and_returns_one_time_nsec_without_activation() { + let core = AppCore::in_memory(RelayConfiguration::default()); + let accounts = InMemoryAccountRepository::default(); + let secrets = InMemorySecretStore::default(); + let journal = InMemoryOperationJournal::default(); + core.bootstrap().expect("bootstrap"); + + let receipt = core + .generate_account(&accounts, &accounts, &secrets, &journal, &FixedClock) + .expect("generate"); + let public_key = receipt.account().public_key(); + assert_eq!(public_key.to_hex().len(), 64); + assert!(secrets.contains(public_key).expect("credential")); + assert_eq!( + accounts.load_selected_account().expect("selection"), + Some(public_key) + ); + assert_eq!(core.snapshot().selected_account(), Some(public_key)); + assert_eq!(core.snapshot().session(), SessionState::SignedOut); + assert!(core.snapshot().active_account().is_none()); + assert_eq!(receipt.generated_nsec().with_exposed_secret(str::len), 63); + assert!(!format!("{:?}", core.snapshot()).contains("nsec1")); + } + + #[test] + fn import_secret_key_accepts_nsec_and_hex_without_exposing_or_activating() { + for input in [ + "nsec1vl029mgpspedva04g90vltkh6fvh240zqtv9k0t9af8935ke9laqsnlfe5", + "7e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7", + ] { + let core = AppCore::in_memory(RelayConfiguration::default()); + let accounts = InMemoryAccountRepository::default(); + let secrets = InMemorySecretStore::default(); + let journal = InMemoryOperationJournal::default(); + core.bootstrap().expect("bootstrap"); + let receipt = core + .import_secret_key( + SecretKeyInput::parse(input.to_owned()).expect("input"), + &accounts, + &accounts, + &secrets, + &journal, + &FixedClock, + ) + .expect("import"); + let public_key = receipt.account().public_key(); + assert!(secrets.contains(public_key).expect("credential")); + assert_eq!(core.snapshot().selected_account(), Some(public_key)); + assert_eq!(core.snapshot().session(), SessionState::SignedOut); + assert!(!format!("{:?}", core.snapshot()).contains(input)); + } + } + + #[test] + fn import_secret_key_rejects_invalid_nsec_checksum_before_persistence() { + let core = AppCore::in_memory(RelayConfiguration::default()); + let accounts = InMemoryAccountRepository::default(); + let secrets = InMemorySecretStore::default(); + let journal = InMemoryOperationJournal::default(); + core.bootstrap().expect("bootstrap"); + let input = SecretKeyInput::parse( + "nsec1qqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqq".to_owned(), + ) + .expect("domain shape"); + let error = core + .import_secret_key(input, &accounts, &accounts, &secrets, &journal, &FixedClock) + .expect_err("invalid import"); + assert_eq!(error.code(), SafeErrorCode::InvalidSecretKey); + assert!(core.snapshot().accounts().is_empty()); + } + + #[test] + fn duplicate_import_preserves_existing_credential_and_snapshot() { + let core = AppCore::in_memory(RelayConfiguration::default()); + let accounts = InMemoryAccountRepository::default(); + let secrets = InMemorySecretStore::default(); + let journal = InMemoryOperationJournal::default(); + core.bootstrap().expect("bootstrap"); + let import = || { + SecretKeyInput::parse( + "7e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7".to_owned(), + ) + .expect("input") + }; + core.import_secret_key( + import(), + &accounts, + &accounts, + &secrets, + &journal, + &FixedClock, + ) + .expect("first import"); + let before = core.snapshot(); + let error = core + .import_secret_key( + import(), + &accounts, + &accounts, + &secrets, + &journal, + &FixedClock, + ) + .expect_err("duplicate"); + assert_eq!(error.code(), SafeErrorCode::AccountAlreadyExists); + assert_eq!(core.snapshot(), before); + assert_eq!(core.snapshot().accounts().len(), 1); + } + + #[test] + fn duplicate_import_repairs_only_explicit_missing_credential_account() { + let core = AppCore::in_memory(RelayConfiguration::default()); + let accounts = InMemoryAccountRepository::default(); + let secrets = InMemorySecretStore::default(); + let journal = InMemoryOperationJournal::default(); + core.bootstrap().expect("bootstrap"); + let input = || { + SecretKeyInput::parse( + "7e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7".to_owned(), + ) + .expect("input") + }; + let imported = core.key_material().import(input()).expect("derive"); + let (public_key, npub, _) = imported.into_parts(); + let missing = AccountSummary::new( + AccountIdentity::verify(public_key, npub.as_str().to_owned()).expect("identity"), + LocalSignerBinding::new(public_key, BindingAvailability::CredentialMissing), + None, + AccountCreatedAt::new(FixedClock.now()), + None, + ) + .expect("missing account"); + accounts.insert_account(&missing).expect("missing metadata"); + accounts + .save_selected_account(Some(public_key)) + .expect("selection"); + core.apply_transition(StateTransition::ReplaceRegistry { + accounts: vec![missing], + selected: Some(public_key), + }) + .expect("registry"); + + let receipt = core + .import_secret_key( + input(), + &accounts, + &accounts, + &secrets, + &journal, + &FixedClock, + ) + .expect("repair"); + assert_eq!( + receipt.account().signer().availability(), + BindingAvailability::Available + ); + assert!(secrets.contains(public_key).expect("credential")); + assert_eq!(core.snapshot().accounts().len(), 1); + } + + #[test] + fn account_transaction_publishes_nothing_when_credential_write_fails() { + let core = AppCore::in_memory(RelayConfiguration::default()); + let accounts = InMemoryAccountRepository::default(); + let secrets = FailureSecretStore::default(); + let journal = InMemoryOperationJournal::default(); + core.bootstrap().expect("bootstrap"); + secrets.fail_next(SecretStoreOperation::Put); + + let error = core + .generate_account(&accounts, &accounts, &secrets, &journal, &FixedClock) + .err() + .expect("credential failure"); + assert_eq!(error.code(), SafeErrorCode::KeyringUnavailable); + assert!(core.snapshot().accounts().is_empty()); + assert!( + journal + .list_pending_operations() + .expect("journal") + .is_empty() + ); + } + + #[test] + fn account_transaction_removes_written_credential_when_metadata_fails() { + let core = AppCore::in_memory(RelayConfiguration::default()); + let accounts = FailingInsertRepository::default(); + let secrets = FailureSecretStore::default(); + let journal = InMemoryOperationJournal::default(); + core.bootstrap().expect("bootstrap"); + + let error = core + .generate_account(&accounts, &accounts, &secrets, &journal, &FixedClock) + .err() + .expect("metadata failure"); + assert_eq!(error.code(), SafeErrorCode::StorageUnavailable); + let calls = secrets.calls(); + assert_eq!(calls[0].operation(), SecretStoreOperation::Put); + assert_eq!(calls[1].operation(), SecretStoreOperation::Delete); + assert_eq!(calls[0].public_key(), calls[1].public_key()); + assert!(core.snapshot().accounts().is_empty()); + assert!( + journal + .list_pending_operations() + .expect("journal") + .is_empty() + ); + } + + #[test] + fn account_transaction_rolls_back_metadata_and_credential_when_selection_fails() { + let core = AppCore::in_memory(RelayConfiguration::default()); + let accounts = FailingSelectionRepository::default(); + let secrets = FailureSecretStore::default(); + let journal = InMemoryOperationJournal::default(); + core.bootstrap().expect("bootstrap"); + accounts.fail_next_selection.store(true, Ordering::SeqCst); + + let error = core + .generate_account(&accounts, &accounts, &secrets, &journal, &FixedClock) + .err() + .expect("selection failure"); + + assert_eq!(error.code(), SafeErrorCode::StorageUnavailable); + assert!(accounts.list_accounts().expect("accounts").is_empty()); + assert_eq!(accounts.load_selected_account().expect("selection"), None); + let calls = secrets.calls(); + assert_eq!(calls[0].operation(), SecretStoreOperation::Put); + assert_eq!(calls[1].operation(), SecretStoreOperation::Delete); + assert_eq!(calls[0].public_key(), calls[1].public_key()); + assert!(core.snapshot().accounts().is_empty()); + assert!( + journal + .list_pending_operations() + .expect("journal") + .is_empty() + ); + } + + #[test] + fn account_transaction_retains_non_secret_journal_when_compensation_fails() { + let core = AppCore::in_memory(RelayConfiguration::default()); + let accounts = FailingInsertRepository::default(); + let secrets = FailureSecretStore::default(); + let journal = InMemoryOperationJournal::default(); + core.bootstrap().expect("bootstrap"); + secrets.fail_next(SecretStoreOperation::Delete); + + let error = core + .generate_account(&accounts, &accounts, &secrets, &journal, &FixedClock) + .err() + .expect("recovery required"); + assert_eq!( + error.code(), + SafeErrorCode::PendingOperationRecoveryRequired + ); + let pending = journal.list_pending_operations().expect("journal"); + assert_eq!(pending.len(), 1); + assert_eq!( + pending[0].phase(), + AccountOperationPhase::CompensationPending + ); + assert!(!format!("{pending:?}").contains("nsec1")); + assert!(core.snapshot().accounts().is_empty()); + } + + #[test] + fn select_account_persists_existing_choice_without_activating() { + let core = AppCore::in_memory(RelayConfiguration::default()); + let accounts = InMemoryAccountRepository::default(); + let secrets = InMemorySecretStore::default(); + let journal = InMemoryOperationJournal::default(); + core.bootstrap().expect("bootstrap"); + let first = core + .generate_account(&accounts, &accounts, &secrets, &journal, &FixedClock) + .expect("first") + .account() + .public_key(); + core.generate_account(&accounts, &accounts, &secrets, &journal, &FixedClock) + .expect("second"); + + let selected = core + .select_account(first, &accounts, &accounts) + .expect("select first"); + assert_eq!(selected.selected_account(), Some(first)); + assert_eq!(selected.session(), SessionState::SignedOut); + assert!(selected.active_account().is_none()); + assert_eq!( + accounts.load_selected_account().expect("saved"), + Some(first) + ); + let missing = core + .select_account( + PublicKey::from_hex( + "e0266e3cfb0d2886f91c73f5f868f3b98273713e5fcd97c081663f5518a4b3af", + ) + .expect("unknown public key"), + &accounts, + &accounts, + ) + .expect_err("missing account"); + assert_eq!(missing.code(), SafeErrorCode::AccountNotFound); + assert_eq!(core.snapshot(), selected); + } + + #[test] + fn remove_account_requires_fresh_single_use_confirmation_and_selects_next_fallback() { + let core = AppCore::in_memory(RelayConfiguration::default()); + let accounts = InMemoryAccountRepository::default(); + let secrets = InMemorySecretStore::default(); + let journal = InMemoryOperationJournal::default(); + core.bootstrap().expect("bootstrap"); + let first = core + .generate_account(&accounts, &accounts, &secrets, &journal, &FixedClock) + .expect("first") + .account() + .public_key(); + let second = core + .generate_account(&accounts, &accounts, &secrets, &journal, &FixedClock) + .expect("second") + .account() + .public_key(); + core.select_account(first, &accounts, &accounts) + .expect("select first"); + let stale = core + .request_account_removal(first, &FixedClock) + .expect("stale token"); + core.select_account(second, &accounts, &accounts) + .expect("change revision"); + let stale_error = core + .confirm_account_removal(stale, &accounts, &accounts, &secrets, &journal, &FixedClock) + .expect_err("stale token"); + assert_eq!(stale_error.code(), SafeErrorCode::InvalidApplicationState); + assert_eq!(core.snapshot().accounts().len(), 2); + + core.select_account(first, &accounts, &accounts) + .expect("reselect first"); + let token = core + .request_account_removal(first, &FixedClock) + .expect("token"); + let removed = core + .confirm_account_removal(token, &accounts, &accounts, &secrets, &journal, &FixedClock) + .expect("remove"); + assert_eq!(removed.accounts().len(), 1); + assert_eq!(removed.selected_account(), Some(second)); + assert!(!secrets.contains(first).expect("credential removed")); + assert_eq!(removed.session(), SessionState::SignedOut); + } + + #[test] + fn removal_preflight_reports_impact_expires_and_can_be_cancelled() { + let core = AppCore::in_memory(RelayConfiguration::default()); + let accounts = InMemoryAccountRepository::default(); + let secrets = InMemorySecretStore::default(); + let journal = InMemoryOperationJournal::default(); + core.bootstrap().expect("bootstrap"); + let account = core + .generate_account(&accounts, &accounts, &secrets, &journal, &FixedClock) + .expect("account") + .account() + .public_key(); + let expired = core + .request_account_removal(account, &FixedClock) + .expect("plan"); + assert!(expired.impact().deletes_local_credential()); + assert!(!expired.impact().signs_out()); + assert!( + core.confirm_account_removal( + expired, &accounts, &accounts, &secrets, &journal, &LateClock, + ) + .is_err() + ); + let cancelled = core + .request_account_removal(account, &FixedClock) + .expect("replacement plan"); + assert!(core.cancel_account_removal(cancelled)); + assert_eq!(core.snapshot().accounts().len(), 1); + } + + #[test] + fn import_rejects_orphan_credentials_and_durable_nonterminal_replays() { + const SECRET: &str = "7e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7"; + let core = AppCore::in_memory(RelayConfiguration::default()); + let accounts = InMemoryAccountRepository::default(); + let secrets = InMemorySecretStore::default(); + let journal = InMemoryOperationJournal::default(); + core.bootstrap().expect("bootstrap"); + let material = core + .key_material() + .import(SecretKeyInput::parse(SECRET.to_owned()).expect("secret")) + .expect("key material"); + let (public_key, _npub, secret) = material.into_parts(); + secrets.put(public_key, secret).expect("orphan credential"); + + assert_eq!( + core.import_secret_key( + SecretKeyInput::parse(SECRET.to_owned()).expect("secret"), + &accounts, + &accounts, + &secrets, + &journal, + &FixedClock, + ) + .expect_err("orphan credential must fail") + .code(), + SafeErrorCode::AccountAlreadyExists + ); + + let pending = durable_operation( + DurableOperationKind::Import, + DurableOperationPhase::IntentRecorded, + public_key, + None, + ); + let request_id = pending.request_id().clone(); + let operations = TestDurableRepository::new(pending); + assert_eq!( + core.import_secret_key_durable( + &request_id, + core.snapshot().revision().value(), + SecretKeyInput::parse(SECRET.to_owned()).expect("secret"), + &accounts, + &accounts, + &secrets, + &operations, + &FixedClock, + ) + .expect_err("unfinished replay must require recovery") + .code(), + SafeErrorCode::PendingOperationRecoveryRequired + ); + } + + #[test] + fn durable_import_covers_new_and_missing_credential_repair_paths() { + const SECRET: &str = "7e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7"; + for repair in [false, true] { + let core = AppCore::in_memory(RelayConfiguration::default()); + let accounts = InMemoryAccountRepository::default(); + let secrets = InMemorySecretStore::default(); + let material = core + .key_material() + .import(SecretKeyInput::parse(SECRET.to_owned()).expect("secret")) + .expect("key material"); + let (public_key, npub, secret) = material.into_parts(); + drop(secret); + if repair { + let account = AccountSummary::new( + AccountIdentity::verify(public_key, npub.as_str().to_owned()) + .expect("identity"), + LocalSignerBinding::new(public_key, BindingAvailability::CredentialMissing), + None, + AccountCreatedAt::new(FixedClock.now()), + None, + ) + .expect("account"); + accounts.insert_account(&account).expect("insert account"); + accounts + .save_selected_account(Some(public_key)) + .expect("selection"); + core.apply_transition(StateTransition::BootstrapRegistry { + accounts: vec![account], + selected: Some(public_key), + }) + .expect("registry"); + } else { + core.bootstrap().expect("bootstrap"); + } + let kind = if repair { + DurableOperationKind::Repair + } else { + DurableOperationKind::Import + }; + let pending = durable_operation( + kind, + DurableOperationPhase::IntentRecorded, + public_key, + repair.then_some(BindingAvailability::CredentialMissing), + ); + let request_id = pending.request_id().clone(); + let operations = TestDurableRepository::fresh(pending); + let receipt = core + .import_secret_key_durable( + &request_id, + core.snapshot().revision().value(), + SecretKeyInput::parse(SECRET.to_owned()).expect("secret"), + &accounts, + &accounts, + &secrets, + &operations, + &FixedClock, + ) + .expect("durable import"); + assert_eq!(receipt.account().public_key(), public_key); + assert_eq!( + operations.operation().phase(), + DurableOperationPhase::Finalized + ); + } + } + + #[test] + fn removal_of_unselected_account_preserves_the_current_selection() { + let core = AppCore::in_memory(RelayConfiguration::default()); + let accounts = InMemoryAccountRepository::default(); + let secrets = InMemorySecretStore::default(); + let journal = InMemoryOperationJournal::default(); + core.bootstrap().expect("bootstrap"); + let first = core + .generate_account(&accounts, &accounts, &secrets, &journal, &FixedClock) + .expect("first") + .account() + .public_key(); + let second = core + .generate_account(&accounts, &accounts, &secrets, &journal, &FixedClock) + .expect("second") + .account() + .public_key(); + let token = core + .request_account_removal(first, &FixedClock) + .expect("removal token"); + let snapshot = core + .confirm_account_removal(token, &accounts, &accounts, &secrets, &journal, &FixedClock) + .expect("remove unselected account"); + assert_eq!(snapshot.selected_account(), Some(second)); + + let missing = crate::test_support::valid_test_public_key(99).expect("missing key"); + assert!(accounts.insert_account(&snapshot.accounts()[0]).is_err()); + assert!(accounts.save_selected_account(Some(missing)).is_err()); + + let third = core + .generate_account(&accounts, &accounts, &secrets, &journal, &FixedClock) + .expect("third") + .account() + .public_key(); + let token = core + .request_account_removal(second, &FixedClock) + .expect("durable removal token"); + let pending = durable_operation( + DurableOperationKind::Remove, + DurableOperationPhase::IntentRecorded, + second, + Some(BindingAvailability::Available), + ); + let request_id = pending.request_id().clone(); + let operations = TestDurableRepository::fresh(pending); + let snapshot = core + .confirm_account_removal_durable( + &request_id, + token, + &accounts, + &accounts, + &secrets, + &operations, + &FixedClock, + ) + .expect("durable unselected removal"); + assert_eq!(snapshot.selected_account(), Some(third)); + } + + #[test] + fn duplicate_missing_binding_with_orphan_credential_fails_closed() { + const SECRET: &str = "7e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7"; + let core = AppCore::in_memory(RelayConfiguration::default()); + let accounts = InMemoryAccountRepository::default(); + let secrets = InMemorySecretStore::default(); + let journal = InMemoryOperationJournal::default(); + let material = core + .key_material() + .import(SecretKeyInput::parse(SECRET.to_owned()).expect("secret")) + .expect("key material"); + let (public_key, npub, secret) = material.into_parts(); + let account = AccountSummary::new( + AccountIdentity::verify(public_key, npub.as_str().to_owned()).expect("identity"), + LocalSignerBinding::new(public_key, BindingAvailability::CredentialMissing), + None, + AccountCreatedAt::new(FixedClock.now()), + None, + ) + .expect("account"); + accounts.insert_account(&account).expect("insert account"); + accounts + .save_selected_account(Some(public_key)) + .expect("selection"); + secrets.put(public_key, secret).expect("credential"); + core.apply_transition(StateTransition::BootstrapRegistry { + accounts: vec![account], + selected: Some(public_key), + }) + .expect("registry"); + + assert_eq!( + core.import_secret_key( + SecretKeyInput::parse(SECRET.to_owned()).expect("secret"), + &accounts, + &accounts, + &secrets, + &journal, + &FixedClock, + ) + .expect_err("orphan credential must fail") + .code(), + SafeErrorCode::AccountAlreadyExists + ); + let pending = durable_operation( + DurableOperationKind::Repair, + DurableOperationPhase::IntentRecorded, + public_key, + Some(BindingAvailability::CredentialMissing), + ); + let request_id = pending.request_id().clone(); + let operations = TestDurableRepository::fresh(pending); + assert_eq!( + core.import_secret_key_durable( + &request_id, + core.snapshot().revision().value(), + SecretKeyInput::parse(SECRET.to_owned()).expect("secret"), + &accounts, + &accounts, + &secrets, + &operations, + &FixedClock, + ) + .expect_err("orphan durable credential must fail") + .code(), + SafeErrorCode::AccountAlreadyExists + ); + } + + #[test] + fn removing_an_active_account_signs_out_for_legacy_and_durable_requests() { + for durable in [false, true] { + let core = AppCore::in_memory(RelayConfiguration::default()); + let accounts = InMemoryAccountRepository::default(); + let secrets = InMemorySecretStore::default(); + let journal = InMemoryOperationJournal::default(); + core.bootstrap().expect("bootstrap"); + let public_key = core + .import_secret_key( + SecretKeyInput::parse( + "7e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7" + .to_owned(), + ) + .expect("secret"), + &accounts, + &accounts, + &secrets, + &journal, + &FixedClock, + ) + .expect("account") + .account() + .public_key(); + core.activate_account( + public_key, + &accounts, + &accounts, + &EmptyProfiles, + &secrets, + &FixedClock, + ) + .expect("activate account"); + let token = core + .request_account_removal(public_key, &FixedClock) + .expect("removal token"); + let snapshot = if durable { + let pending = durable_operation( + DurableOperationKind::Remove, + DurableOperationPhase::IntentRecorded, + public_key, + Some(BindingAvailability::Available), + ); + let request_id = pending.request_id().clone(); + let operations = TestDurableRepository::fresh(pending); + core.confirm_account_removal_durable( + &request_id, + token, + &accounts, + &accounts, + &secrets, + &operations, + &FixedClock, + ) + .expect("durable removal") + } else { + core.confirm_account_removal( + token, + &accounts, + &accounts, + &secrets, + &journal, + &FixedClock, + ) + .expect("removal") + }; + assert_eq!(snapshot.session(), SessionState::SignedOut); + } + } + + #[test] + fn account_transaction_compensates_a_journal_phase_failure() { + let core = AppCore::in_memory(RelayConfiguration::default()); + let accounts = InMemoryAccountRepository::default(); + let secrets = InMemorySecretStore::default(); + let journal = FailingUpdateJournal::default(); + core.bootstrap().expect("bootstrap"); + + assert_eq!( + core.generate_account(&accounts, &accounts, &secrets, &journal, &FixedClock) + .err() + .expect("journal failure must be returned") + .code(), + SafeErrorCode::StorageUnavailable + ); + assert!(accounts.list_accounts().unwrap().is_empty()); + } +} diff --git a/core/crates/studio_application/src/actor.rs b/core/crates/studio_application/src/actor.rs @@ -0,0 +1,787 @@ +use std::num::{NonZeroU64, NonZeroUsize}; +use std::time::Instant; + +use radroots_studio_domain::{ + AccountIdentity, BindingAvailability, LocalSignerBinding, PublicKey, SafeError, SafeErrorCode, + SafeMessage, +}; +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, Debug, Eq, PartialEq)] +pub struct ForegroundSessionBinding { + identity: AccountIdentity, + signer: LocalSignerBinding, + generation: SessionGeneration, +} + +impl ForegroundSessionBinding { + /// Binds one foreground session to a ready local signer and generation. + /// + /// # Errors + /// + /// Returns a safe state error when account and binding differ or when the + /// signer is unavailable. + pub fn new( + identity: AccountIdentity, + signer: LocalSignerBinding, + generation: SessionGeneration, + ) -> Result<Self, SafeError> { + if identity.public_key() != signer.account() + || signer.availability() != BindingAvailability::Available + { + return Err(invalid_foreground_session()); + } + Ok(Self { + identity, + signer, + generation, + }) + } + + #[must_use] + pub const fn identity(&self) -> &AccountIdentity { + &self.identity + } + + #[must_use] + pub const fn signer(&self) -> LocalSignerBinding { + self.signer + } + + #[must_use] + pub const fn generation(&self) -> SessionGeneration { + self.generation + } +} + +const fn invalid_foreground_session() -> SafeError { + SafeError::new( + SafeErrorCode::InvalidApplicationState, + SafeMessage::new("The foreground session binding is invalid."), + ) +} + +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub struct TaskCorrelation { + request_id: RequestId, + account: PublicKey, + binding: LocalSignerBinding, + expected_revision: SnapshotRevision, + session_generation: SessionGeneration, +} + +impl TaskCorrelation { + #[must_use] + pub const fn new( + request_id: RequestId, + account: PublicKey, + binding: LocalSignerBinding, + expected_revision: SnapshotRevision, + session_generation: SessionGeneration, + ) -> Self { + Self { + request_id, + account, + binding, + 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 binding(self) -> LocalSignerBinding { + self.binding + } + + #[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, + CompatibilityChecking, + AcquiringOwnership, + Migrating, + Recovering, + Ready, + Degraded(SafeError), + Blocked(SafeError), + ShuttingDown, + Closed, + Fatal(SafeError), +} + +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum RuntimeCommandClass { + Observe, + MutateLocalState, + UseCredential, + UseRelay, + RetryOpening, + Shutdown, +} + +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub struct LifecycleGate { + lifecycle: RuntimeLifecycle, +} + +impl Default for LifecycleGate { + fn default() -> Self { + Self::opening() + } +} + +impl LifecycleGate { + #[must_use] + pub const fn opening() -> Self { + Self { + lifecycle: RuntimeLifecycle::Opening, + } + } + + #[must_use] + pub const fn lifecycle(self) -> RuntimeLifecycle { + self.lifecycle + } + + #[must_use] + pub const fn allows(self, command: RuntimeCommandClass) -> bool { + match self.lifecycle { + RuntimeLifecycle::Opening + | RuntimeLifecycle::CompatibilityChecking + | RuntimeLifecycle::AcquiringOwnership + | RuntimeLifecycle::Migrating + | RuntimeLifecycle::Recovering => { + matches!( + command, + RuntimeCommandClass::Observe | RuntimeCommandClass::Shutdown + ) + } + RuntimeLifecycle::Ready => !matches!(command, RuntimeCommandClass::RetryOpening), + RuntimeLifecycle::Degraded(_) => !matches!( + command, + RuntimeCommandClass::UseRelay | RuntimeCommandClass::RetryOpening + ), + RuntimeLifecycle::Blocked(_) => matches!( + command, + RuntimeCommandClass::Observe + | RuntimeCommandClass::RetryOpening + | RuntimeCommandClass::Shutdown + ), + RuntimeLifecycle::ShuttingDown => matches!(command, RuntimeCommandClass::Observe), + RuntimeLifecycle::Closed | RuntimeLifecycle::Fatal(_) => false, + } + } + + /// Advances the required open sequence to compatibility checking. + /// + /// # Errors + /// + /// Returns a safe lifecycle error when the stage is out of order. + pub fn begin_compatibility_check(&mut self) -> Result<(), SafeError> { + self.advance( + RuntimeLifecycle::Opening, + RuntimeLifecycle::CompatibilityChecking, + ) + } + + /// Records compatibility acceptance and begins ownership acquisition. + /// + /// # Errors + /// + /// Returns a safe lifecycle error when the stage is out of order. + pub fn compatibility_accepted(&mut self) -> Result<(), SafeError> { + self.advance( + RuntimeLifecycle::CompatibilityChecking, + RuntimeLifecycle::AcquiringOwnership, + ) + } + + /// Records exclusive ownership and begins migration. + /// + /// # Errors + /// + /// Returns a safe lifecycle error when the stage is out of order. + pub fn ownership_acquired(&mut self) -> Result<(), SafeError> { + self.advance( + RuntimeLifecycle::AcquiringOwnership, + RuntimeLifecycle::Migrating, + ) + } + + /// Records migration completion and begins recovery. + /// + /// # Errors + /// + /// Returns a safe lifecycle error when the stage is out of order. + pub fn migration_complete(&mut self) -> Result<(), SafeError> { + self.advance(RuntimeLifecycle::Migrating, RuntimeLifecycle::Recovering) + } + + /// Records recovery completion and admits normal commands. + /// + /// # Errors + /// + /// Returns a safe lifecycle error when the stage is out of order. + pub fn recovery_complete(&mut self) -> Result<(), SafeError> { + self.advance(RuntimeLifecycle::Recovering, RuntimeLifecycle::Ready) + } + + pub fn block(&mut self, error: SafeError) { + self.lifecycle = RuntimeLifecycle::Blocked(error); + } + + pub fn fail(&mut self, error: SafeError) { + self.lifecycle = RuntimeLifecycle::Fatal(error); + } + + /// Moves a ready runtime into a nonfatal degraded state. + /// + /// # Errors + /// + /// Returns a safe lifecycle error when the runtime is not ready. + pub fn degrade(&mut self, error: SafeError) -> Result<(), SafeError> { + self.advance(RuntimeLifecycle::Ready, RuntimeLifecycle::Degraded(error)) + } + + /// Restores local and relay command availability after degradation. + /// + /// # Errors + /// + /// Returns a safe lifecycle error when the runtime is not degraded. + pub fn restore_ready(&mut self) -> Result<(), SafeError> { + if !matches!(self.lifecycle, RuntimeLifecycle::Degraded(_)) { + return Err(invalid_lifecycle_transition()); + } + self.lifecycle = RuntimeLifecycle::Ready; + Ok(()) + } + + /// Begins actor-owned shutdown. + /// + /// # Errors + /// + /// Returns a safe lifecycle error after shutdown or close has begun. + pub fn begin_shutdown(&mut self) -> Result<(), SafeError> { + if matches!( + self.lifecycle, + RuntimeLifecycle::ShuttingDown | RuntimeLifecycle::Closed + ) { + return Err(invalid_lifecycle_transition()); + } + self.lifecycle = RuntimeLifecycle::ShuttingDown; + Ok(()) + } + + /// Completes actor-owned shutdown. + /// + /// # Errors + /// + /// Returns a safe lifecycle error unless shutdown already began. + pub fn finish_shutdown(&mut self) -> Result<(), SafeError> { + self.advance(RuntimeLifecycle::ShuttingDown, RuntimeLifecycle::Closed) + } + + fn advance( + &mut self, + expected: RuntimeLifecycle, + next: RuntimeLifecycle, + ) -> Result<(), SafeError> { + if self.lifecycle != expected { + return Err(invalid_lifecycle_transition()); + } + self.lifecycle = next; + Ok(()) + } +} + +const fn invalid_lifecycle_transition() -> SafeError { + SafeError::new( + radroots_studio_domain::SafeErrorCode::InvalidApplicationState, + radroots_studio_domain::SafeMessage::new("The runtime lifecycle transition is invalid."), + ) +} + +#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)] +pub struct RequestId(NonZeroU64); + +impl RequestId { + #[must_use] + pub const fn new(value: u64) -> Option<Self> { + match NonZeroU64::new(value) { + Some(value) => Some(Self(value)), + None => None, + } + } + + #[must_use] + pub const fn get(self) -> u64 { + self.0.get() + } +} + +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub struct CommandContext { + request_id: RequestId, + expected_revision: Option<SnapshotRevision>, + deadline: Instant, +} + +impl CommandContext { + #[must_use] + pub const fn new( + request_id: RequestId, + expected_revision: Option<SnapshotRevision>, + deadline: Instant, + ) -> Self { + Self { + request_id, + expected_revision, + deadline, + } + } + + #[must_use] + pub const fn request_id(self) -> RequestId { + self.request_id + } + + #[must_use] + pub const fn expected_revision(self) -> Option<SnapshotRevision> { + self.expected_revision + } + + #[must_use] + pub const fn deadline(self) -> Instant { + self.deadline + } + + #[must_use] + pub fn is_expired(self, now: Instant) -> bool { + now >= self.deadline + } +} + +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum CommandRejection { + MailboxSaturated, +} + +#[derive(Clone, Debug, Eq, PartialEq)] +pub enum CommandResult<T> { + Completed(T), + Rejected(CommandRejection), + Conflicted { current_revision: SnapshotRevision }, + TimedOut, + Closed, + Failed(SafeError), +} + +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct CommandReceipt<T> { + request_id: RequestId, + result: CommandResult<T>, +} + +impl<T> CommandReceipt<T> { + #[must_use] + pub const fn new(request_id: RequestId, result: CommandResult<T>) -> Self { + Self { request_id, result } + } + + #[must_use] + pub const fn request_id(&self) -> RequestId { + self.request_id + } + + #[must_use] + pub const fn result(&self) -> &CommandResult<T> { + &self.result + } + + #[must_use] + pub fn into_result(self) -> CommandResult<T> { + self.result + } +} + +pub struct CommandTicket<T> { + request_id: RequestId, + receiver: oneshot::Receiver<CommandReceipt<T>>, +} + +impl<T> CommandTicket<T> { + #[must_use] + pub const fn request_id(&self) -> RequestId { + self.request_id + } + + pub async fn receipt(self) -> CommandReceipt<T> { + self.receiver + .await + .unwrap_or_else(|_| CommandReceipt::new(self.request_id, CommandResult::Closed)) + } +} + +pub enum CommandSubmission<T> { + Accepted(CommandTicket<T>), + Rejected(CommandReceipt<T>), +} + +impl<T> CommandSubmission<T> { + #[must_use] + pub const fn request_id(&self) -> RequestId { + match self { + Self::Accepted(ticket) => ticket.request_id(), + Self::Rejected(receipt) => receipt.request_id(), + } + } +} + +pub struct CommandEnvelope<C, R> { + context: CommandContext, + command: C, + reply: oneshot::Sender<CommandReceipt<R>>, +} + +impl<C, R> CommandEnvelope<C, R> { + #[must_use] + pub const fn context(&self) -> CommandContext { + self.context + } + + #[must_use] + pub const fn command(&self) -> &C { + &self.command + } + + #[must_use] + pub fn into_parts(self) -> (CommandContext, C, oneshot::Sender<CommandReceipt<R>>) { + (self.context, self.command, self.reply) + } +} + +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>>) { + let (sender, receiver) = mpsc::channel(capacity.get()); + (Self { sender }, receiver) + } + + #[must_use] + pub fn available_capacity(&self) -> usize { + self.sender.capacity() + } + + #[must_use] + pub fn submit(&self, context: CommandContext, command: C) -> CommandSubmission<R> { + let request_id = context.request_id(); + if context.is_expired(Instant::now()) { + return CommandSubmission::Rejected(CommandReceipt::new( + request_id, + CommandResult::TimedOut, + )); + } + let (reply, receiver) = oneshot::channel(); + let envelope = CommandEnvelope { + context, + command, + reply, + }; + match self.sender.try_send(envelope) { + Ok(()) => CommandSubmission::Accepted(CommandTicket { + request_id, + receiver, + }), + Err(mpsc::error::TrySendError::Full(_)) => { + CommandSubmission::Rejected(CommandReceipt::new( + request_id, + CommandResult::Rejected(CommandRejection::MailboxSaturated), + )) + } + Err(mpsc::error::TrySendError::Closed(_)) => { + CommandSubmission::Rejected(CommandReceipt::new(request_id, CommandResult::Closed)) + } + } + } +} + +#[cfg(test)] +mod tests { + use std::num::NonZeroUsize; + use std::time::{Duration, Instant}; + + use radroots_studio_domain::{AccountIdentity, BindingAvailability, LocalSignerBinding}; + + use crate::{ + ActorMailbox, CommandContext, CommandReceipt, CommandRejection, CommandResult, + CommandSubmission, ForegroundSessionBinding, LifecycleGate, RequestId, RuntimeCommandClass, + RuntimeLifecycle, SessionGeneration, + }; + + fn context(id: u64) -> CommandContext { + CommandContext::new( + RequestId::new(id).expect("nonzero request"), + None, + Instant::now() + Duration::from_secs(1), + ) + } + + #[test] + fn foreground_session_requires_matching_available_binding_and_generation() { + let public_key = crate::test_support::valid_test_public_key(3).expect("valid public key"); + let identity = AccountIdentity::derive(public_key).expect("identity"); + let generation = SessionGeneration::from_value(4); + let session = ForegroundSessionBinding::new( + identity.clone(), + LocalSignerBinding::new(public_key, BindingAvailability::Available), + generation, + ) + .expect("session"); + assert_eq!(session.identity(), &identity); + assert_eq!(session.signer().account(), public_key); + assert_eq!(session.generation(), generation); + + let correlation = super::TaskCorrelation::new( + RequestId::new(8).expect("request"), + public_key, + session.signer(), + crate::SnapshotRevision::from_value(9), + generation, + ); + assert_eq!(correlation.request_id().get(), 8); + assert_eq!(correlation.account(), public_key); + assert_eq!(correlation.binding(), session.signer()); + assert_eq!(correlation.expected_revision().value(), 9); + assert_eq!(correlation.session_generation(), generation); + + assert!( + ForegroundSessionBinding::new( + identity.clone(), + LocalSignerBinding::new( + radroots_studio_domain::PublicKey::from_hex( + "e0266e3cfb0d2886f91c73f5f868f3b98273713e5fcd97c081663f5518a4b3af", + ) + .expect("different valid public key"), + BindingAvailability::Available, + ), + generation, + ) + .is_err() + ); + assert!( + ForegroundSessionBinding::new( + identity, + LocalSignerBinding::new(public_key, BindingAvailability::CredentialMissing), + generation, + ) + .is_err() + ); + } + + #[tokio::test] + async fn bounded_mailbox_accepts_one_and_rejects_saturation() { + let (mailbox, mut receiver) = + ActorMailbox::<u8, u8>::bounded(NonZeroUsize::new(1).expect("capacity")); + let CommandSubmission::Accepted(ticket) = mailbox.submit(context(1), 7) else { + panic!("first command must be accepted"); + }; + let CommandSubmission::Rejected(rejected) = mailbox.submit(context(2), 8) else { + panic!("second command must be rejected"); + }; + assert_eq!( + rejected.into_result(), + CommandResult::Rejected(CommandRejection::MailboxSaturated) + ); + + let envelope = receiver.recv().await.expect("command"); + assert_eq!(envelope.context().request_id().get(), 1); + assert_eq!(*envelope.command(), 7); + let (context, command, reply) = envelope.into_parts(); + reply + .send(CommandReceipt::new( + context.request_id(), + CommandResult::Completed(command + 1), + )) + .expect("ticket remains open"); + assert_eq!( + ticket.receipt().await.into_result(), + CommandResult::Completed(8) + ); + } + + #[test] + fn expired_and_closed_mailboxes_reject_without_enqueuing() { + let (mailbox, receiver) = + ActorMailbox::<u8, u8>::bounded(NonZeroUsize::new(1).expect("capacity")); + let expired = + CommandContext::new(RequestId::new(1).expect("request"), None, Instant::now()); + let CommandSubmission::Rejected(receipt) = mailbox.submit(expired, 1) else { + panic!("expired command must be rejected"); + }; + assert_eq!(receipt.into_result(), CommandResult::TimedOut); + + drop(receiver); + let CommandSubmission::Rejected(receipt) = mailbox.submit(context(2), 2) else { + panic!("closed mailbox must be rejected"); + }; + assert_eq!(receipt.into_result(), CommandResult::Closed); + } + + #[tokio::test] + async fn dropped_actor_reply_becomes_closed_receipt() { + let (mailbox, mut receiver) = + ActorMailbox::<u8, u8>::bounded(NonZeroUsize::new(1).expect("capacity")); + let CommandSubmission::Accepted(ticket) = mailbox.submit(context(1), 1) else { + panic!("command must be accepted"); + }; + drop(receiver.recv().await.expect("command")); + + assert_eq!(ticket.receipt().await.into_result(), CommandResult::Closed); + } + + #[test] + fn opening_sequence_gates_mutation_until_recovery_completes() { + let mut lifecycle = LifecycleGate::opening(); + for expected in [ + RuntimeLifecycle::Opening, + RuntimeLifecycle::CompatibilityChecking, + RuntimeLifecycle::AcquiringOwnership, + RuntimeLifecycle::Migrating, + RuntimeLifecycle::Recovering, + ] { + assert_eq!(lifecycle.lifecycle(), expected); + assert!(lifecycle.allows(RuntimeCommandClass::Observe)); + assert!(lifecycle.allows(RuntimeCommandClass::Shutdown)); + assert!(!lifecycle.allows(RuntimeCommandClass::MutateLocalState)); + match expected { + RuntimeLifecycle::Opening => { + lifecycle + .begin_compatibility_check() + .expect("compatibility"); + } + RuntimeLifecycle::CompatibilityChecking => { + lifecycle + .compatibility_accepted() + .expect("compatibility accepted"); + } + RuntimeLifecycle::AcquiringOwnership => { + lifecycle.ownership_acquired().expect("ownership"); + } + RuntimeLifecycle::Migrating => { + lifecycle.migration_complete().expect("migration"); + } + RuntimeLifecycle::Recovering => { + lifecycle.recovery_complete().expect("recovery"); + } + _ => unreachable!("opening states only"), + } + } + assert_eq!(lifecycle.lifecycle(), RuntimeLifecycle::Ready); + assert!(lifecycle.allows(RuntimeCommandClass::MutateLocalState)); + assert!(lifecycle.allows(RuntimeCommandClass::UseCredential)); + assert!(lifecycle.allows(RuntimeCommandClass::UseRelay)); + } + + #[test] + fn blocked_degraded_fatal_and_closed_states_fail_safe() { + let problem = radroots_studio_domain::SafeError::new( + radroots_studio_domain::SafeErrorCode::StorageUnavailable, + radroots_studio_domain::SafeMessage::new("The runtime is unavailable."), + ); + let mut blocked = LifecycleGate::opening(); + blocked.block(problem); + assert!(blocked.allows(RuntimeCommandClass::RetryOpening)); + assert!(!blocked.allows(RuntimeCommandClass::MutateLocalState)); + + let mut degraded = LifecycleGate::opening(); + degraded.begin_compatibility_check().expect("compatibility"); + degraded.compatibility_accepted().expect("accepted"); + degraded.ownership_acquired().expect("ownership"); + degraded.migration_complete().expect("migration"); + degraded.recovery_complete().expect("recovery"); + degraded.degrade(problem).expect("degraded"); + assert!(degraded.allows(RuntimeCommandClass::MutateLocalState)); + assert!(!degraded.allows(RuntimeCommandClass::UseRelay)); + degraded.restore_ready().expect("restored"); + assert!(LifecycleGate::opening().restore_ready().is_err()); + + let mut fatal = LifecycleGate::opening(); + fatal.fail(problem); + assert!(!fatal.allows(RuntimeCommandClass::Observe)); + fatal.begin_shutdown().expect("fatal can close"); + fatal.finish_shutdown().expect("closed"); + assert_eq!(fatal.lifecycle(), RuntimeLifecycle::Closed); + assert!(!fatal.allows(RuntimeCommandClass::Shutdown)); + } + + #[test] + fn opening_stages_reject_out_of_order_and_repeated_transitions() { + let mut lifecycle = LifecycleGate::opening(); + assert!(lifecycle.migration_complete().is_err()); + lifecycle.begin_compatibility_check().expect("first stage"); + assert!(lifecycle.begin_compatibility_check().is_err()); + assert!(lifecycle.recovery_complete().is_err()); + } +} diff --git a/core/crates/studio_application/src/app_core.rs b/core/crates/studio_application/src/app_core.rs @@ -0,0 +1,358 @@ +use std::collections::BTreeMap; +use std::sync::{Arc, Mutex, MutexGuard}; + +use radroots_studio_domain::{PublicKey, SafeError, SafeErrorCode, SafeMessage, UnixTimestamp}; + +use crate::{ + AccountRepository, AppSnapshot, AppStateRepository, KeyMaterialProvider, RelayConfiguration, + SnapshotRevision, StateMachine, StateTransition, +}; + +pub struct RemovalConfirmationToken { + id: u64, + public_key: PublicKey, + revision: SnapshotRevision, + expires_at: UnixTimestamp, + impact: RemovalImpact, +} + +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub struct RemovalImpact { + deletes_local_credential: bool, + signs_out: bool, +} + +impl RemovalImpact { + #[must_use] + pub const fn deletes_local_credential(self) -> bool { + self.deletes_local_credential + } + #[must_use] + pub const fn signs_out(self) -> bool { + self.signs_out + } +} + +impl RemovalConfirmationToken { + #[must_use] + pub const fn public_key(&self) -> PublicKey { + self.public_key + } + #[must_use] + pub const fn revision(&self) -> SnapshotRevision { + self.revision + } + #[must_use] + pub const fn expires_at(&self) -> UnixTimestamp { + self.expires_at + } + #[must_use] + pub const fn impact(&self) -> RemovalImpact { + self.impact + } +} + +#[derive(Clone, Copy)] +struct RemovalTokenState { + public_key: PublicKey, + revision: SnapshotRevision, + expires_at: UnixTimestamp, + impact: RemovalImpact, +} + +struct CoreState { + state_machine: StateMachine, + removal_tokens: BTreeMap<u64, RemovalTokenState>, + next_removal_token: u64, +} + +pub struct AppCore { + relay_configuration: RelayConfiguration, + key_material: Arc<dyn KeyMaterialProvider>, + state: Mutex<CoreState>, +} + +impl AppCore { + #[must_use] + pub fn new( + relay_configuration: RelayConfiguration, + key_material: Arc<dyn KeyMaterialProvider>, + ) -> Self { + Self { + relay_configuration, + key_material, + state: Mutex::new(CoreState { + state_machine: StateMachine::booting(), + removal_tokens: BTreeMap::new(), + next_removal_token: 1, + }), + } + } + + #[cfg(test)] + #[must_use] + pub fn in_memory(relay_configuration: RelayConfiguration) -> Self { + Self::new( + relay_configuration, + Arc::new(crate::test_support::TestKeyMaterialProvider::default()), + ) + } + + pub(crate) fn key_material(&self) -> &dyn KeyMaterialProvider { + self.key_material.as_ref() + } + + /// Moves the in-memory core from booting to an empty ready snapshot. + /// + /// # Errors + /// + /// Returns a safe application-state error if the ready snapshot invariant + /// cannot be constructed. + pub fn bootstrap(&self) -> Result<AppSnapshot, SafeError> { + self.apply_transition(StateTransition::Bootstrap) + } + + /// Loads the durable public registry and selection into a signed-out snapshot. + /// + /// # Errors + /// + /// Returns the safe persistence error after publishing a fatal snapshot when + /// durable state cannot be read or violates application invariants. + pub fn bootstrap_from( + &self, + accounts: &(impl AccountRepository + ?Sized), + app_state: &(impl AppStateRepository + ?Sized), + ) -> Result<AppSnapshot, SafeError> { + let loaded = accounts.list_accounts().and_then(|accounts| { + app_state + .load_selected_account() + .map(|selected| (accounts, selected)) + }); + match loaded { + Ok((accounts, selected)) => { + self.apply_transition(StateTransition::BootstrapRegistry { accounts, selected }) + } + Err(error) => { + self.apply_transition(StateTransition::Fatal(error))?; + Err(error) + } + } + } + + #[must_use] + pub fn snapshot(&self) -> AppSnapshot { + self.lock_state().state_machine.snapshot().clone() + } + + pub(crate) fn apply_transition( + &self, + transition: StateTransition, + ) -> Result<AppSnapshot, SafeError> { + self.lock_state() + .state_machine + .apply(transition, &self.relay_configuration) + } + + pub(crate) fn issue_removal_token( + &self, + public_key: PublicKey, + now: UnixTimestamp, + ) -> Result<RemovalConfirmationToken, SafeError> { + let mut state = self.lock_state(); + let Some(account) = state + .state_machine + .snapshot() + .accounts() + .iter() + .find(|account| account.public_key() == public_key) + else { + return Err(account_not_found()); + }; + let deletes_local_credential = account.signer().availability() + != radroots_studio_domain::BindingAvailability::CredentialMissing; + let id = state.next_removal_token; + state.next_removal_token = id.checked_add(1).ok_or_else(invalid_application_state)?; + let revision = state.state_machine.snapshot().revision(); + let expires_at = UnixTimestamp::from_seconds( + now.as_seconds() + .checked_add(300) + .ok_or_else(invalid_application_state)?, + ) + .ok_or_else(invalid_application_state)?; + let impact = RemovalImpact { + deletes_local_credential, + signs_out: state + .state_machine + .snapshot() + .active_account() + .is_some_and(|active| active.account().public_key() == public_key), + }; + state.removal_tokens.insert( + id, + RemovalTokenState { + public_key, + revision, + expires_at, + impact, + }, + ); + Ok(RemovalConfirmationToken { + id, + public_key, + revision, + expires_at, + impact, + }) + } + + #[allow(clippy::needless_pass_by_value)] + pub(crate) fn consume_removal_token( + &self, + token: RemovalConfirmationToken, + now: UnixTimestamp, + ) -> Result<PublicKey, SafeError> { + let RemovalConfirmationToken { + id, + public_key, + revision, + expires_at, + impact, + } = token; + let mut state = self.lock_state(); + let stored = state.removal_tokens.remove(&id); + if stored.is_none_or(|stored| { + stored.public_key != public_key + || stored.revision != revision + || stored.expires_at != expires_at + || stored.impact != impact + }) || state.state_machine.snapshot().revision() != revision + || now.as_seconds() > expires_at.as_seconds() + { + return Err(invalid_application_state()); + } + Ok(public_key) + } + + #[allow(clippy::needless_pass_by_value)] + pub(crate) fn cancel_removal_token(&self, token: RemovalConfirmationToken) -> bool { + self.lock_state().removal_tokens.remove(&token.id).is_some() + } + + fn lock_state(&self) -> MutexGuard<'_, CoreState> { + self.state + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + } +} + +const fn invalid_application_state() -> SafeError { + SafeError::new( + SafeErrorCode::InvalidApplicationState, + SafeMessage::new("The account removal confirmation is no longer valid."), + ) +} + +const fn account_not_found() -> SafeError { + SafeError::new( + SafeErrorCode::AccountNotFound, + SafeMessage::new("The account was not found."), + ) +} + +#[cfg(test)] +mod tests { + use radroots_studio_domain::{ + AccountCreatedAt, AccountIdentity, AccountSummary, BindingAvailability, LocalSignerBinding, + UnixTimestamp, + }; + + use crate::{AppCore, AppLifecycle, RelayConfiguration, StateTransition}; + + #[test] + fn bootstrap_is_idempotent_and_advances_only_once() { + let core = AppCore::in_memory(RelayConfiguration::default()); + let ready = core.bootstrap().expect("bootstrap"); + let repeated = core.bootstrap().expect("idempotent bootstrap"); + + assert_eq!(ready.lifecycle(), AppLifecycle::Ready); + assert_eq!(ready.revision().value(), 1); + assert_eq!(repeated, ready); + } + + #[test] + fn core_instances_never_share_state() { + let first = AppCore::in_memory(RelayConfiguration::default()); + let second = AppCore::in_memory(RelayConfiguration::default()); + + first.bootstrap().expect("first bootstrap"); + + assert_eq!(first.snapshot().revision().value(), 1); + assert_eq!(second.snapshot().revision().value(), 0); + } + + #[test] + fn removal_impact_matches_missing_local_binding() { + let core = AppCore::in_memory(RelayConfiguration::default()); + let public_key = crate::test_support::valid_test_public_key(9).expect("valid public key"); + let account = AccountSummary::new( + AccountIdentity::derive(public_key).expect("identity"), + LocalSignerBinding::new(public_key, BindingAvailability::CredentialMissing), + None, + AccountCreatedAt::new(UnixTimestamp::from_seconds(1).expect("time")), + None, + ) + .expect("account"); + core.apply_transition(StateTransition::BootstrapRegistry { + accounts: vec![account], + selected: Some(public_key), + }) + .expect("registry"); + + let removal = core + .issue_removal_token(public_key, UnixTimestamp::from_seconds(2).expect("time")) + .expect("removal"); + + assert!(!removal.impact().deletes_local_credential()); + assert!(!removal.impact().signs_out()); + } + + #[test] + fn removal_confirmation_rejects_every_tampered_authority_field() { + let core = AppCore::in_memory(RelayConfiguration::default()); + let public_key = crate::test_support::valid_test_public_key(7).expect("public key"); + let other_key = crate::test_support::valid_test_public_key(8).expect("other key"); + let account = AccountSummary::new( + AccountIdentity::derive(public_key).expect("identity"), + LocalSignerBinding::new(public_key, BindingAvailability::Available), + None, + AccountCreatedAt::new(UnixTimestamp::from_seconds(1).expect("time")), + None, + ) + .expect("account"); + core.apply_transition(StateTransition::BootstrapRegistry { + accounts: vec![account], + selected: Some(public_key), + }) + .expect("registry"); + + let now = UnixTimestamp::from_seconds(2).expect("now"); + let mut wrong_key = core.issue_removal_token(public_key, now).expect("token"); + wrong_key.public_key = other_key; + assert!(core.consume_removal_token(wrong_key, now).is_err()); + + let mut wrong_revision = core.issue_removal_token(public_key, now).expect("token"); + wrong_revision.revision = crate::SnapshotRevision::initial(); + assert!(core.consume_removal_token(wrong_revision, now).is_err()); + + let mut wrong_expiry = core.issue_removal_token(public_key, now).expect("token"); + wrong_expiry.expires_at = UnixTimestamp::from_seconds(999).expect("expiry"); + assert!(core.consume_removal_token(wrong_expiry, now).is_err()); + + let mut wrong_impact = core.issue_removal_token(public_key, now).expect("token"); + wrong_impact.impact = super::RemovalImpact { + deletes_local_credential: false, + signs_out: false, + }; + assert!(core.consume_removal_token(wrong_impact, now).is_err()); + } +} diff --git a/core/crates/studio_application/src/change_stream.rs b/core/crates/studio_application/src/change_stream.rs @@ -0,0 +1,239 @@ +use std::collections::BTreeMap; +use std::num::{NonZeroU64, NonZeroUsize}; + +use tokio::sync::mpsc; + +use crate::{AppSnapshot, SnapshotRevision}; + +#[derive(Clone, Copy, Debug, Eq, Ord, PartialEq, PartialOrd)] +pub struct ChangeSubscriptionId(NonZeroU64); + +impl ChangeSubscriptionId { + #[must_use] + pub const fn value(self) -> u64 { + self.0.get() + } +} + +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct SnapshotChange { + snapshot: AppSnapshot, + previous_revision: Option<SnapshotRevision>, +} + +impl SnapshotChange { + #[must_use] + pub const fn revision(&self) -> SnapshotRevision { + self.snapshot.revision() + } + + #[must_use] + pub const fn snapshot(&self) -> &AppSnapshot { + &self.snapshot + } + + #[must_use] + pub fn into_snapshot(self) -> AppSnapshot { + self.snapshot + } + + #[must_use] + pub const fn previous_revision(&self) -> Option<SnapshotRevision> { + self.previous_revision + } + + #[must_use] + pub fn recovers_gap_after(&self, observed: SnapshotRevision) -> bool { + self.previous_revision + .is_some_and(|previous| previous != observed) + } +} + +pub struct SnapshotChangeReceiver { + receiver: mpsc::Receiver<SnapshotChange>, +} + +impl SnapshotChangeReceiver { + pub async fn receive(&mut self) -> Option<SnapshotChange> { + self.receiver.recv().await + } +} + +pub struct OrderedSnapshotChanges { + latest: AppSnapshot, + next_subscription: u64, + subscribers: BTreeMap<ChangeSubscriptionId, mpsc::Sender<SnapshotChange>>, + closed: bool, +} + +impl OrderedSnapshotChanges { + #[must_use] + pub fn new(initial_snapshot: AppSnapshot) -> Self { + Self { + latest: initial_snapshot, + next_subscription: 1, + subscribers: BTreeMap::new(), + closed: false, + } + } + + #[must_use] + pub const fn last_revision(&self) -> SnapshotRevision { + self.latest.revision() + } + + /// Registers a bounded consumer for future changes. + /// + /// # Errors + /// + /// Returns `None` if the subscription identifier space is exhausted. + pub fn subscribe( + &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()); + sender + .try_send(SnapshotChange { + snapshot: self.latest.clone(), + previous_revision: None, + }) + .ok()?; + self.subscribers.insert(id, sender); + Some((id, SnapshotChangeReceiver { receiver })) + } + + #[must_use] + pub fn unsubscribe(&mut self, id: ChangeSubscriptionId) -> bool { + self.subscribers.remove(&id).is_some() + } + + pub fn publish(&mut self, snapshot: AppSnapshot) { + if self.closed || snapshot.revision() <= self.latest.revision() { + return; + } + let change = SnapshotChange { + previous_revision: Some(self.latest.revision()), + snapshot, + }; + self.latest = change.snapshot.clone(); + self.subscribers + .retain(|_, sender| match sender.try_send(change.clone()) { + Ok(()) | Err(mpsc::error::TrySendError::Full(_)) => true, + Err(mpsc::error::TrySendError::Closed(_)) => false, + }); + } + + pub fn close(&mut self) { + self.closed = true; + self.subscribers.clear(); + } +} + +#[cfg(test)] +mod tests { + use std::num::NonZeroUsize; + + use crate::{ + AppSnapshot, OrderedSnapshotChanges, RelayConfiguration, SessionState, SnapshotRevision, + }; + + #[tokio::test] + async fn change_stream_publishes_monotonic_revisions_to_multiple_consumers() { + let mut changes = OrderedSnapshotChanges::new(snapshot(0)); + let (_, mut first) = changes + .subscribe(NonZeroUsize::new(4).expect("capacity")) + .expect("first subscription"); + let (_, mut second) = changes + .subscribe(NonZeroUsize::new(4).expect("capacity")) + .expect("second subscription"); + + assert_eq!( + first.receive().await.expect("initial").revision(), + revision(0) + ); + assert_eq!( + second.receive().await.expect("initial").revision(), + revision(0) + ); + changes.publish(snapshot(1)); + changes.publish(snapshot(1)); + changes.publish(snapshot(2)); + + for receiver in [&mut first, &mut second] { + assert_eq!( + receiver.receive().await.expect("revision 1").revision(), + revision(1) + ); + assert_eq!( + receiver.receive().await.expect("revision 2").revision(), + revision(2) + ); + } + assert_eq!(changes.last_revision(), revision(2)); + } + + #[tokio::test] + async fn slow_consumers_expose_a_revision_gap_without_blocking_publication() { + let mut changes = OrderedSnapshotChanges::new(snapshot(0)); + let (_, mut receiver) = changes + .subscribe(NonZeroUsize::new(1).expect("capacity")) + .expect("subscription"); + + assert_eq!( + receiver.receive().await.expect("initial").revision(), + revision(0) + ); + changes.publish(snapshot(1)); + changes.publish(snapshot(2)); + let first = receiver.receive().await.expect("first"); + assert_eq!(first.revision(), revision(1)); + changes.publish(snapshot(3)); + let recovered = receiver.receive().await.expect("gap recovery"); + assert_eq!(recovered.revision(), revision(3)); + 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) + } + + fn snapshot(value: u64) -> AppSnapshot { + if value == 0 { + AppSnapshot::booting() + } else { + AppSnapshot::ready( + revision(value), + RelayConfiguration::default(), + Vec::new(), + None, + SessionState::SignedOut, + None, + None, + ) + .expect("snapshot") + } + } +} diff --git a/core/crates/studio_application/src/config.rs b/core/crates/studio_application/src/config.rs @@ -0,0 +1,107 @@ +use radroots_studio_domain::{ + RelayDestinationPolicy, SafeError, SafeErrorCode, SafeMessage, normalize_relay_urls, +}; + +use crate::RelayConfiguration; + +pub const RELAY_ENVIRONMENT_VARIABLE: &str = "RADROOTS_NOSTR_RELAYS"; +const DEVELOPMENT_RELAY: &str = "ws://localhost:8080"; + +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum RelayRuntimeMode { + Development, + Packaged, +} + +/// Reads the process relay configuration once through the Rust-owned boundary. +/// +/// # Errors +/// +/// Returns a safe configuration error for missing Unicode or invalid relay data. +pub fn relay_configuration_from_environment( + mode: RelayRuntimeMode, +) -> Result<RelayConfiguration, SafeError> { + let value = match std::env::var(RELAY_ENVIRONMENT_VARIABLE) { + Ok(value) => Some(value), + Err(std::env::VarError::NotPresent) => None, + Err(std::env::VarError::NotUnicode(_)) => return Err(invalid_configuration()), + }; + relay_configuration_from_value(value.as_deref(), mode) +} + +/// Parses an injected comma-separated relay list without mutating process state. +/// +/// # Errors +/// +/// Returns a safe configuration error when an entry is invalid or packaged mode +/// has no configured relay. +pub fn relay_configuration_from_value( + value: Option<&str>, + mode: RelayRuntimeMode, +) -> Result<RelayConfiguration, SafeError> { + let configured = value.unwrap_or_default().trim(); + let (source, policy) = if configured.is_empty() { + match mode { + RelayRuntimeMode::Development => (DEVELOPMENT_RELAY, RelayDestinationPolicy::Local), + RelayRuntimeMode::Packaged => return Err(invalid_configuration()), + } + } else { + (configured, RelayDestinationPolicy::Public) + }; + let normalized = normalize_relay_urls(source.split(',').map(str::trim), policy)?; + if normalized.is_empty() { + return Err(invalid_configuration()); + } + RelayConfiguration::new(normalized) +} + +const fn invalid_configuration() -> SafeError { + SafeError::new( + SafeErrorCode::InvalidRelayConfiguration, + SafeMessage::new("The Nostr relay configuration is invalid."), + ) +} + +#[cfg(test)] +mod tests { + use radroots_studio_domain::SafeErrorCode; + + use super::{RelayRuntimeMode, relay_configuration_from_value}; + + #[test] + fn relay_config_uses_localhost_fallback_only_for_development() { + for value in [None, Some(""), Some(" ")] { + let development = relay_configuration_from_value(value, RelayRuntimeMode::Development) + .expect("development fallback"); + assert_eq!(development.relays()[0].as_str(), "ws://localhost:8080/"); + let packaged = relay_configuration_from_value(value, RelayRuntimeMode::Packaged) + .expect_err("packaged configuration required"); + assert_eq!(packaged.code(), SafeErrorCode::InvalidRelayConfiguration); + } + } + + #[test] + fn relay_config_trims_deduplicates_and_preserves_order() { + let configuration = relay_configuration_from_value( + Some(" wss://relay.one ,wss://relay.two,wss://relay.one/ "), + RelayRuntimeMode::Packaged, + ) + .expect("configuration"); + let relays = configuration + .relays() + .iter() + .map(radroots_studio_domain::RelayUrl::as_str) + .collect::<Vec<_>>(); + assert_eq!(relays, ["wss://relay.one/", "wss://relay.two/"]); + } + + #[test] + fn relay_config_rejects_any_invalid_comma_separated_entry() { + let error = relay_configuration_from_value( + Some("wss://relay.one,https://not-a-relay.test"), + RelayRuntimeMode::Packaged, + ) + .expect_err("invalid entry"); + assert_eq!(error.code(), SafeErrorCode::InvalidRelayConfiguration); + } +} diff --git a/core/crates/studio_application/src/custody.rs b/core/crates/studio_application/src/custody.rs @@ -0,0 +1,288 @@ +use std::num::NonZeroU64; +use std::sync::Mutex; +use std::time::Duration; + +use crate::KeyMaterialProvider; +use radroots_studio_domain::{ + AccountCreatedAt, AccountIdentity, AccountSummary, BindingAvailability, LocalSignerBinding, + Nsec, SafeError, SafeErrorCode, SafeMessage, SecretKeyInput, UnixTimestamp, +}; +pub const GENERATED_KEY_STAGE_TTL: Duration = Duration::from_mins(5); + +#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)] +pub struct RecoveryStageId(NonZeroU64); + +impl RecoveryStageId { + #[must_use] + pub const fn new(value: NonZeroU64) -> Self { + Self(value) + } + + #[must_use] + pub const fn value(self) -> u64 { + self.0.get() + } +} + +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct GeneratedKeyStageView { + account: AccountSummary, + expires_at: UnixTimestamp, +} + +impl GeneratedKeyStageView { + #[must_use] + pub const fn account(&self) -> &AccountSummary { + &self.account + } + + #[must_use] + pub const fn expires_at(&self) -> UnixTimestamp { + self.expires_at + } +} + +pub struct StagedGeneratedKey { + id: RecoveryStageId, + account: AccountSummary, + secret: SecretKeyInput, + expected_revision: u64, + expires_at: UnixTimestamp, +} + +impl StagedGeneratedKey { + #[must_use] + pub fn view(&self) -> GeneratedKeyStageView { + GeneratedKeyStageView { + account: self.account.clone(), + expires_at: self.expires_at, + } + } + + #[must_use] + pub const fn expected_revision(&self) -> u64 { + self.expected_revision + } + + #[must_use] + pub const fn id(&self) -> RecoveryStageId { + self.id + } + + #[must_use] + pub const fn account(&self) -> &AccountSummary { + &self.account + } + + #[must_use] + pub fn into_commit_parts(self) -> (AccountSummary, SecretKeyInput) { + (self.account, self.secret) + } +} + +pub struct GeneratedKeyRecoveryHandle { + id: RecoveryStageId, + view: GeneratedKeyStageView, + recovery_nsec: Mutex<Option<Nsec>>, +} + +impl GeneratedKeyRecoveryHandle { + fn new(id: RecoveryStageId, view: GeneratedKeyStageView, recovery_nsec: Nsec) -> Self { + Self { + id, + view, + recovery_nsec: Mutex::new(Some(recovery_nsec)), + } + } + + #[must_use] + pub const fn id(&self) -> RecoveryStageId { + self.id + } + + #[must_use] + pub const fn view(&self) -> &GeneratedKeyStageView { + &self.view + } + + /// Returns the generated recovery value exactly once. + /// + /// # Errors + /// + /// Returns a safe unavailable error after the value was already consumed. + pub fn take_recovery_nsec(&self) -> Result<Nsec, SafeError> { + self.recovery_nsec + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .take() + .ok_or_else(recovery_not_available) + } +} + +#[derive(Default)] +pub struct GeneratedKeyStage { + pending: Option<StagedGeneratedKey>, +} + +impl GeneratedKeyStage { + /// Replaces an expired stage or creates the only active generated-key stage. + /// + /// # Errors + /// + /// Returns a safe conflict while an unexpired recovery stage is active. + pub fn begin( + &mut self, + key_material: &dyn KeyMaterialProvider, + id: RecoveryStageId, + expected_revision: u64, + now: UnixTimestamp, + ) -> Result<GeneratedKeyRecoveryHandle, SafeError> { + self.expire(now); + if self.pending.is_some() { + return Err(recovery_in_progress()); + } + let generated = key_material.generate()?; + let (public_key, npub, secret, recovery_nsec) = generated.into_parts(); + let account = AccountSummary::new( + AccountIdentity::verify(public_key, npub.as_str().to_owned())?, + LocalSignerBinding::new(public_key, BindingAvailability::Available), + None, + AccountCreatedAt::new(now), + None, + )?; + let ttl = + i64::try_from(GENERATED_KEY_STAGE_TTL.as_secs()).map_err(|_| invalid_stage_expiry())?; + let expires_at = now + .as_seconds() + .checked_add(ttl) + .and_then(UnixTimestamp::from_seconds) + .ok_or_else(invalid_stage_expiry)?; + let pending = StagedGeneratedKey { + id, + account, + secret, + expected_revision, + expires_at, + }; + let view = pending.view(); + self.pending = Some(pending); + Ok(GeneratedKeyRecoveryHandle::new(id, view, recovery_nsec)) + } + + pub fn cancel(&mut self) -> bool { + self.pending.take().is_some() + } + + pub fn expire(&mut self, now: UnixTimestamp) -> bool { + if self + .pending + .as_ref() + .is_some_and(|pending| now >= pending.expires_at) + { + self.pending = None; + true + } else { + false + } + } + + #[must_use] + pub const fn pending(&self) -> Option<&StagedGeneratedKey> { + self.pending.as_ref() + } + + /// Consumes the active, unexpired stage for its commit boundary. + /// + /// # Errors + /// + /// Returns a safe unavailable error when no live stage remains. + pub fn take( + &mut self, + id: RecoveryStageId, + now: UnixTimestamp, + ) -> Result<StagedGeneratedKey, SafeError> { + self.expire(now); + if self.pending.as_ref().map(StagedGeneratedKey::id) != Some(id) { + return Err(recovery_not_available()); + } + self.pending.take().ok_or_else(recovery_not_available) + } +} + +const fn recovery_in_progress() -> SafeError { + SafeError::new( + SafeErrorCode::InvalidApplicationState, + SafeMessage::new("A generated-key recovery step is already in progress."), + ) +} + +const fn recovery_not_available() -> SafeError { + SafeError::new( + SafeErrorCode::InvalidApplicationState, + SafeMessage::new("The generated-key recovery step is no longer available."), + ) +} + +const fn invalid_stage_expiry() -> SafeError { + SafeError::new( + SafeErrorCode::InvalidApplicationState, + SafeMessage::new("The generated-key recovery expiry is invalid."), + ) +} + +#[cfg(test)] +mod tests { + use std::num::NonZeroU64; + + use radroots_studio_domain::UnixTimestamp; + + use super::{GENERATED_KEY_STAGE_TTL, GeneratedKeyStage, RecoveryStageId}; + use crate::test_support::TestKeyMaterialProvider; + + fn time(seconds: i64) -> UnixTimestamp { + UnixTimestamp::from_seconds(seconds).expect("time") + } + + fn id(value: u64) -> RecoveryStageId { + RecoveryStageId::new(NonZeroU64::new(value).expect("id")) + } + + #[test] + fn stage_is_exclusive_cancelable_and_never_publishes_secret_debug() { + let mut stage = GeneratedKeyStage::default(); + let key_material = TestKeyMaterialProvider::default(); + let handle = stage + .begin(&key_material, id(1), 4, time(10)) + .expect("begin"); + let view = handle.view(); + assert_eq!(view.expires_at().as_seconds(), 310); + assert_eq!(stage.pending().expect("pending").expected_revision(), 4); + assert!(stage.begin(&key_material, id(2), 4, time(11)).is_err()); + let nsec = handle.take_recovery_nsec().expect("one-use recovery"); + assert_eq!(nsec.with_exposed_secret(str::len), 63); + assert!(handle.take_recovery_nsec().is_err()); + assert!(stage.cancel()); + assert!(!stage.cancel()); + assert!(format!("{view:?}").contains(view.account().npub().as_str())); + assert!(!format!("{view:?}").contains("nsec1")); + } + + #[test] + fn stage_expires_and_is_destroyed_on_owner_drop() { + let mut stage = GeneratedKeyStage::default(); + let key_material = TestKeyMaterialProvider::default(); + stage + .begin(&key_material, id(1), 0, time(20)) + .expect("begin"); + let expiry = 20 + i64::try_from(GENERATED_KEY_STAGE_TTL.as_secs()).expect("ttl"); + assert!(stage.expire(time(expiry))); + assert!(stage.pending().is_none()); + assert!(stage.take(id(1), time(expiry)).is_err()); + + let mut shutdown_stage = GeneratedKeyStage::default(); + shutdown_stage + .begin(&key_material, id(2), 0, time(30)) + .expect("begin"); + drop(shutdown_stage); + } +} diff --git a/core/crates/studio_application/src/lib.rs b/core/crates/studio_application/src/lib.rs @@ -0,0 +1,58 @@ +#![doc = "Radroots Studio application runtime."] + +pub mod accounts; +pub mod actor; +pub mod app_core; +mod change_stream; +pub mod config; +pub mod custody; +pub mod ports; +mod profile_refresh; +pub mod recovery; +pub mod secrets; +pub mod session; +pub mod snapshot; +pub mod state_machine; + +#[cfg(test)] +mod test_support; + +pub use accounts::{ + GenerateAccountReceipt, ImportAccountReceipt, InMemoryAccountRepository, + InMemoryOperationJournal, +}; +pub use actor::{ + ActorMailbox, CommandContext, CommandEnvelope, CommandReceipt, CommandRejection, CommandResult, + CommandSubmission, CommandTicket, ForegroundSessionBinding, LifecycleGate, RequestId, + RuntimeCommandClass, RuntimeLifecycle, SessionGeneration, TaskCorrelation, +}; +pub use app_core::{AppCore, RemovalConfirmationToken, RemovalImpact}; +pub use change_stream::{ + ChangeSubscriptionId, OrderedSnapshotChanges, SnapshotChange, SnapshotChangeReceiver, +}; +pub use config::{ + RelayRuntimeMode, relay_configuration_from_environment, relay_configuration_from_value, +}; +pub use custody::{ + GENERATED_KEY_STAGE_TTL, GeneratedKeyRecoveryHandle, GeneratedKeyStage, GeneratedKeyStageView, + RecoveryStageId, StagedGeneratedKey, +}; +pub use ports::{ + AccountNamespaceRepository, AccountOperationKind, AccountOperationPhase, AccountPreferenceKey, + AccountRepository, AppStateRepository, BoxFuture, CachedProfile, Clock, + DurableAccountOperation, DurableOperationKind, DurableOperationPhase, DurableOperationReceipt, + DurableOperationRepository, DurableOperationStart, DurableRequestId, DurableTerminalOutcome, + GeneratedKeyMaterial, ImportedKeyMaterial, KeyMaterialProvider, NostrClient, + OperationDiagnostic, OperationId, OperationJournal, OperationPriorState, + PendingAccountOperation, ProfileFetchResult, ProfileRefreshStatus, ProfileRepository, + RelayFetchCompleteness, +}; +pub use profile_refresh::ProfileRefreshPlan; +pub use secrets::{ + FailureSecretStore, InMemorySecretStore, SecretStore, SecretStoreCall, SecretStoreOperation, +}; +pub use snapshot::{ + ActiveAccountSnapshot, AppLifecycle, AppSnapshot, MAX_CONFIGURED_RELAYS, ProfileLoadState, + RelayConfiguration, RelayConnectionState, SessionState, SnapshotRevision, +}; +pub use state_machine::{StateMachine, StateTransition}; diff --git a/core/crates/studio_application/src/ports.rs b/core/crates/studio_application/src/ports.rs @@ -0,0 +1,887 @@ +use std::future::Future; +use std::pin::Pin; +use std::time::Instant; + +use radroots_studio_domain::{ + AccountSummary, BindingAvailability, Kind0ProfileCandidate, Npub, Nsec, PublicKey, RelayUrl, + SafeError, SafeErrorCode, SafeMessage, SecretKeyInput, UnixTimestamp, +}; + +const MAX_DURABLE_REQUEST_ID_BYTES: usize = 128; + +#[derive(Clone, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)] +pub struct DurableRequestId(String); + +impl DurableRequestId { + /// Validates an opaque caller-generated idempotency key. + /// + /// # Errors + /// + /// Returns a safe validation error when the value is empty, oversized, or contains anything + /// other than visible ASCII characters. + pub fn parse(value: impl Into<String>) -> Result<Self, SafeError> { + let value = value.into(); + if value.is_empty() + || value.len() > MAX_DURABLE_REQUEST_ID_BYTES + || !value.bytes().all(|byte| byte.is_ascii_graphic()) + { + return Err(invalid_request_id()); + } + Ok(Self(value)) + } + + #[must_use] + pub fn as_str(&self) -> &str { + &self.0 + } +} + +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum DurableOperationKind { + Create, + Import, + Repair, + Remove, +} + +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum DurableOperationPhase { + IntentRecorded, + CredentialWritten, + MetadataCommitted, + SelectionCommitted, + CompensationPending, + CredentialDeleted, + MetadataDeleted, + Finalized, +} + +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum DurableTerminalOutcome { + Completed, + Cancelled, + Failed, +} + +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub struct OperationPriorState { + selected_account: Option<PublicKey>, + binding_availability: Option<BindingAvailability>, +} + +impl OperationPriorState { + #[must_use] + pub const fn new( + selected_account: Option<PublicKey>, + binding_availability: Option<BindingAvailability>, + ) -> Self { + Self { + selected_account, + binding_availability, + } + } + + #[must_use] + pub const fn selected_account(self) -> Option<PublicKey> { + self.selected_account + } + + #[must_use] + pub const fn binding_availability(self) -> Option<BindingAvailability> { + self.binding_availability + } +} + +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct DurableOperationReceipt { + request_id: DurableRequestId, + account: PublicKey, + outcome: DurableTerminalOutcome, + resulting_revision: Option<u64>, +} + +impl DurableOperationReceipt { + #[must_use] + pub const fn new( + request_id: DurableRequestId, + account: PublicKey, + outcome: DurableTerminalOutcome, + resulting_revision: Option<u64>, + ) -> Self { + Self { + request_id, + account, + outcome, + resulting_revision, + } + } + + #[must_use] + pub const fn request_id(&self) -> &DurableRequestId { + &self.request_id + } + + #[must_use] + pub const fn account(&self) -> PublicKey { + self.account + } + + #[must_use] + pub const fn outcome(&self) -> DurableTerminalOutcome { + self.outcome + } + + #[must_use] + pub const fn resulting_revision(&self) -> Option<u64> { + self.resulting_revision + } +} + +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct DurableAccountOperation { + request_id: DurableRequestId, + kind: DurableOperationKind, + account: PublicKey, + expected_revision: Option<u64>, + phase: DurableOperationPhase, + prior: OperationPriorState, + updated_at: UnixTimestamp, + diagnostic: Option<OperationDiagnostic>, + terminal: Option<DurableOperationReceipt>, +} + +impl DurableAccountOperation { + #[allow(clippy::too_many_arguments)] + #[must_use] + pub const fn new( + request_id: DurableRequestId, + kind: DurableOperationKind, + account: PublicKey, + expected_revision: Option<u64>, + phase: DurableOperationPhase, + prior: OperationPriorState, + updated_at: UnixTimestamp, + diagnostic: Option<OperationDiagnostic>, + terminal: Option<DurableOperationReceipt>, + ) -> Self { + Self { + request_id, + kind, + account, + expected_revision, + phase, + prior, + updated_at, + diagnostic, + terminal, + } + } + + #[must_use] + pub const fn request_id(&self) -> &DurableRequestId { + &self.request_id + } + #[must_use] + pub const fn kind(&self) -> DurableOperationKind { + self.kind + } + #[must_use] + pub const fn account(&self) -> PublicKey { + self.account + } + #[must_use] + pub const fn expected_revision(&self) -> Option<u64> { + self.expected_revision + } + #[must_use] + pub const fn phase(&self) -> DurableOperationPhase { + self.phase + } + #[must_use] + pub const fn prior(&self) -> OperationPriorState { + self.prior + } + #[must_use] + pub const fn updated_at(&self) -> UnixTimestamp { + self.updated_at + } + #[must_use] + pub const fn diagnostic(&self) -> Option<OperationDiagnostic> { + self.diagnostic + } + #[must_use] + pub const fn terminal(&self) -> Option<&DurableOperationReceipt> { + self.terminal.as_ref() + } +} + +#[derive(Clone, Debug, Eq, PartialEq)] +pub enum DurableOperationStart { + Started(DurableAccountOperation), + Existing(DurableAccountOperation), +} + +const fn invalid_request_id() -> SafeError { + SafeError::new( + SafeErrorCode::InvalidApplicationState, + SafeMessage::new("The request identifier is invalid."), + ) +} + +pub type BoxFuture<'a, T> = Pin<Box<dyn Future<Output = T> + Send + 'a>>; + +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum ProfileRefreshStatus { + Success, + Offline, + InvalidData, +} + +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum RelayFetchCompleteness { + Complete, + Partial, +} + +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct ProfileFetchResult { + candidate: Option<Kind0ProfileCandidate>, + completeness: RelayFetchCompleteness, +} + +impl ProfileFetchResult { + #[must_use] + pub const fn complete(candidate: Option<Kind0ProfileCandidate>) -> Self { + Self { + candidate, + completeness: RelayFetchCompleteness::Complete, + } + } + + #[must_use] + pub const fn partial(candidate: Option<Kind0ProfileCandidate>) -> Self { + Self { + candidate, + completeness: RelayFetchCompleteness::Partial, + } + } + + #[must_use] + pub fn into_parts(self) -> (Option<Kind0ProfileCandidate>, RelayFetchCompleteness) { + (self.candidate, self.completeness) + } +} + +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct CachedProfile { + candidate: Kind0ProfileCandidate, + refreshed_at: UnixTimestamp, + refresh_status: ProfileRefreshStatus, +} + +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum AccountPreferenceKey { + NamespaceProbe, +} + +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum AccountOperationKind { + Add, + Import, + Remove, +} + +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum AccountOperationPhase { + IntentRecorded, + CredentialWritten, + MetadataCommitted, + CompensationPending, + CredentialDeleted, + MetadataDeleted, +} + +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum OperationDiagnostic { + StorageUnavailable, + KeyringUnavailable, + CredentialMissing, + CompensationFailed, + Conflict, + Expired, +} + +#[derive(Clone, Copy, Debug, Eq, Hash, PartialEq)] +pub struct OperationId(u64); + +impl OperationId { + #[must_use] + pub const fn from_raw(value: u64) -> Self { + Self(value) + } + + #[must_use] + pub const fn as_raw(self) -> u64 { + self.0 + } +} + +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct PendingAccountOperation { + id: OperationId, + kind: AccountOperationKind, + subject: PublicKey, + phase: AccountOperationPhase, + updated_at: UnixTimestamp, + diagnostic: Option<OperationDiagnostic>, +} + +impl PendingAccountOperation { + #[must_use] + pub const fn new( + id: OperationId, + kind: AccountOperationKind, + subject: PublicKey, + phase: AccountOperationPhase, + updated_at: UnixTimestamp, + diagnostic: Option<OperationDiagnostic>, + ) -> Self { + Self { + id, + kind, + subject, + phase, + updated_at, + diagnostic, + } + } + + #[must_use] + pub const fn id(&self) -> OperationId { + self.id + } + #[must_use] + pub const fn kind(&self) -> AccountOperationKind { + self.kind + } + #[must_use] + pub const fn subject(&self) -> PublicKey { + self.subject + } + #[must_use] + pub const fn phase(&self) -> AccountOperationPhase { + self.phase + } + #[must_use] + pub const fn updated_at(&self) -> UnixTimestamp { + self.updated_at + } + #[must_use] + pub const fn diagnostic(&self) -> Option<OperationDiagnostic> { + self.diagnostic + } +} + +impl CachedProfile { + #[must_use] + pub const fn new( + candidate: Kind0ProfileCandidate, + refreshed_at: UnixTimestamp, + refresh_status: ProfileRefreshStatus, + ) -> Self { + Self { + candidate, + refreshed_at, + refresh_status, + } + } + + #[must_use] + pub const fn candidate(&self) -> &Kind0ProfileCandidate { + &self.candidate + } + + #[must_use] + pub const fn refreshed_at(&self) -> UnixTimestamp { + self.refreshed_at + } + + #[must_use] + pub const fn refresh_status(&self) -> ProfileRefreshStatus { + self.refresh_status + } +} + +pub trait AccountRepository: Send + Sync { + /// Lists saved public account records in deterministic order. + /// + /// # Errors + /// + /// Returns a safe storage error when records cannot be read. + fn list_accounts(&self) -> Result<Vec<AccountSummary>, SafeError>; + /// Finds one saved public account record. + /// + /// # Errors + /// + /// Returns a safe storage error when the lookup cannot complete. + fn find_account(&self, public_key: PublicKey) -> Result<Option<AccountSummary>, SafeError>; + /// Inserts one public account record. + /// + /// # Errors + /// + /// Returns a safe storage error when the durable write fails. + fn insert_account(&self, account: &AccountSummary) -> Result<(), SafeError>; + /// Updates one existing public account record. + /// + /// # Errors + /// + /// Returns a safe storage or account-not-found error when the durable + /// update cannot complete. + fn update_account(&self, account: &AccountSummary) -> Result<(), SafeError>; + /// Removes one public account record. + /// + /// # Errors + /// + /// Returns a safe storage error when the durable delete fails. + fn remove_account(&self, public_key: PublicKey) -> Result<(), SafeError>; +} + +pub trait ProfileRepository: Send + Sync { + /// Loads cached public profile metadata. + /// + /// # Errors + /// + /// Returns a safe storage error when the cache cannot be read. + fn load_profile(&self, public_key: PublicKey) -> Result<Option<CachedProfile>, SafeError>; + /// Saves a verified kind-0 profile candidate. + /// + /// # Errors + /// + /// Returns a safe storage error when the cache cannot be committed. + fn save_profile(&self, profile: &CachedProfile) -> Result<(), SafeError>; + /// Records the result of a profile refresh without replacing cached metadata. + /// + /// # Errors + /// + /// Returns a safe storage error when the cache cannot be committed. + fn record_refresh_status( + &self, + public_key: PublicKey, + refreshed_at: UnixTimestamp, + status: ProfileRefreshStatus, + ) -> Result<(), SafeError>; + /// Removes cached profile metadata for an account. + /// + /// # Errors + /// + /// Returns a safe storage error when the cache cannot be deleted. + fn remove_profile(&self, public_key: PublicKey) -> Result<(), SafeError>; +} + +pub trait AccountNamespaceRepository: Send + Sync { + /// Reads one internal non-secret account-scoped value. + /// + /// # Errors + /// + /// Returns a safe storage error when the value cannot be read. + fn get_value( + &self, + owner: PublicKey, + key: AccountPreferenceKey, + ) -> Result<Option<String>, SafeError>; + /// Writes one internal non-secret account-scoped value. + /// + /// # Errors + /// + /// Returns a safe storage error when the value cannot be committed. + fn set_value( + &self, + owner: PublicKey, + key: AccountPreferenceKey, + value: &str, + ) -> Result<(), SafeError>; + /// Removes all internal values owned by an account. + /// + /// # Errors + /// + /// Returns a safe storage error when cleanup cannot be committed. + fn clear_owner(&self, owner: PublicKey) -> Result<(), SafeError>; +} + +pub trait AppStateRepository: Send + Sync { + /// Loads the persisted selected account. + /// + /// # Errors + /// + /// Returns a safe storage error when application state cannot be read. + fn load_selected_account(&self) -> Result<Option<PublicKey>, SafeError>; + /// Persists the selected account or the empty selection. + /// + /// # Errors + /// + /// Returns a safe storage error when application state cannot be committed. + fn save_selected_account(&self, public_key: Option<PublicKey>) -> Result<(), SafeError>; +} + +pub trait OperationJournal: Send + Sync { + /// Records one cross-resource account operation intent. + /// + /// # Errors + /// + /// Returns a safe storage error when the entry cannot be committed. + fn begin_operation( + &self, + kind: AccountOperationKind, + subject: PublicKey, + updated_at: UnixTimestamp, + ) -> Result<OperationId, SafeError>; + /// Advances an operation to a durable recovery phase. + /// + /// # Errors + /// + /// Returns a safe storage error when the entry cannot be updated. + fn update_operation( + &self, + id: OperationId, + phase: AccountOperationPhase, + updated_at: UnixTimestamp, + diagnostic: Option<OperationDiagnostic>, + ) -> Result<(), SafeError>; + /// Loads all unfinished operations in deterministic order. + /// + /// # Errors + /// + /// Returns a safe storage error when entries cannot be read. + fn list_pending_operations(&self) -> Result<Vec<PendingAccountOperation>, SafeError>; + /// Deletes one fully reconciled operation entry. + /// + /// # Errors + /// + /// Returns a safe storage error when finalization cannot be committed. + fn finalize_operation(&self, id: OperationId) -> Result<(), SafeError>; +} + +pub trait DurableOperationRepository: Send + Sync { + /// Records one idempotent durable operation or returns the existing matching request. + /// + /// # Errors + /// + /// Returns a safe conflict or storage error when the request cannot be recorded. + #[allow(clippy::too_many_arguments)] + fn begin_durable_operation( + &self, + request_id: &DurableRequestId, + kind: DurableOperationKind, + account: PublicKey, + expected_revision: Option<u64>, + prior: OperationPriorState, + updated_at: UnixTimestamp, + ) -> Result<DurableOperationStart, SafeError>; + /// Loads one durable operation by its idempotency key. + /// + /// # Errors + /// + /// Returns a safe storage error when the lookup cannot complete. + fn load_durable_operation( + &self, + request_id: &DurableRequestId, + ) -> Result<Option<DurableAccountOperation>, SafeError>; + /// Advances one operation only from the caller's expected phase. + /// + /// # Errors + /// + /// Returns a safe conflict or storage error when the transition cannot commit. + fn advance_durable_operation( + &self, + request_id: &DurableRequestId, + expected_phase: DurableOperationPhase, + next_phase: DurableOperationPhase, + updated_at: UnixTimestamp, + diagnostic: Option<OperationDiagnostic>, + ) -> Result<DurableAccountOperation, SafeError>; + /// Finalizes one operation and durably retains its recoverable receipt. + /// + /// # Errors + /// + /// Returns a safe conflict or storage error when finalization cannot commit. + fn finalize_durable_operation( + &self, + request_id: &DurableRequestId, + expected_phase: DurableOperationPhase, + outcome: DurableTerminalOutcome, + resulting_revision: Option<u64>, + updated_at: UnixTimestamp, + ) -> Result<DurableOperationReceipt, SafeError>; + /// Lists unfinished operations in deterministic request order. + /// + /// # Errors + /// + /// Returns a safe storage error when operations cannot be read. + fn list_unfinished_durable_operations(&self) + -> Result<Vec<DurableAccountOperation>, SafeError>; +} + +pub trait NostrClient: Send + Sync { + fn fetch_profile<'a>( + &'a self, + public_key: PublicKey, + relays: &'a [RelayUrl], + deadline: Instant, + ) -> BoxFuture<'a, Result<ProfileFetchResult, SafeError>>; +} + +pub struct GeneratedKeyMaterial { + public_key: PublicKey, + npub: Npub, + secret: SecretKeyInput, + nsec: Nsec, +} + +impl GeneratedKeyMaterial { + #[must_use] + pub const fn new( + public_key: PublicKey, + npub: Npub, + secret: SecretKeyInput, + nsec: Nsec, + ) -> Self { + Self { + public_key, + npub, + secret, + nsec, + } + } + + #[must_use] + pub fn into_parts(self) -> (PublicKey, Npub, SecretKeyInput, Nsec) { + (self.public_key, self.npub, self.secret, self.nsec) + } +} + +pub struct ImportedKeyMaterial { + public_key: PublicKey, + npub: Npub, + secret: SecretKeyInput, +} + +impl ImportedKeyMaterial { + #[must_use] + pub const fn new(public_key: PublicKey, npub: Npub, secret: SecretKeyInput) -> Self { + Self { + public_key, + npub, + secret, + } + } + + #[must_use] + pub fn into_parts(self) -> (PublicKey, Npub, SecretKeyInput) { + (self.public_key, self.npub, self.secret) + } +} + +pub trait KeyMaterialProvider: Send + Sync { + /// Generates one keypair from host-provided cryptographic entropy. + /// + /// # Errors + /// + /// Returns a redacted key or entropy error. + fn generate(&self) -> Result<GeneratedKeyMaterial, SafeError>; + + /// Canonicalizes imported secret material and derives its public identity. + /// + /// # Errors + /// + /// Returns a redacted validation error. + fn import(&self, input: SecretKeyInput) -> Result<ImportedKeyMaterial, SafeError>; +} + +pub trait Clock: Send + Sync { + fn now(&self) -> UnixTimestamp; +} + +#[cfg(test)] +mod tests { + use std::time::Instant; + + use std::sync::Mutex; + + use radroots_studio_domain::{AccountSummary, PublicKey, RelayUrl, SafeError, UnixTimestamp}; + + use super::{ + AccountNamespaceRepository, AccountOperationKind, AccountOperationPhase, + AccountPreferenceKey, AccountRepository, AppStateRepository, BoxFuture, CachedProfile, + Clock, DurableOperationReceipt, DurableRequestId, DurableTerminalOutcome, NostrClient, + OperationDiagnostic, OperationId, OperationJournal, PendingAccountOperation, + ProfileFetchResult, ProfileRefreshStatus, ProfileRepository, + }; + + #[test] + fn durable_request_ids_and_terminal_receipts_are_bounded_and_public() { + let request = DurableRequestId::parse("create:desktop:0001").expect("request id"); + let receipt = DurableOperationReceipt::new( + request.clone(), + PublicKey::from_bytes([7; 32]).expect("valid public key"), + DurableTerminalOutcome::Completed, + Some(42), + ); + assert_eq!(receipt.request_id(), &request); + assert_eq!(receipt.resulting_revision(), Some(42)); + for invalid in ["", "contains space", &"x".repeat(129)] { + assert!(DurableRequestId::parse(invalid).is_err()); + } + } + + #[derive(Default)] + struct FakePorts { + selected: Mutex<Option<PublicKey>>, + } + + impl AccountRepository for FakePorts { + fn list_accounts(&self) -> Result<Vec<AccountSummary>, SafeError> { + Ok(Vec::new()) + } + + fn find_account( + &self, + _public_key: PublicKey, + ) -> Result<Option<AccountSummary>, SafeError> { + Ok(None) + } + + fn insert_account(&self, _account: &AccountSummary) -> Result<(), SafeError> { + Ok(()) + } + + fn update_account(&self, _account: &AccountSummary) -> Result<(), SafeError> { + Ok(()) + } + + fn remove_account(&self, _public_key: PublicKey) -> Result<(), SafeError> { + Ok(()) + } + } + + impl ProfileRepository for FakePorts { + fn load_profile(&self, _public_key: PublicKey) -> Result<Option<CachedProfile>, SafeError> { + Ok(None) + } + + fn save_profile(&self, _profile: &CachedProfile) -> Result<(), SafeError> { + Ok(()) + } + + fn record_refresh_status( + &self, + _public_key: PublicKey, + _refreshed_at: UnixTimestamp, + _status: ProfileRefreshStatus, + ) -> Result<(), SafeError> { + Ok(()) + } + + fn remove_profile(&self, _public_key: PublicKey) -> Result<(), SafeError> { + Ok(()) + } + } + + impl AccountNamespaceRepository for FakePorts { + fn get_value( + &self, + _owner: PublicKey, + _key: AccountPreferenceKey, + ) -> Result<Option<String>, SafeError> { + Ok(None) + } + + fn set_value( + &self, + _owner: PublicKey, + _key: AccountPreferenceKey, + _value: &str, + ) -> Result<(), SafeError> { + Ok(()) + } + + fn clear_owner(&self, _owner: PublicKey) -> Result<(), SafeError> { + Ok(()) + } + } + + impl AppStateRepository for FakePorts { + fn load_selected_account(&self) -> Result<Option<PublicKey>, SafeError> { + Ok(*self.selected.lock().expect("selected lock")) + } + + fn save_selected_account(&self, public_key: Option<PublicKey>) -> Result<(), SafeError> { + *self.selected.lock().expect("selected lock") = public_key; + Ok(()) + } + } + + impl OperationJournal for FakePorts { + fn begin_operation( + &self, + _kind: AccountOperationKind, + _subject: PublicKey, + _updated_at: UnixTimestamp, + ) -> Result<OperationId, SafeError> { + Ok(OperationId::from_raw(1)) + } + + fn update_operation( + &self, + _id: OperationId, + _phase: AccountOperationPhase, + _updated_at: UnixTimestamp, + _diagnostic: Option<OperationDiagnostic>, + ) -> Result<(), SafeError> { + Ok(()) + } + + fn list_pending_operations(&self) -> Result<Vec<PendingAccountOperation>, SafeError> { + Ok(Vec::new()) + } + + fn finalize_operation(&self, _id: OperationId) -> Result<(), SafeError> { + Ok(()) + } + } + + impl NostrClient for FakePorts { + fn fetch_profile<'a>( + &'a self, + _public_key: PublicKey, + _relays: &'a [RelayUrl], + _deadline: Instant, + ) -> BoxFuture<'a, Result<ProfileFetchResult, SafeError>> { + Box::pin(async { Ok(ProfileFetchResult::complete(None)) }) + } + } + + impl Clock for FakePorts { + fn now(&self) -> UnixTimestamp { + UnixTimestamp::from_seconds(1).expect("valid fake time") + } + } + + fn assert_send_sync<T: Send + Sync>() {} + + #[test] + fn ports_accept_send_sync_test_fakes() { + assert_send_sync::<FakePorts>(); + + let ports = FakePorts::default(); + ports + .save_selected_account(Some( + PublicKey::from_bytes([7_u8; 32]).expect("valid public key"), + )) + .expect("save selection"); + assert_eq!( + ports.load_selected_account().expect("load selection"), + Some(PublicKey::from_bytes([7_u8; 32]).expect("valid public key")) + ); + assert_eq!(ports.now().as_seconds(), 1); + } +} diff --git a/core/crates/studio_application/src/profile_refresh.rs b/core/crates/studio_application/src/profile_refresh.rs @@ -0,0 +1,659 @@ +use radroots_studio_domain::{PublicKey, RelayUrl, SafeError, SafeErrorCode}; +use std::time::Instant; + +use crate::{ + ActiveAccountSnapshot, AppCore, AppSnapshot, CachedProfile, Clock, NostrClient, + ProfileFetchResult, ProfileLoadState, ProfileRefreshStatus, ProfileRepository, + RelayConnectionState, RelayFetchCompleteness, 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. + /// + /// Cached public metadata remains visible while the asynchronous request is + /// running. Calling this command while signed out is an idempotent no-op. + /// + /// # Errors + /// + /// Returns a safe storage or application-state error. Relay and invalid-data + /// failures are represented as nonfatal snapshot state. + pub async fn refresh_active_profile( + &self, + profiles: &(impl ProfileRepository + ?Sized), + client: &(impl NostrClient + ?Sized), + clock: &(impl Clock + ?Sized), + deadline: Instant, + ) -> Result<AppSnapshot, SafeError> { + self.refresh_profile_for_active_account(profiles, client, clock, deadline) + .await + } + + /// Refreshes the current active account while retaining any cached profile. + /// + /// Stale results are discarded when the account is replaced or signed out. + /// + /// # Errors + /// + /// Returns a safe storage or application-state error. Relay and invalid-data + /// failures are represented as nonfatal snapshot state. + async fn refresh_profile_for_active_account( + &self, + profiles: &(impl ProfileRepository + ?Sized), + client: &(impl NostrClient + ?Sized), + clock: &(impl Clock + ?Sized), + deadline: Instant, + ) -> Result<AppSnapshot, SafeError> { + let Some(plan) = self.begin_profile_refresh()? else { + return Ok(self.snapshot()); + }; + let result = client + .fetch_profile(plan.public_key(), plan.relays(), deadline) + .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(); + let loading = self.apply_transition(StateTransition::UpdateActiveAccount { + expected: public_key, + active_account: Box::new(ActiveAccountSnapshot::new( + active.account().clone(), + RelayConnectionState::Connecting, + ProfileLoadState::Loading, + active.profile().cloned(), + )), + problem: None, + })?; + Ok(Some(ProfileRefreshPlan { + public_key, + active_account: active, + relays: loading.relay_configuration().relays().to_vec(), + expected_revision: loading.revision(), + })) + } + + /// 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<ProfileFetchResult, SafeError>, + profiles: &(impl ProfileRepository + ?Sized), + clock: &(impl Clock + ?Sized), + ) -> Result<AppSnapshot, SafeError> { + if !is_current_active(self, plan.public_key()) { + return Ok(self.snapshot()); + } + + let current_active = self + .snapshot() + .active_account() + .cloned() + .ok_or_else(invalid_profile_completion)?; + + match result { + Ok(fetched) => { + let (candidate, completeness) = fetched.into_parts(); + self.complete_successful_profile_fetch( + plan, + current_active, + candidate, + completeness, + profiles, + clock, + ) + } + Err(error) => { + let status = refresh_status(error); + profiles.record_refresh_status(plan.public_key(), clock.now(), status)?; + self.apply_transition(StateTransition::UpdateActiveAccount { + expected: plan.public_key(), + active_account: Box::new(ActiveAccountSnapshot::new( + current_active.account().clone(), + RelayConnectionState::Degraded, + ProfileLoadState::Error(error), + current_active.profile().cloned(), + )), + problem: Some(error), + }) + } + } + } + + fn complete_successful_profile_fetch( + &self, + plan: &ProfileRefreshPlan, + current_active: ActiveAccountSnapshot, + candidate: Option<radroots_studio_domain::Kind0ProfileCandidate>, + completeness: RelayFetchCompleteness, + profiles: &(impl ProfileRepository + ?Sized), + clock: &(impl Clock + ?Sized), + ) -> Result<AppSnapshot, SafeError> { + let relay_state = match completeness { + RelayFetchCompleteness::Complete => RelayConnectionState::Connected, + RelayFetchCompleteness::Partial => RelayConnectionState::Degraded, + }; + let problem = match completeness { + RelayFetchCompleteness::Complete => None, + RelayFetchCompleteness::Partial => Some(partial_relay_result()), + }; + match candidate { + Some(candidate) => { + let cached = CachedProfile::new( + candidate.clone(), + clock.now(), + ProfileRefreshStatus::Success, + ); + profiles.save_profile(&cached)?; + let winning_profile = profiles.load_profile(plan.public_key())?.map_or_else( + || candidate.metadata().clone(), + |profile| profile.candidate().metadata().clone(), + ); + self.apply_transition(StateTransition::UpdateActiveAccount { + expected: plan.public_key(), + active_account: Box::new(ActiveAccountSnapshot::new( + current_active.account().clone(), + relay_state, + ProfileLoadState::Fresh, + Some(winning_profile), + )), + problem, + }) + } + None => self.apply_transition(StateTransition::UpdateActiveAccount { + expected: plan.public_key(), + active_account: Box::new(ActiveAccountSnapshot::new( + current_active.account().clone(), + relay_state, + if current_active.profile().is_some() { + ProfileLoadState::Cached + } else { + ProfileLoadState::Empty + }, + current_active.profile().cloned(), + )), + problem, + }), + } + } +} + +const fn partial_relay_result() -> SafeError { + SafeError::new( + SafeErrorCode::RelayConnectionFailed, + radroots_studio_domain::SafeMessage::new("One or more Nostr relays did not complete."), + ) +} + +const fn invalid_profile_completion() -> SafeError { + SafeError::new( + SafeErrorCode::InvalidApplicationState, + radroots_studio_domain::SafeMessage::new("The active profile refresh is no longer valid."), + ) +} + +fn is_current_active(core: &AppCore, public_key: PublicKey) -> bool { + core.snapshot() + .active_account() + .is_some_and(|active| active.account().public_key() == public_key) +} + +const fn refresh_status(error: SafeError) -> ProfileRefreshStatus { + match error.code() { + SafeErrorCode::InvalidProfileMetadata | SafeErrorCode::ProfileRefreshFailed => { + ProfileRefreshStatus::InvalidData + } + _ => ProfileRefreshStatus::Offline, + } +} + +#[cfg(test)] +mod tests { + use std::sync::Mutex; + use std::time::{Duration, Instant}; + + use radroots_studio_domain::{ + EventId, Kind0ProfileCandidate, ProfileMetadata, PublicKey, RelayDestinationPolicy, + RelayUrl, SafeError, SafeErrorCode, SafeMessage, SecretKeyInput, UnixTimestamp, + select_latest_kind0, + }; + + use crate::{ + ActiveAccountSnapshot, AppCore, BoxFuture, CachedProfile, Clock, InMemoryAccountRepository, + InMemoryOperationJournal, InMemorySecretStore, NostrClient, ProfileFetchResult, + ProfileLoadState, ProfileRefreshStatus, ProfileRepository, RelayConfiguration, + RelayConnectionState, + }; + + #[derive(Default)] + struct MemoryProfiles(Mutex<Option<CachedProfile>>); + + impl ProfileRepository for MemoryProfiles { + fn load_profile(&self, _public_key: PublicKey) -> Result<Option<CachedProfile>, SafeError> { + Ok(self.0.lock().expect("profiles").clone()) + } + fn save_profile(&self, profile: &CachedProfile) -> Result<(), SafeError> { + let mut cached = self.0.lock().expect("profiles"); + let selected = cached.as_ref().map_or_else( + || profile.clone(), + |current| { + let winner = select_latest_kind0([ + current.candidate().clone(), + profile.candidate().clone(), + ]) + .expect("two candidates"); + if &winner == current.candidate() { + current.clone() + } else { + profile.clone() + } + }, + ); + *cached = Some(selected); + Ok(()) + } + fn record_refresh_status( + &self, + _public_key: PublicKey, + refreshed_at: UnixTimestamp, + status: ProfileRefreshStatus, + ) -> Result<(), SafeError> { + if let Some(profile) = self.0.lock().expect("profiles").as_mut() { + *profile = CachedProfile::new(profile.candidate().clone(), refreshed_at, status); + } + Ok(()) + } + fn remove_profile(&self, _public_key: PublicKey) -> Result<(), SafeError> { + *self.0.lock().expect("profiles") = None; + Ok(()) + } + } + + struct FixedClock; + impl Clock for FixedClock { + fn now(&self) -> UnixTimestamp { + UnixTimestamp::from_seconds(50).expect("time") + } + } + + struct FixedClient(Result<Option<Kind0ProfileCandidate>, SafeError>); + impl NostrClient for FixedClient { + fn fetch_profile<'a>( + &'a self, + _public_key: PublicKey, + _relays: &'a [RelayUrl], + _deadline: std::time::Instant, + ) -> BoxFuture<'a, Result<ProfileFetchResult, SafeError>> { + let result = self.0.clone(); + Box::pin(async move { result.map(ProfileFetchResult::complete) }) + } + } + + struct BlockingClient { + started: tokio::sync::Semaphore, + release: tokio::sync::Semaphore, + result: Result<Option<Kind0ProfileCandidate>, SafeError>, + } + + impl BlockingClient { + fn new(result: Result<Option<Kind0ProfileCandidate>, SafeError>) -> Self { + Self { + started: tokio::sync::Semaphore::new(0), + release: tokio::sync::Semaphore::new(0), + result, + } + } + } + + impl NostrClient for BlockingClient { + fn fetch_profile<'a>( + &'a self, + _public_key: PublicKey, + _relays: &'a [RelayUrl], + _deadline: std::time::Instant, + ) -> BoxFuture<'a, Result<ProfileFetchResult, SafeError>> { + Box::pin(async move { + self.started.add_permits(1); + let permit = self.release.acquire().await.expect("release open"); + permit.forget(); + self.result.clone().map(ProfileFetchResult::complete) + }) + } + } + + fn profile(public_key: PublicKey, name: &str, timestamp: i64) -> Kind0ProfileCandidate { + Kind0ProfileCandidate::new( + EventId::from_bytes([u8::try_from(timestamp).expect("small timestamp"); 32]), + public_key, + UnixTimestamp::from_seconds(timestamp).expect("time"), + ProfileMetadata::new(Some(name.to_owned()), None, None, None, None).expect("profile"), + ) + } + + fn active_core(profiles: &MemoryProfiles, cached_name: Option<&str>) -> (AppCore, PublicKey) { + let relays = RelayConfiguration::new(vec![ + RelayUrl::parse("ws://localhost:8080", RelayDestinationPolicy::Local).expect("relay"), + ]) + .expect("relay configuration"); + let core = AppCore::in_memory(relays); + let accounts = InMemoryAccountRepository::default(); + let secrets = InMemorySecretStore::default(); + let journal = InMemoryOperationJournal::default(); + core.bootstrap().expect("bootstrap"); + let public_key = core + .import_secret_key( + SecretKeyInput::parse( + "7e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7".to_owned(), + ) + .expect("secret"), + &accounts, + &accounts, + &secrets, + &journal, + &FixedClock, + ) + .expect("import") + .account() + .public_key(); + if let Some(name) = cached_name { + profiles + .save_profile(&CachedProfile::new( + profile(public_key, name, 10), + UnixTimestamp::from_seconds(11).expect("time"), + ProfileRefreshStatus::Success, + )) + .expect("cache"); + } + core.activate_account( + public_key, + &accounts, + &accounts, + profiles, + &secrets, + &FixedClock, + ) + .expect("activate"); + (core, public_key) + } + + #[tokio::test] + async fn refresh_transitions_from_cache_through_loading_to_fresh_profile() { + let profiles = MemoryProfiles::default(); + let (core, public_key) = active_core(&profiles, Some("Cached")); + assert_eq!( + core.snapshot() + .active_account() + .map(crate::ActiveAccountSnapshot::profile_state), + Some(ProfileLoadState::Cached) + ); + let plan = core + .begin_profile_refresh() + .expect("begin refresh") + .expect("active refresh"); + let loading = core.snapshot(); + assert_eq!( + loading + .active_account() + .map(crate::ActiveAccountSnapshot::profile_state), + Some(ProfileLoadState::Loading) + ); + assert_eq!( + loading + .active_account() + .map(crate::ActiveAccountSnapshot::relay_state), + Some(RelayConnectionState::Connecting) + ); + let client = FixedClient(Ok(Some(profile(public_key, "Fresh", 20)))); + let result = client + .fetch_profile(plan.public_key(), plan.relays(), deadline()) + .await; + core.complete_profile_refresh(&plan, result, &profiles, &FixedClock) + .expect("complete refresh"); + assert_eq!( + core.snapshot() + .active_account() + .and_then(|active| active.profile()) + .and_then(ProfileMetadata::name), + Some("Fresh") + ); + } + + #[tokio::test] + async fn refresh_failure_preserves_cached_profile_as_nonfatal_state() { + let profiles = MemoryProfiles::default(); + let (core, public_key) = active_core(&profiles, Some("Cached")); + let cached = profile(public_key, "Cached", 10); + let error = SafeError::new( + SafeErrorCode::RelayConnectionFailed, + SafeMessage::new("The relay is offline."), + ); + + let snapshot = core + .refresh_profile_for_active_account( + &profiles, + &FixedClient(Err(error)), + &FixedClock, + deadline(), + ) + .await + .expect("nonfatal refresh"); + + assert_eq!(snapshot.recoverable_problem(), Some(error)); + assert_eq!( + snapshot + .active_account() + .map(crate::ActiveAccountSnapshot::relay_state), + Some(RelayConnectionState::Degraded) + ); + assert_eq!( + profiles + .load_profile(public_key) + .expect("load") + .expect("cache") + .candidate(), + &cached + ); + } + + #[tokio::test] + async fn refresh_discards_stale_completion_after_sign_out() { + let profiles = MemoryProfiles::default(); + let (core, public_key) = active_core(&profiles, Some("Cached")); + let client = BlockingClient::new(Ok(Some(profile(public_key, "Stale", 20)))); + + let refresh = + core.refresh_profile_for_active_account(&profiles, &client, &FixedClock, deadline()); + let sign_out = async { + let permit = client.started.acquire().await.expect("refresh starts"); + permit.forget(); + core.sign_out().expect("sign out"); + client.release.add_permits(1); + }; + let (result, ()) = tokio::join!(refresh, sign_out); + + assert!( + result + .expect("stale result is harmless") + .active_account() + .is_none() + ); + assert_eq!( + profiles + .load_profile(public_key) + .expect("load") + .expect("cached") + .candidate() + .metadata() + .name(), + Some("Cached") + ); + } + + #[tokio::test] + async fn manual_refresh_is_repeatable_and_signed_out_safe() { + let profiles = MemoryProfiles::default(); + let (core, public_key) = active_core(&profiles, None); + let first = core + .refresh_active_profile( + &profiles, + &FixedClient(Ok(Some(profile(public_key, "First", 10)))), + &FixedClock, + deadline(), + ) + .await + .expect("first refresh"); + let second = core + .refresh_active_profile( + &profiles, + &FixedClient(Ok(Some(profile(public_key, "Second", 20)))), + &FixedClock, + deadline(), + ) + .await + .expect("second refresh"); + + assert!(second.revision() > first.revision()); + assert_eq!( + second + .active_account() + .and_then(|active| active.profile()) + .and_then(ProfileMetadata::name), + Some("Second") + ); + let signed_out = core.sign_out().expect("sign out"); + let no_op = core + .refresh_active_profile(&profiles, &FixedClient(Ok(None)), &FixedClock, deadline()) + .await + .expect("signed-out no-op"); + assert_eq!(no_op, signed_out); + } + + fn deadline() -> Instant { + Instant::now() + Duration::from_secs(5) + } + + #[test] + fn partial_relay_success_retains_verified_data_and_marks_degraded_connectivity() { + let profiles = MemoryProfiles::default(); + let (core, public_key) = active_core(&profiles, None); + let plan = core + .begin_profile_refresh() + .expect("begin") + .expect("active"); + let snapshot = core + .complete_profile_refresh( + &plan, + Ok(ProfileFetchResult::partial(Some(profile( + public_key, "Partial", 30, + )))), + &profiles, + &FixedClock, + ) + .expect("partial completion"); + let active = snapshot.active_account().expect("active account"); + assert_eq!(active.relay_state(), RelayConnectionState::Degraded); + assert_eq!(active.profile_state(), ProfileLoadState::Fresh); + assert_eq!( + active.profile().and_then(ProfileMetadata::name), + Some("Partial") + ); + assert_eq!( + snapshot.recoverable_problem().map(SafeError::code), + Some(SafeErrorCode::RelayConnectionFailed) + ); + } + + #[test] + fn overlapping_refreshes_keep_the_newest_event_regardless_of_completion_order() { + let profiles = MemoryProfiles::default(); + let (core, public_key) = active_core(&profiles, Some("Cached")); + let first = core + .begin_profile_refresh() + .expect("first") + .expect("active"); + let second = core + .begin_profile_refresh() + .expect("second") + .expect("active"); + + core.complete_profile_refresh( + &second, + Ok(ProfileFetchResult::complete(Some(profile( + public_key, "Newest", 30, + )))), + &profiles, + &FixedClock, + ) + .expect("newest completes first"); + let final_snapshot = core + .complete_profile_refresh( + &first, + Ok(ProfileFetchResult::complete(Some(profile( + public_key, "Older", 20, + )))), + &profiles, + &FixedClock, + ) + .expect("older completes last"); + + assert_eq!( + final_snapshot + .active_account() + .and_then(ActiveAccountSnapshot::profile) + .and_then(ProfileMetadata::name), + Some("Newest") + ); + assert_eq!( + profiles + .load_profile(public_key) + .expect("cache") + .expect("profile") + .candidate() + .metadata() + .name(), + Some("Newest") + ); + } +} diff --git a/core/crates/studio_application/src/recovery.rs b/core/crates/studio_application/src/recovery.rs @@ -0,0 +1,788 @@ +use radroots_studio_domain::{PublicKey, SafeError}; + +use crate::{ + AccountOperationKind, AccountOperationPhase, AccountRepository, AppCore, AppStateRepository, + Clock, DurableAccountOperation, DurableOperationKind, DurableOperationPhase, + DurableOperationRepository, DurableTerminalOutcome, OperationJournal, SecretStore, +}; + +impl AppCore { + /// Reconciles durable request operations before public state is restored. + /// + /// # Errors + /// + /// Returns a safe credential, persistence, or recovery error while retaining the operation + /// at its last durable phase for a later retry. + pub fn recover_durable_operations( + &self, + accounts: &(impl AccountRepository + ?Sized), + app_state: &(impl AppStateRepository + ?Sized), + secrets: &(impl SecretStore + ?Sized), + operations: &(impl DurableOperationRepository + ?Sized), + clock: &(impl Clock + ?Sized), + ) -> Result<(), SafeError> { + for operation in operations.list_unfinished_durable_operations()? { + match operation.kind() { + DurableOperationKind::Create + | DurableOperationKind::Import + | DurableOperationKind::Repair => recover_durable_addition( + &operation, accounts, app_state, secrets, operations, clock, + )?, + DurableOperationKind::Remove => recover_durable_removal( + &operation, accounts, app_state, secrets, operations, clock, + )?, + } + } + Ok(()) + } + + /// Reconciles non-secret cross-resource journal entries before bootstrap. + /// + /// An empty journal does not access the credential store. + /// + /// # Errors + /// + /// Returns a safe credential, persistence, or recovery error while retaining + /// the unfinished journal entry for a later retry. + pub fn recover_pending_operations( + &self, + accounts: &(impl AccountRepository + ?Sized), + app_state: &(impl AppStateRepository + ?Sized), + secrets: &(impl SecretStore + ?Sized), + journal: &(impl OperationJournal + ?Sized), + clock: &(impl Clock + ?Sized), + ) -> Result<(), SafeError> { + for operation in journal.list_pending_operations()? { + match operation.kind() { + AccountOperationKind::Remove => { + recover_removal(&operation, accounts, app_state, secrets, journal, clock)?; + } + AccountOperationKind::Add | AccountOperationKind::Import => { + recover_addition(&operation, accounts, secrets, journal, clock)?; + } + } + } + Ok(()) + } +} + +fn recover_durable_removal( + operation: &DurableAccountOperation, + accounts: &(impl AccountRepository + ?Sized), + app_state: &(impl AppStateRepository + ?Sized), + secrets: &(impl SecretStore + ?Sized), + operations: &(impl DurableOperationRepository + ?Sized), + clock: &(impl Clock + ?Sized), +) -> Result<(), SafeError> { + let request = operation.request_id(); + let account = operation.account(); + let mut phase = operation.phase(); + if phase == DurableOperationPhase::IntentRecorded { + if secrets.contains(account)? { + secrets.delete(account)?; + } + operations.advance_durable_operation( + request, + phase, + DurableOperationPhase::CredentialDeleted, + clock.now(), + None, + )?; + phase = DurableOperationPhase::CredentialDeleted; + } + if phase == DurableOperationPhase::CredentialDeleted { + if accounts.find_account(account)?.is_some() { + accounts.remove_account(account)?; + } + operations.advance_durable_operation( + request, + phase, + DurableOperationPhase::MetadataDeleted, + clock.now(), + None, + )?; + phase = DurableOperationPhase::MetadataDeleted; + } + if phase == DurableOperationPhase::MetadataDeleted { + app_state.save_selected_account(operation.prior().selected_account())?; + operations.advance_durable_operation( + request, + phase, + DurableOperationPhase::SelectionCommitted, + clock.now(), + None, + )?; + phase = DurableOperationPhase::SelectionCommitted; + } + if phase == DurableOperationPhase::SelectionCommitted { + operations.finalize_durable_operation( + request, + phase, + DurableTerminalOutcome::Completed, + None, + clock.now(), + )?; + } + Ok(()) +} + +fn recover_durable_addition( + operation: &DurableAccountOperation, + accounts: &(impl AccountRepository + ?Sized), + app_state: &(impl AppStateRepository + ?Sized), + secrets: &(impl SecretStore + ?Sized), + operations: &(impl DurableOperationRepository + ?Sized), + clock: &(impl Clock + ?Sized), +) -> Result<(), SafeError> { + let request = operation.request_id(); + let account = operation.account(); + match operation.phase() { + DurableOperationPhase::IntentRecorded => { + if secrets.contains(account)? { + secrets.delete(account)?; + } + operations.finalize_durable_operation( + request, + DurableOperationPhase::IntentRecorded, + DurableTerminalOutcome::Failed, + None, + clock.now(), + )?; + } + DurableOperationPhase::CredentialWritten => { + let metadata = accounts.find_account(account)?; + let committed = metadata.as_ref().is_some_and(|saved| { + saved.signer().availability() + == radroots_studio_domain::BindingAvailability::Available + }); + if committed { + operations.advance_durable_operation( + request, + DurableOperationPhase::CredentialWritten, + DurableOperationPhase::MetadataCommitted, + clock.now(), + None, + )?; + finish_durable_selection(operation, app_state, operations, clock)?; + } else { + operations.advance_durable_operation( + request, + DurableOperationPhase::CredentialWritten, + DurableOperationPhase::CompensationPending, + clock.now(), + None, + )?; + compensate_durable_addition( + operation, accounts, app_state, secrets, operations, clock, + )?; + } + } + DurableOperationPhase::MetadataCommitted => { + finish_durable_selection(operation, app_state, operations, clock)?; + } + DurableOperationPhase::SelectionCommitted => { + operations.finalize_durable_operation( + request, + DurableOperationPhase::SelectionCommitted, + DurableTerminalOutcome::Completed, + None, + clock.now(), + )?; + } + DurableOperationPhase::CompensationPending => { + compensate_durable_addition( + operation, accounts, app_state, secrets, operations, clock, + )?; + } + DurableOperationPhase::CredentialDeleted | DurableOperationPhase::MetadataDeleted => { + operations.finalize_durable_operation( + request, + operation.phase(), + DurableTerminalOutcome::Failed, + None, + clock.now(), + )?; + } + DurableOperationPhase::Finalized => {} + } + Ok(()) +} + +fn finish_durable_selection( + operation: &DurableAccountOperation, + app_state: &(impl AppStateRepository + ?Sized), + operations: &(impl DurableOperationRepository + ?Sized), + clock: &(impl Clock + ?Sized), +) -> Result<(), SafeError> { + app_state.save_selected_account(Some(operation.account()))?; + operations.advance_durable_operation( + operation.request_id(), + DurableOperationPhase::MetadataCommitted, + DurableOperationPhase::SelectionCommitted, + clock.now(), + None, + )?; + operations.finalize_durable_operation( + operation.request_id(), + DurableOperationPhase::SelectionCommitted, + DurableTerminalOutcome::Completed, + None, + clock.now(), + )?; + Ok(()) +} + +fn compensate_durable_addition( + operation: &DurableAccountOperation, + accounts: &(impl AccountRepository + ?Sized), + app_state: &(impl AppStateRepository + ?Sized), + secrets: &(impl SecretStore + ?Sized), + operations: &(impl DurableOperationRepository + ?Sized), + clock: &(impl Clock + ?Sized), +) -> Result<(), SafeError> { + if secrets.contains(operation.account())? { + secrets.delete(operation.account())?; + } + if let Some(availability) = operation.prior().binding_availability() { + if let Some(previous) = accounts.find_account(operation.account())? { + accounts.update_account(&previous.with_binding_availability(availability))?; + } + } else if accounts.find_account(operation.account())?.is_some() { + accounts.remove_account(operation.account())?; + } + app_state.save_selected_account(operation.prior().selected_account())?; + operations.advance_durable_operation( + operation.request_id(), + DurableOperationPhase::CompensationPending, + DurableOperationPhase::CredentialDeleted, + clock.now(), + None, + )?; + operations.finalize_durable_operation( + operation.request_id(), + DurableOperationPhase::CredentialDeleted, + DurableTerminalOutcome::Failed, + None, + clock.now(), + )?; + Ok(()) +} + +fn recover_removal( + operation: &crate::PendingAccountOperation, + accounts: &(impl AccountRepository + ?Sized), + app_state: &(impl AppStateRepository + ?Sized), + secrets: &(impl SecretStore + ?Sized), + journal: &(impl OperationJournal + ?Sized), + clock: &(impl Clock + ?Sized), +) -> Result<(), SafeError> { + let public_key = operation.subject(); + if operation.phase() == AccountOperationPhase::IntentRecorded { + match secrets.delete(public_key) { + Ok(()) => {} + Err(error) + if error.code() == radroots_studio_domain::SafeErrorCode::CredentialMissing => {} + Err(error) => return Err(error), + } + journal.update_operation( + operation.id(), + AccountOperationPhase::CredentialDeleted, + clock.now(), + None, + )?; + } + if matches!( + operation.phase(), + AccountOperationPhase::IntentRecorded | AccountOperationPhase::CredentialDeleted + ) { + let registry = accounts.list_accounts()?; + let selected = removal_fallback(&registry, app_state.load_selected_account()?, public_key); + accounts.remove_account(public_key)?; + app_state.save_selected_account(selected)?; + journal.update_operation( + operation.id(), + AccountOperationPhase::MetadataDeleted, + clock.now(), + None, + )?; + } + journal.finalize_operation(operation.id()) +} + +fn recover_addition( + operation: &crate::PendingAccountOperation, + accounts: &(impl AccountRepository + ?Sized), + secrets: &(impl SecretStore + ?Sized), + journal: &(impl OperationJournal + ?Sized), + clock: &(impl Clock + ?Sized), +) -> Result<(), SafeError> { + let has_metadata = accounts.find_account(operation.subject())?.is_some(); + match operation.phase() { + AccountOperationPhase::CredentialWritten | AccountOperationPhase::CompensationPending + if !has_metadata => + { + match secrets.delete(operation.subject()) { + Ok(()) => {} + Err(error) + if error.code() == radroots_studio_domain::SafeErrorCode::CredentialMissing => { + } + Err(error) => return Err(error), + } + journal.update_operation( + operation.id(), + AccountOperationPhase::MetadataDeleted, + clock.now(), + None, + )?; + } + _ => {} + } + journal.finalize_operation(operation.id()) +} + +fn removal_fallback( + registry: &[radroots_studio_domain::AccountSummary], + selected: Option<PublicKey>, + removed: PublicKey, +) -> Option<PublicKey> { + if selected != Some(removed) { + return selected; + } + let index = registry + .iter() + .position(|account| account.public_key() == removed)?; + registry + .get(index + 1) + .or_else(|| index.checked_sub(1).and_then(|before| registry.get(before))) + .map(radroots_studio_domain::AccountSummary::public_key) +} + +#[cfg(test)] +pub(crate) mod tests { + use std::sync::{Mutex, MutexGuard}; + + use radroots_studio_domain::{ + BindingAvailability, PublicKey, SafeError, SafeErrorCode, SafeMessage, SecretKeyInput, + UnixTimestamp, + }; + + use super::*; + use crate::{ + DurableOperationReceipt, DurableOperationStart, DurableRequestId, FailureSecretStore, + InMemoryAccountRepository, InMemoryOperationJournal, InMemorySecretStore, + RelayConfiguration, SecretStore, SecretStoreOperation, + }; + + struct FixedClock; + + impl Clock for FixedClock { + fn now(&self) -> UnixTimestamp { + UnixTimestamp::from_seconds(10).expect("time") + } + } + + pub(crate) struct TestDurableRepository { + operation: Mutex<DurableAccountOperation>, + return_existing: bool, + } + + impl TestDurableRepository { + pub(crate) fn new(operation: DurableAccountOperation) -> Self { + Self { + operation: Mutex::new(operation), + return_existing: true, + } + } + + pub(crate) fn fresh(operation: DurableAccountOperation) -> Self { + Self { + operation: Mutex::new(operation), + return_existing: false, + } + } + + pub(crate) fn operation(&self) -> MutexGuard<'_, DurableAccountOperation> { + self.operation + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + } + + fn replace( + current: &DurableAccountOperation, + phase: DurableOperationPhase, + diagnostic: Option<crate::OperationDiagnostic>, + terminal: Option<DurableOperationReceipt>, + ) -> DurableAccountOperation { + DurableAccountOperation::new( + current.request_id().clone(), + current.kind(), + current.account(), + current.expected_revision(), + phase, + current.prior(), + current.updated_at(), + diagnostic, + terminal, + ) + } + } + + impl DurableOperationRepository for TestDurableRepository { + fn begin_durable_operation( + &self, + _request_id: &DurableRequestId, + _kind: DurableOperationKind, + _account: PublicKey, + _expected_revision: Option<u64>, + _prior: crate::OperationPriorState, + _updated_at: UnixTimestamp, + ) -> Result<DurableOperationStart, SafeError> { + let operation = self.operation().clone(); + Ok(if self.return_existing { + DurableOperationStart::Existing(operation) + } else { + DurableOperationStart::Started(operation) + }) + } + + fn load_durable_operation( + &self, + request_id: &DurableRequestId, + ) -> Result<Option<DurableAccountOperation>, SafeError> { + if !self.return_existing { + return Ok(None); + } + let operation = self.operation(); + Ok((operation.request_id() == request_id).then(|| operation.clone())) + } + + fn advance_durable_operation( + &self, + request_id: &DurableRequestId, + expected_phase: DurableOperationPhase, + next_phase: DurableOperationPhase, + _updated_at: UnixTimestamp, + diagnostic: Option<crate::OperationDiagnostic>, + ) -> Result<DurableAccountOperation, SafeError> { + let mut operation = self.operation(); + if operation.request_id() != request_id || operation.phase() != expected_phase { + return Err(conflict()); + } + *operation = Self::replace(&operation, next_phase, diagnostic, None); + Ok(operation.clone()) + } + + fn finalize_durable_operation( + &self, + request_id: &DurableRequestId, + expected_phase: DurableOperationPhase, + outcome: DurableTerminalOutcome, + resulting_revision: Option<u64>, + _updated_at: UnixTimestamp, + ) -> Result<DurableOperationReceipt, SafeError> { + let mut operation = self.operation(); + if operation.request_id() != request_id || operation.phase() != expected_phase { + return Err(conflict()); + } + let receipt = DurableOperationReceipt::new( + request_id.clone(), + operation.account(), + outcome, + resulting_revision, + ); + *operation = Self::replace( + &operation, + DurableOperationPhase::Finalized, + operation.diagnostic(), + Some(receipt.clone()), + ); + Ok(receipt) + } + + fn list_unfinished_durable_operations( + &self, + ) -> Result<Vec<DurableAccountOperation>, SafeError> { + Ok(vec![self.operation().clone()]) + } + } + + fn conflict() -> SafeError { + SafeError::new( + SafeErrorCode::InvalidApplicationState, + SafeMessage::new("The test durable operation conflicted."), + ) + } + + fn seeded() -> ( + AppCore, + InMemoryAccountRepository, + InMemorySecretStore, + InMemoryOperationJournal, + PublicKey, + ) { + let core = AppCore::in_memory(RelayConfiguration::default()); + let accounts = InMemoryAccountRepository::default(); + let secrets = InMemorySecretStore::default(); + let journal = InMemoryOperationJournal::default(); + core.bootstrap().expect("bootstrap"); + let receipt = core + .generate_account(&accounts, &accounts, &secrets, &journal, &FixedClock) + .expect("seed account"); + ( + core, + accounts, + secrets, + journal, + receipt.account().public_key(), + ) + } + + pub(crate) fn operation( + kind: DurableOperationKind, + phase: DurableOperationPhase, + account: PublicKey, + prior_availability: Option<BindingAvailability>, + ) -> DurableAccountOperation { + DurableAccountOperation::new( + DurableRequestId::parse(format!("{kind:?}-{phase:?}")).expect("durable request ID"), + kind, + account, + Some(1), + phase, + crate::OperationPriorState::new(None, prior_availability), + FixedClock.now(), + None, + None, + ) + } + + fn run_durable( + core: &AppCore, + accounts: &InMemoryAccountRepository, + secrets: &InMemorySecretStore, + operation: DurableAccountOperation, + ) -> DurableAccountOperation { + let repository = TestDurableRepository::new(operation); + core.recover_durable_operations(accounts, accounts, secrets, &repository, &FixedClock) + .expect("durable recovery"); + repository.operation().clone() + } + + #[test] + fn durable_recovery_exercises_every_removal_phase_and_presence_branch() { + for phase in [ + DurableOperationPhase::IntentRecorded, + DurableOperationPhase::CredentialDeleted, + DurableOperationPhase::MetadataDeleted, + DurableOperationPhase::SelectionCommitted, + DurableOperationPhase::Finalized, + ] { + let (core, accounts, secrets, _journal, public_key) = seeded(); + if phase != DurableOperationPhase::IntentRecorded { + secrets.delete(public_key).expect("delete credential"); + } + if matches!( + phase, + DurableOperationPhase::MetadataDeleted + | DurableOperationPhase::SelectionCommitted + | DurableOperationPhase::Finalized + ) { + accounts.remove_account(public_key).expect("remove account"); + } + let recovered = run_durable( + &core, + &accounts, + &secrets, + operation(DurableOperationKind::Remove, phase, public_key, None), + ); + assert_eq!(recovered.phase(), DurableOperationPhase::Finalized); + } + + let (core, accounts, secrets, _journal, public_key) = seeded(); + secrets.delete(public_key).expect("delete credential"); + accounts.remove_account(public_key).expect("remove account"); + let recovered = run_durable( + &core, + &accounts, + &secrets, + operation( + DurableOperationKind::Remove, + DurableOperationPhase::IntentRecorded, + public_key, + None, + ), + ); + assert_eq!(recovered.phase(), DurableOperationPhase::Finalized); + } + + #[test] + fn durable_recovery_exercises_every_addition_phase_and_compensation_shape() { + for phase in [ + DurableOperationPhase::IntentRecorded, + DurableOperationPhase::CredentialWritten, + DurableOperationPhase::MetadataCommitted, + DurableOperationPhase::SelectionCommitted, + DurableOperationPhase::CredentialDeleted, + DurableOperationPhase::MetadataDeleted, + DurableOperationPhase::Finalized, + ] { + let (core, accounts, secrets, _journal, public_key) = seeded(); + if phase == DurableOperationPhase::IntentRecorded { + accounts + .remove_account(public_key) + .expect("remove metadata"); + } + let recovered = run_durable( + &core, + &accounts, + &secrets, + operation(DurableOperationKind::Create, phase, public_key, None), + ); + assert_eq!(recovered.phase(), DurableOperationPhase::Finalized); + } + + for (prior, retain_metadata, retain_secret) in [ + (Some(BindingAvailability::CredentialMissing), true, true), + (Some(BindingAvailability::CredentialMissing), false, true), + (None, true, true), + (None, false, false), + ] { + let (core, accounts, secrets, _journal, public_key) = seeded(); + if !retain_metadata { + accounts + .remove_account(public_key) + .expect("remove metadata"); + } + if !retain_secret { + secrets.delete(public_key).expect("delete credential"); + } + let recovered = run_durable( + &core, + &accounts, + &secrets, + operation( + DurableOperationKind::Repair, + DurableOperationPhase::CompensationPending, + public_key, + prior, + ), + ); + assert_eq!(recovered.phase(), DurableOperationPhase::Finalized); + } + + let (core, accounts, secrets, _journal, public_key) = seeded(); + accounts + .remove_account(public_key) + .expect("remove metadata"); + let recovered = run_durable( + &core, + &accounts, + &secrets, + operation( + DurableOperationKind::Import, + DurableOperationPhase::CredentialWritten, + public_key, + None, + ), + ); + assert_eq!(recovered.phase(), DurableOperationPhase::Finalized); + + let (core, accounts, secrets, _journal, public_key) = seeded(); + accounts + .remove_account(public_key) + .expect("remove metadata"); + secrets.delete(public_key).expect("delete credential"); + let recovered = run_durable( + &core, + &accounts, + &secrets, + operation( + DurableOperationKind::Create, + DurableOperationPhase::IntentRecorded, + public_key, + None, + ), + ); + assert_eq!(recovered.phase(), DurableOperationPhase::Finalized); + } + + #[test] + fn pending_recovery_exercises_removal_and_addition_presence_branches() { + for credential_present in [true, false] { + let (core, accounts, secrets, journal, public_key) = seeded(); + if !credential_present { + secrets.delete(public_key).expect("delete credential"); + accounts + .save_selected_account(None) + .expect("clear selection"); + } + journal + .begin_operation(AccountOperationKind::Remove, public_key, FixedClock.now()) + .expect("removal intent"); + core.recover_pending_operations(&accounts, &accounts, &secrets, &journal, &FixedClock) + .expect("removal recovery"); + assert!(journal.list_pending_operations().unwrap().is_empty()); + } + + for (kind, metadata_present, credential_present) in [ + (AccountOperationKind::Add, false, true), + (AccountOperationKind::Import, false, false), + (AccountOperationKind::Add, true, true), + ] { + let (core, accounts, secrets, journal, public_key) = seeded(); + if !metadata_present { + accounts + .remove_account(public_key) + .expect("remove metadata"); + } + if !credential_present { + secrets.delete(public_key).expect("delete credential"); + } + let id = journal + .begin_operation(kind, public_key, FixedClock.now()) + .expect("addition intent"); + journal + .update_operation( + id, + AccountOperationPhase::CredentialWritten, + FixedClock.now(), + None, + ) + .expect("credential phase"); + core.recover_pending_operations(&accounts, &accounts, &secrets, &journal, &FixedClock) + .expect("addition recovery"); + assert!(journal.list_pending_operations().unwrap().is_empty()); + } + + let (core, accounts, _secrets, journal, public_key) = seeded(); + accounts + .remove_account(public_key) + .expect("remove metadata"); + let secrets = FailureSecretStore::default(); + secrets + .put( + public_key, + SecretKeyInput::parse( + "7e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7".to_owned(), + ) + .expect("secret"), + ) + .expect("store credential"); + secrets.fail_next(SecretStoreOperation::Delete); + let id = journal + .begin_operation(AccountOperationKind::Add, public_key, FixedClock.now()) + .expect("addition intent"); + journal + .update_operation( + id, + AccountOperationPhase::CredentialWritten, + FixedClock.now(), + None, + ) + .expect("credential phase"); + assert!( + core.recover_pending_operations(&accounts, &accounts, &secrets, &journal, &FixedClock,) + .is_err() + ); + } +} diff --git a/core/crates/studio_application/src/secrets.rs b/core/crates/studio_application/src/secrets.rs @@ -0,0 +1,308 @@ +use std::collections::BTreeMap; +use std::sync::{Mutex, MutexGuard}; + +use radroots_studio_domain::{PublicKey, SafeError, SafeErrorCode, SafeMessage, SecretKeyInput}; +use secrecy::{ExposeSecret, SecretString}; + +pub trait SecretStore: Send + Sync { + /// Stores a credential under its canonical public key without overwriting. + /// + /// # Errors + /// + /// Returns a safe duplicate or keyring error without exposing the credential. + fn put(&self, public_key: PublicKey, secret: SecretKeyInput) -> Result<(), SafeError>; + /// Loads a credential into a non-cloneable redacted boundary value. + /// + /// # Errors + /// + /// Returns a safe missing-credential or keyring error. + fn load(&self, public_key: PublicKey) -> Result<SecretKeyInput, SafeError>; + /// Reports whether a credential exists without exposing it. + /// + /// # Errors + /// + /// Returns a safe keyring error when availability cannot be determined. + fn contains(&self, public_key: PublicKey) -> Result<bool, SafeError>; + /// Deletes a credential without affecting public account metadata. + /// + /// # Errors + /// + /// Returns a safe missing-credential or keyring error. + fn delete(&self, public_key: PublicKey) -> Result<(), SafeError>; +} + +#[derive(Default)] +pub struct InMemorySecretStore { + credentials: Mutex<BTreeMap<PublicKey, SecretString>>, +} + +#[derive(Clone, Copy, Debug, Eq, Ord, PartialEq, PartialOrd)] +pub enum SecretStoreOperation { + Put, + Load, + Contains, + Delete, +} + +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub struct SecretStoreCall { + operation: SecretStoreOperation, + public_key: PublicKey, +} + +impl SecretStoreCall { + #[must_use] + pub const fn operation(self) -> SecretStoreOperation { + self.operation + } + + #[must_use] + pub const fn public_key(self) -> PublicKey { + self.public_key + } +} + +#[derive(Default)] +pub struct FailureSecretStore { + inner: InMemorySecretStore, + remaining_failures: Mutex<BTreeMap<SecretStoreOperation, usize>>, + calls: Mutex<Vec<SecretStoreCall>>, +} + +impl FailureSecretStore { + pub fn fail_next(&self, operation: SecretStoreOperation) { + *self + .remaining_failures + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .entry(operation) + .or_default() += 1; + } + + #[must_use] + pub fn calls(&self) -> Vec<SecretStoreCall> { + self.calls + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .clone() + } + + fn record_and_should_fail( + &self, + operation: SecretStoreOperation, + public_key: PublicKey, + ) -> bool { + self.calls + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .push(SecretStoreCall { + operation, + public_key, + }); + let mut failures = self + .remaining_failures + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + let remaining = failures.entry(operation).or_default(); + let should_fail = *remaining > 0; + *remaining = remaining.saturating_sub(1); + should_fail + } +} + +impl SecretStore for FailureSecretStore { + fn put(&self, public_key: PublicKey, secret: SecretKeyInput) -> Result<(), SafeError> { + if self.record_and_should_fail(SecretStoreOperation::Put, public_key) { + return Err(keyring_unavailable()); + } + self.inner.put(public_key, secret) + } + + fn load(&self, public_key: PublicKey) -> Result<SecretKeyInput, SafeError> { + if self.record_and_should_fail(SecretStoreOperation::Load, public_key) { + return Err(keyring_unavailable()); + } + self.inner.load(public_key) + } + + fn contains(&self, public_key: PublicKey) -> Result<bool, SafeError> { + if self.record_and_should_fail(SecretStoreOperation::Contains, public_key) { + return Err(keyring_unavailable()); + } + self.inner.contains(public_key) + } + + fn delete(&self, public_key: PublicKey) -> Result<(), SafeError> { + if self.record_and_should_fail(SecretStoreOperation::Delete, public_key) { + return Err(keyring_unavailable()); + } + self.inner.delete(public_key) + } +} + +impl InMemorySecretStore { + fn credentials(&self) -> MutexGuard<'_, BTreeMap<PublicKey, SecretString>> { + self.credentials + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + } +} + +impl SecretStore for InMemorySecretStore { + fn put(&self, public_key: PublicKey, secret: SecretKeyInput) -> Result<(), SafeError> { + let mut credentials = self.credentials(); + if credentials.contains_key(&public_key) { + return Err(credential_exists()); + } + let value = secret.with_exposed_secret(ToOwned::to_owned); + credentials.insert(public_key, SecretString::from(value)); + Ok(()) + } + + fn load(&self, public_key: PublicKey) -> Result<SecretKeyInput, SafeError> { + let credentials = self.credentials(); + let secret = credentials + .get(&public_key) + .ok_or_else(credential_missing)?; + SecretKeyInput::parse(secret.expose_secret().to_owned()).map_err(|_| credential_missing()) + } + + fn contains(&self, public_key: PublicKey) -> Result<bool, SafeError> { + Ok(self.credentials().contains_key(&public_key)) + } + + fn delete(&self, public_key: PublicKey) -> Result<(), SafeError> { + self.credentials() + .remove(&public_key) + .map(|_| ()) + .ok_or_else(credential_missing) + } +} + +const fn credential_exists() -> SafeError { + SafeError::new( + SafeErrorCode::AccountAlreadyExists, + SafeMessage::new("The Nostr account credential already exists."), + ) +} + +const fn credential_missing() -> SafeError { + SafeError::new( + SafeErrorCode::CredentialMissing, + SafeMessage::new("The Nostr account credential is missing."), + ) +} + +const fn keyring_unavailable() -> SafeError { + SafeError::new( + SafeErrorCode::KeyringUnavailable, + SafeMessage::new("The operating system credential store is unavailable."), + ) +} + +#[cfg(test)] +mod tests { + use radroots_studio_domain::{PublicKey, SafeErrorCode, SecretKeyInput}; + + use super::{FailureSecretStore, InMemorySecretStore, SecretStore, SecretStoreOperation}; + + const SECRET: &str = "7e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7"; + + #[test] + fn secret_store_puts_loads_checks_and_deletes_redacted_credentials() { + let store = InMemorySecretStore::default(); + let public_key = PublicKey::from_bytes([7; 32]).expect("valid public key"); + assert!(!store.contains(public_key).expect("contains")); + store + .put( + public_key, + SecretKeyInput::parse(SECRET.to_owned()).expect("secret"), + ) + .expect("put"); + assert!(store.contains(public_key).expect("contains")); + let loaded = store.load(public_key).expect("load"); + assert_eq!(loaded.with_exposed_secret(str::len), 64); + store.delete(public_key).expect("delete"); + assert!(!store.contains(public_key).expect("contains")); + } + + #[test] + fn secret_store_rejects_duplicates_and_reports_missing_credentials() { + let store = InMemorySecretStore::default(); + let public_key = PublicKey::from_bytes([7; 32]).expect("valid public key"); + let Err(missing) = store.load(public_key) else { + panic!("missing credential was returned"); + }; + assert_eq!(missing.code(), SafeErrorCode::CredentialMissing); + store + .put( + public_key, + SecretKeyInput::parse(SECRET.to_owned()).expect("secret"), + ) + .expect("put"); + let duplicate = store + .put( + public_key, + SecretKeyInput::parse(SECRET.to_owned()).expect("secret"), + ) + .expect_err("duplicate"); + assert_eq!(duplicate.code(), SafeErrorCode::AccountAlreadyExists); + store.delete(public_key).expect("delete"); + let missing = store.delete(public_key).expect_err("missing delete"); + assert_eq!(missing.code(), SafeErrorCode::CredentialMissing); + } + + #[test] + fn failure_secret_store_injects_each_boundary_without_mutating_state() { + let store = FailureSecretStore::default(); + let public_key = PublicKey::from_bytes([7; 32]).expect("valid public key"); + store.fail_next(SecretStoreOperation::Put); + let error = store + .put( + public_key, + SecretKeyInput::parse(SECRET.to_owned()).expect("secret"), + ) + .expect_err("put failure"); + assert_eq!(error.code(), SafeErrorCode::KeyringUnavailable); + assert!(!store.contains(public_key).expect("not written")); + + store + .put( + public_key, + SecretKeyInput::parse(SECRET.to_owned()).expect("secret"), + ) + .expect("put"); + for operation in [ + SecretStoreOperation::Load, + SecretStoreOperation::Contains, + SecretStoreOperation::Delete, + ] { + store.fail_next(operation); + let error = match operation { + SecretStoreOperation::Load => store.load(public_key).map(|_| ()), + SecretStoreOperation::Contains => store.contains(public_key).map(|_| ()), + SecretStoreOperation::Delete => store.delete(public_key), + SecretStoreOperation::Put => unreachable!("put tested separately"), + } + .expect_err("injected failure"); + assert_eq!(error.code(), SafeErrorCode::KeyringUnavailable); + } + assert!(store.contains(public_key).expect("credential retained")); + } + + #[test] + fn failure_secret_store_call_log_contains_only_public_identity() { + let store = FailureSecretStore::default(); + let public_key = PublicKey::from_bytes([7; 32]).expect("valid public key"); + store + .put( + public_key, + SecretKeyInput::parse(SECRET.to_owned()).expect("secret"), + ) + .expect("put"); + let calls = store.calls(); + assert_eq!(calls[0].operation(), SecretStoreOperation::Put); + assert_eq!(calls[0].public_key(), public_key); + assert!(!format!("{calls:?}").contains(SECRET)); + } +} diff --git a/core/crates/studio_application/src/session.rs b/core/crates/studio_application/src/session.rs @@ -0,0 +1,268 @@ +use crate::{ + AccountRepository, ActiveAccountSnapshot, AppCore, AppSnapshot, AppStateRepository, Clock, + ProfileLoadState, ProfileRepository, RelayConnectionState, SecretStore, StateTransition, +}; +use radroots_studio_domain::{PublicKey, SafeError, SafeErrorCode, SafeMessage}; + +impl AppCore { + /// Drops the active session while retaining accounts, selection, and credentials. + /// + /// # Errors + /// + /// Returns a safe application-state error if the transition cannot be applied. + pub fn sign_out(&self) -> Result<AppSnapshot, SafeError> { + if matches!(self.snapshot().session(), crate::SessionState::SignedOut) { + return Ok(self.snapshot()); + } + self.apply_transition(StateTransition::SignOut) + } + + /// Validates and prepares a saved local account before replacing the active session. + /// + /// # Errors + /// + /// Returns a safe account, credential, profile-cache, persistence, or state + /// error while preserving any previously active session. + pub fn activate_account( + &self, + public_key: PublicKey, + accounts: &(impl AccountRepository + ?Sized), + app_state: &(impl AppStateRepository + ?Sized), + profiles: &(impl ProfileRepository + ?Sized), + secrets: &(impl SecretStore + ?Sized), + clock: &(impl Clock + ?Sized), + ) -> Result<AppSnapshot, SafeError> { + let account = accounts + .find_account(public_key)? + .ok_or_else(account_not_found)?; + self.apply_transition(StateTransition::BeginActivation(public_key))?; + let prepared = (|| { + let credential = secrets.load(public_key)?; + let imported = self.key_material().import(credential)?; + let (derived_public_key, _npub, canonical_secret) = imported.into_parts(); + drop(canonical_secret); + if derived_public_key != public_key { + return Err(invalid_credential()); + } + let cached = profiles.load_profile(public_key)?; + let active = ActiveAccountSnapshot::new( + account.with_last_used_at(clock.now()), + RelayConnectionState::Disconnected, + if cached.is_some() { + ProfileLoadState::Cached + } else { + ProfileLoadState::Empty + }, + cached.map(|profile| profile.candidate().metadata().clone()), + ); + accounts.update_account(active.account())?; + app_state.save_selected_account(Some(public_key))?; + Ok(active) + })(); + match prepared { + Ok(active) => { + self.apply_transition(StateTransition::ActivationSucceeded(Box::new(active))) + } + Err(error) => { + self.apply_transition(StateTransition::ActivationFailed(error))?; + Err(error) + } + } + } +} + +const fn account_not_found() -> SafeError { + SafeError::new( + SafeErrorCode::AccountNotFound, + SafeMessage::new("The account was not found."), + ) +} + +const fn invalid_credential() -> SafeError { + SafeError::new( + SafeErrorCode::InvalidSecretKey, + SafeMessage::new("The Nostr account credential is invalid."), + ) +} + +#[cfg(test)] +mod tests { + use radroots_studio_domain::{PublicKey, SafeError, SecretKeyInput, UnixTimestamp}; + + use crate::{ + AppCore, CachedProfile, Clock, InMemoryAccountRepository, InMemoryOperationJournal, + InMemorySecretStore, ProfileRefreshStatus, ProfileRepository, RelayConfiguration, + SecretStore, SessionState, + }; + + #[derive(Default)] + struct EmptyProfiles; + + impl ProfileRepository for EmptyProfiles { + fn load_profile(&self, _public_key: PublicKey) -> Result<Option<CachedProfile>, SafeError> { + Ok(None) + } + + fn save_profile(&self, _profile: &CachedProfile) -> Result<(), SafeError> { + Ok(()) + } + + fn record_refresh_status( + &self, + _public_key: PublicKey, + _refreshed_at: UnixTimestamp, + _status: ProfileRefreshStatus, + ) -> Result<(), SafeError> { + Ok(()) + } + + fn remove_profile(&self, _public_key: PublicKey) -> Result<(), SafeError> { + Ok(()) + } + } + + struct FixedClock; + + impl Clock for FixedClock { + fn now(&self) -> UnixTimestamp { + UnixTimestamp::from_seconds(30).expect("time") + } + } + + fn input(value: &str) -> SecretKeyInput { + SecretKeyInput::parse(value.to_owned()).expect("input") + } + + #[test] + fn activate_account_switches_only_after_candidate_is_ready() { + let core = AppCore::in_memory(RelayConfiguration::default()); + let accounts = InMemoryAccountRepository::default(); + let secrets = InMemorySecretStore::default(); + let journal = InMemoryOperationJournal::default(); + let profiles = EmptyProfiles; + core.bootstrap().expect("bootstrap"); + let first = core + .import_secret_key( + input("7e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7"), + &accounts, + &accounts, + &secrets, + &journal, + &FixedClock, + ) + .expect("first") + .account() + .public_key(); + let second = core + .import_secret_key( + input("1111111111111111111111111111111111111111111111111111111111111111"), + &accounts, + &accounts, + &secrets, + &journal, + &FixedClock, + ) + .expect("second") + .account() + .public_key(); + core.activate_account( + first, + &accounts, + &accounts, + &profiles, + &secrets, + &FixedClock, + ) + .expect("activate first"); + assert_eq!(core.snapshot().session(), SessionState::Active); + assert_eq!( + core.snapshot() + .active_account() + .map(|active| active.account().public_key()), + Some(first) + ); + + secrets.delete(second).expect("remove second credential"); + let error = core + .activate_account( + second, + &accounts, + &accounts, + &profiles, + &secrets, + &FixedClock, + ) + .expect_err("missing credential"); + assert_eq!( + error.code(), + radroots_studio_domain::SafeErrorCode::CredentialMissing + ); + secrets + .put( + second, + input("7e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7"), + ) + .expect("mismatched credential"); + let invalid = core + .activate_account( + second, + &accounts, + &accounts, + &profiles, + &secrets, + &FixedClock, + ) + .expect_err("mismatched credential"); + assert_eq!( + invalid.code(), + radroots_studio_domain::SafeErrorCode::InvalidSecretKey + ); + assert_eq!(core.snapshot().session(), SessionState::Active); + assert_eq!( + core.snapshot() + .active_account() + .map(|active| active.account().public_key()), + Some(first) + ); + } + + #[test] + fn sign_out_retains_saved_account_selection_and_credential() { + let core = AppCore::in_memory(RelayConfiguration::default()); + let accounts = InMemoryAccountRepository::default(); + let secrets = InMemorySecretStore::default(); + let journal = InMemoryOperationJournal::default(); + let profiles = EmptyProfiles; + core.bootstrap().expect("bootstrap"); + let public_key = core + .import_secret_key( + input("7e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7"), + &accounts, + &accounts, + &secrets, + &journal, + &FixedClock, + ) + .expect("import") + .account() + .public_key(); + core.activate_account( + public_key, + &accounts, + &accounts, + &profiles, + &secrets, + &FixedClock, + ) + .expect("activate"); + + let signed_out = core.sign_out().expect("sign out"); + let repeated = core.sign_out().expect("idempotent sign out"); + assert_eq!(signed_out, repeated); + assert_eq!(signed_out.session(), SessionState::SignedOut); + assert!(signed_out.active_account().is_none()); + assert_eq!(signed_out.accounts().len(), 1); + assert_eq!(signed_out.selected_account(), Some(public_key)); + assert!(secrets.contains(public_key).expect("credential retained")); + } +} diff --git a/core/crates/studio_application/src/snapshot.rs b/core/crates/studio_application/src/snapshot.rs @@ -0,0 +1,479 @@ +use std::collections::HashSet; + +use radroots_studio_domain::{ + AccountSummary, ProfileMetadata, PublicKey, RelayUrl, SafeError, SafeErrorCode, SafeMessage, +}; + +pub const MAX_CONFIGURED_RELAYS: usize = 16; + +#[derive(Clone, Copy, Debug, Default, Eq, Ord, PartialEq, PartialOrd)] +pub struct SnapshotRevision(u64); + +impl SnapshotRevision { + #[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 enum AppLifecycle { + Booting, + Ready, + Fatal(SafeError), +} + +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum SessionState { + SignedOut, + Activating(PublicKey), + Active, + SigningOut, + Failed(SafeError), +} + +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum RelayConnectionState { + Disconnected, + Connecting, + Connected, + Degraded, + Error(SafeError), +} + +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum ProfileLoadState { + Empty, + Loading, + Cached, + Fresh, + Error(SafeError), +} + +#[derive(Clone, Debug, Default, Eq, PartialEq)] +pub struct RelayConfiguration(Vec<RelayUrl>); + +impl RelayConfiguration { + /// Creates a bounded, explicitly classified relay configuration. + /// + /// # Errors + /// + /// Returns a safe configuration error before runtime or network work when + /// the relay count exceeds the Studio policy. + pub fn new(relays: Vec<RelayUrl>) -> Result<Self, SafeError> { + if relays.len() > MAX_CONFIGURED_RELAYS { + return Err(relay_limit_exceeded()); + } + Ok(Self(relays)) + } + + #[must_use] + pub fn relays(&self) -> &[RelayUrl] { + &self.0 + } +} + +const fn relay_limit_exceeded() -> SafeError { + SafeError::new( + SafeErrorCode::InvalidRelayConfiguration, + SafeMessage::new("The Nostr relay configuration exceeds its limit."), + ) +} + +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct ActiveAccountSnapshot { + account: AccountSummary, + relay_state: RelayConnectionState, + profile_state: ProfileLoadState, + profile: Option<ProfileMetadata>, +} + +impl ActiveAccountSnapshot { + #[must_use] + pub const fn new( + account: AccountSummary, + relay_state: RelayConnectionState, + profile_state: ProfileLoadState, + profile: Option<ProfileMetadata>, + ) -> Self { + Self { + account, + relay_state, + profile_state, + profile, + } + } + + #[must_use] + pub const fn account(&self) -> &AccountSummary { + &self.account + } + + #[must_use] + pub const fn relay_state(&self) -> RelayConnectionState { + self.relay_state + } + + #[must_use] + pub const fn profile_state(&self) -> ProfileLoadState { + self.profile_state + } + + #[must_use] + pub const fn profile(&self) -> Option<&ProfileMetadata> { + self.profile.as_ref() + } +} + +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct AppSnapshot { + revision: SnapshotRevision, + lifecycle: AppLifecycle, + relay_configuration: RelayConfiguration, + accounts: Vec<AccountSummary>, + selected_account: Option<PublicKey>, + session: SessionState, + active_account: Option<ActiveAccountSnapshot>, + recoverable_problem: Option<SafeError>, +} + +impl AppSnapshot { + #[must_use] + pub fn booting() -> Self { + Self { + revision: SnapshotRevision::initial(), + lifecycle: AppLifecycle::Booting, + relay_configuration: RelayConfiguration::default(), + accounts: Vec::new(), + selected_account: None, + session: SessionState::SignedOut, + active_account: None, + recoverable_problem: None, + } + } + + #[must_use] + pub fn fatal( + revision: SnapshotRevision, + relay_configuration: RelayConfiguration, + error: SafeError, + ) -> Self { + Self { + revision, + lifecycle: AppLifecycle::Fatal(error), + relay_configuration, + accounts: Vec::new(), + selected_account: None, + session: SessionState::SignedOut, + active_account: None, + recoverable_problem: None, + } + } + + /// Constructs a ready immutable snapshot after validating state invariants. + /// + /// # Errors + /// + /// Returns a safe invalid-state error for duplicate accounts, invalid + /// selection, or inconsistent active-session state. + pub fn ready( + revision: SnapshotRevision, + relay_configuration: RelayConfiguration, + accounts: Vec<AccountSummary>, + selected_account: Option<PublicKey>, + session: SessionState, + active_account: Option<ActiveAccountSnapshot>, + recoverable_problem: Option<SafeError>, + ) -> Result<Self, SafeError> { + validate_snapshot( + &accounts, + selected_account, + session, + active_account.as_ref(), + )?; + Ok(Self { + revision, + lifecycle: AppLifecycle::Ready, + relay_configuration, + accounts, + selected_account, + session, + active_account, + recoverable_problem, + }) + } + + #[must_use] + pub const fn revision(&self) -> SnapshotRevision { + self.revision + } + + #[must_use] + pub const fn lifecycle(&self) -> AppLifecycle { + self.lifecycle + } + + #[must_use] + pub const fn relay_configuration(&self) -> &RelayConfiguration { + &self.relay_configuration + } + + #[must_use] + pub fn accounts(&self) -> &[AccountSummary] { + &self.accounts + } + + #[must_use] + pub const fn selected_account(&self) -> Option<PublicKey> { + self.selected_account + } + + #[must_use] + pub const fn session(&self) -> SessionState { + self.session + } + + #[must_use] + pub const fn active_account(&self) -> Option<&ActiveAccountSnapshot> { + self.active_account.as_ref() + } + + #[must_use] + pub const fn recoverable_problem(&self) -> Option<SafeError> { + self.recoverable_problem + } +} + +fn validate_snapshot( + accounts: &[AccountSummary], + selected_account: Option<PublicKey>, + session: SessionState, + active_account: Option<&ActiveAccountSnapshot>, +) -> Result<(), SafeError> { + let unique_accounts = accounts + .iter() + .map(AccountSummary::public_key) + .collect::<HashSet<_>>(); + if unique_accounts.len() != accounts.len() + || (accounts.is_empty() != selected_account.is_none()) + || selected_account.is_some_and(|key| !unique_accounts.contains(&key)) + || active_account + .is_some_and(|active| !unique_accounts.contains(&active.account().public_key())) + || (matches!(session, SessionState::Active) && active_account.is_none()) + || (matches!(session, SessionState::SignedOut) && active_account.is_some()) + { + return Err(invalid_snapshot()); + } + Ok(()) +} + +const fn invalid_snapshot() -> SafeError { + SafeError::new( + SafeErrorCode::InvalidApplicationState, + SafeMessage::new("The application state is invalid."), + ) +} + +#[cfg(test)] +mod tests { + use radroots_studio_domain::{ + AccountCreatedAt, AccountIdentity, AccountSummary, BindingAvailability, LocalSignerBinding, + RelayDestinationPolicy, RelayUrl, SafeErrorCode, UnixTimestamp, + }; + + use super::{ + ActiveAccountSnapshot, AppLifecycle, AppSnapshot, ProfileLoadState, RelayConfiguration, + RelayConnectionState, SessionState, SnapshotRevision, + }; + + fn account(key_byte: u8) -> AccountSummary { + let public_key = + crate::test_support::valid_test_public_key(key_byte).expect("valid public key"); + AccountSummary::new( + AccountIdentity::derive(public_key).expect("identity"), + LocalSignerBinding::new(public_key, BindingAvailability::Available), + None, + AccountCreatedAt::new(UnixTimestamp::from_seconds(1).expect("valid time")), + None, + ) + .expect("account") + } + + #[test] + fn snapshot_boots_empty_and_secret_free() { + let snapshot = AppSnapshot::booting(); + let debug = format!("{snapshot:?}"); + + assert_eq!(snapshot.revision(), SnapshotRevision::initial()); + assert_eq!(snapshot.lifecycle(), AppLifecycle::Booting); + assert_eq!(snapshot.session(), SessionState::SignedOut); + assert!(snapshot.accounts().is_empty()); + assert!(snapshot.selected_account().is_none()); + assert!(snapshot.active_account().is_none()); + assert!(snapshot.relay_configuration().relays().is_empty()); + assert!(snapshot.recoverable_problem().is_none()); + assert!(!debug.contains("nsec1")); + assert!(!debug.contains(&"11".repeat(32))); + } + + #[test] + fn relay_configuration_rejects_excess_targets_before_runtime_work() { + let relays = (0..=super::MAX_CONFIGURED_RELAYS) + .map(|index| { + RelayUrl::parse( + format!("wss://relay-{index}.example").as_str(), + RelayDestinationPolicy::Public, + ) + .expect("relay") + }) + .collect(); + let error = RelayConfiguration::new(relays).expect_err("relay limit"); + assert_eq!(error.code(), SafeErrorCode::InvalidRelayConfiguration); + } + + #[test] + fn revision_helper_is_monotonic_and_checked() { + assert_eq!( + SnapshotRevision::initial() + .next() + .map(SnapshotRevision::value), + Some(1) + ); + assert_eq!(SnapshotRevision::from_value(u64::MAX).next(), None); + } + + #[test] + fn ready_snapshot_requires_valid_selection_and_active_session() { + let first = account(1); + let second = account(2); + let active = ActiveAccountSnapshot::new( + second.clone(), + RelayConnectionState::Disconnected, + ProfileLoadState::Empty, + None, + ); + let valid = AppSnapshot::ready( + SnapshotRevision::from_value(1), + RelayConfiguration::default(), + vec![first.clone(), second.clone()], + Some(first.public_key()), + SessionState::Active, + Some(active), + None, + ) + .expect("valid ready snapshot"); + + assert_eq!(valid.lifecycle(), AppLifecycle::Ready); + assert_eq!(valid.selected_account(), Some(first.public_key())); + assert_eq!( + valid + .active_account() + .map(|value| value.account().public_key()), + Some(second.public_key()) + ); + + assert!( + AppSnapshot::ready( + SnapshotRevision::initial(), + RelayConfiguration::default(), + vec![first.clone(), first], + Some(second.public_key()), + SessionState::SignedOut, + None, + None, + ) + .is_err() + ); + assert!( + AppSnapshot::ready( + SnapshotRevision::initial(), + RelayConfiguration::default(), + vec![second.clone()], + Some(second.public_key()), + SessionState::Active, + None, + None, + ) + .is_err() + ); + + let missing = account(3); + for result in [ + AppSnapshot::ready( + SnapshotRevision::initial(), + RelayConfiguration::default(), + Vec::new(), + Some(missing.public_key()), + SessionState::SignedOut, + None, + None, + ), + AppSnapshot::ready( + SnapshotRevision::initial(), + RelayConfiguration::default(), + vec![second.clone()], + None, + SessionState::SignedOut, + None, + None, + ), + AppSnapshot::ready( + SnapshotRevision::initial(), + RelayConfiguration::default(), + vec![second.clone()], + Some(missing.public_key()), + SessionState::SignedOut, + None, + None, + ), + AppSnapshot::ready( + SnapshotRevision::initial(), + RelayConfiguration::default(), + vec![second.clone()], + Some(second.public_key()), + SessionState::SignedOut, + Some(ActiveAccountSnapshot::new( + missing, + RelayConnectionState::Disconnected, + ProfileLoadState::Empty, + None, + )), + None, + ), + AppSnapshot::ready( + SnapshotRevision::initial(), + RelayConfiguration::default(), + vec![second.clone()], + Some(second.public_key()), + SessionState::SignedOut, + Some(ActiveAccountSnapshot::new( + second, + RelayConnectionState::Disconnected, + ProfileLoadState::Empty, + None, + )), + None, + ), + ] { + assert!(result.is_err()); + } + } +} diff --git a/core/crates/studio_application/src/state_machine.rs b/core/crates/studio_application/src/state_machine.rs @@ -0,0 +1,616 @@ +use radroots_studio_domain::{AccountSummary, PublicKey, SafeError, SafeErrorCode, SafeMessage}; + +use crate::{ActiveAccountSnapshot, AppLifecycle, AppSnapshot, RelayConfiguration, SessionState}; + +#[derive(Clone, Debug, Eq, PartialEq)] +pub enum StateTransition { + Bootstrap, + BootstrapRegistry { + accounts: Vec<AccountSummary>, + selected: Option<PublicKey>, + }, + Fatal(SafeError), + ReplaceRegistry { + accounts: Vec<AccountSummary>, + selected: Option<PublicKey>, + }, + ReplaceRegistryPreservingSession { + accounts: Vec<AccountSummary>, + selected: Option<PublicKey>, + }, + Select(PublicKey), + BeginActivation(PublicKey), + ActivationSucceeded(Box<ActiveAccountSnapshot>), + ActivationFailed(SafeError), + UpdateActiveAccount { + expected: PublicKey, + active_account: Box<ActiveAccountSnapshot>, + problem: Option<SafeError>, + }, + SignOut, + SetProblem(Option<SafeError>), +} + +#[derive(Clone)] +struct PreviousSession { + session: SessionState, + active_account: Option<ActiveAccountSnapshot>, +} + +pub struct StateMachine { + snapshot: AppSnapshot, + pending_activation: Option<(PublicKey, PreviousSession)>, +} + +impl StateMachine { + #[must_use] + pub fn booting() -> Self { + Self { + snapshot: AppSnapshot::booting(), + pending_activation: None, + } + } + + #[must_use] + pub const fn snapshot(&self) -> &AppSnapshot { + &self.snapshot + } + + /// Applies one deterministic state transition and returns the new snapshot. + /// + /// # Errors + /// + /// Returns a safe application error when the transition violates account, + /// revision, activation, or snapshot invariants. + pub fn apply( + &mut self, + transition: StateTransition, + relay_configuration: &RelayConfiguration, + ) -> Result<AppSnapshot, SafeError> { + let next_revision = self + .snapshot + .revision() + .next() + .ok_or_else(invalid_application_state)?; + + let next = match transition { + StateTransition::Bootstrap => self.bootstrap(next_revision, relay_configuration)?, + StateTransition::BootstrapRegistry { accounts, selected } => { + self.bootstrap_registry(next_revision, relay_configuration, accounts, selected)? + } + StateTransition::Fatal(error) => { + AppSnapshot::fatal(next_revision, relay_configuration.clone(), error) + } + StateTransition::ReplaceRegistry { accounts, selected } => { + self.replace_registry(next_revision, accounts, selected)? + } + StateTransition::ReplaceRegistryPreservingSession { accounts, selected } => { + self.replace_registry_preserving_session(next_revision, accounts, selected)? + } + StateTransition::Select(public_key) => self.select(next_revision, public_key)?, + StateTransition::BeginActivation(public_key) => { + self.begin_activation(next_revision, public_key)? + } + StateTransition::ActivationSucceeded(active_account) => { + self.activation_succeeded(next_revision, *active_account)? + } + StateTransition::ActivationFailed(problem) => { + self.activation_failed(next_revision, problem)? + } + StateTransition::UpdateActiveAccount { + expected, + active_account, + problem, + } => self.update_active_account(next_revision, expected, *active_account, problem)?, + StateTransition::SignOut => self.sign_out(next_revision)?, + StateTransition::SetProblem(problem) => self.copy_ready( + next_revision, + self.snapshot.selected_account(), + self.snapshot.session(), + self.snapshot.active_account().cloned(), + problem, + )?, + }; + self.snapshot = next.clone(); + Ok(next) + } + + fn bootstrap( + &self, + revision: crate::SnapshotRevision, + relay_configuration: &RelayConfiguration, + ) -> Result<AppSnapshot, SafeError> { + if !matches!(self.snapshot.lifecycle(), AppLifecycle::Booting) { + return Ok(self.snapshot.clone()); + } + AppSnapshot::ready( + revision, + relay_configuration.clone(), + Vec::new(), + None, + SessionState::SignedOut, + None, + None, + ) + } + + fn bootstrap_registry( + &self, + revision: crate::SnapshotRevision, + relay_configuration: &RelayConfiguration, + accounts: Vec<AccountSummary>, + selected: Option<PublicKey>, + ) -> Result<AppSnapshot, SafeError> { + if !matches!(self.snapshot.lifecycle(), AppLifecycle::Booting) { + return Ok(self.snapshot.clone()); + } + AppSnapshot::ready( + revision, + relay_configuration.clone(), + accounts, + selected, + SessionState::SignedOut, + None, + None, + ) + } + + fn replace_registry( + &mut self, + revision: crate::SnapshotRevision, + accounts: Vec<AccountSummary>, + selected: Option<PublicKey>, + ) -> Result<AppSnapshot, SafeError> { + self.pending_activation = None; + AppSnapshot::ready( + revision, + self.snapshot.relay_configuration().clone(), + accounts, + selected, + SessionState::SignedOut, + None, + None, + ) + } + + fn replace_registry_preserving_session( + &mut self, + revision: crate::SnapshotRevision, + accounts: Vec<AccountSummary>, + selected: Option<PublicKey>, + ) -> Result<AppSnapshot, SafeError> { + self.pending_activation = None; + AppSnapshot::ready( + revision, + self.snapshot.relay_configuration().clone(), + accounts, + selected, + self.snapshot.session(), + self.snapshot.active_account().cloned(), + None, + ) + } + + fn select( + &self, + revision: crate::SnapshotRevision, + public_key: PublicKey, + ) -> Result<AppSnapshot, SafeError> { + self.require_account(public_key)?; + self.copy_ready( + revision, + Some(public_key), + self.snapshot.session(), + self.snapshot.active_account().cloned(), + None, + ) + } + + fn begin_activation( + &mut self, + revision: crate::SnapshotRevision, + public_key: PublicKey, + ) -> Result<AppSnapshot, SafeError> { + self.require_account(public_key)?; + if self.pending_activation.is_some() { + return Err(invalid_application_state()); + } + self.pending_activation = Some(( + public_key, + PreviousSession { + session: self.snapshot.session(), + active_account: self.snapshot.active_account().cloned(), + }, + )); + self.copy_ready( + revision, + self.snapshot.selected_account(), + SessionState::Activating(public_key), + self.snapshot.active_account().cloned(), + None, + ) + } + + fn activation_succeeded( + &mut self, + revision: crate::SnapshotRevision, + active_account: ActiveAccountSnapshot, + ) -> Result<AppSnapshot, SafeError> { + let Some((target, _previous)) = self.pending_activation.as_ref() else { + return Err(invalid_application_state()); + }; + if active_account.account().public_key() != *target { + return Err(invalid_application_state()); + } + let target = *target; + self.pending_activation = None; + self.copy_ready( + revision, + Some(target), + SessionState::Active, + Some(active_account), + None, + ) + } + + fn activation_failed( + &mut self, + revision: crate::SnapshotRevision, + problem: SafeError, + ) -> Result<AppSnapshot, SafeError> { + let Some((_target, previous)) = self.pending_activation.take() else { + return Err(invalid_application_state()); + }; + self.copy_ready( + revision, + self.snapshot.selected_account(), + previous.session, + previous.active_account, + Some(problem), + ) + } + + fn sign_out(&mut self, revision: crate::SnapshotRevision) -> Result<AppSnapshot, SafeError> { + self.pending_activation = None; + self.copy_ready( + revision, + self.snapshot.selected_account(), + SessionState::SignedOut, + None, + None, + ) + } + + fn update_active_account( + &self, + revision: crate::SnapshotRevision, + expected: PublicKey, + active_account: ActiveAccountSnapshot, + problem: Option<SafeError>, + ) -> Result<AppSnapshot, SafeError> { + if !matches!(self.snapshot.session(), SessionState::Active) + || self + .snapshot + .active_account() + .map(|active| active.account().public_key()) + != Some(expected) + || active_account.account().public_key() != expected + { + return Err(invalid_application_state()); + } + self.copy_ready( + revision, + self.snapshot.selected_account(), + SessionState::Active, + Some(active_account), + problem, + ) + } + + fn require_account(&self, public_key: PublicKey) -> Result<(), SafeError> { + if self + .snapshot + .accounts() + .iter() + .any(|account| account.public_key() == public_key) + { + Ok(()) + } else { + Err(account_not_found()) + } + } + + fn copy_ready( + &self, + revision: crate::SnapshotRevision, + selected_account: Option<PublicKey>, + session: SessionState, + active_account: Option<ActiveAccountSnapshot>, + recoverable_problem: Option<SafeError>, + ) -> Result<AppSnapshot, SafeError> { + AppSnapshot::ready( + revision, + self.snapshot.relay_configuration().clone(), + self.snapshot.accounts().to_vec(), + selected_account, + session, + active_account, + recoverable_problem, + ) + } +} + +const fn invalid_application_state() -> SafeError { + SafeError::new( + SafeErrorCode::InvalidApplicationState, + SafeMessage::new("The application state is invalid."), + ) +} + +const fn account_not_found() -> SafeError { + SafeError::new( + SafeErrorCode::AccountNotFound, + SafeMessage::new("The account was not found."), + ) +} + +#[cfg(test)] +mod tests { + use radroots_studio_domain::{ + AccountCreatedAt, AccountIdentity, AccountSummary, BindingAvailability, LocalSignerBinding, + SafeError, SafeErrorCode, SafeMessage, UnixTimestamp, + }; + + use crate::{ + ActiveAccountSnapshot, ProfileLoadState, RelayConfiguration, RelayConnectionState, + SessionState, StateMachine, StateTransition, + }; + + fn account(key_byte: u8) -> AccountSummary { + let public_key = + crate::test_support::valid_test_public_key(key_byte).expect("valid public key"); + AccountSummary::new( + AccountIdentity::derive(public_key).expect("identity"), + LocalSignerBinding::new(public_key, BindingAvailability::Available), + None, + AccountCreatedAt::new(UnixTimestamp::from_seconds(1).expect("valid time")), + None, + ) + .expect("account") + } + + fn active(account: AccountSummary) -> ActiveAccountSnapshot { + ActiveAccountSnapshot::new( + account, + RelayConnectionState::Disconnected, + ProfileLoadState::Empty, + None, + ) + } + + #[test] + fn state_machine_command_trace_preserves_working_session_on_failed_replacement() { + let first = account(1); + let second = account(2); + let mut machine = StateMachine::booting(); + let relays = RelayConfiguration::default(); + let problem = SafeError::new( + SafeErrorCode::CredentialMissing, + SafeMessage::new("The account credential is missing."), + ); + + machine + .apply(StateTransition::Bootstrap, &relays) + .expect("bootstrap"); + machine + .apply( + StateTransition::ReplaceRegistry { + accounts: vec![first.clone(), second.clone()], + selected: Some(first.public_key()), + }, + &relays, + ) + .expect("load registry"); + machine + .apply( + StateTransition::BeginActivation(first.public_key()), + &relays, + ) + .expect("begin first activation"); + machine + .apply( + StateTransition::ActivationSucceeded(Box::new(active(first.clone()))), + &relays, + ) + .expect("activate first"); + machine + .apply(StateTransition::Select(second.public_key()), &relays) + .expect("select second"); + let pending = machine + .apply( + StateTransition::BeginActivation(second.public_key()), + &relays, + ) + .expect("begin replacement"); + let restored = machine + .apply(StateTransition::ActivationFailed(problem), &relays) + .expect("fail replacement"); + + assert_eq!( + pending.session(), + SessionState::Activating(second.public_key()) + ); + assert_eq!( + pending + .active_account() + .map(|value| value.account().public_key()), + Some(first.public_key()) + ); + assert_eq!(restored.session(), SessionState::Active); + assert_eq!(restored.selected_account(), Some(second.public_key())); + assert_eq!( + restored + .active_account() + .map(|value| value.account().public_key()), + Some(first.public_key()) + ); + assert_eq!(restored.recoverable_problem(), Some(problem)); + assert_eq!(restored.revision().value(), 7); + } + + #[test] + fn state_machine_rejects_missing_targets_and_signs_out_without_deleting() { + let account = account(1); + let mut machine = StateMachine::booting(); + let relays = RelayConfiguration::default(); + machine + .apply(StateTransition::Bootstrap, &relays) + .expect("bootstrap"); + machine + .apply( + StateTransition::ReplaceRegistry { + accounts: vec![account.clone()], + selected: Some(account.public_key()), + }, + &relays, + ) + .expect("load registry"); + + let error = machine + .apply( + StateTransition::Select( + crate::test_support::valid_test_public_key(9).expect("valid public key"), + ), + &relays, + ) + .expect_err("missing account"); + assert_eq!(error.code(), SafeErrorCode::AccountNotFound); + + machine + .apply( + StateTransition::BeginActivation(account.public_key()), + &relays, + ) + .expect("begin activation"); + machine + .apply( + StateTransition::ActivationSucceeded(Box::new(active(account.clone()))), + &relays, + ) + .expect("activate"); + let signed_out = machine + .apply(StateTransition::SignOut, &relays) + .expect("sign out"); + + assert_eq!(signed_out.accounts(), &[account]); + assert_eq!(signed_out.session(), SessionState::SignedOut); + assert!(signed_out.active_account().is_none()); + } + + #[test] + fn activation_state_policy_rejects_every_stale_or_mismatched_transition() { + let first = account(1); + let second = account(2); + let relays = RelayConfiguration::default(); + let problem = SafeError::new( + SafeErrorCode::CredentialMissing, + SafeMessage::new("The account credential is missing."), + ); + let mut machine = StateMachine::booting(); + machine + .apply( + StateTransition::BootstrapRegistry { + accounts: vec![first.clone(), second.clone()], + selected: Some(first.public_key()), + }, + &relays, + ) + .expect("registry"); + let unchanged = machine + .apply( + StateTransition::BootstrapRegistry { + accounts: Vec::new(), + selected: None, + }, + &relays, + ) + .expect("repeated bootstrap is idempotent"); + assert_eq!(unchanged.accounts().len(), 2); + + assert!( + machine + .apply( + StateTransition::ActivationSucceeded(Box::new(active(first.clone()))), + &relays, + ) + .is_err() + ); + assert!( + machine + .apply(StateTransition::ActivationFailed(problem), &relays) + .is_err() + ); + machine + .apply( + StateTransition::BeginActivation(first.public_key()), + &relays, + ) + .expect("begin activation"); + assert!( + machine + .apply( + StateTransition::BeginActivation(second.public_key()), + &relays, + ) + .is_err() + ); + assert!( + machine + .apply( + StateTransition::ActivationSucceeded(Box::new(active(second.clone()))), + &relays, + ) + .is_err() + ); + machine + .apply( + StateTransition::ActivationSucceeded(Box::new(active(first.clone()))), + &relays, + ) + .expect("activate first"); + + for (expected, candidate) in [ + (second.public_key(), first.clone()), + (first.public_key(), second.clone()), + ] { + assert!( + machine + .apply( + StateTransition::UpdateActiveAccount { + expected, + active_account: Box::new(active(candidate)), + problem: None, + }, + &relays, + ) + .is_err() + ); + } + + machine + .apply(StateTransition::SignOut, &relays) + .expect("sign out"); + assert!( + machine + .apply( + StateTransition::UpdateActiveAccount { + expected: first.public_key(), + active_account: Box::new(active(first)), + problem: None, + }, + &relays, + ) + .is_err() + ); + } +} diff --git a/core/crates/studio_application/src/test_support.rs b/core/crates/studio_application/src/test_support.rs @@ -0,0 +1,58 @@ +use std::sync::atomic::{AtomicU8, Ordering}; + +use radroots_studio_domain::{ + Npub, Nsec, PublicKey, SafeError, SafeErrorCode, SafeMessage, SecretKeyInput, +}; + +use crate::{GeneratedKeyMaterial, ImportedKeyMaterial, KeyMaterialProvider}; + +#[derive(Default)] +pub(crate) struct TestKeyMaterialProvider { + next: AtomicU8, +} + +impl KeyMaterialProvider for TestKeyMaterialProvider { + fn generate(&self) -> Result<GeneratedKeyMaterial, SafeError> { + let public_key = (0..=u8::MAX) + .find_map(|_| { + let candidate = self.next.fetch_add(1, Ordering::Relaxed).wrapping_add(9); + PublicKey::from_bytes([candidate; 32]).ok() + }) + .ok_or_else(invalid_secret_key)?; + let secret_byte = public_key.as_bytes()[0]; + Ok(GeneratedKeyMaterial::new( + public_key, + Npub::derive(public_key)?, + SecretKeyInput::parse(format!("{secret_byte:02x}").repeat(32))?, + Nsec::from_encoded( + "nsec1vl029mgpspedva04g90vltkh6fvh240zqtv9k0t9af8935ke9laqsnlfe5".to_owned(), + )?, + )) + } + + fn import(&self, input: SecretKeyInput) -> Result<ImportedKeyMaterial, SafeError> { + let discriminator = input.with_exposed_secret(|value| value.as_bytes()[0]); + if input.with_exposed_secret(|value| value.starts_with("nsec1qq")) { + return Err(invalid_secret_key()); + } + let public_key = valid_test_public_key(discriminator)?; + Ok(ImportedKeyMaterial::new( + public_key, + Npub::derive(public_key)?, + input, + )) + } +} + +pub(crate) fn valid_test_public_key(discriminator: u8) -> Result<PublicKey, SafeError> { + (0..=u8::MAX) + .find_map(|offset| PublicKey::from_bytes([discriminator.wrapping_add(offset); 32]).ok()) + .ok_or_else(invalid_secret_key) +} + +const fn invalid_secret_key() -> SafeError { + SafeError::new( + SafeErrorCode::InvalidSecretKey, + SafeMessage::new("The Nostr secret key is invalid."), + ) +} diff --git a/core/crates/studio_application/tests/redaction.rs b/core/crates/studio_application/tests/redaction.rs @@ -0,0 +1,48 @@ +use radroots_studio_application::{ + AppSnapshot, RelayConfiguration, SessionState, SnapshotRevision, +}; +use radroots_studio_domain::{ + AccountCreatedAt, AccountIdentity, AccountSummary, BindingAvailability, LocalSignerBinding, + PublicKey, SafeError, SafeErrorCode, SafeMessage, UnixTimestamp, +}; + +const SECRET_HEX: &str = "1111111111111111111111111111111111111111111111111111111111111111"; +const SECRET_NSEC: &str = "nsec1vl029mgpspedva04g90vltkh6fvh240zqtv9k0t9af8935ke9laqsnlfe5"; +fn assert_redacted(text: &str) { + assert!(!text.contains(SECRET_HEX)); + assert!(!text.contains(SECRET_NSEC)); + assert!(!text.contains("nsec1")); +} + +#[test] +fn redaction_guards_public_snapshot_and_safe_error_debug() { + let account = AccountSummary::new( + AccountIdentity::derive(PublicKey::from_bytes([7; 32]).expect("valid public key")) + .expect("identity"), + LocalSignerBinding::new( + PublicKey::from_bytes([7; 32]).expect("valid public key"), + BindingAvailability::Available, + ), + None, + AccountCreatedAt::new(UnixTimestamp::from_seconds(1).expect("time")), + None, + ) + .expect("account"); + let snapshot = AppSnapshot::ready( + SnapshotRevision::from_value(1), + RelayConfiguration::default(), + vec![account.clone()], + Some(account.public_key()), + SessionState::SignedOut, + None, + None, + ) + .expect("snapshot"); + let error = SafeError::new( + SafeErrorCode::KeyringUnavailable, + SafeMessage::new("The operating system credential store is unavailable."), + ); + + assert_redacted(&format!("{snapshot:?}")); + assert_redacted(&format!("{error:?} {error}")); +}