lib

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

commit 5d53eb8da344e92453c032609ebef3add0b33fc5
parent 6670b821b3ad5c86b0601b0d10cd8612a6661b6a
Author: triesap <tyson@radroots.org>
Date:   Thu,  6 Aug 2026 10:16:23 +0000

studio-runtime: harden composition and concurrency

- extract runtime composition, preferences, networking, and persistence into canonical package boundaries
- enforce explicit relay policy, bounded work, supervised shutdown, and partial-result semantics
- persist restart-stable installation identity and remove panic-prone global initialization races
- verify Studio architecture, checks, clippy, runtime, FFI, storage, networking, and restart tests

Diffstat:
MCargo.lock | 40++++++++++++++++++++++++++++++++--------
MCargo.toml | 4++++
Mcontracts/coverage.toml | 2++
Mcontracts/crates/catalog.v1.toml | 40++++++++++++++++++++++++++++++++++------
Mcontracts/crates/generated/package_groups.v1.toml | 16++++++++--------
Mcontracts/crates/generated/platform_inventory.v1.toml | 4++--
Mcontracts/crates/generated/release_inventory.v2.toml | 6+++---
Mcontracts/releases/publish_policy.toml | 2++
Mcrates/studio_application/Cargo.toml | 6------
Mcrates/studio_application/src/accounts.rs | 20+++++++++-----------
Mcrates/studio_application/src/app_core.rs | 26++++++++++++++++++++++----
Mcrates/studio_application/src/config.rs | 14++++++++------
Mcrates/studio_application/src/custody.rs | 23++++++++++++++++-------
Mcrates/studio_application/src/lib.rs | 15+++++++++------
Dcrates/studio_application/src/nostr_client.rs | 138-------------------------------------------------------------------------------
Mcrates/studio_application/src/ports.rs | 123+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++------
Mcrates/studio_application/src/profile_refresh.rs | 182+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++------------------
Mcrates/studio_application/src/session.rs | 6++----
Mcrates/studio_application/src/snapshot.rs | 40++++++++++++++++++++++++++++++++++++----
Acrates/studio_application/src/test_support.rs | 48++++++++++++++++++++++++++++++++++++++++++++++++
Mcrates/studio_domain/src/lib.rs | 2+-
Mcrates/studio_domain/src/relay.rs | 107++++++++++++++++++++++++++++++++++++++++++++++++++++++-------------------------
Mcrates/studio_domain/src/time.rs | 2++
Mcrates/studio_ffi/Cargo.toml | 2++
Mcrates/studio_ffi/src/commands.rs | 130++++++++++++++++++++++++++++++++++++++++++++++++++++++++++----------------------
Mcrates/studio_ffi/src/dto.rs | 8+++++++-
Mcrates/studio_ffi/src/observer.rs | 113+++++++++++++++++++++++++++++++++++++++++++------------------------------------
Mcrates/studio_nostr/Cargo.toml | 10++++++++++
Acrates/studio_nostr/src/client.rs | 304+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcrates/studio_nostr/src/keys.rs | 99++++++++++++++++++++++++++++---------------------------------------------------
Mcrates/studio_nostr/src/lib.rs | 4+++-
Acrates/studio_preferences/Cargo.toml | 18++++++++++++++++++
Acrates/studio_preferences/src/lib.rs | 266+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acrates/studio_runtime/Cargo.toml | 29+++++++++++++++++++++++++++++
Acrates/studio_runtime/src/blocking.rs | 114+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acrates/studio_runtime/src/installation.rs | 58++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acrates/studio_runtime/src/lib.rs | 12++++++++++++
Acrates/studio_runtime/src/persistence.rs | 746+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acrates/studio_runtime/src/runtime_actor.rs | 2002+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acrates/studio_runtime/tests/local_relay_e2e.rs | 100+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acrates/studio_runtime/tests/restart_isolation.rs | 95+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcrates/studio_storage/Cargo.toml | 4----
Acrates/studio_storage/migrations/V10__installation_identity.sql | 7+++++++
Dcrates/studio_storage/src/application_adapter.rs | 718-------------------------------------------------------------------------------
Mcrates/studio_storage/src/db.rs | 4++--
Acrates/studio_storage/src/installation.rs | 71+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcrates/studio_storage/src/lib.rs | 5+----
Dcrates/studio_storage/src/runtime_actor.rs | 1693-------------------------------------------------------------------------------
Dcrates/studio_storage/tests/local_relay_e2e.rs | 95-------------------------------------------------------------------------------
Dcrates/studio_storage/tests/restart_isolation.rs | 95-------------------------------------------------------------------------------
Mtools/xtask/src/catalog.rs | 8++++++--
51 files changed, 4613 insertions(+), 3063 deletions(-)

