diff --git a/desktop/src-tauri/src/app_state.rs b/desktop/src-tauri/src/app_state.rs index fc90e6ab1..f96d569dc 100644 --- a/desktop/src-tauri/src/app_state.rs +++ b/desktop/src-tauri/src/app_state.rs @@ -9,12 +9,12 @@ use std::{ use nostr::{Keys, ToBech32}; use tauri::{AppHandle, Manager}; -#[cfg(feature = "mesh-llm")] use tokio::sync::Mutex as AsyncMutex; use crate::huddle::HuddleState; pub(crate) use crate::identity_storage::{IdentityStorage, RecoveryState, ResolvedIdentity}; use crate::managed_agents::config_bridge::SessionConfigCache; +use crate::managed_agents::scope::WorkspaceAgentScope; use crate::managed_agents::{ManagedAgentPairRuntime, ManagedAgentRuntimeKey}; pub struct AppState { @@ -93,7 +93,26 @@ pub struct AppState { /// a newer imported key during concurrent calls. Deliberately separate from /// `keys` so readers (signing, get_identity, etc.) are not blocked during /// keyring I/O. - pub identity_mutation: Mutex<()>, + /// + /// **Layer 1 async lock** — callers may `.await` while holding this guard. + /// Lock order: `identity_mutation` → `workspace_transition` → Mesh + /// `rearm_lock` → `mesh_llm_runtime`. Converted from `Mutex<()>` to + /// `AsyncMutex<()>` so callers in async Tauri commands can hold it across + /// awaits without blocking the executor. + pub identity_mutation: AsyncMutex<()>, + /// Serializes workspace transitions (`apply_workspace` and live identity + /// import). Taken after `identity_mutation` in the lock order. + /// + /// **Layer 1 async lock** — callers may `.await` while holding this guard. + pub workspace_transition: AsyncMutex<()>, + /// The active workspace agent scope — `None` from boot until the first + /// successful `apply_workspace`. Every agent command fails closed on `None`. + /// There is NO fallback to the legacy unscoped root. + /// + /// Protected by the **Layer 2 synchronous commit epoch** (no `.await` while + /// `managed_agents_store_lock` is held). Read outside the lock for + /// non-mutating "capture at entry" use via `capture_active_scope()`. + pub active_agent_scope: Mutex>, /// Set when the boot-time Phase 2 reset attempted a wipe but verification /// failed. The sentinel is preserved so the next relaunch retries. All /// identity-dependent setup is skipped; the frontend shows a reset-failed @@ -211,7 +230,9 @@ pub fn build_app_state() -> AppState { managed_agent_profile_reconcile_enabled: AtomicBool::new(true), shutdown_started: AtomicBool::new(false), managed_agent_runtime_transition: Mutex::new(()), - identity_mutation: Mutex::new(()), + identity_mutation: AsyncMutex::new(()), + workspace_transition: AsyncMutex::new(()), + active_agent_scope: Mutex::new(None), managed_agents_store_lock: Mutex::new(()), channel_templates_store_lock: Mutex::new(()), managed_agent_processes: Mutex::new(HashMap::new()), diff --git a/desktop/src-tauri/src/commands/identity.rs b/desktop/src-tauri/src/commands/identity.rs index 33ecf3cfc..562aa7258 100644 --- a/desktop/src-tauri/src/commands/identity.rs +++ b/desktop/src-tauri/src/commands/identity.rs @@ -226,7 +226,7 @@ pub(crate) fn create_backup_with_log_n( // Serialize against import_identity/persist_current_identity: the blob // must be derived from — and persisted for — one stable identity. Also // caps KDF concurrency at one. - let _mutation_guard = state.identity_mutation.lock().map_err(|e| e.to_string())?; + let _mutation_guard = state.identity_mutation.blocking_lock(); // Recovery mode (lost/locked) → Err, same gate as signing. let keys = state.signing_keys()?; @@ -351,7 +351,7 @@ pub async fn import_identity( // full function body so a concurrent stale persist can't overwrite // this import. let state = app_handle.state::(); - let _mutation_guard = state.identity_mutation.lock().map_err(|e| e.to_string())?; + let _mutation_guard = state.identity_mutation.blocking_lock(); let data_dir = app_handle .path() @@ -473,7 +473,7 @@ pub async fn persist_current_identity( // concurrent import_identity cannot complete between our check and // our persist, which would let the stale ephemeral key overwrite the // imported one. - let _mutation_guard = state.identity_mutation.lock().map_err(|e| e.to_string())?; + let _mutation_guard = state.identity_mutation.blocking_lock(); if !state .identity_lost diff --git a/desktop/src-tauri/src/commands/identity_key_backup_tests.rs b/desktop/src-tauri/src/commands/identity_key_backup_tests.rs index c36af6687..ee8feb751 100644 --- a/desktop/src-tauri/src/commands/identity_key_backup_tests.rs +++ b/desktop/src-tauri/src/commands/identity_key_backup_tests.rs @@ -121,7 +121,7 @@ fn concurrent_identity_swap_vs_backup_is_serialized() { std::thread::spawn(move || { // Mirrors import_identity's locking: mutation guard held // across the key swap. - let _guard = state.identity_mutation.lock().unwrap(); + let _guard = state.identity_mutation.blocking_lock(); *state.keys.lock().unwrap() = key_b; }) }; diff --git a/desktop/src-tauri/src/commands/personas/inbound.rs b/desktop/src-tauri/src/commands/personas/inbound.rs index d7ffecef2..bfc2f0167 100644 --- a/desktop/src-tauri/src/commands/personas/inbound.rs +++ b/desktop/src-tauri/src/commands/personas/inbound.rs @@ -121,11 +121,15 @@ fn reconcile_inbound_persona_event_blocking( // Resolve inbound vs. any pending local edit before touching the store, in // the scope the event ARRIVED on. A workspace switch since arrival leaves // this event to its own community's store — dropping it here is what keeps - // community A's head out of community B's database. + // community A's head out of community B's database. Match on both relay and + // owner: an in-flight old-owner event on the same relay must not land in + // the new owner's active store after an identity switch. + let arrival_owner_pubkey = event.pubkey.to_hex(); let Some(scope) = crate::managed_agents::retention::arrival_retention_scope( &app, &state, &arrival_relay_url, + &arrival_owner_pubkey, )? else { return Ok(()); @@ -261,11 +265,16 @@ fn reconcile_inbound_tombstone( // Resolve against the retained tombstone row (keyed by the target // coordinate, F2c) so a re-received tombstone or one older than a pending - // local edit is a no-op. Scoped to the arrival community, so a workspace - // switch since arrival drops the tombstone instead of retaining it — and - // deleting a record — in the wrong community's store. - let Some(scope) = - crate::managed_agents::retention::arrival_retention_scope(app, state, arrival_relay_url)? + // local edit is a no-op. Scoped to the arrival community + owner, so a + // workspace switch since arrival drops the tombstone instead of retaining + // it — and deleting a record — in the wrong community's or owner's store. + let tombstone_owner_pubkey = event.pubkey.to_hex(); + let Some(scope) = crate::managed_agents::retention::arrival_retention_scope( + app, + state, + arrival_relay_url, + &tombstone_owner_pubkey, + )? else { return Ok(()); }; diff --git a/desktop/src-tauri/src/managed_agents/mod.rs b/desktop/src-tauri/src/managed_agents/mod.rs index 772d707f2..2da852eed 100644 --- a/desktop/src-tauri/src/managed_agents/mod.rs +++ b/desktop/src-tauri/src/managed_agents/mod.rs @@ -30,6 +30,7 @@ pub mod retention; mod runtime; mod runtime_commands; mod runtime_types; +pub(crate) mod scope; pub(crate) mod snapshot_avatar; pub(crate) mod spawn_hash; pub(crate) mod storage; diff --git a/desktop/src-tauri/src/managed_agents/retention.rs b/desktop/src-tauri/src/managed_agents/retention.rs index 7e97fa1f5..2b463517d 100644 --- a/desktop/src-tauri/src/managed_agents/retention.rs +++ b/desktop/src-tauri/src/managed_agents/retention.rs @@ -9,10 +9,10 @@ use std::path::{Path, PathBuf}; use std::time::{Duration, Instant}; use rusqlite::{params, Connection, OptionalExtension}; -use sha2::{Digest, Sha256}; use tauri::AppHandle; use crate::app_state::AppState; +use crate::managed_agents::scope::derive_scope_id; mod legacy_migration; pub use legacy_migration::migrate_legacy_retention_db; @@ -30,20 +30,30 @@ pub struct RetentionScope { } /// Decide whether `scope` — the workspace's active retention scope — is the one -/// that owns an event delivered by `arrival_relay_url`. +/// that owns an event delivered by `arrival_relay_url` from `arrival_owner_pubkey`. /// /// Inbound reconcile resolves its retention database when it PROCESSES an event, /// while the event belongs to the community that DELIVERED it. `None` means a -/// workspace switch happened in between and the caller must drop the event -/// rather than file community A's event into community B's store. +/// workspace switch happened in between (relay or owner changed), and the caller +/// must drop the event rather than file community A's event into community B's store. +/// +/// Matching both relay and owner ensures that an in-flight old-owner event on +/// the same relay cannot land in the new owner's active store after an identity +/// switch. /// /// The comparison goes through the same normalization -/// [`scoped_retention_db_path`] hashes, so "same relay" can never disagree with +/// [`scoped_retention_db_path`] hashes, so "same scope" can never disagree with /// "same database". -pub fn scope_for_arrival(scope: RetentionScope, arrival_relay_url: &str) -> Option { - let same_scope = +pub fn scope_for_arrival( + scope: RetentionScope, + arrival_relay_url: &str, + arrival_owner_pubkey: &str, +) -> Option { + let same_relay = normalized_relay_scope(&scope.relay_url) == normalized_relay_scope(arrival_relay_url); - same_scope.then_some(scope) + let same_owner = scope.owner_keys.public_key().to_hex().to_ascii_lowercase() + == arrival_owner_pubkey.trim().to_ascii_lowercase(); + (same_relay && same_owner).then_some(scope) } /// Relay-URL form that identifies a retention scope: equivalent workspace URLs @@ -54,15 +64,11 @@ fn normalized_relay_scope(relay_url: &str) -> &str { /// Resolve the retention database path for a relay + owner pair. /// -/// The normalized scope is hashed so relay URLs never become path components. -/// Trimming a trailing slash keeps equivalent workspace URLs on one scope. +/// Delegates to [`derive_scope_id`] from the shared scope module so the hash +/// is byte-identical between the retention DB path and the definition store +/// path — "same scope" can never disagree between the two subsystems. pub fn scoped_retention_db_path(base_dir: &Path, relay_url: &str, owner_pubkey: &str) -> PathBuf { - let normalized_relay = normalized_relay_scope(relay_url); - let mut hasher = Sha256::new(); - hasher.update(owner_pubkey.trim().to_ascii_lowercase().as_bytes()); - hasher.update(b"\0"); - hasher.update(normalized_relay.as_bytes()); - let scope_id = hex::encode(hasher.finalize()); + let scope_id = derive_scope_id(relay_url, owner_pubkey); base_dir.join("retention").join(format!("{scope_id}.db")) } @@ -89,7 +95,7 @@ pub fn active_retention_scope(app: &AppHandle, state: &AppState) -> Result Result, String> { Ok(scope_for_arrival( active_retention_scope(app, state)?, arrival_relay_url, + arrival_owner_pubkey, )) } @@ -496,12 +504,13 @@ mod tests { }; let community_a = scoped_retention_db_path(base, "wss://a.example", &owner); - // "Same relay" and "same database" must never disagree: every URL the + // "Same relay + owner" and "same database" must never disagree: every URL the // match accepts has to hash to the scope's own db path, and every URL it // rejects has to hash somewhere else. for equivalent in ["wss://a.example", "wss://a.example/", " wss://a.example "] { assert_eq!( - scope_for_arrival(scope("wss://a.example"), equivalent).map(|scope| scope.db_path), + scope_for_arrival(scope("wss://a.example"), equivalent, &owner) + .map(|scope| scope.db_path), Some(community_a.clone()), "{equivalent}" ); @@ -512,14 +521,23 @@ mod tests { ); } + // Different relay must not match. assert!( - scope_for_arrival(scope("wss://b.example"), "wss://a.example").is_none(), + scope_for_arrival(scope("wss://b.example"), "wss://a.example", &owner).is_none(), "an event from community A must not be filed while community B is active" ); assert_ne!( scoped_retention_db_path(base, "wss://b.example", &owner), community_a ); + + // Different owner on same relay must not match. + let other_keys = nostr::Keys::generate(); + let other_owner = other_keys.public_key().to_hex(); + assert!( + scope_for_arrival(scope("wss://a.example"), "wss://a.example", &other_owner).is_none(), + "an event from a different owner must not be filed into the active scope" + ); } #[test] diff --git a/desktop/src-tauri/src/managed_agents/scope.rs b/desktop/src-tauri/src/managed_agents/scope.rs new file mode 100644 index 000000000..f513adee3 --- /dev/null +++ b/desktop/src-tauri/src/managed_agents/scope.rs @@ -0,0 +1,306 @@ +//! Workspace-scoped agent definition store. +//! +//! Every agent definition (managed-agents.json, teams.json, +//! global-agent-config.json) lives under a `(relay_url, owner_pubkey)` scope +//! so that definitions created in workspace A never appear in workspace B's +//! store or relay. The scope identity uses the same sha256 derivation as the +//! retention database, so "same scope" can never disagree between the two +//! subsystems. +//! +//! # Scoped layout +//! +//! ```text +//! /agents/scopes//managed-agents.json +//! /agents/scopes//teams.json +//! /agents/scopes//global-agent-config.json +//! ``` +//! +//! # Active scope lifecycle +//! +//! `Option` is `None` from boot until the first +//! successful `apply_workspace`. Every agent command fails closed on `None` +//! ("no active workspace"). There is NO fallback to the legacy unscoped root; +//! that fallback would recreate split-brain storage. +//! +//! # Transition model (four stages) +//! +//! See [`crate::commands::workspace`] for the full transition machine. +//! 1. **Prepare (reversible):** derive/init target scope while old scope stays active. +//! 2. **Drain (journaled, compensating):** stop old-scope runtimes; on failure +//! compensate by restarting the journaled set; return `applied: false` with +//! explicit degradation when compensation itself fails. +//! 3. **Commit (infallible critical section):** pure in-memory swaps, no I/O. +//! 4. **Post-commit (non-rollback):** nest regen, event sync, new-scope restore; +//! failures surface as degradation on an `applied: true` result. +//! +//! # Lock architecture (two layers) +//! +//! **Layer 1 — async serialization (Tokio mutexes, awaits OK):** +//! `identity_mutation` → `workspace_transition` → Mesh `rearm_lock` → `mesh_llm_runtime` +//! +//! **Layer 2 — synchronous commit epoch (no `.await` while any guard held):** +//! `managed_agent_runtime_transition` → `managed_agents_store_lock` → +//! `managed_agent_processes` → short commit locks (relay override, keys, +//! active_agent_scope). +//! +//! Generation checks bridge the layers: state read under Layer 1 is +//! revalidated by generation inside the Layer 2 epoch immediately before +//! commit. + +use std::path::PathBuf; +use std::sync::atomic::{AtomicU64, Ordering}; + +use sha2::{Digest, Sha256}; + +/// The single scope authority for a workspace's agent definition store. +/// +/// Immutable after creation — every field is `pub` for read access only; +/// mutations always produce a new `WorkspaceAgentScope`. Callers capture one +/// scope at the entry of any operation that crosses an `.await` or a thread +/// boundary and thread it through every load/save via `_at(scope)` APIs. +/// A stale commit (generation mismatch) must abort; a stale spawn must +/// additionally terminate its child and remove its receipt. +#[derive(Debug, Clone)] +pub struct WorkspaceAgentScope { + /// The sha256 scope identifier — byte-identical to the retention DB's + /// derivation so the two subsystems can never disagree about ownership. + pub scope_id: String, + /// The normalized relay URL for this scope. + pub relay_url: String, + /// The owner's hex pubkey. + pub owner_pubkey: String, + /// The scoped definitions directory: `/scopes//`. + pub definitions_dir: PathBuf, + /// Monotonically increasing generation counter; incremented on every scope + /// change (including identity import that clears the active scope to None). + /// Used by long-running operations to detect a mid-flight workspace switch. + pub generation: u64, +} + +/// Process-lifetime generation counter. Monotonically incremented every time +/// the active scope changes (new scope committed or scope cleared). Operations +/// that cross awaits read this at entry and re-validate before commit. +static SCOPE_GENERATION: AtomicU64 = AtomicU64::new(0); + +/// Increment and return the new generation. Called by the commit stage of +/// `apply_workspace` and by identity import when the active scope is cleared. +pub(crate) fn next_scope_generation() -> u64 { + SCOPE_GENERATION.fetch_add(1, Ordering::AcqRel) + 1 +} + +/// Read the current generation without incrementing. +pub fn current_scope_generation() -> u64 { + SCOPE_GENERATION.load(Ordering::Acquire) +} + +/// Relay-URL normalization used to derive a scope identifier. Must be +/// identical to `normalized_relay_scope` in `retention.rs` so the sha256 +/// output is byte-identical. +pub(crate) fn normalize_relay_for_scope(relay_url: &str) -> &str { + relay_url.trim().trim_end_matches('/') +} + +/// Derive the sha256 scope identifier for a `(relay_url, owner_pubkey)` pair. +/// +/// **This is the canonical derivation** — both the retention DB path +/// (`retention::scoped_retention_db_path`) and the definition scope directory +/// go through this function. The hash encodes the pair so relay URLs never +/// become path components. +/// +/// byte-identical to `retention::scoped_retention_db_path`'s inner hash. +pub fn derive_scope_id(relay_url: &str, owner_pubkey: &str) -> String { + let normalized_relay = normalize_relay_for_scope(relay_url); + let mut hasher = Sha256::new(); + hasher.update(owner_pubkey.trim().to_ascii_lowercase().as_bytes()); + hasher.update(b"\0"); + hasher.update(normalized_relay.as_bytes()); + hex::encode(hasher.finalize()) +} + +/// Resolve the definition scope directory for a `(relay_url, owner_pubkey)` +/// pair under `base_dir`. +/// +/// Layout: `/scopes//` +pub fn scoped_definitions_dir(base_dir: &std::path::Path, scope_id: &str) -> PathBuf { + base_dir.join("scopes").join(scope_id) +} + +impl WorkspaceAgentScope { + /// Construct a new scope, deriving the scope_id and definitions_dir from + /// the relay/owner pair. + pub fn new( + relay_url: String, + owner_pubkey: String, + base_dir: &std::path::Path, + generation: u64, + ) -> Self { + let scope_id = derive_scope_id(&relay_url, &owner_pubkey); + let definitions_dir = scoped_definitions_dir(base_dir, &scope_id); + Self { + scope_id, + relay_url, + owner_pubkey, + definitions_dir, + generation, + } + } + + /// Ensure the definitions directory exists. + pub fn ensure_dir(&self) -> Result<(), String> { + std::fs::create_dir_all(&self.definitions_dir).map_err(|e| { + format!( + "failed to create scope dir {}: {e}", + self.definitions_dir.display() + ) + }) + } + + /// Path to the scoped `managed-agents.json`. + pub fn managed_agents_path(&self) -> PathBuf { + self.definitions_dir.join("managed-agents.json") + } + + /// Path to the scoped `teams.json`. + pub fn teams_path(&self) -> PathBuf { + self.definitions_dir.join("teams.json") + } + + /// Path to the scoped `global-agent-config.json`. + pub fn global_config_path(&self) -> PathBuf { + self.definitions_dir.join("global-agent-config.json") + } +} + +/// Result returned by `apply_workspace` / `import_identity` drain-then-commit. +#[derive(Debug, Clone, serde::Serialize)] +pub struct WorkspaceApplyResult { + /// `true` when the new scope was committed; `false` when drain or + /// compensation failed (old scope is still active). + pub applied: bool, + /// Non-empty when the workspace applied but some post-commit step (nest + /// regen, event sync, runtime restore) failed. The workspace IS active; + /// the degradation is informational. Also populated on drain-failure with + /// the specific runtime(s) that could not be stopped or restored. + pub degraded: Vec, +} + +impl WorkspaceApplyResult { + pub fn success() -> Self { + Self { + applied: true, + degraded: Vec::new(), + } + } + + pub fn with_degradation(mut self, msg: impl Into) -> Self { + self.degraded.push(msg.into()); + self + } + + pub fn drain_failed(msg: impl Into) -> Self { + Self { + applied: false, + degraded: vec![msg.into()], + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + use std::path::Path; + + /// The scope-id derivation must be byte-identical to + /// `retention::scoped_retention_db_path`'s inner hash. + /// + /// `scoped_retention_db_path` computes: + /// sha256(owner.trim().to_ascii_lowercase() + "\0" + normalized_relay) + /// and encodes as hex. We verify parity by computing both and asserting + /// equality. + #[test] + fn test_scope_id_parity_with_retention_hash() { + use sha2::{Digest, Sha256}; + + let relay = "wss://relay.example.com/"; + let owner = "AABBCCDD".repeat(8); // 64-char hex + + // What retention.rs computes: + let normalized = relay.trim().trim_end_matches('/'); + let mut hasher = Sha256::new(); + hasher.update(owner.trim().to_ascii_lowercase().as_bytes()); + hasher.update(b"\0"); + hasher.update(normalized.as_bytes()); + let expected = hex::encode(hasher.finalize()); + + // What our helper computes: + let got = derive_scope_id(relay, &owner); + + assert_eq!( + got, expected, + "scope_id must be byte-identical to retention hash" + ); + } + + #[test] + fn test_scope_id_trailing_slash_normalization() { + let owner = "aa".repeat(32); + assert_eq!( + derive_scope_id("wss://a.example/", &owner), + derive_scope_id("wss://a.example", &owner), + "trailing slash must produce same scope_id" + ); + } + + #[test] + fn test_scope_id_separates_relay_and_owner() { + let owner_a = "aa".repeat(32); + let owner_b = "bb".repeat(32); + let relay_a = "wss://a.example"; + let relay_b = "wss://b.example"; + + assert_ne!( + derive_scope_id(relay_a, &owner_a), + derive_scope_id(relay_b, &owner_a) + ); + assert_ne!( + derive_scope_id(relay_a, &owner_a), + derive_scope_id(relay_a, &owner_b) + ); + assert_eq!( + derive_scope_id(relay_a, &owner_a), + derive_scope_id(relay_a, &owner_a) + ); + } + + #[test] + fn test_scoped_definitions_dir_layout() { + let base = Path::new("/data/agents"); + let scope_id = "abcdef1234"; + let dir = scoped_definitions_dir(base, scope_id); + assert_eq!(dir, base.join("scopes").join(scope_id)); + } + + #[test] + fn test_workspace_agent_scope_paths() { + let base = std::env::temp_dir(); + let scope = + WorkspaceAgentScope::new("wss://relay.example.com".into(), "aa".repeat(32), &base, 0); + assert_eq!( + scope.managed_agents_path(), + scope.definitions_dir.join("managed-agents.json") + ); + assert_eq!(scope.teams_path(), scope.definitions_dir.join("teams.json")); + assert_eq!( + scope.global_config_path(), + scope.definitions_dir.join("global-agent-config.json") + ); + } + + #[test] + fn test_generation_increments_monotonically() { + let before = current_scope_generation(); + let next = next_scope_generation(); + assert_eq!(next, before + 1); + assert_eq!(current_scope_generation(), next); + } +}