app

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

commit 249e12d4c6eaaa71e6fd46b0dc330af684af6c40
parent 97bb016bcf29555bb890a2e3f1effb98890f5502
Author: triesap <tyson@radroots.org>
Date:   Sun,  9 Aug 2026 17:46:00 +0000

storage: own the Studio storage crate

- copy storage source migrations and tests into HarvestCircle core
- register storage as a local workspace member and dependency
- carry forward the source workspace coverage cfg lint contract
- refresh exact cryptographic SQLite and filesystem dependency edges

Diffstat:
Mcore/Cargo.lock | 30++++++++++++++++++++++++------
Mcore/Cargo.toml | 8+++++++-
Acore/crates/studio_storage/Cargo.toml | 33+++++++++++++++++++++++++++++++++
Acore/crates/studio_storage/migrations/V10__installation_identity.sql | 7+++++++
Acore/crates/studio_storage/migrations/V1__initialize.sql | 6++++++
Acore/crates/studio_storage/migrations/V2__accounts.sql | 28++++++++++++++++++++++++++++
Acore/crates/studio_storage/migrations/V3__profile_cache.sql | 12++++++++++++
Acore/crates/studio_storage/migrations/V4__account_namespace.sql | 6++++++
Acore/crates/studio_storage/migrations/V5__operation_journal.sql | 8++++++++
Acore/crates/studio_storage/migrations/V6__normalized_runtime_schema.sql | 100+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acore/crates/studio_storage/migrations/V7__migrate_v5_runtime_data.sql | 68++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acore/crates/studio_storage/migrations/V8__normalized_account_preferences.sql | 10++++++++++
Acore/crates/studio_storage/migrations/V9__durable_operation_receipts.sql | 11+++++++++++
Acore/crates/studio_storage/src/account_namespace.rs | 189+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acore/crates/studio_storage/src/accounts.rs | 516+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acore/crates/studio_storage/src/compatibility.rs | 474+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acore/crates/studio_storage/src/db.rs | 907+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acore/crates/studio_storage/src/installation.rs | 71+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acore/crates/studio_storage/src/journal.rs | 774+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acore/crates/studio_storage/src/lib.rs | 20++++++++++++++++++++
Acore/crates/studio_storage/src/os_keyring.rs | 138+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acore/crates/studio_storage/src/profiles.rs | 254+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acore/crates/studio_storage/src/recovery.rs | 692+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acore/crates/studio_storage/src/repair.rs | 403+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acore/crates/studio_storage/tests/redaction.rs | 52++++++++++++++++++++++++++++++++++++++++++++++++++++
25 files changed, 4810 insertions(+), 7 deletions(-)

