commit 653dd4ee17e9577f4e095aeb0e7bb0321fb7ecbe parent b19555ab50f80899c51bd3ed9ca5bdc2b5d8227d Author: triesap <tyson@radroots.org> Date: Wed, 26 Aug 2026 04:16:25 +0000 refactor(runtime): seal relay and host lifecycle - replace duplicate relay policy and production clients with the pinned Radroots transport - validate bounded FFI inputs before I/O and keep diagnostics and correlations secret-safe - own cancellation-resumable runtime, observer, actor, and bounded keyring shutdown - bind the source lock, compatibility hash, policy guards, and native platform checks Diffstat:
33 files changed, 1678 insertions(+), 857 deletions(-)
diff --git a/AGENTS.md b/AGENTS.md @@ -112,6 +112,14 @@ substitute. - Keep lifecycle and coroutine work structured, cancelable, and scoped. Avoid hidden workers, process-global mutable state, unbounded retries, blocking UI work, and external mutation inferred from environment state. +- Relay parsing, destination policy, DNS admission, and connection ownership + come from the pinned `radroots_transport_nostr` boundary. Product code may + select a governed profile and verify product events, but must not recreate + relay URL policy or open a second production `nostr-sdk` client. +- The native host owns its Tokio runtime for exactly one application-core + lifetime. Runtime close is explicit, idempotent, and cancellation-resumable; + observer work and the bounded keyring worker must finish before terminal + close is reported. Authoritative locks fail closed on poison. - Services-hardening changes use the approved target-state contracts. Do not add compatibility aliases, dual reads, dual writes, or fallback behavior for prototype surfaces removed by the clean-slate refactor. diff --git a/README.md b/README.md @@ -62,6 +62,16 @@ backup capability, closes the live host, uses the governed marker protocol, and reopens recovered state before returning. There is no arbitrary database repair or pathname-only restore authority. +Relay endpoints are explicit inputs validated by the pinned Radroots Nostr +transport policy before any socket work. HarvestCircle owns profile selection +and signature-verified kind-0 interpretation, while the shared transport owns +relay URL, destination, DNS, connection, and bounded-fetch behavior. The +native FFI host owns one runtime per application core and closes it +idempotently. Cancelling a close never reopens command admission, and a later +close call resumes the same shutdown. Operating-system keyring calls run +through a bounded supervised worker rather than directly on an async runtime +worker. + ## Project documentation The consuming Radroots monorepo owns normative HarvestCircle specifications, diff --git a/core/Cargo.lock b/core/Cargo.lock @@ -831,7 +831,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "39cab71617ae0d63f51a36d69f866391735b51691dbda63cf6f96d042b63efeb" dependencies = [ "libc", - "windows-sys 0.61.2", + "windows-sys 0.52.0", ] [[package]] @@ -1131,6 +1131,7 @@ name = "harvestcircle_application" version = "0.1.0-alpha" dependencies = [ "harvestcircle_domain", + "radroots_transport_nostr", "secrecy", "tokio", "uuid", @@ -1143,7 +1144,6 @@ dependencies = [ "bech32", "radroots_identity", "secrecy", - "url", "zeroize", ] @@ -1158,9 +1158,9 @@ dependencies = [ "harvestcircle_product", "harvestcircle_runtime", "harvestcircle_storage", - "nostr", + "nostr 0.44.1", "nostr-relay-builder", - "nostr-sdk", + "nostr-sdk 0.44.0", "quote", "radroots_runtime_paths", "radroots_service_sqlite", @@ -1177,10 +1177,12 @@ version = "0.1.0-alpha" dependencies = [ "harvestcircle_application", "harvestcircle_domain", - "nostr", + "nostr 0.44.1", "nostr-relay-builder", - "nostr-sdk", + "nostr-sdk 0.44.0", "radroots_identity", + "radroots_transport", + "radroots_transport_nostr", "tokio", ] @@ -1200,11 +1202,12 @@ dependencies = [ "harvestcircle_domain", "harvestcircle_nostr", "harvestcircle_storage", - "nostr", + "nostr 0.44.1", "nostr-relay-builder", - "nostr-sdk", + "nostr-sdk 0.44.0", "radroots_runtime_paths", "radroots_service_sqlite", + "radroots_transport_nostr", "tempfile", "tokio", "uuid", @@ -1237,9 +1240,9 @@ dependencies = [ "harvestcircle_domain", "harvestcircle_nostr", "harvestcircle_runtime", - "nostr", + "nostr 0.44.1", "nostr-relay-builder", - "nostr-sdk", + "nostr-sdk 0.44.0", "radroots_runtime_paths", "radroots_service_sqlite", "tokio", @@ -1433,6 +1436,16 @@ dependencies = [ [[package]] name = "idna" +version = "0.5.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "634d9b1461af396cad843f47fdba5597a4f9e6ddd4bfb6ff5d85028c25cb12f6" +dependencies = [ + "unicode-bidi", + "unicode-normalization", +] + +[[package]] +name = "idna" version = "1.1.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "3b0875f23caa03898994f6ddc501886a45c7d3d62d04d2d90788d47be1b1e4de" @@ -1697,21 +1710,67 @@ dependencies = [ ] [[package]] +name = "nostr" +version = "0.44.8" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "40ff7b77ef428b40aa2834a6acbae38a0e104c98b306208ca4b87a420d579a4b" +dependencies = [ + "aes", + "base64", + "bech32", + "bip39", + "bitcoin_hashes", + "cbc", + "chacha20", + "chacha20poly1305", + "getrandom 0.2.17", + "hex", + "instant", + "scrypt", + "secp256k1", + "serde", + "serde_json", + "unicode-normalization", + "url", + "url-fork", +] + +[[package]] +name = "nostr-database" +version = "0.44.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7462c9d8ae5ef6a28d66a192d399ad2530f1f2130b13186296dbb11bdef5b3d1" +dependencies = [ + "lru", + "nostr 0.44.8", + "tokio", +] + +[[package]] name = "nostr-database" version = "0.44.0" source = "git+https://github.com/rust-nostr/nostr.git?rev=5bba5163eb77107f82c4a8262cf29d7f33a73219#5bba5163eb77107f82c4a8262cf29d7f33a73219" dependencies = [ "lru", - "nostr", + "nostr 0.44.1", "tokio", ] [[package]] name = "nostr-gossip" version = "0.44.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ade30de16869618919c6b5efc8258f47b654a98b51541eb77f85e8ec5e3c83a6" +dependencies = [ + "nostr 0.44.8", +] + +[[package]] +name = "nostr-gossip" +version = "0.44.0" source = "git+https://github.com/rust-nostr/nostr.git?rev=5bba5163eb77107f82c4a8262cf29d7f33a73219#5bba5163eb77107f82c4a8262cf29d7f33a73219" dependencies = [ - "nostr", + "nostr 0.44.1", ] [[package]] @@ -1724,8 +1783,8 @@ dependencies = [ "atomic-destructor", "hex", "negentropy", - "nostr", - "nostr-database", + "nostr 0.44.1", + "nostr-database 0.44.0 (git+https://github.com/rust-nostr/nostr.git?rev=5bba5163eb77107f82c4a8262cf29d7f33a73219)", "tokio", "tracing", ] @@ -1741,8 +1800,26 @@ dependencies = [ "hex", "lru", "negentropy", - "nostr", - "nostr-database", + "nostr 0.44.1", + "nostr-database 0.44.0 (git+https://github.com/rust-nostr/nostr.git?rev=5bba5163eb77107f82c4a8262cf29d7f33a73219)", + "tokio", + "tracing", +] + +[[package]] +name = "nostr-relay-pool" +version = "0.44.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c85c54d6ca9aae4ae2bf19a7663ba9db5f45f783f1d24aff55f006386b8b99a1" +dependencies = [ + "async-utility", + "async-wsocket", + "atomic-destructor", + "hex", + "lru", + "negentropy", + "nostr 0.44.8", + "nostr-database 0.44.0 (registry+https://github.com/rust-lang/crates.io-index)", "tokio", "tracing", ] @@ -1753,10 +1830,25 @@ version = "0.44.0" source = "git+https://github.com/rust-nostr/nostr.git?rev=5bba5163eb77107f82c4a8262cf29d7f33a73219#5bba5163eb77107f82c4a8262cf29d7f33a73219" dependencies = [ "async-utility", - "nostr", - "nostr-database", - "nostr-gossip", - "nostr-relay-pool", + "nostr 0.44.1", + "nostr-database 0.44.0 (git+https://github.com/rust-nostr/nostr.git?rev=5bba5163eb77107f82c4a8262cf29d7f33a73219)", + "nostr-gossip 0.44.0 (git+https://github.com/rust-nostr/nostr.git?rev=5bba5163eb77107f82c4a8262cf29d7f33a73219)", + "nostr-relay-pool 0.44.0", + "tokio", + "tracing", +] + +[[package]] +name = "nostr-sdk" +version = "0.44.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "471732576710e779b64f04c55e3f8b5292f865fea228436daf19694f0bf70393" +dependencies = [ + "async-utility", + "nostr 0.44.8", + "nostr-database 0.44.0 (registry+https://github.com/rust-lang/crates.io-index)", + "nostr-gossip 0.44.0 (registry+https://github.com/rust-lang/crates.io-index)", + "nostr-relay-pool 0.44.3", "tokio", "tracing", ] @@ -2043,6 +2135,7 @@ version = "0.1.0-alpha" source = "git+https://github.com/radrootslabs/lib?rev=be9db78e060ebc0000fa7827ac32efa3f6504f53#be9db78e060ebc0000fa7827ac32efa3f6504f53" dependencies = [ "mediatype", + "serde", "sha2", "unicode-general-category", "url", @@ -2103,6 +2196,20 @@ dependencies = [ ] [[package]] +name = "radroots_nostr" +version = "0.1.0-alpha" +source = "git+https://github.com/radrootslabs/lib?rev=be9db78e060ebc0000fa7827ac32efa3f6504f53#be9db78e060ebc0000fa7827ac32efa3f6504f53" +dependencies = [ + "nostr 0.44.8", + "radroots_event", + "radroots_event_codec", + "radroots_identity", + "serde", + "serde_json", + "thiserror 2.0.20", +] + +[[package]] name = "radroots_protocol" version = "0.1.0-alpha" source = "git+https://github.com/radrootslabs/lib?rev=be9db78e060ebc0000fa7827ac32efa3f6504f53#be9db78e060ebc0000fa7827ac32efa3f6504f53" @@ -2172,6 +2279,26 @@ dependencies = [ ] [[package]] +name = "radroots_transport_nostr" +version = "0.1.0-alpha" +source = "git+https://github.com/radrootslabs/lib?rev=be9db78e060ebc0000fa7827ac32efa3f6504f53#be9db78e060ebc0000fa7827ac32efa3f6504f53" +dependencies = [ + "async-wsocket", + "futures", + "nostr-relay-pool 0.44.3", + "nostr-sdk 0.44.1", + "radroots_event_codec", + "radroots_nostr", + "radroots_protocol", + "radroots_transport", + "serde_json", + "sha2", + "tokio", + "tokio-tungstenite", + "url", +] + +[[package]] name = "rand" version = "0.8.7" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -2327,7 +2454,7 @@ dependencies = [ "errno", "libc", "linux-raw-sys", - "windows-sys 0.61.2", + "windows-sys 0.52.0", ] [[package]] @@ -2826,7 +2953,7 @@ dependencies = [ "getrandom 0.3.4", "once_cell", "rustix", - "windows-sys 0.61.2", + "windows-sys 0.52.0", ] [[package]] @@ -3101,6 +3228,12 @@ dependencies = [ ] [[package]] +name = "unicode-bidi" +version = "0.3.18" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5c1cb5db39152898a79168971543b1cb5020dff7fe43c8dc468b0885f5e29df5" + +[[package]] name = "unicode-general-category" version = "1.1.0" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -3265,13 +3398,25 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ff67a8a4397373c3ef660812acab3268222035010ab8680ec4215f38ba3d0eed" dependencies = [ "form_urlencoded", - "idna", + "idna 1.1.0", "percent-encoding", "serde", "serde_derive", ] [[package]] +name = "url-fork" +version = "3.0.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7fa3323c39b8e786154d3000b70ae9af0e9bd746c9791456da0d4a1f68ad89d6" +dependencies = [ + "form_urlencoded", + "idna 0.5.0", + "percent-encoding", + "serde", +] + +[[package]] name = "utf-8" version = "0.7.6" source = "registry+https://github.com/rust-lang/crates.io-index" diff --git a/core/Cargo.toml b/core/Cargo.toml @@ -47,6 +47,8 @@ radroots_identity = { git = "https://github.com/radrootslabs/lib", rev = "be9db7 radroots_runtime_paths = { git = "https://github.com/radrootslabs/lib", rev = "be9db78e060ebc0000fa7827ac32efa3f6504f53", version = "=0.1.0-alpha", default-features = false } radroots_service_sqlite = { git = "https://github.com/radrootslabs/lib", rev = "be9db78e060ebc0000fa7827ac32efa3f6504f53", version = "=0.1.0-alpha", default-features = false } radroots_storage = { git = "https://github.com/radrootslabs/lib", rev = "be9db78e060ebc0000fa7827ac32efa3f6504f53", version = "=0.1.0-alpha", default-features = false } +radroots_transport = { git = "https://github.com/radrootslabs/lib", rev = "be9db78e060ebc0000fa7827ac32efa3f6504f53", version = "=0.1.0-alpha", default-features = false } +radroots_transport_nostr = { git = "https://github.com/radrootslabs/lib", rev = "be9db78e060ebc0000fa7827ac32efa3f6504f53", version = "=0.1.0-alpha", default-features = false } getrandom = { version = "0.2", default-features = false } quote = { version = "1" } sha2 = { version = "0.10", default-features = false } diff --git a/core/compatibility/harvestcircle-ffi-v4.properties b/core/compatibility/harvestcircle-ffi-v4.properties @@ -2,7 +2,7 @@ schema=harvestcircle.ffi.v4 contract.id=harvestcircle-desktop-ffi-v4 contract.major=4 contract.minor=3 -contract.hash=936658b01aca8baf208eba5bc9696fce262e40375edade46873c3106c01046ab +contract.hash=5aebb8c68546fee90611314163db0d944705be6d88dfd1a44dccd202d0e805d2 product.coordinate_digest=bf50f9ea6c2537406de255f025463e670eb6263c295f992f7e4c4db36d957064 snapshot.schema=1 storage.schema.minimum=1 diff --git a/core/crates/harvestcircle_application/Cargo.toml b/core/crates/harvestcircle_application/Cargo.toml @@ -13,6 +13,7 @@ include = ["src/**", "tests/**", "Cargo.toml"] [dependencies] harvestcircle_domain.workspace = true +radroots_transport_nostr.workspace = true secrecy = "=0.10.3" tokio = { version = "=1.47.1", features = ["macros", "rt-multi-thread", "sync", "time"] } uuid.workspace = true diff --git a/core/crates/harvestcircle_application/src/app_core.rs b/core/crates/harvestcircle_application/src/app_core.rs @@ -70,6 +70,7 @@ pub struct AppCore { relay_configuration: RelayConfiguration, key_material: Arc<dyn KeyMaterialProvider>, state: Mutex<CoreState>, + published_snapshot: tokio::sync::watch::Sender<AppSnapshot>, } impl AppCore { @@ -78,14 +79,17 @@ impl AppCore { relay_configuration: RelayConfiguration, key_material: Arc<dyn KeyMaterialProvider>, ) -> Self { + let state_machine = StateMachine::booting(); + let (published_snapshot, _) = tokio::sync::watch::channel(state_machine.snapshot().clone()); Self { relay_configuration, key_material, state: Mutex::new(CoreState { - state_machine: StateMachine::booting(), + state_machine, removal_tokens: BTreeMap::new(), next_removal_token: 1, }), + published_snapshot, } } @@ -145,16 +149,19 @@ impl AppCore { #[must_use] pub fn snapshot(&self) -> AppSnapshot { - self.lock_state().state_machine.snapshot().clone() + self.published_snapshot.borrow().clone() } pub(crate) fn apply_transition( &self, transition: StateTransition, ) -> Result<AppSnapshot, SafeError> { - self.lock_state() + let snapshot = self + .lock_state()? .state_machine - .apply(transition, &self.relay_configuration) + .apply(transition, &self.relay_configuration)?; + self.published_snapshot.send_replace(snapshot.clone()); + Ok(snapshot) } pub(crate) fn issue_removal_token( @@ -162,7 +169,7 @@ impl AppCore { public_key: PublicKey, now: UnixTimestamp, ) -> Result<RemovalConfirmationToken, SafeError> { - let mut state = self.lock_state(); + let mut state = self.lock_state()?; let Some(identity) = state .state_machine .snapshot() @@ -228,7 +235,7 @@ impl AppCore { expires_at, impact, } = token; - let mut state = self.lock_state(); + let mut state = self.lock_state()?; let stored = state.removal_tokens.remove(&id); if stored.is_none_or(|stored| { stored.public_key != public_key @@ -245,14 +252,22 @@ impl AppCore { #[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) + .map(|mut state| state.removal_tokens.remove(&token.id).is_some()) + .unwrap_or(false) } + + fn lock_state(&self) -> Result<MutexGuard<'_, CoreState>, SafeError> { + self.state.lock().map_err(|_| internal_state_unavailable()) + } +} + +const fn internal_state_unavailable() -> SafeError { + SafeError::new( + SafeErrorCode::InvalidApplicationState, + SafeMessage::new("The application state is unavailable."), + ) } const fn invalid_application_state() -> SafeError { diff --git a/core/crates/harvestcircle_application/src/config.rs b/core/crates/harvestcircle_application/src/config.rs @@ -1,29 +1,21 @@ -use std::collections::HashSet; - -use harvestcircle_domain::{ - RelayDestinationPolicy, RelayEndpoint, SafeError, SafeErrorCode, SafeMessage, -}; +use harvestcircle_domain::{SafeError, SafeErrorCode, SafeMessage}; +use radroots_transport_nostr::{RelayAccess, RelayEndpoint, RelayUrlPolicy}; use crate::RelayConfiguration; #[derive(Clone, Debug, Eq, PartialEq)] pub struct RelayEndpointInput { url: String, - destination: RelayDestinationPolicy, + destination: RelayUrlPolicy, read: bool, write: bool, } impl RelayEndpointInput { #[must_use] - pub fn new( - url: impl Into<String>, - destination: RelayDestinationPolicy, - read: bool, - write: bool, - ) -> Self { + pub fn new(url: String, destination: RelayUrlPolicy, read: bool, write: bool) -> Self { Self { - url: url.into(), + url, destination, read, write, @@ -40,15 +32,22 @@ impl RelayEndpointInput { pub fn relay_configuration_from_endpoints( values: &[RelayEndpointInput], ) -> Result<RelayConfiguration, SafeError> { - if values.is_empty() { + if values.is_empty() || values.len() > crate::MAX_CONFIGURED_RELAYS { return Err(invalid_configuration()); } - let mut seen = HashSet::new(); let mut endpoints = Vec::with_capacity(values.len()); for value in values { - let endpoint = - RelayEndpoint::parse(&value.url, value.destination, value.read, value.write)?; - if !seen.insert(endpoint.url().as_str().to_owned()) { + let access = match (value.read, value.write) { + (true, false) => RelayAccess::ReadOnly, + (true, true) => RelayAccess::ReadWrite, + (false, _) => return Err(invalid_configuration()), + }; + let endpoint = RelayEndpoint::new(&value.url, value.destination, access) + .map_err(|_| invalid_configuration())?; + if endpoints + .iter() + .any(|existing: &RelayEndpoint| existing.url() == endpoint.url()) + { return Err(invalid_configuration()); } endpoints.push(endpoint); @@ -65,25 +64,26 @@ const fn invalid_configuration() -> SafeError { #[cfg(test)] mod tests { - use harvestcircle_domain::{RelayDestinationPolicy, SafeErrorCode}; + use harvestcircle_domain::SafeErrorCode; + use radroots_transport_nostr::RelayUrlPolicy; use super::{RelayEndpointInput, relay_configuration_from_endpoints}; #[test] - fn relay_config_requires_explicit_input_and_supports_mixed_destinations() { + fn relay_config_requires_explicit_input_and_one_governed_profile() { let error = relay_configuration_from_endpoints(&[]).expect_err("input required"); assert_eq!(error.code(), SafeErrorCode::InvalidRelayConfiguration); let development = relay_configuration_from_endpoints(&[ RelayEndpointInput::new( - "ws://localhost:8080", - RelayDestinationPolicy::Local, + "ws://localhost:8080".to_owned(), + RelayUrlPolicy::Local, true, true, ), RelayEndpointInput::new( - "wss://relay.example", - RelayDestinationPolicy::Public, + "ws://127.0.0.1:8081".to_owned(), + RelayUrlPolicy::Local, true, true, ), @@ -91,12 +91,25 @@ mod tests { .expect("explicit development relay"); assert_eq!( development.relays()[0].url().as_str(), - "ws://localhost:8080/" - ); - assert_eq!( - development.relays()[1].destination(), - RelayDestinationPolicy::Public + "ws://localhost:8080" ); + assert_eq!(development.relays()[1].policy(), RelayUrlPolicy::Local); + + let mixed = [ + RelayEndpointInput::new( + "ws://localhost:8080".to_owned(), + RelayUrlPolicy::Local, + true, + true, + ), + RelayEndpointInput::new( + "wss://relay.example".to_owned(), + RelayUrlPolicy::Public, + true, + true, + ), + ]; + assert!(relay_configuration_from_endpoints(&mixed).is_err()); } #[test] @@ -104,21 +117,21 @@ mod tests { for values in [ vec![ RelayEndpointInput::new( - "wss://relay.one", - RelayDestinationPolicy::Public, + "wss://relay.one".to_owned(), + RelayUrlPolicy::Public, true, false, ), RelayEndpointInput::new( - "wss://RELAY.one/", - RelayDestinationPolicy::Public, + "wss://RELAY.one/".to_owned(), + RelayUrlPolicy::Public, false, true, ), ], vec![RelayEndpointInput::new( - "wss://relay.one", - RelayDestinationPolicy::Public, + "wss://relay.one".to_owned(), + RelayUrlPolicy::Public, false, false, )], @@ -131,14 +144,14 @@ mod tests { fn relay_config_rejects_any_invalid_comma_separated_entry() { let error = relay_configuration_from_endpoints(&[ RelayEndpointInput::new( - "wss://relay.one", - RelayDestinationPolicy::Public, + "wss://relay.one".to_owned(), + RelayUrlPolicy::Public, true, true, ), RelayEndpointInput::new( - "https://not-a-relay.test", - RelayDestinationPolicy::Public, + "https://not-a-relay.test".to_owned(), + RelayUrlPolicy::Public, true, true, ), diff --git a/core/crates/harvestcircle_application/src/custody.rs b/core/crates/harvestcircle_application/src/custody.rs @@ -113,7 +113,7 @@ impl GeneratedKeyRecoveryHandle { pub fn take_recovery_nsec(&self) -> Result<Nsec, SafeError> { self.recovery_nsec .lock() - .unwrap_or_else(std::sync::PoisonError::into_inner) + .map_err(|_| recovery_not_available())? .take() .ok_or_else(recovery_not_available) } diff --git a/core/crates/harvestcircle_application/src/lib.rs b/core/crates/harvestcircle_application/src/lib.rs @@ -49,6 +49,9 @@ pub use ports::{ PendingIdentityOperation, }; pub use profile_refresh::ProfileRefreshPlan; +pub use radroots_transport_nostr::{ + RelayAccess, RelayEndpoint, RelayProfileKind, RelayUrl, RelayUrlPolicy, +}; pub use secrets::{ FailureSecretStore, InMemorySecretStore, SecretStore, SecretStoreCall, SecretStoreOperation, }; diff --git a/core/crates/harvestcircle_application/src/ports.rs b/core/crates/harvestcircle_application/src/ports.rs @@ -3,9 +3,10 @@ use std::pin::Pin; use std::time::Instant; use harvestcircle_domain::{ - Kind0ProfileCandidate, NostrIdentity, Npub, Nsec, PublicKey, RelayEndpoint, SafeError, - SafeErrorCode, SafeMessage, SecretKeyInput, SignerAvailability, UnixTimestamp, + Kind0ProfileCandidate, NostrIdentity, Npub, Nsec, PublicKey, SafeError, SafeErrorCode, + SafeMessage, SecretKeyInput, SignerAvailability, UnixTimestamp, }; +use radroots_transport_nostr::RelayEndpoint; #[derive(Clone, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)] pub struct DurableRequestId(String); @@ -16,13 +17,23 @@ impl DurableRequestId { /// # Errors /// /// Returns a safe validation error when the value is not canonical UUIDv7 text. - pub fn parse(value: impl Into<String>) -> Result<Self, SafeError> { - let value = value.into(); - let parsed = uuid::Uuid::parse_str(&value).map_err(|_| invalid_request_id())?; + pub fn parse(value: impl AsRef<str>) -> Result<Self, SafeError> { + let value = value.as_ref(); + if value.len() != 36 + || !value.is_ascii() + || value.bytes().enumerate().any(|(index, byte)| match index { + 8 | 13 | 18 | 23 => byte != b'-', + 19 => !matches!(byte, b'8' | b'9' | b'a' | b'b'), + _ => !matches!(byte, b'0'..=b'9' | b'a'..=b'f'), + }) + { + return Err(invalid_request_id()); + } + let parsed = uuid::Uuid::parse_str(value).map_err(|_| invalid_request_id())?; if parsed.get_version_num() != 7 || parsed.hyphenated().to_string() != value { return Err(invalid_request_id()); } - Ok(Self(value)) + Ok(Self(value.to_owned())) } #[must_use] @@ -735,7 +746,8 @@ mod tests { use std::sync::Mutex; - use harvestcircle_domain::{NostrIdentity, PublicKey, RelayEndpoint, SafeError, UnixTimestamp}; + use harvestcircle_domain::{NostrIdentity, PublicKey, SafeError, UnixTimestamp}; + use radroots_transport_nostr::RelayEndpoint; use super::{ AppStateRepository, BoxFuture, CachedProfile, Clock, DurableOperationReceipt, @@ -763,9 +775,12 @@ mod tests { "01890f3e-7b1c-4000-8000-000000000031", "01890F3E-7B1C-7000-8000-000000000031", "01890f3e-7b1c-6000-8000-000000000031", + "01890f3e-7b1c-7000-7000-000000000031", ] { assert!(DurableRequestId::parse(invalid).is_err()); } + let oversized = "0".repeat(1_000_000); + assert!(DurableRequestId::parse(&oversized).is_err()); } #[tokio::test] diff --git a/core/crates/harvestcircle_application/src/profile_refresh.rs b/core/crates/harvestcircle_application/src/profile_refresh.rs @@ -1,4 +1,5 @@ -use harvestcircle_domain::{PublicKey, RelayEndpoint, SafeError, SafeErrorCode}; +use harvestcircle_domain::{PublicKey, SafeError, SafeErrorCode}; +use radroots_transport_nostr::RelayEndpoint; use std::time::Instant; use crate::{ @@ -259,10 +260,10 @@ mod tests { use std::time::{Duration, Instant}; use harvestcircle_domain::{ - EventId, Kind0ProfileCandidate, ProfileMetadata, PublicKey, RelayDestinationPolicy, - RelayEndpoint, SafeError, SafeErrorCode, SafeMessage, SecretKeyInput, UnixTimestamp, - select_latest_kind0, + EventId, Kind0ProfileCandidate, ProfileMetadata, PublicKey, SafeError, SafeErrorCode, + SafeMessage, SecretKeyInput, UnixTimestamp, select_latest_kind0, }; + use radroots_transport_nostr::{RelayAccess, RelayEndpoint, RelayUrlPolicy}; use crate::{ ActiveIdentitySnapshot, AppCore, BoxFuture, CachedProfile, Clock, @@ -394,11 +395,10 @@ mod tests { cached_name: Option<&str>, ) -> (AppCore, PublicKey) { let relays = RelayConfiguration::new(vec![ - RelayEndpoint::parse( + RelayEndpoint::new( "ws://localhost:8080", - RelayDestinationPolicy::Local, - true, - true, + RelayUrlPolicy::Local, + RelayAccess::ReadWrite, ) .expect("relay"), ]) diff --git a/core/crates/harvestcircle_application/src/secrets.rs b/core/crates/harvestcircle_application/src/secrets.rs @@ -71,20 +71,17 @@ pub struct FailureSecretStore { 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; + if let Ok(mut failures) = self.remaining_failures.lock() { + *failures.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() + .map(|calls| calls.clone()) + .unwrap_or_default() } fn record_and_should_fail( @@ -92,17 +89,17 @@ impl FailureSecretStore { 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 Ok(mut calls) = self.calls.lock() else { + return true; + }; + calls.push(SecretStoreCall { + operation, + public_key, + }); + drop(calls); + let Ok(mut failures) = self.remaining_failures.lock() else { + return true; + }; let remaining = failures.entry(operation).or_default(); let should_fail = *remaining > 0; *remaining = remaining.saturating_sub(1); @@ -141,16 +138,14 @@ impl SecretStore for FailureSecretStore { } impl InMemorySecretStore { - fn credentials(&self) -> MutexGuard<'_, BTreeMap<PublicKey, SecretString>> { - self.credentials - .lock() - .unwrap_or_else(std::sync::PoisonError::into_inner) + fn credentials(&self) -> Result<MutexGuard<'_, BTreeMap<PublicKey, SecretString>>, SafeError> { + self.credentials.lock().map_err(|_| keyring_unavailable()) } } impl SecretStore for InMemorySecretStore { fn put(&self, public_key: PublicKey, secret: SecretKeyInput) -> Result<(), SafeError> { - let mut credentials = self.credentials(); + let mut credentials = self.credentials()?; if credentials.contains_key(&public_key) { return Err(credential_exists()); } @@ -160,7 +155,7 @@ impl SecretStore for InMemorySecretStore { } fn load(&self, public_key: PublicKey) -> Result<SecretKeyInput, SafeError> { - let credentials = self.credentials(); + let credentials = self.credentials()?; let secret = credentials .get(&public_key) .ok_or_else(credential_missing)?; @@ -168,11 +163,11 @@ impl SecretStore for InMemorySecretStore { } fn contains(&self, public_key: PublicKey) -> Result<bool, SafeError> { - Ok(self.credentials().contains_key(&public_key)) + Ok(self.credentials()?.contains_key(&public_key)) } fn delete(&self, public_key: PublicKey) -> Result<(), SafeError> { - self.credentials() + self.credentials()? .remove(&public_key) .map(|_| ()) .ok_or_else(credential_missing) diff --git a/core/crates/harvestcircle_application/src/snapshot.rs b/core/crates/harvestcircle_application/src/snapshot.rs @@ -1,8 +1,9 @@ use std::collections::HashSet; use harvestcircle_domain::{ - NostrIdentity, ProfileMetadata, PublicKey, RelayEndpoint, SafeError, SafeErrorCode, SafeMessage, + NostrIdentity, ProfileMetadata, PublicKey, SafeError, SafeErrorCode, SafeMessage, }; +use radroots_transport_nostr::{RelayEndpoint, RelayProfile, RelayProfileKind, RelayUrlPolicy}; pub const MAX_CONFIGURED_RELAYS: usize = 16; @@ -68,8 +69,18 @@ pub enum ProfileLoadState { Error(SafeError), } -#[derive(Clone, Debug, Default, Eq, PartialEq)] -pub struct RelayConfiguration(Vec<RelayEndpoint>); +#[derive(Clone, Default, Eq, PartialEq)] +pub struct RelayConfiguration(Option<RelayProfile>); + +impl core::fmt::Debug for RelayConfiguration { + fn fmt(&self, formatter: &mut core::fmt::Formatter<'_>) -> core::fmt::Result { + formatter + .debug_struct("RelayConfiguration") + .field("profile", &self.0.as_ref().map(RelayProfile::kind)) + .field("relay_count", &self.relays().len()) + .finish() + } +} impl RelayConfiguration { /// Creates a bounded, explicitly classified relay configuration. @@ -79,15 +90,40 @@ impl RelayConfiguration { /// Returns a safe configuration error before runtime or network work when /// the relay count exceeds the HarvestCircle policy. pub fn new(relays: Vec<RelayEndpoint>) -> Result<Self, SafeError> { + if relays.is_empty() { + return Ok(Self::default()); + } if relays.len() > MAX_CONFIGURED_RELAYS { return Err(relay_limit_exceeded()); } - Ok(Self(relays)) + let has_local = relays + .iter() + .any(|relay| relay.policy() == RelayUrlPolicy::Local); + let has_private = relays + .iter() + .any(|relay| relay.policy() == RelayUrlPolicy::PrivateNetwork); + let has_public = relays + .iter() + .any(|relay| relay.policy() == RelayUrlPolicy::Public); + let profile = match (has_local, has_private, has_public) { + (true, false, false) => RelayProfileKind::Simulator, + (false, false, true) => RelayProfileKind::Public, + (false, true, _) => RelayProfileKind::Device, + _ => return Err(relay_limit_exceeded()), + }; + let profile = + RelayProfile::explicit(profile, relays).map_err(|_| relay_limit_exceeded())?; + Ok(Self(Some(profile))) } #[must_use] pub fn relays(&self) -> &[RelayEndpoint] { - &self.0 + self.0.as_ref().map_or(&[], |profile| profile.endpoints()) + } + + #[must_use] + pub const fn profile(&self) -> Option<&RelayProfile> { + self.0.as_ref() } } @@ -295,8 +331,9 @@ const fn invalid_snapshot() -> SafeError { mod tests { use harvestcircle_domain::{ IdentityCreatedAt, LocalKeyringBinding, NostrIdentity, NostrIdentityReference, - RelayDestinationPolicy, RelayEndpoint, SafeErrorCode, SignerAvailability, UnixTimestamp, + SafeErrorCode, SignerAvailability, UnixTimestamp, }; + use radroots_transport_nostr::{RelayAccess, RelayEndpoint, RelayUrlPolicy}; use super::{ ActiveIdentitySnapshot, AppLifecycle, AppSnapshot, ProfileLoadState, RelayConfiguration, @@ -337,11 +374,10 @@ mod tests { fn relay_configuration_rejects_excess_targets_before_runtime_work() { let relays = (0..=super::MAX_CONFIGURED_RELAYS) .map(|index| { - RelayEndpoint::parse( - format!("wss://relay-{index}.example").as_str(), - RelayDestinationPolicy::Public, - true, - true, + RelayEndpoint::new( + format!("wss://relay-{index}.example"), + RelayUrlPolicy::Public, + RelayAccess::ReadWrite, ) .expect("relay") }) diff --git a/core/crates/harvestcircle_domain/Cargo.toml b/core/crates/harvestcircle_domain/Cargo.toml @@ -15,7 +15,6 @@ include = ["src/**", "Cargo.toml"] bech32 = "=0.11.1" radroots_identity.workspace = true secrecy = "=0.10.3" -url = "=2.5.8" zeroize = "=1.9.0" [lints] diff --git a/core/crates/harvestcircle_domain/src/lib.rs b/core/crates/harvestcircle_domain/src/lib.rs @@ -4,7 +4,6 @@ pub mod error; pub mod identity; pub mod key; pub mod profile; -pub mod relay; pub mod time; pub use error::{SafeError, SafeErrorCode, SafeMessage}; @@ -17,5 +16,4 @@ pub use key::{ SecretKeyInput, SecretKeyInputKind, classify_persisted_public_key, }; pub use profile::{EventId, Kind0ProfileCandidate, ProfileMetadata, select_latest_kind0}; -pub use relay::{RelayDestinationPolicy, RelayEndpoint, RelayUrl}; pub use time::UnixTimestamp; diff --git a/core/crates/harvestcircle_domain/src/profile.rs b/core/crates/harvestcircle_domain/src/profile.rs @@ -28,7 +28,7 @@ impl EventId { } let mut bytes = [0_u8; EVENT_ID_BYTES]; - for (index, pair) in value.as_bytes().chunks_exact(2).enumerate() { + for (index, pair) in value.as_bytes().as_chunks::<2>().0.iter().enumerate() { let high = decode_hex(pair[0]).ok_or_else(invalid_profile_metadata)?; let low = decode_hex(pair[1]).ok_or_else(invalid_profile_metadata)?; bytes[index] = (high << 4) | low; diff --git a/core/crates/harvestcircle_domain/src/relay.rs b/core/crates/harvestcircle_domain/src/relay.rs @@ -1,268 +0,0 @@ -//! Validated Nostr relay values. - -use std::fmt::{self, Display, Formatter}; - -use url::{Host, Url}; - -use crate::{SafeError, SafeErrorCode, SafeMessage}; - -#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)] -pub enum RelayDestinationPolicy { - Public, - Local, - PrivateNetwork, -} - -#[derive(Clone, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)] -pub struct RelayUrl { - value: String, - policy: RelayDestinationPolicy, -} - -#[derive(Clone, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)] -pub struct RelayEndpoint { - url: RelayUrl, - read: bool, - write: bool, -} - -impl RelayEndpoint { - /// Validates one explicitly classified relay endpoint and its capabilities. - /// - /// # Errors - /// - /// Returns a safe configuration error when the URL and destination conflict - /// or when the endpoint has neither read nor write capability. - pub fn parse( - value: &str, - destination: RelayDestinationPolicy, - read: bool, - write: bool, - ) -> Result<Self, SafeError> { - if !read && !write { - return Err(invalid_relay()); - } - Ok(Self { - url: RelayUrl::parse(value, destination)?, - read, - write, - }) - } - - #[must_use] - pub const fn url(&self) -> &RelayUrl { - &self.url - } - - #[must_use] - pub const fn destination(&self) -> RelayDestinationPolicy { - self.url.policy() - } - - #[must_use] - pub const fn can_read(&self) -> bool { - self.read - } - - #[must_use] - pub const fn can_write(&self) -> bool { - self.write - } -} - -impl RelayUrl { - /// Parses and normalizes an allowed WebSocket relay URL. - /// - /// # Errors - /// - /// Returns a safe configuration error for empty or malformed input, - /// forbidden schemes, credentials, fragments, or non-loopback `ws://`. - pub fn parse(value: &str, policy: RelayDestinationPolicy) -> Result<Self, SafeError> { - let trimmed = value.trim(); - if trimmed.is_empty() || trimmed.chars().any(char::is_control) { - return Err(invalid_relay()); - } - - let parsed = Url::parse(trimmed).map_err(|_| invalid_relay())?; - if !parsed.username().is_empty() - || parsed.password().is_some() - || parsed.fragment().is_some() - { - return Err(invalid_relay()); - } - - match (policy, parsed.scheme()) { - (RelayDestinationPolicy::Public, "wss") if is_public_destination(&parsed) => {} - (RelayDestinationPolicy::PrivateNetwork, "wss") if is_private_network(&parsed) => {} - (RelayDestinationPolicy::Local, "ws" | "wss") if is_loopback(&parsed) => {} - _ => return Err(invalid_relay()), - } - - if parsed.host().is_none() { - return Err(invalid_relay()); - } - - Ok(Self { - value: parsed.to_string(), - policy, - }) - } - - #[must_use] - pub fn as_str(&self) -> &str { - &self.value - } - - #[must_use] - pub const fn policy(&self) -> RelayDestinationPolicy { - self.policy - } -} - -impl Display for RelayUrl { - fn fmt(&self, formatter: &mut Formatter<'_>) -> fmt::Result { - formatter.write_str(&self.value) - } -} - -fn is_loopback(url: &Url) -> bool { - match url.host() { - Some(Host::Domain(domain)) => domain == "localhost", - Some(Host::Ipv4(address)) => address.octets()[0] == 127, - Some(Host::Ipv6(address)) => address.is_loopback(), - None => false, - } -} - -fn is_private_network(url: &Url) -> bool { - match url.host() { - Some(Host::Ipv4(address)) => address.is_private() || address.is_link_local(), - Some(Host::Ipv6(address)) => { - let first = address.segments()[0]; - first & 0xfe00 == 0xfc00 || first & 0xffc0 == 0xfe80 - } - Some(Host::Domain(_)) | None => false, - } -} - -fn is_public_destination(url: &Url) -> bool { - match url.host() { - Some(Host::Domain(domain)) => domain != "localhost" && !domain.ends_with(".local"), - Some(Host::Ipv4(address)) => { - !address.is_loopback() - && !address.is_private() - && !address.is_link_local() - && !address.is_unspecified() - } - Some(Host::Ipv6(address)) => { - let first = address.segments()[0]; - !address.is_loopback() - && !address.is_unspecified() - && first & 0xfe00 != 0xfc00 - && first & 0xffc0 != 0xfe80 - } - None => false, - } -} - -const fn invalid_relay() -> SafeError { - SafeError::new( - SafeErrorCode::InvalidRelayConfiguration, - SafeMessage::new("The Nostr relay URL is invalid."), - ) -} - -#[cfg(test)] -mod tests { - use super::{RelayDestinationPolicy, RelayEndpoint, RelayUrl}; - use crate::SafeErrorCode; - - #[test] - fn relay_accepts_secure_remote_and_loopback_development_urls() { - for (input, policy, expected) in [ - ( - " wss://Relay.Example/path ", - RelayDestinationPolicy::Public, - "wss://relay.example/path", - ), - ( - "ws://localhost:8080", - RelayDestinationPolicy::Local, - "ws://localhost:8080/", - ), - ( - "ws://127.42.1.9:8080", - RelayDestinationPolicy::Local, - "ws://127.42.1.9:8080/", - ), - ( - "ws://[::1]:8080", - RelayDestinationPolicy::Local, - "ws://[::1]:8080/", - ), - ] { - let relay = RelayUrl::parse(input, policy).expect("allowed relay"); - assert_eq!(relay.as_str(), expected); - assert_eq!(relay.to_string(), expected); - } - } - - #[test] - fn relay_rejects_non_websocket_credentials_fragments_and_remote_plaintext() { - for input in [ - "", - "https://relay.example", - "http://localhost:8080", - "wss://user:password@relay.example", - "wss://relay.example/#fragment", - "ws://relay.example", - "ws://192.168.1.2:8080", - "ws://localhost.evil.example:8080", - "wss://relay.example/\nunsafe", - ] { - let error = RelayUrl::parse(input, RelayDestinationPolicy::Public) - .expect_err("forbidden relay"); - assert_eq!(error.code(), SafeErrorCode::InvalidRelayConfiguration); - } - } - - #[test] - fn relay_endpoint_requires_a_direction_capability() { - let endpoint = RelayEndpoint::parse( - "wss://relay.example", - RelayDestinationPolicy::Public, - true, - false, - ) - .expect("read endpoint"); - assert!(endpoint.can_read()); - assert!(!endpoint.can_write()); - assert_eq!(endpoint.destination(), RelayDestinationPolicy::Public); - assert!( - RelayEndpoint::parse( - "wss://relay.example", - RelayDestinationPolicy::Public, - false, - false, - ) - .is_err() - ); - } - - #[test] - fn relay_destination_policy_is_explicit_and_fail_closed() { - assert!(RelayUrl::parse("ws://localhost:8080", RelayDestinationPolicy::Public).is_err()); - assert!(RelayUrl::parse("wss://relay.example", RelayDestinationPolicy::Local).is_err()); - assert!(RelayUrl::parse("wss://10.0.0.4", RelayDestinationPolicy::Public).is_err()); - assert!( - RelayUrl::parse( - "wss://relay.example", - RelayDestinationPolicy::PrivateNetwork - ) - .is_err() - ); - let private = RelayUrl::parse("wss://10.0.0.4", RelayDestinationPolicy::PrivateNetwork) - .expect("explicit private network"); - assert_eq!(private.policy(), RelayDestinationPolicy::PrivateNetwork); - } -} diff --git a/core/crates/harvestcircle_ffi/src/commands.rs b/core/crates/harvestcircle_ffi/src/commands.rs @@ -2,17 +2,18 @@ use std::collections::BTreeMap; use std::fmt::{self, Display, Formatter}; use std::num::NonZeroUsize; use std::path::{Path, PathBuf}; -use std::sync::atomic::{AtomicBool, Ordering}; -use std::sync::{Arc, Mutex, OnceLock}; +use std::sync::atomic::{AtomicU8, Ordering}; +use std::sync::{Arc, Mutex}; use std::time::{Duration, SystemTime, UNIX_EPOCH}; use directories::BaseDirs; use harvestcircle_application::{ - Clock, DurableRequestId, GeneratedKeyRecoveryHandle, RelayConfiguration, RelayEndpointInput, - RemovalConfirmationToken, relay_configuration_from_endpoints, + Clock, DurableRequestId, GeneratedKeyRecoveryHandle, MAX_CONFIGURED_RELAYS, RelayConfiguration, + RelayEndpointInput, RelayUrlPolicy, RemovalConfirmationToken, SecretStore, + relay_configuration_from_endpoints, }; use harvestcircle_domain::{ - PublicKey, RelayDestinationPolicy, SafeError, SecretKeyInput, UnixTimestamp, + PublicKey, SafeError, SafeErrorCode, SafeMessage, SecretKeyInput, UnixTimestamp, }; use harvestcircle_nostr::SdkNostrClient; use harvestcircle_runtime::{ @@ -37,12 +38,14 @@ use crate::{ SNAPSHOT_SCHEMA_VERSION, SOURCE_FOUNDATION_BASELINE, SOURCE_PROVENANCE_DIGEST, }, dto::error_policy, + host_runtime::HostRuntime, + keyring_worker::BoundedKeyringWorker, }; pub(crate) const ACTOR_MAILBOX_CAPACITY: usize = 64; const MAX_COMMAND_DEADLINE_MILLIS: u64 = 30_000; -#[derive(Clone, Debug, Eq, PartialEq)] +#[derive(Clone, Eq, PartialEq)] #[cfg_attr(not(coverage_nightly), derive(uniffi::Record))] pub struct RequestContextDto { pub request_id: String, @@ -50,13 +53,33 @@ pub struct RequestContextDto { pub deadline_millis: u64, } -#[derive(Clone, Debug, Eq, PartialEq)] +impl fmt::Debug for RequestContextDto { + fn fmt(&self, formatter: &mut Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("RequestContextDto") + .field("request_id", &"<redacted>") + .field("expected_revision", &self.expected_revision) + .field("deadline_millis", &self.deadline_millis) + .finish() + } +} + +#[derive(Clone, Eq, PartialEq)] #[cfg_attr(not(coverage_nightly), derive(uniffi::Record))] pub struct RelayBootstrapInputDto { pub endpoints: Vec<RelayEndpointDto>, } -#[derive(Clone, Debug, Eq, PartialEq)] +impl fmt::Debug for RelayBootstrapInputDto { + fn fmt(&self, formatter: &mut Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("RelayBootstrapInputDto") + .field("endpoint_count", &self.endpoints.len()) + .finish() + } +} + +#[derive(Clone, Eq, PartialEq)] #[cfg_attr(not(coverage_nightly), derive(uniffi::Record))] pub struct RuntimeOpenInputDto { pub development_mode: bool, @@ -64,7 +87,21 @@ pub struct RuntimeOpenInputDto { pub relay_input: RelayBootstrapInputDto, } -#[derive(Clone, Debug, Eq, PartialEq)] +impl fmt::Debug for RuntimeOpenInputDto { + fn fmt(&self, formatter: &mut Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("RuntimeOpenInputDto") + .field("development_mode", &self.development_mode) + .field( + "explicit_data_directory", + &self.explicit_data_directory.as_ref().map(|_| "<redacted>"), + ) + .field("relay_endpoint_count", &self.relay_input.endpoints.len()) + .finish() + } +} + +#[derive(Clone, Eq, PartialEq)] #[cfg_attr(not(coverage_nightly), derive(uniffi::Record))] pub struct IdentityCommandReceiptDto { pub request_id: String, @@ -72,6 +109,17 @@ pub struct IdentityCommandReceiptDto { pub snapshot: AppSnapshotDto, } +impl fmt::Debug for IdentityCommandReceiptDto { + fn fmt(&self, formatter: &mut Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("IdentityCommandReceiptDto") + .field("request_id", &"<redacted>") + .field("committed_revision", &self.committed_revision) + .field("snapshot", &self.snapshot) + .finish() + } +} + #[derive(Clone, Debug, Eq, PartialEq)] #[cfg_attr(not(coverage_nightly), derive(uniffi::Record))] pub struct CompatibilityDescriptor { @@ -162,7 +210,6 @@ pub fn compatibility_descriptor() -> CompatibilityDescriptor { } } -#[derive(Debug)] #[cfg_attr(not(coverage_nightly), derive(uniffi::Error))] pub enum HarvestCircleError { Failure { @@ -175,6 +222,32 @@ pub enum HarvestCircleError { }, } +impl fmt::Debug for HarvestCircleError { + fn fmt(&self, formatter: &mut Formatter<'_>) -> fmt::Result { + match self { + Self::Failure { + code, + category, + retryable, + recovery_action, + correlation_id, + .. + } => formatter + .debug_struct("HarvestCircleError::Failure") + .field("code", code) + .field("category", category) + .field("retryable", retryable) + .field("recovery_action", recovery_action) + .field( + "correlation_id", + &correlation_id.as_ref().map(|_| "<redacted>"), + ) + .field("safe_message", &"<redacted>") + .finish(), + } + } +} + impl Display for HarvestCircleError { fn fmt(&self, formatter: &mut Formatter<'_>) -> fmt::Result { match self { @@ -200,14 +273,14 @@ impl From<SafeError> for HarvestCircleError { } impl HarvestCircleError { - fn correlated(error: SafeError, correlation_id: &str) -> Self { + fn correlated(error: SafeError, correlation_id: &DurableRequestId) -> Self { let (category, retryable, recovery_action) = error_policy(error.code()); Self::Failure { code: error.code().into(), category, retryable, recovery_action, - correlation_id: Some(correlation_id.to_owned()), + correlation_id: Some(correlation_id.as_str().to_owned()), safe_message: error.message().as_str().to_owned(), } } @@ -216,9 +289,13 @@ impl HarvestCircleError { #[cfg_attr(not(coverage_nightly), derive(uniffi::Object))] pub struct GeneratedRecoveryRequest { handle: GeneratedKeyRecoveryHandle, - resolved: AtomicBool, + resolution: AtomicU8, } +const RECOVERY_PENDING: u8 = 0; +const RECOVERY_RESOLVING: u8 = 1; +const RECOVERY_RESOLVED: u8 = 2; + #[cfg_attr(not(coverage_nightly), uniffi::export)] impl GeneratedRecoveryRequest { pub fn identity(&self) -> IdentityDto { @@ -272,14 +349,17 @@ impl RemovalRequest { pub(crate) struct RuntimeCore { pub(crate) actor: RuntimeActorHandle, + pub(crate) runtime: tokio::runtime::Handle, + pub(crate) host_runtime: Option<Arc<HostRuntime>>, + pub(crate) keyring: Option<Arc<BoundedKeyringWorker>>, pub(crate) observers: Mutex< BTreeMap< harvestcircle_application::ChangeSubscriptionId, Option<tokio::task::JoinHandle<()>>, >, >, - pub(crate) closed: AtomicBool, - pub(crate) startup_relay_problem: Option<SafeError>, + pub(crate) close_state: AtomicU8, + pub(crate) close_gate: tokio::sync::Mutex<()>, #[cfg(test)] pub(crate) _test_directory: Option<Arc<tempfile::TempDir>>, } @@ -297,12 +377,18 @@ impl RuntimeCore { } pub(crate) fn effective_lifecycle(&self) -> harvestcircle_application::RuntimeLifecycle { - let lifecycle = self.actor.lifecycle(); - match (lifecycle, self.startup_relay_problem) { - (harvestcircle_application::RuntimeLifecycle::Ready, Some(problem)) => { - harvestcircle_application::RuntimeLifecycle::Degraded(problem) - } - _ => lifecycle, + self.actor.lifecycle() + } + + pub(crate) fn is_open(&self) -> bool { + self.close_state.load(Ordering::Acquire) == 0 + } + + pub(crate) fn ensure_open(&self) -> Result<(), HarvestCircleError> { + if self.is_open() { + Ok(()) + } else { + Err(runtime_closed_error()) } } } @@ -325,8 +411,10 @@ impl HarvestCircleAppCore { expectation: CompatibilityExpectation, input: RuntimeOpenInputDto, ) -> Result<Arc<Self>, HarvestCircleError> { + verify_compatibility(&expectation)?; + let relays = validated_relay_configuration(&input.relay_input)?; let context = application_runtime_context(&input)?; - Self::open_context_compatible(&context, &expectation, input.relay_input) + Self::open_context(&context, relays) } /// Restores durable public application state. @@ -335,6 +423,7 @@ impl HarvestCircleAppCore { /// /// Returns a safe storage, recovery, or application-state error. pub async fn bootstrap(&self) -> Result<AppSnapshotDto, HarvestCircleError> { + self.inner.ensure_open()?; self.inner .actor .bootstrap() @@ -356,6 +445,7 @@ impl HarvestCircleAppCore { pub async fn begin_generated_identity( &self, ) -> Result<Arc<GeneratedRecoveryRequest>, HarvestCircleError> { + self.inner.ensure_open()?; self.inner .actor .begin_generated_key_stage() @@ -363,7 +453,7 @@ impl HarvestCircleAppCore { .map(|handle| { Arc::new(GeneratedRecoveryRequest { handle, - resolved: AtomicBool::new(false), + resolution: AtomicU8::new(RECOVERY_PENDING), }) }) .map_err(HarvestCircleError::from) @@ -380,13 +470,22 @@ impl HarvestCircleAppCore { context: RequestContextDto, request: Arc<GeneratedRecoveryRequest>, ) -> Result<AppSnapshotDto, HarvestCircleError> { - if request.resolved.swap(true, Ordering::AcqRel) { + self.inner.ensure_open()?; + let (request_id, timeout) = validate_request_context(&context)?; + if request + .resolution + .compare_exchange( + RECOVERY_PENDING, + RECOVERY_RESOLVING, + Ordering::AcqRel, + Ordering::Acquire, + ) + .is_err() + { return Err(generated_recovery_expired()); } - let request_id = DurableRequestId::parse(context.request_id.clone()) - .map_err(|error| HarvestCircleError::correlated(error, &context.request_id))?; - let timeout = command_timeout(context.deadline_millis, &context.request_id)?; - self.inner + let result = self + .inner .actor .acknowledge_generated_key_stage( request.handle.id(), @@ -394,7 +493,11 @@ impl HarvestCircleAppCore { harvestcircle_application::SnapshotRevision::from_value(context.expected_revision), timeout, ) - .await + .await; + request + .resolution + .store(RECOVERY_RESOLVED, Ordering::Release); + result .map(|snapshot| self.inner.dto_for(&snapshot)) .map_err(generated_commit_failed) } @@ -408,14 +511,24 @@ impl HarvestCircleAppCore { &self, request: Arc<GeneratedRecoveryRequest>, ) -> Result<bool, HarvestCircleError> { - if request.resolved.swap(true, Ordering::AcqRel) { + self.inner.ensure_open()?; + if request + .resolution + .compare_exchange( + RECOVERY_PENDING, + RECOVERY_RESOLVING, + Ordering::AcqRel, + Ordering::Acquire, + ) + .is_err() + { return Ok(false); } - self.inner - .actor - .cancel_generated_key_stage() - .await - .map_err(HarvestCircleError::from) + let result = self.inner.actor.cancel_generated_key_stage().await; + request + .resolution + .store(RECOVERY_RESOLVED, Ordering::Release); + result.map_err(HarvestCircleError::from) } /// Imports or repairs an identity using a caller-owned idempotency key. @@ -428,15 +541,14 @@ impl HarvestCircleAppCore { context: RequestContextDto, secret_key: Vec<u8>, ) -> Result<IdentityCommandReceiptDto, HarvestCircleError> { - let request_id = DurableRequestId::parse(context.request_id.clone()) - .map_err(|error| HarvestCircleError::correlated(error, &context.request_id))?; - let timeout = command_timeout(context.deadline_millis, &context.request_id)?; + self.inner.ensure_open()?; + let (request_id, timeout) = validate_request_context(&context)?; let input = SecretKeyInput::parse_bytes(secret_key) - .map_err(|error| HarvestCircleError::correlated(error, &context.request_id))?; + .map_err(|error| HarvestCircleError::correlated(error, &request_id))?; self.inner .actor .import_secret_key( - request_id, + request_id.clone(), harvestcircle_application::SnapshotRevision::from_value(context.expected_revision), input, timeout, @@ -450,7 +562,7 @@ impl HarvestCircleAppCore { snapshot, } }) - .map_err(|error| HarvestCircleError::correlated(error, &context.request_id)) + .map_err(|error| HarvestCircleError::correlated(error, &request_id)) } /// Selects one saved identity without activating it. @@ -462,6 +574,7 @@ impl HarvestCircleAppCore { &self, public_key_hex: String, ) -> Result<AppSnapshotDto, HarvestCircleError> { + self.inner.ensure_open()?; let public_key = parse_public_key(&public_key_hex)?; self.inner .actor @@ -480,6 +593,7 @@ impl HarvestCircleAppCore { &self, public_key_hex: String, ) -> Result<AppSnapshotDto, HarvestCircleError> { + self.inner.ensure_open()?; let public_key = parse_public_key(&public_key_hex)?; self.inner .actor @@ -495,6 +609,7 @@ impl HarvestCircleAppCore { /// /// Returns a safe application-state error. pub async fn sign_out(&self) -> Result<AppSnapshotDto, HarvestCircleError> { + self.inner.ensure_open()?; self.inner .actor .sign_out() @@ -509,6 +624,7 @@ impl HarvestCircleAppCore { /// /// Returns a safe storage or application-state error. pub async fn refresh_active_profile(&self) -> Result<AppSnapshotDto, HarvestCircleError> { + self.inner.ensure_open()?; self.inner .actor .refresh_active_profile() @@ -526,6 +642,7 @@ impl HarvestCircleAppCore { &self, public_key_hex: String, ) -> Result<Arc<RemovalRequest>, HarvestCircleError> { + self.inner.ensure_open()?; let public_key = parse_public_key(&public_key_hex)?; self.inner .actor @@ -554,26 +671,25 @@ impl HarvestCircleAppCore { context: RequestContextDto, request: Arc<RemovalRequest>, ) -> Result<AppSnapshotDto, HarvestCircleError> { + self.inner.ensure_open()?; + let (request_id, timeout) = validate_request_context(&context)?; let token = request .token .lock() - .unwrap_or_else(std::sync::PoisonError::into_inner) + .map_err(|_| internal_state_unavailable())? .take() .ok_or_else(confirmation_expired)?; - let request_id = DurableRequestId::parse(context.request_id.clone()) - .map_err(|error| HarvestCircleError::correlated(error, &context.request_id))?; - let timeout = command_timeout(context.deadline_millis, &context.request_id)?; self.inner .actor .confirm_identity_removal( token, - request_id, + request_id.clone(), harvestcircle_application::SnapshotRevision::from_value(context.expected_revision), timeout, ) .await .map(|snapshot| self.inner.dto_for(&snapshot)) - .map_err(HarvestCircleError::from) + .map_err(|error| HarvestCircleError::correlated(error, &request_id)) } } @@ -594,13 +710,15 @@ fn verify_compatibility(expectation: &CompatibilityExpectation) -> Result<(), Ha } impl HarvestCircleAppCore { + #[cfg(test)] fn open_context_compatible( context: &RuntimeContext, expectation: &CompatibilityExpectation, relay_input: RelayBootstrapInputDto, ) -> Result<Arc<Self>, HarvestCircleError> { verify_compatibility(expectation)?; - Self::open_context(context, relay_input) + let relays = validated_relay_configuration(&relay_input)?; + Self::open_context(context, relays) } // The concrete product opener binds operating-system paths, keyrings, and @@ -609,49 +727,44 @@ impl HarvestCircleAppCore { #[cfg_attr(coverage_nightly, coverage(off))] fn open_context( context: &RuntimeContext, - relay_input: RelayBootstrapInputDto, + relays: RelayConfiguration, ) -> Result<Arc<Self>, HarvestCircleError> { - let relay_endpoints = relay_input - .endpoints - .into_iter() - .map(|endpoint| { - RelayEndpointInput::new( - endpoint.url, - match endpoint.destination { - RelayDestinationDto::Local => RelayDestinationPolicy::Local, - RelayDestinationDto::PrivateNetwork => { - RelayDestinationPolicy::PrivateNetwork - } - RelayDestinationDto::Public => RelayDestinationPolicy::Public, - }, - endpoint.read, - endpoint.write, + let runtime = HostRuntime::new().map_err(|()| runtime_unavailable())?; + let runtime_handle = runtime.handle().clone(); + let keyring = BoundedKeyringWorker::new(OsKeyringSecretStore::default()) + .map_err(HarvestCircleError::from)?; + let secrets: Arc<dyn SecretStore> = keyring.clone(); + let build = migration_build_identity()?; + let actor_capacity = actor_mailbox_capacity()?; + let owned_context = context.clone(); + let actor_runtime = runtime_handle.clone(); + let actor = runtime + .block_on(async move { + RuntimeActorHandle::open( + &owned_context, + relays, + RuntimeDependencies::new( + secrets, + Arc::new(SystemClock), + Arc::new(SdkNostrClient::new(Duration::from_secs(5))), + Arc::new(UuidInstallationIdentitySource), + ), + &build, + actor_capacity, + &actor_runtime, ) + .await }) - .collect::<Vec<_>>(); - let (relays, startup_relay_problem) = - local_first_relay_configuration(relay_configuration_from_endpoints(&relay_endpoints)); - let runtime = runtime()?; - let build = migration_build_identity()?; - let actor = runtime.block_on(RuntimeActorHandle::open( - context, - relays, - RuntimeDependencies::new( - Arc::new(OsKeyringSecretStore::default()), - Arc::new(SystemClock), - Arc::new(SdkNostrClient::new(Duration::from_secs(5))), - Arc::new(UuidInstallationIdentitySource), - ), - &build, - actor_mailbox_capacity()?, - runtime.handle(), - ))?; + .map_err(|()| runtime_unavailable())??; Ok(Arc::new(Self { inner: Arc::new(RuntimeCore { actor, + runtime: runtime_handle, + host_runtime: Some(runtime), + keyring: Some(keyring), observers: Mutex::new(BTreeMap::new()), - closed: AtomicBool::new(false), - startup_relay_problem, + close_state: AtomicU8::new(0), + close_gate: tokio::sync::Mutex::new(()), #[cfg(test)] _test_directory: None, }), @@ -659,13 +772,29 @@ impl HarvestCircleAppCore { } } -fn local_first_relay_configuration( - configured: Result<RelayConfiguration, SafeError>, -) -> (RelayConfiguration, Option<SafeError>) { - match configured { - Ok(relays) => (relays, None), - Err(problem) => (RelayConfiguration::default(), Some(problem)), +fn validated_relay_configuration( + relay_input: &RelayBootstrapInputDto, +) -> Result<RelayConfiguration, HarvestCircleError> { + if relay_input.endpoints.len() > MAX_CONFIGURED_RELAYS { + return Err(invalid_relay_configuration()); } + let relay_endpoints = relay_input + .endpoints + .iter() + .map(|endpoint| { + RelayEndpointInput::new( + endpoint.url.clone(), + match endpoint.destination { + RelayDestinationDto::Local => RelayUrlPolicy::Local, + RelayDestinationDto::PrivateNetwork => RelayUrlPolicy::PrivateNetwork, + RelayDestinationDto::Public => RelayUrlPolicy::Public, + }, + endpoint.read, + endpoint.write, + ) + }) + .collect::<Vec<_>>(); + relay_configuration_from_endpoints(&relay_endpoints).map_err(HarvestCircleError::from) } #[derive(Clone, Copy)] @@ -768,34 +897,32 @@ fn parse_public_key(value: &str) -> Result<PublicKey, HarvestCircleError> { PublicKey::from_hex(value).map_err(HarvestCircleError::from) } -fn command_timeout(millis: u64, correlation_id: &str) -> Result<Duration, HarvestCircleError> { +fn validate_request_context( + context: &RequestContextDto, +) -> Result<(DurableRequestId, Duration), HarvestCircleError> { + let request_id = + DurableRequestId::parse(&context.request_id).map_err(HarvestCircleError::from)?; + let timeout = command_timeout(context.deadline_millis, &request_id)?; + Ok((request_id, timeout)) +} + +fn command_timeout( + millis: u64, + correlation_id: &DurableRequestId, +) -> Result<Duration, HarvestCircleError> { if millis == 0 || millis > MAX_COMMAND_DEADLINE_MILLIS { return Err(HarvestCircleError::Failure { code: WireErrorCode::InvalidApplicationState, category: WireErrorCategory::Input, retryable: false, recovery_action: WireRecoveryAction::None, - correlation_id: Some(correlation_id.to_owned()), + correlation_id: Some(correlation_id.as_str().to_owned()), safe_message: "The command deadline is invalid.".to_owned(), }); } Ok(Duration::from_millis(millis)) } -pub(crate) fn runtime() -> Result<&'static tokio::runtime::Runtime, HarvestCircleError> { - static RUNTIME: OnceLock<Result<tokio::runtime::Runtime, ()>> = OnceLock::new(); - RUNTIME - .get_or_init(|| { - tokio::runtime::Builder::new_multi_thread() - .enable_all() - .thread_name("harvestcircle-core") - .build() - .map_err(|_| ()) - }) - .as_ref() - .map_err(|()| runtime_unavailable()) -} - #[cfg(test)] pub(crate) async fn test_actor( relays: RelayConfiguration, @@ -828,7 +955,7 @@ pub(crate) async fn test_actor( ), &build, actor_mailbox_capacity().expect("capacity"), - runtime().expect("runtime").handle(), + &tokio::runtime::Handle::current(), ) .await .expect("test actor"); @@ -850,6 +977,35 @@ fn runtime_unavailable() -> HarvestCircleError { } } +pub(crate) fn internal_state_unavailable() -> HarvestCircleError { + HarvestCircleError::Failure { + code: WireErrorCode::Internal, + category: WireErrorCategory::Internal, + retryable: false, + recovery_action: WireRecoveryAction::RestartApplication, + correlation_id: None, + safe_message: "The application state is unavailable.".to_owned(), + } +} + +pub(crate) fn runtime_closed_error() -> HarvestCircleError { + HarvestCircleError::Failure { + code: WireErrorCode::InvalidApplicationState, + category: WireErrorCategory::Lifecycle, + retryable: false, + recovery_action: WireRecoveryAction::None, + correlation_id: None, + safe_message: "The application runtime is closed.".to_owned(), + } +} + +fn invalid_relay_configuration() -> HarvestCircleError { + HarvestCircleError::from(SafeError::new( + SafeErrorCode::InvalidRelayConfiguration, + SafeMessage::new("The Nostr relay configuration is invalid."), + )) +} + fn path_unavailable() -> HarvestCircleError { HarvestCircleError::Failure { code: WireErrorCode::StorageUnavailable, @@ -911,24 +1067,26 @@ fn compatibility_mismatch() -> HarvestCircleError { #[cfg(test)] #[cfg_attr(coverage_nightly, coverage(off))] mod tests { + use std::error::Error as _; use std::sync::Arc; use harvestcircle_application::{ - RelayConfiguration, RelayEndpointInput, relay_configuration_from_endpoints, + RelayConfiguration, RelayEndpointInput, RelayUrlPolicy, relay_configuration_from_endpoints, }; - use harvestcircle_domain::{RelayDestinationPolicy, SafeError}; + use harvestcircle_domain::SafeError; use harvestcircle_storage::{ CREDENTIAL_SERVICE, CURRENT_SCHEMA_VERSION, HarvestCircleStorageContract, }; use super::{ CompatibilityExpectation, FFI_CONTRACT_HASH, FFI_CONTRACT_ID, FFI_CONTRACT_MAJOR, - FFI_CONTRACT_MINOR, HarvestCircleAppCore, HarvestCircleError, PRODUCT_COORDINATE_DIGEST, - RelayBootstrapInputDto, RequestContextDto, RuntimeCore, RuntimeOpenInputDto, - SNAPSHOT_SCHEMA_VERSION, WireErrorCategory, WireErrorCode, WireRecoveryAction, - actor_mailbox_capacity, application_runtime_context, compatibility_descriptor, - confirmation_expired, generated_commit_failed, local_first_relay_configuration, - path_unavailable, runtime_unavailable, test_actor, verify_compatibility, + FFI_CONTRACT_MINOR, HarvestCircleAppCore, HarvestCircleError, MAX_CONFIGURED_RELAYS, + PRODUCT_COORDINATE_DIGEST, RelayBootstrapInputDto, RelayDestinationDto, RelayEndpointDto, + RequestContextDto, RuntimeCore, RuntimeOpenInputDto, SNAPSHOT_SCHEMA_VERSION, + WireErrorCategory, WireErrorCode, WireRecoveryAction, actor_mailbox_capacity, + application_runtime_context, compatibility_descriptor, confirmation_expired, + generated_commit_failed, path_unavailable, runtime_unavailable, test_actor, + verify_compatibility, }; async fn in_memory_core() -> Arc<HarvestCircleAppCore> { @@ -936,9 +1094,12 @@ mod tests { Arc::new(HarvestCircleAppCore { inner: Arc::new(RuntimeCore { actor, + runtime: tokio::runtime::Handle::current(), + host_runtime: None, + keyring: None, observers: std::sync::Mutex::new(std::collections::BTreeMap::new()), - closed: std::sync::atomic::AtomicBool::new(false), - startup_relay_problem: None, + close_state: std::sync::atomic::AtomicU8::new(0), + close_gate: tokio::sync::Mutex::new(()), _test_directory: Some(directory), }), }) @@ -1054,6 +1215,25 @@ mod tests { assert!(removal.deletes_local_credential()); assert!(!removal.signs_out()); assert!(removal.expires_at_seconds() > 0); + let invalid_confirmation = core + .confirm_identity_removal( + RequestContextDto { + request_id: "secret-invalid-request".to_owned(), + expected_revision: signed_out.revision, + deadline_millis: 5_000, + }, + Arc::clone(&removal), + ) + .await + .expect_err("invalid request context"); + assert!(matches!( + invalid_confirmation, + HarvestCircleError::Failure { + correlation_id: None, + .. + } + )); + assert_eq!(core.snapshot().identities.len(), 1); let removed = core .confirm_identity_removal( RequestContextDto { @@ -1137,6 +1317,82 @@ mod tests { .await .is_err() ); + + let recovery = core + .begin_generated_identity() + .await + .expect("begin after validation failures"); + let invalid = core + .acknowledge_generated_identity( + RequestContextDto { + request_id: "not-a-valid-request-id".to_owned(), + expected_revision: core.snapshot().revision, + deadline_millis: 5_000, + }, + Arc::clone(&recovery), + ) + .await + .expect_err("invalid request"); + assert!(matches!( + invalid, + HarvestCircleError::Failure { + correlation_id: None, + .. + } + )); + core.acknowledge_generated_identity( + RequestContextDto { + request_id: "01890f3e-7b1c-7000-8000-000000000049".to_owned(), + expected_revision: core.snapshot().revision, + deadline_millis: 5_000, + }, + recovery, + ) + .await + .expect("valid retry retains one-shot recovery"); + } + + #[test] + fn input_and_error_debug_are_type_safe_and_redacted() { + let request_secret = "01890f3e-7b1c-7000-8000-00000000dead"; + let path_secret = "/Users/private/secret-data"; + let relay_secret = "wss://user:secret@example.invalid/private"; + let request = RequestContextDto { + request_id: request_secret.to_owned(), + expected_revision: 9, + deadline_millis: 1_000, + }; + let relay = crate::RelayEndpointDto { + url: relay_secret.to_owned(), + destination: crate::RelayDestinationDto::PrivateNetwork, + read: true, + write: true, + }; + let open = RuntimeOpenInputDto { + development_mode: true, + explicit_data_directory: Some(path_secret.to_owned()), + relay_input: RelayBootstrapInputDto { + endpoints: vec![relay.clone()], + }, + }; + let error = HarvestCircleError::Failure { + code: WireErrorCode::Internal, + category: WireErrorCategory::Internal, + retryable: false, + recovery_action: WireRecoveryAction::RestartApplication, + correlation_id: Some(request_secret.to_owned()), + safe_message: "A safe public message.".to_owned(), + }; + let rendered = format!("{request:?} {relay:?} {open:?} {error:?}"); + for secret in [ + request_secret, + path_secret, + relay_secret, + "A safe public message.", + ] { + assert!(!rendered.contains(secret)); + } + assert!(error.source().is_none()); } #[test] @@ -1250,6 +1506,33 @@ mod tests { assert!(verify_compatibility(&incompatible).is_err()); } + let oversized_relay_error = HarvestCircleAppCore::open_compatible( + compatible.clone(), + RuntimeOpenInputDto { + development_mode: true, + explicit_data_directory: Some("relative/path-must-not-be-read".to_owned()), + relay_input: RelayBootstrapInputDto { + endpoints: (0..=MAX_CONFIGURED_RELAYS) + .map(|index| RelayEndpointDto { + url: format!("wss://relay-{index}.example"), + destination: RelayDestinationDto::Public, + read: true, + write: true, + }) + .collect(), + }, + }, + ) + .err() + .expect("relay bound must reject before path inspection"); + assert!(matches!( + oversized_relay_error, + HarvestCircleError::Failure { + code: WireErrorCode::InvalidRelayConfiguration, + .. + } + )); + let directory = tempfile::tempdir().expect("directory"); let canonical = directory .path() @@ -1397,52 +1680,58 @@ mod tests { } #[test] - fn invalid_relay_configuration_preserves_local_startup_as_degraded() { - let problem = SafeError::new( - harvestcircle_domain::SafeErrorCode::InvalidRelayConfiguration, - harvestcircle_domain::SafeMessage::new("The Nostr relay configuration is invalid."), - ); - let (relays, degraded) = local_first_relay_configuration(Err(problem)); - - assert!(relays.relays().is_empty()); - assert_eq!(degraded, Some(problem)); + fn invalid_relay_configuration_fails_before_runtime_mutation() { + assert!(relay_configuration_from_endpoints(&[]).is_err()); } #[test] - fn injected_relay_input_is_explicit_mixed_and_fail_closed() { - let mixed = relay_configuration_from_endpoints(&[ + fn injected_relay_input_is_explicit_profile_bound_and_fail_closed() { + let local = relay_configuration_from_endpoints(&[ RelayEndpointInput::new( - "ws://localhost:8080", - RelayDestinationPolicy::Local, + "ws://localhost:8080".to_owned(), + RelayUrlPolicy::Local, true, true, ), RelayEndpointInput::new( - "wss://relay.example", - RelayDestinationPolicy::Public, + "ws://127.0.0.1:8081".to_owned(), + RelayUrlPolicy::Local, true, true, ), ]) - .expect("explicit mixed relays"); - assert_eq!(mixed.relays()[0].url().as_str(), "ws://localhost:8080/"); - assert_eq!(mixed.relays()[1].url().as_str(), "wss://relay.example/"); + .expect("explicit local profile"); + assert_eq!(local.relays()[0].url().as_str(), "ws://localhost:8080"); + assert_eq!(local.relays()[1].url().as_str(), "ws://127.0.0.1:8081"); for input in [ Vec::new(), vec![RelayEndpointInput::new( - "https://not-a-relay.example", - RelayDestinationPolicy::Public, + "https://not-a-relay.example".to_owned(), + RelayUrlPolicy::Public, true, true, )], + vec![ + RelayEndpointInput::new( + "ws://localhost:8080".to_owned(), + RelayUrlPolicy::Local, + true, + true, + ), + RelayEndpointInput::new( + "wss://relay.example".to_owned(), + RelayUrlPolicy::Public, + true, + true, + ), + ], ] { - let (relays, degraded) = - local_first_relay_configuration(relay_configuration_from_endpoints(&input)); - assert!(relays.relays().is_empty()); assert_eq!( - degraded.map(|problem| problem.code()), - Some(harvestcircle_domain::SafeErrorCode::InvalidRelayConfiguration) + relay_configuration_from_endpoints(&input) + .expect_err("invalid profile") + .code(), + harvestcircle_domain::SafeErrorCode::InvalidRelayConfiguration ); } } diff --git a/core/crates/harvestcircle_ffi/src/dto.rs b/core/crates/harvestcircle_ffi/src/dto.rs @@ -1,10 +1,11 @@ +use std::fmt::{self, Formatter}; + use harvestcircle_application::{ ActiveIdentitySnapshot, AppLifecycle, AppSnapshot, ProfileLoadState, RelayConnectionState, - RuntimeLifecycle, SessionState, + RelayEndpoint, RelayUrlPolicy, RuntimeLifecycle, SessionState, }; use harvestcircle_domain::{ - NostrIdentity, ProfileMetadata, RelayDestinationPolicy, RelayEndpoint, SafeError, - SafeErrorCode, SignerAvailability, + NostrIdentity, ProfileMetadata, SafeError, SafeErrorCode, SignerAvailability, }; use harvestcircle_nostr::{NostrReferenceKind, NostrReferenceParse}; @@ -171,7 +172,7 @@ pub enum RelayDestinationDto { Public, } -#[derive(Clone, Debug, Eq, PartialEq)] +#[derive(Clone, Eq, PartialEq)] #[cfg_attr(not(coverage_nightly), derive(uniffi::Record))] pub struct RelayEndpointDto { pub url: String, @@ -180,6 +181,18 @@ pub struct RelayEndpointDto { pub write: bool, } +impl fmt::Debug for RelayEndpointDto { + fn fmt(&self, formatter: &mut Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("RelayEndpointDto") + .field("url", &"<redacted>") + .field("destination", &self.destination) + .field("read", &self.read) + .field("write", &self.write) + .finish() + } +} + #[derive(Clone, Copy, Debug, Eq, PartialEq)] #[cfg_attr(not(coverage_nightly), derive(uniffi::Enum))] pub enum SignerBindingKindDto { @@ -334,19 +347,20 @@ impl From<&RelayEndpoint> for RelayEndpointDto { fn from(endpoint: &RelayEndpoint) -> Self { Self { url: endpoint.url().as_str().to_owned(), - destination: endpoint.destination().into(), - read: endpoint.can_read(), - write: endpoint.can_write(), + destination: endpoint.policy().into(), + read: endpoint.access().can_read(), + write: endpoint.access().can_write(), } } } -impl From<RelayDestinationPolicy> for RelayDestinationDto { - fn from(destination: RelayDestinationPolicy) -> Self { +impl From<RelayUrlPolicy> for RelayDestinationDto { + fn from(destination: RelayUrlPolicy) -> Self { match destination { - RelayDestinationPolicy::Local => Self::Local, - RelayDestinationPolicy::PrivateNetwork => Self::PrivateNetwork, - RelayDestinationPolicy::Public => Self::Public, + RelayUrlPolicy::Local => Self::Local, + RelayUrlPolicy::PrivateNetwork => Self::PrivateNetwork, + RelayUrlPolicy::Public => Self::Public, + _ => unreachable!("validated relay policy is outside the governed v1 vocabulary"), } } } diff --git a/core/crates/harvestcircle_ffi/src/host_runtime.rs b/core/crates/harvestcircle_ffi/src/host_runtime.rs @@ -0,0 +1,166 @@ +use std::future::Future; +use std::sync::{Arc, Mutex}; +use std::thread::JoinHandle; +use std::time::Duration; + +use tokio::runtime::{Builder, Handle}; +use tokio::sync::{oneshot, watch}; + +pub(crate) struct HostRuntime { + handle: Handle, + shutdown: Mutex<Option<oneshot::Sender<()>>>, + completion: watch::Receiver<bool>, + thread: Mutex<Option<JoinHandle<()>>>, +} + +#[cfg(test)] +pub(crate) struct CompletionGatedHostRuntime { + pub(crate) runtime: Arc<HostRuntime>, + pub(crate) entered: std::sync::mpsc::Receiver<()>, + pub(crate) release: std::sync::mpsc::SyncSender<()>, +} + +impl HostRuntime { + pub(crate) fn new() -> Result<Arc<Self>, ()> { + Self::new_inner(None) + } + + fn new_inner( + completion_gate: Option<( + std::sync::mpsc::SyncSender<()>, + std::sync::mpsc::Receiver<()>, + )>, + ) -> Result<Arc<Self>, ()> { + let (startup_sender, startup_receiver) = std::sync::mpsc::sync_channel(1); + let (shutdown_sender, shutdown_receiver) = oneshot::channel(); + let (completion_sender, completion_receiver) = watch::channel(false); + let thread = std::thread::Builder::new() + .name("harvestcircle-host-runtime".to_owned()) + .spawn(move || { + let Ok(runtime) = Builder::new_multi_thread() + .enable_all() + .thread_name("harvestcircle-runtime-worker") + .build() + else { + let _ = startup_sender.send(Err(())); + let _ = completion_sender.send(true); + return; + }; + if startup_sender.send(Ok(runtime.handle().clone())).is_err() { + let _ = completion_sender.send(true); + return; + } + runtime.block_on(async { + let _ = shutdown_receiver.await; + }); + runtime.shutdown_timeout(Duration::from_secs(5)); + if let Some((entered, release)) = completion_gate { + let _ = entered.send(()); + let _ = release.recv(); + } + let _ = completion_sender.send(true); + }) + .map_err(|_| ())?; + let handle = startup_receiver.recv().map_err(|_| ())??; + Ok(Arc::new(Self { + handle, + shutdown: Mutex::new(Some(shutdown_sender)), + completion: completion_receiver, + thread: Mutex::new(Some(thread)), + })) + } + + #[cfg(test)] + pub(crate) fn new_completion_gated_for_test() -> Result<CompletionGatedHostRuntime, ()> { + let (entered_sender, entered_receiver) = std::sync::mpsc::sync_channel(1); + let (release_sender, release_receiver) = std::sync::mpsc::sync_channel(1); + let runtime = Self::new_inner(Some((entered_sender, release_receiver)))?; + Ok(CompletionGatedHostRuntime { + runtime, + entered: entered_receiver, + release: release_sender, + }) + } + + pub(crate) fn handle(&self) -> &Handle { + &self.handle + } + + pub(crate) fn block_on<F>(&self, future: F) -> Result<F::Output, ()> + where + F: Future + Send + 'static, + F::Output: Send + 'static, + { + let (sender, receiver) = std::sync::mpsc::sync_channel(1); + self.handle.spawn(async move { + let _ = sender.send(future.await); + }); + receiver.recv().map_err(|_| ()) + } + + pub(crate) async fn shutdown(&self) -> Result<(), ()> { + let sender = self.shutdown.lock().map_err(|_| ())?.take(); + if let Some(sender) = sender { + let _ = sender.send(()); + } + + let mut completion = self.completion.clone(); + while !*completion.borrow() { + completion.changed().await.map_err(|_| ())?; + } + + let thread = self.thread.lock().map_err(|_| ())?.take(); + if let Some(thread) = thread { + tokio::task::spawn_blocking(move || thread.join().map_err(|_| ())) + .await + .map_err(|_| ())??; + } + Ok(()) + } +} + +impl Drop for HostRuntime { + fn drop(&mut self) { + if let Ok(sender) = self.shutdown.get_mut() + && let Some(sender) = sender.take() + { + let _ = sender.send(()); + } + } +} + +#[cfg(test)] +mod tests { + use std::sync::Arc; + + use super::HostRuntime; + + #[test] + fn host_runtime_executes_work_and_shuts_down_explicitly() { + let host = HostRuntime::new().expect("host runtime"); + assert_eq!(host.block_on(async { 7 }).expect("runtime result"), 7); + let test_runtime = tokio::runtime::Runtime::new().expect("test runtime"); + test_runtime.block_on(host.shutdown()).expect("shutdown"); + test_runtime + .block_on(host.shutdown()) + .expect("idempotent shutdown"); + } + + #[tokio::test] + async fn cancelled_shutdown_is_resumable() { + let gated = HostRuntime::new_completion_gated_for_test().expect("host runtime"); + let host = gated.runtime; + let entered_receiver = gated.entered; + let release_sender = gated.release; + let first_host = Arc::clone(&host); + let first = tokio::spawn(async move { first_host.shutdown().await }); + tokio::task::spawn_blocking(move || entered_receiver.recv()) + .await + .expect("entered join") + .expect("shutdown entered completion gate"); + first.abort(); + assert!(first.await.is_err()); + release_sender.send(()).expect("release shutdown"); + host.shutdown().await.expect("resumed shutdown"); + } +} diff --git a/core/crates/harvestcircle_ffi/src/keyring_worker.rs b/core/crates/harvestcircle_ffi/src/keyring_worker.rs @@ -0,0 +1,200 @@ +use std::sync::{Arc, Mutex}; +use std::thread::JoinHandle; + +use harvestcircle_application::SecretStore; +use harvestcircle_domain::{PublicKey, SafeError, SafeErrorCode, SafeMessage, SecretKeyInput}; +use tokio::sync::watch; + +const KEYRING_QUEUE_CAPACITY: usize = 8; + +enum Request { + Put( + PublicKey, + SecretKeyInput, + std::sync::mpsc::SyncSender<Result<(), SafeError>>, + ), + Load( + PublicKey, + std::sync::mpsc::SyncSender<Result<SecretKeyInput, SafeError>>, + ), + Contains( + PublicKey, + std::sync::mpsc::SyncSender<Result<bool, SafeError>>, + ), + Delete( + PublicKey, + std::sync::mpsc::SyncSender<Result<(), SafeError>>, + ), + Close, +} + +pub(crate) struct BoundedKeyringWorker { + sender: Mutex<Option<std::sync::mpsc::SyncSender<Request>>>, + completion: watch::Receiver<bool>, + thread: Mutex<Option<JoinHandle<()>>>, +} + +impl BoundedKeyringWorker { + pub(crate) fn new(store: impl SecretStore + 'static) -> Result<Arc<Self>, SafeError> { + let (sender, receiver) = std::sync::mpsc::sync_channel(KEYRING_QUEUE_CAPACITY); + let (completion_sender, completion_receiver) = watch::channel(false); + let thread = std::thread::Builder::new() + .name("harvestcircle-keyring-worker".to_owned()) + .spawn(move || { + while let Ok(request) = receiver.recv() { + match request { + Request::Put(public_key, secret, response) => { + let _ = response.send(store.put(public_key, secret)); + } + Request::Load(public_key, response) => { + let _ = response.send(store.load(public_key)); + } + Request::Contains(public_key, response) => { + let _ = response.send(store.contains(public_key)); + } + Request::Delete(public_key, response) => { + let _ = response.send(store.delete(public_key)); + } + Request::Close => break, + } + } + let _ = completion_sender.send(true); + }) + .map_err(|_| worker_unavailable())?; + Ok(Arc::new(Self { + sender: Mutex::new(Some(sender)), + completion: completion_receiver, + thread: Mutex::new(Some(thread)), + })) + } + + fn submit<T>( + &self, + request: impl FnOnce(std::sync::mpsc::SyncSender<T>) -> Request, + ) -> Result<T, SafeError> { + let (response_sender, response_receiver) = std::sync::mpsc::sync_channel(1); + let sender_guard = self.sender.lock().map_err(|_| worker_unavailable())?; + let Some(sender) = sender_guard.as_ref() else { + return Err(worker_unavailable()); + }; + sender + .try_send(request(response_sender)) + .map_err(|_| worker_unavailable())?; + drop(sender_guard); + response_receiver.recv().map_err(|_| worker_unavailable()) + } + + pub(crate) async fn close(&self) -> Result<(), SafeError> { + let sender = self.sender.lock().map_err(|_| worker_unavailable())?.take(); + if let Some(sender) = sender { + signal_close(sender); + } + let mut completion = self.completion.clone(); + while !*completion.borrow() { + completion + .changed() + .await + .map_err(|_| worker_unavailable())?; + } + let thread = self.thread.lock().map_err(|_| worker_unavailable())?.take(); + if let Some(thread) = thread { + tokio::task::spawn_blocking(move || thread.join().map_err(|_| worker_unavailable())) + .await + .map_err(|_| worker_unavailable())??; + } + Ok(()) + } +} + +fn signal_close(sender: std::sync::mpsc::SyncSender<Request>) { + match sender.try_send(Request::Close) { + Ok(()) + | Err(std::sync::mpsc::TrySendError::Full(Request::Close)) + | Err(std::sync::mpsc::TrySendError::Disconnected(Request::Close)) => {} + Err( + std::sync::mpsc::TrySendError::Full(_) | std::sync::mpsc::TrySendError::Disconnected(_), + ) => { + unreachable!("close signaling constructs only close requests") + } + } +} + +impl SecretStore for BoundedKeyringWorker { + fn put(&self, public_key: PublicKey, secret: SecretKeyInput) -> Result<(), SafeError> { + self.submit(|response| Request::Put(public_key, secret, response))? + } + + fn load(&self, public_key: PublicKey) -> Result<SecretKeyInput, SafeError> { + self.submit(|response| Request::Load(public_key, response))? + } + + fn contains(&self, public_key: PublicKey) -> Result<bool, SafeError> { + self.submit(|response| Request::Contains(public_key, response))? + } + + fn delete(&self, public_key: PublicKey) -> Result<(), SafeError> { + self.submit(|response| Request::Delete(public_key, response))? + } +} + +impl Drop for BoundedKeyringWorker { + fn drop(&mut self) { + if let Ok(sender) = self.sender.get_mut() + && let Some(sender) = sender.take() + { + let _ = sender.try_send(Request::Close); + } + } +} + +const fn worker_unavailable() -> SafeError { + SafeError::new( + SafeErrorCode::KeyringUnavailable, + SafeMessage::new("The operating system credential store is unavailable."), + ) +} + +#[cfg(test)] +mod tests { + use harvestcircle_application::{InMemorySecretStore, SecretStore}; + use harvestcircle_domain::{PublicKey, SecretKeyInput}; + + use super::{BoundedKeyringWorker, Request, signal_close}; + + fn public_key() -> PublicKey { + PublicKey::from_hex("7e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7") + .expect("public key") + } + + #[tokio::test] + async fn worker_round_trips_without_exposing_secret_material() { + let worker = BoundedKeyringWorker::new(InMemorySecretStore::default()).expect("worker"); + let secret = SecretKeyInput::parse( + "0000000000000000000000000000000000000000000000000000000000000001".to_owned(), + ) + .expect("secret"); + worker.put(public_key(), secret).expect("put"); + assert!(worker.contains(public_key()).expect("contains")); + let loaded = worker.load(public_key()).expect("load"); + assert_eq!(loaded.with_exposed_secret(str::len), 64); + worker.delete(public_key()).expect("delete"); + worker.close().await.expect("close"); + assert!(worker.contains(public_key()).is_err()); + } + + #[test] + fn close_signal_never_blocks_on_a_full_bounded_queue() { + let (sender, receiver) = std::sync::mpsc::sync_channel(1); + let (response, _response_receiver) = std::sync::mpsc::sync_channel(1); + assert!( + sender + .try_send(Request::Contains(public_key(), response)) + .is_ok() + ); + + signal_close(sender); + + assert!(matches!(receiver.recv(), Ok(Request::Contains(_, _)))); + assert!(receiver.recv().is_err()); + } +} diff --git a/core/crates/harvestcircle_ffi/src/lib.rs b/core/crates/harvestcircle_ffi/src/lib.rs @@ -4,6 +4,8 @@ mod commands; mod contract; mod dto; +mod host_runtime; +mod keyring_worker; mod observer; pub use commands::{ diff --git a/core/crates/harvestcircle_ffi/src/observer.rs b/core/crates/harvestcircle_ffi/src/observer.rs @@ -39,19 +39,20 @@ pub struct ObserverSubscription { #[cfg_attr(not(coverage_nightly), uniffi::export)] impl ObserverSubscription { pub async fn unsubscribe(&self) { - let id = self - .id - .lock() - .unwrap_or_else(std::sync::PoisonError::into_inner) - .take(); + let id = { + let Ok(mut retained_id) = self.id.lock() else { + return; + }; + retained_id.take() + }; let (Some(core), Some(id)) = (self.core.upgrade(), id) else { return; }; let task = { - core.observers - .lock() - .unwrap_or_else(std::sync::PoisonError::into_inner) - .remove(&id) + let Ok(mut observers) = core.observers.lock() else { + return; + }; + observers.remove(&id) }; if let Some(Some(task)) = task { task.abort(); @@ -72,7 +73,7 @@ impl HarvestCircleAppCore { &self, observer: Box<dyn HarvestCircleChangeObserver>, ) -> Result<Arc<ObserverSubscription>, HarvestCircleError> { - if self.inner.closed.load(Ordering::Acquire) { + if !self.inner.is_open() { return Err(closed_error()); } let mut subscription = self @@ -89,8 +90,8 @@ impl HarvestCircleAppCore { .inner .observers .lock() - .unwrap_or_else(std::sync::PoisonError::into_inner); - if self.inner.closed.load(Ordering::Acquire) || observers.len() >= MAX_OBSERVERS { + .map_err(|_| observer_registration_error())?; + if !self.inner.is_open() || observers.len() >= MAX_OBSERVERS { false } else { observers.insert(id, None); @@ -105,7 +106,7 @@ impl HarvestCircleAppCore { .map_err(HarvestCircleError::from)?; return Err(observer_registration_error()); } - let task = crate::commands::runtime()?.spawn(async move { + let task = self.inner.runtime.spawn(async move { while let Some(change) = subscription.receive().await { let Some(runtime_core) = runtime_core.upgrade() else { break; @@ -125,19 +126,16 @@ impl HarvestCircleAppCore { } if let Some(runtime_core) = runtime_core.upgrade() { let _ = runtime_core.actor.unsubscribe_changes(id).await; - runtime_core - .observers - .lock() - .unwrap_or_else(std::sync::PoisonError::into_inner) - .remove(&id); + if let Ok(mut observers) = runtime_core.observers.lock() { + observers.remove(&id); + } } }); let retained = { - let mut observers = self - .inner - .observers - .lock() - .unwrap_or_else(std::sync::PoisonError::into_inner); + let mut observers = self.inner.observers.lock().map_err(|_| { + task.abort(); + observer_registration_error() + })?; if let Some(slot) = observers.get_mut(&id) { *slot = Some(task); true @@ -156,13 +154,22 @@ impl HarvestCircleAppCore { })) } - /// Stops observer delivery and waits for actor-owned shutdown. + /// Stops admission and waits for observer, actor, keyring, and runtime shutdown. + /// + /// Once shutdown begins, dropping or cancelling the calling future does + /// not reopen admission. A later call resumes the same close sequence and + /// successful calls are idempotent. /// /// # Errors /// /// Returns a safe closed or timeout error when shutdown cannot complete. pub async fn shutdown_v2(&self) -> Result<ShutdownReceiptDto, HarvestCircleError> { - if self.inner.closed.swap(true, Ordering::AcqRel) { + let _ = self + .inner + .close_state + .compare_exchange(0, 1, Ordering::AcqRel, Ordering::Acquire); + let _close = self.inner.close_gate.lock().await; + if self.inner.close_state.load(Ordering::Acquire) == 2 { return Ok(ShutdownReceiptDto { final_revision: self.inner.actor.snapshot().revision().value(), closed: true, @@ -173,19 +180,30 @@ impl HarvestCircleAppCore { .inner .observers .lock() - .unwrap_or_else(std::sync::PoisonError::into_inner), + .map_err(|_| crate::commands::internal_state_unavailable())?, ); - for (_, task) in handles { - if let Some(task) = task { - task.abort(); - let _ = task.await; - } + let tasks = handles + .into_values() + .flatten() + .collect::<Vec<tokio::task::JoinHandle<()>>>(); + for task in &tasks { + task.abort(); + } + for task in tasks { + let _ = task.await; } self.inner .actor .close() .await .map_err(HarvestCircleError::from)?; + if let Some(keyring) = self.inner.keyring.as_ref() { + keyring.close().await.map_err(HarvestCircleError::from)?; + } + if let Some(runtime) = self.inner.host_runtime.as_ref() { + runtime.shutdown().await.map_err(|()| closed_error())?; + } + self.inner.close_state.store(2, Ordering::Release); Ok(ShutdownReceiptDto { final_revision: self.inner.actor.snapshot().revision().value(), closed: true, @@ -194,14 +212,7 @@ impl HarvestCircleAppCore { } fn closed_error() -> HarvestCircleError { - HarvestCircleError::Failure { - code: crate::WireErrorCode::InvalidApplicationState, - category: crate::WireErrorCategory::Lifecycle, - retryable: false, - recovery_action: crate::WireRecoveryAction::None, - correlation_id: None, - safe_message: "The application runtime is closed.".to_owned(), - } + crate::commands::runtime_closed_error() } fn observer_registration_error() -> HarvestCircleError { @@ -221,13 +232,14 @@ mod tests { use std::sync::{Arc, Mutex}; use std::time::Duration; - use harvestcircle_application::RelayConfiguration; - use harvestcircle_domain::{RelayDestinationPolicy, RelayEndpoint}; + use harvestcircle_application::{ + RelayAccess, RelayConfiguration, RelayEndpoint, RelayUrlPolicy, + }; use nostr::{EventBuilder, Keys, Metadata}; use nostr_relay_builder::MockRelay; use nostr_sdk::Client; - use crate::commands::{RuntimeCore, runtime, test_actor}; + use crate::commands::{RuntimeCore, test_actor}; use crate::{ AppSnapshotDto, HarvestCircleAppCore, HarvestCircleChangeObserver, ProfileLoadStateDto, SnapshotChangeDto, @@ -267,17 +279,42 @@ mod tests { Arc::new(HarvestCircleAppCore { inner: Arc::new(RuntimeCore { actor, + runtime: tokio::runtime::Handle::current(), + host_runtime: None, + keyring: None, + observers: Mutex::new(std::collections::BTreeMap::new()), + close_state: std::sync::atomic::AtomicU8::new(0), + close_gate: tokio::sync::Mutex::new(()), + _test_directory: Some(directory), + }), + }) + } + + async fn core_with_host_runtime( + host_runtime: Arc<crate::host_runtime::HostRuntime>, + ) -> Arc<HarvestCircleAppCore> { + let (actor, directory) = test_actor(RelayConfiguration::default()).await; + Arc::new(HarvestCircleAppCore { + inner: Arc::new(RuntimeCore { + actor, + runtime: tokio::runtime::Handle::current(), + host_runtime: Some(host_runtime), + keyring: None, observers: Mutex::new(std::collections::BTreeMap::new()), - closed: std::sync::atomic::AtomicBool::new(false), - startup_relay_problem: None, + close_state: std::sync::atomic::AtomicU8::new(0), + close_gate: tokio::sync::Mutex::new(()), _test_directory: Some(directory), }), }) } + fn test_runtime() -> tokio::runtime::Runtime { + tokio::runtime::Runtime::new().expect("test runtime") + } + #[test] fn callbacks_allow_reentry_and_stop_after_subscription_close() { - runtime().expect("runtime").block_on(async { + test_runtime().block_on(async { let core = core().await; let observer = Arc::new(RecordingObserver::default()); *observer.core.lock().expect("core") = Some(Arc::clone(&core)); @@ -301,7 +338,7 @@ mod tests { #[test] fn core_close_deregisters_all_observers_and_rejects_new_subscriptions() { - runtime().expect("runtime").block_on(async { + test_runtime().block_on(async { let core = core().await; let observer = Arc::new(RecordingObserver::default()); let subscription = core @@ -341,9 +378,31 @@ mod tests { }); } + #[tokio::test] + async fn cancelled_host_close_remains_non_admitting_and_resumes() { + let gated = crate::host_runtime::HostRuntime::new_completion_gated_for_test() + .expect("host runtime"); + let core = core_with_host_runtime(gated.runtime).await; + let closing_core = Arc::clone(&core); + let closing = tokio::spawn(async move { closing_core.shutdown_v2().await }); + tokio::task::spawn_blocking(move || gated.entered.recv()) + .await + .expect("entered join") + .expect("close reached host completion gate"); + closing.abort(); + assert!(closing.await.is_err()); + assert!( + core.subscribe_changes_v2(Box::new(PanickingObserver)) + .await + .is_err() + ); + gated.release.send(()).expect("release host close"); + assert!(core.shutdown_v2().await.expect("resumed close").closed); + } + #[test] fn subscription_unsubscribe_tolerates_a_dropped_runtime_core() { - runtime().expect("runtime").block_on(async { + test_runtime().block_on(async { let core = core().await; let observer = Arc::new(RecordingObserver::default()); let subscription = core @@ -359,7 +418,7 @@ mod tests { #[test] fn observer_registration_is_bounded_and_callback_panics_are_contained() { - runtime().expect("runtime").block_on(async { + test_runtime().block_on(async { let core = core().await; let panic_subscription = core .subscribe_changes_v2(Box::new(PanickingObserver)) @@ -409,11 +468,10 @@ mod tests { let core = core_with_relays( RelayConfiguration::new(vec![ - RelayEndpoint::parse( + RelayEndpoint::new( relay_url.as_str(), - RelayDestinationPolicy::Local, - true, - true, + RelayUrlPolicy::Local, + RelayAccess::ReadWrite, ) .expect("relay endpoint"), ]) diff --git a/core/crates/harvestcircle_nostr/Cargo.toml b/core/crates/harvestcircle_nostr/Cargo.toml @@ -13,14 +13,16 @@ include = ["src/**", "Cargo.toml"] [dependencies] nostr = { git = "https://github.com/rust-nostr/nostr.git", rev = "5bba5163eb77107f82c4a8262cf29d7f33a73219", package = "nostr" } -nostr-sdk = { git = "https://github.com/rust-nostr/nostr.git", rev = "5bba5163eb77107f82c4a8262cf29d7f33a73219", package = "nostr-sdk" } harvestcircle_application.workspace = true harvestcircle_domain.workspace = true radroots_identity.workspace = true +radroots_transport.workspace = true +radroots_transport_nostr.workspace = true tokio = { version = "=1.47.1", features = ["sync", "time"] } [dev-dependencies] nostr-relay-builder = { git = "https://github.com/rust-nostr/nostr.git", rev = "5bba5163eb77107f82c4a8262cf29d7f33a73219", package = "nostr-relay-builder" } +nostr-sdk = { git = "https://github.com/rust-nostr/nostr.git", rev = "5bba5163eb77107f82c4a8262cf29d7f33a73219", package = "nostr-sdk" } tokio = { version = "=1.47.1", features = ["macros", "rt-multi-thread", "sync", "time"] } [lints] diff --git a/core/crates/harvestcircle_nostr/src/client.rs b/core/crates/harvestcircle_nostr/src/client.rs @@ -1,25 +1,50 @@ -use std::time::{Duration, Instant}; - -use harvestcircle_domain::{ - PublicKey, RelayEndpoint, SafeError, SafeErrorCode, SafeMessage, select_latest_kind0, -}; -use nostr::{Filter, JsonUtil, Kind, PublicKey as NostrPublicKey}; -use nostr_sdk::Client; +use core::fmt; +use std::sync::atomic::{AtomicU64, Ordering}; +use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH}; use harvestcircle_application::{ BoxFuture, MAX_CONFIGURED_RELAYS, NostrClient, ProfileFetchResult, }; +use harvestcircle_domain::{PublicKey, SafeError, SafeErrorCode, SafeMessage, select_latest_kind0}; +use radroots_transport::outcome::FetchTargetState; +use radroots_transport::source::{FetchBounds, FetchSelector}; +use radroots_transport::{EventSource, FetchRequest, TargetSet}; +use radroots_transport_nostr::{Config, NostrTransport, RelayEndpoint, RelayProfile}; + +const MAX_PROFILE_EVENTS_PER_FETCH: u16 = 64; pub struct SdkNostrClient { timeout: Duration, + next_request: AtomicU64, } -const MAX_PROFILE_EVENTS_PER_RELAY: usize = 64; +impl fmt::Debug for SdkNostrClient { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("SdkNostrClient") + .field("timeout", &self.timeout) + .field("transport", &"[sealed]") + .finish() + } +} impl SdkNostrClient { #[must_use] pub const fn new(timeout: Duration) -> Self { - Self { timeout } + Self { + timeout, + next_request: AtomicU64::new(1), + } + } + + fn request_id(&self) -> Result<String, SafeError> { + let sequence = self + .next_request + .fetch_update(Ordering::AcqRel, Ordering::Acquire, |value| { + value.checked_add(1) + }) + .map_err(|_| relay_connection_failed())?; + Ok(format!("harvest-profile-{sequence}")) } } @@ -31,65 +56,78 @@ impl NostrClient for SdkNostrClient { deadline: Instant, ) -> BoxFuture<'a, Result<ProfileFetchResult, SafeError>> { Box::pin(async move { - let readable_relays = relays - .iter() - .filter(|endpoint| endpoint.can_read()) - .collect::<Vec<_>>(); - if readable_relays.is_empty() { + if relays.is_empty() || relays.len() > MAX_CONFIGURED_RELAYS { return Err(invalid_relay_configuration()); } - if relays.len() > MAX_CONFIGURED_RELAYS { + let readable = relays + .iter() + .filter(|endpoint| endpoint.access().can_read()) + .cloned() + .collect::<Vec<_>>(); + if readable.is_empty() { return Err(invalid_relay_configuration()); } - let author = NostrPublicKey::from_slice(public_key.as_bytes()) + let profile = RelayProfile::explicit(profile_kind(relays)?, relays.iter().cloned()) .map_err(|_| invalid_relay_configuration())?; - let deadline = deadline.min(Instant::now() + self.timeout); - let mut candidates = Vec::new(); - let mut successful_relays = 0usize; - let filter = Filter::new() - .author(author) - .kind(Kind::Metadata) - .limit(MAX_PROFILE_EVENTS_PER_RELAY); - for relay in &readable_relays { - let relay_url = relay.url().as_str(); - let remaining = deadline.saturating_duration_since(Instant::now()); - if remaining.is_zero() { - break; - } - let client = Client::default(); - let result = match tokio::time::timeout_at(deadline.into(), async { - client - .add_relay(relay_url) - .await - .map_err(|_| relay_connection_failed())?; - client - .try_connect_relay(relay_url, remaining) - .await - .map_err(|_| relay_connection_failed())?; - client - .fetch_events_from([relay_url], filter.clone(), remaining) - .await - .map_err(|_| relay_connection_failed()) - }) - .await - { - Ok(result) => result, - Err(_) => Err(relay_connection_failed()), - }; - client.shutdown().await; - if let Ok(events) = result { - successful_relays += 1; - for event in events { - candidates.push(crate::parse_verified_kind0(&event.as_json(), public_key)?); - } - } + let timeout_ms = u64::try_from(self.timeout.as_millis()) + .ok() + .filter(|value| *value > 0) + .ok_or_else(invalid_relay_configuration)?; + let config = Config::from_profile(profile) + .with_timeouts(timeout_ms, timeout_ms, timeout_ms) + .and_then(|config| config.with_max_connections(readable.len())) + .map_err(|_| invalid_relay_configuration())?; + let transport = NostrTransport::new(config); + + let targets = readable + .iter() + .map(|endpoint| endpoint.url().to_target()) + .collect::<Result<Vec<_>, _>>() + .map_err(|_| invalid_relay_configuration())?; + let targets = TargetSet::new(targets).map_err(|_| invalid_relay_configuration())?; + let remaining = deadline.saturating_duration_since(Instant::now()); + if remaining.is_zero() { + return Err(relay_connection_failed()); } - if successful_relays == 0 { + let deadline_unix_ms = current_unix_ms()? + .checked_add( + u64::try_from(remaining.as_millis()).map_err(|_| relay_connection_failed())?, + ) + .ok_or_else(relay_connection_failed)?; + let bounds = FetchBounds::new(MAX_PROFILE_EVENTS_PER_FETCH, deadline_unix_ms) + .map_err(|_| invalid_relay_configuration())?; + let author = radroots_identity::PublicKey::from_bytes(*public_key.as_bytes()) + .map_err(|_| invalid_relay_configuration())?; + let selector = FetchSelector::all() + .with_kinds(vec![0]) + .and_then(|selector| selector.with_authors(vec![author])) + .map_err(|_| invalid_relay_configuration())?; + let request = FetchRequest::new(self.request_id()?, targets, bounds) + .map_err(|_| invalid_relay_configuration())? + .with_selector(selector); + let page = transport + .fetch(request) + .await + .map_err(|_| relay_connection_failed())?; + + let successful = page + .target_outcomes() + .iter() + .filter(|outcome| outcome.state() == FetchTargetState::Complete) + .count(); + if successful == 0 { return Err(relay_connection_failed()); } + let mut candidates = Vec::with_capacity(page.events().len()); + for observed in page.events() { + candidates.push(crate::parse_verified_kind0( + observed.event().raw_json(), + public_key, + )?); + } let candidate = select_latest_kind0(candidates); - if successful_relays == readable_relays.len() { + if successful == readable.len() { Ok(ProfileFetchResult::complete(candidate)) } else { Ok(ProfileFetchResult::partial(candidate)) @@ -98,10 +136,38 @@ impl NostrClient for SdkNostrClient { } } +fn profile_kind( + relays: &[RelayEndpoint], +) -> Result<radroots_transport_nostr::RelayProfileKind, SafeError> { + let has_local = relays + .iter() + .any(|relay| relay.policy() == radroots_transport_nostr::RelayUrlPolicy::Local); + let has_private = relays + .iter() + .any(|relay| relay.policy() == radroots_transport_nostr::RelayUrlPolicy::PrivateNetwork); + let has_public = relays + .iter() + .any(|relay| relay.policy() == radroots_transport_nostr::RelayUrlPolicy::Public); + match (has_local, has_private, has_public) { + (true, false, false) => Ok(radroots_transport_nostr::RelayProfileKind::Simulator), + (false, false, true) => Ok(radroots_transport_nostr::RelayProfileKind::Public), + (false, true, _) => Ok(radroots_transport_nostr::RelayProfileKind::Device), + _ => Err(invalid_relay_configuration()), + } +} + +fn current_unix_ms() -> Result<u64, SafeError> { + SystemTime::now() + .duration_since(UNIX_EPOCH) + .ok() + .and_then(|duration| u64::try_from(duration.as_millis()).ok()) + .ok_or_else(relay_connection_failed) +} + const fn invalid_relay_configuration() -> SafeError { SafeError::new( SafeErrorCode::InvalidRelayConfiguration, - SafeMessage::new("No Nostr relay is configured."), + SafeMessage::new("The Nostr relay configuration is invalid."), ) } @@ -116,27 +182,22 @@ const fn relay_connection_failed() -> SafeError { mod tests { use std::time::Duration; - use harvestcircle_domain::{ - PublicKey, RelayDestinationPolicy, RelayEndpoint, RelayUrl, SafeErrorCode, - }; + use harvestcircle_application::{NostrClient, RelayFetchCompleteness}; + use harvestcircle_domain::{PublicKey, SafeErrorCode}; use nostr::{EventBuilder, Keys, Metadata}; use nostr_relay_builder::MockRelay; use nostr_sdk::Client; - - use harvestcircle_application::NostrClient; + use radroots_transport_nostr::{RelayAccess, RelayEndpoint, RelayUrlPolicy}; use crate::SdkNostrClient; #[tokio::test] - async fn sdk_client_fetches_verified_profile_from_ephemeral_local_relay() { + async fn governed_transport_fetches_verified_profile_from_local_relay() { let relay = MockRelay::run().await.expect("local relay"); let relay_url = relay.url().await; let keys = Keys::generate(); let publisher = Client::new(keys.clone()); - publisher - .add_relay(relay_url.clone()) - .await - .expect("add relay"); + publisher.add_relay(relay_url.clone()).await.expect("relay"); publisher.connect().await; publisher.wait_for_connection(Duration::from_secs(2)).await; publisher @@ -147,148 +208,96 @@ mod tests { .expect("publish metadata"); let adapter = SdkNostrClient::new(Duration::from_secs(2)); - let domain_relay = endpoint(relay_url.as_str(), RelayDestinationPolicy::Local); - let public_key = - PublicKey::from_bytes(keys.public_key().to_bytes()).expect("valid public key"); + let public_key = PublicKey::from_bytes(keys.public_key().to_bytes()).expect("public key"); let fetched = adapter .fetch_profile( public_key, - &[domain_relay], + &[endpoint(relay_url.as_str(), RelayUrlPolicy::Local)], std::time::Instant::now() + Duration::from_secs(2), ) .await - .expect("fetch profile"); + .expect("profile"); let (profile, completeness) = fetched.into_parts(); - let profile = profile.expect("published profile"); - - assert_eq!(profile.author(), public_key); - assert_eq!(profile.metadata().preferred_name(), Some("Farm Identity")); assert_eq!( - completeness, - harvestcircle_application::RelayFetchCompleteness::Complete + profile + .expect("published profile") + .metadata() + .preferred_name(), + Some("Farm Identity") ); + assert_eq!(completeness, RelayFetchCompleteness::Complete); publisher.shutdown().await; relay.shutdown(); } #[tokio::test] - async fn sdk_client_rejects_empty_configuration_without_network_access() { - let error = SdkNostrClient::new(Duration::from_millis(10)) + async fn governed_transport_rejects_invalid_profiles_before_network_access() { + let adapter = SdkNostrClient::new(Duration::from_millis(10)); + let public_key = PublicKey::from_bytes([7; 32]).expect("public key"); + let empty = adapter .fetch_profile( - PublicKey::from_bytes([7; 32]).expect("valid public key"), + public_key, &[], std::time::Instant::now() + Duration::from_millis(10), ) .await .expect_err("empty relay list"); + assert_eq!(empty.code(), SafeErrorCode::InvalidRelayConfiguration); - assert_eq!(error.code(), SafeErrorCode::InvalidRelayConfiguration); - - let write_only = RelayEndpoint::parse( - "wss://relay.example.test", - RelayDestinationPolicy::Public, - false, - true, - ) - .expect("write-only relay"); - let error = SdkNostrClient::new(Duration::from_millis(10)) - .fetch_profile( - PublicKey::from_bytes([7; 32]).expect("valid public key"), - &[write_only], - std::time::Instant::now() + Duration::from_millis(10), - ) - .await - .expect_err("read capability required"); - assert_eq!(error.code(), SafeErrorCode::InvalidRelayConfiguration); - - let relay = endpoint("wss://relay.example.test", RelayDestinationPolicy::Public); - let too_many = vec![relay; harvestcircle_application::MAX_CONFIGURED_RELAYS + 1]; - let error = SdkNostrClient::new(Duration::from_millis(10)) + let mixed = [ + endpoint("ws://127.0.0.1:7777", RelayUrlPolicy::Local), + endpoint("wss://relay.example", RelayUrlPolicy::Public), + ]; + let error = adapter .fetch_profile( - PublicKey::from_bytes([7; 32]).expect("valid public key"), - &too_many, + public_key, + &mixed, std::time::Instant::now() + Duration::from_millis(10), ) .await - .expect_err("oversized relay list"); + .expect_err("mixed trust profiles"); assert_eq!(error.code(), SafeErrorCode::InvalidRelayConfiguration); } #[tokio::test] - async fn sdk_client_fails_when_no_configured_relay_completes() { - let relay = endpoint("ws://127.0.0.1:1", RelayDestinationPolicy::Local); + async fn governed_transport_fails_when_no_relay_completes() { let error = SdkNostrClient::new(Duration::from_millis(25)) .fetch_profile( - PublicKey::from_bytes([7; 32]).expect("valid public key"), - &[relay], + PublicKey::from_bytes([7; 32]).expect("public key"), + &[endpoint("ws://127.0.0.1:1", RelayUrlPolicy::Local)], std::time::Instant::now() + Duration::from_millis(50), ) .await - .expect_err("all relays unavailable"); + .expect_err("unavailable relay"); assert_eq!(error.code(), SafeErrorCode::RelayConnectionFailed); } - #[tokio::test] - async fn sdk_client_reports_partial_when_one_configured_relay_is_unavailable() { - let relay = MockRelay::run().await.expect("local relay"); - let relay_url = relay.url().await; - let keys = Keys::generate(); - let publisher = Client::new(keys.clone()); - publisher - .add_relay(relay_url.clone()) - .await - .expect("add relay"); - publisher.connect().await; - publisher - .send_event_builder(EventBuilder::metadata(&Metadata::new().name("Partial"))) - .await - .expect("publish metadata"); - - let configured = [ - endpoint(relay_url.as_str(), RelayDestinationPolicy::Local), - endpoint("ws://127.0.0.1:1", RelayDestinationPolicy::Local), - ]; - let fetched = SdkNostrClient::new(Duration::from_millis(250)) - .fetch_profile( - PublicKey::from_bytes(keys.public_key().to_bytes()).expect("valid public key"), - &configured, - std::time::Instant::now() + Duration::from_secs(1), - ) - .await - .expect("partial fetch"); - let (candidate, completeness) = fetched.into_parts(); - assert!(candidate.is_some()); + #[test] + fn debug_and_policy_are_safe_and_lib_owned() { + let adapter = SdkNostrClient::new(Duration::from_secs(1)); assert_eq!( - completeness, - harvestcircle_application::RelayFetchCompleteness::Partial + format!("{adapter:?}"), + "SdkNostrClient { timeout: 1s, transport: \"[sealed]\" }" ); - publisher.shutdown().await; - relay.shutdown(); - } - - #[test] - fn harvestcircle_relay_policy_remains_domain_owned_and_fail_closed() { - assert!(RelayUrl::parse("wss://relay.example", RelayDestinationPolicy::Public).is_ok()); - assert!(RelayUrl::parse("ws://127.0.0.1:7777", RelayDestinationPolicy::Local).is_ok()); assert!( - RelayUrl::parse( - "wss://10.0.0.1:7777", - RelayDestinationPolicy::PrivateNetwork + RelayEndpoint::new( + "wss://relay.example", + RelayUrlPolicy::Public, + RelayAccess::ReadOnly, ) .is_ok() ); - assert!(RelayUrl::parse("ws://relay.example", RelayDestinationPolicy::Public).is_err()); - assert_eq!( - super::invalid_relay_configuration().code(), - SafeErrorCode::InvalidRelayConfiguration - ); - assert_eq!( - super::relay_connection_failed().code(), - SafeErrorCode::RelayConnectionFailed + assert!( + RelayEndpoint::new( + "ws://relay.example", + RelayUrlPolicy::Public, + RelayAccess::ReadOnly, + ) + .is_err() ); } - fn endpoint(value: &str, destination: RelayDestinationPolicy) -> RelayEndpoint { - RelayEndpoint::parse(value, destination, true, true).expect("relay endpoint") + fn endpoint(value: &str, policy: RelayUrlPolicy) -> RelayEndpoint { + RelayEndpoint::new(value, policy, RelayAccess::ReadWrite).expect("relay endpoint") } } diff --git a/core/crates/harvestcircle_runtime/Cargo.toml b/core/crates/harvestcircle_runtime/Cargo.toml @@ -25,6 +25,7 @@ uuid.workspace = true nostr = { git = "https://github.com/rust-nostr/nostr.git", rev = "5bba5163eb77107f82c4a8262cf29d7f33a73219", package = "nostr" } nostr-relay-builder = { git = "https://github.com/rust-nostr/nostr.git", rev = "5bba5163eb77107f82c4a8262cf29d7f33a73219", package = "nostr-relay-builder" } nostr-sdk = { git = "https://github.com/rust-nostr/nostr.git", rev = "5bba5163eb77107f82c4a8262cf29d7f33a73219", package = "nostr-sdk" } +radroots_transport_nostr.workspace = true tempfile = "=3.23.0" [lints] diff --git a/core/crates/harvestcircle_runtime/src/runtime_actor.rs b/core/crates/harvestcircle_runtime/src/runtime_actor.rs @@ -310,10 +310,10 @@ impl RuntimeActorHandle { #[must_use] pub fn lifecycle(&self) -> RuntimeLifecycle { - self.lifecycle - .lock() - .unwrap_or_else(std::sync::PoisonError::into_inner) - .lifecycle() + self.lifecycle.lock().map_or_else( + |_| RuntimeLifecycle::Fatal(runtime_state_unavailable()), + |lifecycle| lifecycle.lifecycle(), + ) } #[must_use] @@ -325,8 +325,8 @@ impl RuntimeActorHandle { pub fn foreground_session(&self) -> Option<ActiveSessionBinding> { self.foreground_session .lock() - .unwrap_or_else(std::sync::PoisonError::into_inner) - .clone() + .map(|session| session.clone()) + .unwrap_or(None) } #[must_use] @@ -598,7 +598,7 @@ impl RuntimeActorHandle { let actor_task = self .actor_task .lock() - .unwrap_or_else(std::sync::PoisonError::into_inner) + .map_err(|_| runtime_state_unavailable())? .take(); if let Some(actor_task) = actor_task { let remaining = deadline.saturating_duration_since(Instant::now()); @@ -883,11 +883,10 @@ impl RuntimeActor { if context.is_expired(Instant::now()) { return Some(CommandResult::TimedOut); } - let lifecycle = self - .lifecycle - .lock() - .unwrap_or_else(std::sync::PoisonError::into_inner) - .to_owned(); + let lifecycle = match self.lifecycle.lock() { + Ok(lifecycle) => lifecycle.to_owned(), + Err(_) => return Some(CommandResult::Failed(runtime_state_unavailable())), + }; if matches!(lifecycle.lifecycle(), RuntimeLifecycle::Closed) { return Some(CommandResult::Closed); } @@ -1102,11 +1101,16 @@ impl RuntimeActor { reply: oneshot::Sender<CommandReceipt<RuntimeCommandValue>>, completion_sender: mpsc::Sender<ProfileCompletion>, ) { - let foreground = self - .published_foreground_session - .lock() - .unwrap_or_else(std::sync::PoisonError::into_inner) - .clone(); + let foreground = match self.published_foreground_session.lock() { + Ok(foreground) => foreground.clone(), + Err(_) => { + let _ = reply.send(CommandReceipt::new( + context.request_id(), + CommandResult::Failed(runtime_state_unavailable()), + )); + return; + } + }; let plan = match self.adapter.core().begin_profile_refresh() { Ok(Some(plan)) => plan, Ok(None) => { @@ -1165,10 +1169,7 @@ impl RuntimeActor { }, ); if let Some(previous) = previous { - self.lifecycle - .lock() - .unwrap_or_else(std::sync::PoisonError::into_inner) - .fail(request_space_exhausted()); + self.fail_lifecycle(request_space_exhausted()); previous.handle.abort(); let _ = previous.handle.await; let _ = previous.reply.send(CommandReceipt::new( @@ -1181,10 +1182,9 @@ impl RuntimeActor { async fn close_actor(&mut self) -> CommandResult<RuntimeCommandValue> { let begin = { - let mut lifecycle = self - .lifecycle - .lock() - .unwrap_or_else(std::sync::PoisonError::into_inner); + let Ok(mut lifecycle) = self.lifecycle.lock() else { + return CommandResult::Failed(runtime_state_unavailable()); + }; lifecycle.begin_shutdown() }; match begin { @@ -1192,23 +1192,19 @@ impl RuntimeActor { self.generated_key_stage.cancel(); self.cancel_profile_tasks(None).await; if let Err(error) = self.adapter.close().await { - self.lifecycle - .lock() - .unwrap_or_else(std::sync::PoisonError::into_inner) - .fail(error); + self.fail_lifecycle(error); return CommandResult::Failed(error); } self.changes.close(); - *self - .published_foreground_session - .lock() - .unwrap_or_else(std::sync::PoisonError::into_inner) = None; - match self - .lifecycle - .lock() - .unwrap_or_else(std::sync::PoisonError::into_inner) - .finish_shutdown() - { + let Ok(mut foreground) = self.published_foreground_session.lock() else { + return CommandResult::Failed(runtime_state_unavailable()); + }; + *foreground = None; + drop(foreground); + let Ok(mut lifecycle) = self.lifecycle.lock() else { + return CommandResult::Failed(runtime_state_unavailable()); + }; + match lifecycle.finish_shutdown() { Ok(()) => CommandResult::Completed(RuntimeCommandValue::Closed), Err(error) => CommandResult::Failed(error), } @@ -1223,11 +1219,17 @@ impl RuntimeActor { }; let _ = task.handle.await; let current = self.adapter.core().snapshot(); - let foreground = self - .published_foreground_session - .lock() - .unwrap_or_else(std::sync::PoisonError::into_inner) - .clone(); + let foreground = match self.published_foreground_session.lock() { + Ok(foreground) => foreground.clone(), + Err(_) => { + self.fail_lifecycle(runtime_state_unavailable()); + let _ = task.reply.send(CommandReceipt::new( + task.correlation.request_id(), + CommandResult::Failed(runtime_state_unavailable()), + )); + return; + } + }; let correlated = task.correlation.session_generation() == self.session_generation && foreground.is_some_and(|binding| { binding.generation() == task.correlation.session_generation() @@ -1268,10 +1270,7 @@ impl RuntimeActor { async fn advance_session_generation(&mut self) { let Some(next) = self.session_generation.next() else { - self.lifecycle - .lock() - .unwrap_or_else(std::sync::PoisonError::into_inner) - .fail(request_space_exhausted()); + self.fail_lifecycle(request_space_exhausted()); self.cancel_profile_tasks(None).await; return; }; @@ -1299,17 +1298,20 @@ impl RuntimeActor { let session = match session.transpose() { Ok(session) => session, Err(error) => { - self.lifecycle - .lock() - .unwrap_or_else(std::sync::PoisonError::into_inner) - .fail(error); + self.fail_lifecycle(error); None } }; - *self - .published_foreground_session - .lock() - .unwrap_or_else(std::sync::PoisonError::into_inner) = session; + match self.published_foreground_session.lock() { + Ok(mut foreground) => *foreground = session, + Err(_) => self.fail_lifecycle(runtime_state_unavailable()), + } + } + + fn fail_lifecycle(&self, error: SafeError) { + if let Ok(mut lifecycle) = self.lifecycle.lock() { + lifecycle.fail(error); + } } async fn cancel_profile_tasks(&mut self, snapshot: Option<&AppSnapshot>) { @@ -1435,6 +1437,13 @@ const fn runtime_closed() -> SafeError { ) } +const fn runtime_state_unavailable() -> SafeError { + SafeError::new( + SafeErrorCode::InvalidApplicationState, + SafeMessage::new("The application runtime state is unavailable."), + ) +} + const fn command_unavailable() -> SafeError { SafeError::new( SafeErrorCode::InvalidApplicationState, @@ -1479,9 +1488,10 @@ mod tests { SecretStore, SecretStoreOperation, SessionGeneration, SessionState, SnapshotRevision, }; use harvestcircle_domain::{ - LocalKeyringBinding, NostrIdentityReference, PublicKey, RelayDestinationPolicy, - RelayEndpoint, SafeError, SafeErrorCode, SecretKeyInput, SignerAvailability, UnixTimestamp, + LocalKeyringBinding, NostrIdentityReference, PublicKey, SafeError, SafeErrorCode, + SecretKeyInput, SignerAvailability, UnixTimestamp, }; + use radroots_transport_nostr::{RelayAccess, RelayEndpoint, RelayUrlPolicy}; use super::{ DEFAULT_COMMAND_TIMEOUT, RuntimeActorHandle, RuntimeDependencies, command_unavailable, @@ -1945,11 +1955,10 @@ mod tests { let client = Arc::new(BlockingNostr::new()); let actor = RuntimeActorHandle::in_memory( RelayConfiguration::new(vec![ - RelayEndpoint::parse( + RelayEndpoint::new( "ws://localhost:8080", - RelayDestinationPolicy::Local, - true, - true, + RelayUrlPolicy::Local, + RelayAccess::ReadWrite, ) .expect("relay"), ]) @@ -1996,11 +2005,10 @@ mod tests { let client = Arc::new(BlockingNostr::new()); let actor = RuntimeActorHandle::in_memory( RelayConfiguration::new(vec![ - RelayEndpoint::parse( + RelayEndpoint::new( "ws://localhost:8080", - RelayDestinationPolicy::Local, - true, - true, + RelayUrlPolicy::Local, + RelayAccess::ReadWrite, ) .expect("relay"), ]) @@ -2340,11 +2348,10 @@ mod tests { let client = Arc::new(BlockingNostr::new()); let actor = RuntimeActorHandle::in_memory( RelayConfiguration::new(vec![ - RelayEndpoint::parse( + RelayEndpoint::new( "ws://localhost:8080", - RelayDestinationPolicy::Local, - true, - true, + RelayUrlPolicy::Local, + RelayAccess::ReadWrite, ) .expect("relay"), ]) diff --git a/core/crates/harvestcircle_runtime/tests/local_relay_e2e.rs b/core/crates/harvestcircle_runtime/tests/local_relay_e2e.rs @@ -4,7 +4,7 @@ use harvestcircle_application::{ Clock, DurableRequestId, InMemorySecretStore, ProfileLoadState, ProfileRepository, RelayConfiguration, RelayConnectionState, SecretStore, SessionState, }; -use harvestcircle_domain::{RelayDestinationPolicy, RelayEndpoint, SecretKeyInput, UnixTimestamp}; +use harvestcircle_domain::{SecretKeyInput, UnixTimestamp}; use harvestcircle_nostr::SdkNostrClient; use harvestcircle_runtime::PersistentAppCore; use nostr::{EventBuilder, Keys, Metadata}; @@ -15,6 +15,7 @@ use radroots_runtime_paths::{ RadrootsPlatform, RuntimeContext, RuntimeContextBootstrap, RuntimeContextSource, ServiceId, }; use radroots_service_sqlite::MigrationBuildIdentity; +use radroots_transport_nostr::{RelayAccess, RelayEndpoint, RelayUrlPolicy}; const SECRET_HEX: &str = "7e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7"; @@ -79,11 +80,10 @@ async fn local_relay_e2e_imports_activates_refreshes_and_caches_profile() { .await .expect("publish profile"); - let relay = RelayEndpoint::parse( + let relay = RelayEndpoint::new( relay_url.as_str(), - RelayDestinationPolicy::Local, - true, - true, + RelayUrlPolicy::Local, + RelayAccess::ReadWrite, ) .expect("relay endpoint"); let adapter = PersistentAppCore::open( diff --git a/core/crates/harvestcircle_storage/src/installation.rs b/core/crates/harvestcircle_storage/src/installation.rs @@ -83,7 +83,7 @@ fn decode_hex(value: &str) -> Result<[u8; 16], SafeError> { return Err(invalid_installation_identity()); } let mut output = [0_u8; 16]; - for (index, pair) in value.as_bytes().chunks_exact(2).enumerate() { + for (index, pair) in value.as_bytes().as_chunks::<2>().0.iter().enumerate() { let high = hex_nibble(pair[0]).ok_or_else(invalid_installation_identity)?; let low = hex_nibble(pair[1]).ok_or_else(invalid_installation_identity)?; output[index] = (high << 4) | low; diff --git a/core/crates/harvestcircle_test_bridge/src/lib.rs b/core/crates/harvestcircle_test_bridge/src/lib.rs @@ -10,12 +10,10 @@ use std::time::Duration; use harvestcircle_application::{ AppLifecycle, AppSnapshot, Clock, DurableRequestId, GeneratedKeyRecoveryHandle, - InMemorySecretStore, RelayConfiguration, RemovalConfirmationToken, SecretStore, SessionState, - SnapshotRevision, -}; -use harvestcircle_domain::{ - PublicKey, RelayDestinationPolicy, RelayEndpoint, SafeError, SecretKeyInput, UnixTimestamp, + InMemorySecretStore, RelayAccess, RelayConfiguration, RelayEndpoint, RelayUrlPolicy, + RemovalConfirmationToken, SecretStore, SessionState, SnapshotRevision, }; +use harvestcircle_domain::{PublicKey, SafeError, SecretKeyInput, UnixTimestamp}; use harvestcircle_nostr::SdkNostrClient; use harvestcircle_runtime::{ InstallationIdentity, InstallationIdentitySource, RuntimeActorHandle, @@ -551,7 +549,10 @@ async fn open_actor( clock: Arc<FixedClock>, runtime: &tokio::runtime::Handle, ) -> Result<RuntimeActorHandle, TestBridgeError> { - let relay = RelayEndpoint::parse(relay_url, RelayDestinationPolicy::Local, true, true)?; + let relay = RelayEndpoint::new(relay_url, RelayUrlPolicy::Local, RelayAccess::ReadWrite) + .map_err(|_| TestBridgeError::Failure { + safe_message: "The local relay configuration is invalid.".to_owned(), + })?; let dependencies = RuntimeDependencies::new( secrets, clock, diff --git a/radroots.lib.source-lock.v1.toml b/radroots.lib.source-lock.v1.toml @@ -5,4 +5,4 @@ architecture = "radroots.crates.release.v2" workspace_catalog_sha256 = "deca0c080deae187ff8186c0708903e42f41ea57f77c5f91581e23aa561164a4" version = "0.1.0-alpha" source_archive_sha256 = "aec2fe198b200f40af81424fbec70a9a8f22b0b38455bc6c81b7eb3be4241748" -lockfile_sha256 = "1bbaae4bd586b936120bb8b2cfab8a5463e7e15aa3d6de0ccdf85f831835467f" +lockfile_sha256 = "1d5187f6394e470a027d3643bc1b74beaf92c5168718783e69f2e3f67dfeb1c0" diff --git a/tools/xtask/src/lib.rs b/tools/xtask/src/lib.rs @@ -231,6 +231,97 @@ fn repo_audit(root: &Path, inventory: &Inventory, findings: &mut Vec<String>) { } } git_source_policy(root, findings); + native_runtime_boundary(root, findings); +} + +fn native_runtime_boundary(root: &Path, findings: &mut Vec<String>) { + let domain_lib = root.join("core/crates/harvestcircle_domain/src/lib.rs"); + if !domain_lib.is_file() { + return; + } + + if root + .join("core/crates/harvestcircle_domain/src/relay.rs") + .exists() + || read_text(root, "core/crates/harvestcircle_domain/src/lib.rs").contains("mod relay") + { + findings + .push("harvestcircle_domain: duplicate relay policy surface is forbidden".to_owned()); + } + + let nostr_manifest = read_text(root, "core/crates/harvestcircle_nostr/Cargo.toml"); + let production_manifest = nostr_manifest + .split_once("[dev-dependencies]") + .map_or(nostr_manifest.as_str(), |(production, _)| production); + if production_manifest.contains("nostr-sdk") { + findings.push( + "harvestcircle_nostr: production nostr-sdk connection authority is forbidden" + .to_owned(), + ); + } + let nostr_client = read_text(root, "core/crates/harvestcircle_nostr/src/client.rs"); + for required in [ + "radroots_transport_nostr::{Config, NostrTransport, RelayEndpoint, RelayProfile}", + "parse_verified_kind0", + "FetchBounds::new(MAX_PROFILE_EVENTS_PER_FETCH", + ] { + if !nostr_client.contains(required) { + findings.push(format!( + "harvestcircle_nostr: governed transport boundary is missing {required}" + )); + } + } + + for (path, forbidden) in [ + ("core/crates/harvestcircle_ffi/src/commands.rs", "OnceLock"), + ( + "core/crates/harvestcircle_ffi/src/commands.rs", + "PoisonError::into_inner", + ), + ( + "core/crates/harvestcircle_ffi/src/observer.rs", + "PoisonError::into_inner", + ), + ( + "core/crates/harvestcircle_application/src/app_core.rs", + "PoisonError::into_inner", + ), + ( + "core/crates/harvestcircle_application/src/custody.rs", + "PoisonError::into_inner", + ), + ( + "core/crates/harvestcircle_application/src/secrets.rs", + "PoisonError::into_inner", + ), + ] { + if read_text(root, path).contains(forbidden) { + findings.push(format!("{path}: forbidden runtime boundary {forbidden}")); + } + } + + let runtime = read_text(root, "core/crates/harvestcircle_ffi/src/host_runtime.rs"); + let keyring = read_text(root, "core/crates/harvestcircle_ffi/src/keyring_worker.rs"); + for (source, required, owner) in [ + (&runtime, "pub(crate) struct HostRuntime", "host runtime"), + (&runtime, "pub(crate) async fn shutdown", "host runtime"), + ( + &keyring, + "const KEYRING_QUEUE_CAPACITY: usize = 8", + "keyring worker", + ), + ( + &keyring, + "pub(crate) struct BoundedKeyringWorker", + "keyring worker", + ), + ] { + if !source.contains(required) { + findings.push(format!( + "harvestcircle_ffi: {owner} contract is missing {required}" + )); + } + } } fn namespace_audit(root: &Path, inventory: &Inventory, findings: &mut Vec<String>) { @@ -757,6 +848,15 @@ fn provenance_check(root: &Path, inventory: &Inventory, findings: &mut Vec<Strin format!( "radroots_service_sqlite = {{ git = \"https://github.com/radrootslabs/lib\", rev = \"{LIB_REVISION}\", version = \"=0.1.0-alpha\", default-features = false }}" ), + format!( + "radroots_storage = {{ git = \"https://github.com/radrootslabs/lib\", rev = \"{LIB_REVISION}\", version = \"=0.1.0-alpha\", default-features = false }}" + ), + format!( + "radroots_transport = {{ git = \"https://github.com/radrootslabs/lib\", rev = \"{LIB_REVISION}\", version = \"=0.1.0-alpha\", default-features = false }}" + ), + format!( + "radroots_transport_nostr = {{ git = \"https://github.com/radrootslabs/lib\", rev = \"{LIB_REVISION}\", version = \"=0.1.0-alpha\", default-features = false }}" + ), ] { if cargo .lines() @@ -789,7 +889,7 @@ fn provenance_check(root: &Path, inventory: &Inventory, findings: &mut Vec<Strin "workspace_catalog_sha256 = \"deca0c080deae187ff8186c0708903e42f41ea57f77c5f91581e23aa561164a4\"\n", "version = \"0.1.0-alpha\"\n", "source_archive_sha256 = \"aec2fe198b200f40af81424fbec70a9a8f22b0b38455bc6c81b7eb3be4241748\"\n", - "lockfile_sha256 = \"1bbaae4bd586b936120bb8b2cfab8a5463e7e15aa3d6de0ccdf85f831835467f\"\n", + "lockfile_sha256 = \"1d5187f6394e470a027d3643bc1b74beaf92c5168718783e69f2e3f67dfeb1c0\"\n", ); if read_text(root, SOURCE_LOCK_PATH) != expected_source_lock { findings.push(format!("{SOURCE_LOCK_PATH}: exact Lib source lock changed"));