diff --git a/Cargo.lock b/Cargo.lock @@ -3965,11 +3965,7 @@ dependencies = [ name = "radroots_studio_application" version = "0.1.0-alpha" dependencies = [ - "nostr 0.44.1", - "nostr-relay-builder", - "nostr-sdk 0.44.0", "radroots_studio_domain", - "radroots_studio_nostr", "secrecy", "tokio", ] @@ -3994,6 +3990,8 @@ dependencies = [ "nostr-sdk 0.44.0", "radroots_studio_application", "radroots_studio_domain", + "radroots_studio_nostr", + "radroots_studio_runtime", "radroots_studio_storage", "tempfile", "tokio", @@ -4005,24 +4003,50 @@ name = "radroots_studio_nostr" version = "0.1.0-alpha" dependencies = [ "nostr 0.44.1", + "nostr-relay-builder", + "nostr-sdk 0.44.0", + "radroots_identity", + "radroots_studio_application", "radroots_studio_domain", + "radroots_transport", + "radroots_transport_nostr", + "tokio", ] [[package]] -name = "radroots_studio_storage" +name = "radroots_studio_preferences" +version = "0.1.0-alpha" +dependencies = [ + "url", +] + +[[package]] +name = "radroots_studio_runtime" version = "0.1.0-alpha" dependencies = [ - "fs2", - "keyring 4.1.6", "nostr 0.44.1", "nostr-relay-builder", "nostr-sdk 0.44.0", "radroots_studio_application", "radroots_studio_domain", + "radroots_studio_nostr", + "radroots_studio_storage", + "tempfile", + "tokio", + "uuid", +] + +[[package]] +name = "radroots_studio_storage" +version = "0.1.0-alpha" +dependencies = [ + "fs2", + "keyring 4.1.6", + "radroots_studio_application", + "radroots_studio_domain", "refinery", "rusqlite", "tempfile", - "tokio", "zeroize", ] diff --git a/Cargo.toml b/Cargo.toml @@ -44,6 +44,8 @@ members = [ "crates/studio_domain", "crates/studio_ffi", "crates/studio_nostr", + "crates/studio_preferences", + "crates/studio_runtime", "crates/studio_storage", "crates/studio_uniffi_bindgen", "crates/core_bindings", @@ -166,6 +168,8 @@ radroots_studio_application = { path = "crates/studio_application", version = "= radroots_studio_domain = { path = "crates/studio_domain", version = "=0.1.0-alpha" } radroots_studio_ffi = { path = "crates/studio_ffi", version = "=0.1.0-alpha" } radroots_studio_nostr = { path = "crates/studio_nostr", version = "=0.1.0-alpha" } +radroots_studio_preferences = { path = "crates/studio_preferences", version = "=0.1.0-alpha" } +radroots_studio_runtime = { path = "crates/studio_runtime", version = "=0.1.0-alpha" } radroots_studio_storage = { path = "crates/studio_storage", version = "=0.1.0-alpha" } radroots_studio_uniffi_bindgen = { path = "crates/studio_uniffi_bindgen", version = "=0.1.0-alpha" } radroots_simplex_app_store = { path = "crates/simplex_app_store", version = "=0.1.0-alpha", default-features = false } diff --git a/contracts/coverage.toml b/contracts/coverage.toml @@ -92,6 +92,8 @@ crates = [ "radroots_studio_domain", "radroots_studio_ffi", "radroots_studio_nostr", + "radroots_studio_preferences", + "radroots_studio_runtime", "radroots_studio_storage", "radroots_studio_uniffi_bindgen", "radroots_sync", diff --git a/contracts/crates/catalog.v1.toml b/contracts/crates/catalog.v1.toml @@ -7,7 +7,7 @@ rust_version = "1.97.1" edition = "2024" resolver = "3" public_package_count = 19 -package_count = 60 +package_count = 61 digest_algorithm = "sha256-raw-bytes-v1" source_tree_digest_algorithm = "sha256-git-ls-tree-r-v1" source_provenance_policy = "immutable_after_source_retirement" @@ -1404,6 +1404,35 @@ compatibility = ["product", "data", "package_private"] replaces = ["radroots-studio-domain"] [[package]] +name = "radroots_studio_preferences" +path = "crates/studio_preferences" +state = "active" +tier = "application_domain" +visibility = "private_runtime" +publish = false +version = "0.1.0-alpha" +license = "MPL-2.0" +platforms = ["native"] +groups = ["coverage_required", "studio"] +owners = ["studio"] +permitted_dependency_tiers = [ + "foundation", + "domain", + "spi", + "adapter", + "orchestration", + "sdk", + "facade", + "application_domain", +] +source_repository = "https://github.com/radrootslabs/_radroots" +source_revision = "6074a4745be361f21bb47d4778c74a14b2d57954" +source_path = "studio_app/studio_app_core/crates/core" +source_tree_sha256 = "0237265710a676ce1db0fb3091e1e5d9a5831339e00076ea2a33b96d6343834d" +compatibility = ["product", "preferences", "package_private"] +replaces = ["studio_app_core"] + +[[package]] name = "radroots_studio_application" path = "crates/studio_application" state = "active" @@ -1498,14 +1527,14 @@ replaces = ["radroots-studio-storage"] [[package]] name = "radroots_studio_runtime" path = "crates/studio_runtime" -state = "reserved_refactor" +state = "active" tier = "runtime_composition" visibility = "private_runtime" publish = false version = "0.1.0-alpha" license = "GPL-3.0-only" platforms = ["native"] -groups = ["studio"] +groups = ["coverage_required", "studio"] owners = ["studio"] permitted_dependency_tiers = [ "foundation", @@ -1522,10 +1551,9 @@ permitted_dependency_tiers = [ ] source_repository = "https://github.com/radrootslabs/studio_app" source_revision = "2b5fe5d8321fd0e248a9d304651b1ba54ee6a180" -source_path = "core/crates/ffi" -source_tree_sha256 = "0bb1d02eb963a606a288d9c35fb418de7bfd914dc5e2c4a38112af64c643151c" +source_path = "core/crates/storage" +source_tree_sha256 = "f3addecbd7f6dbccda4a438597b4e433443a34b48a6c4ef9efa90e49e8f05c4b" compatibility = ["product", "lifecycle", "package_private"] -removal_gate = "studio_runtime_extraction_verified" replaces = [] [[package]] diff --git a/contracts/crates/generated/package_groups.v1.toml b/contracts/crates/generated/package_groups.v1.toml @@ -1,16 +1,16 @@ schema = "radroots.workspace.package-groups.v1" -catalog_sha256 = "990cae5fc037497618bf5324de9b955c08077ef134ce885b46b34b629da93136" +catalog_sha256 = "8c27cebf6825f9ed74e122513c661f6dd31837dddb39d0988014121d8ab05e75" [[group]] id = "boundaries" packages = ["radroots_core_bindings", "radroots_event_bindings", "radroots_event_codec_wasm", "radroots_identity_bindings", "radroots_mobile_bindgen", "radroots_mobile_ffi", "radroots_mobile_wasm", "radroots_replica_schema_bindings", "radroots_replica_store_wasm", "radroots_replica_sync_wasm", "radroots_sdk_ffi", "radroots_studio_ffi", "radroots_studio_uniffi_bindgen", "radroots_trade_bindings"] -active_packages = ["radroots_core_bindings", "radroots_event_bindings", "radroots_event_codec_wasm", "radroots_identity_bindings", "radroots_mobile_bindgen", "radroots_mobile_ffi", "radroots_mobile_wasm", "radroots_replica_schema_bindings", "radroots_replica_store_wasm", "radroots_replica_sync_wasm", "radroots_sdk_ffi", "radroots_trade_bindings"] -reserved_packages = ["radroots_studio_ffi", "radroots_studio_uniffi_bindgen"] +active_packages = ["radroots_core_bindings", "radroots_event_bindings", "radroots_event_codec_wasm", "radroots_identity_bindings", "radroots_mobile_bindgen", "radroots_mobile_ffi", "radroots_mobile_wasm", "radroots_replica_schema_bindings", "radroots_replica_store_wasm", "radroots_replica_sync_wasm", "radroots_sdk_ffi", "radroots_studio_ffi", "radroots_studio_uniffi_bindgen", "radroots_trade_bindings"] +reserved_packages = [] [[group]] id = "coverage_required" -packages = ["radroots", "radroots_blossom", "radroots_core", "radroots_core_bindings", "radroots_event", "radroots_event_bindings", "radroots_event_codec", "radroots_event_codec_wasm", "radroots_geonames", "radroots_identity", "radroots_identity_bindings", "radroots_mesh", "radroots_mesh_agent_client", "radroots_mesh_agent_proto", "radroots_mobile_bindgen", "radroots_mobile_core", "radroots_mobile_ffi", "radroots_mobile_wasm", "radroots_nostr", "radroots_nostr_connect", "radroots_nostrdb", "radroots_protocol", "radroots_replica_schema", "radroots_replica_schema_bindings", "radroots_replica_store", "radroots_replica_store_wasm", "radroots_replica_sync", "radroots_replica_sync_wasm", "radroots_runtime_distribution", "radroots_runtime_manager", "radroots_runtime_paths", "radroots_sdk", "radroots_sdk_ffi", "radroots_sdk_sql_wasm_runtime", "radroots_secrets", "radroots_signing", "radroots_sql_core", "radroots_storage", "radroots_storage_sqlite", "radroots_sync", "radroots_test_fixtures", "radroots_trade", "radroots_trade_bindings", "radroots_transport", "radroots_transport_nostr", "radroots_transport_reticulum", "xtask"] -active_packages = ["radroots", "radroots_blossom", "radroots_core", "radroots_core_bindings", "radroots_event", "radroots_event_bindings", "radroots_event_codec", "radroots_event_codec_wasm", "radroots_geonames", "radroots_identity", "radroots_identity_bindings", "radroots_mesh", "radroots_mesh_agent_client", "radroots_mesh_agent_proto", "radroots_mobile_bindgen", "radroots_mobile_core", "radroots_mobile_ffi", "radroots_mobile_wasm", "radroots_nostr", "radroots_nostr_connect", "radroots_nostrdb", "radroots_protocol", "radroots_replica_schema", "radroots_replica_schema_bindings", "radroots_replica_store", "radroots_replica_store_wasm", "radroots_replica_sync", "radroots_replica_sync_wasm", "radroots_runtime_distribution", "radroots_runtime_manager", "radroots_runtime_paths", "radroots_sdk", "radroots_sdk_ffi", "radroots_sdk_sql_wasm_runtime", "radroots_secrets", "radroots_signing", "radroots_sql_core", "radroots_storage", "radroots_storage_sqlite", "radroots_sync", "radroots_test_fixtures", "radroots_trade", "radroots_trade_bindings", "radroots_transport", "radroots_transport_nostr", "radroots_transport_reticulum", "xtask"] +packages = ["radroots", "radroots_blossom", "radroots_core", "radroots_core_bindings", "radroots_event", "radroots_event_bindings", "radroots_event_codec", "radroots_event_codec_wasm", "radroots_geonames", "radroots_identity", "radroots_identity_bindings", "radroots_mesh", "radroots_mesh_agent_client", "radroots_mesh_agent_proto", "radroots_mobile_bindgen", "radroots_mobile_core", "radroots_mobile_ffi", "radroots_mobile_wasm", "radroots_nostr", "radroots_nostr_connect", "radroots_nostrdb", "radroots_protocol", "radroots_replica_schema", "radroots_replica_schema_bindings", "radroots_replica_store", "radroots_replica_store_wasm", "radroots_replica_sync", "radroots_replica_sync_wasm", "radroots_runtime_distribution", "radroots_runtime_manager", "radroots_runtime_paths", "radroots_sdk", "radroots_sdk_ffi", "radroots_sdk_sql_wasm_runtime", "radroots_secrets", "radroots_signing", "radroots_sql_core", "radroots_storage", "radroots_storage_sqlite", "radroots_studio_application", "radroots_studio_domain", "radroots_studio_ffi", "radroots_studio_nostr", "radroots_studio_preferences", "radroots_studio_runtime", "radroots_studio_storage", "radroots_studio_uniffi_bindgen", "radroots_sync", "radroots_test_fixtures", "radroots_trade", "radroots_trade_bindings", "radroots_transport", "radroots_transport_nostr", "radroots_transport_reticulum", "xtask"] +active_packages = ["radroots", "radroots_blossom", "radroots_core", "radroots_core_bindings", "radroots_event", "radroots_event_bindings", "radroots_event_codec", "radroots_event_codec_wasm", "radroots_geonames", "radroots_identity", "radroots_identity_bindings", "radroots_mesh", "radroots_mesh_agent_client", "radroots_mesh_agent_proto", "radroots_mobile_bindgen", "radroots_mobile_core", "radroots_mobile_ffi", "radroots_mobile_wasm", "radroots_nostr", "radroots_nostr_connect", "radroots_nostrdb", "radroots_protocol", "radroots_replica_schema", "radroots_replica_schema_bindings", "radroots_replica_store", "radroots_replica_store_wasm", "radroots_replica_sync", "radroots_replica_sync_wasm", "radroots_runtime_distribution", "radroots_runtime_manager", "radroots_runtime_paths", "radroots_sdk", "radroots_sdk_ffi", "radroots_sdk_sql_wasm_runtime", "radroots_secrets", "radroots_signing", "radroots_sql_core", "radroots_storage", "radroots_storage_sqlite", "radroots_studio_application", "radroots_studio_domain", "radroots_studio_ffi", "radroots_studio_nostr", "radroots_studio_preferences", "radroots_studio_runtime", "radroots_studio_storage", "radroots_studio_uniffi_bindgen", "radroots_sync", "radroots_test_fixtures", "radroots_trade", "radroots_trade_bindings", "radroots_transport", "radroots_transport_nostr", "radroots_transport_reticulum", "xtask"] reserved_packages = [] [[group]] @@ -45,9 +45,9 @@ reserved_packages = [] [[group]] id = "studio" -packages = ["radroots_studio_application", "radroots_studio_domain", "radroots_studio_ffi", "radroots_studio_nostr", "radroots_studio_runtime", "radroots_studio_storage", "radroots_studio_uniffi_bindgen"] -active_packages = [] -reserved_packages = ["radroots_studio_application", "radroots_studio_domain", "radroots_studio_ffi", "radroots_studio_nostr", "radroots_studio_runtime", "radroots_studio_storage", "radroots_studio_uniffi_bindgen"] +packages = ["radroots_studio_application", "radroots_studio_domain", "radroots_studio_ffi", "radroots_studio_nostr", "radroots_studio_preferences", "radroots_studio_runtime", "radroots_studio_storage", "radroots_studio_uniffi_bindgen"] +active_packages = ["radroots_studio_application", "radroots_studio_domain", "radroots_studio_ffi", "radroots_studio_nostr", "radroots_studio_preferences", "radroots_studio_runtime", "radroots_studio_storage", "radroots_studio_uniffi_bindgen"] +reserved_packages = [] [[group]] id = "tools" diff --git a/contracts/crates/generated/platform_inventory.v1.toml b/contracts/crates/generated/platform_inventory.v1.toml @@ -1,5 +1,5 @@ schema = "radroots.workspace.platform-inventory.v1" -catalog_sha256 = "990cae5fc037497618bf5324de9b955c08077ef134ce885b46b34b629da93136" +catalog_sha256 = "8c27cebf6825f9ed74e122513c661f6dd31837dddb39d0988014121d8ab05e75" [[platform]] id = "android" @@ -23,7 +23,7 @@ packages = ["radroots_studio_ffi"] [[platform]] id = "native" -packages = ["radroots_core_bindings", "radroots_event_bindings", "radroots_geonames", "radroots_identity_bindings", "radroots_mesh", "radroots_mesh_agent_client", "radroots_mobile_bindgen", "radroots_mobile_core", "radroots_nostrdb", "radroots_replica_schema_bindings", "radroots_replica_store", "radroots_replica_sync", "radroots_runtime_distribution", "radroots_runtime_manager", "radroots_runtime_paths", "radroots_sdk_ffi", "radroots_secrets", "radroots_simplex_app_store", "radroots_simplex_smp_transport", "radroots_sql_core", "radroots_storage", "radroots_storage_sqlite", "radroots_studio_application", "radroots_studio_domain", "radroots_studio_nostr", "radroots_studio_runtime", "radroots_studio_storage", "radroots_studio_uniffi_bindgen", "radroots_sync", "radroots_trade_bindings", "radroots_transport_nostr", "radroots_transport_reticulum", "xtask"] +packages = ["radroots_core_bindings", "radroots_event_bindings", "radroots_geonames", "radroots_identity_bindings", "radroots_mesh", "radroots_mesh_agent_client", "radroots_mobile_bindgen", "radroots_mobile_core", "radroots_nostrdb", "radroots_replica_schema_bindings", "radroots_replica_store", "radroots_replica_sync", "radroots_runtime_distribution", "radroots_runtime_manager", "radroots_runtime_paths", "radroots_sdk_ffi", "radroots_secrets", "radroots_simplex_app_store", "radroots_simplex_smp_transport", "radroots_sql_core", "radroots_storage", "radroots_storage_sqlite", "radroots_studio_application", "radroots_studio_domain", "radroots_studio_nostr", "radroots_studio_preferences", "radroots_studio_runtime", "radroots_studio_storage", "radroots_studio_uniffi_bindgen", "radroots_sync", "radroots_trade_bindings", "radroots_transport_nostr", "radroots_transport_reticulum", "xtask"] [[platform]] id = "wasm32" diff --git a/contracts/crates/generated/release_inventory.v2.toml b/contracts/crates/generated/release_inventory.v2.toml @@ -1,7 +1,7 @@ schema = "radroots.workspace.release-inventory.v2" -catalog_sha256 = "990cae5fc037497618bf5324de9b955c08077ef134ce885b46b34b629da93136" +catalog_sha256 = "8c27cebf6825f9ed74e122513c661f6dd31837dddb39d0988014121d8ab05e75" architecture = "radroots.crates.release.v2" version = "0.1.0-alpha" public_packages = ["radroots", "radroots_blossom", "radroots_core", "radroots_event", "radroots_event_codec", "radroots_geonames", "radroots_identity", "radroots_nostr", "radroots_nostr_connect", "radroots_protocol", "radroots_sdk", "radroots_secrets", "radroots_signing", "radroots_storage", "radroots_storage_sqlite", "radroots_sync", "radroots_trade", "radroots_transport", "radroots_transport_nostr"] -private_packages = ["radroots_core_bindings", "radroots_event_bindings", "radroots_event_codec_wasm", "radroots_identity_bindings", "radroots_mesh", "radroots_mesh_agent_client", "radroots_mesh_agent_proto", "radroots_mobile_bindgen", "radroots_mobile_core", "radroots_mobile_ffi", "radroots_mobile_wasm", "radroots_nostrdb", "radroots_replica_schema", "radroots_replica_schema_bindings", "radroots_replica_store", "radroots_replica_store_wasm", "radroots_replica_sync", "radroots_replica_sync_wasm", "radroots_runtime_distribution", "radroots_runtime_manager", "radroots_runtime_paths", "radroots_sdk_ffi", "radroots_sdk_sql_wasm_runtime", "radroots_simplex_agent_proto", "radroots_simplex_app_store", "radroots_simplex_chat_proto", "radroots_simplex_smp_crypto", "radroots_simplex_smp_proto", "radroots_simplex_smp_transport", "radroots_sql_core", "radroots_studio_application", "radroots_studio_domain", "radroots_studio_ffi", "radroots_studio_nostr", "radroots_studio_runtime", "radroots_studio_storage", "radroots_studio_uniffi_bindgen", "radroots_test_fixtures", "radroots_trade_bindings", "radroots_transport_reticulum", "xtask"] -reserved_packages = ["radroots_studio_application", "radroots_studio_domain", "radroots_studio_ffi", "radroots_studio_nostr", "radroots_studio_runtime", "radroots_studio_storage", "radroots_studio_uniffi_bindgen"] +private_packages = ["radroots_core_bindings", "radroots_event_bindings", "radroots_event_codec_wasm", "radroots_identity_bindings", "radroots_mesh", "radroots_mesh_agent_client", "radroots_mesh_agent_proto", "radroots_mobile_bindgen", "radroots_mobile_core", "radroots_mobile_ffi", "radroots_mobile_wasm", "radroots_nostrdb", "radroots_replica_schema", "radroots_replica_schema_bindings", "radroots_replica_store", "radroots_replica_store_wasm", "radroots_replica_sync", "radroots_replica_sync_wasm", "radroots_runtime_distribution", "radroots_runtime_manager", "radroots_runtime_paths", "radroots_sdk_ffi", "radroots_sdk_sql_wasm_runtime", "radroots_simplex_agent_proto", "radroots_simplex_app_store", "radroots_simplex_chat_proto", "radroots_simplex_smp_crypto", "radroots_simplex_smp_proto", "radroots_simplex_smp_transport", "radroots_sql_core", "radroots_studio_application", "radroots_studio_domain", "radroots_studio_ffi", "radroots_studio_nostr", "radroots_studio_preferences", "radroots_studio_runtime", "radroots_studio_storage", "radroots_studio_uniffi_bindgen", "radroots_test_fixtures", "radroots_trade_bindings", "radroots_transport_reticulum", "xtask"] +reserved_packages = [] diff --git a/contracts/releases/publish_policy.toml b/contracts/releases/publish_policy.toml @@ -61,6 +61,8 @@ private = [ "radroots_studio_application", "radroots_studio_domain", "radroots_studio_nostr", + "radroots_studio_preferences", + "radroots_studio_runtime", "radroots_studio_storage", ] build_codegen = [ diff --git a/crates/studio_application/Cargo.toml b/crates/studio_application/Cargo.toml @@ -12,15 +12,9 @@ publish = false include = ["src/**", "tests/**", "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" } radroots_studio_domain.workspace = true -radroots_studio_nostr.workspace = true secrecy = "=0.10.3" tokio = { version = "=1.47.1", features = ["macros", "rt-multi-thread", "sync", "time"] } -[dev-dependencies] -nostr-relay-builder = { git = "https://github.com/rust-nostr/nostr.git", rev = "5bba5163eb77107f82c4a8262cf29d7f33a73219", package = "nostr-relay-builder" } - [lints] workspace = true diff --git a/crates/studio_application/src/accounts.rs b/crates/studio_application/src/accounts.rs @@ -1,11 +1,5 @@ use std::sync::{Mutex, MutexGuard}; -use radroots_studio_domain::{ - AccountCreatedAt, AccountIdentity, AccountSummary, BindingAvailability, LocalSignerBinding, - Nsec, PublicKey, SafeError, SafeErrorCode, SafeMessage, SecretKeyInput, -}; -use radroots_studio_nostr::{generate_local_keypair, import_secret}; - use crate::{ AccountOperationKind, AccountOperationPhase, AccountRepository, AppCore, AppStateRepository, Clock, DurableOperationKind, DurableOperationPhase, DurableOperationRepository, @@ -13,6 +7,10 @@ use crate::{ OperationId, OperationJournal, OperationPriorState, PendingAccountOperation, RemovalConfirmationToken, SecretStore, StagedGeneratedKey, StateTransition, }; +use radroots_studio_domain::{ + AccountCreatedAt, AccountIdentity, AccountSummary, BindingAvailability, LocalSignerBinding, + Nsec, PublicKey, SafeError, SafeErrorCode, SafeMessage, SecretKeyInput, +}; pub struct GenerateAccountReceipt { account: AccountSummary, @@ -97,7 +95,7 @@ impl AppCore { clock: &(impl Clock + ?Sized), ) -> Result<GenerateAccountReceipt, SafeError> { self.require_revision(expected_revision)?; - let generated = generate_local_keypair()?; + let generated = self.key_material().generate()?; let (public_key, npub, secret, nsec) = generated.into_parts(); let account = AccountSummary::new( AccountIdentity::verify(public_key, npub.as_str().to_owned())?, @@ -156,7 +154,7 @@ impl AppCore { }; } self.require_revision(expected_revision)?; - let imported = import_secret(input)?; + let imported = self.key_material().import(input)?; let (public_key, npub, secret) = imported.into_parts(); let previous = accounts.find_account(public_key)?; if let Some(existing) = &previous @@ -508,7 +506,7 @@ impl AppCore { journal: &(impl OperationJournal + ?Sized), clock: &(impl Clock + ?Sized), ) -> Result<GenerateAccountReceipt, SafeError> { - let generated = generate_local_keypair()?; + let generated = self.key_material().generate()?; let (public_key, npub, secret, nsec) = generated.into_parts(); let account = AccountSummary::new( AccountIdentity::verify(public_key, npub.as_str().to_owned())?, @@ -553,7 +551,7 @@ impl AppCore { journal: &(impl OperationJournal + ?Sized), clock: &(impl Clock + ?Sized), ) -> Result<ImportAccountReceipt, SafeError> { - let imported = import_secret(input)?; + let imported = self.key_material().import(input)?; let (public_key, npub, secret) = imported.into_parts(); if let Some(existing) = accounts.find_account(public_key)? { if existing.signer().availability() != BindingAvailability::CredentialMissing @@ -1160,7 +1158,7 @@ mod tests { ) .expect("input") }; - let imported = radroots_studio_nostr::import_secret(input()).expect("derive"); + let imported = core.key_material().import(input()).expect("derive"); let (public_key, npub, _) = imported.into_parts(); let missing = AccountSummary::new( AccountIdentity::verify(public_key, npub.as_str().to_owned()).expect("identity"), diff --git a/crates/studio_application/src/app_core.rs b/crates/studio_application/src/app_core.rs @@ -1,11 +1,11 @@ use std::collections::BTreeMap; -use std::sync::{Mutex, MutexGuard}; +use std::sync::{Arc, Mutex, MutexGuard}; use radroots_studio_domain::{PublicKey, SafeError, SafeErrorCode, SafeMessage, UnixTimestamp}; use crate::{ - AccountRepository, AppSnapshot, AppStateRepository, RelayConfiguration, SnapshotRevision, - StateMachine, StateTransition, + AccountRepository, AppSnapshot, AppStateRepository, KeyMaterialProvider, RelayConfiguration, + SnapshotRevision, StateMachine, StateTransition, }; pub struct RemovalConfirmationToken { @@ -68,14 +68,19 @@ struct CoreState { pub struct AppCore { relay_configuration: RelayConfiguration, + key_material: Arc<dyn KeyMaterialProvider>, state: Mutex<CoreState>, } impl AppCore { #[must_use] - pub fn in_memory(relay_configuration: RelayConfiguration) -> Self { + pub fn new( + relay_configuration: RelayConfiguration, + key_material: Arc<dyn KeyMaterialProvider>, + ) -> Self { Self { relay_configuration, + key_material, state: Mutex::new(CoreState { state_machine: StateMachine::booting(), removal_tokens: BTreeMap::new(), @@ -84,6 +89,19 @@ impl AppCore { } } + #[cfg(test)] + #[must_use] + pub fn in_memory(relay_configuration: RelayConfiguration) -> Self { + Self::new( + relay_configuration, + Arc::new(crate::test_support::TestKeyMaterialProvider::default()), + ) + } + + pub(crate) fn key_material(&self) -> &dyn KeyMaterialProvider { + self.key_material.as_ref() + } + /// Moves the in-memory core from booting to an empty ready snapshot. /// /// # Errors diff --git a/crates/studio_application/src/config.rs b/crates/studio_application/src/config.rs @@ -1,4 +1,6 @@ -use radroots_studio_domain::{SafeError, SafeErrorCode, SafeMessage, normalize_relay_urls}; +use radroots_studio_domain::{ + RelayDestinationPolicy, SafeError, SafeErrorCode, SafeMessage, normalize_relay_urls, +}; use crate::RelayConfiguration; @@ -38,19 +40,19 @@ pub fn relay_configuration_from_value( mode: RelayRuntimeMode, ) -> Result<RelayConfiguration, SafeError> { let configured = value.unwrap_or_default().trim(); - let source = if configured.is_empty() { + let (source, policy) = if configured.is_empty() { match mode { - RelayRuntimeMode::Development => DEVELOPMENT_RELAY, + RelayRuntimeMode::Development => (DEVELOPMENT_RELAY, RelayDestinationPolicy::Local), RelayRuntimeMode::Packaged => return Err(invalid_configuration()), } } else { - configured + (configured, RelayDestinationPolicy::Public) }; - let normalized = normalize_relay_urls(source.split(',').map(str::trim))?; + let normalized = normalize_relay_urls(source.split(',').map(str::trim), policy)?; if normalized.is_empty() { return Err(invalid_configuration()); } - Ok(RelayConfiguration::new(normalized)) + RelayConfiguration::new(normalized) } const fn invalid_configuration() -> SafeError { diff --git a/crates/studio_application/src/custody.rs b/crates/studio_application/src/custody.rs @@ -2,12 +2,11 @@ use std::num::NonZeroU64; use std::sync::Mutex; use std::time::Duration; +use crate::KeyMaterialProvider; use radroots_studio_domain::{ AccountCreatedAt, AccountIdentity, AccountSummary, BindingAvailability, LocalSignerBinding, Nsec, SafeError, SafeErrorCode, SafeMessage, SecretKeyInput, UnixTimestamp, }; -use radroots_studio_nostr::generate_local_keypair; - pub const GENERATED_KEY_STAGE_TTL: Duration = Duration::from_mins(5); #[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)] @@ -133,6 +132,7 @@ impl GeneratedKeyStage { /// Returns a safe conflict while an unexpired recovery stage is active. pub fn begin( &mut self, + key_material: &dyn KeyMaterialProvider, id: RecoveryStageId, expected_revision: u64, now: UnixTimestamp, @@ -141,7 +141,7 @@ impl GeneratedKeyStage { if self.pending.is_some() { return Err(recovery_in_progress()); } - let generated = generate_local_keypair()?; + let generated = key_material.generate()?; let (public_key, npub, secret, recovery_nsec) = generated.into_parts(); let account = AccountSummary::new( AccountIdentity::verify(public_key, npub.as_str().to_owned())?, @@ -237,6 +237,7 @@ mod tests { use radroots_studio_domain::UnixTimestamp; use super::{GENERATED_KEY_STAGE_TTL, GeneratedKeyStage, RecoveryStageId}; + use crate::test_support::TestKeyMaterialProvider; fn time(seconds: i64) -> UnixTimestamp { UnixTimestamp::from_seconds(seconds).expect("time") @@ -249,11 +250,14 @@ mod tests { #[test] fn stage_is_exclusive_cancelable_and_never_publishes_secret_debug() { let mut stage = GeneratedKeyStage::default(); - let handle = stage.begin(id(1), 4, time(10)).expect("begin"); + let key_material = TestKeyMaterialProvider::default(); + let handle = stage + .begin(&key_material, id(1), 4, time(10)) + .expect("begin"); let view = handle.view(); assert_eq!(view.expires_at().as_seconds(), 310); assert_eq!(stage.pending().expect("pending").expected_revision(), 4); - assert!(stage.begin(id(2), 4, time(11)).is_err()); + assert!(stage.begin(&key_material, id(2), 4, time(11)).is_err()); let nsec = handle.take_recovery_nsec().expect("one-use recovery"); assert_eq!(nsec.with_exposed_secret(str::len), 63); assert!(handle.take_recovery_nsec().is_err()); @@ -266,14 +270,19 @@ mod tests { #[test] fn stage_expires_and_is_destroyed_on_owner_drop() { let mut stage = GeneratedKeyStage::default(); - stage.begin(id(1), 0, time(20)).expect("begin"); + let key_material = TestKeyMaterialProvider::default(); + stage + .begin(&key_material, id(1), 0, time(20)) + .expect("begin"); let expiry = 20 + i64::try_from(GENERATED_KEY_STAGE_TTL.as_secs()).expect("ttl"); assert!(stage.expire(time(expiry))); assert!(stage.pending().is_none()); assert!(stage.take(id(1), time(expiry)).is_err()); let mut shutdown_stage = GeneratedKeyStage::default(); - shutdown_stage.begin(id(2), 0, time(30)).expect("begin"); + shutdown_stage + .begin(&key_material, id(2), 0, time(30)) + .expect("begin"); drop(shutdown_stage); } } diff --git a/crates/studio_application/src/lib.rs b/crates/studio_application/src/lib.rs @@ -6,7 +6,6 @@ pub mod app_core; mod change_stream; pub mod config; pub mod custody; -pub mod nostr_client; pub mod ports; mod profile_refresh; pub mod recovery; @@ -15,6 +14,9 @@ pub mod session; pub mod snapshot; pub mod state_machine; +#[cfg(test)] +mod test_support; + pub use accounts::{ GenerateAccountReceipt, ImportAccountReceipt, InMemoryAccountRepository, InMemoryOperationJournal, @@ -35,21 +37,22 @@ pub use custody::{ GENERATED_KEY_STAGE_TTL, GeneratedKeyRecoveryHandle, GeneratedKeyStage, GeneratedKeyStageView, RecoveryStageId, StagedGeneratedKey, }; -pub use nostr_client::SdkNostrClient; pub use ports::{ AccountNamespaceRepository, AccountOperationKind, AccountOperationPhase, AccountPreferenceKey, AccountRepository, AppStateRepository, BoxFuture, CachedProfile, Clock, DurableAccountOperation, DurableOperationKind, DurableOperationPhase, DurableOperationReceipt, DurableOperationRepository, DurableOperationStart, DurableRequestId, DurableTerminalOutcome, - NostrClient, OperationDiagnostic, OperationId, OperationJournal, OperationPriorState, - PendingAccountOperation, ProfileRefreshStatus, ProfileRepository, + GeneratedKeyMaterial, ImportedKeyMaterial, KeyMaterialProvider, NostrClient, + OperationDiagnostic, OperationId, OperationJournal, OperationPriorState, + PendingAccountOperation, ProfileFetchResult, ProfileRefreshStatus, ProfileRepository, + RelayFetchCompleteness, }; pub use profile_refresh::ProfileRefreshPlan; pub use secrets::{ FailureSecretStore, InMemorySecretStore, SecretStore, SecretStoreCall, SecretStoreOperation, }; pub use snapshot::{ - ActiveAccountSnapshot, AppLifecycle, AppSnapshot, ProfileLoadState, RelayConfiguration, - RelayConnectionState, SessionState, SnapshotRevision, + ActiveAccountSnapshot, AppLifecycle, AppSnapshot, MAX_CONFIGURED_RELAYS, ProfileLoadState, + RelayConfiguration, RelayConnectionState, SessionState, SnapshotRevision, }; pub use state_machine::{StateMachine, StateTransition}; diff --git a/crates/studio_application/src/nostr_client.rs b/crates/studio_application/src/nostr_client.rs @@ -1,138 +0,0 @@ -use std::time::Duration; - -use nostr::{Filter, JsonUtil, Kind, PublicKey as NostrPublicKey}; -use nostr_sdk::ClientBuilder; -use radroots_studio_domain::{ - Kind0ProfileCandidate, PublicKey, RelayUrl, SafeError, SafeErrorCode, SafeMessage, - select_latest_kind0, -}; - -use crate::{BoxFuture, NostrClient}; - -pub struct SdkNostrClient { - timeout: Duration, -} - -impl SdkNostrClient { - #[must_use] - pub const fn new(timeout: Duration) -> Self { - Self { timeout } - } -} - -impl NostrClient for SdkNostrClient { - fn fetch_profile<'a>( - &'a self, - public_key: PublicKey, - relays: &'a [RelayUrl], - ) -> BoxFuture<'a, Result<Option<Kind0ProfileCandidate>, SafeError>> { - Box::pin(async move { - if relays.is_empty() { - return Err(invalid_relay_configuration()); - } - - let client = ClientBuilder::new().build(); - for relay in relays { - client - .add_relay(relay.as_str()) - .await - .map_err(|_| relay_connection_failed())?; - } - client.connect().await; - client.wait_for_connection(self.timeout).await; - - let author = NostrPublicKey::from_slice(public_key.as_bytes()) - .map_err(|_| profile_refresh_failed())?; - let filter = Filter::new().author(author).kind(Kind::Metadata).limit(64); - let fetched = client.fetch_events(filter, self.timeout).await; - client.shutdown().await; - let events = fetched.map_err(|_| relay_connection_failed())?; - - let mut candidates = Vec::with_capacity(events.len()); - for event in events.iter() { - candidates.push(radroots_studio_nostr::parse_verified_kind0( - &event.as_json(), - public_key, - )?); - } - Ok(select_latest_kind0(candidates)) - }) - } -} - -const fn invalid_relay_configuration() -> SafeError { - SafeError::new( - SafeErrorCode::InvalidRelayConfiguration, - SafeMessage::new("No Nostr relay is configured."), - ) -} - -const fn relay_connection_failed() -> SafeError { - SafeError::new( - SafeErrorCode::RelayConnectionFailed, - SafeMessage::new("The Nostr relays could not be reached."), - ) -} - -const fn profile_refresh_failed() -> SafeError { - SafeError::new( - SafeErrorCode::ProfileRefreshFailed, - SafeMessage::new("The Nostr profile could not be refreshed."), - ) -} - -#[cfg(test)] -mod tests { - use std::time::Duration; - - use nostr::{EventBuilder, Keys, Metadata}; - use nostr_relay_builder::MockRelay; - use nostr_sdk::Client; - use radroots_studio_domain::{PublicKey, RelayUrl, SafeErrorCode}; - - use crate::{NostrClient, SdkNostrClient}; - - #[tokio::test] - async fn sdk_client_fetches_verified_profile_from_ephemeral_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.connect().await; - publisher.wait_for_connection(Duration::from_secs(2)).await; - publisher - .send_event_builder(EventBuilder::metadata( - &Metadata::new().name("Farmer").display_name("Farm Account"), - )) - .await - .expect("publish metadata"); - - let adapter = SdkNostrClient::new(Duration::from_secs(2)); - let domain_relay = RelayUrl::parse(relay_url.as_str()).expect("domain relay URL"); - let public_key = PublicKey::from_bytes(keys.public_key().to_bytes()); - let profile = adapter - .fetch_profile(public_key, &[domain_relay]) - .await - .expect("fetch profile") - .expect("published profile"); - - assert_eq!(profile.author(), public_key); - assert_eq!(profile.metadata().preferred_name(), Some("Farm Account")); - 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)) - .fetch_profile(PublicKey::from_bytes([1; 32]), &[]) - .await - .expect_err("empty relay list"); - - assert_eq!(error.code(), SafeErrorCode::InvalidRelayConfiguration); - } -} diff --git a/crates/studio_application/src/ports.rs b/crates/studio_application/src/ports.rs @@ -1,9 +1,10 @@ use std::future::Future; use std::pin::Pin; +use std::time::Instant; use radroots_studio_domain::{ - AccountSummary, BindingAvailability, Kind0ProfileCandidate, PublicKey, RelayUrl, SafeError, - SafeErrorCode, SafeMessage, UnixTimestamp, + AccountSummary, BindingAvailability, Kind0ProfileCandidate, Npub, Nsec, PublicKey, RelayUrl, + SafeError, SafeErrorCode, SafeMessage, SecretKeyInput, UnixTimestamp, }; const MAX_DURABLE_REQUEST_ID_BYTES: usize = 128; @@ -236,6 +237,41 @@ pub enum ProfileRefreshStatus { InvalidData, } +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum RelayFetchCompleteness { + Complete, + Partial, +} + +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct ProfileFetchResult { + candidate: Option<Kind0ProfileCandidate>, + completeness: RelayFetchCompleteness, +} + +impl ProfileFetchResult { + #[must_use] + pub const fn complete(candidate: Option<Kind0ProfileCandidate>) -> Self { + Self { + candidate, + completeness: RelayFetchCompleteness::Complete, + } + } + + #[must_use] + pub const fn partial(candidate: Option<Kind0ProfileCandidate>) -> Self { + Self { + candidate, + completeness: RelayFetchCompleteness::Partial, + } + } + + #[must_use] + pub fn into_parts(self) -> (Option<Kind0ProfileCandidate>, RelayFetchCompleteness) { + (self.candidate, self.completeness) + } +} + #[derive(Clone, Debug, Eq, PartialEq)] pub struct CachedProfile { candidate: Kind0ProfileCandidate, @@ -590,7 +626,75 @@ pub trait NostrClient: Send + Sync { &'a self, public_key: PublicKey, relays: &'a [RelayUrl], - ) -> BoxFuture<'a, Result<Option<Kind0ProfileCandidate>, SafeError>>; + deadline: Instant, + ) -> BoxFuture<'a, Result<ProfileFetchResult, SafeError>>; +} + +pub struct GeneratedKeyMaterial { + public_key: PublicKey, + npub: Npub, + secret: SecretKeyInput, + nsec: Nsec, +} + +impl GeneratedKeyMaterial { + #[must_use] + pub const fn new( + public_key: PublicKey, + npub: Npub, + secret: SecretKeyInput, + nsec: Nsec, + ) -> Self { + Self { + public_key, + npub, + secret, + nsec, + } + } + + #[must_use] + pub fn into_parts(self) -> (PublicKey, Npub, SecretKeyInput, Nsec) { + (self.public_key, self.npub, self.secret, self.nsec) + } +} + +pub struct ImportedKeyMaterial { + public_key: PublicKey, + npub: Npub, + secret: SecretKeyInput, +} + +impl ImportedKeyMaterial { + #[must_use] + pub const fn new(public_key: PublicKey, npub: Npub, secret: SecretKeyInput) -> Self { + Self { + public_key, + npub, + secret, + } + } + + #[must_use] + pub fn into_parts(self) -> (PublicKey, Npub, SecretKeyInput) { + (self.public_key, self.npub, self.secret) + } +} + +pub trait KeyMaterialProvider: Send + Sync { + /// Generates one keypair from host-provided cryptographic entropy. + /// + /// # Errors + /// + /// Returns a redacted key or entropy error. + fn generate(&self) -> Result<GeneratedKeyMaterial, SafeError>; + + /// Canonicalizes imported secret material and derives its public identity. + /// + /// # Errors + /// + /// Returns a redacted validation error. + fn import(&self, input: SecretKeyInput) -> Result<ImportedKeyMaterial, SafeError>; } pub trait Clock: Send + Sync { @@ -599,18 +703,18 @@ pub trait Clock: Send + Sync { #[cfg(test)] mod tests { + use std::time::Instant; + use std::sync::Mutex; - use radroots_studio_domain::{ - AccountSummary, Kind0ProfileCandidate, PublicKey, RelayUrl, SafeError, UnixTimestamp, - }; + use radroots_studio_domain::{AccountSummary, PublicKey, RelayUrl, SafeError, UnixTimestamp}; use super::{ AccountNamespaceRepository, AccountOperationKind, AccountOperationPhase, AccountPreferenceKey, AccountRepository, AppStateRepository, BoxFuture, CachedProfile, Clock, DurableOperationReceipt, DurableRequestId, DurableTerminalOutcome, NostrClient, OperationDiagnostic, OperationId, OperationJournal, PendingAccountOperation, - ProfileRefreshStatus, ProfileRepository, + ProfileFetchResult, ProfileRefreshStatus, ProfileRepository, }; #[test] @@ -750,8 +854,9 @@ mod tests { &'a self, _public_key: PublicKey, _relays: &'a [RelayUrl], - ) -> BoxFuture<'a, Result<Option<Kind0ProfileCandidate>, SafeError>> { - Box::pin(async { Ok(None) }) + _deadline: Instant, + ) -> BoxFuture<'a, Result<ProfileFetchResult, SafeError>> { + Box::pin(async { Ok(ProfileFetchResult::complete(None)) }) } } diff --git a/crates/studio_application/src/profile_refresh.rs b/crates/studio_application/src/profile_refresh.rs @@ -1,9 +1,10 @@ use radroots_studio_domain::{PublicKey, RelayUrl, SafeError, SafeErrorCode}; +use std::time::Instant; use crate::{ ActiveAccountSnapshot, AppCore, AppSnapshot, CachedProfile, Clock, NostrClient, - ProfileLoadState, ProfileRefreshStatus, ProfileRepository, RelayConnectionState, - SnapshotRevision, StateTransition, + ProfileFetchResult, ProfileLoadState, ProfileRefreshStatus, ProfileRepository, + RelayConnectionState, RelayFetchCompleteness, SnapshotRevision, StateTransition, }; #[derive(Clone, Debug, Eq, PartialEq)] @@ -51,8 +52,9 @@ impl AppCore { profiles: &(impl ProfileRepository + ?Sized), client: &(impl NostrClient + ?Sized), clock: &(impl Clock + ?Sized), + deadline: Instant, ) -> Result<AppSnapshot, SafeError> { - self.refresh_profile_for_active_account(profiles, client, clock) + self.refresh_profile_for_active_account(profiles, client, clock, deadline) .await } @@ -69,11 +71,14 @@ impl AppCore { profiles: &(impl ProfileRepository + ?Sized), client: &(impl NostrClient + ?Sized), clock: &(impl Clock + ?Sized), + deadline: Instant, ) -> Result<AppSnapshot, SafeError> { let Some(plan) = self.begin_profile_refresh()? else { return Ok(self.snapshot()); }; - let result = client.fetch_profile(plan.public_key(), plan.relays()).await; + let result = client + .fetch_profile(plan.public_key(), plan.relays(), deadline) + .await; self.complete_profile_refresh(&plan, result, profiles, clock) } @@ -114,7 +119,7 @@ impl AppCore { pub fn complete_profile_refresh( &self, plan: &ProfileRefreshPlan, - result: Result<Option<radroots_studio_domain::Kind0ProfileCandidate>, SafeError>, + result: Result<ProfileFetchResult, SafeError>, profiles: &(impl ProfileRepository + ?Sized), clock: &(impl Clock + ?Sized), ) -> Result<AppSnapshot, SafeError> { @@ -129,7 +134,53 @@ impl AppCore { .ok_or_else(invalid_profile_completion)?; match result { - Ok(Some(candidate)) => { + Ok(fetched) => { + let (candidate, completeness) = fetched.into_parts(); + self.complete_successful_profile_fetch( + plan, + current_active, + candidate, + completeness, + profiles, + clock, + ) + } + Err(error) => { + let status = refresh_status(error); + profiles.record_refresh_status(plan.public_key(), clock.now(), status)?; + self.apply_transition(StateTransition::UpdateActiveAccount { + expected: plan.public_key(), + active_account: Box::new(ActiveAccountSnapshot::new( + current_active.account().clone(), + RelayConnectionState::Degraded, + ProfileLoadState::Error(error), + current_active.profile().cloned(), + )), + problem: Some(error), + }) + } + } + } + + fn complete_successful_profile_fetch( + &self, + plan: &ProfileRefreshPlan, + current_active: ActiveAccountSnapshot, + candidate: Option<radroots_studio_domain::Kind0ProfileCandidate>, + completeness: RelayFetchCompleteness, + profiles: &(impl ProfileRepository + ?Sized), + clock: &(impl Clock + ?Sized), + ) -> Result<AppSnapshot, SafeError> { + let relay_state = match completeness { + RelayFetchCompleteness::Complete => RelayConnectionState::Connected, + RelayFetchCompleteness::Partial => RelayConnectionState::Degraded, + }; + let problem = match completeness { + RelayFetchCompleteness::Complete => None, + RelayFetchCompleteness::Partial => Some(partial_relay_result()), + }; + match candidate { + Some(candidate) => { let cached = CachedProfile::new( candidate.clone(), clock.now(), @@ -144,18 +195,18 @@ impl AppCore { expected: plan.public_key(), active_account: Box::new(ActiveAccountSnapshot::new( current_active.account().clone(), - RelayConnectionState::Connected, + relay_state, ProfileLoadState::Fresh, Some(winning_profile), )), - problem: None, + problem, }) } - Ok(None) => self.apply_transition(StateTransition::UpdateActiveAccount { + None => self.apply_transition(StateTransition::UpdateActiveAccount { expected: plan.public_key(), active_account: Box::new(ActiveAccountSnapshot::new( current_active.account().clone(), - RelayConnectionState::Connected, + relay_state, if current_active.profile().is_some() { ProfileLoadState::Cached } else { @@ -163,26 +214,19 @@ impl AppCore { }, current_active.profile().cloned(), )), - problem: None, + problem, }), - Err(error) => { - let status = refresh_status(error); - profiles.record_refresh_status(plan.public_key(), clock.now(), status)?; - self.apply_transition(StateTransition::UpdateActiveAccount { - expected: plan.public_key(), - active_account: Box::new(ActiveAccountSnapshot::new( - current_active.account().clone(), - RelayConnectionState::Degraded, - ProfileLoadState::Error(error), - current_active.profile().cloned(), - )), - problem: Some(error), - }) - } } } } +const fn partial_relay_result() -> SafeError { + SafeError::new( + SafeErrorCode::RelayConnectionFailed, + radroots_studio_domain::SafeMessage::new("One or more Nostr relays did not complete."), + ) +} + const fn invalid_profile_completion() -> SafeError { SafeError::new( SafeErrorCode::InvalidApplicationState, @@ -208,16 +252,19 @@ const fn refresh_status(error: SafeError) -> ProfileRefreshStatus { #[cfg(test)] mod tests { use std::sync::Mutex; + use std::time::{Duration, Instant}; use radroots_studio_domain::{ - EventId, Kind0ProfileCandidate, ProfileMetadata, PublicKey, RelayUrl, SafeError, - SafeErrorCode, SafeMessage, SecretKeyInput, UnixTimestamp, select_latest_kind0, + EventId, Kind0ProfileCandidate, ProfileMetadata, PublicKey, RelayDestinationPolicy, + RelayUrl, SafeError, SafeErrorCode, SafeMessage, SecretKeyInput, UnixTimestamp, + select_latest_kind0, }; use crate::{ ActiveAccountSnapshot, AppCore, BoxFuture, CachedProfile, Clock, InMemoryAccountRepository, - InMemoryOperationJournal, InMemorySecretStore, NostrClient, ProfileLoadState, - ProfileRefreshStatus, ProfileRepository, RelayConfiguration, RelayConnectionState, + InMemoryOperationJournal, InMemorySecretStore, NostrClient, ProfileFetchResult, + ProfileLoadState, ProfileRefreshStatus, ProfileRepository, RelayConfiguration, + RelayConnectionState, }; #[derive(Default)] @@ -277,9 +324,10 @@ mod tests { &'a self, _public_key: PublicKey, _relays: &'a [RelayUrl], - ) -> BoxFuture<'a, Result<Option<Kind0ProfileCandidate>, SafeError>> { + _deadline: std::time::Instant, + ) -> BoxFuture<'a, Result<ProfileFetchResult, SafeError>> { let result = self.0.clone(); - Box::pin(async move { result }) + Box::pin(async move { result.map(ProfileFetchResult::complete) }) } } @@ -304,12 +352,13 @@ mod tests { &'a self, _public_key: PublicKey, _relays: &'a [RelayUrl], - ) -> BoxFuture<'a, Result<Option<Kind0ProfileCandidate>, SafeError>> { + _deadline: std::time::Instant, + ) -> BoxFuture<'a, Result<ProfileFetchResult, SafeError>> { Box::pin(async move { self.started.add_permits(1); let permit = self.release.acquire().await.expect("release open"); permit.forget(); - self.result.clone() + self.result.clone().map(ProfileFetchResult::complete) }) } } @@ -324,8 +373,10 @@ mod tests { } fn active_core(profiles: &MemoryProfiles, cached_name: Option<&str>) -> (AppCore, PublicKey) { - let relays = - RelayConfiguration::new(vec![RelayUrl::parse("ws://localhost:8080").expect("relay")]); + let relays = RelayConfiguration::new(vec![ + RelayUrl::parse("ws://localhost:8080", RelayDestinationPolicy::Local).expect("relay"), + ]) + .expect("relay configuration"); let core = AppCore::in_memory(relays); let accounts = InMemoryAccountRepository::default(); let secrets = InMemorySecretStore::default(); @@ -395,7 +446,9 @@ mod tests { Some(RelayConnectionState::Connecting) ); let client = FixedClient(Ok(Some(profile(public_key, "Fresh", 20)))); - let result = client.fetch_profile(plan.public_key(), plan.relays()).await; + let result = client + .fetch_profile(plan.public_key(), plan.relays(), deadline()) + .await; core.complete_profile_refresh(&plan, result, &profiles, &FixedClock) .expect("complete refresh"); assert_eq!( @@ -418,7 +471,12 @@ mod tests { ); let snapshot = core - .refresh_profile_for_active_account(&profiles, &FixedClient(Err(error)), &FixedClock) + .refresh_profile_for_active_account( + &profiles, + &FixedClient(Err(error)), + &FixedClock, + deadline(), + ) .await .expect("nonfatal refresh"); @@ -445,7 +503,8 @@ mod tests { let (core, public_key) = active_core(&profiles, Some("Cached")); let client = BlockingClient::new(Ok(Some(profile(public_key, "Stale", 20)))); - let refresh = core.refresh_profile_for_active_account(&profiles, &client, &FixedClock); + let refresh = + core.refresh_profile_for_active_account(&profiles, &client, &FixedClock, deadline()); let sign_out = async { let permit = client.started.acquire().await.expect("refresh starts"); permit.forget(); @@ -481,6 +540,7 @@ mod tests { &profiles, &FixedClient(Ok(Some(profile(public_key, "First", 10)))), &FixedClock, + deadline(), ) .await .expect("first refresh"); @@ -489,6 +549,7 @@ mod tests { &profiles, &FixedClient(Ok(Some(profile(public_key, "Second", 20)))), &FixedClock, + deadline(), ) .await .expect("second refresh"); @@ -503,12 +564,47 @@ mod tests { ); let signed_out = core.sign_out().expect("sign out"); let no_op = core - .refresh_active_profile(&profiles, &FixedClient(Ok(None)), &FixedClock) + .refresh_active_profile(&profiles, &FixedClient(Ok(None)), &FixedClock, deadline()) .await .expect("signed-out no-op"); assert_eq!(no_op, signed_out); } + fn deadline() -> Instant { + Instant::now() + Duration::from_secs(5) + } + + #[test] + fn partial_relay_success_retains_verified_data_and_marks_degraded_connectivity() { + let profiles = MemoryProfiles::default(); + let (core, public_key) = active_core(&profiles, None); + let plan = core + .begin_profile_refresh() + .expect("begin") + .expect("active"); + let snapshot = core + .complete_profile_refresh( + &plan, + Ok(ProfileFetchResult::partial(Some(profile( + public_key, "Partial", 30, + )))), + &profiles, + &FixedClock, + ) + .expect("partial completion"); + let active = snapshot.active_account().expect("active account"); + assert_eq!(active.relay_state(), RelayConnectionState::Degraded); + assert_eq!(active.profile_state(), ProfileLoadState::Fresh); + assert_eq!( + active.profile().and_then(ProfileMetadata::name), + Some("Partial") + ); + assert_eq!( + snapshot.recoverable_problem().map(SafeError::code), + Some(SafeErrorCode::RelayConnectionFailed) + ); + } + #[test] fn overlapping_refreshes_keep_the_newest_event_regardless_of_completion_order() { let profiles = MemoryProfiles::default(); @@ -524,7 +620,9 @@ mod tests { core.complete_profile_refresh( &second, - Ok(Some(profile(public_key, "Newest", 30))), + Ok(ProfileFetchResult::complete(Some(profile( + public_key, "Newest", 30, + )))), &profiles, &FixedClock, ) @@ -532,7 +630,9 @@ mod tests { let final_snapshot = core .complete_profile_refresh( &first, - Ok(Some(profile(public_key, "Older", 20))), + Ok(ProfileFetchResult::complete(Some(profile( + public_key, "Older", 20, + )))), &profiles, &FixedClock, ) diff --git a/crates/studio_application/src/session.rs b/crates/studio_application/src/session.rs @@ -1,10 +1,8 @@ -use radroots_studio_domain::{PublicKey, SafeError, SafeErrorCode, SafeMessage}; -use radroots_studio_nostr::import_secret; - use crate::{ AccountRepository, ActiveAccountSnapshot, AppCore, AppSnapshot, AppStateRepository, Clock, ProfileLoadState, ProfileRepository, RelayConnectionState, SecretStore, StateTransition, }; +use radroots_studio_domain::{PublicKey, SafeError, SafeErrorCode, SafeMessage}; impl AppCore { /// Drops the active session while retaining accounts, selection, and credentials. @@ -40,7 +38,7 @@ impl AppCore { self.apply_transition(StateTransition::BeginActivation(public_key))?; let prepared = (|| { let credential = secrets.load(public_key)?; - let imported = import_secret(credential)?; + let imported = self.key_material().import(credential)?; let (derived_public_key, _npub, canonical_secret) = imported.into_parts(); drop(canonical_secret); if derived_public_key != public_key { diff --git a/crates/studio_application/src/snapshot.rs b/crates/studio_application/src/snapshot.rs @@ -4,6 +4,8 @@ use radroots_studio_domain::{ AccountSummary, ProfileMetadata, PublicKey, RelayUrl, SafeError, SafeErrorCode, SafeMessage, }; +pub const MAX_CONFIGURED_RELAYS: usize = 16; + #[derive(Clone, Copy, Debug, Default, Eq, Ord, PartialEq, PartialOrd)] pub struct SnapshotRevision(u64); @@ -70,9 +72,17 @@ pub enum ProfileLoadState { pub struct RelayConfiguration(Vec<RelayUrl>); impl RelayConfiguration { - #[must_use] - pub fn new(relays: Vec<RelayUrl>) -> Self { - Self(relays) + /// Creates a bounded, explicitly classified relay configuration. + /// + /// # Errors + /// + /// Returns a safe configuration error before runtime or network work when + /// the relay count exceeds the Studio policy. + pub fn new(relays: Vec<RelayUrl>) -> Result<Self, SafeError> { + if relays.len() > MAX_CONFIGURED_RELAYS { + return Err(relay_limit_exceeded()); + } + Ok(Self(relays)) } #[must_use] @@ -81,6 +91,13 @@ impl RelayConfiguration { } } +const fn relay_limit_exceeded() -> SafeError { + SafeError::new( + SafeErrorCode::InvalidRelayConfiguration, + SafeMessage::new("The Nostr relay configuration exceeds its limit."), + ) +} + #[derive(Clone, Debug, Eq, PartialEq)] pub struct ActiveAccountSnapshot { account: AccountSummary, @@ -279,7 +296,7 @@ const fn invalid_snapshot() -> SafeError { mod tests { use radroots_studio_domain::{ AccountCreatedAt, AccountIdentity, AccountSummary, BindingAvailability, LocalSignerBinding, - PublicKey, UnixTimestamp, + PublicKey, RelayDestinationPolicy, RelayUrl, SafeErrorCode, UnixTimestamp, }; use super::{ @@ -317,6 +334,21 @@ mod tests { } #[test] + fn relay_configuration_rejects_excess_targets_before_runtime_work() { + let relays = (0..=super::MAX_CONFIGURED_RELAYS) + .map(|index| { + RelayUrl::parse( + format!("wss://relay-{index}.example").as_str(), + RelayDestinationPolicy::Public, + ) + .expect("relay") + }) + .collect(); + let error = RelayConfiguration::new(relays).expect_err("relay limit"); + assert_eq!(error.code(), SafeErrorCode::InvalidRelayConfiguration); + } + + #[test] fn revision_helper_is_monotonic_and_checked() { assert_eq!( SnapshotRevision::initial() diff --git a/crates/studio_application/src/test_support.rs b/crates/studio_application/src/test_support.rs @@ -0,0 +1,48 @@ +use std::sync::atomic::{AtomicU8, Ordering}; + +use radroots_studio_domain::{ + Npub, Nsec, PublicKey, SafeError, SafeErrorCode, SafeMessage, SecretKeyInput, +}; + +use crate::{GeneratedKeyMaterial, ImportedKeyMaterial, KeyMaterialProvider}; + +#[derive(Default)] +pub(crate) struct TestKeyMaterialProvider { + next: AtomicU8, +} + +impl KeyMaterialProvider for TestKeyMaterialProvider { + fn generate(&self) -> Result<GeneratedKeyMaterial, SafeError> { + let candidate = self.next.fetch_add(1, Ordering::Relaxed).wrapping_add(9); + let public_key = PublicKey::from_bytes([candidate; 32]); + let secret_byte = public_key.as_bytes()[0]; + Ok(GeneratedKeyMaterial::new( + public_key, + Npub::derive(public_key)?, + SecretKeyInput::parse(format!("{secret_byte:02x}").repeat(32))?, + Nsec::from_encoded( + "nsec1vl029mgpspedva04g90vltkh6fvh240zqtv9k0t9af8935ke9laqsnlfe5".to_owned(), + )?, + )) + } + + fn import(&self, input: SecretKeyInput) -> Result<ImportedKeyMaterial, SafeError> { + let discriminator = input.with_exposed_secret(|value| value.as_bytes()[0]); + if input.with_exposed_secret(|value| value.starts_with("nsec1qq")) { + return Err(invalid_secret_key()); + } + let public_key = PublicKey::from_bytes([discriminator; 32]); + Ok(ImportedKeyMaterial::new( + public_key, + Npub::derive(public_key)?, + input, + )) + } +} + +const fn invalid_secret_key() -> SafeError { + SafeError::new( + SafeErrorCode::InvalidSecretKey, + SafeMessage::new("The Nostr secret key is invalid."), + ) +} diff --git a/crates/studio_domain/src/lib.rs b/crates/studio_domain/src/lib.rs @@ -16,5 +16,5 @@ pub use key::{ MAX_SECRET_KEY_INPUT_BYTES, Npub, Nsec, PublicKey, SecretKeyInput, SecretKeyInputKind, }; pub use profile::{EventId, Kind0ProfileCandidate, ProfileMetadata, select_latest_kind0}; -pub use relay::{RelayUrl, normalize_relay_urls}; +pub use relay::{RelayDestinationPolicy, RelayUrl, normalize_relay_urls}; pub use time::UnixTimestamp; diff --git a/crates/studio_domain/src/relay.rs b/crates/studio_domain/src/relay.rs @@ -2,14 +2,23 @@ use std::collections::HashSet; use std::fmt::{self, Display, Formatter}; -use std::str::FromStr; 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(String); +pub struct RelayUrl { + value: String, + policy: RelayDestinationPolicy, +} impl RelayUrl { /// Parses and normalizes an allowed WebSocket relay URL. @@ -18,7 +27,7 @@ impl RelayUrl { /// /// Returns a safe configuration error for empty or malformed input, /// forbidden schemes, credentials, fragments, or non-loopback `ws://`. - pub fn parse(value: &str) -> Result<Self, SafeError> { + 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()); @@ -32,9 +41,9 @@ impl RelayUrl { return Err(invalid_relay()); } - match parsed.scheme() { - "wss" => {} - "ws" if is_loopback(&parsed) => {} + match (policy, parsed.scheme()) { + (RelayDestinationPolicy::Public | RelayDestinationPolicy::PrivateNetwork, "wss") => {} + (RelayDestinationPolicy::Local, "ws" | "wss") if is_loopback(&parsed) => {} _ => return Err(invalid_relay()), } @@ -42,26 +51,26 @@ impl RelayUrl { return Err(invalid_relay()); } - Ok(Self(parsed.to_string())) + Ok(Self { + value: parsed.to_string(), + policy, + }) } #[must_use] pub fn as_str(&self) -> &str { - &self.0 + &self.value } -} -impl Display for RelayUrl { - fn fmt(&self, formatter: &mut Formatter<'_>) -> fmt::Result { - formatter.write_str(&self.0) + #[must_use] + pub const fn policy(&self) -> RelayDestinationPolicy { + self.policy } } -impl FromStr for RelayUrl { - type Err = SafeError; - - fn from_str(value: &str) -> Result<Self, Self::Err> { - Self::parse(value) +impl Display for RelayUrl { + fn fmt(&self, formatter: &mut Formatter<'_>) -> fmt::Result { + formatter.write_str(&self.value) } } @@ -70,7 +79,10 @@ impl FromStr for RelayUrl { /// # Errors /// /// Returns the first safe relay validation error. -pub fn normalize_relay_urls<I, S>(values: I) -> Result<Vec<RelayUrl>, SafeError> +pub fn normalize_relay_urls<I, S>( + values: I, + policy: RelayDestinationPolicy, +) -> Result<Vec<RelayUrl>, SafeError> where I: IntoIterator<Item = S>, S: AsRef<str>, @@ -78,7 +90,7 @@ where let mut seen = HashSet::new(); let mut relays = Vec::new(); for value in values { - let relay = RelayUrl::parse(value.as_ref())?; + let relay = RelayUrl::parse(value.as_ref(), policy)?; if seen.insert(relay.clone()) { relays.push(relay); } @@ -104,18 +116,34 @@ const fn invalid_relay() -> SafeError { #[cfg(test)] mod tests { - use super::{RelayUrl, normalize_relay_urls}; + use super::{RelayDestinationPolicy, RelayUrl, normalize_relay_urls}; use crate::SafeErrorCode; #[test] fn relay_accepts_secure_remote_and_loopback_development_urls() { - for (input, expected) in [ - (" wss://Relay.Example/path ", "wss://relay.example/path"), - ("ws://localhost:8080", "ws://localhost:8080/"), - ("ws://127.42.1.9:8080", "ws://127.42.1.9:8080/"), - ("ws://[::1]:8080", "ws://[::1]:8080/"), + 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).expect("allowed relay"); + let relay = RelayUrl::parse(input, policy).expect("allowed relay"); assert_eq!(relay.as_str(), expected); assert_eq!(relay.to_string(), expected); } @@ -134,19 +162,23 @@ mod tests { "ws://localhost.evil.example:8080", "wss://relay.example/\nunsafe", ] { - let error = RelayUrl::parse(input).expect_err("forbidden relay"); + let error = RelayUrl::parse(input, RelayDestinationPolicy::Public) + .expect_err("forbidden relay"); assert_eq!(error.code(), SafeErrorCode::InvalidRelayConfiguration); } } #[test] fn relay_deduplication_preserves_normalized_first_seen_order() { - let relays = normalize_relay_urls([ - "wss://relay.example", - " wss://second.example/path ", - "wss://RELAY.example/", - "wss://second.example/path", - ]) + let relays = normalize_relay_urls( + [ + "wss://relay.example", + " wss://second.example/path ", + "wss://RELAY.example/", + "wss://second.example/path", + ], + RelayDestinationPolicy::Public, + ) .expect("valid relays"); assert_eq!( @@ -154,4 +186,13 @@ mod tests { vec!["wss://relay.example/", "wss://second.example/path"] ); } + + #[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()); + let private = RelayUrl::parse("wss://10.0.0.4", RelayDestinationPolicy::PrivateNetwork) + .expect("explicit private network"); + assert_eq!(private.policy(), RelayDestinationPolicy::PrivateNetwork); + } } diff --git a/crates/studio_domain/src/time.rs b/crates/studio_domain/src/time.rs @@ -4,6 +4,8 @@ pub struct UnixTimestamp(i64); impl UnixTimestamp { + pub const UNIX_EPOCH: Self = Self(0); + #[must_use] pub const fn from_seconds(seconds: i64) -> Option<Self> { if seconds < 0 { diff --git a/crates/studio_ffi/Cargo.toml b/crates/studio_ffi/Cargo.toml @@ -18,6 +18,8 @@ crate-type = ["cdylib", "rlib"] directories = "=6.0.0" radroots_studio_application.workspace = true radroots_studio_domain.workspace = true +radroots_studio_nostr.workspace = true +radroots_studio_runtime.workspace = true radroots_studio_storage.workspace = true tokio = { version = "=1.47.1", features = ["macros", "rt-multi-thread", "sync", "time"] } uniffi = "=0.32.0" diff --git a/crates/studio_ffi/src/commands.rs b/crates/studio_ffi/src/commands.rs @@ -9,10 +9,14 @@ use std::time::{Duration, SystemTime, UNIX_EPOCH}; use directories::ProjectDirs; use radroots_studio_application::{ Clock, DurableRequestId, GeneratedKeyRecoveryHandle, RelayConfiguration, RelayRuntimeMode, - RemovalConfirmationToken, SdkNostrClient, relay_configuration_from_environment, + RemovalConfirmationToken, relay_configuration_from_environment, }; use radroots_studio_domain::{PublicKey, SafeError, SecretKeyInput, UnixTimestamp}; -use radroots_studio_storage::{OsKeyringSecretStore, RuntimeActorHandle}; +use radroots_studio_nostr::SdkNostrClient; +use radroots_studio_runtime::{ + RuntimeActorHandle, RuntimeDependencies, UuidInstallationIdentitySource, +}; +use radroots_studio_storage::OsKeyringSecretStore; use crate::{ AccountDto, AppSnapshotDto, WireErrorCategory, WireErrorCode, WireRecoveryAction, @@ -282,14 +286,25 @@ impl StudioAppCore { /// A failed commit must be recovered by importing the already-saved recovery key. pub async fn acknowledge_generated_account_v2( &self, + context: RequestContextDto, request: Arc<GeneratedRecoveryRequest>, ) -> Result<AppSnapshotDto, StudioError> { if request.resolved.swap(true, Ordering::AcqRel) { return Err(generated_recovery_expired()); } + let request_id = DurableRequestId::parse(context.request_id.clone()) + .map_err(|error| StudioError::correlated(error, &context.request_id))?; + let timeout = command_timeout(context.deadline_millis, &context.request_id)?; self.inner .actor - .acknowledge_generated_key_stage(request.handle.id()) + .acknowledge_generated_key_stage( + request.handle.id(), + request_id, + radroots_studio_application::SnapshotRevision::from_value( + context.expected_revision, + ), + timeout, + ) .await .map(|snapshot| self.inner.dto_for(&snapshot)) .map_err(generated_commit_failed) @@ -331,7 +346,7 @@ impl StudioAppCore { .map_err(|error| StudioError::correlated(error, &context.request_id))?; self.inner .actor - .import_secret_key_request( + .import_secret_key( request_id, radroots_studio_application::SnapshotRevision::from_value( context.expected_revision, @@ -449,6 +464,7 @@ impl StudioAppCore { /// Returns a safe confirmation, credential, recovery, or storage error. pub async fn confirm_account_removal( &self, + context: RequestContextDto, request: Arc<RemovalRequest>, ) -> Result<AppSnapshotDto, StudioError> { let token = request @@ -457,9 +473,19 @@ impl StudioAppCore { .unwrap_or_else(std::sync::PoisonError::into_inner) .take() .ok_or_else(confirmation_expired)?; + let request_id = DurableRequestId::parse(context.request_id.clone()) + .map_err(|error| StudioError::correlated(error, &context.request_id))?; + let timeout = command_timeout(context.deadline_millis, &context.request_id)?; self.inner .actor - .confirm_account_removal(token) + .confirm_account_removal( + token, + request_id, + radroots_studio_application::SnapshotRevision::from_value( + context.expected_revision, + ), + timeout, + ) .await .map(|snapshot| self.inner.dto_for(&snapshot)) .map_err(StudioError::from) @@ -499,15 +525,19 @@ impl StudioAppCore { }; let (relays, startup_relay_problem) = local_first_relay_configuration(relay_configuration_from_environment(mode)); - let actor = RuntimeActorHandle::open( + let runtime = runtime()?; + let actor = runtime.block_on(RuntimeActorHandle::open( path, relays, - Arc::new(OsKeyringSecretStore::default()), - Arc::new(SystemClock), - Arc::new(SdkNostrClient::new(Duration::from_secs(5))), - NonZeroUsize::new(ACTOR_MAILBOX_CAPACITY).expect("nonzero actor mailbox capacity"), - runtime().handle(), - )?; + RuntimeDependencies::new( + Arc::new(OsKeyringSecretStore::default()), + Arc::new(SystemClock), + Arc::new(SdkNostrClient::new(Duration::from_secs(5))), + Arc::new(UuidInstallationIdentitySource), + ), + actor_mailbox_capacity()?, + runtime.handle(), + ))?; Ok(Arc::new(Self { inner: Arc::new(RuntimeCore { actor, @@ -538,7 +568,7 @@ impl Clock for SystemClock { .map_or(0, |duration| { i64::try_from(duration.as_secs()).unwrap_or(i64::MAX) }); - UnixTimestamp::from_seconds(seconds).expect("system time is nonnegative") + UnixTimestamp::from_seconds(seconds).unwrap_or(UnixTimestamp::UNIX_EPOCH) } } @@ -574,15 +604,33 @@ fn command_timeout(millis: u64, correlation_id: &str) -> Result<Duration, Studio Ok(Duration::from_millis(millis)) } -pub(crate) fn runtime() -> &'static tokio::runtime::Runtime { - static RUNTIME: OnceLock<tokio::runtime::Runtime> = OnceLock::new(); - RUNTIME.get_or_init(|| { - tokio::runtime::Builder::new_multi_thread() - .enable_all() - .thread_name("radroots-studio-core") - .build() - .expect("Tokio runtime construction") - }) +pub(crate) fn runtime() -> Result<&'static tokio::runtime::Runtime, StudioError> { + static RUNTIME: OnceLock<Result<tokio::runtime::Runtime, ()>> = OnceLock::new(); + RUNTIME + .get_or_init(|| { + tokio::runtime::Builder::new_multi_thread() + .enable_all() + .thread_name("radroots-studio-core") + .build() + .map_err(|_| ()) + }) + .as_ref() + .map_err(|()| runtime_unavailable()) +} + +fn actor_mailbox_capacity() -> Result<NonZeroUsize, StudioError> { + NonZeroUsize::new(ACTOR_MAILBOX_CAPACITY).ok_or_else(runtime_unavailable) +} + +fn runtime_unavailable() -> StudioError { + StudioError::Failure { + code: WireErrorCode::InvalidApplicationState, + category: WireErrorCategory::Lifecycle, + retryable: true, + recovery_action: WireRecoveryAction::RestartApplication, + correlation_id: None, + safe_message: "The application runtime is unavailable.".to_owned(), + } } fn path_unavailable() -> StudioError { @@ -648,9 +696,12 @@ mod tests { use std::num::NonZeroUsize; use std::sync::Arc; - use radroots_studio_application::{InMemorySecretStore, RelayConfiguration, SdkNostrClient}; + use radroots_studio_application::{InMemorySecretStore, RelayConfiguration}; use radroots_studio_domain::SafeError; - use radroots_studio_storage::RuntimeActorHandle; + use radroots_studio_nostr::SdkNostrClient; + use radroots_studio_runtime::{ + RuntimeActorHandle, RuntimeDependencies, UuidInstallationIdentitySource, + }; use radroots_studio_storage::{CREDENTIAL_SERVICE, CURRENT_SCHEMA_VERSION}; @@ -662,15 +713,19 @@ mod tests { verify_compatibility, }; - fn in_memory_core() -> Arc<StudioAppCore> { + async fn in_memory_core() -> Arc<StudioAppCore> { let actor = RuntimeActorHandle::in_memory( RelayConfiguration::default(), - Arc::new(InMemorySecretStore::default()), - Arc::new(SystemClock), - Arc::new(SdkNostrClient::new(std::time::Duration::from_millis(10))), + RuntimeDependencies::new( + Arc::new(InMemorySecretStore::default()), + Arc::new(SystemClock), + Arc::new(SdkNostrClient::new(std::time::Duration::from_millis(10))), + Arc::new(UuidInstallationIdentitySource), + ), NonZeroUsize::new(ACTOR_MAILBOX_CAPACITY).expect("capacity"), - runtime().handle(), + runtime().expect("runtime").handle(), ) + .await .expect("in-memory actor"); Arc::new(StudioAppCore { inner: Arc::new(RuntimeCore { @@ -684,7 +739,7 @@ mod tests { #[tokio::test] async fn exported_bootstrap_and_snapshot_are_revisioned() { - let core = in_memory_core(); + let core = in_memory_core().await; let bootstrapped = core.bootstrap().await.expect("bootstrap"); let current = core.snapshot(); @@ -694,7 +749,7 @@ mod tests { #[tokio::test] async fn request_context_import_replays_one_committed_receipt() { - let core = in_memory_core(); + let core = in_memory_core().await; let initial = core.snapshot(); let context = RequestContextDto { request_id: "ffi-test-import-1".to_owned(), @@ -718,7 +773,7 @@ mod tests { #[tokio::test] async fn generated_recovery_handle_is_one_use_and_acknowledgement_gated() { - let core = in_memory_core(); + let core = in_memory_core().await; let initial = core.snapshot(); let recovery = core .begin_generated_account_v2() @@ -729,13 +784,18 @@ mod tests { let nsec = recovery.take_recovery_nsec().expect("one-use nsec"); assert!(nsec.starts_with("nsec1")); assert!(recovery.take_recovery_nsec().is_err()); + let context = RequestContextDto { + request_id: "ffi-test-generate-1".to_owned(), + expected_revision: initial.revision, + deadline_millis: 5_000, + }; let committed = core - .acknowledge_generated_account_v2(Arc::clone(&recovery)) + .acknowledge_generated_account_v2(context.clone(), Arc::clone(&recovery)) .await .expect("acknowledge"); assert_eq!(committed.accounts.len(), 1); let repeated = core - .acknowledge_generated_account_v2(recovery) + .acknowledge_generated_account_v2(context, recovery) .await .expect_err("repeated acknowledgement"); assert!(matches!( @@ -806,7 +866,7 @@ mod tests { assert_eq!(property("baseline.id"), Some("studio-runtime-v5")); assert_eq!(property("schema.version"), Some("5")); - assert_eq!(CURRENT_SCHEMA_VERSION, 9); + assert_eq!(CURRENT_SCHEMA_VERSION, 10); assert_eq!(property("ffi.contract"), Some("legacy-unversioned-v1")); assert_eq!(property("ffi.snapshot.schema"), Some("1")); assert_eq!(property("ffi.runtime.version"), Some("0.1.0-alpha")); diff --git a/crates/studio_ffi/src/dto.rs b/crates/studio_ffi/src/dto.rs @@ -399,13 +399,19 @@ impl From<ProfileLoadState> for ProfileLoadStateDto { #[cfg(test)] mod tests { + use std::sync::Arc; + use radroots_studio_application::{AppCore, RelayConfiguration}; + use radroots_studio_nostr::NostrKeyMaterialProvider; use super::AppSnapshotDto; #[test] fn snapshot_dto_is_revisioned_public_and_secret_free() { - let core = AppCore::in_memory(RelayConfiguration::default()); + let core = AppCore::new( + RelayConfiguration::default(), + Arc::new(NostrKeyMaterialProvider), + ); let snapshot = core.bootstrap().expect("bootstrap"); let dto = AppSnapshotDto::from(&snapshot); let debug = format!("{dto:?}"); diff --git a/crates/studio_ffi/src/observer.rs b/crates/studio_ffi/src/observer.rs @@ -7,10 +7,7 @@ use radroots_studio_application::ChangeSubscriptionId; use crate::commands::RuntimeCore; use crate::{AppSnapshotDto, StudioAppCore, StudioError}; -const OBSERVER_CHANGE_CAPACITY: NonZeroUsize = match NonZeroUsize::new(64) { - Some(capacity) => capacity, - None => unreachable!(), -}; +const OBSERVER_CHANGE_CAPACITY: NonZeroUsize = NonZeroUsize::MIN.saturating_add(63); #[derive(Clone, Debug, Eq, PartialEq, uniffi::Record)] pub struct SnapshotChangeDto { @@ -37,7 +34,7 @@ pub struct ObserverSubscription { #[uniffi::export] impl ObserverSubscription { - pub fn unsubscribe(&self) { + pub async fn unsubscribe(&self) { let id = self .id .lock() @@ -46,24 +43,17 @@ impl ObserverSubscription { let (Some(core), Some(id)) = (self.core.upgrade(), id) else { return; }; - if let Some(task) = core - .observers - .lock() - .unwrap_or_else(std::sync::PoisonError::into_inner) - .remove(&id) - { + let task = { + core.observers + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .remove(&id) + }; + if let Some(task) = task { task.abort(); + let _ = task.await; } - let actor = core.actor.clone(); - crate::commands::runtime().spawn(async move { - let _ = actor.unsubscribe_changes(id).await; - }); - } -} - -impl Drop for ObserverSubscription { - fn drop(&mut self) { - self.unsubscribe(); + let _ = core.actor.unsubscribe_changes(id).await; } } @@ -89,9 +79,12 @@ impl StudioAppCore { .map_err(StudioError::from)?; let id = subscription.id(); let observer: Arc<dyn StudioChangeObserver> = Arc::from(observer); - let runtime_core = Arc::clone(&self.inner); - let task = crate::commands::runtime().spawn(async move { + let runtime_core = Arc::downgrade(&self.inner); + let task = crate::commands::runtime()?.spawn(async move { while let Some(change) = subscription.receive().await { + let Some(runtime_core) = runtime_core.upgrade() else { + break; + }; observer.on_change(SnapshotChangeDto { snapshot: AppSnapshotDto::from_runtime( change.snapshot(), @@ -132,6 +125,7 @@ impl StudioAppCore { ); for (_, task) in handles { task.abort(); + let _ = task.await; } self.inner.actor.close().await.map_err(StudioError::from)?; Ok(ShutdownReceiptDto { @@ -161,9 +155,12 @@ mod tests { use nostr::{EventBuilder, Keys, Metadata}; use nostr_relay_builder::MockRelay; use nostr_sdk::Client; - use radroots_studio_application::{InMemorySecretStore, RelayConfiguration, SdkNostrClient}; - use radroots_studio_domain::RelayUrl; - use radroots_studio_storage::RuntimeActorHandle; + use radroots_studio_application::{InMemorySecretStore, RelayConfiguration}; + use radroots_studio_domain::{RelayDestinationPolicy, RelayUrl}; + use radroots_studio_nostr::SdkNostrClient; + use radroots_studio_runtime::{ + RuntimeActorHandle, RuntimeDependencies, UuidInstallationIdentitySource, + }; use crate::commands::{ACTOR_MAILBOX_CAPACITY, RuntimeCore, SystemClock, runtime}; use crate::{ @@ -188,19 +185,23 @@ mod tests { } } - fn core() -> Arc<StudioAppCore> { - core_with_relays(RelayConfiguration::default()) + async fn core() -> Arc<StudioAppCore> { + core_with_relays(RelayConfiguration::default()).await } - fn core_with_relays(relays: RelayConfiguration) -> Arc<StudioAppCore> { + async fn core_with_relays(relays: RelayConfiguration) -> Arc<StudioAppCore> { let actor = RuntimeActorHandle::in_memory( relays, - Arc::new(InMemorySecretStore::default()), - Arc::new(SystemClock), - Arc::new(SdkNostrClient::new(std::time::Duration::from_millis(10))), + RuntimeDependencies::new( + Arc::new(InMemorySecretStore::default()), + Arc::new(SystemClock), + Arc::new(SdkNostrClient::new(std::time::Duration::from_millis(10))), + Arc::new(UuidInstallationIdentitySource), + ), NonZeroUsize::new(ACTOR_MAILBOX_CAPACITY).expect("capacity"), - runtime().handle(), + runtime().expect("runtime").handle(), ) + .await .expect("actor"); Arc::new(StudioAppCore { inner: Arc::new(RuntimeCore { @@ -214,8 +215,8 @@ mod tests { #[test] fn callbacks_allow_reentry_and_stop_after_subscription_close() { - runtime().block_on(async { - let core = core(); + runtime().expect("runtime").block_on(async { + let core = core().await; let observer = Arc::new(RecordingObserver::default()); *observer.core.lock().expect("core") = Some(Arc::clone(&core)); let subscription = core @@ -230,7 +231,7 @@ mod tests { .await .expect("idempotent bootstrap"); assert_eq!(observer.snapshots.lock().expect("snapshots").len(), 1); - subscription.unsubscribe(); + subscription.unsubscribe().await; core.inner.actor.sign_out().await.expect("sign out"); assert_eq!(observer.snapshots.lock().expect("snapshots").len(), 1); }); @@ -238,20 +239,23 @@ mod tests { #[test] fn core_close_deregisters_all_observers_and_rejects_new_subscriptions() { - let core = core(); - let observer = Arc::new(RecordingObserver::default()); - let _subscription = runtime() - .block_on(core.subscribe_changes_v2(Box::new(ArcObserver(observer.clone())))) - .expect("subscribe"); + runtime().expect("runtime").block_on(async { + let core = core().await; + let observer = Arc::new(RecordingObserver::default()); + let _subscription = core + .subscribe_changes_v2(Box::new(ArcObserver(observer.clone()))) + .await + .expect("subscribe"); - runtime().block_on(core.shutdown_v2()).expect("shutdown"); + core.shutdown_v2().await.expect("shutdown"); - assert!( - runtime() - .block_on(core.subscribe_changes_v2(Box::new(ArcObserver(observer)))) - .is_err() - ); - assert!(core.inner.observers.lock().expect("observers").is_empty()); + assert!( + core.subscribe_changes_v2(Box::new(ArcObserver(observer))) + .await + .is_err() + ); + assert!(core.inner.observers.lock().expect("observers").is_empty()); + }); } #[tokio::test] @@ -272,9 +276,14 @@ mod tests { .await .expect("publish profile"); - let core = core_with_relays(RelayConfiguration::new(vec![ - RelayUrl::parse(relay_url.as_str()).expect("relay URL"), - ])); + let core = core_with_relays( + RelayConfiguration::new(vec![ + RelayUrl::parse(relay_url.as_str(), RelayDestinationPolicy::Local) + .expect("relay URL"), + ]) + .expect("relay configuration"), + ) + .await; core.bootstrap().await.expect("bootstrap"); let observer = Arc::new(RecordingObserver::default()); *observer.core.lock().expect("core") = Some(Arc::clone(&core)); @@ -310,7 +319,7 @@ mod tests { == Some("FFI Profile") }) })); - subscription.unsubscribe(); + subscription.unsubscribe().await; let count = observer.snapshots.lock().expect("snapshots").len(); core.sign_out().await.expect("sign out"); assert_eq!(observer.snapshots.lock().expect("snapshots").len(), count); diff --git a/crates/studio_nostr/Cargo.toml b/crates/studio_nostr/Cargo.toml @@ -13,7 +13,17 @@ 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" } +radroots_studio_application.workspace = true radroots_studio_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" } +tokio = { version = "=1.47.1", features = ["macros", "rt-multi-thread", "sync", "time"] } [lints] workspace = true diff --git a/crates/studio_nostr/src/client.rs b/crates/studio_nostr/src/client.rs @@ -0,0 +1,304 @@ +use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH}; + +use radroots_studio_domain::{ + PublicKey, RelayDestinationPolicy, RelayUrl, SafeError, SafeErrorCode, SafeMessage, + select_latest_kind0, +}; +use radroots_transport::{ + EventSource, FetchRequest, Target, TargetSet, + outcome::FetchTargetState, + source::{FetchBounds, FetchSelector}, +}; +use radroots_transport_nostr::{Config, NostrTransport, RelayUrlPolicy}; + +use radroots_studio_application::{ + BoxFuture, MAX_CONFIGURED_RELAYS, NostrClient, ProfileFetchResult, +}; + +pub struct SdkNostrClient { + timeout: Duration, +} + +const MAX_PROFILE_EVENTS_PER_RELAY: usize = 64; + +impl SdkNostrClient { + #[must_use] + pub const fn new(timeout: Duration) -> Self { + Self { timeout } + } +} + +impl NostrClient for SdkNostrClient { + fn fetch_profile<'a>( + &'a self, + public_key: PublicKey, + relays: &'a [RelayUrl], + deadline: Instant, + ) -> BoxFuture<'a, Result<ProfileFetchResult, SafeError>> { + Box::pin(async move { + if relays.is_empty() { + return Err(invalid_relay_configuration()); + } + if relays.len() > MAX_CONFIGURED_RELAYS { + return Err(invalid_relay_configuration()); + } + + let author = radroots_identity::PublicKey::from_bytes(*public_key.as_bytes()) + .map_err(|_| profile_refresh_failed())?; + let deadline = deadline.min(Instant::now() + self.timeout); + let mut candidates = Vec::new(); + let mut successful_relays = 0usize; + for policy in [ + RelayDestinationPolicy::Public, + RelayDestinationPolicy::Local, + RelayDestinationPolicy::PrivateNetwork, + ] { + let policy_relays = relays + .iter() + .filter(|relay| relay.policy() == policy) + .collect::<Vec<_>>(); + if policy_relays.is_empty() { + continue; + } + let config = Config::new( + canonical_policy(policy), + policy_relays.iter().map(|relay| relay.as_str()), + ) + .and_then(|config| { + let timeout_ms = + timeout_millis(deadline.saturating_duration_since(Instant::now())); + config.with_timeouts(timeout_ms, timeout_ms, timeout_ms) + }) + .map_err(|_| invalid_relay_configuration())?; + let targets = policy_relays + .iter() + .map(|relay| Target::nostr_relay(relay.as_str())) + .collect::<Result<Vec<_>, _>>() + .map_err(|_| invalid_relay_configuration())?; + let request = FetchRequest::new( + format!("studio-profile-{policy:?}"), + TargetSet::new(targets).map_err(|_| invalid_relay_configuration())?, + FetchBounds::new( + MAX_PROFILE_EVENTS_PER_RELAY as u16, + unix_deadline(deadline.saturating_duration_since(Instant::now()))?, + ) + .map_err(|_| invalid_relay_configuration())?, + ) + .map_err(|_| invalid_relay_configuration())? + .with_selector( + FetchSelector::all() + .with_kinds(vec![0]) + .and_then(|selector| selector.with_authors(vec![author])) + .map_err(|_| invalid_relay_configuration())?, + ); + let page = tokio::time::timeout_at( + deadline.into(), + NostrTransport::new(config).fetch(request), + ) + .await + .map_err(|_| relay_connection_failed())? + .map_err(|_| relay_connection_failed())?; + successful_relays += page + .target_outcomes() + .iter() + .filter(|outcome| { + matches!( + outcome.state(), + FetchTargetState::Complete | FetchTargetState::Partial + ) + }) + .count(); + for observed in page.events() { + candidates.push(crate::parse_verified_kind0( + observed.event().raw_json(), + public_key, + )?); + } + } + if successful_relays == 0 { + return Err(relay_connection_failed()); + } + let candidate = select_latest_kind0(candidates); + if successful_relays == relays.len() { + Ok(ProfileFetchResult::complete(candidate)) + } else { + Ok(ProfileFetchResult::partial(candidate)) + } + }) + } +} + +fn unix_deadline(timeout: Duration) -> Result<u64, SafeError> { + let now = SystemTime::now() + .duration_since(UNIX_EPOCH) + .map_err(|_| relay_connection_failed())?; + let deadline = now + .checked_add(timeout) + .ok_or_else(relay_connection_failed)?; + u64::try_from(deadline.as_millis()) + .ok() + .filter(|deadline| *deadline > 0) + .ok_or_else(relay_connection_failed) +} + +fn timeout_millis(timeout: Duration) -> u64 { + u64::try_from(timeout.as_millis()) + .unwrap_or(u64::MAX) + .clamp(1, 120_000) +} + +const fn canonical_policy(policy: RelayDestinationPolicy) -> RelayUrlPolicy { + match policy { + RelayDestinationPolicy::Public => RelayUrlPolicy::Public, + RelayDestinationPolicy::Local => RelayUrlPolicy::Local, + RelayDestinationPolicy::PrivateNetwork => RelayUrlPolicy::PrivateNetwork, + } +} + +const fn invalid_relay_configuration() -> SafeError { + SafeError::new( + SafeErrorCode::InvalidRelayConfiguration, + SafeMessage::new("No Nostr relay is configured."), + ) +} + +const fn relay_connection_failed() -> SafeError { + SafeError::new( + SafeErrorCode::RelayConnectionFailed, + SafeMessage::new("The Nostr relays could not be reached."), + ) +} + +const fn profile_refresh_failed() -> SafeError { + SafeError::new( + SafeErrorCode::ProfileRefreshFailed, + SafeMessage::new("The Nostr profile could not be refreshed."), + ) +} + +#[cfg(test)] +mod tests { + use std::time::Duration; + + use nostr::{EventBuilder, Keys, Metadata}; + use nostr_relay_builder::MockRelay; + use nostr_sdk::Client; + use radroots_studio_domain::{PublicKey, RelayDestinationPolicy, RelayUrl, SafeErrorCode}; + + use radroots_studio_application::NostrClient; + + use crate::SdkNostrClient; + + #[tokio::test] + async fn sdk_client_fetches_verified_profile_from_ephemeral_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.connect().await; + publisher.wait_for_connection(Duration::from_secs(2)).await; + publisher + .send_event_builder(EventBuilder::metadata( + &Metadata::new().name("Farmer").display_name("Farm Account"), + )) + .await + .expect("publish metadata"); + + let adapter = SdkNostrClient::new(Duration::from_secs(2)); + let domain_relay = RelayUrl::parse(relay_url.as_str(), RelayDestinationPolicy::Local) + .expect("domain relay URL"); + let public_key = PublicKey::from_bytes(keys.public_key().to_bytes()); + let fetched = adapter + .fetch_profile( + public_key, + &[domain_relay], + std::time::Instant::now() + Duration::from_secs(2), + ) + .await + .expect("fetch 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 Account")); + assert_eq!( + completeness, + radroots_studio_application::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)) + .fetch_profile( + PublicKey::from_bytes([1; 32]), + &[], + std::time::Instant::now() + Duration::from_millis(10), + ) + .await + .expect_err("empty relay list"); + + assert_eq!(error.code(), SafeErrorCode::InvalidRelayConfiguration); + } + + #[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 = [ + RelayUrl::parse(relay_url.as_str(), RelayDestinationPolicy::Local).expect("live relay"), + RelayUrl::parse("ws://127.0.0.1:1", RelayDestinationPolicy::Local) + .expect("unavailable relay"), + ]; + let fetched = SdkNostrClient::new(Duration::from_millis(250)) + .fetch_profile( + PublicKey::from_bytes(keys.public_key().to_bytes()), + &configured, + std::time::Instant::now() + Duration::from_secs(1), + ) + .await + .expect("partial fetch"); + let (candidate, completeness) = fetched.into_parts(); + assert!(candidate.is_some()); + assert_eq!( + completeness, + radroots_studio_application::RelayFetchCompleteness::Partial + ); + publisher.shutdown().await; + relay.shutdown(); + } + + #[test] + fn studio_policy_maps_exactly_to_the_canonical_transport_policy() { + assert_eq!( + super::canonical_policy(RelayDestinationPolicy::Public), + radroots_transport_nostr::RelayUrlPolicy::Public + ); + assert_eq!( + super::canonical_policy(RelayDestinationPolicy::Local), + radroots_transport_nostr::RelayUrlPolicy::Local + ); + assert_eq!( + super::canonical_policy(RelayDestinationPolicy::PrivateNetwork), + radroots_transport_nostr::RelayUrlPolicy::PrivateNetwork + ); + } +} diff --git a/crates/studio_nostr/src/keys.rs b/crates/studio_nostr/src/keys.rs @@ -1,34 +1,11 @@ use nostr::{Keys, ToBech32}; +use radroots_studio_application::{GeneratedKeyMaterial, ImportedKeyMaterial, KeyMaterialProvider}; use radroots_studio_domain::{ Npub, Nsec, PublicKey, SafeError, SafeErrorCode, SafeMessage, SecretKeyInput, }; -pub struct GeneratedKeyMaterial { - public_key: PublicKey, - npub: Npub, - secret: SecretKeyInput, - nsec: Nsec, -} - -impl GeneratedKeyMaterial { - #[must_use] - pub fn into_parts(self) -> (PublicKey, Npub, SecretKeyInput, Nsec) { - (self.public_key, self.npub, self.secret, self.nsec) - } -} - -pub struct ImportedKeyMaterial { - public_key: PublicKey, - npub: Npub, - secret: SecretKeyInput, -} - -impl ImportedKeyMaterial { - #[must_use] - pub fn into_parts(self) -> (PublicKey, Npub, SecretKeyInput) { - (self.public_key, self.npub, self.secret) - } -} +#[derive(Clone, Copy, Debug, Default)] +pub struct NostrKeyMaterialProvider; /// Generates one cryptographically random local Nostr keypair. /// @@ -36,40 +13,27 @@ impl ImportedKeyMaterial { /// /// Returns a safe key error if an upstream encoding cannot be represented by /// the stricter Radroots domain boundary. -pub fn generate_local_keypair() -> Result<GeneratedKeyMaterial, SafeError> { - let keys = Keys::generate(); - let (public_key, npub, secret, nsec) = encode_keys(&keys)?; - Ok(GeneratedKeyMaterial { - public_key, - npub, - secret, - nsec, - }) -} +impl KeyMaterialProvider for NostrKeyMaterialProvider { + fn generate(&self) -> Result<GeneratedKeyMaterial, SafeError> { + let keys = Keys::generate(); + let (public_key, npub, secret, nsec) = encode_keys(&keys)?; + Ok(GeneratedKeyMaterial::new(public_key, npub, secret, nsec)) + } -/// Parses nsec or canonical secret hex and derives public Nostr identity. -/// -/// # Errors -/// -/// Returns a safe invalid-secret-key error for checksum, scalar, or encoding -/// failures without exposing the rejected input. -pub fn import_secret(input: SecretKeyInput) -> Result<ImportedKeyMaterial, SafeError> { - let keys = input - .with_exposed_secret(Keys::parse) - .map_err(|_| invalid_secret_key())?; - drop(input); - let public_key = PublicKey::from_bytes(keys.public_key().to_bytes()); - let npub = keys - .public_key() - .to_bech32() - .map_err(|_| invalid_public_key()) - .and_then(Npub::from_encoded)?; - let secret = SecretKeyInput::parse(keys.secret_key().to_secret_hex())?; - Ok(ImportedKeyMaterial { - public_key, - npub, - secret, - }) + fn import(&self, input: SecretKeyInput) -> Result<ImportedKeyMaterial, SafeError> { + let keys = input + .with_exposed_secret(Keys::parse) + .map_err(|_| invalid_secret_key())?; + drop(input); + let public_key = PublicKey::from_bytes(keys.public_key().to_bytes()); + let npub = keys + .public_key() + .to_bech32() + .map_err(|_| invalid_public_key()) + .and_then(Npub::from_encoded)?; + let secret = SecretKeyInput::parse(keys.secret_key().to_secret_hex())?; + Ok(ImportedKeyMaterial::new(public_key, npub, secret)) + } } fn encode_keys(keys: &Keys) -> Result<(PublicKey, Npub, SecretKeyInput, Nsec), SafeError> { @@ -106,7 +70,9 @@ const fn invalid_public_key() -> SafeError { mod tests { use radroots_studio_domain::{SafeErrorCode, SecretKeyInput}; - use super::{generate_local_keypair, import_secret}; + use radroots_studio_application::KeyMaterialProvider; + + use super::NostrKeyMaterialProvider; const SECRET_HEX: &str = "7e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7"; const NSEC: &str = "nsec1vl029mgpspedva04g90vltkh6fvh240zqtv9k0t9af8935ke9laqsnlfe5"; @@ -116,7 +82,7 @@ mod tests { #[test] fn keys_generate_valid_redacted_material() { - let generated = generate_local_keypair().expect("generated"); + let generated = NostrKeyMaterialProvider.generate().expect("generated"); let (public_key, npub, secret, nsec) = generated.into_parts(); assert_eq!(public_key.to_hex().len(), 64); assert!(npub.as_str().starts_with("npub1")); @@ -128,9 +94,11 @@ mod tests { #[test] fn keys_import_known_nsec_and_hex_vectors() { - let from_nsec = import_secret(SecretKeyInput::parse(NSEC.to_owned()).expect("nsec")) + let from_nsec = NostrKeyMaterialProvider + .import(SecretKeyInput::parse(NSEC.to_owned()).expect("nsec")) .expect("import nsec"); - let from_hex = import_secret(SecretKeyInput::parse(SECRET_HEX.to_owned()).expect("hex")) + let from_hex = NostrKeyMaterialProvider + .import(SecretKeyInput::parse(SECRET_HEX.to_owned()).expect("hex")) .expect("import hex"); let (nsec_public, nsec_npub, _) = from_nsec.into_parts(); let (hex_public, hex_npub, _) = from_hex.into_parts(); @@ -146,7 +114,10 @@ mod tests { "nsec1qqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqq".to_owned(), ) .expect("domain shape"); - let error = import_secret(input).err().expect("invalid checksum"); + let error = NostrKeyMaterialProvider + .import(input) + .err() + .expect("invalid checksum"); assert_eq!(error.code(), SafeErrorCode::InvalidSecretKey); } } diff --git a/crates/studio_nostr/src/lib.rs b/crates/studio_nostr/src/lib.rs @@ -1,7 +1,9 @@ #![doc = "Radroots Studio Nostr protocol adapters."] +pub mod client; pub mod keys; pub mod profile; -pub use keys::{GeneratedKeyMaterial, ImportedKeyMaterial, generate_local_keypair, import_secret}; +pub use client::SdkNostrClient; +pub use keys::NostrKeyMaterialProvider; pub use profile::parse_verified_kind0; diff --git a/crates/studio_preferences/Cargo.toml b/crates/studio_preferences/Cargo.toml @@ -0,0 +1,18 @@ +[package] +name = "radroots_studio_preferences" +description = "Private preference model for Radroots Studio" +version = "0.1.0-alpha" +edition.workspace = true +authors.workspace = true +rust-version.workspace = true +license = "MPL-2.0" +repository.workspace = true +homepage.workspace = true +publish = false +include = ["src/**", "Cargo.toml"] + +[dependencies] +url = "=2.5.8" + +[lints] +workspace = true diff --git a/crates/studio_preferences/src/lib.rs b/crates/studio_preferences/src/lib.rs @@ -0,0 +1,266 @@ +// SPDX-License-Identifier: MPL-2.0 +//! UI-neutral Studio preference state. +//! +//! This module carries forward the uniquely required preference behavior from +//! source commit `6074a4745be361f21bb47d4778c74a14b2d57954`. It intentionally +//! excludes that source's process-global state, sample account, and FFI layer. + +use url::Url; + +pub const PREFERENCES_SCHEMA_VERSION: u32 = 1; +const MAX_SUMMARY_BYTES: usize = 256; +const MAX_SERVER_URL_BYTES: usize = 2_048; + +#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)] +pub enum UpdateChannel { + #[default] + Stable, + Preview, +} + +#[derive(Clone, Debug, PartialEq, Eq)] +pub struct StudioPreferences { + pub allow_incoming_connections: bool, + pub use_radroots_dns: bool, + pub use_radroots_subnets: bool, + pub launch_at_login: bool, + pub hide_dock_icon: bool, + pub vpn_on_demand_enabled: bool, + pub run_as_exit_node: bool, + pub allow_local_network_access: bool, + pub automatically_check_for_updates: bool, + pub update_channel: UpdateChannel, + pub last_update_check_summary: String, + pub alternate_server_url: String, +} + +impl Default for StudioPreferences { + fn default() -> Self { + Self { + allow_incoming_connections: true, + use_radroots_dns: true, + use_radroots_subnets: true, + launch_at_login: true, + hide_dock_icon: false, + vpn_on_demand_enabled: false, + run_as_exit_node: false, + allow_local_network_access: false, + automatically_check_for_updates: true, + update_channel: UpdateChannel::Stable, + last_update_check_summary: String::new(), + alternate_server_url: String::new(), + } + } +} + +#[derive(Clone, Debug, PartialEq, Eq)] +pub struct PreferencesState { + schema_version: u32, + revision: u64, + preferences: StudioPreferences, +} + +impl Default for PreferencesState { + fn default() -> Self { + Self { + schema_version: PREFERENCES_SCHEMA_VERSION, + revision: 1, + preferences: StudioPreferences::default(), + } + } +} + +impl PreferencesState { + #[must_use] + pub const fn schema_version(&self) -> u32 { + self.schema_version + } + + #[must_use] + pub const fn revision(&self) -> u64 { + self.revision + } + + #[must_use] + pub const fn preferences(&self) -> &StudioPreferences { + &self.preferences + } + + pub fn apply(&mut self, change: PreferenceChange) -> Result<bool, PreferencesError> { + let mut candidate = self.preferences.clone(); + change.apply_to(&mut candidate)?; + if candidate == self.preferences { + return Ok(false); + } + self.revision = self + .revision + .checked_add(1) + .ok_or(PreferencesError::RevisionExhausted)?; + self.preferences = candidate; + Ok(true) + } +} + +#[derive(Clone, Debug, PartialEq, Eq)] +pub enum PreferenceChange { + AllowIncomingConnections(bool), + UseRadrootsDns(bool), + UseRadrootsSubnets(bool), + LaunchAtLogin(bool), + HideDockIcon(bool), + VpnOnDemandEnabled(bool), + RunAsExitNode(bool), + AllowLocalNetworkAccess(bool), + AutomaticallyCheckForUpdates(bool), + UpdateChannel(UpdateChannel), + LastUpdateCheckSummary(String), + AlternateServerUrl(String), +} + +impl PreferenceChange { + fn apply_to(self, preferences: &mut StudioPreferences) -> Result<(), PreferencesError> { + match self { + Self::AllowIncomingConnections(value) => { + preferences.allow_incoming_connections = value; + } + Self::UseRadrootsDns(value) => preferences.use_radroots_dns = value, + Self::UseRadrootsSubnets(value) => preferences.use_radroots_subnets = value, + Self::LaunchAtLogin(value) => preferences.launch_at_login = value, + Self::HideDockIcon(value) => preferences.hide_dock_icon = value, + Self::VpnOnDemandEnabled(value) => preferences.vpn_on_demand_enabled = value, + Self::RunAsExitNode(value) => preferences.run_as_exit_node = value, + Self::AllowLocalNetworkAccess(value) => { + preferences.allow_local_network_access = value; + } + Self::AutomaticallyCheckForUpdates(value) => { + preferences.automatically_check_for_updates = value; + } + Self::UpdateChannel(value) => preferences.update_channel = value, + Self::LastUpdateCheckSummary(value) => { + preferences.last_update_check_summary = validated_summary(value)?; + } + Self::AlternateServerUrl(value) => { + preferences.alternate_server_url = validated_server_url(value)?; + } + } + Ok(()) + } +} + +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum PreferencesError { + InvalidSummary, + InvalidAlternateServerUrl, + RevisionExhausted, +} + +fn validated_summary(value: String) -> Result<String, PreferencesError> { + let value = value.trim(); + if value.len() > MAX_SUMMARY_BYTES || value.chars().any(char::is_control) { + return Err(PreferencesError::InvalidSummary); + } + Ok(value.to_owned()) +} + +fn validated_server_url(value: String) -> Result<String, PreferencesError> { + let value = value.trim(); + if value.is_empty() { + return Ok(String::new()); + } + if value.len() > MAX_SERVER_URL_BYTES || value.chars().any(char::is_control) { + return Err(PreferencesError::InvalidAlternateServerUrl); + } + let url = Url::parse(value).map_err(|_| PreferencesError::InvalidAlternateServerUrl)?; + if url.scheme() != "https" + || url.host_str().is_none() + || !url.username().is_empty() + || url.password().is_some() + || url.fragment().is_some() + { + return Err(PreferencesError::InvalidAlternateServerUrl); + } + Ok(url.to_string()) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn defaults_preserve_the_reviewed_boolean_policy_without_sample_identity() { + let state = PreferencesState::default(); + assert_eq!(state.schema_version(), PREFERENCES_SCHEMA_VERSION); + assert_eq!(state.revision(), 1); + assert!(state.preferences().allow_incoming_connections); + assert!(state.preferences().use_radroots_dns); + assert!(state.preferences().use_radroots_subnets); + assert!(state.preferences().launch_at_login); + assert!(state.preferences().automatically_check_for_updates); + assert_eq!(state.preferences().update_channel, UpdateChannel::Stable); + assert!(state.preferences().last_update_check_summary.is_empty()); + assert!(state.preferences().alternate_server_url.is_empty()); + } + + #[test] + fn revisions_advance_only_when_a_valid_canonical_value_changes() { + let mut state = PreferencesState::default(); + assert!( + !state + .apply(PreferenceChange::HideDockIcon(false)) + .expect("unchanged value") + ); + assert_eq!(state.revision(), 1); + assert!( + state + .apply(PreferenceChange::HideDockIcon(true)) + .expect("changed value") + ); + assert_eq!(state.revision(), 2); + assert!(state.preferences().hide_dock_icon); + } + + #[test] + fn alternate_server_is_trimmed_canonical_and_credential_free() { + let mut state = PreferencesState::default(); + state + .apply(PreferenceChange::AlternateServerUrl( + " https://example.com/api ".to_owned(), + )) + .expect("valid URL"); + assert_eq!( + state.preferences().alternate_server_url, + "https://example.com/api" + ); + for invalid in [ + "http://example.com", + "https://user@example.com", + "https://example.com/#fragment", + "not a URL", + ] { + assert_eq!( + state.apply(PreferenceChange::AlternateServerUrl(invalid.to_owned())), + Err(PreferencesError::InvalidAlternateServerUrl) + ); + } + } + + #[test] + fn summary_is_bounded_trimmed_and_control_free() { + let mut state = PreferencesState::default(); + state + .apply(PreferenceChange::LastUpdateCheckSummary( + " Checked today ".to_owned(), + )) + .expect("valid summary"); + assert_eq!( + state.preferences().last_update_check_summary, + "Checked today" + ); + assert_eq!( + state.apply(PreferenceChange::LastUpdateCheckSummary( + "bad\nvalue".to_owned() + )), + Err(PreferencesError::InvalidSummary) + ); + } +} diff --git a/crates/studio_runtime/Cargo.toml b/crates/studio_runtime/Cargo.toml @@ -0,0 +1,29 @@ +[package] +name = "radroots_studio_runtime" +description = "Private supervised composition runtime for Radroots Studio" +version = "0.1.0-alpha" +edition.workspace = true +authors.workspace = true +rust-version.workspace = true +license = "GPL-3.0-only" +repository.workspace = true +homepage.workspace = true +publish = false +include = ["src/**", "tests/**", "Cargo.toml"] + +[dependencies] +radroots_studio_application.workspace = true +radroots_studio_domain.workspace = true +radroots_studio_nostr.workspace = true +radroots_studio_storage.workspace = true +tokio = { version = "=1.47.1", features = ["macros", "rt-multi-thread", "sync", "time"] } +uuid.workspace = true + +[dev-dependencies] +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" } +tempfile = "=3.23.0" + +[lints] +workspace = true diff --git a/crates/studio_runtime/src/blocking.rs b/crates/studio_runtime/src/blocking.rs @@ -0,0 +1,114 @@ +use std::sync::Arc; +use std::time::Instant; + +use tokio::runtime::Handle; +use tokio::sync::Semaphore; + +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub(crate) enum BlockingExecutionError { + DeadlineElapsed, + Saturated, + TaskFailed, +} + +#[derive(Clone)] +pub(crate) struct BoundedBlockingExecutor { + permits: Arc<Semaphore>, + runtime: Handle, +} + +impl BoundedBlockingExecutor { + pub(crate) fn new(capacity: usize, runtime: &Handle) -> Self { + Self { + permits: Arc::new(Semaphore::new(capacity)), + runtime: runtime.clone(), + } + } + + pub(crate) async fn execute<T, F>( + &self, + deadline: Instant, + operation: F, + ) -> Result<T, BlockingExecutionError> + where + T: Send + 'static, + F: FnOnce() -> T + Send + 'static, + { + if Instant::now() >= deadline { + return Err(BlockingExecutionError::DeadlineElapsed); + } + let permit = self + .permits + .clone() + .try_acquire_owned() + .map_err(|_| BlockingExecutionError::Saturated)?; + self.runtime + .spawn_blocking(move || { + let _permit = permit; + operation() + }) + .await + .map_err(|_| BlockingExecutionError::TaskFailed) + } +} + +#[cfg(test)] +mod tests { + use std::sync::{Arc, Condvar, Mutex}; + use std::time::{Duration, Instant}; + + use tokio::sync::oneshot; + + use super::{BlockingExecutionError, BoundedBlockingExecutor}; + + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn executor_rejects_saturation_without_starting_excess_work() { + let executor = BoundedBlockingExecutor::new(1, &tokio::runtime::Handle::current()); + let release = Arc::new((Mutex::new(false), Condvar::new())); + let first_release = Arc::clone(&release); + let (started, started_rx) = oneshot::channel(); + let first_executor = executor.clone(); + let first = tokio::spawn(async move { + first_executor + .execute(Instant::now() + Duration::from_secs(5), move || { + let _ = started.send(()); + let (lock, ready) = &*first_release; + let mut released = lock.lock().expect("release lock"); + while !*released { + released = ready.wait(released).expect("release wait"); + } + 7 + }) + .await + }); + started_rx.await.expect("first work starts"); + + let second = executor + .execute(Instant::now() + Duration::from_secs(5), || 9) + .await; + assert_eq!(second, Err(BlockingExecutionError::Saturated)); + + let (lock, ready) = &*release; + *lock.lock().expect("release lock") = true; + ready.notify_all(); + assert_eq!(first.await.expect("first join"), Ok(7)); + } + + #[tokio::test] + async fn executor_rejects_expired_work_before_spawn() { + let executor = BoundedBlockingExecutor::new(1, &tokio::runtime::Handle::current()); + let result = executor.execute(Instant::now(), || 1).await; + assert_eq!(result, Err(BlockingExecutionError::DeadlineElapsed)); + } + + #[tokio::test] + async fn executor_classifies_panicked_work_without_panicking_the_actor() { + let executor = BoundedBlockingExecutor::new(1, &tokio::runtime::Handle::current()); + let result = executor + .execute::<(), _>(Instant::now() + Duration::from_secs(1), || { + panic!("test-only blocking task failure"); + }) + .await; + assert_eq!(result, Err(BlockingExecutionError::TaskFailed)); + } +} diff --git a/crates/studio_runtime/src/installation.rs b/crates/studio_runtime/src/installation.rs @@ -0,0 +1,58 @@ +use radroots_studio_domain::{SafeError, SafeErrorCode, SafeMessage}; + +#[derive(Clone, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)] +pub struct InstallationIdentity(String); + +impl InstallationIdentity { + pub fn parse(value: impl Into<String>) -> Result<Self, SafeError> { + let value = value.into(); + if value.len() != 32 + || !value + .bytes() + .all(|byte| byte.is_ascii_digit() || (b'a'..=b'f').contains(&byte)) + { + return Err(invalid_installation_identity()); + } + Ok(Self(value)) + } + + #[must_use] + pub fn as_str(&self) -> &str { + self.0.as_str() + } +} + +pub trait InstallationIdentitySource: Send + Sync { + fn generate(&self) -> Result<InstallationIdentity, SafeError>; +} + +#[derive(Clone, Copy, Debug, Default)] +pub struct UuidInstallationIdentitySource; + +impl InstallationIdentitySource for UuidInstallationIdentitySource { + fn generate(&self) -> Result<InstallationIdentity, SafeError> { + InstallationIdentity::parse(uuid::Uuid::new_v4().simple().to_string()) + } +} + +const fn invalid_installation_identity() -> SafeError { + SafeError::new( + SafeErrorCode::InvalidApplicationState, + SafeMessage::new("The installation identity is invalid."), + ) +} + +#[cfg(test)] +mod tests { + use super::{InstallationIdentity, InstallationIdentitySource, UuidInstallationIdentitySource}; + + #[test] + fn installation_identity_is_fixed_width_lowercase_hex() { + let identity = UuidInstallationIdentitySource.generate().expect("identity"); + assert_eq!(identity.as_str().len(), 32); + assert!(InstallationIdentity::parse(identity.as_str()).is_ok()); + for denied in ["", "AAaabbccddeeff001122334455667788", "not-an-identity"] { + assert!(InstallationIdentity::parse(denied).is_err()); + } + } +} diff --git a/crates/studio_runtime/src/lib.rs b/crates/studio_runtime/src/lib.rs @@ -0,0 +1,12 @@ +#![doc = "Radroots Studio supervised runtime composition."] + +mod blocking; +mod installation; +mod persistence; +mod runtime_actor; + +pub use installation::{ + InstallationIdentity, InstallationIdentitySource, UuidInstallationIdentitySource, +}; +pub use persistence::PersistentAppCore; +pub use runtime_actor::{RuntimeActorHandle, RuntimeChangeSubscription, RuntimeDependencies}; diff --git a/crates/studio_runtime/src/persistence.rs b/crates/studio_runtime/src/persistence.rs @@ -0,0 +1,746 @@ +use std::path::Path; +use std::sync::Arc; + +use radroots_studio_application::{ + AppCore, AppSnapshot, Clock, DurableRequestId, GenerateAccountReceipt, ImportAccountReceipt, + KeyMaterialProvider, RelayConfiguration, RemovalConfirmationToken, SecretStore, + StagedGeneratedKey, +}; +use radroots_studio_domain::{PublicKey, SafeError, SecretKeyInput}; +use radroots_studio_nostr::NostrKeyMaterialProvider; + +use radroots_studio_storage::Database; + +use crate::{InstallationIdentity, InstallationIdentitySource}; + +pub struct PersistentAppCore { + core: AppCore, + database: Database, + key_material: Arc<dyn KeyMaterialProvider>, +} + +impl PersistentAppCore { + pub(crate) fn initialize_installation_identity( + &self, + source: &dyn InstallationIdentitySource, + ) -> Result<InstallationIdentity, SafeError> { + if let Some(existing) = self.database.load_installation_id()? { + return InstallationIdentity::parse(existing); + } + let candidate = source.generate()?; + InstallationIdentity::parse( + self.database + .initialize_installation_id(candidate.as_str())?, + ) + } + + /// Commits an acknowledged generated-key stage through the durable coordinator. + /// + /// # Errors + /// + /// Returns a safe conflict, credential, storage, or recovery error. + pub fn commit_staged_generated_key( + &self, + request_id: &DurableRequestId, + staged: StagedGeneratedKey, + secrets: &(impl SecretStore + ?Sized), + clock: &(impl Clock + ?Sized), + ) -> Result<ImportAccountReceipt, SafeError> { + self.core.commit_staged_generated_key( + request_id, + staged, + &self.database, + &self.database, + secrets, + &self.database, + clock, + ) + } + + /// Opens the application database without accessing credentials or relays. + /// + /// # Errors + /// + /// Returns a safe storage error when the database cannot be opened or migrated. + pub fn open(path: &Path, relay_configuration: RelayConfiguration) -> Result<Self, SafeError> { + let key_material: Arc<dyn KeyMaterialProvider> = Arc::new(NostrKeyMaterialProvider); + Ok(Self { + core: AppCore::new(relay_configuration, Arc::clone(&key_material)), + database: Database::open(path)?, + key_material, + }) + } + + /// Creates an isolated persistent-core adapter for tests. + /// + /// # Errors + /// + /// Returns a safe storage error when the database cannot be initialized. + pub fn in_memory(relay_configuration: RelayConfiguration) -> Result<Self, SafeError> { + let key_material: Arc<dyn KeyMaterialProvider> = Arc::new(NostrKeyMaterialProvider); + Ok(Self { + core: AppCore::new(relay_configuration, Arc::clone(&key_material)), + database: Database::in_memory()?, + key_material, + }) + } + + pub(crate) fn key_material(&self) -> &dyn KeyMaterialProvider { + self.key_material.as_ref() + } + + /// Restores public accounts and selection while keeping the session signed out. + /// + /// # Errors + /// + /// Returns a safe storage or application-state error after publishing a fatal + /// snapshot when durable state cannot be restored. + pub fn bootstrap( + &self, + secrets: &(impl SecretStore + ?Sized), + clock: &(impl Clock + ?Sized), + ) -> Result<AppSnapshot, SafeError> { + self.core.recover_durable_operations( + &self.database, + &self.database, + secrets, + &self.database, + clock, + )?; + self.core.recover_pending_operations( + &self.database, + &self.database, + secrets, + &self.database, + clock, + )?; + self.core.bootstrap_from(&self.database, &self.database) + } + + /// Generates and durably persists one selected, signed-out local account. + /// + /// # Errors + /// + /// Returns a safe credential, storage, key, or application-state error. + pub fn generate_account( + &self, + secrets: &(impl SecretStore + ?Sized), + clock: &(impl Clock + ?Sized), + ) -> Result<GenerateAccountReceipt, SafeError> { + self.core.generate_account( + &self.database, + &self.database, + secrets, + &self.database, + clock, + ) + } + + /// Imports and durably persists one selected, signed-out local account. + /// + /// # Errors + /// + /// Returns a safe credential, storage, key, or application-state error. + pub fn import_secret_key( + &self, + input: SecretKeyInput, + secrets: &(impl SecretStore + ?Sized), + clock: &(impl Clock + ?Sized), + ) -> Result<ImportAccountReceipt, SafeError> { + self.core.import_secret_key( + input, + &self.database, + &self.database, + secrets, + &self.database, + clock, + ) + } + + /// Generates an account through the durable request coordinator. + /// + /// # Errors + /// + /// Returns a safe conflict, credential, storage, or application-state error. + pub fn generate_account_durable( + &self, + request_id: &DurableRequestId, + expected_revision: u64, + secrets: &(impl SecretStore + ?Sized), + clock: &(impl Clock + ?Sized), + ) -> Result<GenerateAccountReceipt, SafeError> { + self.core.generate_account_durable( + request_id, + expected_revision, + &self.database, + &self.database, + secrets, + &self.database, + clock, + ) + } + + /// Imports or repairs an account through the durable request coordinator. + /// + /// # Errors + /// + /// Returns a safe conflict, validation, credential, storage, or state error. + pub fn import_secret_key_durable( + &self, + request_id: &DurableRequestId, + expected_revision: u64, + input: SecretKeyInput, + secrets: &(impl SecretStore + ?Sized), + clock: &(impl Clock + ?Sized), + ) -> Result<ImportAccountReceipt, SafeError> { + self.core.import_secret_key_durable( + request_id, + expected_revision, + input, + &self.database, + &self.database, + secrets, + &self.database, + clock, + ) + } + + /// Persists and publishes one saved-account selection without activation. + /// + /// # Errors + /// + /// Returns a safe account, storage, or application-state error. + pub fn select_account(&self, public_key: PublicKey) -> Result<AppSnapshot, SafeError> { + self.core + .select_account(public_key, &self.database, &self.database) + } + + /// Activates a saved account after validating its credential and cached profile. + /// + /// # Errors + /// + /// Returns a safe account, credential, storage, or application-state error. + pub fn activate_account( + &self, + public_key: PublicKey, + secrets: &(impl SecretStore + ?Sized), + clock: &(impl Clock + ?Sized), + ) -> Result<AppSnapshot, SafeError> { + self.core.activate_account( + public_key, + &self.database, + &self.database, + &self.database, + secrets, + clock, + ) + } + + /// Signs out while retaining durable account data and credentials. + /// + /// # Errors + /// + /// Returns a safe application-state error if sign out cannot complete. + pub fn sign_out(&self) -> Result<AppSnapshot, SafeError> { + self.core.sign_out() + } + + /// Issues a revision-bound, single-use account-removal confirmation. + /// + /// # Errors + /// + /// Returns a safe error when the target account is not saved. + pub fn request_account_removal( + &self, + public_key: PublicKey, + clock: &(impl Clock + ?Sized), + ) -> Result<RemovalConfirmationToken, SafeError> { + self.core.request_account_removal(public_key, clock) + } + + /// Permanently removes one confirmed account and its credential. + /// + /// # Errors + /// + /// Returns a safe confirmation, credential, storage, recovery, or state error. + pub fn confirm_account_removal( + &self, + token: RemovalConfirmationToken, + secrets: &(impl SecretStore + ?Sized), + clock: &(impl Clock + ?Sized), + ) -> Result<AppSnapshot, SafeError> { + self.core.confirm_account_removal( + token, + &self.database, + &self.database, + secrets, + &self.database, + clock, + ) + } + + /// Executes a confirmed removal through the durable request coordinator. + /// + /// # Errors + /// + /// Returns a safe expiry, conflict, credential, storage, recovery, or state error. + pub fn confirm_account_removal_durable( + &self, + request_id: &DurableRequestId, + token: RemovalConfirmationToken, + secrets: &(impl SecretStore + ?Sized), + clock: &(impl Clock + ?Sized), + ) -> Result<AppSnapshot, SafeError> { + self.core.confirm_account_removal_durable( + request_id, + token, + &self.database, + &self.database, + secrets, + &self.database, + clock, + ) + } + + #[must_use] + pub const fn core(&self) -> &AppCore { + &self.core + } + + #[must_use] + pub const fn database(&self) -> &Database { + &self.database + } +} + +#[cfg(test)] +mod tests { + use std::fs; + + use radroots_studio_application::{ + AccountOperationKind, AccountOperationPhase, AccountRepository, AppLifecycle, + AppStateRepository, Clock, DurableOperationKind, DurableOperationPhase, + DurableOperationRepository, DurableRequestId, DurableTerminalOutcome, FailureSecretStore, + InMemorySecretStore, OperationJournal, OperationPriorState, RelayConfiguration, + SecretStore, SecretStoreOperation, SessionState, + }; + use radroots_studio_domain::{ + AccountCreatedAt, AccountIdentity, AccountSummary, BindingAvailability, LocalSignerBinding, + PublicKey, SafeErrorCode, SecretKeyInput, UnixTimestamp, + }; + use tempfile::tempdir; + + use super::PersistentAppCore; + + fn account() -> AccountSummary { + let public_key = PublicKey::from_bytes([4; 32]); + AccountSummary::new( + AccountIdentity::derive(public_key).expect("identity"), + LocalSignerBinding::new(public_key, BindingAvailability::Available), + None, + AccountCreatedAt::new(UnixTimestamp::from_seconds(1).expect("time")), + None, + ) + .expect("account") + } + + struct FixedClock; + + impl Clock for FixedClock { + fn now(&self) -> UnixTimestamp { + UnixTimestamp::from_seconds(25).expect("time") + } + } + + #[test] + fn persistent_bootstrap_handles_fresh_and_existing_signed_out_state() { + let directory = tempdir().expect("directory"); + let path = directory.path().join("studio.sqlite3"); + let public_key = account().public_key(); + let secrets = InMemorySecretStore::default(); + { + let adapter = PersistentAppCore::open(&path, RelayConfiguration::default()) + .expect("open adapter"); + let fresh = adapter + .bootstrap(&secrets, &FixedClock) + .expect("fresh bootstrap"); + assert!(fresh.accounts().is_empty()); + adapter + .database() + .insert_account(&account()) + .expect("account"); + adapter + .database() + .save_selected_account(Some(public_key)) + .expect("selection"); + } + + let adapter = + PersistentAppCore::open(&path, RelayConfiguration::default()).expect("reopen adapter"); + let restored = adapter.bootstrap(&secrets, &FixedClock).expect("restore"); + assert_eq!(restored.lifecycle(), AppLifecycle::Ready); + assert_eq!(restored.accounts().len(), 1); + assert_eq!(restored.selected_account(), Some(public_key)); + assert_eq!(restored.session(), SessionState::SignedOut); + assert!(restored.active_account().is_none()); + } + + #[test] + fn corrupt_database_fails_safely_without_recreation() { + let directory = tempdir().expect("directory"); + let path = directory.path().join("studio.sqlite3"); + fs::write(&path, b"not a sqlite database").expect("corrupt file"); + + let error = PersistentAppCore::open(&path, RelayConfiguration::default()) + .err() + .expect("safe failure"); + assert_eq!(error.code(), SafeErrorCode::StorageCorrupt); + assert_eq!( + fs::read(&path).expect("unchanged file"), + b"not a sqlite database" + ); + } + + #[test] + fn persisted_generate_and_import_survive_restart_without_secret_bytes() { + let directory = tempdir().expect("directory"); + let path = directory.path().join("studio.sqlite3"); + let secrets = InMemorySecretStore::default(); + let selected; + { + let adapter = + PersistentAppCore::open(&path, RelayConfiguration::default()).expect("adapter"); + adapter.bootstrap(&secrets, &FixedClock).expect("bootstrap"); + let generated = adapter + .generate_account(&secrets, &FixedClock) + .expect("generate"); + assert!( + secrets + .contains(generated.account().public_key()) + .expect("generated credential") + ); + let imported = adapter + .import_secret_key( + SecretKeyInput::parse( + "7e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7" + .to_owned(), + ) + .expect("secret"), + &secrets, + &FixedClock, + ) + .expect("import"); + selected = imported.account().public_key(); + assert_eq!(adapter.core().snapshot().accounts().len(), 2); + } + + let bytes = fs::read(&path).expect("database bytes"); + assert!(!bytes.windows(5).any(|value| value == b"nsec1")); + assert!(!bytes.windows(64).any(|value| { + value == b"7e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7" + })); + let reopened = + PersistentAppCore::open(&path, RelayConfiguration::default()).expect("reopen"); + let restored = reopened.bootstrap(&secrets, &FixedClock).expect("restore"); + assert_eq!(restored.accounts().len(), 2); + assert_eq!(restored.selected_account(), Some(selected)); + assert_eq!(restored.session(), SessionState::SignedOut); + } + + #[test] + fn durable_import_commits_each_phase_and_recovers_the_terminal_receipt() { + let adapter = PersistentAppCore::in_memory(RelayConfiguration::default()).expect("adapter"); + let secrets = InMemorySecretStore::default(); + let snapshot = adapter.bootstrap(&secrets, &FixedClock).expect("bootstrap"); + let request = DurableRequestId::parse("import:adapter:1").expect("request"); + let imported = adapter + .import_secret_key_durable( + &request, + snapshot.revision().value(), + SecretKeyInput::parse( + "7e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7".to_owned(), + ) + .expect("secret"), + &secrets, + &FixedClock, + ) + .expect("durable import"); + let operation = adapter + .database() + .load_durable_operation(&request) + .expect("operation") + .expect("durable record"); + let receipt = operation.terminal().expect("terminal receipt"); + assert_eq!(receipt.account(), imported.account().public_key()); + assert_eq!( + receipt.resulting_revision(), + Some(adapter.core().snapshot().revision().value()) + ); + } + + #[test] + fn durable_recovery_preserves_repair_metadata_and_deletes_orphan_credentials() { + let adapter = PersistentAppCore::in_memory(RelayConfiguration::default()).expect("adapter"); + let secrets = InMemorySecretStore::default(); + let missing = account().with_binding_availability(BindingAvailability::CredentialMissing); + adapter + .database() + .insert_account(&missing) + .expect("account"); + adapter + .database() + .save_selected_account(Some(missing.public_key())) + .expect("selection"); + let request = DurableRequestId::parse("repair:recovery:1").expect("request"); + adapter + .database() + .begin_durable_operation( + &request, + DurableOperationKind::Repair, + missing.public_key(), + Some(0), + OperationPriorState::new( + Some(missing.public_key()), + Some(BindingAvailability::CredentialMissing), + ), + FixedClock.now(), + ) + .expect("intent"); + secrets + .put( + missing.public_key(), + SecretKeyInput::parse( + "7e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7".to_owned(), + ) + .expect("secret"), + ) + .expect("credential"); + adapter + .database() + .advance_durable_operation( + &request, + DurableOperationPhase::IntentRecorded, + DurableOperationPhase::CredentialWritten, + FixedClock.now(), + None, + ) + .expect("credential phase"); + + adapter.bootstrap(&secrets, &FixedClock).expect("recovery"); + let repaired = adapter + .database() + .find_account(missing.public_key()) + .expect("lookup") + .expect("preserved account"); + assert_eq!( + repaired.signer().availability(), + BindingAvailability::CredentialMissing + ); + assert!(!secrets.contains(missing.public_key()).expect("credential")); + assert_eq!( + adapter + .database() + .load_durable_operation(&request) + .expect("operation") + .expect("record") + .terminal() + .expect("receipt") + .outcome(), + DurableTerminalOutcome::Failed + ); + } + + #[test] + fn durable_recovery_covers_response_loss_and_irreversible_removal_windows() { + let secrets = InMemorySecretStore::default(); + let adapter = PersistentAppCore::in_memory(RelayConfiguration::default()).expect("adapter"); + let saved = account(); + adapter.database().insert_account(&saved).expect("account"); + let import = DurableRequestId::parse("import:response-loss:1").expect("request"); + adapter + .database() + .begin_durable_operation( + &import, + DurableOperationKind::Import, + saved.public_key(), + Some(0), + OperationPriorState::new(None, None), + FixedClock.now(), + ) + .expect("intent"); + adapter + .database() + .advance_durable_operation( + &import, + DurableOperationPhase::IntentRecorded, + DurableOperationPhase::CredentialWritten, + FixedClock.now(), + None, + ) + .expect("credential"); + adapter + .database() + .advance_durable_operation( + &import, + DurableOperationPhase::CredentialWritten, + DurableOperationPhase::MetadataCommitted, + FixedClock.now(), + None, + ) + .expect("metadata"); + let restored = adapter + .bootstrap(&secrets, &FixedClock) + .expect("response recovery"); + assert_eq!(restored.selected_account(), Some(saved.public_key())); + assert_eq!( + adapter + .database() + .load_durable_operation(&import) + .expect("operation") + .expect("record") + .terminal() + .expect("receipt") + .outcome(), + DurableTerminalOutcome::Completed + ); + + let removal_adapter = + PersistentAppCore::in_memory(RelayConfiguration::default()).expect("remove adapter"); + removal_adapter + .database() + .insert_account(&saved) + .expect("remove account"); + removal_adapter + .database() + .save_selected_account(Some(saved.public_key())) + .expect("remove selection"); + let removal = DurableRequestId::parse("remove:response-loss:1").expect("request"); + removal_adapter + .database() + .begin_durable_operation( + &removal, + DurableOperationKind::Remove, + saved.public_key(), + Some(0), + OperationPriorState::new(None, Some(BindingAvailability::Available)), + FixedClock.now(), + ) + .expect("remove intent"); + removal_adapter + .database() + .advance_durable_operation( + &removal, + DurableOperationPhase::IntentRecorded, + DurableOperationPhase::CredentialDeleted, + FixedClock.now(), + None, + ) + .expect("credential deleted"); + let removed = removal_adapter + .bootstrap(&secrets, &FixedClock) + .expect("removal recovery"); + assert!(removed.accounts().is_empty()); + assert_eq!(removed.selected_account(), None); + } + + #[test] + fn bootstrap_recovery_completes_credential_deleted_removal_and_fallback() { + let directory = tempdir().expect("directory"); + let path = directory.path().join("studio.sqlite3"); + let secrets = InMemorySecretStore::default(); + let first; + let removed; + { + let adapter = + PersistentAppCore::open(&path, RelayConfiguration::default()).expect("adapter"); + adapter.bootstrap(&secrets, &FixedClock).expect("bootstrap"); + first = adapter + .generate_account(&secrets, &FixedClock) + .expect("first") + .account() + .public_key(); + removed = adapter + .generate_account(&secrets, &FixedClock) + .expect("removed") + .account() + .public_key(); + let operation = adapter + .database() + .begin_operation(AccountOperationKind::Remove, removed, FixedClock.now()) + .expect("intent"); + secrets.delete(removed).expect("credential deletion"); + adapter + .database() + .update_operation( + operation, + AccountOperationPhase::CredentialDeleted, + FixedClock.now(), + None, + ) + .expect("phase"); + } + + let reopened = + PersistentAppCore::open(&path, RelayConfiguration::default()).expect("reopen"); + let restored = reopened + .bootstrap(&secrets, &FixedClock) + .expect("recover and bootstrap"); + assert_eq!(restored.accounts().len(), 1); + assert_eq!(restored.selected_account(), Some(first)); + assert_eq!(restored.session(), SessionState::SignedOut); + assert!( + reopened + .database() + .list_pending_operations() + .expect("journal") + .is_empty() + ); + assert!( + reopened + .database() + .find_account(removed) + .expect("removed") + .is_none() + ); + } + + #[test] + fn bootstrap_skips_keyring_when_journal_empty_and_retains_failed_intent() { + let empty = PersistentAppCore::in_memory(RelayConfiguration::default()).expect("empty"); + let unavailable = FailureSecretStore::default(); + unavailable.fail_next(SecretStoreOperation::Delete); + empty + .bootstrap(&unavailable, &FixedClock) + .expect("empty journal does not access keyring"); + + let adapter = PersistentAppCore::in_memory(RelayConfiguration::default()).expect("adapter"); + adapter + .database() + .insert_account(&account()) + .expect("account"); + adapter + .database() + .save_selected_account(Some(account().public_key())) + .expect("selection"); + adapter + .database() + .begin_operation( + AccountOperationKind::Remove, + account().public_key(), + FixedClock.now(), + ) + .expect("intent"); + let failing = FailureSecretStore::default(); + failing.fail_next(SecretStoreOperation::Delete); + let error = adapter + .bootstrap(&failing, &FixedClock) + .expect_err("keyring unavailable"); + assert_eq!(error.code(), SafeErrorCode::KeyringUnavailable); + let pending = adapter + .database() + .list_pending_operations() + .expect("pending"); + assert_eq!(pending.len(), 1); + assert_eq!(pending[0].phase(), AccountOperationPhase::IntentRecorded); + } +} diff --git a/crates/studio_runtime/src/runtime_actor.rs b/crates/studio_runtime/src/runtime_actor.rs @@ -0,0 +1,2002 @@ +use std::collections::BTreeMap; +use std::num::{NonZeroU64, NonZeroUsize}; +use std::path::Path; +use std::sync::atomic::{AtomicU64, Ordering}; +use std::sync::{Arc, Mutex}; +use std::time::{Duration, Instant}; + +use radroots_studio_application::{ + ActorMailbox, AppSnapshot, ChangeSubscriptionId, Clock, CommandContext, CommandEnvelope, + CommandReceipt, CommandResult, CommandSubmission, DurableRequestId, ForegroundSessionBinding, + GenerateAccountReceipt, GeneratedKeyRecoveryHandle, GeneratedKeyStage, ImportAccountReceipt, + LifecycleGate, NostrClient, OrderedSnapshotChanges, ProfileFetchResult, ProfileRefreshPlan, + RecoveryStageId, RelayConfiguration, RemovalConfirmationToken, RequestId, RuntimeCommandClass, + RuntimeLifecycle, SecretStore, SessionGeneration, SnapshotChange, SnapshotChangeReceiver, + SnapshotRevision, StagedGeneratedKey, TaskCorrelation, +}; +use radroots_studio_domain::{ + AccountIdentity, BindingAvailability, LocalSignerBinding, PublicKey, SafeError, SafeErrorCode, + SafeMessage, SecretKeyInput, +}; +use tokio::runtime::Handle; +use tokio::sync::{mpsc, oneshot, watch}; + +use crate::blocking::{BlockingExecutionError, BoundedBlockingExecutor}; +use crate::{InstallationIdentity, InstallationIdentitySource, PersistentAppCore}; + +const DEFAULT_COMMAND_TIMEOUT: Duration = Duration::from_secs(30); +const DEFAULT_TASK_CAPACITY: usize = 64; +const DEFAULT_BLOCKING_CAPACITY: usize = 4; + +enum RuntimeCommand { + Snapshot, + GenerateAccount { + durable_request: DurableRequestId, + expected_revision: u64, + }, + BeginGeneratedKeyStage, + AcknowledgeGeneratedKeyStage { + id: RecoveryStageId, + durable_request: DurableRequestId, + }, + CancelGeneratedKeyStage, + ImportSecretKey { + input: SecretKeyInput, + durable_request: DurableRequestId, + expected_revision: u64, + }, + SelectAccount(PublicKey), + ActivateAccount(PublicKey), + SignOut, + RefreshActiveProfile, + RequestAccountRemoval(PublicKey), + ConfirmAccountRemoval { + token: RemovalConfirmationToken, + durable_request: DurableRequestId, + }, + SubscribeChanges(NonZeroUsize), + UnsubscribeChanges(ChangeSubscriptionId), + Close, +} + +enum RuntimeCommandValue { + Snapshot(Box<AppSnapshot>), + Generated(GenerateAccountReceipt), + GeneratedKeyStage(GeneratedKeyRecoveryHandle), + GeneratedKeyStageCancelled(bool), + Imported(ImportAccountReceipt), + RemovalRequest(RemovalConfirmationToken), + Subscription(RuntimeChangeSubscription), + Unsubscribed(bool), + Closed, +} + +impl RuntimeCommand { + const fn class(&self) -> RuntimeCommandClass { + match self { + Self::Snapshot | Self::SubscribeChanges(_) | Self::UnsubscribeChanges(_) => { + RuntimeCommandClass::Observe + } + Self::GenerateAccount { .. } + | Self::BeginGeneratedKeyStage + | Self::AcknowledgeGeneratedKeyStage { .. } + | Self::ImportSecretKey { .. } + | Self::ActivateAccount(_) + | Self::ConfirmAccountRemoval { .. } => RuntimeCommandClass::UseCredential, + Self::SelectAccount(_) + | Self::SignOut + | Self::RequestAccountRemoval(_) + | Self::CancelGeneratedKeyStage => RuntimeCommandClass::MutateLocalState, + Self::RefreshActiveProfile => RuntimeCommandClass::UseRelay, + Self::Close => RuntimeCommandClass::Shutdown, + } + } + + const fn resolves_revision_through_durable_replay(&self) -> bool { + matches!( + self, + Self::GenerateAccount { .. } | Self::ImportSecretKey { .. } + ) + } +} + +struct RuntimeActor { + adapter: Arc<PersistentAppCore>, + secrets: Arc<dyn SecretStore>, + clock: Arc<dyn Clock>, + nostr: Arc<dyn NostrClient>, + lifecycle: Arc<Mutex<LifecycleGate>>, + runtime: Handle, + blocking: BoundedBlockingExecutor, + session_generation: SessionGeneration, + published_session_generation: Arc<AtomicU64>, + profile_tasks: BTreeMap<RequestId, PendingProfileTask>, + changes: OrderedSnapshotChanges, + published_foreground_session: Arc<Mutex<Option<ForegroundSessionBinding>>>, + generated_key_stage: GeneratedKeyStage, +} + +struct PendingProfileTask { + correlation: TaskCorrelation, + plan: ProfileRefreshPlan, + deadline: Instant, + reply: oneshot::Sender<CommandReceipt<RuntimeCommandValue>>, + handle: tokio::task::JoinHandle<()>, +} + +struct ProfileCompletion { + request_id: RequestId, + result: Result<ProfileFetchResult, SafeError>, +} + +#[derive(Clone)] +pub struct RuntimeActorHandle { + mailbox: ActorMailbox<RuntimeCommand, RuntimeCommandValue>, + adapter: Arc<PersistentAppCore>, + lifecycle: Arc<Mutex<LifecycleGate>>, + next_request: Arc<AtomicU64>, + session_generation: Arc<AtomicU64>, + foreground_session: Arc<Mutex<Option<ForegroundSessionBinding>>>, + installation_identity: InstallationIdentity, + actor_task: Arc<Mutex<Option<tokio::task::JoinHandle<()>>>>, + actor_exit: watch::Receiver<bool>, +} + +#[derive(Clone)] +pub struct RuntimeDependencies { + secrets: Arc<dyn SecretStore>, + clock: Arc<dyn Clock>, + nostr: Arc<dyn NostrClient>, + installation_source: Arc<dyn InstallationIdentitySource>, +} + +impl RuntimeDependencies { + #[must_use] + pub fn new( + secrets: Arc<dyn SecretStore>, + clock: Arc<dyn Clock>, + nostr: Arc<dyn NostrClient>, + installation_source: Arc<dyn InstallationIdentitySource>, + ) -> Self { + Self { + secrets, + clock, + nostr, + installation_source, + } + } +} + +pub struct RuntimeChangeSubscription { + id: ChangeSubscriptionId, + receiver: SnapshotChangeReceiver, +} + +impl RuntimeChangeSubscription { + #[must_use] + pub const fn id(&self) -> ChangeSubscriptionId { + self.id + } + + pub async fn receive(&mut self) -> Option<SnapshotChange> { + self.receiver.receive().await + } +} + +impl RuntimeActorHandle { + /// Opens, migrates, recovers, and starts one actor-owned file-backed runtime. + /// + /// # Errors + /// + /// Returns a safe storage, recovery, or lifecycle error before the actor is + /// published when opening cannot reach ready state. + pub async fn open( + path: &Path, + relay_configuration: RelayConfiguration, + dependencies: RuntimeDependencies, + capacity: NonZeroUsize, + runtime: &Handle, + ) -> Result<Self, SafeError> { + let blocking = BoundedBlockingExecutor::new(DEFAULT_BLOCKING_CAPACITY, runtime); + let path = path.to_path_buf(); + let adapter = blocking + .execute(Instant::now() + DEFAULT_COMMAND_TIMEOUT, move || { + PersistentAppCore::open(&path, relay_configuration) + }) + .await + .map_err(blocking_execution_failed)??; + Self::start(adapter, dependencies, capacity, runtime, blocking).await + } + + /// Starts one isolated actor-owned in-memory runtime for tests. + /// + /// # Errors + /// + /// Returns a safe storage, recovery, or lifecycle error before publication. + pub async fn in_memory( + relay_configuration: RelayConfiguration, + dependencies: RuntimeDependencies, + capacity: NonZeroUsize, + runtime: &Handle, + ) -> Result<Self, SafeError> { + let blocking = BoundedBlockingExecutor::new(DEFAULT_BLOCKING_CAPACITY, runtime); + let adapter = blocking + .execute(Instant::now() + DEFAULT_COMMAND_TIMEOUT, move || { + PersistentAppCore::in_memory(relay_configuration) + }) + .await + .map_err(blocking_execution_failed)??; + Self::start(adapter, dependencies, capacity, runtime, blocking).await + } + + async fn start( + adapter: PersistentAppCore, + dependencies: RuntimeDependencies, + capacity: NonZeroUsize, + runtime: &Handle, + blocking: BoundedBlockingExecutor, + ) -> Result<Self, SafeError> { + let mut gate = LifecycleGate::opening(); + gate.begin_compatibility_check()?; + gate.compatibility_accepted()?; + gate.ownership_acquired()?; + gate.migration_complete()?; + let RuntimeDependencies { + secrets, + clock, + nostr, + installation_source, + } = dependencies; + let adapter = Arc::new(adapter); + let bootstrap_adapter = Arc::clone(&adapter); + let bootstrap_secrets = Arc::clone(&secrets); + let bootstrap_clock = Arc::clone(&clock); + let installation_identity = blocking + .execute(Instant::now() + DEFAULT_COMMAND_TIMEOUT, move || { + bootstrap_adapter + .bootstrap(bootstrap_secrets.as_ref(), bootstrap_clock.as_ref())?; + bootstrap_adapter.initialize_installation_identity(installation_source.as_ref()) + }) + .await + .map_err(blocking_execution_failed)??; + gate.recovery_complete()?; + + let lifecycle = Arc::new(Mutex::new(gate)); + let (mailbox, receiver) = ActorMailbox::bounded(capacity); + let session_generation = Arc::new(AtomicU64::new(SessionGeneration::initial().value())); + let foreground_session = Arc::new(Mutex::new(None)); + let changes = OrderedSnapshotChanges::new(adapter.core().snapshot()); + let actor = RuntimeActor { + adapter: Arc::clone(&adapter), + secrets, + clock, + nostr, + lifecycle: Arc::clone(&lifecycle), + runtime: runtime.clone(), + blocking, + session_generation: SessionGeneration::initial(), + published_session_generation: Arc::clone(&session_generation), + profile_tasks: BTreeMap::new(), + changes, + published_foreground_session: Arc::clone(&foreground_session), + generated_key_stage: GeneratedKeyStage::default(), + }; + let (actor_exit_sender, actor_exit) = watch::channel(false); + let actor_task = runtime.spawn(async move { + actor.run(receiver).await; + let _ = actor_exit_sender.send(true); + }); + Ok(Self { + mailbox, + adapter, + lifecycle, + next_request: Arc::new(AtomicU64::new(1)), + session_generation, + foreground_session, + installation_identity, + actor_task: Arc::new(Mutex::new(Some(actor_task))), + actor_exit, + }) + } + + #[must_use] + pub fn lifecycle(&self) -> RuntimeLifecycle { + self.lifecycle + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .lifecycle() + } + + #[must_use] + pub fn session_generation(&self) -> SessionGeneration { + SessionGeneration::from_value(self.session_generation.load(Ordering::Acquire)) + } + + #[must_use] + pub fn foreground_session(&self) -> Option<ForegroundSessionBinding> { + self.foreground_session + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .clone() + } + + #[must_use] + pub fn installation_identity(&self) -> &InstallationIdentity { + &self.installation_identity + } + + #[must_use] + pub fn snapshot(&self) -> AppSnapshot { + self.adapter.core().snapshot() + } + + /// Returns the ready snapshot through the actor command boundary. + /// + /// # Errors + /// + /// Returns a typed safe actor error. + pub async fn bootstrap(&self) -> Result<AppSnapshot, SafeError> { + Self::expect_snapshot(self.dispatch(RuntimeCommand::Snapshot, None).await?) + } + + /// Generates one account through the serialized actor boundary. + /// + /// # Errors + /// + /// Returns a safe account, storage, keyring, timeout, or actor error. + pub async fn generate_account( + &self, + request: DurableRequestId, + expected_revision: SnapshotRevision, + timeout: Duration, + ) -> Result<GenerateAccountReceipt, SafeError> { + match self + .dispatch_durable( + RuntimeCommand::GenerateAccount { + durable_request: request, + expected_revision: expected_revision.value(), + }, + expected_revision, + timeout, + ) + .await? + { + RuntimeCommandValue::Generated(receipt) => Ok(receipt), + _ => Err(invalid_actor_response()), + } + } + + /// Begins the only actor-owned generated-key recovery stage. + /// + /// # Errors + /// + /// Returns a safe conflict, timeout, key-generation, or actor error. + pub async fn begin_generated_key_stage(&self) -> Result<GeneratedKeyRecoveryHandle, SafeError> { + match self + .dispatch(RuntimeCommand::BeginGeneratedKeyStage, None) + .await? + { + RuntimeCommandValue::GeneratedKeyStage(view) => Ok(view), + _ => Err(invalid_actor_response()), + } + } + + /// Acknowledges recovery and commits the staged account and credential once. + /// + /// # Errors + /// + /// Returns a safe unavailable, conflict, keyring, storage, timeout, or actor error. + pub async fn acknowledge_generated_key_stage( + &self, + id: RecoveryStageId, + request: DurableRequestId, + expected_revision: SnapshotRevision, + timeout: Duration, + ) -> Result<AppSnapshot, SafeError> { + let value = self + .dispatch_durable( + RuntimeCommand::AcknowledgeGeneratedKeyStage { + id, + durable_request: request, + }, + expected_revision, + timeout, + ) + .await?; + Self::expect_snapshot(value) + } + + /// Cancels and zeroizes the active generated-key stage, if present. + /// + /// # Errors + /// + /// Returns a safe timeout or actor error. + pub async fn cancel_generated_key_stage(&self) -> Result<bool, SafeError> { + match self + .dispatch(RuntimeCommand::CancelGeneratedKeyStage, None) + .await? + { + RuntimeCommandValue::GeneratedKeyStageCancelled(cancelled) => Ok(cancelled), + _ => Err(invalid_actor_response()), + } + } + + /// Imports one account through the serialized actor boundary. + /// + /// # Errors + /// + /// Returns a safe account, storage, keyring, timeout, or actor error. + pub async fn import_secret_key( + &self, + request: DurableRequestId, + expected_revision: SnapshotRevision, + input: SecretKeyInput, + timeout: Duration, + ) -> Result<ImportAccountReceipt, SafeError> { + match self + .dispatch_durable( + RuntimeCommand::ImportSecretKey { + input, + durable_request: request, + expected_revision: expected_revision.value(), + }, + expected_revision, + timeout, + ) + .await? + { + RuntimeCommandValue::Imported(receipt) => Ok(receipt), + _ => Err(invalid_actor_response()), + } + } + + /// Selects one account through the serialized actor boundary. + /// + /// # Errors + /// + /// Returns a safe account, storage, timeout, or actor error. + pub async fn select_account(&self, public_key: PublicKey) -> Result<AppSnapshot, SafeError> { + let value = self + .dispatch(RuntimeCommand::SelectAccount(public_key), None) + .await?; + Self::expect_snapshot(value) + } + + /// Activates one account through the serialized actor boundary. + /// + /// # Errors + /// + /// Returns a safe account, credential, storage, timeout, or actor error. + pub async fn activate_account(&self, public_key: PublicKey) -> Result<AppSnapshot, SafeError> { + let value = self + .dispatch(RuntimeCommand::ActivateAccount(public_key), None) + .await?; + Self::expect_snapshot(value) + } + + /// Signs out through the serialized actor boundary. + /// + /// # Errors + /// + /// Returns a safe timeout or actor error. + pub async fn sign_out(&self) -> Result<AppSnapshot, SafeError> { + let value = self.dispatch(RuntimeCommand::SignOut, None).await?; + Self::expect_snapshot(value) + } + + /// Refreshes the active profile through the serialized actor boundary. + /// + /// # Errors + /// + /// Returns a safe relay, storage, timeout, or actor error. + pub async fn refresh_active_profile(&self) -> Result<AppSnapshot, SafeError> { + let value = self + .dispatch(RuntimeCommand::RefreshActiveProfile, None) + .await?; + Self::expect_snapshot(value) + } + + /// Creates one removal request through the serialized actor boundary. + /// + /// # Errors + /// + /// Returns a safe account, timeout, or actor error. + pub async fn request_account_removal( + &self, + public_key: PublicKey, + ) -> Result<RemovalConfirmationToken, SafeError> { + match self + .dispatch(RuntimeCommand::RequestAccountRemoval(public_key), None) + .await? + { + RuntimeCommandValue::RemovalRequest(token) => Ok(token), + _ => Err(invalid_actor_response()), + } + } + + /// Confirms one removal through the serialized actor boundary. + /// + /// # Errors + /// + /// Returns a safe account, credential, storage, timeout, or actor error. + pub async fn confirm_account_removal( + &self, + token: RemovalConfirmationToken, + request: DurableRequestId, + expected_revision: SnapshotRevision, + timeout: Duration, + ) -> Result<AppSnapshot, SafeError> { + let value = self + .dispatch_durable( + RuntimeCommand::ConfirmAccountRemoval { + token, + durable_request: request, + }, + expected_revision, + timeout, + ) + .await?; + Self::expect_snapshot(value) + } + + /// Closes command admission and cancels supervised work. + /// + /// # Errors + /// + /// Returns a safe timeout or actor error. Repeated calls return closed. + pub async fn close(&self) -> Result<(), SafeError> { + self.close_with_timeout(DEFAULT_COMMAND_TIMEOUT).await + } + + /// Closes the runtime within the supplied command deadline. + /// + /// # Errors + /// + /// Returns a safe timeout or actor error. An expired queued close cannot + /// later change runtime state. + pub async fn close_with_timeout(&self, timeout: Duration) -> Result<(), SafeError> { + let deadline = Instant::now() + timeout; + if matches!(self.lifecycle(), RuntimeLifecycle::Closed) { + return self.await_actor_exit(deadline).await; + } + let raw_request = self.next_request.fetch_add(1, Ordering::Relaxed); + let request_id = RequestId::new(raw_request).ok_or_else(request_space_exhausted)?; + let command_result = match self + .dispatch_with_deadline(RuntimeCommand::Close, None, request_id, deadline) + .await + { + Ok(RuntimeCommandValue::Closed) => Ok(()), + Ok(_) => Err(invalid_actor_response()), + Err(error) => Err(error), + }; + let exit_result = self.await_actor_exit(deadline).await; + if matches!(self.lifecycle(), RuntimeLifecycle::Closed) { + exit_result + } else { + command_result.and(exit_result) + } + } + + async fn await_actor_exit(&self, deadline: Instant) -> Result<(), SafeError> { + let mut actor_exit = self.actor_exit.clone(); + if !*actor_exit.borrow() { + let remaining = deadline.saturating_duration_since(Instant::now()); + tokio::time::timeout(remaining, actor_exit.wait_for(|exited| *exited)) + .await + .map_err(|_| command_timed_out())? + .map_err(|_| runtime_closed())?; + } + let actor_task = self + .actor_task + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .take(); + if let Some(actor_task) = actor_task { + let remaining = deadline.saturating_duration_since(Instant::now()); + tokio::time::timeout(remaining, actor_task) + .await + .map_err(|_| command_timed_out())? + .map_err(|_| runtime_closed())?; + } + Ok(()) + } + + /// Atomically registers a bounded ordered change consumer with its initial snapshot. + /// + /// # Errors + /// + /// Returns a safe actor or subscription error. + pub async fn subscribe_changes( + &self, + capacity: NonZeroUsize, + ) -> Result<RuntimeChangeSubscription, SafeError> { + match self + .dispatch(RuntimeCommand::SubscribeChanges(capacity), None) + .await? + { + RuntimeCommandValue::Subscription(subscription) => Ok(subscription), + _ => Err(invalid_actor_response()), + } + } + + /// Removes a change consumer through the serialized actor boundary. + /// + /// # Errors + /// + /// Returns a safe actor error. + pub async fn unsubscribe_changes(&self, id: ChangeSubscriptionId) -> Result<bool, SafeError> { + match self + .dispatch(RuntimeCommand::UnsubscribeChanges(id), None) + .await? + { + RuntimeCommandValue::Unsubscribed(removed) => Ok(removed), + _ => Err(invalid_actor_response()), + } + } + + async fn dispatch( + &self, + command: RuntimeCommand, + expected_revision: Option<SnapshotRevision>, + ) -> Result<RuntimeCommandValue, SafeError> { + let raw_request = self.next_request.fetch_add(1, Ordering::Relaxed); + let request_id = RequestId::new(raw_request).ok_or_else(request_space_exhausted)?; + self.dispatch_with_deadline( + command, + expected_revision, + request_id, + Instant::now() + DEFAULT_COMMAND_TIMEOUT, + ) + .await + } + + async fn dispatch_durable( + &self, + command: RuntimeCommand, + expected_revision: SnapshotRevision, + timeout: Duration, + ) -> Result<RuntimeCommandValue, SafeError> { + let raw_request = self.next_request.fetch_add(1, Ordering::Relaxed); + let request_id = RequestId::new(raw_request).ok_or_else(request_space_exhausted)?; + self.dispatch_with_deadline( + command, + Some(expected_revision), + request_id, + Instant::now() + timeout, + ) + .await + } + + async fn dispatch_with_deadline( + &self, + command: RuntimeCommand, + expected_revision: Option<SnapshotRevision>, + request_id: RequestId, + deadline: Instant, + ) -> Result<RuntimeCommandValue, SafeError> { + let context = CommandContext::new(request_id, expected_revision, deadline); + let receipt = match self.mailbox.submit(context, command) { + CommandSubmission::Accepted(ticket) => { + let remaining = deadline.saturating_duration_since(Instant::now()); + match tokio::time::timeout(remaining, ticket.receipt()).await { + Ok(receipt) => receipt, + Err(_) => CommandReceipt::new(request_id, CommandResult::TimedOut), + } + } + CommandSubmission::Rejected(receipt) => receipt, + }; + match receipt.into_result() { + CommandResult::Completed(value) => Ok(value), + CommandResult::Conflicted { .. } => Err(command_conflicted()), + CommandResult::Rejected(_) => Err(command_rejected()), + CommandResult::TimedOut => Err(command_timed_out()), + CommandResult::Closed => Err(runtime_closed()), + CommandResult::Failed(error) => Err(error), + } + } + + #[cfg(test)] + async fn import_secret_key_test( + &self, + input: SecretKeyInput, + ) -> Result<ImportAccountReceipt, SafeError> { + let request_number = self.next_request.fetch_add(1, Ordering::Relaxed); + self.import_secret_key( + DurableRequestId::parse(format!("test:import:{request_number}"))?, + self.snapshot().revision(), + input, + DEFAULT_COMMAND_TIMEOUT, + ) + .await + } + + #[cfg(test)] + async fn acknowledge_generated_key_stage_test( + &self, + id: RecoveryStageId, + ) -> Result<AppSnapshot, SafeError> { + let request_number = self.next_request.fetch_add(1, Ordering::Relaxed); + self.acknowledge_generated_key_stage( + id, + DurableRequestId::parse(format!("test:generate:{request_number}"))?, + self.snapshot().revision(), + DEFAULT_COMMAND_TIMEOUT, + ) + .await + } + + #[cfg(test)] + async fn confirm_account_removal_test( + &self, + token: RemovalConfirmationToken, + ) -> Result<AppSnapshot, SafeError> { + let request_number = self.next_request.fetch_add(1, Ordering::Relaxed); + self.confirm_account_removal( + token, + DurableRequestId::parse(format!("test:remove:{request_number}"))?, + self.snapshot().revision(), + DEFAULT_COMMAND_TIMEOUT, + ) + .await + } + + #[cfg(test)] + async fn import_secret_key_with_timeout( + &self, + input: SecretKeyInput, + timeout: Duration, + ) -> Result<ImportAccountReceipt, SafeError> { + let raw_request = self.next_request.fetch_add(1, Ordering::Relaxed); + let request_id = RequestId::new(raw_request).ok_or_else(request_space_exhausted)?; + let expected_revision = self.adapter.core().snapshot().revision(); + let durable_request = DurableRequestId::parse(format!("test:timeout:{raw_request}"))?; + match self + .dispatch_with_deadline( + RuntimeCommand::ImportSecretKey { + input, + durable_request, + expected_revision: expected_revision.value(), + }, + Some(expected_revision), + request_id, + Instant::now() + timeout, + ) + .await? + { + RuntimeCommandValue::Imported(receipt) => Ok(receipt), + _ => Err(invalid_actor_response()), + } + } + + fn expect_snapshot(value: RuntimeCommandValue) -> Result<AppSnapshot, SafeError> { + match value { + RuntimeCommandValue::Snapshot(snapshot) => Ok(*snapshot), + _ => Err(invalid_actor_response()), + } + } +} + +impl RuntimeActor { + async fn run( + mut self, + mut receiver: mpsc::Receiver<CommandEnvelope<RuntimeCommand, RuntimeCommandValue>>, + ) { + let (completion_sender, mut completions) = mpsc::channel(DEFAULT_TASK_CAPACITY); + loop { + tokio::select! { + envelope = receiver.recv() => { + let Some(envelope) = envelope else { + break; + }; + if !self.handle_command(envelope, &completion_sender).await { + break; + } + } + completion = completions.recv(), if !self.profile_tasks.is_empty() => { + if let Some(completion) = completion { + self.complete_profile_task(completion).await; + } + } + } + } + self.cancel_profile_tasks(None).await; + } + + async fn handle_command( + &mut self, + envelope: CommandEnvelope<RuntimeCommand, RuntimeCommandValue>, + completion_sender: &mpsc::Sender<ProfileCompletion>, + ) -> bool { + let (context, command, reply) = envelope.into_parts(); + if let Some(result) = self.preflight(context, &command) { + let _ = reply.send(CommandReceipt::new(context.request_id(), result)); + return true; + } + if matches!(command, RuntimeCommand::RefreshActiveProfile) { + self.start_profile_task(context, reply, completion_sender.clone()) + .await; + return true; + } + if matches!(command, RuntimeCommand::Close) { + let result = self.close_actor().await; + let closed = matches!(result, CommandResult::Completed(_)); + let _ = reply.send(CommandReceipt::new(context.request_id(), result)); + return !closed; + } + let changes_session = matches!( + command, + RuntimeCommand::ActivateAccount(_) + | RuntimeCommand::SignOut + | RuntimeCommand::ConfirmAccountRemoval { .. } + ); + let begins_generated_recovery = matches!(&command, RuntimeCommand::BeginGeneratedKeyStage); + let result = self.execute_command(context, command).await; + if begins_generated_recovery && matches!(&result, CommandResult::Completed(_)) { + let snapshot = self.adapter.core().snapshot(); + self.cancel_profile_tasks(Some(&snapshot)).await; + } + if changes_session && matches!(result, CommandResult::Completed(_)) { + self.advance_session_generation().await; + self.synchronize_foreground_session(); + } + if matches!(result, CommandResult::Completed(_)) { + self.changes.publish(self.adapter.core().snapshot()); + } + let _ = reply.send(CommandReceipt::new(context.request_id(), result)); + true + } + + fn preflight( + &self, + context: CommandContext, + command: &RuntimeCommand, + ) -> Option<CommandResult<RuntimeCommandValue>> { + 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(); + if matches!(lifecycle.lifecycle(), RuntimeLifecycle::Closed) { + return Some(CommandResult::Closed); + } + if !lifecycle.allows(command.class()) { + return Some(CommandResult::Failed(command_unavailable())); + } + if self.generated_key_stage.pending().is_some() + && !matches!( + command, + RuntimeCommand::Snapshot + | RuntimeCommand::AcknowledgeGeneratedKeyStage { .. } + | RuntimeCommand::CancelGeneratedKeyStage + | RuntimeCommand::SubscribeChanges(_) + | RuntimeCommand::UnsubscribeChanges(_) + | RuntimeCommand::Close + ) + { + return Some(CommandResult::Failed(generated_recovery_route_active())); + } + let current_revision = self.adapter.core().snapshot().revision(); + if !command.resolves_revision_through_durable_replay() + && context + .expected_revision() + .is_some_and(|expected| expected != current_revision) + { + return Some(CommandResult::Conflicted { current_revision }); + } + None + } + + async fn execute_command( + &mut self, + context: CommandContext, + command: RuntimeCommand, + ) -> CommandResult<RuntimeCommandValue> { + let result = match command { + RuntimeCommand::Snapshot => Ok(RuntimeCommandValue::Snapshot(Box::new( + self.adapter.core().snapshot(), + ))), + RuntimeCommand::GenerateAccount { + durable_request, + expected_revision, + } => { + self.run_blocking(context.deadline(), move |adapter, secrets, clock| { + adapter + .generate_account_durable( + &durable_request, + expected_revision, + secrets.as_ref(), + clock.as_ref(), + ) + .map(RuntimeCommandValue::Generated) + }) + .await + } + RuntimeCommand::BeginGeneratedKeyStage => { + match NonZeroU64::new(context.request_id().get()).map(RecoveryStageId::new) { + Some(stage_id) => self + .generated_key_stage + .begin( + self.adapter.key_material(), + stage_id, + self.adapter.core().snapshot().revision().value(), + self.clock.now(), + ) + .map(RuntimeCommandValue::GeneratedKeyStage), + None => Err(request_space_exhausted()), + } + } + RuntimeCommand::AcknowledgeGeneratedKeyStage { + id, + durable_request, + } => match self.generated_key_stage.take(id, self.clock.now()) { + Ok(staged) => { + self.run_blocking(context.deadline(), move |adapter, secrets, clock| { + commit_generated_key_stage( + adapter.as_ref(), + secrets.as_ref(), + clock.as_ref(), + &durable_request, + staged, + ) + }) + .await + } + Err(error) => Err(error), + }, + RuntimeCommand::CancelGeneratedKeyStage => Ok( + RuntimeCommandValue::GeneratedKeyStageCancelled(self.generated_key_stage.cancel()), + ), + RuntimeCommand::ImportSecretKey { + input, + durable_request, + expected_revision, + } => { + self.run_blocking(context.deadline(), move |adapter, secrets, clock| { + adapter + .import_secret_key_durable( + &durable_request, + expected_revision, + input, + secrets.as_ref(), + clock.as_ref(), + ) + .map(RuntimeCommandValue::Imported) + }) + .await + } + RuntimeCommand::SelectAccount(public_key) => { + self.run_blocking(context.deadline(), move |adapter, _, _| { + adapter + .select_account(public_key) + .map(Box::new) + .map(RuntimeCommandValue::Snapshot) + }) + .await + } + RuntimeCommand::ActivateAccount(public_key) => { + self.run_blocking(context.deadline(), move |adapter, secrets, clock| { + adapter + .activate_account(public_key, secrets.as_ref(), clock.as_ref()) + .map(Box::new) + .map(RuntimeCommandValue::Snapshot) + }) + .await + } + RuntimeCommand::SignOut => self + .adapter + .sign_out() + .map(Box::new) + .map(RuntimeCommandValue::Snapshot), + RuntimeCommand::RequestAccountRemoval(public_key) => self + .adapter + .request_account_removal(public_key, self.clock.as_ref()) + .map(RuntimeCommandValue::RemovalRequest), + RuntimeCommand::ConfirmAccountRemoval { + token, + durable_request, + } => { + self.run_blocking(context.deadline(), move |adapter, secrets, clock| { + adapter + .confirm_account_removal_durable( + &durable_request, + token, + secrets.as_ref(), + clock.as_ref(), + ) + .map(Box::new) + .map(RuntimeCommandValue::Snapshot) + }) + .await + } + RuntimeCommand::SubscribeChanges(capacity) => self + .changes + .subscribe(capacity) + .map(|(id, receiver)| { + RuntimeCommandValue::Subscription(RuntimeChangeSubscription { id, receiver }) + }) + .ok_or_else(observer_registration_failed), + RuntimeCommand::UnsubscribeChanges(id) => Ok(RuntimeCommandValue::Unsubscribed( + self.changes.unsubscribe(id), + )), + RuntimeCommand::Close | RuntimeCommand::RefreshActiveProfile => { + Err(invalid_actor_response()) + } + }; + result.map_or_else(CommandResult::Failed, CommandResult::Completed) + } + + async fn run_blocking<F>( + &self, + deadline: Instant, + operation: F, + ) -> Result<RuntimeCommandValue, SafeError> + where + F: FnOnce( + Arc<PersistentAppCore>, + Arc<dyn SecretStore>, + Arc<dyn Clock>, + ) -> Result<RuntimeCommandValue, SafeError> + + Send + + 'static, + { + let adapter = Arc::clone(&self.adapter); + let secrets = Arc::clone(&self.secrets); + let clock = Arc::clone(&self.clock); + self.blocking + .execute(deadline, move || operation(adapter, secrets, clock)) + .await + .map_err(blocking_execution_failed)? + } + + async fn start_profile_task( + &mut self, + context: CommandContext, + 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 plan = match self.adapter.core().begin_profile_refresh() { + Ok(Some(plan)) => plan, + Ok(None) => { + let _ = reply.send(CommandReceipt::new( + context.request_id(), + CommandResult::Completed(RuntimeCommandValue::Snapshot(Box::new( + self.adapter.core().snapshot(), + ))), + )); + return; + } + Err(error) => { + let _ = reply.send(CommandReceipt::new( + context.request_id(), + CommandResult::Failed(error), + )); + return; + } + }; + let Some(foreground) = foreground.filter(|binding| { + binding.identity().public_key() == plan.public_key() + && binding.generation() == self.session_generation + }) else { + let _ = reply.send(CommandReceipt::new( + context.request_id(), + CommandResult::Failed(stale_profile_binding()), + )); + return; + }; + let correlation = TaskCorrelation::new( + context.request_id(), + plan.public_key(), + foreground.signer(), + plan.expected_revision(), + self.session_generation, + ); + let client = Arc::clone(&self.nostr); + let relays = plan.relays().to_vec(); + let request_id = context.request_id(); + let handle = self.runtime.spawn(async move { + let result = client + .fetch_profile(correlation.account(), &relays, context.deadline()) + .await; + let _ = completion_sender + .send(ProfileCompletion { request_id, result }) + .await; + }); + let previous = self.profile_tasks.insert( + request_id, + PendingProfileTask { + correlation, + plan, + deadline: context.deadline(), + reply, + handle, + }, + ); + if let Some(previous) = previous { + self.lifecycle + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .fail(request_space_exhausted()); + previous.handle.abort(); + let _ = previous.handle.await; + let _ = previous.reply.send(CommandReceipt::new( + previous.correlation.request_id(), + CommandResult::Failed(request_space_exhausted()), + )); + self.cancel_profile_tasks(None).await; + } + } + + async fn close_actor(&mut self) -> CommandResult<RuntimeCommandValue> { + let transition = (|| { + let mut lifecycle = self + .lifecycle + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + lifecycle.begin_shutdown()?; + lifecycle.finish_shutdown() + })(); + match transition { + Ok(()) => { + self.generated_key_stage.cancel(); + self.cancel_profile_tasks(None).await; + self.changes.close(); + *self + .published_foreground_session + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) = None; + CommandResult::Completed(RuntimeCommandValue::Closed) + } + Err(error) => CommandResult::Failed(error), + } + } + + async fn complete_profile_task(&mut self, completion: ProfileCompletion) { + let Some(task) = self.profile_tasks.remove(&completion.request_id) else { + return; + }; + 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 correlated = task.correlation.session_generation() == self.session_generation + && foreground.is_some_and(|binding| { + binding.generation() == task.correlation.session_generation() + && binding.identity().public_key() == task.correlation.account() + && binding.signer() == task.correlation.binding() + }) + && current + .active_account() + .is_some_and(|active| active.account().public_key() == task.correlation.account()); + let result = if correlated { + let plan = task.plan.clone(); + let completed = self + .run_blocking(task.deadline, move |adapter, _, clock| { + adapter + .core() + .complete_profile_refresh( + &plan, + completion.result, + adapter.database(), + clock.as_ref(), + ) + .map(Box::new) + .map(RuntimeCommandValue::Snapshot) + }) + .await; + completed.map_or_else(CommandResult::Failed, CommandResult::Completed) + } else { + CommandResult::Completed(RuntimeCommandValue::Snapshot(Box::new(current))) + }; + if matches!(result, CommandResult::Completed(_)) { + self.changes.publish(self.adapter.core().snapshot()); + } + let _ = task + .reply + .send(CommandReceipt::new(task.correlation.request_id(), result)); + } + + 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.cancel_profile_tasks(None).await; + return; + }; + self.session_generation = next; + self.published_session_generation + .store(next.value(), Ordering::Release); + let snapshot = self.adapter.core().snapshot(); + self.cancel_profile_tasks(Some(&snapshot)).await; + } + + fn synchronize_foreground_session(&mut self) { + let session = self + .adapter + .core() + .snapshot() + .active_account() + .map(|active| { + let public_key = active.account().public_key(); + ForegroundSessionBinding::new( + AccountIdentity::derive(public_key)?, + LocalSignerBinding::new(public_key, BindingAvailability::Available), + self.session_generation, + ) + }); + let session = match session.transpose() { + Ok(session) => session, + Err(error) => { + self.lifecycle + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .fail(error); + None + } + }; + *self + .published_foreground_session + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) = session; + } + + async fn cancel_profile_tasks(&mut self, snapshot: Option<&AppSnapshot>) { + let tasks = std::mem::take(&mut self.profile_tasks); + for (_, task) in tasks { + task.handle.abort(); + let _ = task.handle.await; + let receipt_result = snapshot.map_or(CommandResult::Closed, |snapshot| { + CommandResult::Completed(RuntimeCommandValue::Snapshot(Box::new(snapshot.clone()))) + }); + let _ = task.reply.send(CommandReceipt::new( + task.correlation.request_id(), + receipt_result, + )); + } + } +} + +fn commit_generated_key_stage( + adapter: &PersistentAppCore, + secrets: &dyn SecretStore, + clock: &dyn Clock, + request: &DurableRequestId, + staged: StagedGeneratedKey, +) -> Result<RuntimeCommandValue, SafeError> { + adapter.commit_staged_generated_key(request, staged, secrets, clock)?; + Ok(RuntimeCommandValue::Snapshot(Box::new( + adapter.core().snapshot(), + ))) +} + +const fn blocking_execution_failed(error: BlockingExecutionError) -> SafeError { + match error { + BlockingExecutionError::DeadlineElapsed => command_timed_out(), + BlockingExecutionError::Saturated => command_rejected(), + BlockingExecutionError::TaskFailed => SafeError::new( + SafeErrorCode::InvalidApplicationState, + SafeMessage::new("The runtime blocking worker failed."), + ), + } +} + +const fn request_space_exhausted() -> SafeError { + SafeError::new( + SafeErrorCode::InvalidApplicationState, + SafeMessage::new("The runtime request identifier space is exhausted."), + ) +} + +const fn stale_profile_binding() -> SafeError { + SafeError::new( + SafeErrorCode::InvalidApplicationState, + SafeMessage::new("The active account binding changed before profile refresh."), + ) +} + +const fn command_conflicted() -> SafeError { + SafeError::new( + SafeErrorCode::InvalidApplicationState, + SafeMessage::new("The command conflicts with newer application state."), + ) +} + +const fn command_rejected() -> SafeError { + SafeError::new( + SafeErrorCode::InvalidApplicationState, + SafeMessage::new("The runtime is busy. Try again."), + ) +} + +const fn command_timed_out() -> SafeError { + SafeError::new( + SafeErrorCode::InvalidApplicationState, + SafeMessage::new("The runtime command timed out."), + ) +} + +const fn runtime_closed() -> SafeError { + SafeError::new( + SafeErrorCode::InvalidApplicationState, + SafeMessage::new("The application runtime is closed."), + ) +} + +const fn command_unavailable() -> SafeError { + SafeError::new( + SafeErrorCode::InvalidApplicationState, + SafeMessage::new("The command is unavailable in the current runtime state."), + ) +} + +const fn generated_recovery_route_active() -> SafeError { + SafeError::new( + SafeErrorCode::InvalidApplicationState, + SafeMessage::new("Complete or cancel generated-key recovery before another action."), + ) +} + +const fn invalid_actor_response() -> SafeError { + SafeError::new( + SafeErrorCode::InvalidApplicationState, + SafeMessage::new("The runtime returned an invalid command response."), + ) +} + +const fn observer_registration_failed() -> SafeError { + SafeError::new( + SafeErrorCode::ObserverRegistrationFailed, + SafeMessage::new("The application change subscription could not be registered."), + ) +} + +#[cfg(test)] +mod tests { + use std::num::NonZeroUsize; + use std::sync::atomic::{AtomicBool, Ordering}; + use std::sync::{Arc, Condvar, Mutex}; + use std::time::{Duration, Instant}; + + use radroots_studio_application::{ + BoxFuture, Clock, FailureSecretStore, InMemorySecretStore, NostrClient, ProfileFetchResult, + RelayConfiguration, RuntimeLifecycle, SecretStore, SecretStoreOperation, SessionState, + }; + use radroots_studio_domain::{ + PublicKey, RelayDestinationPolicy, RelayUrl, SafeError, SafeErrorCode, SecretKeyInput, + UnixTimestamp, + }; + + use super::{RuntimeActorHandle, RuntimeDependencies}; + use crate::{InstallationIdentity, InstallationIdentitySource, UuidInstallationIdentitySource}; + + struct FixedInstallationIdentity(&'static str); + + impl InstallationIdentitySource for FixedInstallationIdentity { + fn generate(&self) -> Result<InstallationIdentity, SafeError> { + InstallationIdentity::parse(self.0) + } + } + + fn dependencies( + secrets: Arc<dyn SecretStore>, + nostr: Arc<dyn NostrClient>, + ) -> RuntimeDependencies { + RuntimeDependencies::new( + secrets, + Arc::new(FixedClock), + nostr, + Arc::new(UuidInstallationIdentitySource), + ) + } + + struct FixedClock; + + impl Clock for FixedClock { + fn now(&self) -> UnixTimestamp { + UnixTimestamp::from_seconds(50).expect("time") + } + } + + struct OfflineNostr; + + impl NostrClient for OfflineNostr { + fn fetch_profile<'a>( + &'a self, + _public_key: PublicKey, + _relays: &'a [RelayUrl], + _deadline: Instant, + ) -> BoxFuture<'a, Result<ProfileFetchResult, SafeError>> { + Box::pin(async { Ok(ProfileFetchResult::complete(None)) }) + } + } + + struct BlockingNostr { + started: tokio::sync::Semaphore, + release: tokio::sync::Semaphore, + } + + impl BlockingNostr { + fn new() -> Self { + Self { + started: tokio::sync::Semaphore::new(0), + release: tokio::sync::Semaphore::new(0), + } + } + } + + impl NostrClient for BlockingNostr { + fn fetch_profile<'a>( + &'a self, + _public_key: PublicKey, + _relays: &'a [RelayUrl], + _deadline: Instant, + ) -> BoxFuture<'a, Result<ProfileFetchResult, SafeError>> { + Box::pin(async move { + self.started.add_permits(1); + let permit = self.release.acquire().await.expect("release"); + permit.forget(); + Ok(ProfileFetchResult::complete(None)) + }) + } + } + + struct BlockingSecretStore { + inner: InMemorySecretStore, + block_next_put: AtomicBool, + put_started: AtomicBool, + released: Mutex<bool>, + release_signal: Condvar, + } + + impl BlockingSecretStore { + fn new() -> Self { + Self { + inner: InMemorySecretStore::default(), + block_next_put: AtomicBool::new(true), + put_started: AtomicBool::new(false), + released: Mutex::new(false), + release_signal: Condvar::new(), + } + } + + async fn wait_until_put_started(&self) { + while !self.put_started.load(Ordering::Acquire) { + tokio::task::yield_now().await; + } + } + + fn release(&self) { + *self + .released + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) = true; + self.release_signal.notify_all(); + } + } + + impl SecretStore for BlockingSecretStore { + fn put(&self, public_key: PublicKey, secret: SecretKeyInput) -> Result<(), SafeError> { + if self.block_next_put.swap(false, Ordering::AcqRel) { + self.put_started.store(true, Ordering::Release); + let released = self + .released + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + drop( + self.release_signal + .wait_while(released, |released| !*released) + .unwrap_or_else(std::sync::PoisonError::into_inner), + ); + } + self.inner.put(public_key, secret) + } + + fn load(&self, public_key: PublicKey) -> Result<SecretKeyInput, SafeError> { + self.inner.load(public_key) + } + + fn contains(&self, public_key: PublicKey) -> Result<bool, SafeError> { + self.inner.contains(public_key) + } + + fn delete(&self, public_key: PublicKey) -> Result<(), SafeError> { + self.inner.delete(public_key) + } + } + + async fn actor() -> (RuntimeActorHandle, Arc<InMemorySecretStore>) { + let secrets = Arc::new(InMemorySecretStore::default()); + let secret_port: Arc<dyn SecretStore> = secrets.clone(); + let actor = RuntimeActorHandle::in_memory( + RelayConfiguration::default(), + dependencies(secret_port, Arc::new(OfflineNostr)), + NonZeroUsize::new(8).expect("capacity"), + &tokio::runtime::Handle::current(), + ) + .await + .expect("actor"); + (actor, secrets) + } + + #[tokio::test(flavor = "multi_thread")] + async fn installation_identity_survives_file_backed_runtime_restart() { + let directory = tempfile::tempdir().expect("temporary directory"); + let path = directory.path().join("studio.sqlite3"); + let first = RuntimeActorHandle::open( + &path, + RelayConfiguration::default(), + RuntimeDependencies::new( + Arc::new(InMemorySecretStore::default()), + Arc::new(FixedClock), + Arc::new(OfflineNostr), + Arc::new(FixedInstallationIdentity( + "11aabbccddeeff001122334455667788", + )), + ), + NonZeroUsize::new(8).expect("capacity"), + &tokio::runtime::Handle::current(), + ) + .await + .expect("first runtime"); + assert_eq!( + first.installation_identity().as_str(), + "11aabbccddeeff001122334455667788" + ); + first.close().await.expect("first close"); + drop(first); + + let second = RuntimeActorHandle::open( + &path, + RelayConfiguration::default(), + RuntimeDependencies::new( + Arc::new(InMemorySecretStore::default()), + Arc::new(FixedClock), + Arc::new(OfflineNostr), + Arc::new(FixedInstallationIdentity( + "22aabbccddeeff001122334455667788", + )), + ), + NonZeroUsize::new(8).expect("capacity"), + &tokio::runtime::Handle::current(), + ) + .await + .expect("second runtime"); + assert_eq!( + second.installation_identity().as_str(), + "11aabbccddeeff001122334455667788" + ); + second.close().await.expect("second close"); + } + + #[tokio::test(flavor = "multi_thread")] + async fn account_mutations_run_serially_through_one_ready_actor() { + let (actor, secrets) = actor().await; + assert_eq!(actor.lifecycle(), RuntimeLifecycle::Ready); + + let imported = actor + .import_secret_key_test( + SecretKeyInput::parse( + "7e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7".to_owned(), + ) + .expect("input"), + ) + .await + .expect("import"); + let public_key = imported.account().public_key(); + let activated = actor.activate_account(public_key).await.expect("activate"); + assert_eq!(activated.session(), SessionState::Active); + let foreground = actor.foreground_session().expect("foreground session"); + assert_eq!(foreground.identity().public_key(), public_key); + assert_eq!(foreground.signer().account(), public_key); + assert_eq!(foreground.generation(), actor.session_generation()); + assert!(secrets.contains(public_key).expect("credential")); + + let signed_out = actor.sign_out().await.expect("sign out"); + assert_eq!(signed_out.session(), SessionState::SignedOut); + assert!(actor.foreground_session().is_none()); + let removal = actor + .request_account_removal(public_key) + .await + .expect("removal request"); + let removed = actor + .confirm_account_removal_test(removal) + .await + .expect("remove"); + assert!(removed.accounts().is_empty()); + assert!(!secrets.contains(public_key).expect("credential removed")); + } + + #[tokio::test(flavor = "multi_thread")] + async fn generated_key_stage_is_exclusive_cancelable_and_snapshot_free() { + let (actor, secrets) = actor().await; + let initial = actor.snapshot(); + let stage = actor + .begin_generated_key_stage() + .await + .expect("generated key stage"); + + assert!(actor.begin_generated_key_stage().await.is_err()); + assert_eq!(actor.snapshot(), initial); + assert!( + !secrets + .contains(stage.view().account().public_key()) + .expect("keyring") + ); + assert!(actor.sign_out().await.is_err()); + assert_eq!(actor.snapshot(), initial); + assert!(actor.cancel_generated_key_stage().await.expect("cancel")); + assert!( + !actor + .cancel_generated_key_stage() + .await + .expect("cancel empty") + ); + assert_eq!(actor.snapshot(), initial); + + actor + .begin_generated_key_stage() + .await + .expect("replacement stage"); + actor.close().await.expect("close clears stage"); + assert_eq!(actor.lifecycle(), RuntimeLifecycle::Closed); + } + + #[tokio::test(flavor = "multi_thread")] + async fn recovery_handle_is_one_use_and_acknowledgement_commits_once() { + let (actor, secrets) = actor().await; + let initial = actor.snapshot(); + let handle = actor + .begin_generated_key_stage() + .await + .expect("generated key stage"); + let public_key = handle.view().account().public_key(); + let recovery = handle.take_recovery_nsec().expect("recovery material"); + assert_eq!(recovery.with_exposed_secret(str::len), 63); + assert!(handle.take_recovery_nsec().is_err()); + assert_eq!(actor.snapshot(), initial); + assert!(!secrets.contains(public_key).expect("not committed")); + + let committed = actor + .acknowledge_generated_key_stage_test(handle.id()) + .await + .expect("acknowledge"); + assert_eq!(committed.accounts().len(), 1); + assert_eq!(committed.selected_account(), Some(public_key)); + assert!(secrets.contains(public_key).expect("credential committed")); + assert!( + actor + .acknowledge_generated_key_stage_test(handle.id()) + .await + .is_err() + ); + } + + #[tokio::test(flavor = "multi_thread")] + async fn failed_generated_commit_consumes_the_stage_without_poisoning_the_actor() { + let secrets = Arc::new(FailureSecretStore::default()); + secrets.fail_next(SecretStoreOperation::Put); + let secret_port: Arc<dyn SecretStore> = secrets.clone(); + let actor = RuntimeActorHandle::in_memory( + RelayConfiguration::default(), + dependencies(secret_port, Arc::new(OfflineNostr)), + NonZeroUsize::new(8).expect("capacity"), + &tokio::runtime::Handle::current(), + ) + .await + .expect("actor"); + let handle = actor + .begin_generated_key_stage() + .await + .expect("generated key stage"); + + let error = actor + .acknowledge_generated_key_stage_test(handle.id()) + .await + .expect_err("injected keyring failure"); + + assert_eq!(error.code(), SafeErrorCode::KeyringUnavailable); + assert!(actor.snapshot().accounts().is_empty()); + actor + .begin_generated_key_stage() + .await + .expect("fresh recovery after terminal failure"); + assert!(actor.cancel_generated_key_stage().await.expect("cancel")); + } + + #[tokio::test(flavor = "multi_thread")] + async fn session_generation_cancels_correlated_profile_work_on_sign_out() { + let client = Arc::new(BlockingNostr::new()); + let actor = RuntimeActorHandle::in_memory( + RelayConfiguration::new(vec![ + RelayUrl::parse("ws://localhost:8080", RelayDestinationPolicy::Local) + .expect("relay"), + ]) + .expect("relay configuration"), + dependencies(Arc::new(InMemorySecretStore::default()), client.clone()), + NonZeroUsize::new(8).expect("capacity"), + &tokio::runtime::Handle::current(), + ) + .await + .expect("actor"); + let imported = actor + .import_secret_key_test( + SecretKeyInput::parse( + "7e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7".to_owned(), + ) + .expect("input"), + ) + .await + .expect("import"); + actor + .activate_account(imported.account().public_key()) + .await + .expect("activate"); + assert_eq!(actor.session_generation().value(), 1); + + let refresh_actor = actor.clone(); + let refresh = tokio::spawn(async move { refresh_actor.refresh_active_profile().await }); + let started = client.started.acquire().await.expect("refresh started"); + started.forget(); + let signed_out = actor.sign_out().await.expect("sign out"); + let cancelled = refresh + .await + .expect("refresh task") + .expect("safe cancellation"); + + assert_eq!(actor.session_generation().value(), 2); + assert_eq!(signed_out.session(), SessionState::SignedOut); + assert_eq!(cancelled.session(), SessionState::SignedOut); + assert!(cancelled.active_account().is_none()); + } + + #[tokio::test(flavor = "multi_thread", worker_threads = 4)] + async fn bounded_runtime_rejects_saturation_without_dropping_accepted_commands() { + let secrets = Arc::new(BlockingSecretStore::new()); + let actor = RuntimeActorHandle::in_memory( + RelayConfiguration::default(), + dependencies(secrets.clone(), Arc::new(OfflineNostr)), + NonZeroUsize::new(1).expect("capacity"), + &tokio::runtime::Handle::current(), + ) + .await + .expect("actor"); + + let first_actor = actor.clone(); + let first = tokio::spawn(async move { + first_actor + .import_secret_key_test(secret( + "7e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7", + )) + .await + }); + secrets.wait_until_put_started().await; + + let second_actor = actor.clone(); + let second = tokio::spawn(async move { + second_actor + .import_secret_key_test(secret( + "0000000000000000000000000000000000000000000000000000000000000001", + )) + .await + }); + while actor.mailbox.available_capacity() != 0 { + assert!( + !second.is_finished(), + "second command must enter the mailbox" + ); + tokio::task::yield_now().await; + } + let rejected = actor + .import_secret_key_test(secret( + "0000000000000000000000000000000000000000000000000000000000000002", + )) + .await + .expect_err("full mailbox must reject"); + assert_eq!( + rejected.message().as_str(), + "The runtime is busy. Try again." + ); + + secrets.release(); + first.await.expect("first task").expect("first command"); + let second = second + .await + .expect("second task") + .expect_err("accepted stale revision conflicts explicitly"); + assert_eq!( + second.message().as_str(), + "The account operation conflicts with the current application state." + ); + assert_eq!( + actor.bootstrap().await.expect("snapshot").accounts().len(), + 1 + ); + } + + #[tokio::test(flavor = "multi_thread", worker_threads = 4)] + async fn queued_command_expiry_returns_timeout_and_prevents_late_mutation() { + let secrets = Arc::new(BlockingSecretStore::new()); + let actor = RuntimeActorHandle::in_memory( + RelayConfiguration::default(), + dependencies(secrets.clone(), Arc::new(OfflineNostr)), + NonZeroUsize::new(1).expect("capacity"), + &tokio::runtime::Handle::current(), + ) + .await + .expect("actor"); + + let first_actor = actor.clone(); + let first = tokio::spawn(async move { + first_actor + .import_secret_key_test(secret( + "7e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7", + )) + .await + }); + secrets.wait_until_put_started().await; + + let expired = actor + .import_secret_key_with_timeout( + secret("0000000000000000000000000000000000000000000000000000000000000001"), + Duration::from_millis(10), + ) + .await + .expect_err("queued command must time out"); + assert_eq!(expired.message().as_str(), "The runtime command timed out."); + + secrets.release(); + first.await.expect("first task").expect("first command"); + assert_eq!( + actor.bootstrap().await.expect("snapshot").accounts().len(), + 1 + ); + } + + #[tokio::test(flavor = "multi_thread")] + async fn close_is_terminal_and_every_later_command_is_rejected_as_closed() { + let (actor, _) = actor().await; + actor.close().await.expect("close"); + assert_eq!(actor.lifecycle(), RuntimeLifecycle::Closed); + + let error = actor.bootstrap().await.expect_err("bootstrap after close"); + assert_eq!( + error.message().as_str(), + "The application runtime is closed." + ); + actor.close().await.expect("repeated close is idempotent"); + } + + #[tokio::test(flavor = "multi_thread")] + async fn actor_subscription_atomically_delivers_initial_then_ordered_changes() { + let (actor, _) = actor().await; + let mut subscription = actor + .subscribe_changes(NonZeroUsize::new(4).expect("capacity")) + .await + .expect("subscribe"); + let initial = subscription.receive().await.expect("initial snapshot"); + assert_eq!(initial.revision(), actor.snapshot().revision()); + assert!(initial.previous_revision().is_none()); + + actor + .import_secret_key_test(secret( + "7e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7", + )) + .await + .expect("import"); + let changed = subscription.receive().await.expect("change"); + assert!(changed.revision() > initial.revision()); + assert_eq!(changed.previous_revision(), Some(initial.revision())); + assert!( + actor + .unsubscribe_changes(subscription.id()) + .await + .expect("unsubscribe") + ); + } + + #[tokio::test(flavor = "multi_thread", worker_threads = 4)] + async fn expired_queued_shutdown_does_not_close_runtime_later() { + let secrets = Arc::new(BlockingSecretStore::new()); + let actor = RuntimeActorHandle::in_memory( + RelayConfiguration::default(), + dependencies(secrets.clone(), Arc::new(OfflineNostr)), + NonZeroUsize::new(1).expect("capacity"), + &tokio::runtime::Handle::current(), + ) + .await + .expect("actor"); + let import_actor = actor.clone(); + let import = tokio::spawn(async move { + import_actor + .import_secret_key_test(secret( + "7e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7", + )) + .await + }); + secrets.wait_until_put_started().await; + + let timeout = actor + .close_with_timeout(Duration::from_millis(10)) + .await + .expect_err("queued shutdown must expire"); + assert_eq!(timeout.message().as_str(), "The runtime command timed out."); + secrets.release(); + import.await.expect("import task").expect("import"); + assert_eq!(actor.lifecycle(), RuntimeLifecycle::Ready); + assert_eq!( + actor + .bootstrap() + .await + .expect("still open") + .accounts() + .len(), + 1 + ); + actor.close().await.expect("later close"); + } + + #[tokio::test(flavor = "multi_thread", worker_threads = 4)] + async fn shutdown_cancels_in_flight_work_and_terminates_publication() { + let client = Arc::new(BlockingNostr::new()); + let actor = RuntimeActorHandle::in_memory( + RelayConfiguration::new(vec![ + RelayUrl::parse("ws://localhost:8080", RelayDestinationPolicy::Local) + .expect("relay"), + ]) + .expect("relay configuration"), + dependencies(Arc::new(InMemorySecretStore::default()), client.clone()), + NonZeroUsize::new(8).expect("capacity"), + &tokio::runtime::Handle::current(), + ) + .await + .expect("actor"); + let imported = actor + .import_secret_key_test(secret( + "7e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7", + )) + .await + .expect("import"); + actor + .activate_account(imported.account().public_key()) + .await + .expect("activate"); + let mut changes = actor + .subscribe_changes(NonZeroUsize::new(4).expect("capacity")) + .await + .expect("subscribe"); + changes.receive().await.expect("initial"); + + let refresh_actor = actor.clone(); + let refresh = tokio::spawn(async move { refresh_actor.refresh_active_profile().await }); + let started = client.started.acquire().await.expect("refresh started"); + started.forget(); + actor.close().await.expect("close"); + + let cancelled = refresh + .await + .expect("refresh task") + .expect_err("refresh closes"); + assert_eq!( + cancelled.message().as_str(), + "The application runtime is closed." + ); + assert!(changes.receive().await.is_none()); + assert_eq!(actor.lifecycle(), RuntimeLifecycle::Closed); + } + + fn secret(value: &str) -> SecretKeyInput { + SecretKeyInput::parse(value.to_owned()).expect("valid test secret") + } +} diff --git a/crates/studio_runtime/tests/local_relay_e2e.rs b/crates/studio_runtime/tests/local_relay_e2e.rs @@ -0,0 +1,100 @@ +use std::time::Duration; + +use nostr::{EventBuilder, Keys, Metadata}; +use nostr_relay_builder::MockRelay; +use nostr_sdk::Client; +use radroots_studio_application::{ + Clock, InMemorySecretStore, ProfileLoadState, ProfileRepository, RelayConfiguration, + RelayConnectionState, SecretStore, SessionState, +}; +use radroots_studio_domain::{RelayDestinationPolicy, RelayUrl, SecretKeyInput, UnixTimestamp}; +use radroots_studio_nostr::SdkNostrClient; +use radroots_studio_runtime::PersistentAppCore; + +const SECRET_HEX: &str = "7e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7"; + +struct FixedClock; + +impl Clock for FixedClock { + fn now(&self) -> UnixTimestamp { + UnixTimestamp::from_seconds(100).expect("fixed timestamp") + } +} + +#[tokio::test] +async fn local_relay_e2e_imports_activates_refreshes_and_caches_profile() { + let local_relay = MockRelay::run().await.expect("local relay"); + let relay_url = local_relay.url().await; + let keys = Keys::parse(SECRET_HEX).expect("known secret key"); + let publisher = Client::new(keys); + publisher + .add_relay(relay_url.clone()) + .await + .expect("publisher relay"); + publisher.connect().await; + publisher.wait_for_connection(Duration::from_secs(2)).await; + publisher + .send_event_builder(EventBuilder::metadata( + &Metadata::new() + .name("farmer") + .display_name("Farm Account") + .about("Local food profile"), + )) + .await + .expect("publish profile"); + + let relay = + RelayUrl::parse(relay_url.as_str(), RelayDestinationPolicy::Local).expect("relay URL"); + let adapter = PersistentAppCore::in_memory( + RelayConfiguration::new(vec![relay]).expect("relay configuration"), + ) + .expect("persistent adapter"); + let secrets = InMemorySecretStore::default(); + adapter.bootstrap(&secrets, &FixedClock).expect("bootstrap"); + let imported = adapter + .import_secret_key( + SecretKeyInput::parse(SECRET_HEX.to_owned()).expect("secret input"), + &secrets, + &FixedClock, + ) + .expect("import account"); + let public_key = imported.account().public_key(); + assert!(secrets.contains(public_key).expect("credential exists")); + adapter + .activate_account(public_key, &secrets, &FixedClock) + .expect("activate account"); + + let refreshed = adapter + .core() + .refresh_active_profile( + adapter.database(), + &SdkNostrClient::new(Duration::from_secs(2)), + &FixedClock, + std::time::Instant::now() + Duration::from_secs(2), + ) + .await + .expect("refresh profile"); + + assert_eq!(refreshed.session(), SessionState::Active); + let active = refreshed.active_account().expect("active account"); + assert_eq!(active.relay_state(), RelayConnectionState::Connected); + assert_eq!(active.profile_state(), ProfileLoadState::Fresh); + assert_eq!( + active.profile().and_then(|profile| profile.display_name()), + Some("Farm Account") + ); + let cached = adapter + .database() + .load_profile(public_key) + .expect("load cache") + .expect("cached profile"); + assert_eq!( + cached.candidate().metadata().preferred_name(), + Some("Farm Account") + ); + let public_debug = format!("{refreshed:?}"); + assert!(!public_debug.contains(SECRET_HEX)); + assert!(!public_debug.contains("nsec1")); + publisher.shutdown().await; + local_relay.shutdown(); +} diff --git a/crates/studio_runtime/tests/restart_isolation.rs b/crates/studio_runtime/tests/restart_isolation.rs @@ -0,0 +1,95 @@ +use std::fs; + +use radroots_studio_application::{ + AccountNamespaceRepository, AccountPreferenceKey, Clock, InMemorySecretStore, + RelayConfiguration, SessionState, +}; +use radroots_studio_domain::{SecretKeyInput, UnixTimestamp}; +use radroots_studio_runtime::PersistentAppCore; +use tempfile::tempdir; + +const SECRET_A: &str = "7e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7"; +const SECRET_B: &str = "0101010101010101010101010101010101010101010101010101010101010101"; + +struct FixedClock; + +impl Clock for FixedClock { + fn now(&self) -> UnixTimestamp { + UnixTimestamp::from_seconds(200).expect("fixed timestamp") + } +} + +#[test] +fn restart_restores_selection_and_keeps_account_namespaces_isolated() { + let directory = tempdir().expect("temporary directory"); + let path = directory.path().join("studio.sqlite3"); + let secrets = InMemorySecretStore::default(); + let (owner_a, owner_b); + + { + let adapter = PersistentAppCore::open(&path, RelayConfiguration::default()) + .expect("persistent adapter"); + adapter.bootstrap(&secrets, &FixedClock).expect("bootstrap"); + owner_a = adapter + .import_secret_key( + SecretKeyInput::parse(SECRET_A.to_owned()).expect("secret A"), + &secrets, + &FixedClock, + ) + .expect("account A") + .account() + .public_key(); + owner_b = adapter + .import_secret_key( + SecretKeyInput::parse(SECRET_B.to_owned()).expect("secret B"), + &secrets, + &FixedClock, + ) + .expect("account B") + .account() + .public_key(); + adapter + .database() + .set_value(owner_a, AccountPreferenceKey::NamespaceProbe, "account-a") + .expect("namespace A"); + adapter + .database() + .set_value(owner_b, AccountPreferenceKey::NamespaceProbe, "account-b") + .expect("namespace B"); + adapter.select_account(owner_b).expect("select B"); + } + + let reopened = + PersistentAppCore::open(&path, RelayConfiguration::default()).expect("reopen adapter"); + let restored = reopened.bootstrap(&secrets, &FixedClock).expect("restore"); + assert_eq!(restored.accounts().len(), 2); + assert_eq!(restored.selected_account(), Some(owner_b)); + assert_eq!(restored.session(), SessionState::SignedOut); + assert_eq!( + reopened + .database() + .get_value(owner_a, AccountPreferenceKey::NamespaceProbe) + .expect("read A"), + Some("account-a".to_owned()) + ); + assert_eq!( + reopened + .database() + .get_value(owner_b, AccountPreferenceKey::NamespaceProbe) + .expect("read B"), + Some("account-b".to_owned()) + ); + + let database = fs::read(path).expect("database bytes"); + assert!( + !database + .windows(SECRET_A.len()) + .any(|bytes| bytes == SECRET_A.as_bytes()) + ); + assert!( + !database + .windows(SECRET_B.len()) + .any(|bytes| bytes == SECRET_B.as_bytes()) + ); + assert!(!database.windows(5).any(|bytes| bytes == b"nsec1")); +} diff --git a/crates/studio_storage/Cargo.toml b/crates/studio_storage/Cargo.toml @@ -18,13 +18,9 @@ radroots_studio_application.workspace = true radroots_studio_domain.workspace = true refinery = { version = "=0.9.2", default-features = false, features = ["rusqlite"] } rusqlite = { version = "=0.39.0", features = ["bundled"] } -tokio = { version = "=1.47.1", features = ["macros", "rt-multi-thread", "sync", "time"] } zeroize = "=1.9.0" [dev-dependencies] -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" } tempfile = "=3.23.0" [lints] diff --git a/crates/studio_storage/migrations/V10__installation_identity.sql b/crates/studio_storage/migrations/V10__installation_identity.sql @@ -0,0 +1,7 @@ +CREATE TABLE installation_identity ( + singleton INTEGER PRIMARY KEY CHECK (singleton = 1), + installation_id TEXT NOT NULL CHECK ( + length(installation_id) = 32 + AND installation_id NOT GLOB '*[^0-9a-f]*' + ) +) STRICT; diff --git a/crates/studio_storage/src/application_adapter.rs b/crates/studio_storage/src/application_adapter.rs @@ -1,718 +0,0 @@ -use std::path::Path; - -use radroots_studio_application::{ - AppCore, AppSnapshot, Clock, DurableRequestId, GenerateAccountReceipt, ImportAccountReceipt, - RelayConfiguration, RemovalConfirmationToken, SecretStore, StagedGeneratedKey, -}; -use radroots_studio_domain::{PublicKey, SafeError, SecretKeyInput}; - -use crate::Database; - -pub struct PersistentAppCore { - core: AppCore, - database: Database, -} - -impl PersistentAppCore { - /// Commits an acknowledged generated-key stage through the durable coordinator. - /// - /// # Errors - /// - /// Returns a safe conflict, credential, storage, or recovery error. - pub fn commit_staged_generated_key( - &self, - request_id: &DurableRequestId, - staged: StagedGeneratedKey, - secrets: &(impl SecretStore + ?Sized), - clock: &(impl Clock + ?Sized), - ) -> Result<ImportAccountReceipt, SafeError> { - self.core.commit_staged_generated_key( - request_id, - staged, - &self.database, - &self.database, - secrets, - &self.database, - clock, - ) - } - - /// Opens the application database without accessing credentials or relays. - /// - /// # Errors - /// - /// Returns a safe storage error when the database cannot be opened or migrated. - pub fn open(path: &Path, relay_configuration: RelayConfiguration) -> Result<Self, SafeError> { - Ok(Self { - core: AppCore::in_memory(relay_configuration), - database: Database::open(path)?, - }) - } - - /// Creates an isolated persistent-core adapter for tests. - /// - /// # Errors - /// - /// Returns a safe storage error when the database cannot be initialized. - pub fn in_memory(relay_configuration: RelayConfiguration) -> Result<Self, SafeError> { - Ok(Self { - core: AppCore::in_memory(relay_configuration), - database: Database::in_memory()?, - }) - } - - /// Restores public accounts and selection while keeping the session signed out. - /// - /// # Errors - /// - /// Returns a safe storage or application-state error after publishing a fatal - /// snapshot when durable state cannot be restored. - pub fn bootstrap( - &self, - secrets: &(impl SecretStore + ?Sized), - clock: &(impl Clock + ?Sized), - ) -> Result<AppSnapshot, SafeError> { - self.core.recover_durable_operations( - &self.database, - &self.database, - secrets, - &self.database, - clock, - )?; - self.core.recover_pending_operations( - &self.database, - &self.database, - secrets, - &self.database, - clock, - )?; - self.core.bootstrap_from(&self.database, &self.database) - } - - /// Generates and durably persists one selected, signed-out local account. - /// - /// # Errors - /// - /// Returns a safe credential, storage, key, or application-state error. - pub fn generate_account( - &self, - secrets: &(impl SecretStore + ?Sized), - clock: &(impl Clock + ?Sized), - ) -> Result<GenerateAccountReceipt, SafeError> { - self.core.generate_account( - &self.database, - &self.database, - secrets, - &self.database, - clock, - ) - } - - /// Imports and durably persists one selected, signed-out local account. - /// - /// # Errors - /// - /// Returns a safe credential, storage, key, or application-state error. - pub fn import_secret_key( - &self, - input: SecretKeyInput, - secrets: &(impl SecretStore + ?Sized), - clock: &(impl Clock + ?Sized), - ) -> Result<ImportAccountReceipt, SafeError> { - self.core.import_secret_key( - input, - &self.database, - &self.database, - secrets, - &self.database, - clock, - ) - } - - /// Generates an account through the durable request coordinator. - /// - /// # Errors - /// - /// Returns a safe conflict, credential, storage, or application-state error. - pub fn generate_account_durable( - &self, - request_id: &DurableRequestId, - expected_revision: u64, - secrets: &(impl SecretStore + ?Sized), - clock: &(impl Clock + ?Sized), - ) -> Result<GenerateAccountReceipt, SafeError> { - self.core.generate_account_durable( - request_id, - expected_revision, - &self.database, - &self.database, - secrets, - &self.database, - clock, - ) - } - - /// Imports or repairs an account through the durable request coordinator. - /// - /// # Errors - /// - /// Returns a safe conflict, validation, credential, storage, or state error. - pub fn import_secret_key_durable( - &self, - request_id: &DurableRequestId, - expected_revision: u64, - input: SecretKeyInput, - secrets: &(impl SecretStore + ?Sized), - clock: &(impl Clock + ?Sized), - ) -> Result<ImportAccountReceipt, SafeError> { - self.core.import_secret_key_durable( - request_id, - expected_revision, - input, - &self.database, - &self.database, - secrets, - &self.database, - clock, - ) - } - - /// Persists and publishes one saved-account selection without activation. - /// - /// # Errors - /// - /// Returns a safe account, storage, or application-state error. - pub fn select_account(&self, public_key: PublicKey) -> Result<AppSnapshot, SafeError> { - self.core - .select_account(public_key, &self.database, &self.database) - } - - /// Activates a saved account after validating its credential and cached profile. - /// - /// # Errors - /// - /// Returns a safe account, credential, storage, or application-state error. - pub fn activate_account( - &self, - public_key: PublicKey, - secrets: &(impl SecretStore + ?Sized), - clock: &(impl Clock + ?Sized), - ) -> Result<AppSnapshot, SafeError> { - self.core.activate_account( - public_key, - &self.database, - &self.database, - &self.database, - secrets, - clock, - ) - } - - /// Signs out while retaining durable account data and credentials. - /// - /// # Errors - /// - /// Returns a safe application-state error if sign out cannot complete. - pub fn sign_out(&self) -> Result<AppSnapshot, SafeError> { - self.core.sign_out() - } - - /// Issues a revision-bound, single-use account-removal confirmation. - /// - /// # Errors - /// - /// Returns a safe error when the target account is not saved. - pub fn request_account_removal( - &self, - public_key: PublicKey, - clock: &(impl Clock + ?Sized), - ) -> Result<RemovalConfirmationToken, SafeError> { - self.core.request_account_removal(public_key, clock) - } - - /// Permanently removes one confirmed account and its credential. - /// - /// # Errors - /// - /// Returns a safe confirmation, credential, storage, recovery, or state error. - pub fn confirm_account_removal( - &self, - token: RemovalConfirmationToken, - secrets: &(impl SecretStore + ?Sized), - clock: &(impl Clock + ?Sized), - ) -> Result<AppSnapshot, SafeError> { - self.core.confirm_account_removal( - token, - &self.database, - &self.database, - secrets, - &self.database, - clock, - ) - } - - /// Executes a confirmed removal through the durable request coordinator. - /// - /// # Errors - /// - /// Returns a safe expiry, conflict, credential, storage, recovery, or state error. - pub fn confirm_account_removal_durable( - &self, - request_id: &DurableRequestId, - token: RemovalConfirmationToken, - secrets: &(impl SecretStore + ?Sized), - clock: &(impl Clock + ?Sized), - ) -> Result<AppSnapshot, SafeError> { - self.core.confirm_account_removal_durable( - request_id, - token, - &self.database, - &self.database, - secrets, - &self.database, - clock, - ) - } - - #[must_use] - pub const fn core(&self) -> &AppCore { - &self.core - } - - #[must_use] - pub const fn database(&self) -> &Database { - &self.database - } -} - -#[cfg(test)] -mod tests { - use std::fs; - - use radroots_studio_application::{ - AccountOperationKind, AccountOperationPhase, AccountRepository, AppLifecycle, - AppStateRepository, Clock, DurableOperationKind, DurableOperationPhase, - DurableOperationRepository, DurableRequestId, DurableTerminalOutcome, FailureSecretStore, - InMemorySecretStore, OperationJournal, OperationPriorState, RelayConfiguration, - SecretStore, SecretStoreOperation, SessionState, - }; - use radroots_studio_domain::{ - AccountCreatedAt, AccountIdentity, AccountSummary, BindingAvailability, LocalSignerBinding, - PublicKey, SafeErrorCode, SecretKeyInput, UnixTimestamp, - }; - use tempfile::tempdir; - - use super::PersistentAppCore; - - fn account() -> AccountSummary { - let public_key = PublicKey::from_bytes([4; 32]); - AccountSummary::new( - AccountIdentity::derive(public_key).expect("identity"), - LocalSignerBinding::new(public_key, BindingAvailability::Available), - None, - AccountCreatedAt::new(UnixTimestamp::from_seconds(1).expect("time")), - None, - ) - .expect("account") - } - - struct FixedClock; - - impl Clock for FixedClock { - fn now(&self) -> UnixTimestamp { - UnixTimestamp::from_seconds(25).expect("time") - } - } - - #[test] - fn persistent_bootstrap_handles_fresh_and_existing_signed_out_state() { - let directory = tempdir().expect("directory"); - let path = directory.path().join("studio.sqlite3"); - let public_key = account().public_key(); - let secrets = InMemorySecretStore::default(); - { - let adapter = PersistentAppCore::open(&path, RelayConfiguration::default()) - .expect("open adapter"); - let fresh = adapter - .bootstrap(&secrets, &FixedClock) - .expect("fresh bootstrap"); - assert!(fresh.accounts().is_empty()); - adapter - .database() - .insert_account(&account()) - .expect("account"); - adapter - .database() - .save_selected_account(Some(public_key)) - .expect("selection"); - } - - let adapter = - PersistentAppCore::open(&path, RelayConfiguration::default()).expect("reopen adapter"); - let restored = adapter.bootstrap(&secrets, &FixedClock).expect("restore"); - assert_eq!(restored.lifecycle(), AppLifecycle::Ready); - assert_eq!(restored.accounts().len(), 1); - assert_eq!(restored.selected_account(), Some(public_key)); - assert_eq!(restored.session(), SessionState::SignedOut); - assert!(restored.active_account().is_none()); - } - - #[test] - fn corrupt_database_fails_safely_without_recreation() { - let directory = tempdir().expect("directory"); - let path = directory.path().join("studio.sqlite3"); - fs::write(&path, b"not a sqlite database").expect("corrupt file"); - - let error = PersistentAppCore::open(&path, RelayConfiguration::default()) - .err() - .expect("safe failure"); - assert_eq!(error.code(), SafeErrorCode::StorageCorrupt); - assert_eq!( - fs::read(&path).expect("unchanged file"), - b"not a sqlite database" - ); - } - - #[test] - fn persisted_generate_and_import_survive_restart_without_secret_bytes() { - let directory = tempdir().expect("directory"); - let path = directory.path().join("studio.sqlite3"); - let secrets = InMemorySecretStore::default(); - let selected; - { - let adapter = - PersistentAppCore::open(&path, RelayConfiguration::default()).expect("adapter"); - adapter.bootstrap(&secrets, &FixedClock).expect("bootstrap"); - let generated = adapter - .generate_account(&secrets, &FixedClock) - .expect("generate"); - assert!( - secrets - .contains(generated.account().public_key()) - .expect("generated credential") - ); - let imported = adapter - .import_secret_key( - SecretKeyInput::parse( - "7e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7" - .to_owned(), - ) - .expect("secret"), - &secrets, - &FixedClock, - ) - .expect("import"); - selected = imported.account().public_key(); - assert_eq!(adapter.core().snapshot().accounts().len(), 2); - } - - let bytes = fs::read(&path).expect("database bytes"); - assert!(!bytes.windows(5).any(|value| value == b"nsec1")); - assert!(!bytes.windows(64).any(|value| { - value == b"7e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7" - })); - let reopened = - PersistentAppCore::open(&path, RelayConfiguration::default()).expect("reopen"); - let restored = reopened.bootstrap(&secrets, &FixedClock).expect("restore"); - assert_eq!(restored.accounts().len(), 2); - assert_eq!(restored.selected_account(), Some(selected)); - assert_eq!(restored.session(), SessionState::SignedOut); - } - - #[test] - fn durable_import_commits_each_phase_and_recovers_the_terminal_receipt() { - let adapter = PersistentAppCore::in_memory(RelayConfiguration::default()).expect("adapter"); - let secrets = InMemorySecretStore::default(); - let snapshot = adapter.bootstrap(&secrets, &FixedClock).expect("bootstrap"); - let request = DurableRequestId::parse("import:adapter:1").expect("request"); - let imported = adapter - .import_secret_key_durable( - &request, - snapshot.revision().value(), - SecretKeyInput::parse( - "7e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7".to_owned(), - ) - .expect("secret"), - &secrets, - &FixedClock, - ) - .expect("durable import"); - let operation = adapter - .database() - .load_durable_operation(&request) - .expect("operation") - .expect("durable record"); - let receipt = operation.terminal().expect("terminal receipt"); - assert_eq!(receipt.account(), imported.account().public_key()); - assert_eq!( - receipt.resulting_revision(), - Some(adapter.core().snapshot().revision().value()) - ); - } - - #[test] - fn durable_recovery_preserves_repair_metadata_and_deletes_orphan_credentials() { - let adapter = PersistentAppCore::in_memory(RelayConfiguration::default()).expect("adapter"); - let secrets = InMemorySecretStore::default(); - let missing = account().with_binding_availability(BindingAvailability::CredentialMissing); - adapter - .database() - .insert_account(&missing) - .expect("account"); - adapter - .database() - .save_selected_account(Some(missing.public_key())) - .expect("selection"); - let request = DurableRequestId::parse("repair:recovery:1").expect("request"); - adapter - .database() - .begin_durable_operation( - &request, - DurableOperationKind::Repair, - missing.public_key(), - Some(0), - OperationPriorState::new( - Some(missing.public_key()), - Some(BindingAvailability::CredentialMissing), - ), - FixedClock.now(), - ) - .expect("intent"); - secrets - .put( - missing.public_key(), - SecretKeyInput::parse( - "7e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7".to_owned(), - ) - .expect("secret"), - ) - .expect("credential"); - adapter - .database() - .advance_durable_operation( - &request, - DurableOperationPhase::IntentRecorded, - DurableOperationPhase::CredentialWritten, - FixedClock.now(), - None, - ) - .expect("credential phase"); - - adapter.bootstrap(&secrets, &FixedClock).expect("recovery"); - let repaired = adapter - .database() - .find_account(missing.public_key()) - .expect("lookup") - .expect("preserved account"); - assert_eq!( - repaired.signer().availability(), - BindingAvailability::CredentialMissing - ); - assert!(!secrets.contains(missing.public_key()).expect("credential")); - assert_eq!( - adapter - .database() - .load_durable_operation(&request) - .expect("operation") - .expect("record") - .terminal() - .expect("receipt") - .outcome(), - DurableTerminalOutcome::Failed - ); - } - - #[test] - fn durable_recovery_covers_response_loss_and_irreversible_removal_windows() { - let secrets = InMemorySecretStore::default(); - let adapter = PersistentAppCore::in_memory(RelayConfiguration::default()).expect("adapter"); - let saved = account(); - adapter.database().insert_account(&saved).expect("account"); - let import = DurableRequestId::parse("import:response-loss:1").expect("request"); - adapter - .database() - .begin_durable_operation( - &import, - DurableOperationKind::Import, - saved.public_key(), - Some(0), - OperationPriorState::new(None, None), - FixedClock.now(), - ) - .expect("intent"); - adapter - .database() - .advance_durable_operation( - &import, - DurableOperationPhase::IntentRecorded, - DurableOperationPhase::CredentialWritten, - FixedClock.now(), - None, - ) - .expect("credential"); - adapter - .database() - .advance_durable_operation( - &import, - DurableOperationPhase::CredentialWritten, - DurableOperationPhase::MetadataCommitted, - FixedClock.now(), - None, - ) - .expect("metadata"); - let restored = adapter - .bootstrap(&secrets, &FixedClock) - .expect("response recovery"); - assert_eq!(restored.selected_account(), Some(saved.public_key())); - assert_eq!( - adapter - .database() - .load_durable_operation(&import) - .expect("operation") - .expect("record") - .terminal() - .expect("receipt") - .outcome(), - DurableTerminalOutcome::Completed - ); - - let removal_adapter = - PersistentAppCore::in_memory(RelayConfiguration::default()).expect("remove adapter"); - removal_adapter - .database() - .insert_account(&saved) - .expect("remove account"); - removal_adapter - .database() - .save_selected_account(Some(saved.public_key())) - .expect("remove selection"); - let removal = DurableRequestId::parse("remove:response-loss:1").expect("request"); - removal_adapter - .database() - .begin_durable_operation( - &removal, - DurableOperationKind::Remove, - saved.public_key(), - Some(0), - OperationPriorState::new(None, Some(BindingAvailability::Available)), - FixedClock.now(), - ) - .expect("remove intent"); - removal_adapter - .database() - .advance_durable_operation( - &removal, - DurableOperationPhase::IntentRecorded, - DurableOperationPhase::CredentialDeleted, - FixedClock.now(), - None, - ) - .expect("credential deleted"); - let removed = removal_adapter - .bootstrap(&secrets, &FixedClock) - .expect("removal recovery"); - assert!(removed.accounts().is_empty()); - assert_eq!(removed.selected_account(), None); - } - - #[test] - fn bootstrap_recovery_completes_credential_deleted_removal_and_fallback() { - let directory = tempdir().expect("directory"); - let path = directory.path().join("studio.sqlite3"); - let secrets = InMemorySecretStore::default(); - let first; - let removed; - { - let adapter = - PersistentAppCore::open(&path, RelayConfiguration::default()).expect("adapter"); - adapter.bootstrap(&secrets, &FixedClock).expect("bootstrap"); - first = adapter - .generate_account(&secrets, &FixedClock) - .expect("first") - .account() - .public_key(); - removed = adapter - .generate_account(&secrets, &FixedClock) - .expect("removed") - .account() - .public_key(); - let operation = adapter - .database() - .begin_operation(AccountOperationKind::Remove, removed, FixedClock.now()) - .expect("intent"); - secrets.delete(removed).expect("credential deletion"); - adapter - .database() - .update_operation( - operation, - AccountOperationPhase::CredentialDeleted, - FixedClock.now(), - None, - ) - .expect("phase"); - } - - let reopened = - PersistentAppCore::open(&path, RelayConfiguration::default()).expect("reopen"); - let restored = reopened - .bootstrap(&secrets, &FixedClock) - .expect("recover and bootstrap"); - assert_eq!(restored.accounts().len(), 1); - assert_eq!(restored.selected_account(), Some(first)); - assert_eq!(restored.session(), SessionState::SignedOut); - assert!( - reopened - .database() - .list_pending_operations() - .expect("journal") - .is_empty() - ); - assert!( - reopened - .database() - .find_account(removed) - .expect("removed") - .is_none() - ); - } - - #[test] - fn bootstrap_skips_keyring_when_journal_empty_and_retains_failed_intent() { - let empty = PersistentAppCore::in_memory(RelayConfiguration::default()).expect("empty"); - let unavailable = FailureSecretStore::default(); - unavailable.fail_next(SecretStoreOperation::Delete); - empty - .bootstrap(&unavailable, &FixedClock) - .expect("empty journal does not access keyring"); - - let adapter = PersistentAppCore::in_memory(RelayConfiguration::default()).expect("adapter"); - adapter - .database() - .insert_account(&account()) - .expect("account"); - adapter - .database() - .save_selected_account(Some(account().public_key())) - .expect("selection"); - adapter - .database() - .begin_operation( - AccountOperationKind::Remove, - account().public_key(), - FixedClock.now(), - ) - .expect("intent"); - let failing = FailureSecretStore::default(); - failing.fail_next(SecretStoreOperation::Delete); - let error = adapter - .bootstrap(&failing, &FixedClock) - .expect_err("keyring unavailable"); - assert_eq!(error.code(), SafeErrorCode::KeyringUnavailable); - let pending = adapter - .database() - .list_pending_operations() - .expect("pending"); - assert_eq!(pending.len(), 1); - assert_eq!(pending[0].phase(), AccountOperationPhase::IntentRecorded); - } -} diff --git a/crates/studio_storage/src/db.rs b/crates/studio_storage/src/db.rs @@ -9,7 +9,7 @@ use radroots_studio_domain::{AccountIdentity, PublicKey, SafeError, SafeErrorCod use refinery::embed_migrations; use rusqlite::{Connection, OpenFlags}; -pub const CURRENT_SCHEMA_VERSION: u32 = 9; +pub const CURRENT_SCHEMA_VERSION: u32 = 10; mod migrations { use super::embed_migrations; @@ -383,7 +383,7 @@ mod tests { } let database = Database::open(&path).expect("migrated database"); - assert_eq!(database.schema_version().expect("version"), 9); + assert_eq!(database.schema_version().expect("version"), 10); assert_eq!(database.list_accounts().expect("accounts").len(), 1); assert_eq!( database.load_selected_account().expect("selection"), diff --git a/crates/studio_storage/src/installation.rs b/crates/studio_storage/src/installation.rs @@ -0,0 +1,71 @@ +use radroots_studio_domain::{SafeError, SafeErrorCode, SafeMessage}; +use rusqlite::OptionalExtension; + +use crate::Database; + +impl Database { + pub fn load_installation_id(&self) -> Result<Option<String>, SafeError> { + self.connection() + .query_row( + "SELECT installation_id FROM installation_identity WHERE singleton = 1", + [], + |row| row.get(0), + ) + .optional() + .map_err(|_| installation_storage_error()) + } + + pub fn initialize_installation_id(&self, candidate: &str) -> Result<String, SafeError> { + let connection = self.connection(); + connection + .execute( + "INSERT INTO installation_identity (singleton, installation_id) VALUES (1, ?1) ON CONFLICT(singleton) DO NOTHING", + [candidate], + ) + .map_err(|_| installation_storage_error())?; + connection + .query_row( + "SELECT installation_id FROM installation_identity WHERE singleton = 1", + [], + |row| row.get(0), + ) + .map_err(|_| installation_storage_error()) + } +} + +const fn installation_storage_error() -> SafeError { + SafeError::new( + SafeErrorCode::StorageUnavailable, + SafeMessage::new("The installation identity is unavailable."), + ) +} + +#[cfg(test)] +mod tests { + use crate::Database; + + #[test] + fn installation_identity_is_insert_once_and_stable() { + let database = Database::in_memory().expect("database"); + assert_eq!(database.load_installation_id().expect("empty"), None); + let first = database + .initialize_installation_id("11aabbccddeeff001122334455667788") + .expect("first identity"); + let second = database + .initialize_installation_id("22aabbccddeeff001122334455667788") + .expect("existing identity"); + assert_eq!(first, "11aabbccddeeff001122334455667788"); + assert_eq!(second, first); + assert_eq!(database.load_installation_id().expect("load"), Some(first)); + } + + #[test] + fn installation_identity_rejects_invalid_values() { + let database = Database::in_memory().expect("database"); + assert!( + database + .initialize_installation_id("not-an-identity") + .is_err() + ); + } +} diff --git a/crates/studio_storage/src/lib.rs b/crates/studio_storage/src/lib.rs @@ -2,14 +2,11 @@ pub mod account_namespace; pub mod accounts; -pub mod application_adapter; pub mod db; +mod installation; pub mod journal; pub mod os_keyring; pub mod profiles; -pub mod runtime_actor; -pub use application_adapter::PersistentAppCore; pub use db::{CURRENT_SCHEMA_VERSION, Database}; pub use os_keyring::{CREDENTIAL_SERVICE, OsKeyringSecretStore}; -pub use runtime_actor::RuntimeActorHandle; diff --git a/crates/studio_storage/src/runtime_actor.rs b/crates/studio_storage/src/runtime_actor.rs @@ -1,1693 +0,0 @@ -use std::collections::BTreeMap; -use std::num::{NonZeroU64, NonZeroUsize}; -use std::path::Path; -use std::sync::atomic::{AtomicU64, Ordering}; -use std::sync::{Arc, Mutex}; -use std::time::{Duration, Instant}; - -use radroots_studio_application::{ - ActorMailbox, AppSnapshot, ChangeSubscriptionId, Clock, CommandContext, CommandEnvelope, - CommandReceipt, CommandResult, CommandSubmission, ForegroundSessionBinding, - GenerateAccountReceipt, GeneratedKeyRecoveryHandle, GeneratedKeyStage, ImportAccountReceipt, - LifecycleGate, NostrClient, OrderedSnapshotChanges, ProfileRefreshPlan, RecoveryStageId, - RelayConfiguration, RemovalConfirmationToken, RequestId, RuntimeCommandClass, RuntimeLifecycle, - SecretStore, SessionGeneration, SnapshotChange, SnapshotChangeReceiver, SnapshotRevision, - TaskCorrelation, -}; -use radroots_studio_domain::{ - AccountIdentity, BindingAvailability, Kind0ProfileCandidate, LocalSignerBinding, PublicKey, - SafeError, SafeErrorCode, SafeMessage, SecretKeyInput, -}; -use tokio::runtime::Handle; -use tokio::sync::{mpsc, oneshot}; - -use crate::PersistentAppCore; - -const DEFAULT_COMMAND_TIMEOUT: Duration = Duration::from_secs(30); -const DEFAULT_TASK_CAPACITY: usize = 64; - -enum RuntimeCommand { - Snapshot, - GenerateAccount, - BeginGeneratedKeyStage, - AcknowledgeGeneratedKeyStage(RecoveryStageId), - CancelGeneratedKeyStage, - ImportSecretKey { - input: SecretKeyInput, - durable_request: Option<radroots_studio_application::DurableRequestId>, - durable_expected_revision: Option<u64>, - }, - SelectAccount(PublicKey), - ActivateAccount(PublicKey), - SignOut, - RefreshActiveProfile, - RequestAccountRemoval(PublicKey), - ConfirmAccountRemoval(RemovalConfirmationToken), - SubscribeChanges(NonZeroUsize), - UnsubscribeChanges(ChangeSubscriptionId), - Close, -} - -enum RuntimeCommandValue { - Snapshot(Box<AppSnapshot>), - Generated(GenerateAccountReceipt), - GeneratedKeyStage(GeneratedKeyRecoveryHandle), - GeneratedKeyStageCancelled(bool), - Imported(ImportAccountReceipt), - RemovalRequest(RemovalConfirmationToken), - Subscription(RuntimeChangeSubscription), - Unsubscribed(bool), - Closed, -} - -impl RuntimeCommand { - const fn class(&self) -> RuntimeCommandClass { - match self { - Self::Snapshot | Self::SubscribeChanges(_) | Self::UnsubscribeChanges(_) => { - RuntimeCommandClass::Observe - } - Self::GenerateAccount - | Self::BeginGeneratedKeyStage - | Self::AcknowledgeGeneratedKeyStage(_) - | Self::ImportSecretKey { .. } - | Self::ActivateAccount(_) - | Self::ConfirmAccountRemoval(_) => RuntimeCommandClass::UseCredential, - Self::SelectAccount(_) - | Self::SignOut - | Self::RequestAccountRemoval(_) - | Self::CancelGeneratedKeyStage => RuntimeCommandClass::MutateLocalState, - Self::RefreshActiveProfile => RuntimeCommandClass::UseRelay, - Self::Close => RuntimeCommandClass::Shutdown, - } - } -} - -struct RuntimeActor { - adapter: Arc<PersistentAppCore>, - secrets: Arc<dyn SecretStore>, - clock: Arc<dyn Clock>, - nostr: Arc<dyn NostrClient>, - lifecycle: Arc<Mutex<LifecycleGate>>, - runtime: Handle, - session_generation: SessionGeneration, - published_session_generation: Arc<AtomicU64>, - profile_tasks: BTreeMap<RequestId, PendingProfileTask>, - changes: OrderedSnapshotChanges, - published_foreground_session: Arc<Mutex<Option<ForegroundSessionBinding>>>, - durable_request_namespace: String, - generated_key_stage: GeneratedKeyStage, -} - -struct PendingProfileTask { - correlation: TaskCorrelation, - plan: ProfileRefreshPlan, - reply: oneshot::Sender<CommandReceipt<RuntimeCommandValue>>, - handle: tokio::task::JoinHandle<()>, -} - -struct ProfileCompletion { - request_id: RequestId, - result: Result<Option<Kind0ProfileCandidate>, SafeError>, -} - -#[derive(Clone)] -pub struct RuntimeActorHandle { - mailbox: ActorMailbox<RuntimeCommand, RuntimeCommandValue>, - adapter: Arc<PersistentAppCore>, - lifecycle: Arc<Mutex<LifecycleGate>>, - runtime: Handle, - next_request: Arc<AtomicU64>, - session_generation: Arc<AtomicU64>, - foreground_session: Arc<Mutex<Option<ForegroundSessionBinding>>>, -} - -pub struct RuntimeChangeSubscription { - id: ChangeSubscriptionId, - receiver: SnapshotChangeReceiver, -} - -impl RuntimeChangeSubscription { - #[must_use] - pub const fn id(&self) -> ChangeSubscriptionId { - self.id - } - - pub async fn receive(&mut self) -> Option<SnapshotChange> { - self.receiver.receive().await - } -} - -impl RuntimeActorHandle { - /// Opens, migrates, recovers, and starts one actor-owned file-backed runtime. - /// - /// # Errors - /// - /// Returns a safe storage, recovery, or lifecycle error before the actor is - /// published when opening cannot reach ready state. - pub fn open( - path: &Path, - relay_configuration: RelayConfiguration, - secrets: Arc<dyn SecretStore>, - clock: Arc<dyn Clock>, - nostr: Arc<dyn NostrClient>, - capacity: NonZeroUsize, - runtime: &Handle, - ) -> Result<Self, SafeError> { - Self::start( - PersistentAppCore::open(path, relay_configuration)?, - secrets, - clock, - nostr, - capacity, - runtime, - ) - } - - /// Starts one isolated actor-owned in-memory runtime for tests. - /// - /// # Errors - /// - /// Returns a safe storage, recovery, or lifecycle error before publication. - pub fn in_memory( - relay_configuration: RelayConfiguration, - secrets: Arc<dyn SecretStore>, - clock: Arc<dyn Clock>, - nostr: Arc<dyn NostrClient>, - capacity: NonZeroUsize, - runtime: &Handle, - ) -> Result<Self, SafeError> { - Self::start( - PersistentAppCore::in_memory(relay_configuration)?, - secrets, - clock, - nostr, - capacity, - runtime, - ) - } - - fn start( - adapter: PersistentAppCore, - secrets: Arc<dyn SecretStore>, - clock: Arc<dyn Clock>, - nostr: Arc<dyn NostrClient>, - capacity: NonZeroUsize, - runtime: &Handle, - ) -> Result<Self, SafeError> { - let mut gate = LifecycleGate::opening(); - gate.begin_compatibility_check()?; - gate.compatibility_accepted()?; - gate.ownership_acquired()?; - gate.migration_complete()?; - adapter.bootstrap(secrets.as_ref(), clock.as_ref())?; - gate.recovery_complete()?; - - let adapter = Arc::new(adapter); - let lifecycle = Arc::new(Mutex::new(gate)); - let (mailbox, receiver) = ActorMailbox::bounded(capacity); - let session_generation = Arc::new(AtomicU64::new(SessionGeneration::initial().value())); - let foreground_session = Arc::new(Mutex::new(None)); - let changes = OrderedSnapshotChanges::new(adapter.core().snapshot()); - let durable_request_namespace = format!( - "runtime:{}:{}", - std::process::id(), - std::time::SystemTime::now() - .duration_since(std::time::UNIX_EPOCH) - .map_or(0, |duration| duration.as_nanos()) - ); - let actor = RuntimeActor { - adapter: Arc::clone(&adapter), - secrets, - clock, - nostr, - lifecycle: Arc::clone(&lifecycle), - runtime: runtime.clone(), - session_generation: SessionGeneration::initial(), - published_session_generation: Arc::clone(&session_generation), - profile_tasks: BTreeMap::new(), - changes, - published_foreground_session: Arc::clone(&foreground_session), - durable_request_namespace, - generated_key_stage: GeneratedKeyStage::default(), - }; - drop(runtime.spawn(actor.run(receiver))); - Ok(Self { - mailbox, - adapter, - lifecycle, - runtime: runtime.clone(), - next_request: Arc::new(AtomicU64::new(1)), - session_generation, - foreground_session, - }) - } - - #[must_use] - pub fn lifecycle(&self) -> RuntimeLifecycle { - self.lifecycle - .lock() - .unwrap_or_else(std::sync::PoisonError::into_inner) - .lifecycle() - } - - #[must_use] - pub fn session_generation(&self) -> SessionGeneration { - SessionGeneration::from_value(self.session_generation.load(Ordering::Acquire)) - } - - #[must_use] - pub fn foreground_session(&self) -> Option<ForegroundSessionBinding> { - self.foreground_session - .lock() - .unwrap_or_else(std::sync::PoisonError::into_inner) - .clone() - } - - #[must_use] - pub fn snapshot(&self) -> AppSnapshot { - self.adapter.core().snapshot() - } - - /// Returns the ready snapshot through the actor command boundary. - /// - /// # Errors - /// - /// Returns a typed safe actor error. - pub async fn bootstrap(&self) -> Result<AppSnapshot, SafeError> { - Self::expect_snapshot(self.dispatch(RuntimeCommand::Snapshot, None).await?) - } - - /// Generates one account through the serialized actor boundary. - /// - /// # Errors - /// - /// Returns a safe account, storage, keyring, timeout, or actor error. - pub async fn generate_account(&self) -> Result<GenerateAccountReceipt, SafeError> { - match self.dispatch(RuntimeCommand::GenerateAccount, None).await? { - RuntimeCommandValue::Generated(receipt) => Ok(receipt), - _ => Err(invalid_actor_response()), - } - } - - /// Begins the only actor-owned generated-key recovery stage. - /// - /// # Errors - /// - /// Returns a safe conflict, timeout, key-generation, or actor error. - pub async fn begin_generated_key_stage(&self) -> Result<GeneratedKeyRecoveryHandle, SafeError> { - match self - .dispatch(RuntimeCommand::BeginGeneratedKeyStage, None) - .await? - { - RuntimeCommandValue::GeneratedKeyStage(view) => Ok(view), - _ => Err(invalid_actor_response()), - } - } - - /// Acknowledges recovery and commits the staged account and credential once. - /// - /// # Errors - /// - /// Returns a safe unavailable, conflict, keyring, storage, timeout, or actor error. - pub async fn acknowledge_generated_key_stage( - &self, - id: RecoveryStageId, - ) -> Result<AppSnapshot, SafeError> { - let value = self - .dispatch(RuntimeCommand::AcknowledgeGeneratedKeyStage(id), None) - .await?; - Self::expect_snapshot(value) - } - - /// Cancels and zeroizes the active generated-key stage, if present. - /// - /// # Errors - /// - /// Returns a safe timeout or actor error. - pub async fn cancel_generated_key_stage(&self) -> Result<bool, SafeError> { - match self - .dispatch(RuntimeCommand::CancelGeneratedKeyStage, None) - .await? - { - RuntimeCommandValue::GeneratedKeyStageCancelled(cancelled) => Ok(cancelled), - _ => Err(invalid_actor_response()), - } - } - - /// Imports one account through the serialized actor boundary. - /// - /// # Errors - /// - /// Returns a safe account, storage, keyring, timeout, or actor error. - pub async fn import_secret_key( - &self, - input: SecretKeyInput, - ) -> Result<ImportAccountReceipt, SafeError> { - match self - .dispatch( - RuntimeCommand::ImportSecretKey { - input, - durable_request: None, - durable_expected_revision: None, - }, - None, - ) - .await? - { - RuntimeCommandValue::Imported(receipt) => Ok(receipt), - _ => Err(invalid_actor_response()), - } - } - - /// Imports or repairs with a caller-owned durable request and deadline. - /// - /// # Errors - /// - /// Returns a safe validation, conflict, timeout, persistence, or actor error. - pub async fn import_secret_key_request( - &self, - request: radroots_studio_application::DurableRequestId, - expected_revision: SnapshotRevision, - input: SecretKeyInput, - timeout: Duration, - ) -> Result<ImportAccountReceipt, SafeError> { - let raw_request = self.next_request.fetch_add(1, Ordering::Relaxed); - let request_id = RequestId::new(raw_request).ok_or_else(request_space_exhausted)?; - match self - .dispatch_with_deadline( - RuntimeCommand::ImportSecretKey { - input, - durable_request: Some(request), - durable_expected_revision: Some(expected_revision.value()), - }, - None, - request_id, - Instant::now() + timeout, - ) - .await? - { - RuntimeCommandValue::Imported(receipt) => Ok(receipt), - _ => Err(invalid_actor_response()), - } - } - - /// Selects one account through the serialized actor boundary. - /// - /// # Errors - /// - /// Returns a safe account, storage, timeout, or actor error. - pub async fn select_account(&self, public_key: PublicKey) -> Result<AppSnapshot, SafeError> { - let value = self - .dispatch(RuntimeCommand::SelectAccount(public_key), None) - .await?; - Self::expect_snapshot(value) - } - - /// Activates one account through the serialized actor boundary. - /// - /// # Errors - /// - /// Returns a safe account, credential, storage, timeout, or actor error. - pub async fn activate_account(&self, public_key: PublicKey) -> Result<AppSnapshot, SafeError> { - let value = self - .dispatch(RuntimeCommand::ActivateAccount(public_key), None) - .await?; - Self::expect_snapshot(value) - } - - /// Signs out through the serialized actor boundary. - /// - /// # Errors - /// - /// Returns a safe timeout or actor error. - pub async fn sign_out(&self) -> Result<AppSnapshot, SafeError> { - let value = self.dispatch(RuntimeCommand::SignOut, None).await?; - Self::expect_snapshot(value) - } - - /// Refreshes the active profile through the serialized actor boundary. - /// - /// # Errors - /// - /// Returns a safe relay, storage, timeout, or actor error. - pub async fn refresh_active_profile(&self) -> Result<AppSnapshot, SafeError> { - let value = self - .dispatch(RuntimeCommand::RefreshActiveProfile, None) - .await?; - Self::expect_snapshot(value) - } - - /// Creates one removal request through the serialized actor boundary. - /// - /// # Errors - /// - /// Returns a safe account, timeout, or actor error. - pub async fn request_account_removal( - &self, - public_key: PublicKey, - ) -> Result<RemovalConfirmationToken, SafeError> { - match self - .dispatch(RuntimeCommand::RequestAccountRemoval(public_key), None) - .await? - { - RuntimeCommandValue::RemovalRequest(token) => Ok(token), - _ => Err(invalid_actor_response()), - } - } - - /// Confirms one removal through the serialized actor boundary. - /// - /// # Errors - /// - /// Returns a safe account, credential, storage, timeout, or actor error. - pub async fn confirm_account_removal( - &self, - token: RemovalConfirmationToken, - ) -> Result<AppSnapshot, SafeError> { - let value = self - .dispatch(RuntimeCommand::ConfirmAccountRemoval(token), None) - .await?; - Self::expect_snapshot(value) - } - - /// Closes command admission and cancels supervised work. - /// - /// # Errors - /// - /// Returns a safe timeout or actor error. Repeated calls return closed. - pub async fn close(&self) -> Result<(), SafeError> { - self.close_with_timeout(DEFAULT_COMMAND_TIMEOUT).await - } - - /// Closes the runtime within the supplied command deadline. - /// - /// # Errors - /// - /// Returns a safe timeout or actor error. An expired queued close cannot - /// later change runtime state. - pub async fn close_with_timeout(&self, timeout: Duration) -> Result<(), SafeError> { - let raw_request = self.next_request.fetch_add(1, Ordering::Relaxed); - let request_id = RequestId::new(raw_request).ok_or_else(request_space_exhausted)?; - match self - .dispatch_with_deadline( - RuntimeCommand::Close, - None, - request_id, - Instant::now() + timeout, - ) - .await? - { - RuntimeCommandValue::Closed => Ok(()), - _ => Err(invalid_actor_response()), - } - } - - /// Atomically registers a bounded ordered change consumer with its initial snapshot. - /// - /// # Errors - /// - /// Returns a safe actor or subscription error. - pub async fn subscribe_changes( - &self, - capacity: NonZeroUsize, - ) -> Result<RuntimeChangeSubscription, SafeError> { - match self - .dispatch(RuntimeCommand::SubscribeChanges(capacity), None) - .await? - { - RuntimeCommandValue::Subscription(subscription) => Ok(subscription), - _ => Err(invalid_actor_response()), - } - } - - /// Removes a change consumer through the serialized actor boundary. - /// - /// # Errors - /// - /// Returns a safe actor error. - pub async fn unsubscribe_changes(&self, id: ChangeSubscriptionId) -> Result<bool, SafeError> { - match self - .dispatch(RuntimeCommand::UnsubscribeChanges(id), None) - .await? - { - RuntimeCommandValue::Unsubscribed(removed) => Ok(removed), - _ => Err(invalid_actor_response()), - } - } - - async fn dispatch( - &self, - command: RuntimeCommand, - expected_revision: Option<SnapshotRevision>, - ) -> Result<RuntimeCommandValue, SafeError> { - let raw_request = self.next_request.fetch_add(1, Ordering::Relaxed); - let request_id = RequestId::new(raw_request).ok_or_else(request_space_exhausted)?; - self.dispatch_with_deadline( - command, - expected_revision, - request_id, - Instant::now() + DEFAULT_COMMAND_TIMEOUT, - ) - .await - } - - async fn dispatch_with_deadline( - &self, - command: RuntimeCommand, - expected_revision: Option<SnapshotRevision>, - request_id: RequestId, - deadline: Instant, - ) -> Result<RuntimeCommandValue, SafeError> { - let context = CommandContext::new(request_id, expected_revision, deadline); - let receipt = match self.mailbox.submit(context, command) { - CommandSubmission::Accepted(ticket) => { - let remaining = deadline.saturating_duration_since(Instant::now()); - let waiting = self - .runtime - .spawn(async move { tokio::time::timeout(remaining, ticket.receipt()).await }); - match waiting.await { - Ok(Ok(receipt)) => receipt, - Ok(Err(_)) => CommandReceipt::new(request_id, CommandResult::TimedOut), - Err(_) => CommandReceipt::new(request_id, CommandResult::Closed), - } - } - CommandSubmission::Rejected(receipt) => receipt, - }; - match receipt.into_result() { - CommandResult::Completed(value) => Ok(value), - CommandResult::Conflicted { .. } => Err(command_conflicted()), - CommandResult::Rejected(_) => Err(command_rejected()), - CommandResult::TimedOut => Err(command_timed_out()), - CommandResult::Closed => Err(runtime_closed()), - CommandResult::Failed(error) => Err(error), - } - } - - #[cfg(test)] - async fn import_secret_key_with_timeout( - &self, - input: SecretKeyInput, - timeout: Duration, - ) -> Result<ImportAccountReceipt, SafeError> { - let raw_request = self.next_request.fetch_add(1, Ordering::Relaxed); - let request_id = RequestId::new(raw_request).ok_or_else(request_space_exhausted)?; - match self - .dispatch_with_deadline( - RuntimeCommand::ImportSecretKey { - input, - durable_request: None, - durable_expected_revision: None, - }, - None, - request_id, - Instant::now() + timeout, - ) - .await? - { - RuntimeCommandValue::Imported(receipt) => Ok(receipt), - _ => Err(invalid_actor_response()), - } - } - - fn expect_snapshot(value: RuntimeCommandValue) -> Result<AppSnapshot, SafeError> { - match value { - RuntimeCommandValue::Snapshot(snapshot) => Ok(*snapshot), - _ => Err(invalid_actor_response()), - } - } -} - -impl RuntimeActor { - async fn run( - mut self, - mut receiver: mpsc::Receiver<CommandEnvelope<RuntimeCommand, RuntimeCommandValue>>, - ) { - let (completion_sender, mut completions) = mpsc::channel(DEFAULT_TASK_CAPACITY); - loop { - tokio::select! { - envelope = receiver.recv() => { - let Some(envelope) = envelope else { - break; - }; - if !self.handle_command(envelope, &completion_sender) { - break; - } - } - completion = completions.recv(), if !self.profile_tasks.is_empty() => { - if let Some(completion) = completion { - self.complete_profile_task(completion); - } - } - } - } - self.cancel_profile_tasks(None); - } - - fn handle_command( - &mut self, - envelope: CommandEnvelope<RuntimeCommand, RuntimeCommandValue>, - completion_sender: &mpsc::Sender<ProfileCompletion>, - ) -> bool { - let (context, command, reply) = envelope.into_parts(); - if let Some(result) = self.preflight(context, &command) { - let _ = reply.send(CommandReceipt::new(context.request_id(), result)); - return true; - } - if matches!(command, RuntimeCommand::RefreshActiveProfile) { - self.start_profile_task(context, reply, completion_sender.clone()); - return true; - } - if matches!(command, RuntimeCommand::Close) { - let result = self.close_actor(); - let closed = matches!(result, CommandResult::Completed(_)); - let _ = reply.send(CommandReceipt::new(context.request_id(), result)); - return !closed; - } - let changes_session = matches!( - command, - RuntimeCommand::ActivateAccount(_) - | RuntimeCommand::SignOut - | RuntimeCommand::ConfirmAccountRemoval(_) - ); - let begins_generated_recovery = matches!(&command, RuntimeCommand::BeginGeneratedKeyStage); - let result = self.execute_sync(context, command); - if begins_generated_recovery && matches!(&result, CommandResult::Completed(_)) { - let snapshot = self.adapter.core().snapshot(); - self.cancel_profile_tasks(Some(&snapshot)); - } - if changes_session && matches!(result, CommandResult::Completed(_)) { - self.advance_session_generation(); - self.synchronize_foreground_session(); - } - if matches!(result, CommandResult::Completed(_)) { - self.changes.publish(self.adapter.core().snapshot()); - } - let _ = reply.send(CommandReceipt::new(context.request_id(), result)); - true - } - - fn preflight( - &self, - context: CommandContext, - command: &RuntimeCommand, - ) -> Option<CommandResult<RuntimeCommandValue>> { - 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(); - if matches!(lifecycle.lifecycle(), RuntimeLifecycle::Closed) { - return Some(CommandResult::Closed); - } - if !lifecycle.allows(command.class()) { - return Some(CommandResult::Failed(command_unavailable())); - } - if self.generated_key_stage.pending().is_some() - && !matches!( - command, - RuntimeCommand::Snapshot - | RuntimeCommand::AcknowledgeGeneratedKeyStage(_) - | RuntimeCommand::CancelGeneratedKeyStage - | RuntimeCommand::SubscribeChanges(_) - | RuntimeCommand::UnsubscribeChanges(_) - | RuntimeCommand::Close - ) - { - return Some(CommandResult::Failed(generated_recovery_route_active())); - } - let current_revision = self.adapter.core().snapshot().revision(); - if context - .expected_revision() - .is_some_and(|expected| expected != current_revision) - { - return Some(CommandResult::Conflicted { current_revision }); - } - None - } - - fn execute_sync( - &mut self, - context: CommandContext, - command: RuntimeCommand, - ) -> CommandResult<RuntimeCommandValue> { - let durable_request = radroots_studio_application::DurableRequestId::parse(format!( - "{}:{}", - self.durable_request_namespace, - context.request_id().get() - )); - let expected_revision = context - .expected_revision() - .unwrap_or_else(|| self.adapter.core().snapshot().revision()) - .value(); - let result = match command { - RuntimeCommand::Snapshot => Ok(RuntimeCommandValue::Snapshot(Box::new( - self.adapter.core().snapshot(), - ))), - RuntimeCommand::GenerateAccount => durable_request.and_then(|request| { - self.adapter - .generate_account_durable( - &request, - expected_revision, - self.secrets.as_ref(), - self.clock.as_ref(), - ) - .map(RuntimeCommandValue::Generated) - }), - RuntimeCommand::BeginGeneratedKeyStage => self - .generated_key_stage - .begin( - RecoveryStageId::new( - NonZeroU64::new(context.request_id().get()) - .expect("request IDs are always non-zero"), - ), - expected_revision, - self.clock.now(), - ) - .map(RuntimeCommandValue::GeneratedKeyStage), - RuntimeCommand::AcknowledgeGeneratedKeyStage(id) => { - durable_request.and_then(|request| self.commit_generated_key_stage(&request, id)) - } - RuntimeCommand::CancelGeneratedKeyStage => Ok( - RuntimeCommandValue::GeneratedKeyStageCancelled(self.generated_key_stage.cancel()), - ), - RuntimeCommand::ImportSecretKey { - input, - durable_request: caller_request, - durable_expected_revision, - } => self.import_secret_key_command( - input, - caller_request, - durable_request, - durable_expected_revision.unwrap_or(expected_revision), - ), - RuntimeCommand::SelectAccount(public_key) => self - .adapter - .select_account(public_key) - .map(Box::new) - .map(RuntimeCommandValue::Snapshot), - RuntimeCommand::ActivateAccount(public_key) => self - .adapter - .activate_account(public_key, self.secrets.as_ref(), self.clock.as_ref()) - .map(Box::new) - .map(RuntimeCommandValue::Snapshot), - RuntimeCommand::SignOut => self - .adapter - .sign_out() - .map(Box::new) - .map(RuntimeCommandValue::Snapshot), - RuntimeCommand::RequestAccountRemoval(public_key) => self - .adapter - .request_account_removal(public_key, self.clock.as_ref()) - .map(RuntimeCommandValue::RemovalRequest), - RuntimeCommand::ConfirmAccountRemoval(token) => durable_request.and_then(|request| { - self.adapter - .confirm_account_removal_durable( - &request, - token, - self.secrets.as_ref(), - self.clock.as_ref(), - ) - .map(Box::new) - .map(RuntimeCommandValue::Snapshot) - }), - RuntimeCommand::SubscribeChanges(capacity) => self - .changes - .subscribe(capacity) - .map(|(id, receiver)| { - RuntimeCommandValue::Subscription(RuntimeChangeSubscription { id, receiver }) - }) - .ok_or_else(observer_registration_failed), - RuntimeCommand::UnsubscribeChanges(id) => Ok(RuntimeCommandValue::Unsubscribed( - self.changes.unsubscribe(id), - )), - RuntimeCommand::Close | RuntimeCommand::RefreshActiveProfile => { - Err(invalid_actor_response()) - } - }; - result.map_or_else(CommandResult::Failed, CommandResult::Completed) - } - - fn commit_generated_key_stage( - &mut self, - request: &radroots_studio_application::DurableRequestId, - id: RecoveryStageId, - ) -> Result<RuntimeCommandValue, SafeError> { - let staged = self.generated_key_stage.take(id, self.clock.now())?; - self.adapter.commit_staged_generated_key( - request, - staged, - self.secrets.as_ref(), - self.clock.as_ref(), - )?; - Ok(RuntimeCommandValue::Snapshot(Box::new( - self.adapter.core().snapshot(), - ))) - } - - fn import_secret_key_command( - &self, - input: SecretKeyInput, - caller_request: Option<radroots_studio_application::DurableRequestId>, - fallback_request: Result<radroots_studio_application::DurableRequestId, SafeError>, - expected_revision: u64, - ) -> Result<RuntimeCommandValue, SafeError> { - let request = caller_request.map_or(fallback_request, Ok)?; - self.adapter - .import_secret_key_durable( - &request, - expected_revision, - input, - self.secrets.as_ref(), - self.clock.as_ref(), - ) - .map(RuntimeCommandValue::Imported) - } - - fn start_profile_task( - &mut self, - context: CommandContext, - 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 plan = match self.adapter.core().begin_profile_refresh() { - Ok(Some(plan)) => plan, - Ok(None) => { - let _ = reply.send(CommandReceipt::new( - context.request_id(), - CommandResult::Completed(RuntimeCommandValue::Snapshot(Box::new( - self.adapter.core().snapshot(), - ))), - )); - return; - } - Err(error) => { - let _ = reply.send(CommandReceipt::new( - context.request_id(), - CommandResult::Failed(error), - )); - return; - } - }; - let Some(foreground) = foreground.filter(|binding| { - binding.identity().public_key() == plan.public_key() - && binding.generation() == self.session_generation - }) else { - let _ = reply.send(CommandReceipt::new( - context.request_id(), - CommandResult::Failed(stale_profile_binding()), - )); - return; - }; - let correlation = TaskCorrelation::new( - context.request_id(), - plan.public_key(), - foreground.signer(), - plan.expected_revision(), - self.session_generation, - ); - let client = Arc::clone(&self.nostr); - let relays = plan.relays().to_vec(); - let request_id = context.request_id(); - let handle = self.runtime.spawn(async move { - let result = client.fetch_profile(correlation.account(), &relays).await; - let _ = completion_sender - .send(ProfileCompletion { request_id, result }) - .await; - }); - let previous = self.profile_tasks.insert( - request_id, - PendingProfileTask { - correlation, - plan, - reply, - handle, - }, - ); - debug_assert!(previous.is_none(), "request identifiers are unique"); - } - - fn close_actor(&mut self) -> CommandResult<RuntimeCommandValue> { - let transition = (|| { - let mut lifecycle = self - .lifecycle - .lock() - .unwrap_or_else(std::sync::PoisonError::into_inner); - lifecycle.begin_shutdown()?; - lifecycle.finish_shutdown() - })(); - match transition { - Ok(()) => { - self.generated_key_stage.cancel(); - self.cancel_profile_tasks(None); - self.changes.close(); - *self - .published_foreground_session - .lock() - .unwrap_or_else(std::sync::PoisonError::into_inner) = None; - CommandResult::Completed(RuntimeCommandValue::Closed) - } - Err(error) => CommandResult::Failed(error), - } - } - - fn complete_profile_task(&mut self, completion: ProfileCompletion) { - let Some(task) = self.profile_tasks.remove(&completion.request_id) else { - return; - }; - let current = self.adapter.core().snapshot(); - let foreground = self - .published_foreground_session - .lock() - .unwrap_or_else(std::sync::PoisonError::into_inner) - .clone(); - let correlated = task.correlation.session_generation() == self.session_generation - && foreground.is_some_and(|binding| { - binding.generation() == task.correlation.session_generation() - && binding.identity().public_key() == task.correlation.account() - && binding.signer() == task.correlation.binding() - }) - && current - .active_account() - .is_some_and(|active| active.account().public_key() == task.correlation.account()); - let result = if correlated { - self.adapter - .core() - .complete_profile_refresh( - &task.plan, - completion.result, - self.adapter.database(), - self.clock.as_ref(), - ) - .map(Box::new) - .map(RuntimeCommandValue::Snapshot) - .map_or_else(CommandResult::Failed, CommandResult::Completed) - } else { - CommandResult::Completed(RuntimeCommandValue::Snapshot(Box::new(current))) - }; - if matches!(result, CommandResult::Completed(_)) { - self.changes.publish(self.adapter.core().snapshot()); - } - let _ = task - .reply - .send(CommandReceipt::new(task.correlation.request_id(), result)); - } - - fn advance_session_generation(&mut self) { - let Some(next) = self.session_generation.next() else { - self.lifecycle - .lock() - .unwrap_or_else(std::sync::PoisonError::into_inner) - .fail(request_space_exhausted()); - self.cancel_profile_tasks(None); - return; - }; - self.session_generation = next; - self.published_session_generation - .store(next.value(), Ordering::Release); - let snapshot = self.adapter.core().snapshot(); - self.cancel_profile_tasks(Some(&snapshot)); - } - - fn synchronize_foreground_session(&mut self) { - let session = self - .adapter - .core() - .snapshot() - .active_account() - .map(|active| { - let public_key = active.account().public_key(); - ForegroundSessionBinding::new( - AccountIdentity::derive(public_key)?, - LocalSignerBinding::new(public_key, BindingAvailability::Available), - self.session_generation, - ) - }); - let session = match session.transpose() { - Ok(session) => session, - Err(error) => { - self.lifecycle - .lock() - .unwrap_or_else(std::sync::PoisonError::into_inner) - .fail(error); - None - } - }; - *self - .published_foreground_session - .lock() - .unwrap_or_else(std::sync::PoisonError::into_inner) = session; - } - - fn cancel_profile_tasks(&mut self, snapshot: Option<&AppSnapshot>) { - let tasks = std::mem::take(&mut self.profile_tasks); - for (_, task) in tasks { - task.handle.abort(); - let receipt_result = snapshot.map_or(CommandResult::Closed, |snapshot| { - CommandResult::Completed(RuntimeCommandValue::Snapshot(Box::new(snapshot.clone()))) - }); - let _ = task.reply.send(CommandReceipt::new( - task.correlation.request_id(), - receipt_result, - )); - } - } -} - -const fn request_space_exhausted() -> SafeError { - SafeError::new( - SafeErrorCode::InvalidApplicationState, - SafeMessage::new("The runtime request identifier space is exhausted."), - ) -} - -const fn stale_profile_binding() -> SafeError { - SafeError::new( - SafeErrorCode::InvalidApplicationState, - SafeMessage::new("The active account binding changed before profile refresh."), - ) -} - -const fn command_conflicted() -> SafeError { - SafeError::new( - SafeErrorCode::InvalidApplicationState, - SafeMessage::new("The command conflicts with newer application state."), - ) -} - -const fn command_rejected() -> SafeError { - SafeError::new( - SafeErrorCode::InvalidApplicationState, - SafeMessage::new("The runtime is busy. Try again."), - ) -} - -const fn command_timed_out() -> SafeError { - SafeError::new( - SafeErrorCode::InvalidApplicationState, - SafeMessage::new("The runtime command timed out."), - ) -} - -const fn runtime_closed() -> SafeError { - SafeError::new( - SafeErrorCode::InvalidApplicationState, - SafeMessage::new("The application runtime is closed."), - ) -} - -const fn command_unavailable() -> SafeError { - SafeError::new( - SafeErrorCode::InvalidApplicationState, - SafeMessage::new("The command is unavailable in the current runtime state."), - ) -} - -const fn generated_recovery_route_active() -> SafeError { - SafeError::new( - SafeErrorCode::InvalidApplicationState, - SafeMessage::new("Complete or cancel generated-key recovery before another action."), - ) -} - -const fn invalid_actor_response() -> SafeError { - SafeError::new( - SafeErrorCode::InvalidApplicationState, - SafeMessage::new("The runtime returned an invalid command response."), - ) -} - -const fn observer_registration_failed() -> SafeError { - SafeError::new( - SafeErrorCode::ObserverRegistrationFailed, - SafeMessage::new("The application change subscription could not be registered."), - ) -} - -#[cfg(test)] -mod tests { - use std::num::NonZeroUsize; - use std::sync::atomic::{AtomicBool, Ordering}; - use std::sync::{Arc, Condvar, Mutex}; - use std::time::Duration; - - use radroots_studio_application::{ - BoxFuture, Clock, FailureSecretStore, InMemorySecretStore, NostrClient, RelayConfiguration, - RuntimeLifecycle, SecretStore, SecretStoreOperation, SessionState, - }; - use radroots_studio_domain::{ - Kind0ProfileCandidate, PublicKey, RelayUrl, SafeError, SafeErrorCode, SecretKeyInput, - UnixTimestamp, - }; - - use super::RuntimeActorHandle; - - struct FixedClock; - - impl Clock for FixedClock { - fn now(&self) -> UnixTimestamp { - UnixTimestamp::from_seconds(50).expect("time") - } - } - - struct OfflineNostr; - - impl NostrClient for OfflineNostr { - fn fetch_profile<'a>( - &'a self, - _public_key: PublicKey, - _relays: &'a [RelayUrl], - ) -> BoxFuture<'a, Result<Option<Kind0ProfileCandidate>, SafeError>> { - Box::pin(async { Ok(None) }) - } - } - - struct BlockingNostr { - started: tokio::sync::Semaphore, - release: tokio::sync::Semaphore, - } - - impl BlockingNostr { - fn new() -> Self { - Self { - started: tokio::sync::Semaphore::new(0), - release: tokio::sync::Semaphore::new(0), - } - } - } - - impl NostrClient for BlockingNostr { - fn fetch_profile<'a>( - &'a self, - _public_key: PublicKey, - _relays: &'a [RelayUrl], - ) -> BoxFuture<'a, Result<Option<Kind0ProfileCandidate>, SafeError>> { - Box::pin(async move { - self.started.add_permits(1); - let permit = self.release.acquire().await.expect("release"); - permit.forget(); - Ok(None) - }) - } - } - - struct BlockingSecretStore { - inner: InMemorySecretStore, - block_next_put: AtomicBool, - put_started: AtomicBool, - released: Mutex<bool>, - release_signal: Condvar, - } - - impl BlockingSecretStore { - fn new() -> Self { - Self { - inner: InMemorySecretStore::default(), - block_next_put: AtomicBool::new(true), - put_started: AtomicBool::new(false), - released: Mutex::new(false), - release_signal: Condvar::new(), - } - } - - async fn wait_until_put_started(&self) { - while !self.put_started.load(Ordering::Acquire) { - tokio::task::yield_now().await; - } - } - - fn release(&self) { - *self - .released - .lock() - .unwrap_or_else(std::sync::PoisonError::into_inner) = true; - self.release_signal.notify_all(); - } - } - - impl SecretStore for BlockingSecretStore { - fn put(&self, public_key: PublicKey, secret: SecretKeyInput) -> Result<(), SafeError> { - if self.block_next_put.swap(false, Ordering::AcqRel) { - self.put_started.store(true, Ordering::Release); - let released = self - .released - .lock() - .unwrap_or_else(std::sync::PoisonError::into_inner); - drop( - self.release_signal - .wait_while(released, |released| !*released) - .unwrap_or_else(std::sync::PoisonError::into_inner), - ); - } - self.inner.put(public_key, secret) - } - - fn load(&self, public_key: PublicKey) -> Result<SecretKeyInput, SafeError> { - self.inner.load(public_key) - } - - fn contains(&self, public_key: PublicKey) -> Result<bool, SafeError> { - self.inner.contains(public_key) - } - - fn delete(&self, public_key: PublicKey) -> Result<(), SafeError> { - self.inner.delete(public_key) - } - } - - fn actor() -> (RuntimeActorHandle, Arc<InMemorySecretStore>) { - let secrets = Arc::new(InMemorySecretStore::default()); - let secret_port: Arc<dyn SecretStore> = secrets.clone(); - let actor = RuntimeActorHandle::in_memory( - RelayConfiguration::default(), - secret_port, - Arc::new(FixedClock), - Arc::new(OfflineNostr), - NonZeroUsize::new(8).expect("capacity"), - &tokio::runtime::Handle::current(), - ) - .expect("actor"); - (actor, secrets) - } - - #[tokio::test(flavor = "multi_thread")] - async fn account_mutations_run_serially_through_one_ready_actor() { - let (actor, secrets) = actor(); - assert_eq!(actor.lifecycle(), RuntimeLifecycle::Ready); - - let imported = actor - .import_secret_key( - SecretKeyInput::parse( - "7e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7".to_owned(), - ) - .expect("input"), - ) - .await - .expect("import"); - let public_key = imported.account().public_key(); - let activated = actor.activate_account(public_key).await.expect("activate"); - assert_eq!(activated.session(), SessionState::Active); - let foreground = actor.foreground_session().expect("foreground session"); - assert_eq!(foreground.identity().public_key(), public_key); - assert_eq!(foreground.signer().account(), public_key); - assert_eq!(foreground.generation(), actor.session_generation()); - assert!(secrets.contains(public_key).expect("credential")); - - let signed_out = actor.sign_out().await.expect("sign out"); - assert_eq!(signed_out.session(), SessionState::SignedOut); - assert!(actor.foreground_session().is_none()); - let removal = actor - .request_account_removal(public_key) - .await - .expect("removal request"); - let removed = actor - .confirm_account_removal(removal) - .await - .expect("remove"); - assert!(removed.accounts().is_empty()); - assert!(!secrets.contains(public_key).expect("credential removed")); - } - - #[tokio::test(flavor = "multi_thread")] - async fn generated_key_stage_is_exclusive_cancelable_and_snapshot_free() { - let (actor, secrets) = actor(); - let initial = actor.snapshot(); - let stage = actor - .begin_generated_key_stage() - .await - .expect("generated key stage"); - - assert!(actor.begin_generated_key_stage().await.is_err()); - assert_eq!(actor.snapshot(), initial); - assert!( - !secrets - .contains(stage.view().account().public_key()) - .expect("keyring") - ); - assert!(actor.sign_out().await.is_err()); - assert_eq!(actor.snapshot(), initial); - assert!(actor.cancel_generated_key_stage().await.expect("cancel")); - assert!( - !actor - .cancel_generated_key_stage() - .await - .expect("cancel empty") - ); - assert_eq!(actor.snapshot(), initial); - - actor - .begin_generated_key_stage() - .await - .expect("replacement stage"); - actor.close().await.expect("close clears stage"); - assert_eq!(actor.lifecycle(), RuntimeLifecycle::Closed); - } - - #[tokio::test(flavor = "multi_thread")] - async fn recovery_handle_is_one_use_and_acknowledgement_commits_once() { - let (actor, secrets) = actor(); - let initial = actor.snapshot(); - let handle = actor - .begin_generated_key_stage() - .await - .expect("generated key stage"); - let public_key = handle.view().account().public_key(); - let recovery = handle.take_recovery_nsec().expect("recovery material"); - assert_eq!(recovery.with_exposed_secret(str::len), 63); - assert!(handle.take_recovery_nsec().is_err()); - assert_eq!(actor.snapshot(), initial); - assert!(!secrets.contains(public_key).expect("not committed")); - - let committed = actor - .acknowledge_generated_key_stage(handle.id()) - .await - .expect("acknowledge"); - assert_eq!(committed.accounts().len(), 1); - assert_eq!(committed.selected_account(), Some(public_key)); - assert!(secrets.contains(public_key).expect("credential committed")); - assert!( - actor - .acknowledge_generated_key_stage(handle.id()) - .await - .is_err() - ); - } - - #[tokio::test(flavor = "multi_thread")] - async fn failed_generated_commit_consumes_the_stage_without_poisoning_the_actor() { - let secrets = Arc::new(FailureSecretStore::default()); - secrets.fail_next(SecretStoreOperation::Put); - let secret_port: Arc<dyn SecretStore> = secrets.clone(); - let actor = RuntimeActorHandle::in_memory( - RelayConfiguration::default(), - secret_port, - Arc::new(FixedClock), - Arc::new(OfflineNostr), - NonZeroUsize::new(8).expect("capacity"), - &tokio::runtime::Handle::current(), - ) - .expect("actor"); - let handle = actor - .begin_generated_key_stage() - .await - .expect("generated key stage"); - - let error = actor - .acknowledge_generated_key_stage(handle.id()) - .await - .expect_err("injected keyring failure"); - - assert_eq!(error.code(), SafeErrorCode::KeyringUnavailable); - assert!(actor.snapshot().accounts().is_empty()); - actor - .begin_generated_key_stage() - .await - .expect("fresh recovery after terminal failure"); - assert!(actor.cancel_generated_key_stage().await.expect("cancel")); - } - - #[tokio::test(flavor = "multi_thread")] - async fn session_generation_cancels_correlated_profile_work_on_sign_out() { - let client = Arc::new(BlockingNostr::new()); - let actor = RuntimeActorHandle::in_memory( - RelayConfiguration::new(vec![RelayUrl::parse("ws://localhost:8080").expect("relay")]), - Arc::new(InMemorySecretStore::default()), - Arc::new(FixedClock), - client.clone(), - NonZeroUsize::new(8).expect("capacity"), - &tokio::runtime::Handle::current(), - ) - .expect("actor"); - let imported = actor - .import_secret_key( - SecretKeyInput::parse( - "7e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7".to_owned(), - ) - .expect("input"), - ) - .await - .expect("import"); - actor - .activate_account(imported.account().public_key()) - .await - .expect("activate"); - assert_eq!(actor.session_generation().value(), 1); - - let refresh_actor = actor.clone(); - let refresh = tokio::spawn(async move { refresh_actor.refresh_active_profile().await }); - let started = client.started.acquire().await.expect("refresh started"); - started.forget(); - let signed_out = actor.sign_out().await.expect("sign out"); - let cancelled = refresh - .await - .expect("refresh task") - .expect("safe cancellation"); - - assert_eq!(actor.session_generation().value(), 2); - assert_eq!(signed_out.session(), SessionState::SignedOut); - assert_eq!(cancelled.session(), SessionState::SignedOut); - assert!(cancelled.active_account().is_none()); - } - - #[tokio::test(flavor = "multi_thread", worker_threads = 4)] - async fn bounded_runtime_rejects_saturation_without_dropping_accepted_commands() { - let secrets = Arc::new(BlockingSecretStore::new()); - let actor = RuntimeActorHandle::in_memory( - RelayConfiguration::default(), - secrets.clone(), - Arc::new(FixedClock), - Arc::new(OfflineNostr), - NonZeroUsize::new(1).expect("capacity"), - &tokio::runtime::Handle::current(), - ) - .expect("actor"); - - let first_actor = actor.clone(); - let first = tokio::spawn(async move { - first_actor - .import_secret_key(secret( - "7e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7", - )) - .await - }); - secrets.wait_until_put_started().await; - - let second_actor = actor.clone(); - let second = tokio::spawn(async move { - second_actor - .import_secret_key(secret( - "0000000000000000000000000000000000000000000000000000000000000001", - )) - .await - }); - while actor.mailbox.available_capacity() != 0 { - assert!( - !second.is_finished(), - "second command must enter the mailbox" - ); - tokio::task::yield_now().await; - } - let rejected = actor - .import_secret_key(secret( - "0000000000000000000000000000000000000000000000000000000000000002", - )) - .await - .expect_err("full mailbox must reject"); - assert_eq!( - rejected.message().as_str(), - "The runtime is busy. Try again." - ); - - secrets.release(); - first.await.expect("first task").expect("first command"); - second.await.expect("second task").expect("second command"); - assert_eq!( - actor.bootstrap().await.expect("snapshot").accounts().len(), - 2 - ); - } - - #[tokio::test(flavor = "multi_thread", worker_threads = 4)] - async fn queued_command_expiry_returns_timeout_and_prevents_late_mutation() { - let secrets = Arc::new(BlockingSecretStore::new()); - let actor = RuntimeActorHandle::in_memory( - RelayConfiguration::default(), - secrets.clone(), - Arc::new(FixedClock), - Arc::new(OfflineNostr), - NonZeroUsize::new(1).expect("capacity"), - &tokio::runtime::Handle::current(), - ) - .expect("actor"); - - let first_actor = actor.clone(); - let first = tokio::spawn(async move { - first_actor - .import_secret_key(secret( - "7e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7", - )) - .await - }); - secrets.wait_until_put_started().await; - - let expired = actor - .import_secret_key_with_timeout( - secret("0000000000000000000000000000000000000000000000000000000000000001"), - Duration::from_millis(10), - ) - .await - .expect_err("queued command must time out"); - assert_eq!(expired.message().as_str(), "The runtime command timed out."); - - secrets.release(); - first.await.expect("first task").expect("first command"); - assert_eq!( - actor.bootstrap().await.expect("snapshot").accounts().len(), - 1 - ); - } - - #[tokio::test(flavor = "multi_thread")] - async fn close_is_terminal_and_every_later_command_is_rejected_as_closed() { - let (actor, _) = actor(); - actor.close().await.expect("close"); - assert_eq!(actor.lifecycle(), RuntimeLifecycle::Closed); - - for error in [ - actor.bootstrap().await.expect_err("bootstrap after close"), - actor.close().await.expect_err("repeated close"), - ] { - assert_eq!( - error.message().as_str(), - "The application runtime is closed." - ); - } - } - - #[tokio::test(flavor = "multi_thread")] - async fn actor_subscription_atomically_delivers_initial_then_ordered_changes() { - let (actor, _) = actor(); - let mut subscription = actor - .subscribe_changes(NonZeroUsize::new(4).expect("capacity")) - .await - .expect("subscribe"); - let initial = subscription.receive().await.expect("initial snapshot"); - assert_eq!(initial.revision(), actor.snapshot().revision()); - assert!(initial.previous_revision().is_none()); - - actor - .import_secret_key(secret( - "7e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7", - )) - .await - .expect("import"); - let changed = subscription.receive().await.expect("change"); - assert!(changed.revision() > initial.revision()); - assert_eq!(changed.previous_revision(), Some(initial.revision())); - assert!( - actor - .unsubscribe_changes(subscription.id()) - .await - .expect("unsubscribe") - ); - } - - #[tokio::test(flavor = "multi_thread", worker_threads = 4)] - async fn expired_queued_shutdown_does_not_close_runtime_later() { - let secrets = Arc::new(BlockingSecretStore::new()); - let actor = RuntimeActorHandle::in_memory( - RelayConfiguration::default(), - secrets.clone(), - Arc::new(FixedClock), - Arc::new(OfflineNostr), - NonZeroUsize::new(1).expect("capacity"), - &tokio::runtime::Handle::current(), - ) - .expect("actor"); - let import_actor = actor.clone(); - let import = tokio::spawn(async move { - import_actor - .import_secret_key(secret( - "7e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7", - )) - .await - }); - secrets.wait_until_put_started().await; - - let timeout = actor - .close_with_timeout(Duration::from_millis(10)) - .await - .expect_err("queued shutdown must expire"); - assert_eq!(timeout.message().as_str(), "The runtime command timed out."); - secrets.release(); - import.await.expect("import task").expect("import"); - assert_eq!(actor.lifecycle(), RuntimeLifecycle::Ready); - assert_eq!( - actor - .bootstrap() - .await - .expect("still open") - .accounts() - .len(), - 1 - ); - actor.close().await.expect("later close"); - } - - #[tokio::test(flavor = "multi_thread", worker_threads = 4)] - async fn shutdown_cancels_in_flight_work_and_terminates_publication() { - let client = Arc::new(BlockingNostr::new()); - let actor = RuntimeActorHandle::in_memory( - RelayConfiguration::new(vec![RelayUrl::parse("ws://localhost:8080").expect("relay")]), - Arc::new(InMemorySecretStore::default()), - Arc::new(FixedClock), - client.clone(), - NonZeroUsize::new(8).expect("capacity"), - &tokio::runtime::Handle::current(), - ) - .expect("actor"); - let imported = actor - .import_secret_key(secret( - "7e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7", - )) - .await - .expect("import"); - actor - .activate_account(imported.account().public_key()) - .await - .expect("activate"); - let mut changes = actor - .subscribe_changes(NonZeroUsize::new(4).expect("capacity")) - .await - .expect("subscribe"); - changes.receive().await.expect("initial"); - - let refresh_actor = actor.clone(); - let refresh = tokio::spawn(async move { refresh_actor.refresh_active_profile().await }); - let started = client.started.acquire().await.expect("refresh started"); - started.forget(); - actor.close().await.expect("close"); - - let cancelled = refresh - .await - .expect("refresh task") - .expect_err("refresh closes"); - assert_eq!( - cancelled.message().as_str(), - "The application runtime is closed." - ); - assert!(changes.receive().await.is_none()); - assert_eq!(actor.lifecycle(), RuntimeLifecycle::Closed); - } - - fn secret(value: &str) -> SecretKeyInput { - SecretKeyInput::parse(value.to_owned()).expect("valid test secret") - } -} diff --git a/crates/studio_storage/tests/local_relay_e2e.rs b/crates/studio_storage/tests/local_relay_e2e.rs @@ -1,95 +0,0 @@ -use std::time::Duration; - -use nostr::{EventBuilder, Keys, Metadata}; -use nostr_relay_builder::MockRelay; -use nostr_sdk::Client; -use radroots_studio_application::{ - Clock, InMemorySecretStore, ProfileLoadState, ProfileRepository, RelayConfiguration, - RelayConnectionState, SdkNostrClient, SecretStore, SessionState, -}; -use radroots_studio_domain::{RelayUrl, SecretKeyInput, UnixTimestamp}; -use radroots_studio_storage::PersistentAppCore; - -const SECRET_HEX: &str = "7e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7"; - -struct FixedClock; - -impl Clock for FixedClock { - fn now(&self) -> UnixTimestamp { - UnixTimestamp::from_seconds(100).expect("fixed timestamp") - } -} - -#[tokio::test] -async fn local_relay_e2e_imports_activates_refreshes_and_caches_profile() { - let local_relay = MockRelay::run().await.expect("local relay"); - let relay_url = local_relay.url().await; - let keys = Keys::parse(SECRET_HEX).expect("known secret key"); - let publisher = Client::new(keys); - publisher - .add_relay(relay_url.clone()) - .await - .expect("publisher relay"); - publisher.connect().await; - publisher.wait_for_connection(Duration::from_secs(2)).await; - publisher - .send_event_builder(EventBuilder::metadata( - &Metadata::new() - .name("farmer") - .display_name("Farm Account") - .about("Local food profile"), - )) - .await - .expect("publish profile"); - - let relay = RelayUrl::parse(relay_url.as_str()).expect("relay URL"); - let adapter = PersistentAppCore::in_memory(RelayConfiguration::new(vec![relay])) - .expect("persistent adapter"); - let secrets = InMemorySecretStore::default(); - adapter.bootstrap(&secrets, &FixedClock).expect("bootstrap"); - let imported = adapter - .import_secret_key( - SecretKeyInput::parse(SECRET_HEX.to_owned()).expect("secret input"), - &secrets, - &FixedClock, - ) - .expect("import account"); - let public_key = imported.account().public_key(); - assert!(secrets.contains(public_key).expect("credential exists")); - adapter - .activate_account(public_key, &secrets, &FixedClock) - .expect("activate account"); - - let refreshed = adapter - .core() - .refresh_active_profile( - adapter.database(), - &SdkNostrClient::new(Duration::from_secs(2)), - &FixedClock, - ) - .await - .expect("refresh profile"); - - assert_eq!(refreshed.session(), SessionState::Active); - let active = refreshed.active_account().expect("active account"); - assert_eq!(active.relay_state(), RelayConnectionState::Connected); - assert_eq!(active.profile_state(), ProfileLoadState::Fresh); - assert_eq!( - active.profile().and_then(|profile| profile.display_name()), - Some("Farm Account") - ); - let cached = adapter - .database() - .load_profile(public_key) - .expect("load cache") - .expect("cached profile"); - assert_eq!( - cached.candidate().metadata().preferred_name(), - Some("Farm Account") - ); - let public_debug = format!("{refreshed:?}"); - assert!(!public_debug.contains(SECRET_HEX)); - assert!(!public_debug.contains("nsec1")); - publisher.shutdown().await; - local_relay.shutdown(); -} diff --git a/crates/studio_storage/tests/restart_isolation.rs b/crates/studio_storage/tests/restart_isolation.rs @@ -1,95 +0,0 @@ -use std::fs; - -use radroots_studio_application::{ - AccountNamespaceRepository, AccountPreferenceKey, Clock, InMemorySecretStore, - RelayConfiguration, SessionState, -}; -use radroots_studio_domain::{SecretKeyInput, UnixTimestamp}; -use radroots_studio_storage::PersistentAppCore; -use tempfile::tempdir; - -const SECRET_A: &str = "7e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7"; -const SECRET_B: &str = "0101010101010101010101010101010101010101010101010101010101010101"; - -struct FixedClock; - -impl Clock for FixedClock { - fn now(&self) -> UnixTimestamp { - UnixTimestamp::from_seconds(200).expect("fixed timestamp") - } -} - -#[test] -fn restart_restores_selection_and_keeps_account_namespaces_isolated() { - let directory = tempdir().expect("temporary directory"); - let path = directory.path().join("studio.sqlite3"); - let secrets = InMemorySecretStore::default(); - let (owner_a, owner_b); - - { - let adapter = PersistentAppCore::open(&path, RelayConfiguration::default()) - .expect("persistent adapter"); - adapter.bootstrap(&secrets, &FixedClock).expect("bootstrap"); - owner_a = adapter - .import_secret_key( - SecretKeyInput::parse(SECRET_A.to_owned()).expect("secret A"), - &secrets, - &FixedClock, - ) - .expect("account A") - .account() - .public_key(); - owner_b = adapter - .import_secret_key( - SecretKeyInput::parse(SECRET_B.to_owned()).expect("secret B"), - &secrets, - &FixedClock, - ) - .expect("account B") - .account() - .public_key(); - adapter - .database() - .set_value(owner_a, AccountPreferenceKey::NamespaceProbe, "account-a") - .expect("namespace A"); - adapter - .database() - .set_value(owner_b, AccountPreferenceKey::NamespaceProbe, "account-b") - .expect("namespace B"); - adapter.select_account(owner_b).expect("select B"); - } - - let reopened = - PersistentAppCore::open(&path, RelayConfiguration::default()).expect("reopen adapter"); - let restored = reopened.bootstrap(&secrets, &FixedClock).expect("restore"); - assert_eq!(restored.accounts().len(), 2); - assert_eq!(restored.selected_account(), Some(owner_b)); - assert_eq!(restored.session(), SessionState::SignedOut); - assert_eq!( - reopened - .database() - .get_value(owner_a, AccountPreferenceKey::NamespaceProbe) - .expect("read A"), - Some("account-a".to_owned()) - ); - assert_eq!( - reopened - .database() - .get_value(owner_b, AccountPreferenceKey::NamespaceProbe) - .expect("read B"), - Some("account-b".to_owned()) - ); - - let database = fs::read(path).expect("database bytes"); - assert!( - !database - .windows(SECRET_A.len()) - .any(|bytes| bytes == SECRET_A.as_bytes()) - ); - assert!( - !database - .windows(SECRET_B.len()) - .any(|bytes| bytes == SECRET_B.as_bytes()) - ); - assert!(!database.windows(5).any(|bytes| bytes == b"nsec1")); -} diff --git a/tools/xtask/src/catalog.rs b/tools/xtask/src/catalog.rs @@ -280,8 +280,12 @@ fn validate_catalog(catalog: &Catalog) -> Result<(), String> { "private_fixture", "private_tool", ]); - let allowed_licenses = - BTreeSet::from(["MIT OR Apache-2.0", "GPL-3.0-only", "GPL-3.0-or-later"]); + let allowed_licenses = BTreeSet::from([ + "MIT OR Apache-2.0", + "GPL-3.0-only", + "GPL-3.0-or-later", + "MPL-2.0", + ]); let mut names = BTreeSet::new(); let mut paths = BTreeSet::new(); let mut public = BTreeSet::new();