mirror of
https://github.com/block/buzz.git
synced 2026-08-18 06:50:31 +02:00
feat(desktop): add workspace-scoped agent definition store scope model
Introduce the `WorkspaceAgentScope` type and the scaffolding for a
four-stage workspace transition state machine (Phase 1 foundation).
## Scope model (managed_agents/scope.rs)
- `WorkspaceAgentScope { scope_id, relay_url, owner_pubkey, definitions_dir,
generation }` — the single scope authority for a workspace's agent
definition store. Immutable; callers capture one scope at operation
entry and thread it through `_at(scope)` APIs.
- `derive_scope_id(relay_url, owner_pubkey)` — the canonical sha256
derivation, byte-identical to the retention DB derivation. Both
subsystems now go through one shared helper so "same scope" can never
disagree between definitions and retention.
- `next_scope_generation()` / `current_scope_generation()` — global
monotonic counter incremented on every scope change or identity-import
clear. Long-running operations read at entry and revalidate before
commit; a stale commit aborts.
- `WorkspaceApplyResult { applied, degraded }` — typed result for the
four-stage transition machine (prepare / drain / commit / post-commit).
- Scoped layout: `agents/scopes/<scope_id>/{managed-agents.json,
teams.json, global-agent-config.json}`.
## AppState additions (app_state.rs)
- `identity_mutation: AsyncMutex<()>` (was `Mutex<()>`) — Layer 1 async
lock; callers may `.await` while holding it. Converted so the workspace
transition machine can hold it across awaits without blocking the
executor.
- `workspace_transition: AsyncMutex<()>` — serializes workspace
transitions (`apply_workspace` and live identity import). Lock order:
identity_mutation → workspace_transition → Mesh rearm → mesh_llm_runtime.
- `active_agent_scope: Mutex<Option<WorkspaceAgentScope>>` — `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.
## Retention parity (retention.rs)
- `scoped_retention_db_path` now delegates to `derive_scope_id` instead
of inlining its own sha256, making the hash provably identical.
- `scope_for_arrival` / `arrival_retention_scope` extended to match on
both relay AND owner pubkey. An in-flight old-owner event on the same
relay can no longer land in the new owner's active store after an
identity switch.
## Inbound reconcile (commands/personas/inbound.rs)
- Both `arrival_retention_scope` call sites pass the event's pubkey as
the owner dimension, closing the identity-switch cross-contamination gap.
## Caller updates
- `identity.rs`: three `identity_mutation.lock().map_err()` callers
converted to `.blocking_lock()` (Tokio async mutex's sync-context
variant, safe from `spawn_blocking` threads).
- `identity_key_backup_tests.rs`: test thread mirror updated to match.
Co-authored-by: Will Pfleger <pfleger.will@gmail.com>
Signed-off-by: Will Pfleger <pfleger.will@gmail.com>
This commit is contained in:
co-authored by
Will Pfleger
parent
a5dbdf5e61
commit
7fe049359c
@@ -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<Option<WorkspaceAgentScope>>,
|
||||
/// 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()),
|
||||
|
||||
@@ -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::<AppState>();
|
||||
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
|
||||
|
||||
@@ -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;
|
||||
})
|
||||
};
|
||||
|
||||
@@ -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(());
|
||||
};
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -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<RetentionScope> {
|
||||
let same_scope =
|
||||
pub fn scope_for_arrival(
|
||||
scope: RetentionScope,
|
||||
arrival_relay_url: &str,
|
||||
arrival_owner_pubkey: &str,
|
||||
) -> Option<RetentionScope> {
|
||||
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<Reten
|
||||
}
|
||||
|
||||
/// Snapshot the active relay + owner, but only when it is the scope that owns
|
||||
/// events delivered by `arrival_relay_url`.
|
||||
/// events delivered by `arrival_relay_url` from `arrival_owner_pubkey`.
|
||||
///
|
||||
/// Resolving the scope and matching it in one step is what closes the gap: the
|
||||
/// returned scope is both the one that will be written to and the one the event
|
||||
@@ -99,10 +105,12 @@ pub fn arrival_retention_scope(
|
||||
app: &AppHandle,
|
||||
state: &AppState,
|
||||
arrival_relay_url: &str,
|
||||
arrival_owner_pubkey: &str,
|
||||
) -> Result<Option<RetentionScope>, 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]
|
||||
|
||||
@@ -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
|
||||
//! <app-data>/agents/scopes/<scope_id>/managed-agents.json
|
||||
//! <app-data>/agents/scopes/<scope_id>/teams.json
|
||||
//! <app-data>/agents/scopes/<scope_id>/global-agent-config.json
|
||||
//! ```
|
||||
//!
|
||||
//! # Active scope lifecycle
|
||||
//!
|
||||
//! `Option<WorkspaceAgentScope>` 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: `<agents-base>/scopes/<scope_id>/`.
|
||||
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: `<base_dir>/scopes/<scope_id>/`
|
||||
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<String>,
|
||||
}
|
||||
|
||||
impl WorkspaceApplyResult {
|
||||
pub fn success() -> Self {
|
||||
Self {
|
||||
applied: true,
|
||||
degraded: Vec::new(),
|
||||
}
|
||||
}
|
||||
|
||||
pub fn with_degradation(mut self, msg: impl Into<String>) -> Self {
|
||||
self.degraded.push(msg.into());
|
||||
self
|
||||
}
|
||||
|
||||
pub fn drain_failed(msg: impl Into<String>) -> 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);
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user