diff --git a/core/Cargo.lock b/core/Cargo.lock @@ -1993,7 +1993,7 @@ dependencies = [ "radroots_studio_domain 0.1.0-alpha (git+https://github.com/radrootslabs/lib?rev=09065a610d95e57acdc895a14c07580fa099e7c3)", "radroots_studio_nostr", "radroots_studio_runtime", - "radroots_studio_storage", + "radroots_studio_storage 0.1.0-alpha (git+https://github.com/radrootslabs/lib?rev=09065a610d95e57acdc895a14c07580fa099e7c3)", "sha2", "syn 2.0.119", "tokio", @@ -2030,7 +2030,7 @@ dependencies = [ "radroots_studio_application 0.1.0-alpha (git+https://github.com/radrootslabs/lib?rev=09065a610d95e57acdc895a14c07580fa099e7c3)", "radroots_studio_domain 0.1.0-alpha (git+https://github.com/radrootslabs/lib?rev=09065a610d95e57acdc895a14c07580fa099e7c3)", "radroots_studio_nostr", - "radroots_studio_storage", + "radroots_studio_storage 0.1.0-alpha (git+https://github.com/radrootslabs/lib?rev=09065a610d95e57acdc895a14c07580fa099e7c3)", "tokio", "uuid", ] @@ -2045,7 +2045,25 @@ dependencies = [ "radroots_studio_nostr", "radroots_studio_preferences", "radroots_studio_runtime", - "radroots_studio_storage", + "radroots_studio_storage 0.1.0-alpha", +] + +[[package]] +name = "radroots_studio_storage" +version = "0.1.0-alpha" +dependencies = [ + "fs2", + "getrandom 0.2.17", + "hmac", + "keyring", + "radroots_studio_application 0.1.0-alpha", + "radroots_studio_domain 0.1.0-alpha", + "refinery", + "rusqlite", + "rustix", + "sha2", + "tempfile", + "zeroize", ] [[package]] @@ -2695,12 +2713,12 @@ dependencies = [ [[package]] name = "tempfile" -version = "3.27.0" +version = "3.23.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "32497e9a4c7b38532efcdebeef879707aa9f794296a4f0244f6f69e9bc8574bd" +checksum = "2d31c77bdf42a745371d260a26ca7163f1e0924b64afa0b688e61b5a9fa02f16" dependencies = [ "fastrand", - "getrandom 0.4.3", + "getrandom 0.3.4", "once_cell", "rustix", "windows-sys 0.61.2", diff --git a/core/Cargo.toml b/core/Cargo.toml @@ -4,6 +4,7 @@ members = [ "crates/studio_application", "crates/studio_domain", "crates/studio_preferences", + "crates/studio_storage", ] resolver = "3" @@ -18,6 +19,7 @@ authors = ["Tyson Lupul <tyson@radroots.org>"] [workspace.lints.rust] unsafe_code = "forbid" +unexpected_cfgs = { level = "warn", check-cfg = ['cfg(coverage_nightly)'] } [workspace.lints.clippy] all = "deny" @@ -30,5 +32,9 @@ radroots_studio_ffi = { git = "https://github.com/radrootslabs/lib", rev = "0906 radroots_studio_nostr = { git = "https://github.com/radrootslabs/lib", rev = "09065a610d95e57acdc895a14c07580fa099e7c3", version = "=0.1.0-alpha" } radroots_studio_preferences = { path = "crates/studio_preferences", version = "=0.1.0-alpha" } radroots_studio_runtime = { git = "https://github.com/radrootslabs/lib", rev = "09065a610d95e57acdc895a14c07580fa099e7c3", version = "=0.1.0-alpha" } -radroots_studio_storage = { git = "https://github.com/radrootslabs/lib", rev = "09065a610d95e57acdc895a14c07580fa099e7c3", version = "=0.1.0-alpha" } +radroots_studio_storage = { path = "crates/studio_storage", version = "=0.1.0-alpha" } radroots_identity = { git = "https://github.com/radrootslabs/lib", rev = "09065a610d95e57acdc895a14c07580fa099e7c3", version = "=0.1.0-alpha", default-features = false } +getrandom = { version = "0.2", default-features = false } +hmac = { version = "0.12", default-features = false } +rustix = { version = "1", features = ["fs", "process", "std"] } +sha2 = { version = "0.10", default-features = false } diff --git a/core/crates/studio_storage/Cargo.toml b/core/crates/studio_storage/Cargo.toml @@ -0,0 +1,33 @@ +[package] +name = "radroots_studio_storage" +description = "Private persistence and keyring adapters 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/**", "migrations/**", "Cargo.toml"] + +[dependencies] +fs2 = "=0.4.3" +keyring = "=4.1.6" +radroots_studio_application.workspace = true +radroots_studio_domain.workspace = true +refinery = { version = "=0.9.2", default-features = false, features = ["rusqlite"] } +getrandom.workspace = true +hmac.workspace = true +rusqlite = { version = "=0.39.0", features = ["backup", "bundled"] } +sha2.workspace = true +zeroize = "=1.9.0" + +[target.'cfg(unix)'.dependencies] +rustix.workspace = true + +[dev-dependencies] +tempfile = "=3.23.0" + +[lints] +workspace = true diff --git a/core/crates/studio_storage/migrations/V10__installation_identity.sql b/core/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/core/crates/studio_storage/migrations/V1__initialize.sql b/core/crates/studio_storage/migrations/V1__initialize.sql @@ -0,0 +1,6 @@ +CREATE TABLE application_schema ( + singleton INTEGER PRIMARY KEY CHECK (singleton = 1), + schema_version INTEGER NOT NULL CHECK (schema_version >= 1) +); + +INSERT INTO application_schema (singleton, schema_version) VALUES (1, 1); diff --git a/core/crates/studio_storage/migrations/V2__accounts.sql b/core/crates/studio_storage/migrations/V2__accounts.sql @@ -0,0 +1,28 @@ +CREATE TABLE accounts ( + pubkey TEXT PRIMARY KEY NOT NULL CHECK ( + length(pubkey) = 64 AND pubkey = lower(pubkey) + ), + npub TEXT NOT NULL CHECK (length(npub) = 63), + signer_kind TEXT NOT NULL CHECK ( + signer_kind IN ('local_secret', 'watch_only', 'remote_nip46') + ), + key_availability TEXT NOT NULL CHECK ( + key_availability IN ( + 'available', + 'credential_missing', + 'store_unavailable', + 'not_required' + ) + ), + label TEXT, + created_at INTEGER NOT NULL CHECK (created_at >= 0), + last_used_at INTEGER CHECK (last_used_at >= 0) +); + +CREATE TABLE app_state ( + singleton INTEGER PRIMARY KEY CHECK (singleton = 1), + selected_pubkey TEXT REFERENCES accounts(pubkey) ON DELETE SET NULL +); + +INSERT INTO app_state (singleton, selected_pubkey) VALUES (1, NULL); +UPDATE application_schema SET schema_version = 2 WHERE singleton = 1; diff --git a/core/crates/studio_storage/migrations/V3__profile_cache.sql b/core/crates/studio_storage/migrations/V3__profile_cache.sql @@ -0,0 +1,12 @@ +CREATE TABLE profile_cache ( + subject_pubkey TEXT PRIMARY KEY NOT NULL REFERENCES accounts(pubkey) ON DELETE CASCADE, + event_id TEXT NOT NULL, + event_created_at INTEGER NOT NULL, + name TEXT, + display_name TEXT, + nip05 TEXT, + about TEXT, + picture TEXT, + refreshed_at INTEGER NOT NULL, + refresh_status TEXT NOT NULL CHECK (refresh_status IN ('success', 'offline', 'invalid_data')) +) STRICT; diff --git a/core/crates/studio_storage/migrations/V4__account_namespace.sql b/core/crates/studio_storage/migrations/V4__account_namespace.sql @@ -0,0 +1,6 @@ +CREATE TABLE account_namespace ( + owner_pubkey TEXT NOT NULL REFERENCES accounts(pubkey) ON DELETE CASCADE, + preference_key TEXT NOT NULL CHECK (preference_key IN ('namespace_probe')), + preference_value TEXT NOT NULL CHECK (length(preference_value) <= 4096), + PRIMARY KEY (owner_pubkey, preference_key) +) STRICT; diff --git a/core/crates/studio_storage/migrations/V5__operation_journal.sql b/core/crates/studio_storage/migrations/V5__operation_journal.sql @@ -0,0 +1,8 @@ +CREATE TABLE operation_journal ( + operation_id INTEGER PRIMARY KEY AUTOINCREMENT, + operation_kind TEXT NOT NULL CHECK (operation_kind IN ('add', 'import', 'remove')), + subject_pubkey TEXT NOT NULL, + phase TEXT NOT NULL CHECK (phase IN ('intent_recorded', 'credential_written', 'metadata_committed', 'compensation_pending', 'credential_deleted', 'metadata_deleted')), + updated_at INTEGER NOT NULL, + diagnostic_code TEXT CHECK (diagnostic_code IN ('storage_unavailable', 'keyring_unavailable', 'credential_missing', 'compensation_failed')) +) STRICT; diff --git a/core/crates/studio_storage/migrations/V6__normalized_runtime_schema.sql b/core/crates/studio_storage/migrations/V6__normalized_runtime_schema.sql @@ -0,0 +1,100 @@ +CREATE TABLE account_identities ( + public_key TEXT PRIMARY KEY NOT NULL CHECK ( + length(public_key) = 64 AND public_key = lower(public_key) + ), + npub TEXT NOT NULL UNIQUE CHECK (length(npub) = 63), + label TEXT CHECK (label IS NULL OR length(label) BETWEEN 1 AND 80), + created_at INTEGER NOT NULL CHECK (created_at >= 0), + last_used_at INTEGER CHECK (last_used_at IS NULL OR last_used_at >= 0) +) STRICT; + +CREATE TABLE local_signer_bindings ( + account_public_key TEXT NOT NULL, + binding_public_key TEXT NOT NULL, + binding_kind TEXT NOT NULL CHECK (binding_kind = 'local_secret'), + availability TEXT NOT NULL CHECK ( + availability IN ('available', 'credential_missing', 'store_unavailable') + ), + PRIMARY KEY (account_public_key, binding_public_key), + UNIQUE (account_public_key, binding_kind), + FOREIGN KEY (account_public_key) REFERENCES account_identities(public_key) ON DELETE CASCADE, + CHECK (account_public_key = binding_public_key) +) STRICT; + +CREATE TABLE runtime_state ( + singleton INTEGER PRIMARY KEY CHECK (singleton = 1), + selected_public_key TEXT REFERENCES account_identities(public_key) ON DELETE SET NULL, + active_account_public_key TEXT, + active_binding_public_key TEXT, + session_generation INTEGER NOT NULL DEFAULT 0 CHECK (session_generation >= 0), + FOREIGN KEY (active_account_public_key, active_binding_public_key) + REFERENCES local_signer_bindings(account_public_key, binding_public_key) + ON DELETE SET NULL, + CHECK ( + (active_account_public_key IS NULL AND active_binding_public_key IS NULL) + OR + (active_account_public_key IS NOT NULL AND active_binding_public_key IS NOT NULL) + ) +) STRICT; + +INSERT INTO runtime_state (singleton) VALUES (1); + +CREATE TABLE profile_cache_v6 ( + subject_public_key TEXT PRIMARY KEY NOT NULL + REFERENCES account_identities(public_key) ON DELETE CASCADE, + event_id TEXT NOT NULL CHECK (length(event_id) = 64 AND event_id = lower(event_id)), + event_created_at INTEGER NOT NULL CHECK (event_created_at >= 0), + name TEXT, + display_name TEXT, + nip05 TEXT, + about TEXT, + picture TEXT, + refreshed_at INTEGER NOT NULL CHECK (refreshed_at >= 0), + refresh_status TEXT NOT NULL CHECK ( + refresh_status IN ('success', 'offline', 'invalid_data') + ) +) STRICT; + +CREATE TABLE durable_operations ( + request_id TEXT PRIMARY KEY NOT NULL CHECK (length(request_id) BETWEEN 1 AND 128), + operation_kind TEXT NOT NULL CHECK ( + operation_kind IN ('create', 'import', 'repair', 'remove') + ), + account_public_key TEXT NOT NULL CHECK ( + length(account_public_key) = 64 AND account_public_key = lower(account_public_key) + ), + binding_public_key TEXT NOT NULL CHECK (binding_public_key = account_public_key), + expected_revision INTEGER CHECK (expected_revision IS NULL OR expected_revision >= 0), + phase TEXT NOT NULL CHECK ( + phase IN ( + 'intent_recorded', + 'credential_written', + 'metadata_committed', + 'selection_committed', + 'compensation_pending', + 'credential_deleted', + 'metadata_deleted', + 'finalized' + ) + ), + terminal_outcome TEXT CHECK ( + terminal_outcome IS NULL OR terminal_outcome IN ('completed', 'cancelled', 'failed') + ), + prior_selected_public_key TEXT, + updated_at INTEGER NOT NULL CHECK (updated_at >= 0), + diagnostic_code TEXT CHECK ( + diagnostic_code IS NULL OR diagnostic_code IN ( + 'storage_unavailable', + 'keyring_unavailable', + 'credential_missing', + 'compensation_failed', + 'conflict', + 'expired' + ) + ), + CHECK ( + (phase = 'finalized' AND terminal_outcome IS NOT NULL) + OR + (phase <> 'finalized' AND terminal_outcome IS NULL) + ) +) STRICT; diff --git a/core/crates/studio_storage/migrations/V7__migrate_v5_runtime_data.sql b/core/crates/studio_storage/migrations/V7__migrate_v5_runtime_data.sql @@ -0,0 +1,68 @@ +INSERT INTO account_identities ( + public_key, + npub, + label, + created_at, + last_used_at +) +SELECT pubkey, npub, label, created_at, last_used_at +FROM accounts; + +INSERT INTO local_signer_bindings ( + account_public_key, + binding_public_key, + binding_kind, + availability +) +SELECT pubkey, pubkey, 'local_secret', key_availability +FROM accounts; + +UPDATE runtime_state +SET selected_public_key = ( + SELECT selected_pubkey FROM app_state WHERE singleton = 1 +) +WHERE singleton = 1; + +INSERT INTO profile_cache_v6 ( + subject_public_key, + event_id, + event_created_at, + name, + display_name, + nip05, + about, + picture, + refreshed_at, + refresh_status +) +SELECT + subject_pubkey, + event_id, + event_created_at, + name, + display_name, + nip05, + about, + picture, + refreshed_at, + refresh_status +FROM profile_cache; + +INSERT INTO durable_operations ( + request_id, + operation_kind, + account_public_key, + binding_public_key, + phase, + updated_at, + diagnostic_code +) +SELECT + 'legacy-v5-' || operation_id, + CASE operation_kind WHEN 'add' THEN 'create' ELSE operation_kind END, + subject_pubkey, + subject_pubkey, + phase, + updated_at, + diagnostic_code +FROM operation_journal; diff --git a/core/crates/studio_storage/migrations/V8__normalized_account_preferences.sql b/core/crates/studio_storage/migrations/V8__normalized_account_preferences.sql @@ -0,0 +1,10 @@ +CREATE TABLE account_preferences ( + owner_public_key TEXT NOT NULL REFERENCES account_identities(public_key) ON DELETE CASCADE, + preference_key TEXT NOT NULL CHECK (preference_key = 'namespace_probe'), + preference_value TEXT NOT NULL CHECK (length(preference_value) <= 4096), + PRIMARY KEY (owner_public_key, preference_key) +) STRICT; + +INSERT INTO account_preferences (owner_public_key, preference_key, preference_value) +SELECT owner_pubkey, preference_key, preference_value +FROM account_namespace; diff --git a/core/crates/studio_storage/migrations/V9__durable_operation_receipts.sql b/core/crates/studio_storage/migrations/V9__durable_operation_receipts.sql @@ -0,0 +1,11 @@ +ALTER TABLE durable_operations ADD COLUMN prior_binding_availability TEXT CHECK ( + prior_binding_availability IS NULL OR prior_binding_availability IN ( + 'available', + 'credential_missing', + 'store_unavailable' + ) +); + +ALTER TABLE durable_operations ADD COLUMN resulting_revision INTEGER CHECK ( + resulting_revision IS NULL OR resulting_revision >= 0 +); diff --git a/core/crates/studio_storage/src/account_namespace.rs b/core/crates/studio_storage/src/account_namespace.rs @@ -0,0 +1,189 @@ +use radroots_studio_application::{AccountNamespaceRepository, AccountPreferenceKey}; +use radroots_studio_domain::{PublicKey, SafeError, SafeErrorCode, SafeMessage}; +use rusqlite::{OptionalExtension, params}; + +use crate::Database; + +const MAX_VALUE_CHARS: usize = 4_096; + +impl AccountNamespaceRepository for Database { + fn get_value( + &self, + owner: PublicKey, + key: AccountPreferenceKey, + ) -> Result<Option<String>, SafeError> { + self.connection() + .query_row( + "SELECT preference_value FROM account_preferences \ + WHERE owner_public_key = ?1 AND preference_key = ?2", + params![owner.to_hex(), encode_key(key)], + |row| row.get(0), + ) + .optional() + .map_err(|_| storage_error()) + } + + fn set_value( + &self, + owner: PublicKey, + key: AccountPreferenceKey, + value: &str, + ) -> Result<(), SafeError> { + if value.chars().count() > MAX_VALUE_CHARS || value.chars().any(char::is_control) { + return Err(invalid_preference()); + } + self.connection() + .execute( + "INSERT INTO account_preferences (owner_public_key, preference_key, preference_value) \ + VALUES (?1, ?2, ?3) ON CONFLICT(owner_public_key, preference_key) DO UPDATE SET \ + preference_value = excluded.preference_value", + params![owner.to_hex(), encode_key(key), value], + ) + .map(|_| ()) + .map_err(|_| storage_error()) + } + + fn clear_owner(&self, owner: PublicKey) -> Result<(), SafeError> { + self.connection() + .execute( + "DELETE FROM account_preferences WHERE owner_public_key = ?1", + [owner.to_hex()], + ) + .map(|_| ()) + .map_err(|_| storage_error()) + } +} + +const fn encode_key(key: AccountPreferenceKey) -> &'static str { + match key { + AccountPreferenceKey::NamespaceProbe => "namespace_probe", + } +} + +const fn storage_error() -> SafeError { + SafeError::new( + SafeErrorCode::StorageUnavailable, + SafeMessage::new("The account preference is unavailable."), + ) +} + +const fn invalid_preference() -> SafeError { + SafeError::new( + SafeErrorCode::InvalidAccountMetadata, + SafeMessage::new("The account preference is invalid."), + ) +} + +#[cfg(test)] +mod tests { + use radroots_studio_application::{ + AccountNamespaceRepository, AccountPreferenceKey, AccountRepository, AppStateRepository, + }; + use radroots_studio_domain::{ + AccountCreatedAt, AccountIdentity, AccountSummary, BindingAvailability, LocalSignerBinding, + PublicKey, UnixTimestamp, + }; + + use crate::Database; + + fn public_key(byte: u8) -> PublicKey { + let value = match byte { + 1 => "585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df", + 2 => "e0266e3cfb0d2886f91c73f5f868f3b98273713e5fcd97c081663f5518a4b3af", + _ => "7e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7", + }; + PublicKey::from_hex(value).expect("valid public key") + } + + fn account(byte: u8) -> AccountSummary { + let public_key = public_key(byte); + AccountSummary::new( + AccountIdentity::derive(public_key).expect("identity"), + LocalSignerBinding::new(public_key, BindingAvailability::Available), + None, + AccountCreatedAt::new(UnixTimestamp::from_seconds(i64::from(byte)).expect("time")), + None, + ) + .expect("account") + } + + #[test] + fn namespace_partitions_same_typed_key_by_owner_and_selection() { + let database = Database::in_memory().expect("database"); + let owner_a = public_key(1); + let owner_b = public_key(2); + database.insert_account(&account(1)).expect("account a"); + database.insert_account(&account(2)).expect("account b"); + database + .set_value(owner_a, AccountPreferenceKey::NamespaceProbe, "A") + .expect("set a"); + database + .set_value(owner_b, AccountPreferenceKey::NamespaceProbe, "B") + .expect("set b"); + + database + .save_selected_account(Some(owner_b)) + .expect("select b"); + let selected = database + .load_selected_account() + .expect("selection") + .expect("selected owner"); + assert_eq!( + database + .get_value(selected, AccountPreferenceKey::NamespaceProbe) + .expect("selected value"), + Some("B".to_owned()) + ); + assert_eq!( + database + .get_value(owner_a, AccountPreferenceKey::NamespaceProbe) + .expect("owner a value"), + Some("A".to_owned()) + ); + } + + #[test] + fn namespace_updates_and_cascades_with_owner_removal() { + let database = Database::in_memory().expect("database"); + let owner = public_key(3); + database.insert_account(&account(3)).expect("account"); + database + .set_value(owner, AccountPreferenceKey::NamespaceProbe, "before") + .expect("set"); + database + .set_value(owner, AccountPreferenceKey::NamespaceProbe, "after") + .expect("update"); + assert_eq!( + database + .get_value(owner, AccountPreferenceKey::NamespaceProbe) + .expect("value"), + Some("after".to_owned()) + ); + + database.remove_account(owner).expect("remove"); + assert_eq!( + database + .get_value(owner, AccountPreferenceKey::NamespaceProbe) + .expect("deleted value"), + None + ); + } + + #[test] + fn namespace_rejects_oversized_and_control_character_values() { + let database = Database::in_memory().expect("database"); + let owner = public_key(3); + database.insert_account(&account(3)).expect("account"); + let oversized = "a".repeat(super::MAX_VALUE_CHARS + 1); + assert!( + database + .set_value(owner, AccountPreferenceKey::NamespaceProbe, &oversized) + .is_err() + ); + assert!( + database + .set_value(owner, AccountPreferenceKey::NamespaceProbe, "line\nbreak") + .is_err() + ); + } +} diff --git a/core/crates/studio_storage/src/accounts.rs b/core/crates/studio_storage/src/accounts.rs @@ -0,0 +1,516 @@ +use radroots_studio_application::{AccountRepository, AppStateRepository}; +use radroots_studio_domain::{ + AccountCreatedAt, AccountIdentity, AccountLabel, AccountSummary, BindingAvailability, + LocalSignerBinding, PublicKey, SafeError, SafeErrorCode, SafeMessage, UnixTimestamp, +}; +use rusqlite::{OptionalExtension, Row, params}; + +use crate::Database; + +impl AccountRepository for Database { + fn list_accounts(&self) -> Result<Vec<AccountSummary>, SafeError> { + let connection = self.connection(); + let mut statement = connection + .prepare( + "SELECT identity.public_key, identity.npub, binding.binding_kind, \ + binding.availability, identity.label, identity.created_at, identity.last_used_at \ + FROM account_identities AS identity \ + JOIN local_signer_bindings AS binding \ + ON binding.account_public_key = identity.public_key \ + ORDER BY identity.created_at ASC, identity.public_key ASC", + ) + .map_err(|_| storage_error())?; + let rows = statement + .query_map([], decode_account) + .map_err(|_| storage_error())?; + rows.map(|row| row.map_err(|_| corrupt_storage_error())) + .collect() + } + + fn find_account(&self, public_key: PublicKey) -> Result<Option<AccountSummary>, SafeError> { + self.connection() + .query_row( + "SELECT identity.public_key, identity.npub, binding.binding_kind, \ + binding.availability, identity.label, identity.created_at, identity.last_used_at \ + FROM account_identities AS identity \ + JOIN local_signer_bindings AS binding \ + ON binding.account_public_key = identity.public_key \ + WHERE identity.public_key = ?1", + [public_key.to_hex()], + decode_account, + ) + .optional() + .map_err(|_| storage_error()) + } + + fn insert_account(&self, account: &AccountSummary) -> Result<(), SafeError> { + let encoded = EncodedAccount::from(account); + let mut connection = self.connection(); + let transaction = connection.transaction().map_err(|_| storage_error())?; + let result = transaction.execute( + "INSERT INTO account_identities (public_key, npub, label, created_at, last_used_at) \ + VALUES (?1, ?2, ?3, ?4, ?5)", + params![ + encoded.public_key, + encoded.npub, + encoded.label, + encoded.created_at, + encoded.last_used_at + ], + ); + match result { + Ok(1) => {} + Err(error) if is_constraint_violation(&error) => return Err(account_exists()), + Ok(_) | Err(_) => return Err(storage_error()), + } + if transaction + .execute( + "INSERT INTO local_signer_bindings (account_public_key, binding_public_key, \ + binding_kind, availability) VALUES (?1, ?1, ?2, ?3)", + params![ + encoded.public_key, + encoded.signer_kind, + encoded.key_availability + ], + ) + .map_err(|_| storage_error())? + != 1 + { + return Err(storage_error()); + } + transaction.commit().map_err(|_| storage_error()) + } + + fn update_account(&self, account: &AccountSummary) -> Result<(), SafeError> { + let encoded = EncodedAccount::from(account); + let mut connection = self.connection(); + let transaction = connection.transaction().map_err(|_| storage_error())?; + let identity_rows = transaction + .execute( + "UPDATE account_identities SET npub = ?2, label = ?5, created_at = ?6, \ + last_used_at = ?7 WHERE public_key = ?1", + params![ + encoded.public_key, + encoded.npub, + encoded.signer_kind, + encoded.key_availability, + encoded.label, + encoded.created_at, + encoded.last_used_at, + ], + ) + .map_err(|_| storage_error())?; + if identity_rows == 0 { + return Err(account_not_found()); + } + if identity_rows != 1 { + return Err(storage_error()); + } + let binding_rows = transaction + .execute( + "UPDATE local_signer_bindings SET binding_kind = ?2, availability = ?3 \ + WHERE account_public_key = ?1 AND binding_public_key = ?1", + params![ + encoded.public_key, + encoded.signer_kind, + encoded.key_availability + ], + ) + .map_err(|_| storage_error())?; + if binding_rows != 1 { + return Err(corrupt_storage_error()); + } + transaction.commit().map_err(|_| storage_error()) + } + + fn remove_account(&self, public_key: PublicKey) -> Result<(), SafeError> { + match self.connection().execute( + "DELETE FROM account_identities WHERE public_key = ?1", + [public_key.to_hex()], + ) { + Ok(1) => Ok(()), + Ok(0) => Err(account_not_found()), + Ok(_) | Err(_) => Err(storage_error()), + } + } +} + +impl AppStateRepository for Database { + fn load_selected_account(&self) -> Result<Option<PublicKey>, SafeError> { + let value = self + .connection() + .query_row( + "SELECT selected_public_key FROM runtime_state WHERE singleton = 1", + [], + |row| row.get::<_, Option<String>>(0), + ) + .map_err(|_| corrupt_storage_error())?; + value + .map(|hex| PublicKey::from_hex(&hex).map_err(|_| corrupt_storage_error())) + .transpose() + } + + fn save_selected_account(&self, public_key: Option<PublicKey>) -> Result<(), SafeError> { + let mut connection = self.connection(); + let transaction = connection.transaction().map_err(|_| storage_error())?; + if let Some(public_key) = public_key { + let exists = transaction + .query_row( + "SELECT EXISTS(SELECT 1 FROM account_identities WHERE public_key = ?1)", + [public_key.to_hex()], + |row| row.get::<_, bool>(0), + ) + .map_err(|_| storage_error())?; + if !exists { + return Err(account_not_found()); + } + } + let rows = transaction + .execute( + "UPDATE runtime_state SET selected_public_key = ?1 WHERE singleton = 1", + [public_key.map(PublicKey::to_hex)], + ) + .map_err(|_| storage_error())?; + if rows != 1 { + return Err(corrupt_storage_error()); + } + transaction.commit().map_err(|_| storage_error()) + } +} + +struct EncodedAccount { + public_key: String, + npub: String, + signer_kind: &'static str, + key_availability: &'static str, + label: Option<String>, + created_at: i64, + last_used_at: Option<i64>, +} + +impl From<&AccountSummary> for EncodedAccount { + fn from(account: &AccountSummary) -> Self { + Self { + public_key: account.public_key().to_hex(), + npub: account.npub().as_str().to_owned(), + signer_kind: "local_secret", + key_availability: encode_key_availability(account.signer().availability()), + label: account.label().map(|label| label.as_str().to_owned()), + created_at: account.created_at().timestamp().as_seconds(), + last_used_at: account.last_used_at().map(UnixTimestamp::as_seconds), + } + } +} + +fn decode_account(row: &Row<'_>) -> rusqlite::Result<AccountSummary> { + let public_key = + PublicKey::from_hex(row.get::<_, String>(0)?.as_str()).map_err(|_| invalid_column(0))?; + let npub: String = row.get(1)?; + if row.get::<_, String>(2)?.as_str() != "local_secret" { + return Err(invalid_column(2)); + } + let key_availability = decode_key_availability(row.get::<_, String>(3)?.as_str())?; + let label = row + .get::<_, Option<String>>(4)? + .map(|value| AccountLabel::parse(&value).map_err(|_| invalid_column(4))) + .transpose()?; + let created_at = UnixTimestamp::from_seconds(row.get(5)?).ok_or_else(|| invalid_column(5))?; + let last_used_at = row + .get::<_, Option<i64>>(6)? + .map(|value| UnixTimestamp::from_seconds(value).ok_or_else(|| invalid_column(6))) + .transpose()?; + + AccountSummary::new( + AccountIdentity::verify(public_key, npub).map_err(|_| invalid_column(1))?, + LocalSignerBinding::new(public_key, key_availability), + label, + AccountCreatedAt::new(created_at), + last_used_at, + ) + .map_err(|_| invalid_column(0)) +} + +const fn encode_key_availability(value: BindingAvailability) -> &'static str { + match value { + BindingAvailability::Available => "available", + BindingAvailability::CredentialMissing => "credential_missing", + BindingAvailability::StoreUnavailable => "store_unavailable", + } +} + +fn decode_key_availability(value: &str) -> rusqlite::Result<BindingAvailability> { + match value { + "available" => Ok(BindingAvailability::Available), + "credential_missing" => Ok(BindingAvailability::CredentialMissing), + "store_unavailable" => Ok(BindingAvailability::StoreUnavailable), + _ => Err(invalid_column(3)), + } +} + +fn invalid_column(index: usize) -> rusqlite::Error { + rusqlite::Error::InvalidColumnType( + index, + "public account metadata".to_owned(), + rusqlite::types::Type::Text, + ) +} + +fn is_constraint_violation(error: &rusqlite::Error) -> bool { + matches!( + error, + rusqlite::Error::SqliteFailure( + rusqlite::ffi::Error { + code: rusqlite::ErrorCode::ConstraintViolation, + .. + }, + _ + ) + ) +} + +const fn storage_error() -> SafeError { + SafeError::new( + SafeErrorCode::StorageUnavailable, + SafeMessage::new("The application database is unavailable."), + ) +} + +const fn corrupt_storage_error() -> SafeError { + SafeError::new( + SafeErrorCode::StorageCorrupt, + SafeMessage::new("The application database could not be read."), + ) +} + +const fn account_exists() -> SafeError { + SafeError::new( + SafeErrorCode::AccountAlreadyExists, + SafeMessage::new("The Nostr account is already saved."), + ) +} + +const fn account_not_found() -> SafeError { + SafeError::new( + SafeErrorCode::AccountNotFound, + SafeMessage::new("The account was not found."), + ) +} + +#[cfg(test)] +mod tests { + use std::fs; + + use radroots_studio_application::{AccountRepository, AppStateRepository}; + use radroots_studio_domain::{ + AccountCreatedAt, AccountIdentity, AccountLabel, AccountSummary, BindingAvailability, + LocalSignerBinding, PublicKey, SafeErrorCode, UnixTimestamp, + }; + use tempfile::tempdir; + + use crate::Database; + + fn public_key(key_byte: u8) -> PublicKey { + let value = match key_byte { + 1 => "585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df", + 2 => "e0266e3cfb0d2886f91c73f5f868f3b98273713e5fcd97c081663f5518a4b3af", + _ => "7e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7", + }; + PublicKey::from_hex(value).expect("valid public key") + } + + fn account(key_byte: u8, created_at: i64) -> AccountSummary { + let public_key = public_key(key_byte); + AccountSummary::new( + AccountIdentity::derive(public_key).expect("identity"), + LocalSignerBinding::new(public_key, BindingAvailability::Available), + Some(AccountLabel::parse("Farm account").expect("valid label")), + AccountCreatedAt::new( + UnixTimestamp::from_seconds(created_at).expect("valid timestamp"), + ), + None, + ) + .expect("account") + } + + #[test] + fn accounts_insert_list_update_and_reject_duplicates() { + let database = Database::in_memory().expect("database"); + let first = account(1, 20); + let second = account(2, 10); + + database.insert_account(&first).expect("insert first"); + database.insert_account(&second).expect("insert second"); + let duplicate = database.insert_account(&first).expect_err("duplicate"); + + assert_eq!(duplicate.code(), SafeErrorCode::AccountAlreadyExists); + assert_eq!( + database.list_accounts().expect("list"), + vec![second, first.clone()] + ); + assert_eq!( + database.find_account(first.public_key()).expect("find"), + Some(first) + ); + } + + #[test] + fn accounts_and_selection_survive_restart_without_secret_text() { + let directory = tempdir().expect("temporary directory"); + let path = directory.path().join("studio.sqlite3"); + let account = account(3, 30); + + { + let database = Database::open(&path).expect("database"); + database.insert_account(&account).expect("insert"); + database + .save_selected_account(Some(account.public_key())) + .expect("select"); + } + let reopened = Database::open(&path).expect("reopen"); + + assert_eq!( + reopened.list_accounts().expect("list"), + vec![account.clone()] + ); + assert_eq!( + reopened.load_selected_account().expect("selection"), + Some(account.public_key()) + ); + let bytes = fs::read(path).expect("database bytes"); + assert!(!String::from_utf8_lossy(&bytes).contains("nsec1known-test-secret")); + } + + #[test] + fn selection_requires_an_existing_account_and_clears_on_delete() { + let database = Database::in_memory().expect("database"); + let account = account(4, 40); + + let missing = database + .save_selected_account(Some(account.public_key())) + .expect_err("missing account"); + assert_eq!(missing.code(), SafeErrorCode::AccountNotFound); + + database.insert_account(&account).expect("insert"); + database + .save_selected_account(Some(account.public_key())) + .expect("select"); + database + .remove_account(account.public_key()) + .expect("remove"); + + assert_eq!(database.load_selected_account().expect("selection"), None); + } + + #[test] + fn account_mutations_reject_missing_and_corrupt_rows() { + let database = Database::in_memory().expect("database"); + let missing = account(3, 30); + assert_eq!( + database + .update_account(&missing) + .expect_err("missing update") + .code(), + SafeErrorCode::AccountNotFound + ); + assert_eq!( + database + .remove_account(missing.public_key()) + .expect_err("missing removal") + .code(), + SafeErrorCode::AccountNotFound + ); + assert_eq!( + database.find_account(missing.public_key()).expect("find"), + None + ); + + database.insert_account(&missing).expect("insert"); + database.update_account(&missing).expect("update"); + database + .connection() + .execute( + "DELETE FROM local_signer_bindings WHERE account_public_key = ?1", + [missing.public_key().to_hex()], + ) + .expect("delete binding"); + assert_eq!( + database + .update_account(&missing) + .expect_err("missing binding must fail") + .code(), + SafeErrorCode::StorageCorrupt + ); + database + .connection() + .execute( + "INSERT INTO local_signer_bindings (account_public_key, binding_public_key, binding_kind, availability) VALUES (?1, ?1, 'local_secret', 'available')", + [missing.public_key().to_hex()], + ) + .expect("restore binding"); + database + .connection() + .pragma_update(None, "ignore_check_constraints", "ON") + .expect("disable check constraints for corruption fixture"); + database + .connection() + .execute( + "UPDATE local_signer_bindings SET binding_kind = 'remote' WHERE account_public_key = ?1", + [missing.public_key().to_hex()], + ) + .expect("corrupt binding kind"); + assert_eq!( + database + .list_accounts() + .expect_err("corrupt binding must fail") + .code(), + SafeErrorCode::StorageCorrupt + ); + + let database = Database::in_memory().expect("database"); + database.insert_account(&missing).expect("insert"); + database + .connection() + .pragma_update(None, "ignore_check_constraints", "ON") + .expect("disable check constraints for corruption fixture"); + database + .connection() + .execute( + "UPDATE local_signer_bindings SET availability = 'invalid' WHERE account_public_key = ?1", + [missing.public_key().to_hex()], + ) + .expect("corrupt availability"); + assert_eq!( + database + .find_account(missing.public_key()) + .expect_err("corrupt availability must fail") + .code(), + SafeErrorCode::StorageUnavailable + ); + + let database = Database::in_memory().expect("database"); + database + .connection() + .execute("DELETE FROM runtime_state", []) + .expect("delete runtime singleton"); + assert_eq!( + database + .save_selected_account(None) + .expect_err("missing runtime singleton must fail") + .code(), + SafeErrorCode::StorageCorrupt + ); + + let read_only = Database::in_memory().expect("read-only database"); + read_only + .connection() + .pragma_update(None, "query_only", "ON") + .expect("enable query-only mode"); + assert_eq!( + read_only + .insert_account(&missing) + .expect_err("non-constraint insertion failure must fail closed") + .code(), + SafeErrorCode::StorageUnavailable + ); + } +} diff --git a/core/crates/studio_storage/src/compatibility.rs b/core/crates/studio_storage/src/compatibility.rs @@ -0,0 +1,474 @@ +use std::path::Path; + +use radroots_studio_domain::{ + AccountIdentity, PersistedPublicKeyClassification, SafeError, SafeErrorCode, SafeMessage, + classify_persisted_public_key, +}; +use rusqlite::{Connection, OpenFlags}; +use sha2::{Digest, Sha256}; + +use crate::CURRENT_SCHEMA_VERSION; + +const KNOWN_TABLES: &[(&str, u32)] = &[ + ("application_schema", 1), + ("accounts", 2), + ("app_state", 2), + ("profile_cache", 3), + ("account_namespace", 4), + ("operation_journal", 5), + ("account_identities", 6), + ("local_signer_bindings", 6), + ("runtime_state", 6), + ("profile_cache_v6", 6), + ("durable_operations", 6), + ("account_preferences", 8), + ("installation_identity", 10), +]; + +const PUBLIC_KEY_COLUMNS: &[(&str, &str)] = &[ + ("accounts", "pubkey"), + ("app_state", "selected_pubkey"), + ("profile_cache", "subject_pubkey"), + ("account_namespace", "owner_pubkey"), + ("operation_journal", "subject_pubkey"), + ("account_identities", "public_key"), + ("local_signer_bindings", "account_public_key"), + ("local_signer_bindings", "binding_public_key"), + ("runtime_state", "selected_public_key"), + ("runtime_state", "active_account_public_key"), + ("runtime_state", "active_binding_public_key"), + ("profile_cache_v6", "subject_public_key"), + ("durable_operations", "account_public_key"), + ("durable_operations", "binding_public_key"), + ("durable_operations", "prior_selected_public_key"), + ("account_preferences", "owner_public_key"), +]; + +#[derive(Clone, Debug, Eq, PartialEq)] +pub enum DatabasePreflight { + Fresh, + Ready { + schema_version: u32, + }, + Quarantined { + schema_version: u32, + issues: Vec<PersistedIdentityIssue>, + }, +} + +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum PersistedIdentityIssueKind { + MalformedEncoding, + NonCanonicalEncoding, + InvalidCurvePoint, + DisplayIdentityMismatch, +} + +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct PersistedIdentityIssue { + table: &'static str, + column: &'static str, + row_id: i64, + kind: PersistedIdentityIssueKind, + fingerprint: [u8; 32], +} + +impl PersistedIdentityIssue { + #[must_use] + pub const fn table(&self) -> &'static str { + self.table + } + + #[must_use] + pub const fn column(&self) -> &'static str { + self.column + } + + #[must_use] + pub const fn row_id(&self) -> i64 { + self.row_id + } + + #[must_use] + pub const fn kind(&self) -> PersistedIdentityIssueKind { + self.kind + } + + #[must_use] + pub const fn fingerprint(&self) -> &[u8; 32] { + &self.fingerprint + } +} + +pub(crate) fn preflight(path: &Path) -> Result<DatabasePreflight, SafeError> { + match std::fs::symlink_metadata(path) { + Ok(metadata) if metadata.file_type().is_symlink() || !metadata.is_file() => { + return Err(corrupt_storage_error()); + } + Ok(_) => {} + Err(error) if error.kind() == std::io::ErrorKind::NotFound => { + return Ok(DatabasePreflight::Fresh); + } + Err(_) => return Err(corrupt_storage_error()), + } + let flags = OpenFlags::SQLITE_OPEN_READ_ONLY + | OpenFlags::SQLITE_OPEN_NO_MUTEX + | OpenFlags::SQLITE_OPEN_NOFOLLOW; + let connection = + Connection::open_with_flags(path, flags).map_err(|_| corrupt_storage_error())?; + connection + .pragma_update(None, "trusted_schema", "OFF") + .map_err(|_| corrupt_storage_error())?; + let integrity: String = connection + .pragma_query_value(None, "quick_check", |row| row.get(0)) + .map_err(|_| corrupt_storage_error())?; + if integrity != "ok" { + return Err(corrupt_storage_error()); + } + let schema_version = schema_version(&connection)?; + if schema_version == 0 || schema_version > CURRENT_SCHEMA_VERSION { + return Err(unsupported_schema_error()); + } + validate_schema_inventory(&connection, schema_version)?; + + let mut issues = Vec::new(); + for &(table, column) in PUBLIC_KEY_COLUMNS { + if column_exists(&connection, table, column)? { + scan_public_key_column(&connection, table, column, &mut issues)?; + } + } + scan_display_identities(&connection, "accounts", "pubkey", "npub", &mut issues)?; + scan_display_identities( + &connection, + "account_identities", + "public_key", + "npub", + &mut issues, + )?; + issues.sort_by_key(|issue| (issue.table, issue.column, issue.row_id)); + if issues.is_empty() { + Ok(DatabasePreflight::Ready { schema_version }) + } else { + Ok(DatabasePreflight::Quarantined { + schema_version, + issues, + }) + } +} + +fn schema_version(connection: &Connection) -> Result<u32, SafeError> { + if !table_exists(connection, "refinery_schema_history")? { + return Err(unsupported_schema_error()); + } + connection + .query_row( + "SELECT COALESCE(MAX(version), 0) FROM refinery_schema_history", + [], + |row| row.get(0), + ) + .map_err(|_| corrupt_storage_error()) +} + +fn validate_schema_inventory(connection: &Connection, version: u32) -> Result<(), SafeError> { + for &(table, introduced) in KNOWN_TABLES { + let present = table_exists(connection, table)?; + if present != (version >= introduced) { + return Err(corrupt_storage_error()); + } + } + let mut statement = connection + .prepare( + "SELECT name FROM sqlite_master WHERE type = 'table' AND name NOT LIKE 'sqlite_%' AND name <> 'refinery_schema_history'", + ) + .map_err(|_| corrupt_storage_error())?; + let names = statement + .query_map([], |row| row.get::<_, String>(0)) + .map_err(|_| corrupt_storage_error())?; + for name in names { + let name = name.map_err(|_| corrupt_storage_error())?; + if !KNOWN_TABLES.iter().any(|(known, _)| *known == name) { + return Err(corrupt_storage_error()); + } + } + Ok(()) +} + +fn scan_public_key_column( + connection: &Connection, + table: &'static str, + column: &'static str, + issues: &mut Vec<PersistedIdentityIssue>, +) -> Result<(), SafeError> { + let sql = format!("SELECT rowid, {column} FROM {table} WHERE {column} IS NOT NULL"); + let mut statement = connection + .prepare(&sql) + .map_err(|_| corrupt_storage_error())?; + let rows = statement + .query_map([], |row| { + Ok((row.get::<_, i64>(0)?, row.get::<_, String>(1)?)) + }) + .map_err(|_| corrupt_storage_error())?; + for row in rows { + let (row_id, value) = row.map_err(|_| corrupt_storage_error())?; + let kind = match classify_persisted_public_key(&value) { + PersistedPublicKeyClassification::Canonical(_) => continue, + PersistedPublicKeyClassification::MalformedEncoding => { + PersistedIdentityIssueKind::MalformedEncoding + } + PersistedPublicKeyClassification::NonCanonicalEncoding => { + PersistedIdentityIssueKind::NonCanonicalEncoding + } + PersistedPublicKeyClassification::InvalidCurvePoint => { + PersistedIdentityIssueKind::InvalidCurvePoint + } + }; + issues.push(issue(table, column, row_id, kind, &value)); + } + Ok(()) +} + +fn scan_display_identities( + connection: &Connection, + table: &'static str, + key_column: &'static str, + npub_column: &'static str, + issues: &mut Vec<PersistedIdentityIssue>, +) -> Result<(), SafeError> { + if !column_exists(connection, table, key_column)? + || !column_exists(connection, table, npub_column)? + { + return Ok(()); + } + let sql = format!("SELECT rowid, {key_column}, {npub_column} FROM {table}"); + let mut statement = connection + .prepare(&sql) + .map_err(|_| corrupt_storage_error())?; + let rows = statement + .query_map([], |row| { + Ok(( + row.get::<_, i64>(0)?, + row.get::<_, String>(1)?, + row.get::<_, String>(2)?, + )) + }) + .map_err(|_| corrupt_storage_error())?; + for row in rows { + let (row_id, key, npub) = row.map_err(|_| corrupt_storage_error())?; + let PersistedPublicKeyClassification::Canonical(public_key) = + classify_persisted_public_key(&key) + else { + continue; + }; + if AccountIdentity::verify(public_key, npub.clone()).is_err() { + issues.push(issue( + table, + npub_column, + row_id, + PersistedIdentityIssueKind::DisplayIdentityMismatch, + &npub, + )); + } + } + Ok(()) +} + +fn issue( + table: &'static str, + column: &'static str, + row_id: i64, + kind: PersistedIdentityIssueKind, + value: &str, +) -> PersistedIdentityIssue { + PersistedIdentityIssue { + table, + column, + row_id, + kind, + fingerprint: Sha256::digest(value.as_bytes()).into(), + } +} + +fn table_exists(connection: &Connection, table: &str) -> Result<bool, SafeError> { + connection + .query_row( + "SELECT EXISTS(SELECT 1 FROM sqlite_master WHERE type = 'table' AND name = ?1)", + [table], + |row| row.get(0), + ) + .map_err(|_| corrupt_storage_error()) +} + +fn column_exists(connection: &Connection, table: &str, column: &str) -> Result<bool, SafeError> { + if !table_exists(connection, table)? { + return Ok(false); + } + let sql = format!("SELECT EXISTS(SELECT 1 FROM pragma_table_info('{table}') WHERE name = ?1)"); + connection + .query_row(&sql, [column], |row| row.get(0)) + .map_err(|_| corrupt_storage_error()) +} + +const fn corrupt_storage_error() -> SafeError { + SafeError::new( + SafeErrorCode::StorageCorrupt, + SafeMessage::new("The application database could not be read."), + ) +} + +const fn unsupported_schema_error() -> SafeError { + SafeError::new( + SafeErrorCode::UnsupportedSchemaVersion, + SafeMessage::new("The application database schema is not supported."), + ) +} + +pub(crate) const fn quarantined_storage_error() -> SafeError { + SafeError::new( + SafeErrorCode::StorageQuarantined, + SafeMessage::new("The application database requires authenticated repair."), + ) +} + +#[cfg(test)] +mod tests { + use rusqlite::{Connection, params}; + use tempfile::tempdir; + + use super::{ + DatabasePreflight, PersistedIdentityIssueKind, column_exists, preflight, + scan_display_identities, scan_public_key_column, + }; + use crate::Database; + use radroots_studio_domain::{AccountIdentity, PublicKey, SafeErrorCode}; + + #[test] + fn preflight_rejects_non_files_missing_schema_zero_version_and_unknown_tables() { + let directory = tempdir().expect("temporary directory"); + let missing = directory.path().join("missing.sqlite3"); + assert_eq!( + preflight(&missing).expect("fresh preflight"), + DatabasePreflight::Fresh + ); + assert_eq!( + preflight(directory.path()) + .expect_err("directory must fail") + .code(), + SafeErrorCode::StorageCorrupt + ); + let regular_parent = directory.path().join("regular-parent"); + std::fs::write(&regular_parent, b"not a directory").expect("write regular parent"); + assert_eq!( + preflight(&regular_parent.join("nested.sqlite3")) + .expect_err("non-directory parent must fail") + .code(), + SafeErrorCode::StorageCorrupt + ); + + let no_schema = directory.path().join("no-schema.sqlite3"); + drop(Connection::open(&no_schema).expect("blank sqlite database")); + assert_eq!( + preflight(&no_schema) + .expect_err("missing schema history") + .code(), + SafeErrorCode::UnsupportedSchemaVersion + ); + + let zero_schema = directory.path().join("zero-schema.sqlite3"); + let connection = Connection::open(&zero_schema).expect("zero schema database"); + connection + .execute( + "CREATE TABLE refinery_schema_history (version INTEGER NOT NULL)", + [], + ) + .expect("schema history"); + connection + .execute( + "INSERT INTO refinery_schema_history (version) VALUES (0)", + [], + ) + .expect("zero version"); + drop(connection); + assert_eq!( + preflight(&zero_schema) + .expect_err("zero schema version") + .code(), + SafeErrorCode::UnsupportedSchemaVersion + ); + + let unknown = directory.path().join("unknown-table.sqlite3"); + drop(Database::open(&unknown).expect("current database")); + let connection = Connection::open(&unknown).expect("open current database"); + connection + .execute("CREATE TABLE ungoverned_table (value INTEGER)", []) + .expect("unknown table"); + drop(connection); + assert_eq!( + preflight(&unknown) + .expect_err("unknown table must fail") + .code(), + SafeErrorCode::StorageCorrupt + ); + } + + #[test] + fn identity_scans_classify_all_persisted_key_and_display_failures() { + let connection = Connection::open_in_memory().expect("database"); + connection + .execute("CREATE TABLE identities (public_key TEXT, npub TEXT)", []) + .expect("identity table"); + let canonical = + PublicKey::from_hex("585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df") + .expect("canonical key"); + let npub = AccountIdentity::derive(canonical) + .expect("identity") + .npub() + .as_str() + .to_owned(); + let values = [ + (canonical.to_hex(), npub), + (canonical.to_hex().to_uppercase(), "invalid-npub".to_owned()), + ("bad".to_owned(), "invalid-npub".to_owned()), + ("00".repeat(32), "invalid-npub".to_owned()), + (canonical.to_hex(), "invalid-npub".to_owned()), + ]; + for (public_key, npub) in values { + connection + .execute( + "INSERT INTO identities (public_key, npub) VALUES (?1, ?2)", + params![public_key, npub], + ) + .expect("identity row"); + } + + assert!(column_exists(&connection, "identities", "public_key").expect("column")); + assert!(!column_exists(&connection, "missing", "public_key").expect("missing table")); + assert!(!column_exists(&connection, "identities", "missing").expect("missing column")); + let mut issues = Vec::new(); + scan_public_key_column(&connection, "identities", "public_key", &mut issues) + .expect("scan public keys"); + scan_display_identities(&connection, "identities", "public_key", "npub", &mut issues) + .expect("scan display identities"); + scan_display_identities(&connection, "missing", "public_key", "npub", &mut issues) + .expect("skip missing table"); + connection + .execute("CREATE TABLE key_only (public_key TEXT)", []) + .expect("key-only table"); + scan_display_identities(&connection, "key_only", "public_key", "npub", &mut issues) + .expect("skip missing display column"); + + for kind in [ + PersistedIdentityIssueKind::MalformedEncoding, + PersistedIdentityIssueKind::NonCanonicalEncoding, + PersistedIdentityIssueKind::InvalidCurvePoint, + PersistedIdentityIssueKind::DisplayIdentityMismatch, + ] { + assert!(issues.iter().any(|issue| issue.kind() == kind)); + } + for issue in &issues { + assert_eq!(issue.table(), "identities"); + assert!(matches!(issue.column(), "public_key" | "npub")); + assert!(issue.row_id() > 0); + assert_ne!(issue.fingerprint(), &[0_u8; 32]); + } + } +} diff --git a/core/crates/studio_storage/src/db.rs b/core/crates/studio_storage/src/db.rs @@ -0,0 +1,907 @@ +use std::fs::{self, File, OpenOptions}; +use std::ops::{Deref, DerefMut}; +use std::path::{Path, PathBuf}; +use std::sync::{Mutex, MutexGuard}; +use std::time::Duration; + +use fs2::FileExt; +use radroots_studio_domain::{SafeError, SafeErrorCode, SafeMessage}; +use refinery::embed_migrations; +use rusqlite::{Connection, OpenFlags}; + +use crate::compatibility::{DatabasePreflight, preflight, quarantined_storage_error}; +use crate::recovery::MigrationRecovery; +use crate::repair::{ + QuarantineExportReceipt, RepairAuthorization, RepairCandidate, authenticate_candidate, + export_quarantined, install_candidate, +}; + +pub const CURRENT_SCHEMA_VERSION: u32 = 10; + +mod migrations { + use super::embed_migrations; + + embed_migrations!("migrations"); +} + +pub struct Database { + connection: Mutex<Connection>, + path: Option<PathBuf>, + _ownership: Option<WritableOwnership>, +} + +pub(crate) struct DatabaseConnection<'a> { + connection: MutexGuard<'a, Connection>, + path: Option<&'a Path>, +} + +struct WritableOwnership { + _file: File, +} + +impl Database { + /// Opens, configures, and migrates a file-backed `SQLite` database. + /// + /// # Errors + /// + /// Returns a safe storage error when the file, connection configuration, + /// permission update, or migration cannot complete. + pub fn open(path: &Path) -> Result<Self, SafeError> { + let preflight = preflight(path)?; + if matches!(&preflight, DatabasePreflight::Quarantined { .. }) { + return Err(quarantined_storage_error()); + } + let parent = path.parent().ok_or_else(storage_error)?; + create_secure_directory(parent)?; + restrict_sqlite_sidecars(path)?; + let ownership = WritableOwnership::acquire(path)?; + let recovery_source_schema = match &preflight { + DatabasePreflight::Ready { schema_version } + if *schema_version < CURRENT_SCHEMA_VERSION => + { + Some(*schema_version) + } + _ => None, + }; + let recovery = match preflight { + DatabasePreflight::Ready { schema_version } + if schema_version < CURRENT_SCHEMA_VERSION => + { + Some(MigrationRecovery::prepare( + path, + schema_version, + CURRENT_SCHEMA_VERSION, + )?) + } + DatabasePreflight::Fresh | DatabasePreflight::Ready { .. } => None, + DatabasePreflight::Quarantined { .. } => unreachable!("handled above"), + }; + let flags = OpenFlags::SQLITE_OPEN_READ_WRITE + | OpenFlags::SQLITE_OPEN_CREATE + | OpenFlags::SQLITE_OPEN_NO_MUTEX + | OpenFlags::SQLITE_OPEN_NOFOLLOW; + let mut connection = + Connection::open_with_flags(path, flags).map_err(|_| storage_error())?; + configure(&connection).map_err(|_| corrupt_storage_error())?; + if migrations::migrations::runner() + .run(&mut connection) + .is_err() + { + drop(connection); + if let Some(source_schema) = recovery_source_schema { + MigrationRecovery::restore(path, source_schema, CURRENT_SCHEMA_VERSION)?; + } + return Err(corrupt_storage_error()); + } + let schema_version = connection + .query_row( + "SELECT COALESCE(MAX(version), 0) FROM refinery_schema_history", + [], + |row| row.get::<_, u32>(0), + ) + .map_err(|_| corrupt_storage_error())?; + if schema_version != CURRENT_SCHEMA_VERSION { + return Err(corrupt_storage_error()); + } + restrict_file_permissions(path)?; + restrict_sqlite_sidecars(path)?; + if let Some(recovery) = recovery { + recovery.finish(schema_version)?; + } + Ok(Self { + connection: Mutex::new(connection), + path: Some(path.to_path_buf()), + _ownership: Some(ownership), + }) + } + + /// Opens and migrates an isolated in-memory `SQLite` database. + /// + /// # Errors + /// + /// Returns a safe storage error when configuration or migration fails. + pub fn in_memory() -> Result<Self, SafeError> { + let mut connection = Connection::open_in_memory().map_err(|_| storage_error())?; + configure(&connection)?; + migrations::migrations::runner() + .run(&mut connection) + .map_err(|_| corrupt_storage_error())?; + Ok(Self { + connection: Mutex::new(connection), + path: None, + _ownership: None, + }) + } + + /// Inspects schema and persisted identities without mutating the database. + /// + /// # Errors + /// + /// Returns a safe corrupt or unsupported-schema error when the database + /// cannot be classified. + pub fn preflight(path: &Path) -> Result<DatabasePreflight, SafeError> { + preflight(path) + } + + /// Verifies the authenticated, immutable backup retained for a migration. + /// + /// # Errors + /// + /// Returns a safe backup error when any manifest, digest, authentication + /// tag, schema identity, or SQLite integrity check fails. + pub fn verify_migration_backup(path: &Path, source_schema: u32) -> Result<(), SafeError> { + MigrationRecovery::verify_evidence(path, source_schema, CURRENT_SCHEMA_VERSION) + } + + /// Restores an authenticated pre-migration backup while retaining the + /// displaced database as recovery evidence. + /// + /// # Errors + /// + /// Returns a safe storage or backup error without replacing the database + /// when authentication or the atomic replacement fails. + pub fn restore_migration_backup(path: &Path, source_schema: u32) -> Result<(), SafeError> { + let _ownership = WritableOwnership::acquire(path)?; + MigrationRecovery::restore(path, source_schema, CURRENT_SCHEMA_VERSION) + } + + /// Exports a quarantined database without mutating it and authenticates + /// the resulting SQLite artifact with a caller-owned repair capability. + /// + /// # Errors + /// + /// Returns a safe state, authorization, or storage error. + pub fn export_quarantined( + path: &Path, + destination: &Path, + authorization: &RepairAuthorization, + ) -> Result<QuarantineExportReceipt, SafeError> { + export_quarantined(path, destination, authorization) + } + + /// Validates and authenticates a canonical repaired database candidate. + /// + /// # Errors + /// + /// Returns a safe compatibility or storage error for an invalid candidate. + pub fn authenticate_repair_candidate( + path: &Path, + authorization: &RepairAuthorization, + ) -> Result<RepairCandidate, SafeError> { + authenticate_candidate(path, authorization) + } + + /// Atomically installs an authenticated candidate over a quarantined + /// database while retaining the original as immutable evidence. + /// + /// # Errors + /// + /// Returns a safe authorization, ownership, or storage error without + /// replacing the target when any gate fails. + pub fn install_repair_candidate( + path: &Path, + candidate: &RepairCandidate, + authorization: &RepairAuthorization, + ) -> Result<(), SafeError> { + let _ownership = WritableOwnership::acquire(path)?; + install_candidate(path, candidate, authorization) + } + + /// Returns the highest successfully applied migration version. + /// + /// # Errors + /// + /// Returns a safe storage error when migration history cannot be read. + pub fn schema_version(&self) -> Result<u32, SafeError> { + self.connection() + .query_row( + "SELECT COALESCE(MAX(version), 0) FROM refinery_schema_history", + [], + |row| row.get(0), + ) + .map_err(|_| corrupt_storage_error()) + } + + pub(crate) fn connection(&self) -> DatabaseConnection<'_> { + DatabaseConnection { + connection: self + .connection + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner), + path: self.path.as_deref(), + } + } +} + +impl Deref for DatabaseConnection<'_> { + type Target = Connection; + + fn deref(&self) -> &Self::Target { + &self.connection + } +} + +impl DerefMut for DatabaseConnection<'_> { + fn deref_mut(&mut self) -> &mut Self::Target { + &mut self.connection + } +} + +impl Drop for DatabaseConnection<'_> { + fn drop(&mut self) { + if let Some(path) = self.path { + let _ = restrict_sqlite_sidecars(path); + } + } +} + +impl WritableOwnership { + fn acquire(database_path: &Path) -> Result<Self, SafeError> { + let lock_path = database_path.with_extension("sqlite3.lock"); + let mut options = OpenOptions::new(); + options.read(true).write(true).create(true).truncate(false); + #[cfg(unix)] + { + use std::os::unix::fs::OpenOptionsExt; + options.custom_flags( + (rustix::fs::OFlags::NOFOLLOW | rustix::fs::OFlags::CLOEXEC).bits() as i32, + ); + } + let file = options.open(&lock_path).map_err(|_| storage_error())?; + restrict_file_permissions(&lock_path)?; + file.try_lock_exclusive().map_err(|_| ownership_error())?; + Ok(Self { _file: file }) + } +} + +fn create_secure_directory(path: &Path) -> Result<(), SafeError> { + let mut existing = path; + loop { + match fs::symlink_metadata(existing) { + Ok(metadata) => { + if metadata.file_type().is_symlink() || !metadata.is_dir() { + return Err(storage_error()); + } + break; + } + Err(error) if error.kind() == std::io::ErrorKind::NotFound => { + existing = existing.parent().ok_or_else(storage_error)?; + } + Err(_) => return Err(storage_error()), + } + } + fs::create_dir_all(path).map_err(|_| storage_error())?; + let metadata = fs::symlink_metadata(path).map_err(|_| storage_error())?; + if metadata.file_type().is_symlink() || !metadata.is_dir() { + return Err(storage_error()); + } + restrict_directory_permissions(path) +} + +fn configure(connection: &Connection) -> Result<(), SafeError> { + connection + .pragma_update(None, "foreign_keys", "ON") + .and_then(|()| connection.pragma_update(None, "trusted_schema", "OFF")) + .and_then(|()| connection.pragma_update(None, "journal_mode", "WAL")) + .and_then(|()| connection.pragma_update(None, "synchronous", "FULL")) + .and_then(|()| connection.pragma_update(None, "secure_delete", "ON")) + .and_then(|()| connection.pragma_update(None, "wal_autocheckpoint", 1_000)) + .and_then(|()| connection.busy_timeout(Duration::from_secs(5))) + .map_err(|_| storage_error()) +} + +fn restrict_sqlite_sidecars(path: &Path) -> Result<(), SafeError> { + for suffix in ["-wal", "-shm"] { + let sidecar = PathBuf::from(format!("{}{suffix}", path.display())); + match fs::symlink_metadata(&sidecar) { + Ok(metadata) if metadata.file_type().is_symlink() || !metadata.is_file() => { + return Err(storage_error()); + } + Ok(_) => restrict_file_permissions(&sidecar)?, + Err(error) if error.kind() == std::io::ErrorKind::NotFound => {} + Err(_) => return Err(storage_error()), + } + } + Ok(()) +} + +#[cfg(unix)] +pub(crate) fn restrict_file_permissions(path: &Path) -> Result<(), SafeError> { + use std::os::unix::fs::PermissionsExt; + + fs::set_permissions(path, fs::Permissions::from_mode(0o600)).map_err(|_| storage_error()) +} + +#[cfg(unix)] +pub(crate) fn restrict_directory_permissions(path: &Path) -> Result<(), SafeError> { + use std::os::unix::fs::PermissionsExt; + + fs::set_permissions(path, fs::Permissions::from_mode(0o700)).map_err(|_| storage_error()) +} + +#[cfg(not(unix))] +pub(crate) fn restrict_file_permissions(_path: &Path) -> Result<(), SafeError> { + Ok(()) +} + +#[cfg(not(unix))] +pub(crate) fn restrict_directory_permissions(_path: &Path) -> Result<(), SafeError> { + Ok(()) +} + +const fn storage_error() -> SafeError { + SafeError::new( + SafeErrorCode::StorageUnavailable, + SafeMessage::new("The application database is unavailable."), + ) +} + +const fn corrupt_storage_error() -> SafeError { + SafeError::new( + SafeErrorCode::StorageCorrupt, + SafeMessage::new("The application database could not be read."), + ) +} + +const fn ownership_error() -> SafeError { + SafeError::new( + SafeErrorCode::StorageUnavailable, + SafeMessage::new("The application database is already in use."), + ) +} + +#[cfg(test)] +mod tests { + use std::fs; + use std::io::Write; + use std::path::Path; + use std::process::Command; + + use tempfile::tempdir; + + use radroots_studio_application::{AccountRepository, AppStateRepository}; + use radroots_studio_domain::{PublicKey, SafeErrorCode}; + use refinery::Target; + use rusqlite::Connection; + + use super::{ + CURRENT_SCHEMA_VERSION, Database, configure, create_secure_directory, migrations, + restrict_sqlite_sidecars, + }; + use crate::{DatabasePreflight, PersistedIdentityIssueKind, RepairAuthorization}; + + #[test] + fn migration_opens_fresh_memory_database_once() { + let database = Database::in_memory().expect("open memory database"); + + assert_eq!( + database.schema_version().expect("schema version"), + CURRENT_SCHEMA_VERSION + ); + assert_eq!( + database.schema_version().expect("repeat schema version"), + CURRENT_SCHEMA_VERSION + ); + } + + #[test] + fn database_path_guards_reject_files_as_directories_and_sidecars() { + let directory = tempdir().expect("temporary directory"); + let regular = directory.path().join("regular"); + fs::write(&regular, b"file").expect("write regular file"); + assert!(create_secure_directory(&regular).is_err()); + + let database = directory.path().join("studio.sqlite3"); + fs::write(&database, b"database").expect("write database file"); + fs::create_dir(directory.path().join("studio.sqlite3-wal")) + .expect("create invalid WAL sidecar"); + assert!(restrict_sqlite_sidecars(&database).is_err()); + } + + #[test] + fn sqlite_connection_enforces_trust_durability_and_busy_policy() { + let database = Database::in_memory().expect("open memory database"); + let connection = database.connection(); + + assert_eq!( + connection + .pragma_query_value(None, "foreign_keys", |row| row.get::<_, u8>(0)) + .expect("foreign keys"), + 1 + ); + assert_eq!( + connection + .pragma_query_value(None, "trusted_schema", |row| row.get::<_, u8>(0)) + .expect("trusted schema"), + 0 + ); + assert_eq!( + connection + .pragma_query_value(None, "synchronous", |row| row.get::<_, u8>(0)) + .expect("synchronous"), + 2 + ); + assert_eq!( + connection + .pragma_query_value(None, "busy_timeout", |row| row.get::<_, i64>(0)) + .expect("busy timeout"), + 5_000 + ); + } + + #[test] + fn normalized_schema_is_strict_and_enforces_same_account_bindings() { + let database = Database::in_memory().expect("open memory database"); + let connection = database.connection(); + let strict_tables: i64 = connection + .query_row( + "SELECT COUNT(*) FROM pragma_table_list WHERE name IN ('account_identities', 'local_signer_bindings', 'runtime_state', 'profile_cache_v6', 'durable_operations') AND strict = 1", + [], + |row| row.get(0), + ) + .expect("strict table inventory"); + assert_eq!(strict_tables, 5); + + connection + .execute( + "INSERT INTO account_identities (public_key, npub, created_at) VALUES (?1, ?2, 1)", + [ + "07".repeat(32), + "npub1qurswpc8qurswpc8qurswpc8qurswpc8qurswpc8qurswpc8qursnvjvl7".to_owned(), + ], + ) + .expect("identity"); + assert!( + connection + .execute( + "INSERT INTO local_signer_bindings (account_public_key, binding_public_key, binding_kind, availability) VALUES (?1, ?2, 'local_secret', 'available')", + ["07".repeat(32), "08".repeat(32)], + ) + .is_err() + ); + } + + #[test] + fn v5_data_migrates_append_only_with_identity_profile_and_selection() { + let directory = tempdir().expect("temporary directory"); + let path = directory.path().join("studio.sqlite3"); + let public_key = "07".repeat(32); + { + let mut connection = Connection::open(&path).expect("legacy database"); + configure(&connection).expect("configuration"); + migrations::migrations::runner() + .set_target(Target::Version(5)) + .run(&mut connection) + .expect("V5 schema"); + connection + .execute( + "INSERT INTO accounts (pubkey, npub, signer_kind, key_availability, created_at) VALUES (?1, ?2, 'local_secret', 'available', 10)", + [&public_key, "npub1qurswpc8qurswpc8qurswpc8qurswpc8qurswpc8qurswpc8qursnvjvl7"], + ) + .expect("legacy account"); + connection + .execute( + "UPDATE app_state SET selected_pubkey = ?1 WHERE singleton = 1", + [&public_key], + ) + .expect("legacy selection"); + connection + .execute( + "INSERT INTO profile_cache (subject_pubkey, event_id, event_created_at, name, refreshed_at, refresh_status) VALUES (?1, ?2, 11, 'Farm', 12, 'success')", + [&public_key, &"01".repeat(32)], + ) + .expect("legacy profile"); + } + + let database = Database::open(&path).expect("migrated database"); + 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"), + Some(PublicKey::from_bytes([7; 32]).expect("valid public key")) + ); + let connection = database.connection(); + let migrated: (i64, i64, i64) = connection + .query_row( + "SELECT (SELECT COUNT(*) FROM account_identities), (SELECT COUNT(*) FROM local_signer_bindings), (SELECT COUNT(*) FROM profile_cache_v6)", + [], + |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?)), + ) + .expect("migrated inventory"); + assert_eq!(migrated, (1, 1, 1)); + drop(connection); + drop(database); + + Database::verify_migration_backup(&path, 5).expect("authenticated backup"); + Database::restore_migration_backup(&path, 5).expect("authenticated restore"); + assert_eq!( + Database::preflight(&path).expect("restored preflight"), + DatabasePreflight::Ready { schema_version: 5 } + ); + let retried = Database::open(&path).expect("idempotent migration retry"); + assert_eq!(retried.schema_version().expect("retried version"), 10); + drop(retried); + + let backup = directory + .path() + .join("studio.sqlite3.recovery/migration-v5-to-v10.sqlite3"); + fs::OpenOptions::new() + .append(true) + .open(backup) + .expect("open backup") + .write_all(b"tamper") + .expect("tamper backup"); + let error = + Database::verify_migration_backup(&path, 5).expect_err("tampered backup must fail"); + assert_eq!(error.code(), SafeErrorCode::StorageBackupInvalid); + } + + #[test] + fn corrupt_v5_identity_fails_before_migration_without_recreation() { + let directory = tempdir().expect("temporary directory"); + let path = directory.path().join("studio.sqlite3"); + { + let mut connection = Connection::open(&path).expect("legacy database"); + configure(&connection).expect("configuration"); + migrations::migrations::runner() + .set_target(Target::Version(5)) + .run(&mut connection) + .expect("V5 schema"); + connection + .execute( + "INSERT INTO accounts (pubkey, npub, signer_kind, key_availability, created_at) VALUES (?1, ?2, 'local_secret', 'available', 10)", + ["07".repeat(32), "npub10elfcs4fr0l0r8af98jlmgdh9c8tcxjvz9qkw038js35mp4dma8qzvjptg".to_owned()], + ) + .expect("mismatched legacy account"); + } + + assert!(Database::open(&path).is_err()); + let connection = Connection::open(&path).expect("inspect legacy database"); + let version: u32 = connection + .query_row( + "SELECT MAX(version) FROM refinery_schema_history", + [], + |row| row.get(0), + ) + .expect("legacy version"); + let accounts: i64 = connection + .query_row("SELECT COUNT(*) FROM accounts", [], |row| row.get(0)) + .expect("legacy accounts"); + assert_eq!((version, accounts), (5, 1)); + } + + #[test] + fn invalid_curve_identity_is_quarantined_without_mutation() { + let directory = tempdir().expect("temporary directory"); + let path = directory.path().join("studio.sqlite3"); + { + let mut connection = Connection::open(&path).expect("legacy database"); + configure(&connection).expect("configuration"); + migrations::migrations::runner() + .set_target(Target::Version(5)) + .run(&mut connection) + .expect("V5 schema"); + connection + .execute( + "INSERT INTO accounts (pubkey, npub, signer_kind, key_availability, created_at) VALUES (?1, ?2, 'local_secret', 'available', 10)", + ["00".repeat(32), "npub1qurswpc8qurswpc8qurswpc8qurswpc8qurswpc8qurswpc8qursnvjvl7".to_owned()], + ) + .expect("invalid-curve fixture"); + connection + .execute_batch("PRAGMA wal_checkpoint(TRUNCATE)") + .expect("checkpoint"); + } + let before = fs::read(&path).expect("before bytes"); + + let DatabasePreflight::Quarantined { + schema_version, + issues, + } = Database::preflight(&path).expect("classified preflight") + else { + panic!("invalid identity was not quarantined"); + }; + assert_eq!(schema_version, 5); + assert!(issues.iter().any(|issue| { + issue.table() == "accounts" + && issue.column() == "pubkey" + && issue.kind() == PersistedIdentityIssueKind::InvalidCurvePoint + })); + let error = Database::open(&path) + .err() + .expect("quarantined open must fail"); + assert_eq!(error.code(), SafeErrorCode::StorageQuarantined); + assert_eq!(fs::read(&path).expect("after bytes"), before); + assert!(!path.with_extension("sqlite3.lock").exists()); + + let authorization = RepairAuthorization::from_bytes(vec![0x41; 32]) + .unwrap_or_else(|_| panic!("repair authorization")); + let export_path = directory.path().join("quarantine-export.sqlite3"); + let export = Database::export_quarantined(&path, &export_path, &authorization) + .expect("authenticated quarantine export"); + assert_eq!(export.path(), export_path); + assert_eq!(export.sha256().len(), 64); + assert_eq!(export.authentication_tag().len(), 64); + assert_eq!(fs::read(&path).expect("post-export bytes"), before); + + let candidate_path = directory.path().join("repaired.sqlite3"); + drop(Database::open(&candidate_path).expect("canonical repair candidate")); + let candidate = Database::authenticate_repair_candidate(&candidate_path, &authorization) + .expect("authenticate candidate"); + let wrong_authorization = RepairAuthorization::from_bytes(vec![0x42; 32]) + .unwrap_or_else(|_| panic!("wrong authorization shape")); + let error = Database::install_repair_candidate(&path, &candidate, &wrong_authorization) + .expect_err("wrong repair authorization"); + assert_eq!(error.code(), SafeErrorCode::RepairUnauthorized); + assert_eq!(fs::read(&path).expect("unauthorized bytes"), before); + + Database::install_repair_candidate(&path, &candidate, &authorization) + .expect("authenticated repair install"); + assert!(matches!( + Database::preflight(&path).expect("repaired preflight"), + DatabasePreflight::Ready { + schema_version: CURRENT_SCHEMA_VERSION + } + )); + assert!( + directory + .path() + .join("studio.sqlite3.quarantined-evidence") + .is_file() + ); + } + + #[test] + fn newer_and_mixed_schema_inventory_fail_before_mutation() { + let directory = tempdir().expect("temporary directory"); + let newer_path = directory.path().join("newer.sqlite3"); + { + let database = Database::open(&newer_path).expect("current database"); + database + .connection() + .execute( + "UPDATE refinery_schema_history SET version = ?1 WHERE version = ?2", + [CURRENT_SCHEMA_VERSION + 1, CURRENT_SCHEMA_VERSION], + ) + .expect("future schema row"); + } + let newer_before = fs::read(&newer_path).expect("newer bytes"); + let error = Database::preflight(&newer_path).expect_err("newer schema"); + assert_eq!(error.code(), SafeErrorCode::UnsupportedSchemaVersion); + assert_eq!(fs::read(&newer_path).expect("newer after"), newer_before); + + let mixed_path = directory.path().join("mixed.sqlite3"); + { + let mut connection = Connection::open(&mixed_path).expect("legacy database"); + configure(&connection).expect("configuration"); + migrations::migrations::runner() + .set_target(Target::Version(5)) + .run(&mut connection) + .expect("V5 schema"); + connection + .execute("CREATE TABLE installation_identity (singleton INTEGER)", []) + .expect("mixed table"); + } + let mixed_before = fs::read(&mixed_path).expect("mixed bytes"); + let error = Database::preflight(&mixed_path).expect_err("mixed schema"); + assert_eq!(error.code(), SafeErrorCode::StorageCorrupt); + assert_eq!(fs::read(&mixed_path).expect("mixed after"), mixed_before); + } + + #[test] + fn failed_v5_copy_rolls_back_the_active_migration() { + let directory = tempdir().expect("temporary directory"); + let path = directory.path().join("studio.sqlite3"); + let public_key = "07".repeat(32); + { + let mut connection = Connection::open(&path).expect("legacy database"); + configure(&connection).expect("configuration"); + migrations::migrations::runner() + .set_target(Target::Version(5)) + .run(&mut connection) + .expect("V5 schema"); + connection + .execute( + "INSERT INTO accounts (pubkey, npub, signer_kind, key_availability, created_at) VALUES (?1, ?2, 'local_secret', 'available', 10)", + [&public_key, "npub1qurswpc8qurswpc8qurswpc8qurswpc8qurswpc8qurswpc8qursnvjvl7"], + ) + .expect("legacy account"); + connection + .execute( + "INSERT INTO profile_cache (subject_pubkey, event_id, event_created_at, refreshed_at, refresh_status) VALUES (?1, 'invalid', 11, 12, 'success')", + [&public_key], + ) + .expect("legacy corrupt profile"); + } + + assert!(Database::open(&path).is_err()); + let connection = Connection::open(&path).expect("inspect interrupted migration"); + let version: u32 = connection + .query_row( + "SELECT MAX(version) FROM refinery_schema_history", + [], + |row| row.get(0), + ) + .expect("migration version"); + assert_eq!(version, 5); + assert!(!connection + .query_row( + "SELECT EXISTS(SELECT 1 FROM sqlite_master WHERE type = 'table' AND name = 'account_identities')", + [], + |row| row.get::<_, bool>(0), + ) + .expect("normalized table inventory")); + } + + #[test] + fn foreign_keys_reject_orphan_normalized_records() { + let database = Database::in_memory().expect("database"); + let connection = database.connection(); + assert!( + connection + .execute( + "INSERT INTO local_signer_bindings (account_public_key, binding_public_key, binding_kind, availability) VALUES (?1, ?1, 'local_secret', 'available')", + ["09".repeat(32)], + ) + .is_err() + ); + } + + #[test] + fn second_process_cannot_acquire_writable_ownership() { + let directory = tempdir().expect("temporary directory"); + let path = directory.path().join("studio.sqlite3"); + let _owner = Database::open(&path).expect("parent owner"); + let status = Command::new(std::env::current_exe().expect("test executable")) + .arg("--exact") + .arg("db::tests::writable_ownership_child_probe") + .arg("--nocapture") + .env("RADROOTS_STUDIO_LOCK_PROBE_PATH", &path) + .status() + .expect("child process"); + assert!(status.success()); + } + + #[test] + fn writable_ownership_child_probe() { + let Ok(path) = std::env::var("RADROOTS_STUDIO_LOCK_PROBE_PATH") else { + return; + }; + assert!(Database::open(Path::new(&path)).is_err()); + } + + #[test] + fn migration_persists_schema_version_across_file_reopen() { + let directory = tempdir().expect("temporary directory"); + let path = directory.path().join("studio.sqlite3"); + + { + let database = Database::open(&path).expect("open file database"); + assert_eq!( + database.schema_version().expect("schema version"), + CURRENT_SCHEMA_VERSION + ); + } + let reopened = Database::open(&path).expect("reopen file database"); + assert_eq!( + reopened.schema_version().expect("schema version"), + CURRENT_SCHEMA_VERSION + ); + assert!(fs::metadata(path).expect("database metadata").len() > 0); + } + + #[test] + fn writable_ownership_rejects_a_second_runtime_and_releases_on_drop() { + let directory = tempdir().expect("temporary directory"); + let path = directory.path().join("studio.sqlite3"); + let first = Database::open(&path).expect("first owner"); + let Err(error) = Database::open(&path) else { + panic!("second owner must fail"); + }; + assert_eq!( + error.message().as_str(), + "The application database is already in use." + ); + drop(first); + Database::open(&path).expect("ownership released"); + } + + #[cfg(unix)] + #[test] + fn migration_attempts_owner_only_database_permissions() { + use std::os::unix::fs::PermissionsExt; + + let directory = tempdir().expect("temporary directory"); + let path = directory.path().join("studio.sqlite3"); + let database = Database::open(&path).expect("open file database"); + let mode = fs::metadata(&path) + .expect("database metadata") + .permissions() + .mode() + & 0o777; + + assert_eq!(mode, 0o600); + let directory_mode = fs::metadata(directory.path()) + .expect("directory metadata") + .permissions() + .mode() + & 0o777; + assert_eq!(directory_mode, 0o700); + + let connection = database.connection(); + connection + .execute_batch("CREATE TABLE sidecar_probe (value INTEGER) STRICT; INSERT INTO sidecar_probe VALUES (1);") + .expect("write through WAL"); + drop(connection); + for suffix in ["-wal", "-shm"] { + let sidecar = std::path::PathBuf::from(format!("{}{suffix}", path.display())); + let sidecar_mode = fs::metadata(sidecar) + .expect("sidecar metadata") + .permissions() + .mode() + & 0o777; + assert_eq!(sidecar_mode, 0o600); + } + } + + #[cfg(unix)] + #[test] + fn database_lock_sidecar_and_recovery_symlinks_fail_closed() { + use std::os::unix::fs::symlink; + + let directory = tempdir().expect("temporary directory"); + let victim = directory.path().join("victim"); + fs::write(&victim, b"unchanged").expect("victim"); + + let database_link = directory.path().join("database-link.sqlite3"); + symlink(&victim, &database_link).expect("database symlink"); + assert!(Database::open(&database_link).is_err()); + assert_eq!(fs::read(&victim).expect("victim bytes"), b"unchanged"); + + let lock_path = directory.path().join("locked.sqlite3"); + symlink(&victim, lock_path.with_extension("sqlite3.lock")).expect("lock symlink"); + assert!(Database::open(&lock_path).is_err()); + assert_eq!(fs::read(&victim).expect("victim bytes"), b"unchanged"); + + let sidecar_path = directory.path().join("sidecar.sqlite3"); + let wal = std::path::PathBuf::from(format!("{}-wal", sidecar_path.display())); + symlink(&victim, wal).expect("WAL symlink"); + assert!(Database::open(&sidecar_path).is_err()); + assert_eq!(fs::read(&victim).expect("victim bytes"), b"unchanged"); + + let legacy_path = directory.path().join("legacy.sqlite3"); + { + let mut connection = Connection::open(&legacy_path).expect("legacy database"); + configure(&connection).expect("configuration"); + migrations::migrations::runner() + .set_target(Target::Version(5)) + .run(&mut connection) + .expect("V5 schema"); + } + symlink( + directory.path().join("not-present"), + directory.path().join("legacy.sqlite3.recovery"), + ) + .expect("recovery symlink"); + assert!(Database::open(&legacy_path).is_err()); + } +} diff --git a/core/crates/studio_storage/src/installation.rs b/core/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/core/crates/studio_storage/src/journal.rs b/core/crates/studio_storage/src/journal.rs @@ -0,0 +1,774 @@ +use radroots_studio_application::{ + AccountOperationKind, AccountOperationPhase, DurableAccountOperation, DurableOperationKind, + DurableOperationPhase, DurableOperationReceipt, DurableOperationRepository, + DurableOperationStart, DurableRequestId, DurableTerminalOutcome, OperationDiagnostic, + OperationId, OperationJournal, OperationPriorState, PendingAccountOperation, +}; +use radroots_studio_domain::{ + BindingAvailability, PublicKey, SafeError, SafeErrorCode, SafeMessage, UnixTimestamp, +}; +use rusqlite::{OptionalExtension, Row, params}; + +use crate::Database; + +impl DurableOperationRepository for Database { + fn begin_durable_operation( + &self, + request_id: &DurableRequestId, + kind: DurableOperationKind, + account: PublicKey, + expected_revision: Option<u64>, + prior: OperationPriorState, + updated_at: UnixTimestamp, + ) -> Result<DurableOperationStart, SafeError> { + let encoded_expected_revision = expected_revision + .map(i64::try_from) + .transpose() + .map_err(|_| operation_conflict())?; + let mut connection = self.connection(); + let transaction = connection.transaction().map_err(|_| storage_error())?; + let inserted = transaction + .execute( + "INSERT OR IGNORE INTO durable_operations (request_id, operation_kind, \ + account_public_key, binding_public_key, expected_revision, phase, \ + prior_selected_public_key, updated_at, prior_binding_availability) \ + VALUES (?1, ?2, ?3, ?3, ?4, 'intent_recorded', ?5, ?6, ?7)", + params![ + request_id.as_str(), + encode_durable_kind(kind), + account.to_hex(), + encoded_expected_revision, + prior.selected_account().map(PublicKey::to_hex), + updated_at.as_seconds(), + prior + .binding_availability() + .map(encode_binding_availability), + ], + ) + .map_err(|_| storage_error())?; + let operation = + query_durable_operation(&transaction, request_id)?.ok_or_else(corrupt_storage_error)?; + if operation.kind() != kind + || operation.account() != account + || operation.expected_revision() != expected_revision + || operation.prior() != prior + { + return Err(operation_conflict()); + } + transaction.commit().map_err(|_| storage_error())?; + Ok(if inserted == 1 { + DurableOperationStart::Started(operation) + } else { + DurableOperationStart::Existing(operation) + }) + } + + fn load_durable_operation( + &self, + request_id: &DurableRequestId, + ) -> Result<Option<DurableAccountOperation>, SafeError> { + query_durable_operation(&self.connection(), request_id) + } + + fn advance_durable_operation( + &self, + request_id: &DurableRequestId, + expected_phase: DurableOperationPhase, + next_phase: DurableOperationPhase, + updated_at: UnixTimestamp, + diagnostic: Option<OperationDiagnostic>, + ) -> Result<DurableAccountOperation, SafeError> { + let mut connection = self.connection(); + let transaction = connection.transaction().map_err(|_| storage_error())?; + let rows = transaction + .execute( + "UPDATE durable_operations SET phase = ?3, updated_at = ?4, diagnostic_code = ?5 \ + WHERE request_id = ?1 AND phase = ?2 AND terminal_outcome IS NULL", + params![ + request_id.as_str(), + encode_durable_phase(expected_phase), + encode_durable_phase(next_phase), + updated_at.as_seconds(), + diagnostic.map(encode_diagnostic), + ], + ) + .map_err(|_| storage_error())?; + if rows != 1 { + return Err(operation_conflict()); + } + let operation = + query_durable_operation(&transaction, request_id)?.ok_or_else(corrupt_storage_error)?; + transaction.commit().map_err(|_| storage_error())?; + Ok(operation) + } + + fn finalize_durable_operation( + &self, + request_id: &DurableRequestId, + expected_phase: DurableOperationPhase, + outcome: DurableTerminalOutcome, + resulting_revision: Option<u64>, + updated_at: UnixTimestamp, + ) -> Result<DurableOperationReceipt, SafeError> { + if let Some(existing) = self.load_durable_operation(request_id)? + && let Some(receipt) = existing.terminal() + { + return if receipt.outcome() == outcome + && receipt.resulting_revision() == resulting_revision + { + Ok(receipt.clone()) + } else { + Err(operation_conflict()) + }; + } + let resulting_revision = resulting_revision + .map(i64::try_from) + .transpose() + .map_err(|_| operation_conflict())?; + let rows = self + .connection() + .execute( + "UPDATE durable_operations SET phase = 'finalized', terminal_outcome = ?3, \ + resulting_revision = ?4, updated_at = ?5 \ + WHERE request_id = ?1 AND phase = ?2 AND terminal_outcome IS NULL", + params![ + request_id.as_str(), + encode_durable_phase(expected_phase), + encode_terminal_outcome(outcome), + resulting_revision, + updated_at.as_seconds(), + ], + ) + .map_err(|_| storage_error())?; + if rows != 1 { + return Err(operation_conflict()); + } + self.load_durable_operation(request_id)? + .and_then(|operation| operation.terminal().cloned()) + .ok_or_else(corrupt_storage_error) + } + + fn list_unfinished_durable_operations( + &self, + ) -> Result<Vec<DurableAccountOperation>, SafeError> { + let connection = self.connection(); + let mut statement = connection + .prepare(&format!( + "{DURABLE_OPERATION_SELECT} WHERE terminal_outcome IS NULL ORDER BY request_id ASC" + )) + .map_err(|_| storage_error())?; + let rows = statement + .query_map([], decode_durable_operation) + .map_err(|_| storage_error())?; + rows.map(|row| row.map_err(|_| corrupt_storage_error())) + .collect() + } +} + +const DURABLE_OPERATION_SELECT: &str = "SELECT request_id, operation_kind, account_public_key, \ + expected_revision, phase, prior_selected_public_key, updated_at, diagnostic_code, \ + terminal_outcome, prior_binding_availability, resulting_revision FROM durable_operations"; + +fn query_durable_operation( + connection: &rusqlite::Connection, + request_id: &DurableRequestId, +) -> Result<Option<DurableAccountOperation>, SafeError> { + connection + .query_row( + &format!("{DURABLE_OPERATION_SELECT} WHERE request_id = ?1"), + [request_id.as_str()], + decode_durable_operation, + ) + .optional() + .map_err(|_| corrupt_storage_error()) +} + +fn decode_durable_operation(row: &Row<'_>) -> rusqlite::Result<DurableAccountOperation> { + let request_id = + DurableRequestId::parse(row.get::<_, String>(0)?).map_err(|_| invalid_column(0))?; + let kind = decode_durable_kind(row.get::<_, String>(1)?.as_str())?; + let account = + PublicKey::from_hex(row.get::<_, String>(2)?.as_str()).map_err(|_| invalid_column(2))?; + let expected_revision = row + .get::<_, Option<i64>>(3)? + .map(|value| u64::try_from(value).map_err(|_| invalid_column(3))) + .transpose()?; + let phase = decode_durable_phase(row.get::<_, String>(4)?.as_str())?; + let prior_selected = row + .get::<_, Option<String>>(5)? + .map(|value| PublicKey::from_hex(&value).map_err(|_| invalid_column(5))) + .transpose()?; + let updated_at = UnixTimestamp::from_seconds(row.get(6)?).ok_or_else(|| invalid_column(6))?; + let diagnostic = row + .get::<_, Option<String>>(7)? + .map(|value| decode_diagnostic(&value)) + .transpose()?; + let outcome = row + .get::<_, Option<String>>(8)? + .map(|value| decode_terminal_outcome(&value)) + .transpose()?; + let prior_availability = row + .get::<_, Option<String>>(9)? + .map(|value| decode_binding_availability(&value)) + .transpose()?; + let resulting_revision = row + .get::<_, Option<i64>>(10)? + .map(|value| u64::try_from(value).map_err(|_| invalid_column(10))) + .transpose()?; + let terminal = outcome.map(|outcome| { + DurableOperationReceipt::new(request_id.clone(), account, outcome, resulting_revision) + }); + Ok(DurableAccountOperation::new( + request_id, + kind, + account, + expected_revision, + phase, + OperationPriorState::new(prior_selected, prior_availability), + updated_at, + diagnostic, + terminal, + )) +} + +impl OperationJournal for Database { + fn begin_operation( + &self, + kind: AccountOperationKind, + subject: PublicKey, + updated_at: UnixTimestamp, + ) -> Result<OperationId, SafeError> { + let connection = self.connection(); + connection + .execute( + "INSERT INTO operation_journal (operation_kind, subject_pubkey, phase, \ + updated_at) VALUES (?1, ?2, 'intent_recorded', ?3)", + params![encode_kind(kind), subject.to_hex(), updated_at.as_seconds()], + ) + .map_err(|_| storage_error())?; + let id = + u64::try_from(connection.last_insert_rowid()).map_err(|_| corrupt_storage_error())?; + Ok(OperationId::from_raw(id)) + } + + fn update_operation( + &self, + id: OperationId, + phase: AccountOperationPhase, + updated_at: UnixTimestamp, + diagnostic: Option<OperationDiagnostic>, + ) -> Result<(), SafeError> { + let encoded_id = i64::try_from(id.as_raw()).map_err(|_| corrupt_storage_error())?; + match self.connection().execute( + "UPDATE operation_journal SET phase = ?2, updated_at = ?3, diagnostic_code = ?4 \ + WHERE operation_id = ?1", + params![ + encoded_id, + encode_phase(phase), + updated_at.as_seconds(), + diagnostic.map(encode_diagnostic) + ], + ) { + Ok(1) => Ok(()), + Ok(0) => Err(operation_not_found()), + Ok(_) | Err(_) => Err(storage_error()), + } + } + + fn list_pending_operations(&self) -> Result<Vec<PendingAccountOperation>, SafeError> { + let connection = self.connection(); + let mut statement = connection + .prepare( + "SELECT operation_id, operation_kind, subject_pubkey, phase, updated_at, \ + diagnostic_code FROM operation_journal ORDER BY operation_id ASC", + ) + .map_err(|_| storage_error())?; + let rows = statement + .query_map([], decode_operation) + .map_err(|_| storage_error())?; + rows.map(|row| row.map_err(|_| corrupt_storage_error())) + .collect() + } + + fn finalize_operation(&self, id: OperationId) -> Result<(), SafeError> { + let encoded_id = i64::try_from(id.as_raw()).map_err(|_| corrupt_storage_error())?; + self.connection() + .execute( + "DELETE FROM operation_journal WHERE operation_id = ?1", + [encoded_id], + ) + .map(|_| ()) + .map_err(|_| storage_error()) + } +} + +fn decode_operation(row: &Row<'_>) -> rusqlite::Result<PendingAccountOperation> { + let id = u64::try_from(row.get::<_, i64>(0)?).map_err(|_| invalid_column(0))?; + let kind = decode_kind(row.get::<_, String>(1)?.as_str())?; + let subject = + PublicKey::from_hex(row.get::<_, String>(2)?.as_str()).map_err(|_| invalid_column(2))?; + let phase = decode_phase(row.get::<_, String>(3)?.as_str())?; + let updated_at = UnixTimestamp::from_seconds(row.get(4)?).ok_or_else(|| invalid_column(4))?; + let diagnostic = row + .get::<_, Option<String>>(5)? + .map(|value| decode_diagnostic(&value)) + .transpose()?; + Ok(PendingAccountOperation::new( + OperationId::from_raw(id), + kind, + subject, + phase, + updated_at, + diagnostic, + )) +} + +const fn encode_durable_kind(value: DurableOperationKind) -> &'static str { + match value { + DurableOperationKind::Create => "create", + DurableOperationKind::Import => "import", + DurableOperationKind::Repair => "repair", + DurableOperationKind::Remove => "remove", + } +} + +fn decode_durable_kind(value: &str) -> rusqlite::Result<DurableOperationKind> { + match value { + "create" => Ok(DurableOperationKind::Create), + "import" => Ok(DurableOperationKind::Import), + "repair" => Ok(DurableOperationKind::Repair), + "remove" => Ok(DurableOperationKind::Remove), + _ => Err(invalid_column(1)), + } +} + +const fn encode_durable_phase(value: DurableOperationPhase) -> &'static str { + match value { + DurableOperationPhase::IntentRecorded => "intent_recorded", + DurableOperationPhase::CredentialWritten => "credential_written", + DurableOperationPhase::MetadataCommitted => "metadata_committed", + DurableOperationPhase::SelectionCommitted => "selection_committed", + DurableOperationPhase::CompensationPending => "compensation_pending", + DurableOperationPhase::CredentialDeleted => "credential_deleted", + DurableOperationPhase::MetadataDeleted => "metadata_deleted", + DurableOperationPhase::Finalized => "finalized", + } +} + +fn decode_durable_phase(value: &str) -> rusqlite::Result<DurableOperationPhase> { + match value { + "intent_recorded" => Ok(DurableOperationPhase::IntentRecorded), + "credential_written" => Ok(DurableOperationPhase::CredentialWritten), + "metadata_committed" => Ok(DurableOperationPhase::MetadataCommitted), + "selection_committed" => Ok(DurableOperationPhase::SelectionCommitted), + "compensation_pending" => Ok(DurableOperationPhase::CompensationPending), + "credential_deleted" => Ok(DurableOperationPhase::CredentialDeleted), + "metadata_deleted" => Ok(DurableOperationPhase::MetadataDeleted), + "finalized" => Ok(DurableOperationPhase::Finalized), + _ => Err(invalid_column(4)), + } +} + +const fn encode_terminal_outcome(value: DurableTerminalOutcome) -> &'static str { + match value { + DurableTerminalOutcome::Completed => "completed", + DurableTerminalOutcome::Cancelled => "cancelled", + DurableTerminalOutcome::Failed => "failed", + } +} + +fn decode_terminal_outcome(value: &str) -> rusqlite::Result<DurableTerminalOutcome> { + match value { + "completed" => Ok(DurableTerminalOutcome::Completed), + "cancelled" => Ok(DurableTerminalOutcome::Cancelled), + "failed" => Ok(DurableTerminalOutcome::Failed), + _ => Err(invalid_column(8)), + } +} + +const fn encode_binding_availability(value: BindingAvailability) -> &'static str { + match value { + BindingAvailability::Available => "available", + BindingAvailability::CredentialMissing => "credential_missing", + BindingAvailability::StoreUnavailable => "store_unavailable", + } +} + +fn decode_binding_availability(value: &str) -> rusqlite::Result<BindingAvailability> { + match value { + "available" => Ok(BindingAvailability::Available), + "credential_missing" => Ok(BindingAvailability::CredentialMissing), + "store_unavailable" => Ok(BindingAvailability::StoreUnavailable), + _ => Err(invalid_column(9)), + } +} + +const fn encode_kind(value: AccountOperationKind) -> &'static str { + match value { + AccountOperationKind::Add => "add", + AccountOperationKind::Import => "import", + AccountOperationKind::Remove => "remove", + } +} + +fn decode_kind(value: &str) -> rusqlite::Result<AccountOperationKind> { + match value { + "add" => Ok(AccountOperationKind::Add), + "import" => Ok(AccountOperationKind::Import), + "remove" => Ok(AccountOperationKind::Remove), + _ => Err(invalid_column(1)), + } +} + +const fn encode_phase(value: AccountOperationPhase) -> &'static str { + match value { + AccountOperationPhase::IntentRecorded => "intent_recorded", + AccountOperationPhase::CredentialWritten => "credential_written", + AccountOperationPhase::MetadataCommitted => "metadata_committed", + AccountOperationPhase::CompensationPending => "compensation_pending", + AccountOperationPhase::CredentialDeleted => "credential_deleted", + AccountOperationPhase::MetadataDeleted => "metadata_deleted", + } +} + +fn decode_phase(value: &str) -> rusqlite::Result<AccountOperationPhase> { + match value { + "intent_recorded" => Ok(AccountOperationPhase::IntentRecorded), + "credential_written" => Ok(AccountOperationPhase::CredentialWritten), + "metadata_committed" => Ok(AccountOperationPhase::MetadataCommitted), + "compensation_pending" => Ok(AccountOperationPhase::CompensationPending), + "credential_deleted" => Ok(AccountOperationPhase::CredentialDeleted), + "metadata_deleted" => Ok(AccountOperationPhase::MetadataDeleted), + _ => Err(invalid_column(3)), + } +} + +const fn encode_diagnostic(value: OperationDiagnostic) -> &'static str { + match value { + OperationDiagnostic::StorageUnavailable => "storage_unavailable", + OperationDiagnostic::KeyringUnavailable => "keyring_unavailable", + OperationDiagnostic::CredentialMissing => "credential_missing", + OperationDiagnostic::CompensationFailed => "compensation_failed", + OperationDiagnostic::Conflict => "conflict", + OperationDiagnostic::Expired => "expired", + } +} + +fn decode_diagnostic(value: &str) -> rusqlite::Result<OperationDiagnostic> { + match value { + "storage_unavailable" => Ok(OperationDiagnostic::StorageUnavailable), + "keyring_unavailable" => Ok(OperationDiagnostic::KeyringUnavailable), + "credential_missing" => Ok(OperationDiagnostic::CredentialMissing), + "compensation_failed" => Ok(OperationDiagnostic::CompensationFailed), + "conflict" => Ok(OperationDiagnostic::Conflict), + "expired" => Ok(OperationDiagnostic::Expired), + _ => Err(invalid_column(5)), + } +} + +fn invalid_column(index: usize) -> rusqlite::Error { + rusqlite::Error::InvalidColumnType( + index, + "account operation journal".to_owned(), + rusqlite::types::Type::Text, + ) +} + +const fn storage_error() -> SafeError { + SafeError::new( + SafeErrorCode::StorageUnavailable, + SafeMessage::new("The account recovery journal is unavailable."), + ) +} + +const fn corrupt_storage_error() -> SafeError { + SafeError::new( + SafeErrorCode::StorageCorrupt, + SafeMessage::new("The account recovery journal could not be read."), + ) +} + +const fn operation_not_found() -> SafeError { + SafeError::new( + SafeErrorCode::PendingOperationRecoveryRequired, + SafeMessage::new("The account recovery operation was not found."), + ) +} + +const fn operation_conflict() -> SafeError { + SafeError::new( + SafeErrorCode::InvalidApplicationState, + SafeMessage::new("The durable account operation conflicts with existing state."), + ) +} + +#[cfg(test)] +mod tests { + use radroots_studio_application::{ + AccountOperationKind, AccountOperationPhase, DurableOperationKind, DurableOperationPhase, + DurableOperationRepository, DurableOperationStart, DurableRequestId, + DurableTerminalOutcome, OperationDiagnostic, OperationJournal, OperationPriorState, + }; + use radroots_studio_domain::{BindingAvailability, PublicKey, UnixTimestamp}; + + use crate::Database; + + fn public_key(discriminator: u8) -> PublicKey { + let value = match discriminator { + 7 => "0707070707070707070707070707070707070707070707070707070707070707", + 8 => "585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df", + _ => "e0266e3cfb0d2886f91c73f5f868f3b98273713e5fcd97c081663f5518a4b3af", + }; + PublicKey::from_hex(value).expect("valid public key") + } + + #[test] + fn journal_creates_advances_loads_and_finalizes_pending_operations() { + let database = Database::in_memory().expect("database"); + let subject = public_key(7); + let id = database + .begin_operation( + AccountOperationKind::Import, + subject, + UnixTimestamp::from_seconds(10).expect("time"), + ) + .expect("begin"); + database + .update_operation( + id, + AccountOperationPhase::CompensationPending, + UnixTimestamp::from_seconds(11).expect("time"), + Some(OperationDiagnostic::KeyringUnavailable), + ) + .expect("advance"); + + let pending = database.list_pending_operations().expect("pending"); + assert_eq!(pending.len(), 1); + assert_eq!(pending[0].subject(), subject); + assert_eq!(pending[0].kind(), AccountOperationKind::Import); + assert_eq!( + pending[0].phase(), + AccountOperationPhase::CompensationPending + ); + assert_eq!( + pending[0].diagnostic(), + Some(OperationDiagnostic::KeyringUnavailable) + ); + + database.finalize_operation(id).expect("finalize"); + assert!( + database + .list_pending_operations() + .expect("pending") + .is_empty() + ); + } + + #[test] + fn journal_schema_and_rows_exclude_secret_payload_columns() { + let database = Database::in_memory().expect("database"); + database + .begin_operation( + AccountOperationKind::Remove, + public_key(8), + UnixTimestamp::from_seconds(12).expect("time"), + ) + .expect("begin"); + let connection = database.connection(); + let schema: String = connection + .query_row( + "SELECT sql FROM sqlite_master WHERE name = 'operation_journal'", + [], + |row| row.get(0), + ) + .expect("schema"); + assert!(!schema.contains("secret")); + assert!(!schema.contains("payload")); + } + + #[test] + fn durable_repository_replays_matching_requests_and_retains_terminal_receipts() { + let database = Database::in_memory().expect("database"); + let request = DurableRequestId::parse("import:test:1").expect("request"); + let account = public_key(9); + let prior = OperationPriorState::new( + Some(public_key(8)), + Some(BindingAvailability::CredentialMissing), + ); + let started = database + .begin_durable_operation( + &request, + DurableOperationKind::Repair, + account, + Some(4), + prior, + UnixTimestamp::from_seconds(10).expect("time"), + ) + .expect("begin"); + assert!(matches!(started, DurableOperationStart::Started(_))); + let replay = database + .begin_durable_operation( + &request, + DurableOperationKind::Repair, + account, + Some(4), + prior, + UnixTimestamp::from_seconds(11).expect("time"), + ) + .expect("replay"); + assert!(matches!(replay, DurableOperationStart::Existing(_))); + assert!( + database + .begin_durable_operation( + &request, + DurableOperationKind::Remove, + account, + Some(4), + prior, + UnixTimestamp::from_seconds(11).expect("time"), + ) + .is_err() + ); + let missing_request = DurableRequestId::parse("import:test:missing").expect("request"); + assert!( + database + .finalize_durable_operation( + &missing_request, + DurableOperationPhase::IntentRecorded, + DurableTerminalOutcome::Completed, + None, + UnixTimestamp::from_seconds(17).expect("time"), + ) + .is_err() + ); + assert!( + database + .begin_durable_operation( + &request, + DurableOperationKind::Repair, + public_key(8), + Some(4), + prior, + UnixTimestamp::from_seconds(11).expect("time"), + ) + .is_err() + ); + assert!( + database + .begin_durable_operation( + &request, + DurableOperationKind::Repair, + account, + Some(5), + prior, + UnixTimestamp::from_seconds(11).expect("time"), + ) + .is_err() + ); + assert!( + database + .begin_durable_operation( + &request, + DurableOperationKind::Repair, + account, + Some(4), + OperationPriorState::new(None, None), + UnixTimestamp::from_seconds(11).expect("time"), + ) + .is_err() + ); + assert!( + database + .advance_durable_operation( + &request, + DurableOperationPhase::CredentialDeleted, + DurableOperationPhase::Finalized, + UnixTimestamp::from_seconds(11).expect("time"), + None, + ) + .is_err() + ); + database + .advance_durable_operation( + &request, + DurableOperationPhase::IntentRecorded, + DurableOperationPhase::CredentialWritten, + UnixTimestamp::from_seconds(12).expect("time"), + None, + ) + .expect("advance"); + let receipt = database + .finalize_durable_operation( + &request, + DurableOperationPhase::CredentialWritten, + DurableTerminalOutcome::Completed, + Some(5), + UnixTimestamp::from_seconds(13).expect("time"), + ) + .expect("finalize"); + assert_eq!(receipt.resulting_revision(), Some(5)); + assert_eq!( + database + .finalize_durable_operation( + &request, + DurableOperationPhase::CredentialWritten, + DurableTerminalOutcome::Completed, + Some(5), + UnixTimestamp::from_seconds(14).expect("time"), + ) + .expect("receipt replay"), + receipt + ); + assert!( + database + .finalize_durable_operation( + &request, + DurableOperationPhase::CredentialWritten, + DurableTerminalOutcome::Cancelled, + Some(5), + UnixTimestamp::from_seconds(14).expect("time"), + ) + .is_err() + ); + assert!( + database + .finalize_durable_operation( + &request, + DurableOperationPhase::CredentialWritten, + DurableTerminalOutcome::Completed, + Some(6), + UnixTimestamp::from_seconds(14).expect("time"), + ) + .is_err() + ); + let overflow_request = DurableRequestId::parse("import:test:overflow").expect("request"); + database + .begin_durable_operation( + &overflow_request, + DurableOperationKind::Import, + account, + None, + OperationPriorState::new(None, None), + UnixTimestamp::from_seconds(15).expect("time"), + ) + .expect("begin overflow operation"); + assert!( + database + .finalize_durable_operation( + &overflow_request, + DurableOperationPhase::IntentRecorded, + DurableTerminalOutcome::Completed, + Some(u64::MAX), + UnixTimestamp::from_seconds(16).expect("time"), + ) + .is_err() + ); + assert!( + database + .list_unfinished_durable_operations() + .expect("unfinished") + .iter() + .any(|operation| operation.request_id() == &overflow_request) + ); + } +} diff --git a/core/crates/studio_storage/src/lib.rs b/core/crates/studio_storage/src/lib.rs @@ -0,0 +1,20 @@ +#![doc = "Radroots Studio persistence adapters."] +#![cfg_attr(coverage_nightly, feature(coverage_attribute))] + +pub mod account_namespace; +pub mod accounts; +mod compatibility; +pub mod db; +mod installation; +pub mod journal; +// The operating-system credential store requires an explicit, ignored host smoke test. +#[cfg_attr(coverage_nightly, coverage(off))] +pub mod os_keyring; +pub mod profiles; +mod recovery; +mod repair; + +pub use compatibility::{DatabasePreflight, PersistedIdentityIssue, PersistedIdentityIssueKind}; +pub use db::{CURRENT_SCHEMA_VERSION, Database}; +pub use os_keyring::{CREDENTIAL_SERVICE, OsKeyringSecretStore}; +pub use repair::{QuarantineExportReceipt, RepairAuthorization, RepairCandidate}; diff --git a/core/crates/studio_storage/src/os_keyring.rs b/core/crates/studio_storage/src/os_keyring.rs @@ -0,0 +1,138 @@ +use std::sync::{Mutex, MutexGuard}; + +use keyring::{Entry, Error as KeyringError}; +use radroots_studio_application::SecretStore; +use radroots_studio_domain::{PublicKey, SafeError, SafeErrorCode, SafeMessage, SecretKeyInput}; +use zeroize::Zeroizing; + +pub const CREDENTIAL_SERVICE: &str = "org.radroots.studio.nostr"; + +#[derive(Default)] +pub struct OsKeyringSecretStore { + operation_lock: Mutex<()>, +} + +impl OsKeyringSecretStore { + fn entry(public_key: PublicKey) -> Result<Entry, SafeError> { + Entry::new(CREDENTIAL_SERVICE, &public_key.to_hex()).map_err(|_| keyring_unavailable()) + } + + fn operation(&self) -> MutexGuard<'_, ()> { + self.operation_lock + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + } +} + +impl SecretStore for OsKeyringSecretStore { + fn put(&self, public_key: PublicKey, secret: SecretKeyInput) -> Result<(), SafeError> { + let _operation = self.operation(); + let entry = Self::entry(public_key)?; + match entry.get_password() { + Ok(password) => { + drop(Zeroizing::new(password)); + return Err(credential_exists()); + } + Err(KeyringError::NoEntry) => {} + Err(_) => return Err(keyring_unavailable()), + } + secret + .with_exposed_secret(|value| entry.set_password(value)) + .map_err(|_| keyring_unavailable()) + } + + fn load(&self, public_key: PublicKey) -> Result<SecretKeyInput, SafeError> { + let _operation = self.operation(); + let password = Self::entry(public_key)? + .get_password() + .map_err(|error| map_read_error(&error))?; + SecretKeyInput::parse(password) + } + + fn contains(&self, public_key: PublicKey) -> Result<bool, SafeError> { + let _operation = self.operation(); + match Self::entry(public_key)?.get_password() { + Ok(password) => { + drop(Zeroizing::new(password)); + Ok(true) + } + Err(KeyringError::NoEntry) => Ok(false), + Err(_) => Err(keyring_unavailable()), + } + } + + fn delete(&self, public_key: PublicKey) -> Result<(), SafeError> { + let _operation = self.operation(); + Self::entry(public_key)? + .delete_credential() + .map_err(|error| map_read_error(&error)) + } +} + +const fn map_read_error(error: &KeyringError) -> SafeError { + match error { + KeyringError::NoEntry => credential_missing(), + _ => keyring_unavailable(), + } +} + +const fn credential_exists() -> SafeError { + SafeError::new( + SafeErrorCode::AccountAlreadyExists, + SafeMessage::new("The Nostr account credential already exists."), + ) +} + +const fn credential_missing() -> SafeError { + SafeError::new( + SafeErrorCode::CredentialMissing, + SafeMessage::new("The Nostr account credential is missing."), + ) +} + +const fn keyring_unavailable() -> SafeError { + SafeError::new( + SafeErrorCode::KeyringUnavailable, + SafeMessage::new("The operating system credential store is unavailable."), + ) +} + +#[cfg(test)] +mod tests { + use radroots_studio_application::SecretStore; + use radroots_studio_domain::{PublicKey, SecretKeyInput}; + + use super::{CREDENTIAL_SERVICE, OsKeyringSecretStore}; + + #[test] + fn keyring_coordinates_are_stable_and_public() { + let public_key = + PublicKey::from_hex("7e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7") + .expect("valid public key"); + assert_eq!(CREDENTIAL_SERVICE, "org.radroots.studio.nostr"); + assert_eq!( + public_key.to_hex(), + "7e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7" + ); + } + + #[test] + #[ignore = "mutates the current user's operating-system credential store"] + fn real_keyring_smoke_round_trips_and_deletes() { + let store = OsKeyringSecretStore::default(); + let public_key = + PublicKey::from_hex("7e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7") + .expect("valid public key"); + let _ = store.delete(public_key); + store + .put( + public_key, + SecretKeyInput::parse("11".repeat(32)).expect("secret"), + ) + .expect("keyring put"); + assert!(store.contains(public_key).expect("keyring contains")); + let loaded = store.load(public_key).expect("keyring load"); + assert_eq!(loaded.with_exposed_secret(str::len), 64); + store.delete(public_key).expect("keyring delete"); + } +} diff --git a/core/crates/studio_storage/src/profiles.rs b/core/crates/studio_storage/src/profiles.rs @@ -0,0 +1,254 @@ +use radroots_studio_application::{CachedProfile, ProfileRefreshStatus, ProfileRepository}; +use radroots_studio_domain::{ + EventId, Kind0ProfileCandidate, ProfileMetadata, PublicKey, SafeError, SafeErrorCode, + SafeMessage, UnixTimestamp, +}; +use rusqlite::{OptionalExtension, Row, params}; + +use crate::Database; + +impl ProfileRepository for Database { + fn load_profile(&self, public_key: PublicKey) -> Result<Option<CachedProfile>, SafeError> { + self.connection() + .query_row( + "SELECT event_id, event_created_at, name, display_name, nip05, about, picture, \ + refreshed_at, refresh_status FROM profile_cache_v6 WHERE subject_public_key = ?1", + [public_key.to_hex()], + |row| decode_profile(row, public_key), + ) + .optional() + .map_err(|_| corrupt_storage_error()) + } + + fn save_profile(&self, profile: &CachedProfile) -> Result<(), SafeError> { + let candidate = profile.candidate(); + let metadata = candidate.metadata(); + self.connection() + .execute( + "INSERT INTO profile_cache_v6 (subject_public_key, event_id, event_created_at, name, \ + display_name, nip05, about, picture, refreshed_at, refresh_status) \ + VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10) \ + ON CONFLICT(subject_public_key) DO UPDATE SET \ + event_id = excluded.event_id, event_created_at = excluded.event_created_at, \ + name = excluded.name, display_name = excluded.display_name, nip05 = excluded.nip05, \ + about = excluded.about, picture = excluded.picture, \ + refreshed_at = excluded.refreshed_at, refresh_status = excluded.refresh_status \ + WHERE excluded.event_created_at > profile_cache_v6.event_created_at \ + OR (excluded.event_created_at = profile_cache_v6.event_created_at \ + AND excluded.event_id < profile_cache_v6.event_id)", + params![ + candidate.author().to_hex(), + candidate.event_id().to_hex(), + candidate.created_at().as_seconds(), + metadata.name(), + metadata.display_name(), + metadata.nip05(), + metadata.about(), + metadata.picture(), + profile.refreshed_at().as_seconds(), + encode_refresh_status(profile.refresh_status()), + ], + ) + .map(|_| ()) + .map_err(|_| storage_error()) + } + + fn record_refresh_status( + &self, + public_key: PublicKey, + refreshed_at: UnixTimestamp, + status: ProfileRefreshStatus, + ) -> Result<(), SafeError> { + self.connection() + .execute( + "UPDATE profile_cache_v6 SET refreshed_at = ?2, refresh_status = ?3 \ + WHERE subject_public_key = ?1", + params![ + public_key.to_hex(), + refreshed_at.as_seconds(), + encode_refresh_status(status) + ], + ) + .map(|_| ()) + .map_err(|_| storage_error()) + } + + fn remove_profile(&self, public_key: PublicKey) -> Result<(), SafeError> { + self.connection() + .execute( + "DELETE FROM profile_cache_v6 WHERE subject_public_key = ?1", + [public_key.to_hex()], + ) + .map(|_| ()) + .map_err(|_| storage_error()) + } +} + +fn decode_profile(row: &Row<'_>, author: PublicKey) -> rusqlite::Result<CachedProfile> { + let event_id = + EventId::from_hex(row.get::<_, String>(0)?.as_str()).map_err(|_| invalid_column(0))?; + let created_at = UnixTimestamp::from_seconds(row.get(1)?).ok_or_else(|| invalid_column(1))?; + let metadata = ProfileMetadata::new( + row.get(2)?, + row.get(3)?, + row.get(4)?, + row.get(5)?, + row.get(6)?, + ) + .map_err(|_| invalid_column(2))?; + let refreshed_at = UnixTimestamp::from_seconds(row.get(7)?).ok_or_else(|| invalid_column(7))?; + let refresh_status = decode_refresh_status(row.get::<_, String>(8)?.as_str())?; + Ok(CachedProfile::new( + Kind0ProfileCandidate::new(event_id, author, created_at, metadata), + refreshed_at, + refresh_status, + )) +} + +const fn encode_refresh_status(status: ProfileRefreshStatus) -> &'static str { + match status { + ProfileRefreshStatus::Success => "success", + ProfileRefreshStatus::Offline => "offline", + ProfileRefreshStatus::InvalidData => "invalid_data", + } +} + +fn decode_refresh_status(value: &str) -> rusqlite::Result<ProfileRefreshStatus> { + match value { + "success" => Ok(ProfileRefreshStatus::Success), + "offline" => Ok(ProfileRefreshStatus::Offline), + "invalid_data" => Ok(ProfileRefreshStatus::InvalidData), + _ => Err(invalid_column(8)), + } +} + +fn invalid_column(index: usize) -> rusqlite::Error { + rusqlite::Error::InvalidColumnType( + index, + "cached Nostr profile".to_owned(), + rusqlite::types::Type::Text, + ) +} + +const fn storage_error() -> SafeError { + SafeError::new( + SafeErrorCode::StorageUnavailable, + SafeMessage::new("The profile cache is unavailable."), + ) +} + +const fn corrupt_storage_error() -> SafeError { + SafeError::new( + SafeErrorCode::StorageCorrupt, + SafeMessage::new("The profile cache could not be read."), + ) +} + +#[cfg(test)] +mod tests { + use radroots_studio_application::{ + AccountRepository, CachedProfile, ProfileRefreshStatus, ProfileRepository, + }; + use radroots_studio_domain::{ + AccountCreatedAt, AccountIdentity, AccountSummary, BindingAvailability, EventId, + Kind0ProfileCandidate, LocalSignerBinding, ProfileMetadata, PublicKey, UnixTimestamp, + }; + + use crate::Database; + + fn public_key() -> PublicKey { + PublicKey::from_bytes([7; 32]).expect("valid public key") + } + + fn account(public_key: PublicKey) -> AccountSummary { + 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") + } + + fn profile(public_key: PublicKey, id: u8, created_at: i64, name: &str) -> CachedProfile { + CachedProfile::new( + Kind0ProfileCandidate::new( + EventId::from_bytes([id; 32]), + public_key, + UnixTimestamp::from_seconds(created_at).expect("time"), + ProfileMetadata::new(Some(name.to_owned()), None, None, None, None) + .expect("metadata"), + ), + UnixTimestamp::from_seconds(created_at + 1).expect("refresh time"), + ProfileRefreshStatus::Success, + ) + } + + #[test] + fn profile_cache_round_trips_and_records_refresh_status() { + let database = Database::in_memory().expect("database"); + let public_key = public_key(); + database + .insert_account(&account(public_key)) + .expect("account"); + database + .save_profile(&profile(public_key, 1, 10, "Farm")) + .expect("save profile"); + database + .record_refresh_status( + public_key, + UnixTimestamp::from_seconds(20).expect("time"), + ProfileRefreshStatus::Offline, + ) + .expect("record status"); + + let loaded = database + .load_profile(public_key) + .expect("load profile") + .expect("cached profile"); + assert_eq!(loaded.candidate().metadata().name(), Some("Farm")); + assert_eq!(loaded.refreshed_at().as_seconds(), 20); + assert_eq!(loaded.refresh_status(), ProfileRefreshStatus::Offline); + } + + #[test] + fn profile_cache_keeps_newest_then_lowest_event_id() { + let database = Database::in_memory().expect("database"); + let public_key = public_key(); + database + .insert_account(&account(public_key)) + .expect("account"); + database + .save_profile(&profile(public_key, 9, 20, "High ID")) + .expect("initial"); + database + .save_profile(&profile(public_key, 1, 20, "Low ID")) + .expect("equal newer candidate"); + database + .save_profile(&profile(public_key, 0, 10, "Older")) + .expect("older candidate"); + + let loaded = database + .load_profile(public_key) + .expect("load") + .expect("profile"); + assert_eq!(loaded.candidate().metadata().name(), Some("Low ID")); + assert_eq!(loaded.candidate().event_id(), EventId::from_bytes([1; 32])); + } + + #[test] + fn profile_cache_cascades_with_account_removal() { + let database = Database::in_memory().expect("database"); + let public_key = public_key(); + database + .insert_account(&account(public_key)) + .expect("account"); + database + .save_profile(&profile(public_key, 1, 10, "Farm")) + .expect("profile"); + database.remove_account(public_key).expect("remove account"); + + assert_eq!(database.load_profile(public_key).expect("load"), None); + } +} diff --git a/core/crates/studio_storage/src/recovery.rs b/core/crates/studio_storage/src/recovery.rs @@ -0,0 +1,692 @@ +use std::fs::{self, File, OpenOptions}; +use std::io::{Read, Write}; +use std::path::{Path, PathBuf}; + +use hmac::{Hmac, Mac}; +use radroots_studio_domain::{SafeError, SafeErrorCode, SafeMessage}; +use rusqlite::{Connection, MAIN_DB, OpenFlags}; +use sha2::{Digest, Sha256}; +use zeroize::Zeroizing; + +use crate::db::{restrict_directory_permissions, restrict_file_permissions}; + +type HmacSha256 = Hmac<Sha256>; + +const RECOVERY_DIRECTORY_SUFFIX: &str = "recovery"; +const AUTHENTICATION_KEY_FILENAME: &str = "authentication-key-v1"; +const MANIFEST_FORMAT: &str = "radroots-studio-migration-recovery-v1"; + +pub(crate) struct MigrationRecovery { + directory: PathBuf, + backup: PathBuf, + marker: PathBuf, + source_schema: u32, + target_schema: u32, + digest: String, + tag: String, + state: String, +} + +impl MigrationRecovery { + pub(crate) fn prepare( + database_path: &Path, + source_schema: u32, + target_schema: u32, + ) -> Result<Self, SafeError> { + let directory = recovery_directory(database_path)?; + create_recovery_directory(&directory)?; + let key = load_or_create_authentication_key(&directory)?; + let stem = format!("migration-v{source_schema}-to-v{target_schema}"); + let backup = directory.join(format!("{stem}.sqlite3")); + let marker = directory.join(format!("{stem}.marker")); + + if marker.try_exists().map_err(|_| storage_error())? { + let mut recovery = Self::load_existing( + directory, + backup, + marker, + source_schema, + target_schema, + &key, + )?; + if recovery.state == "complete" { + recovery.tag = authentication_tag( + &key, + source_schema, + target_schema, + &recovery.digest, + "prepared", + )?; + recovery.state = "prepared".to_owned(); + recovery.write_marker("prepared", &key)?; + } + return Ok(recovery); + } + if backup.try_exists().map_err(|_| storage_error())? { + return Err(backup_invalid()); + } + + create_verified_backup(database_path, &backup)?; + let digest = file_digest(&backup)?; + let tag = authentication_tag(&key, source_schema, target_schema, &digest, "prepared")?; + let recovery = Self { + directory, + backup, + marker, + source_schema, + target_schema, + digest, + tag, + state: "prepared".to_owned(), + }; + recovery.write_marker("prepared", &key)?; + recovery.verify_backup(&key, "prepared")?; + Ok(recovery) + } + + pub(crate) fn finish(self, current_schema: u32) -> Result<(), SafeError> { + if current_schema != self.target_schema { + return Err(backup_invalid()); + } + let key = load_authentication_key(&self.directory)?; + self.verify_backup(&key, "prepared")?; + self.write_marker("complete", &key) + } + + pub(crate) fn verify_evidence( + database_path: &Path, + source_schema: u32, + target_schema: u32, + ) -> Result<(), SafeError> { + let directory = recovery_directory(database_path)?; + let key = load_authentication_key(&directory)?; + let stem = format!("migration-v{source_schema}-to-v{target_schema}"); + Self::load_existing( + directory.clone(), + directory.join(format!("{stem}.sqlite3")), + directory.join(format!("{stem}.marker")), + source_schema, + target_schema, + &key, + ) + .map(|_| ()) + } + + pub(crate) fn restore( + database_path: &Path, + source_schema: u32, + target_schema: u32, + ) -> Result<(), SafeError> { + let directory = recovery_directory(database_path)?; + let key = load_authentication_key(&directory)?; + let stem = format!("migration-v{source_schema}-to-v{target_schema}"); + let recovery = Self::load_existing( + directory.clone(), + directory.join(format!("{stem}.sqlite3")), + directory.join(format!("{stem}.marker")), + source_schema, + target_schema, + &key, + )?; + recovery.verify_backup(&key, &recovery.state)?; + replace_with_backup(database_path, &recovery.backup) + } + + fn load_existing( + directory: PathBuf, + backup: PathBuf, + marker: PathBuf, + source_schema: u32, + target_schema: u32, + key: &[u8], + ) -> Result<Self, SafeError> { + let manifest = read_bounded_file(&marker, 4_096)?; + let manifest = std::str::from_utf8(&manifest).map_err(|_| backup_invalid())?; + let mut lines = manifest.lines(); + if lines.next() != Some(MANIFEST_FORMAT) + || parse_field(&mut lines, "source_schema")? != source_schema.to_string() + || parse_field(&mut lines, "target_schema")? != target_schema.to_string() + || parse_field(&mut lines, "backup")? + != backup + .file_name() + .ok_or_else(backup_invalid)? + .to_string_lossy() + || lines.clone().count() != 3 + { + return Err(backup_invalid()); + } + let digest = parse_field(&mut lines, "sha256")?; + let state = parse_field(&mut lines, "state")?; + let tag = parse_field(&mut lines, "hmac_sha256")?; + if !matches!(state.as_str(), "prepared" | "complete") { + return Err(backup_invalid()); + } + let recovery = Self { + directory, + backup, + marker, + source_schema, + target_schema, + digest, + tag, + state, + }; + recovery.verify_backup(key, &recovery.state)?; + Ok(recovery) + } + + fn verify_backup(&self, key: &[u8], state: &str) -> Result<(), SafeError> { + if file_digest(&self.backup)? != self.digest { + return Err(backup_invalid()); + } + let expected = authentication_tag( + key, + self.source_schema, + self.target_schema, + &self.digest, + state, + )?; + let expected = decode_hex_32(&expected)?; + let actual = decode_hex_32(&self.tag)?; + if !constant_time_eq(&expected, &actual) { + return Err(backup_invalid()); + } + let flags = OpenFlags::SQLITE_OPEN_READ_ONLY + | OpenFlags::SQLITE_OPEN_NO_MUTEX + | OpenFlags::SQLITE_OPEN_NOFOLLOW; + let connection = + Connection::open_with_flags(&self.backup, flags).map_err(|_| backup_invalid())?; + let integrity: String = connection + .pragma_query_value(None, "quick_check", |row| row.get(0)) + .map_err(|_| backup_invalid())?; + if integrity != "ok" { + return Err(backup_invalid()); + } + Ok(()) + } + + fn write_marker(&self, state: &str, key: &[u8]) -> Result<(), SafeError> { + let tag = authentication_tag( + key, + self.source_schema, + self.target_schema, + &self.digest, + state, + )?; + let content = format!( + "{MANIFEST_FORMAT}\nsource_schema={}\ntarget_schema={}\nbackup={}\nsha256={}\nstate={state}\nhmac_sha256={tag}\n", + self.source_schema, + self.target_schema, + self.backup + .file_name() + .ok_or_else(backup_invalid)? + .to_string_lossy(), + self.digest, + ); + atomic_secure_write(&self.marker, content.as_bytes()) + } +} + +fn recovery_directory(database_path: &Path) -> Result<PathBuf, SafeError> { + let filename = database_path + .file_name() + .ok_or_else(storage_error)? + .to_string_lossy(); + Ok(database_path.with_file_name(format!("{filename}.{RECOVERY_DIRECTORY_SUFFIX}"))) +} + +fn create_recovery_directory(directory: &Path) -> Result<(), SafeError> { + match fs::symlink_metadata(directory) { + Ok(metadata) if metadata.file_type().is_symlink() || !metadata.is_dir() => { + return Err(storage_error()); + } + Ok(_) => {} + Err(error) if error.kind() == std::io::ErrorKind::NotFound => { + fs::create_dir(directory).map_err(|_| storage_error())?; + } + Err(_) => return Err(storage_error()), + } + restrict_directory_permissions(directory) +} + +fn load_or_create_authentication_key(directory: &Path) -> Result<Zeroizing<Vec<u8>>, SafeError> { + let path = directory.join(AUTHENTICATION_KEY_FILENAME); + if path.try_exists().map_err(|_| storage_error())? { + return load_authentication_key(directory); + } + let mut key = Zeroizing::new(vec![0_u8; 32]); + getrandom::getrandom(&mut key).map_err(|_| storage_error())?; + let mut options = OpenOptions::new(); + options.write(true).create_new(true); + #[cfg(unix)] + { + use std::os::unix::fs::OpenOptionsExt; + options.mode(0o600).custom_flags( + i32::try_from((rustix::fs::OFlags::NOFOLLOW | rustix::fs::OFlags::CLOEXEC).bits()) + .map_err(|_| storage_error())?, + ); + } + match options.open(&path) { + Ok(mut file) => { + file.write_all(&key).map_err(|_| storage_error())?; + file.sync_all().map_err(|_| storage_error())?; + restrict_file_permissions(&path)?; + Ok(key) + } + Err(error) if error.kind() == std::io::ErrorKind::AlreadyExists => { + load_authentication_key(directory) + } + Err(_) => Err(storage_error()), + } +} + +fn load_authentication_key(directory: &Path) -> Result<Zeroizing<Vec<u8>>, SafeError> { + let path = directory.join(AUTHENTICATION_KEY_FILENAME); + let key = read_bounded_file(&path, 32)?; + if key.len() != 32 { + return Err(backup_invalid()); + } + Ok(Zeroizing::new(key)) +} + +fn create_verified_backup(source: &Path, destination: &Path) -> Result<(), SafeError> { + let flags = OpenFlags::SQLITE_OPEN_READ_ONLY + | OpenFlags::SQLITE_OPEN_NO_MUTEX + | OpenFlags::SQLITE_OPEN_NOFOLLOW; + let connection = Connection::open_with_flags(source, flags).map_err(|_| backup_invalid())?; + connection + .backup(MAIN_DB, destination, None) + .map_err(|_| backup_invalid())?; + restrict_file_permissions(destination)?; + File::open(destination) + .and_then(|file| file.sync_all()) + .map_err(|_| backup_invalid()) +} + +fn replace_with_backup(database_path: &Path, backup: &Path) -> Result<(), SafeError> { + let parent = database_path.parent().ok_or_else(storage_error)?; + let mut suffix = [0_u8; 8]; + getrandom::getrandom(&mut suffix).map_err(|_| storage_error())?; + let replacement = parent.join(format!(".database-restore-{}.tmp", hex(&suffix))); + let displaced = parent.join(format!(".database-displaced-{}.sqlite3", hex(&suffix))); + let mut source = secure_read(backup)?; + let mut options = OpenOptions::new(); + options.write(true).create_new(true); + #[cfg(unix)] + { + use std::os::unix::fs::OpenOptionsExt; + options.mode(0o600).custom_flags( + i32::try_from((rustix::fs::OFlags::NOFOLLOW | rustix::fs::OFlags::CLOEXEC).bits()) + .map_err(|_| storage_error())?, + ); + } + let result = (|| { + let mut destination = options.open(&replacement).map_err(|_| storage_error())?; + std::io::copy(&mut source, &mut destination).map_err(|_| storage_error())?; + destination.sync_all().map_err(|_| storage_error())?; + restrict_file_permissions(&replacement)?; + fs::rename(database_path, &displaced).map_err(|_| storage_error())?; + if fs::rename(&replacement, database_path).is_err() { + let _ = fs::rename(&displaced, database_path); + return Err(storage_error()); + } + File::open(parent) + .and_then(|directory| directory.sync_all()) + .map_err(|_| storage_error()) + })(); + if result.is_err() { + let _ = fs::remove_file(&replacement); + } + result +} + +fn secure_read(path: &Path) -> Result<File, SafeError> { + let metadata = fs::symlink_metadata(path).map_err(|_| backup_invalid())?; + if metadata.file_type().is_symlink() || !metadata.is_file() { + return Err(backup_invalid()); + } + let mut options = OpenOptions::new(); + options.read(true); + #[cfg(unix)] + { + use std::os::unix::fs::OpenOptionsExt; + options.custom_flags( + i32::try_from((rustix::fs::OFlags::NOFOLLOW | rustix::fs::OFlags::CLOEXEC).bits()) + .map_err(|_| backup_invalid())?, + ); + } + options.open(path).map_err(|_| backup_invalid()) +} + +fn file_digest(path: &Path) -> Result<String, SafeError> { + let metadata = fs::symlink_metadata(path).map_err(|_| backup_invalid())?; + if metadata.file_type().is_symlink() || !metadata.is_file() { + return Err(backup_invalid()); + } + let mut file = secure_read(path)?; + let mut digest = Sha256::new(); + let mut buffer = [0_u8; 64 * 1024]; + loop { + let read = file.read(&mut buffer).map_err(|_| backup_invalid())?; + if read == 0 { + break; + } + digest.update(&buffer[..read]); + } + Ok(hex(&digest.finalize())) +} + +fn read_bounded_file(path: &Path, limit: usize) -> Result<Vec<u8>, SafeError> { + let metadata = fs::symlink_metadata(path).map_err(|_| backup_invalid())?; + if metadata.file_type().is_symlink() + || !metadata.is_file() + || usize::try_from(metadata.len()).map_err(|_| backup_invalid())? > limit + { + return Err(backup_invalid()); + } + let mut options = OpenOptions::new(); + options.read(true); + #[cfg(unix)] + { + use std::os::unix::fs::OpenOptionsExt; + options.custom_flags( + i32::try_from((rustix::fs::OFlags::NOFOLLOW | rustix::fs::OFlags::CLOEXEC).bits()) + .map_err(|_| backup_invalid())?, + ); + } + let file = options.open(path).map_err(|_| backup_invalid())?; + let mut bytes = Vec::with_capacity(usize::try_from(metadata.len()).unwrap_or(0)); + file.take(u64::try_from(limit).map_err(|_| backup_invalid())? + 1) + .read_to_end(&mut bytes) + .map_err(|_| backup_invalid())?; + if bytes.len() > limit { + return Err(backup_invalid()); + } + Ok(bytes) +} + +fn atomic_secure_write(path: &Path, bytes: &[u8]) -> Result<(), SafeError> { + let parent = path.parent().ok_or_else(storage_error)?; + let mut suffix = [0_u8; 8]; + getrandom::getrandom(&mut suffix).map_err(|_| storage_error())?; + let temporary = parent.join(format!(".marker-{}.tmp", hex(&suffix))); + let mut options = OpenOptions::new(); + options.write(true).create_new(true); + #[cfg(unix)] + { + use std::os::unix::fs::OpenOptionsExt; + options.mode(0o600).custom_flags( + i32::try_from((rustix::fs::OFlags::NOFOLLOW | rustix::fs::OFlags::CLOEXEC).bits()) + .map_err(|_| storage_error())?, + ); + } + let result = (|| { + let mut file = options.open(&temporary).map_err(|_| storage_error())?; + file.write_all(bytes).map_err(|_| storage_error())?; + file.sync_all().map_err(|_| storage_error())?; + restrict_file_permissions(&temporary)?; + fs::rename(&temporary, path).map_err(|_| storage_error())?; + File::open(parent) + .and_then(|directory| directory.sync_all()) + .map_err(|_| storage_error()) + })(); + if result.is_err() { + let _ = fs::remove_file(&temporary); + } + result +} + +fn authentication_tag( + key: &[u8], + source_schema: u32, + target_schema: u32, + digest: &str, + state: &str, +) -> Result<String, SafeError> { + let mut mac = HmacSha256::new_from_slice(key).map_err(|_| backup_invalid())?; + mac.update(MANIFEST_FORMAT.as_bytes()); + mac.update(&source_schema.to_be_bytes()); + mac.update(&target_schema.to_be_bytes()); + mac.update(digest.as_bytes()); + mac.update(state.as_bytes()); + Ok(hex(&mac.finalize().into_bytes())) +} + +fn parse_field<'a>( + lines: &mut impl Iterator<Item = &'a str>, + name: &str, +) -> Result<String, SafeError> { + lines + .next() + .and_then(|line| line.strip_prefix(name)) + .and_then(|value| value.strip_prefix('=')) + .map(str::to_owned) + .ok_or_else(backup_invalid) +} + +fn decode_hex_32(value: &str) -> Result<[u8; 32], SafeError> { + if value.len() != 64 { + return Err(backup_invalid()); + } + let mut bytes = [0_u8; 32]; + for (index, pair) in value.as_bytes().chunks_exact(2).enumerate() { + let high = hex_nibble(pair[0]).ok_or_else(backup_invalid)?; + let low = hex_nibble(pair[1]).ok_or_else(backup_invalid)?; + bytes[index] = (high << 4) | low; + } + Ok(bytes) +} + +const fn hex_nibble(byte: u8) -> Option<u8> { + match byte { + b'0'..=b'9' => Some(byte - b'0'), + b'a'..=b'f' => Some(byte - b'a' + 10), + _ => None, + } +} + +fn constant_time_eq(left: &[u8; 32], right: &[u8; 32]) -> bool { + left.iter() + .zip(right) + .fold(0_u8, |difference, (left, right)| { + difference | (left ^ right) + }) + == 0 +} + +fn hex(bytes: &[u8]) -> String { + bytes.iter().map(|byte| format!("{byte:02x}")).collect() +} + +const fn storage_error() -> SafeError { + SafeError::new( + SafeErrorCode::StorageUnavailable, + SafeMessage::new("The application database recovery path is unavailable."), + ) +} + +const fn backup_invalid() -> SafeError { + SafeError::new( + SafeErrorCode::StorageBackupInvalid, + SafeMessage::new("The application database recovery backup is invalid."), + ) +} + +#[cfg(test)] +mod tests { + use std::fs; + use std::path::Path; + + use rusqlite::Connection; + use tempfile::tempdir; + + use super::{ + AUTHENTICATION_KEY_FILENAME, MANIFEST_FORMAT, MigrationRecovery, atomic_secure_write, + constant_time_eq, create_recovery_directory, decode_hex_32, file_digest, hex, hex_nibble, + load_authentication_key, load_or_create_authentication_key, parse_field, read_bounded_file, + recovery_directory, replace_with_backup, secure_read, + }; + + fn sqlite_database(path: &Path) { + let connection = Connection::open(path).expect("open sqlite database"); + connection + .execute("CREATE TABLE durable_probe (value INTEGER NOT NULL)", []) + .expect("create probe table"); + connection + .execute("INSERT INTO durable_probe (value) VALUES (7)", []) + .expect("insert probe row"); + } + + #[test] + fn migration_recovery_authenticates_finishes_reopens_and_restores() { + let directory = tempdir().expect("temporary directory"); + let database = directory.path().join("studio.sqlite3"); + sqlite_database(&database); + + let recovery = MigrationRecovery::prepare(&database, 5, 10).expect("prepare recovery"); + MigrationRecovery::verify_evidence(&database, 5, 10).expect("prepared evidence"); + assert!(recovery.finish(9).is_err()); + + MigrationRecovery::prepare(&database, 5, 10) + .expect("reopen prepared recovery") + .finish(10) + .expect("finish recovery"); + MigrationRecovery::verify_evidence(&database, 5, 10).expect("complete evidence"); + + MigrationRecovery::prepare(&database, 5, 10) + .expect("reopen complete recovery") + .finish(10) + .expect("finish reopened recovery"); + fs::write(&database, b"not sqlite").expect("corrupt active database"); + MigrationRecovery::restore(&database, 5, 10).expect("restore authenticated backup"); + let connection = Connection::open(&database).expect("open restored database"); + let value: i64 = connection + .query_row("SELECT value FROM durable_probe", [], |row| row.get(0)) + .expect("restored row"); + assert_eq!(value, 7); + } + + #[test] + fn recovery_manifest_rejects_every_tampered_authority_field() { + let directory = tempdir().expect("temporary directory"); + let database = directory.path().join("studio.sqlite3"); + sqlite_database(&database); + let recovery = MigrationRecovery::prepare(&database, 5, 10).expect("prepare recovery"); + let original = fs::read_to_string(&recovery.marker).expect("read marker"); + let backup_name = recovery + .backup + .file_name() + .expect("backup name") + .to_string_lossy(); + let cases = [ + original.replacen(MANIFEST_FORMAT, "wrong-format", 1), + original.replacen("source_schema=5", "source_schema=4", 1), + original.replacen("target_schema=10", "target_schema=11", 1), + original.replacen(&format!("backup={backup_name}"), "backup=other.sqlite3", 1), + original.replacen("sha256=", "unexpected=value\nsha256=", 1), + original.replacen("state=prepared", "state=invalid", 1), + original.replacen("sha256=", "sha256=00", 1), + original.replacen("hmac_sha256=", "hmac_sha256=gg", 1), + { + let mut lines = original.lines().map(str::to_owned).collect::<Vec<_>>(); + let tag = lines + .iter_mut() + .find(|line| line.starts_with("hmac_sha256=")) + .expect("tag field"); + let replacement = if tag.ends_with('0') { '1' } else { '0' }; + tag.pop(); + tag.push(replacement); + format!("{}\n", lines.join("\n")) + }, + ]; + for tampered in cases { + fs::write(&recovery.marker, tampered).expect("write tampered marker"); + assert!(MigrationRecovery::verify_evidence(&database, 5, 10).is_err()); + } + fs::write(&recovery.marker, original).expect("restore marker"); + MigrationRecovery::verify_evidence(&database, 5, 10).expect("restored evidence"); + } + + #[test] + fn recovery_helpers_reject_invalid_paths_sizes_and_encodings() { + let directory = tempdir().expect("temporary directory"); + let regular = directory.path().join("regular"); + fs::write(&regular, b"abc").expect("write regular file"); + let child = directory.path().join("child"); + fs::create_dir(&child).expect("create child directory"); + + assert!(recovery_directory(Path::new("/")).is_err()); + assert!(create_recovery_directory(&regular).is_err()); + assert!(create_recovery_directory(&regular.join("nested")).is_err()); + assert!(secure_read(&child).is_err()); + assert!(file_digest(&child).is_err()); + assert!(read_bounded_file(&child, 4).is_err()); + assert!(read_bounded_file(&regular, 2).is_err()); + assert_eq!( + read_bounded_file(&regular, 3).expect("bounded read"), + b"abc" + ); + + let missing_key_dir = directory.path().join("missing-key"); + fs::create_dir(&missing_key_dir).expect("create missing key directory"); + assert!(load_authentication_key(&missing_key_dir).is_err()); + fs::write( + missing_key_dir.join(AUTHENTICATION_KEY_FILENAME), + [0_u8; 31], + ) + .expect("write short key"); + assert!(load_authentication_key(&missing_key_dir).is_err()); + + assert!(decode_hex_32("00").is_err()); + assert!(decode_hex_32(&format!("g0{}", "00".repeat(31))).is_err()); + assert!(decode_hex_32(&format!("0g{}", "00".repeat(31))).is_err()); + let zeros = decode_hex_32(&"00".repeat(32)).expect("decode zeros"); + assert!(constant_time_eq(&zeros, &[0_u8; 32])); + assert!(!constant_time_eq(&zeros, &[1_u8; 32])); + assert_eq!(hex(&[0, 15, 255]), "000fff"); + assert_eq!(hex_nibble(b'9'), Some(9)); + assert_eq!(hex_nibble(b'f'), Some(15)); + assert_eq!(hex_nibble(b'G'), None); + + let mut valid = ["field=value"].into_iter(); + assert_eq!(parse_field(&mut valid, "field").expect("field"), "value"); + let mut invalid = ["other=value"].into_iter(); + assert!(parse_field(&mut invalid, "field").is_err()); + let mut missing = std::iter::empty(); + assert!(parse_field(&mut missing, "field").is_err()); + + let absent_parent = directory.path().join("absent").join("marker"); + assert!(atomic_secure_write(&absent_parent, b"marker").is_err()); + + let orphan_database = directory.path().join("orphan.sqlite3"); + sqlite_database(&orphan_database); + let orphan_directory = recovery_directory(&orphan_database).expect("recovery directory"); + create_recovery_directory(&orphan_directory).expect("create recovery directory"); + load_or_create_authentication_key(&orphan_directory).expect("authentication key"); + fs::write( + orphan_directory.join("migration-v5-to-v10.sqlite3"), + b"orphan backup", + ) + .expect("orphan backup"); + assert!(MigrationRecovery::prepare(&orphan_database, 5, 10).is_err()); + + let missing_database = directory.path().join("missing-database.sqlite3"); + assert!(replace_with_backup(&missing_database, &regular).is_err()); + } + + #[cfg(unix)] + #[test] + fn recovery_helpers_reject_symlink_inputs() { + use std::os::unix::fs::symlink; + + let directory = tempdir().expect("temporary directory"); + let regular = directory.path().join("regular"); + fs::write(&regular, b"abc").expect("write regular file"); + let link = directory.path().join("link"); + symlink(&regular, &link).expect("create symlink"); + assert!(secure_read(&link).is_err()); + assert!(file_digest(&link).is_err()); + assert!(read_bounded_file(&link, 3).is_err()); + assert!(create_recovery_directory(&link).is_err()); + } +} diff --git a/core/crates/studio_storage/src/repair.rs b/core/crates/studio_storage/src/repair.rs @@ -0,0 +1,403 @@ +use std::fs::{self, File, OpenOptions}; +use std::io::Read; +use std::path::{Path, PathBuf}; + +use hmac::{Hmac, Mac}; +use radroots_studio_domain::{SafeError, SafeErrorCode, SafeMessage}; +use rusqlite::{Connection, MAIN_DB, OpenFlags}; +use sha2::{Digest, Sha256}; +use zeroize::Zeroizing; + +use crate::compatibility::{DatabasePreflight, preflight}; +use crate::db::{CURRENT_SCHEMA_VERSION, restrict_file_permissions}; + +type HmacSha256 = Hmac<Sha256>; +const EXPORT_DOMAIN: &[u8] = b"radroots-studio-quarantine-export-v1"; +const REPAIR_DOMAIN: &[u8] = b"radroots-studio-repair-candidate-v1"; + +pub struct RepairAuthorization(Zeroizing<[u8; 32]>); + +impl RepairAuthorization { + /// Moves an exact 256-bit caller authorization secret into zeroizing storage. + /// + /// # Errors + /// + /// Returns a safe authorization error for every other input length. + pub fn from_bytes(bytes: Vec<u8>) -> Result<Self, SafeError> { + let bytes = Zeroizing::new(bytes); + let value = <[u8; 32]>::try_from(bytes.as_slice()).map_err(|_| unauthorized())?; + Ok(Self(Zeroizing::new(value))) + } + + fn expose(&self) -> &[u8; 32] { + &self.0 + } +} + +pub struct QuarantineExportReceipt { + path: PathBuf, + sha256: String, + authentication_tag: String, +} + +impl QuarantineExportReceipt { + #[must_use] + pub fn path(&self) -> &Path { + &self.path + } + + #[must_use] + pub fn sha256(&self) -> &str { + &self.sha256 + } + + #[must_use] + pub fn authentication_tag(&self) -> &str { + &self.authentication_tag + } +} + +pub struct RepairCandidate { + path: PathBuf, + sha256: String, + authentication_tag: String, +} + +impl RepairCandidate { + #[must_use] + pub fn path(&self) -> &Path { + &self.path + } +} + +pub(crate) fn export_quarantined( + source: &Path, + destination: &Path, + authorization: &RepairAuthorization, +) -> Result<QuarantineExportReceipt, SafeError> { + if !matches!(preflight(source)?, DatabasePreflight::Quarantined { .. }) { + return Err(not_quarantined()); + } + ensure_new_destination(destination)?; + let flags = OpenFlags::SQLITE_OPEN_READ_ONLY + | OpenFlags::SQLITE_OPEN_NO_MUTEX + | OpenFlags::SQLITE_OPEN_NOFOLLOW; + let connection = Connection::open_with_flags(source, flags).map_err(|_| storage_error())?; + if connection.backup(MAIN_DB, destination, None).is_err() { + let _ = fs::remove_file(destination); + return Err(storage_error()); + } + restrict_file_permissions(destination)?; + File::open(destination) + .and_then(|file| file.sync_all()) + .map_err(|_| storage_error())?; + let sha256 = digest_file(destination)?; + let authentication_tag = authenticate(authorization, EXPORT_DOMAIN, &sha256)?; + Ok(QuarantineExportReceipt { + path: destination.to_path_buf(), + sha256, + authentication_tag, + }) +} + +pub(crate) fn authenticate_candidate( + path: &Path, + authorization: &RepairAuthorization, +) -> Result<RepairCandidate, SafeError> { + if !matches!( + preflight(path)?, + DatabasePreflight::Ready { schema_version } if schema_version <= CURRENT_SCHEMA_VERSION + ) { + return Err(storage_error()); + } + let sha256 = digest_file(path)?; + let authentication_tag = authenticate(authorization, REPAIR_DOMAIN, &sha256)?; + Ok(RepairCandidate { + path: path.to_path_buf(), + sha256, + authentication_tag, + }) +} + +pub(crate) fn install_candidate( + target: &Path, + candidate: &RepairCandidate, + authorization: &RepairAuthorization, +) -> Result<(), SafeError> { + if !matches!(preflight(target)?, DatabasePreflight::Quarantined { .. }) { + return Err(not_quarantined()); + } + let digest = digest_file(&candidate.path)?; + if digest != candidate.sha256 + || authenticate(authorization, REPAIR_DOMAIN, &digest)? != candidate.authentication_tag + { + return Err(unauthorized()); + } + if !matches!(preflight(&candidate.path)?, DatabasePreflight::Ready { .. }) { + return Err(storage_error()); + } + let parent = target.parent().ok_or_else(storage_error)?; + let replacement = parent.join(".authenticated-repair.tmp"); + if replacement.try_exists().map_err(|_| storage_error())? { + return Err(storage_error()); + } + copy_secure(&candidate.path, &replacement)?; + let retained = parent.join("studio.sqlite3.quarantined-evidence"); + if retained.try_exists().map_err(|_| storage_error())? { + let _ = fs::remove_file(&replacement); + return Err(storage_error()); + } + fs::rename(target, &retained).map_err(|_| storage_error())?; + if fs::rename(&replacement, target).is_err() { + let _ = fs::rename(&retained, target); + let _ = fs::remove_file(&replacement); + return Err(storage_error()); + } + File::open(parent) + .and_then(|directory| directory.sync_all()) + .map_err(|_| storage_error()) +} + +fn ensure_new_destination(path: &Path) -> Result<(), SafeError> { + if path.try_exists().map_err(|_| storage_error())? { + return Err(storage_error()); + } + let parent = path.parent().ok_or_else(storage_error)?; + let metadata = fs::symlink_metadata(parent).map_err(|_| storage_error())?; + if metadata.file_type().is_symlink() || !metadata.is_dir() { + return Err(storage_error()); + } + Ok(()) +} + +fn copy_secure(source: &Path, destination_path: &Path) -> Result<(), SafeError> { + let mut source = secure_read(source)?; + let mut options = OpenOptions::new(); + options.write(true).create_new(true); + #[cfg(unix)] + { + use std::os::unix::fs::OpenOptionsExt; + options.mode(0o600).custom_flags( + i32::try_from((rustix::fs::OFlags::NOFOLLOW | rustix::fs::OFlags::CLOEXEC).bits()) + .map_err(|_| storage_error())?, + ); + } + let mut destination = options + .open(destination_path) + .map_err(|_| storage_error())?; + std::io::copy(&mut source, &mut destination).map_err(|_| storage_error())?; + destination.sync_all().map_err(|_| storage_error())?; + restrict_file_permissions(destination_path) +} + +fn secure_read(path: &Path) -> Result<File, SafeError> { + let metadata = fs::symlink_metadata(path).map_err(|_| storage_error())?; + if metadata.file_type().is_symlink() || !metadata.is_file() { + return Err(storage_error()); + } + let mut options = OpenOptions::new(); + options.read(true); + #[cfg(unix)] + { + use std::os::unix::fs::OpenOptionsExt; + options.custom_flags( + i32::try_from((rustix::fs::OFlags::NOFOLLOW | rustix::fs::OFlags::CLOEXEC).bits()) + .map_err(|_| storage_error())?, + ); + } + options.open(path).map_err(|_| storage_error()) +} + +fn digest_file(path: &Path) -> Result<String, SafeError> { + let mut file = secure_read(path)?; + let mut digest = Sha256::new(); + let mut buffer = [0_u8; 64 * 1024]; + loop { + let read = file.read(&mut buffer).map_err(|_| storage_error())?; + if read == 0 { + break; + } + digest.update(&buffer[..read]); + } + Ok(hex(&digest.finalize())) +} + +fn authenticate( + authorization: &RepairAuthorization, + domain: &[u8], + digest: &str, +) -> Result<String, SafeError> { + let mut hmac = + HmacSha256::new_from_slice(authorization.expose()).map_err(|_| unauthorized())?; + hmac.update(domain); + hmac.update(digest.as_bytes()); + Ok(hex(&hmac.finalize().into_bytes())) +} + +fn hex(bytes: &[u8]) -> String { + bytes.iter().map(|byte| format!("{byte:02x}")).collect() +} + +const fn unauthorized() -> SafeError { + SafeError::new( + SafeErrorCode::RepairUnauthorized, + SafeMessage::new("The database repair authorization is invalid."), + ) +} + +const fn not_quarantined() -> SafeError { + SafeError::new( + SafeErrorCode::InvalidApplicationState, + SafeMessage::new("The database is not in quarantine."), + ) +} + +const fn storage_error() -> SafeError { + SafeError::new( + SafeErrorCode::StorageUnavailable, + SafeMessage::new("The database repair operation could not be completed."), + ) +} + +#[cfg(test)] +mod tests { + use std::fs; + use std::io::Write; + + use rusqlite::Connection; + use tempfile::tempdir; + + use super::{ + REPAIR_DOMAIN, RepairAuthorization, RepairCandidate, authenticate, authenticate_candidate, + copy_secure, digest_file, ensure_new_destination, export_quarantined, hex, + install_candidate, secure_read, + }; + use crate::Database; + + fn quarantined_database(path: &std::path::Path) { + drop(Database::open(path).expect("current database")); + let connection = Connection::open(path).expect("open database"); + connection + .execute( + "INSERT INTO account_identities (public_key, npub, created_at) VALUES (?1, ?2, 1)", + [ + "00".repeat(32), + "npub1qurswpc8qurswpc8qurswpc8qurswpc8qurswpc8qurswpc8qursnvjvl7".to_owned(), + ], + ) + .expect("invalid identity fixture"); + connection + .execute( + "INSERT INTO local_signer_bindings (account_public_key, binding_public_key, binding_kind, availability) VALUES (?1, ?1, 'local_secret', 'available')", + ["00".repeat(32)], + ) + .expect("binding fixture"); + } + + #[test] + fn repair_authority_and_candidate_reject_invalid_states() { + assert!(RepairAuthorization::from_bytes(vec![0_u8; 31]).is_err()); + assert!(RepairAuthorization::from_bytes(vec![0_u8; 33]).is_err()); + let authorization = RepairAuthorization::from_bytes(vec![0x41; 32]).expect("authorization"); + assert_eq!( + authenticate(&authorization, b"domain", "digest") + .expect("authentication tag") + .len(), + 64 + ); + assert_eq!(hex(&[0, 15, 255]), "000fff"); + + let directory = tempdir().expect("temporary directory"); + let ready = directory.path().join("ready.sqlite3"); + drop(Database::open(&ready).expect("ready database")); + let candidate = authenticate_candidate(&ready, &authorization).expect("candidate"); + assert_eq!(candidate.path(), ready); + + let missing = directory.path().join("missing.sqlite3"); + assert!(authenticate_candidate(&missing, &authorization).is_err()); + let export = directory.path().join("export.sqlite3"); + assert!(export_quarantined(&ready, &export, &authorization).is_err()); + assert!(install_candidate(&ready, &candidate, &authorization).is_err()); + } + + #[test] + fn repair_file_boundaries_reject_existing_non_file_and_missing_parent_paths() { + let directory = tempdir().expect("temporary directory"); + let regular = directory.path().join("regular"); + fs::write(&regular, b"repair material").expect("write regular file"); + let child = directory.path().join("child"); + fs::create_dir(&child).expect("create child directory"); + + assert!(ensure_new_destination(&regular).is_err()); + assert!(ensure_new_destination(&regular.join("nested")).is_err()); + assert!(secure_read(&child).is_err()); + assert_eq!(digest_file(&regular).expect("digest").len(), 64); + + let copied = directory.path().join("copied"); + copy_secure(&regular, &copied).expect("secure copy"); + assert_eq!(fs::read(&copied).expect("copied bytes"), b"repair material"); + assert!(copy_secure(&regular, &copied).is_err()); + assert!(copy_secure(&child, &directory.path().join("invalid-copy")).is_err()); + } + + #[test] + fn repair_installation_rejects_tampering_quarantined_candidates_and_staging_collisions() { + let directory = tempdir().expect("temporary directory"); + let authorization = RepairAuthorization::from_bytes(vec![0x41; 32]).expect("authorization"); + let target = directory.path().join("studio.sqlite3"); + quarantined_database(&target); + let candidate_path = directory.path().join("candidate.sqlite3"); + drop(Database::open(&candidate_path).expect("candidate database")); + let candidate = authenticate_candidate(&candidate_path, &authorization).expect("candidate"); + + fs::OpenOptions::new() + .append(true) + .open(&candidate_path) + .expect("open candidate") + .write_all(b"tamper") + .expect("tamper candidate"); + assert!(install_candidate(&target, &candidate, &authorization).is_err()); + + let quarantined_candidate_path = directory.path().join("quarantined-candidate.sqlite3"); + quarantined_database(&quarantined_candidate_path); + let digest = digest_file(&quarantined_candidate_path).expect("candidate digest"); + let quarantined_candidate = RepairCandidate { + path: quarantined_candidate_path, + sha256: digest.clone(), + authentication_tag: authenticate(&authorization, REPAIR_DOMAIN, &digest) + .expect("candidate tag"), + }; + assert!(install_candidate(&target, &quarantined_candidate, &authorization).is_err()); + + let candidate_path = directory.path().join("candidate-two.sqlite3"); + drop(Database::open(&candidate_path).expect("candidate database")); + let candidate = authenticate_candidate(&candidate_path, &authorization).expect("candidate"); + let replacement = directory.path().join(".authenticated-repair.tmp"); + fs::write(&replacement, b"occupied").expect("occupied replacement"); + assert!(install_candidate(&target, &candidate, &authorization).is_err()); + fs::remove_file(&replacement).expect("remove occupied replacement"); + + let retained = directory.path().join("studio.sqlite3.quarantined-evidence"); + fs::write(&retained, b"occupied").expect("occupied retained evidence"); + assert!(install_candidate(&target, &candidate, &authorization).is_err()); + assert!(!replacement.exists()); + } + + #[cfg(unix)] + #[test] + fn repair_file_boundaries_reject_symlinks() { + use std::os::unix::fs::symlink; + + let directory = tempdir().expect("temporary directory"); + let regular = directory.path().join("regular"); + fs::write(&regular, b"repair material").expect("write regular file"); + let link = directory.path().join("link"); + symlink(&regular, &link).expect("create file symlink"); + assert!(secure_read(&link).is_err()); + assert!(digest_file(&link).is_err()); + + let directory_link = directory.path().join("directory-link"); + symlink(directory.path(), &directory_link).expect("create directory symlink"); + assert!(ensure_new_destination(&directory_link.join("export")).is_err()); + } +} diff --git a/core/crates/studio_storage/tests/redaction.rs b/core/crates/studio_storage/tests/redaction.rs @@ -0,0 +1,52 @@ +use std::fs; + +use radroots_studio_application::{AccountOperationKind, AccountRepository, OperationJournal}; +use radroots_studio_domain::{ + AccountCreatedAt, AccountIdentity, AccountSummary, BindingAvailability, LocalSignerBinding, + PublicKey, UnixTimestamp, +}; +use radroots_studio_storage::Database; +use tempfile::tempdir; + +const SECRET_HEX: &str = "1111111111111111111111111111111111111111111111111111111111111111"; +const SECRET_NSEC: &str = "nsec1vl029mgpspedva04g90vltkh6fvh240zqtv9k0t9af8935ke9laqsnlfe5"; +fn assert_redacted(bytes: &[u8]) { + assert!( + !bytes + .windows(SECRET_HEX.len()) + .any(|value| value == SECRET_HEX.as_bytes()) + ); + assert!( + !bytes + .windows(SECRET_NSEC.len()) + .any(|value| value == SECRET_NSEC.as_bytes()) + ); + assert!(!bytes.windows(5).any(|value| value == b"nsec1")); +} + +#[test] +fn redaction_guards_sqlite_schema_and_non_secret_records() { + let directory = tempdir().expect("directory"); + let path = directory.path().join("studio.sqlite3"); + { + let database = Database::open(&path).expect("database"); + let public_key = PublicKey::from_bytes([7; 32]).expect("valid public key"); + let account = AccountSummary::new( + AccountIdentity::derive(public_key).expect("identity"), + LocalSignerBinding::new(public_key, BindingAvailability::Available), + None, + AccountCreatedAt::new(UnixTimestamp::from_seconds(1).expect("time")), + None, + ) + .expect("account"); + database.insert_account(&account).expect("account"); + database + .begin_operation( + AccountOperationKind::Add, + account.public_key(), + UnixTimestamp::from_seconds(2).expect("time"), + ) + .expect("journal"); + } + assert_redacted(&fs::read(path).expect("database bytes")); +}