feat(desktop): add scoped _at() storage APIs and retarget event-sync chokepoints

## Scoped storage APIs

Add path-based `_at(definitions_dir)` variants alongside every
`app: &AppHandle` storage chokepoint. These are the primary API for
long-running operations that have captured a `WorkspaceAgentScope` at
their entry point:

- `storage.rs`: `load_agent_store_at`, `load_managed_agents_at`,
  `load_agent_definitions_at`, `save_managed_agents_at`,
  `save_agent_definitions_at`, `managed_agents_store_path_at`; the
  internal `write_agent_store` now delegates to `write_agent_store_to_path`
  which is shared with the new scoped write path.
- `teams.rs`: `teams_store_path_at`, `load_teams_at`, `save_teams_at`.
- `global_config/mod.rs`: `global_config_path_at`,
  `load_global_agent_config_at`, `save_global_agent_config_at`;
  the load path is factored into `load_global_agent_config_from_path`.

## AppState scope helpers

- `capture_active_scope()` — snapshot of current `Option<WorkspaceAgentScope>`.
  Callers crossing `.await` or thread boundaries capture at entry.
- `commit_active_scope(scope)` — infallible commit-stage setter (Layer 2).
- `clear_active_scope()` — clear + generation bump for identity import
  drain and prepare-stage rollback.

## Event-sync retarget

`run_event_sync`, `spawn_event_sync`, `migrate_personas_to_events`,
`migrate_teams_to_events`, and `reconcile_agents_to_events` all gain a
`definitions_dir: &Path` / `PathBuf` parameter. They no longer resolve the
base dir from `AppHandle` — the caller passes the scoped definitions dir
directly, closing the bypass that read from the legacy unscoped root.

## apply_workspace scope commit

After applying relay + keys, `apply_workspace` derives a
`WorkspaceAgentScope` from the effective (relay, owner) pair and commits
it via `commit_active_scope`. The immediately following `spawn_event_sync`
call reads the committed scope via `capture_active_scope()`, so event sync
for this apply uses the scoped definitions dir. A legacy-root fallback is
preserved during the Phase 1→2 transition period for pre-apply boot callers.

Co-authored-by: Will Pfleger <pfleger.will@gmail.com>
Signed-off-by: Will Pfleger <pfleger.will@gmail.com>
This commit is contained in:
npub1mn7jgtj4w2pd0g0zeuhxsa6jy6p0rewxz4kujt98my82ahfmp72sxjexk7
2026-08-02 22:53:30 -04:00
co-authored by Will Pfleger
parent 7fe049359c
commit 108a31e9c1
7 changed files with 288 additions and 34 deletions
+34
View File
@@ -267,6 +267,40 @@ impl AppState {
self.huddle_state.lock().map_err(|e| e.to_string())
}
/// Capture a snapshot of the active workspace agent scope.
///
/// Operations that cross an `.await` or a thread boundary MUST capture the
/// scope at entry (before the first await) and thread it through every
/// `_at(scope)` API. A stale commit (generation mismatch) must abort.
///
/// Returns `None` when no workspace has been applied yet — callers must
/// fail closed on `None`; there is NO fallback to the legacy unscoped root.
pub fn capture_active_scope(&self) -> Option<WorkspaceAgentScope> {
self.active_agent_scope.lock().ok().and_then(|g| g.clone())
}
/// Set the active workspace agent scope. Called by the infallible commit
/// stage of the workspace transition machine.
///
/// Must be called while holding `managed_agents_store_lock` (Layer 2) with
/// no `.await` pending — the commit stage is the only caller.
pub(crate) fn commit_active_scope(&self, scope: WorkspaceAgentScope) {
if let Ok(mut g) = self.active_agent_scope.lock() {
*g = Some(scope);
}
}
/// Clear the active workspace agent scope and bump the generation.
///
/// Called by live identity import (drain → persist → clear → re-apply)
/// and by the prepare stage when rolling back after a failed pipeline.
pub(crate) fn clear_active_scope(&self) {
if let Ok(mut g) = self.active_agent_scope.lock() {
*g = None;
}
crate::managed_agents::scope::next_scope_generation();
}
pub fn get_session_cache(&self, key: &ManagedAgentRuntimeKey) -> Option<SessionConfigCache> {
self.session_config_cache.lock().ok()?.get(key).cloned()
}
+40 -1
View File
@@ -166,7 +166,7 @@ pub async fn apply_workspace(
// ── Apply all state changes (nothing below can fail) ──────────────────
{
let mut override_guard = state.relay_url_override.lock().map_err(|e| e.to_string())?;
*override_guard = Some(relay_url);
*override_guard = Some(relay_url.clone());
}
// Reset the Rust-side admission gate when switching workspace/community,
// matching `resetRateLimitGate()` on the TS side (useCommunityInit.ts:38).
@@ -184,6 +184,32 @@ pub async fn apply_workspace(
.managed_agent_profile_reconcile_enabled
.store(!agent_managed_profiles.unwrap_or(false), Ordering::Release);
// ── Commit the active workspace agent scope ───────────────────────────
// Derive the scope from the now-applied relay + owner keys and store it
// as the active scope. All subsequent store reads/writes and event-sync
// calls use the scoped definitions directory.
{
let owner_pubkey = state
.keys
.lock()
.map_err(|e| e.to_string())?
.public_key()
.to_hex();
let base_dir = crate::managed_agents::managed_agents_base_dir(&app).unwrap_or_default();
let generation = crate::managed_agents::scope::next_scope_generation();
let scope = crate::managed_agents::scope::WorkspaceAgentScope::new(
relay_url,
owner_pubkey,
&base_dir,
generation,
);
// Ensure the scoped dir exists so first-apply callers find it.
if let Err(e) = scope.ensure_dir() {
eprintln!("buzz-desktop: failed to create scope dir: {e}");
}
state.commit_active_scope(scope);
}
// ── Filesystem side-effect (non-fatal) ────────────────────────────────
// Persist the *effective* repos_dir (None when the candidate failed
// validation) for the backend to read at boot, then re-point REPOS to
@@ -222,10 +248,23 @@ pub async fn apply_workspace(
// stranded tombstones and archive requests publish on this boot
// instead of being abandoned by the storage cutover.
migrate_legacy_retention_into(&restore_app, &scope);
// Use the scoped definitions directory for event sync.
// After apply_workspace commits the scope, capture_active_scope()
// returns the scoped dir; fall back to the legacy unscoped root
// only when no scope has been committed (pre-Phase-2 transition).
let definitions_dir = state
.capture_active_scope()
.map(|s| s.definitions_dir)
.unwrap_or_else(|| {
crate::managed_agents::managed_agents_base_dir(&restore_app).unwrap_or_default()
});
crate::event_sync::spawn_event_sync(
restore_app.clone(),
scope.owner_keys,
scope.db_path,
definitions_dir,
)
}
Err(error) => {
+27 -21
View File
@@ -13,10 +13,23 @@ use std::path::Path;
/// `sync_team_personas` wrote in [`crate::migration::run_boot_migrations`]
/// (see its `# Ordering` guard). Event signing needs the resolved owner keys,
/// so this runs after identity resolution, not in the boot migrations.
pub fn run_event_sync(app: &tauri::AppHandle, owner_keys: &nostr::Keys, db_path: &Path) {
migrate_personas_to_events(app, owner_keys, db_path);
migrate_teams_to_events(app, owner_keys, db_path);
crate::managed_agents::reconcile::reconcile_agents_to_events(app, owner_keys, db_path);
///
/// `definitions_dir` is the scoped definitions directory for this workspace
/// (`WorkspaceAgentScope::definitions_dir`). Reads personas/teams/agents from
/// that directory rather than the legacy unscoped `agents/` root.
pub fn run_event_sync(
app: &tauri::AppHandle,
owner_keys: &nostr::Keys,
db_path: &Path,
definitions_dir: &Path,
) {
migrate_personas_to_events(definitions_dir, owner_keys, db_path);
migrate_teams_to_events(definitions_dir, owner_keys, db_path);
crate::managed_agents::reconcile::reconcile_agents_to_events(
definitions_dir,
owner_keys,
db_path,
);
}
/// Spawn the best-effort event reconcile off the synchronous Tauri setup path.
@@ -29,10 +42,11 @@ pub fn spawn_event_sync(
app: tauri::AppHandle,
owner_keys: nostr::Keys,
db_path: std::path::PathBuf,
definitions_dir: std::path::PathBuf,
) {
tauri::async_runtime::spawn(async move {
if let Err(e) = tauri::async_runtime::spawn_blocking(move || {
run_event_sync(&app, &owner_keys, &db_path);
run_event_sync(&app, &owner_keys, &db_path, &definitions_dir);
})
.await
{
@@ -61,14 +75,10 @@ pub fn spawn_event_sync(
/// `pending_sync = 1` for later relay publish. Migration succeeds on local
/// write, not relay acknowledgment. Every retained row is a real signed
/// event — there is no placeholder path.
pub fn migrate_personas_to_events(app: &tauri::AppHandle, keys: &nostr::Keys, db_path: &Path) {
use crate::managed_agents::managed_agents_base_dir;
let Ok(base_dir) = managed_agents_base_dir(app) else {
return;
};
match migrate_personas_in_dir_at(&base_dir, keys, db_path) {
///
/// `definitions_dir` is the scoped definitions directory (`WorkspaceAgentScope::definitions_dir`).
pub fn migrate_personas_to_events(definitions_dir: &Path, keys: &nostr::Keys, db_path: &Path) {
match migrate_personas_in_dir_at(definitions_dir, keys, db_path) {
Ok(0) => {}
Ok(migrated) => {
eprintln!(
@@ -219,14 +229,10 @@ fn migrate_personas_in_dir_at(
///
/// Must run after the persisted identity is resolved (it signs each event with
/// the owner's keys).
pub fn migrate_teams_to_events(app: &tauri::AppHandle, keys: &nostr::Keys, db_path: &Path) {
use crate::managed_agents::managed_agents_base_dir;
let Ok(base_dir) = managed_agents_base_dir(app) else {
return;
};
match migrate_teams_in_dir_at(&base_dir, keys, db_path) {
///
/// `definitions_dir` is the scoped definitions directory (`WorkspaceAgentScope::definitions_dir`).
pub fn migrate_teams_to_events(definitions_dir: &Path, keys: &nostr::Keys, db_path: &Path) {
match migrate_teams_in_dir_at(definitions_dir, keys, db_path) {
Ok(0) => {}
Ok(migrated) => {
eprintln!("buzz-desktop: team-event-migration: {migrated} teams migrated to retention");
@@ -178,15 +178,33 @@ fn global_config_path(app: &AppHandle) -> Result<std::path::PathBuf, String> {
Ok(managed_agents_base_dir(app)?.join("global-agent-config.json"))
}
/// Scoped variant: resolve `global-agent-config.json` under a workspace scope's
/// definitions directory.
pub(crate) fn global_config_path_at(definitions_dir: &std::path::Path) -> std::path::PathBuf {
definitions_dir.join("global-agent-config.json")
}
/// Load the global agent config from disk.
///
/// Returns the default (all-empty) config if the file does not exist yet.
pub fn load_global_agent_config(app: &AppHandle) -> Result<GlobalAgentConfig, String> {
let path = global_config_path(app)?;
load_global_agent_config_from_path(&path)
}
/// Scoped variant: load global agent config from the given definitions dir.
pub(crate) fn load_global_agent_config_at(
definitions_dir: &std::path::Path,
) -> Result<GlobalAgentConfig, String> {
let path = global_config_path_at(definitions_dir);
load_global_agent_config_from_path(&path)
}
fn load_global_agent_config_from_path(path: &std::path::Path) -> Result<GlobalAgentConfig, String> {
if !path.exists() {
return Ok(GlobalAgentConfig::default());
}
let content = std::fs::read_to_string(&path)
let content = std::fs::read_to_string(path)
.map_err(|e| format!("failed to read global agent config: {e}"))?;
serde_json::from_str(&content).map_err(|e| format!("failed to parse global agent config: {e}"))
}
@@ -207,6 +225,26 @@ pub fn save_global_agent_config(app: &AppHandle, config: &GlobalAgentConfig) ->
atomic_write_json_restricted(&path, &payload)
}
/// Scoped variant: save global agent config into the given definitions dir.
pub(crate) fn save_global_agent_config_at(
definitions_dir: &std::path::Path,
config: &GlobalAgentConfig,
) -> Result<(), String> {
let mut config = config.clone();
strip_empty_env_vars(&mut config);
normalize_global_config_fields(&mut config);
let path = global_config_path_at(definitions_dir);
// Ensure the directory exists (scoped dirs are created lazily).
if let Some(parent) = path.parent() {
std::fs::create_dir_all(parent)
.map_err(|e| format!("failed to create scoped store dir: {e}"))?;
}
let payload = serde_json::to_vec_pretty(&config)
.map_err(|e| format!("failed to serialize global agent config: {e}"))?;
atomic_write_json_restricted(&path, &payload)
}
/// Resolve the effective model and provider for an agent.
///
/// Delegates to `effective_config::resolve_effective_config` which enforces
@@ -32,16 +32,14 @@ use nostr::JsonUtil;
/// Reconcile `managed-agents.json` into kind:30177 events in the retention
/// store. Boot-time entry point, called from `event_sync::run_event_sync`
/// after the persona and team legs.
///
/// `definitions_dir` is the scoped definitions directory (`WorkspaceAgentScope::definitions_dir`).
pub(crate) fn reconcile_agents_to_events(
app: &tauri::AppHandle,
definitions_dir: &Path,
keys: &nostr::Keys,
db_path: &Path,
) {
let Ok(base_dir) = super::managed_agents_base_dir(app) else {
return;
};
match reconcile_agents_in_dir_at(&base_dir, keys, db_path) {
match reconcile_agents_in_dir_at(definitions_dir, keys, db_path) {
Ok(0) => {}
Ok(reconciled) => {
eprintln!(
@@ -46,6 +46,13 @@ pub(crate) fn managed_agents_store_path(app: &AppHandle) -> Result<PathBuf, Stri
Ok(managed_agents_base_dir(app)?.join("managed-agents.json"))
}
/// Scoped variant: resolve `managed-agents.json` under a workspace scope's
/// definitions directory. For callers that have captured a
/// `WorkspaceAgentScope` at their operation entry point.
pub(crate) fn managed_agents_store_path_at(definitions_dir: &std::path::Path) -> PathBuf {
definitions_dir.join("managed-agents.json")
}
fn managed_agents_logs_dir(app: &AppHandle) -> Result<PathBuf, String> {
let dir = managed_agents_base_dir(app)?.join("logs");
fs::create_dir_all(&dir).map_err(|error| format!("failed to create logs dir: {error}"))?;
@@ -238,12 +245,18 @@ pub(crate) fn spawn_key_refusal(record: &ManagedAgentRecord) -> Option<String> {
/// with fail-loud parse handling. Internal seam; public readers filter.
fn load_agent_store(app: &AppHandle) -> Result<Vec<ManagedAgentRecord>, String> {
let path = managed_agents_store_path(app)?;
load_agent_store_at(&path)
}
/// Path-based variant of [`load_agent_store`]. Used by scoped callers that
/// have already resolved the correct store path from a [`WorkspaceAgentScope`].
pub(crate) fn load_agent_store_at(path: &Path) -> Result<Vec<ManagedAgentRecord>, String> {
if !path.exists() {
return Ok(Vec::new());
}
let content = fs::read_to_string(&path)
.map_err(|error| format!("failed to read agent store: {error}"))?;
let content =
fs::read_to_string(path).map_err(|error| format!("failed to read agent store: {error}"))?;
serde_json::from_str(&content).map_err(|error| {
// Fail loudly and preserve the evidence: a later in-app save rewrites
// this file wholesale, which would silently destroy a malformed hand
@@ -251,7 +264,7 @@ fn load_agent_store(app: &AppHandle) -> Result<Vec<ManagedAgentRecord>, String>
// reconcile): the broken content survives as `.invalid` for the user
// to recover, and the parse error propagates instead of being
// swallowed into an empty store.
backup_invalid_store(&path);
backup_invalid_store(path);
format!("failed to parse agent store (preserved as .invalid): {error}")
})
}
@@ -266,6 +279,18 @@ pub fn load_managed_agents(app: &AppHandle) -> Result<Vec<ManagedAgentRecord>, S
Ok(records)
}
/// Scoped variant: load keyed agent instances from the given definitions dir.
/// For callers that have captured a `WorkspaceAgentScope`.
pub(crate) fn load_managed_agents_at(
definitions_dir: &Path,
) -> Result<Vec<ManagedAgentRecord>, String> {
let path = managed_agents_store_path_at(definitions_dir);
let mut records = load_agent_store_at(&path)?;
records.retain(|record| !record.pubkey.is_empty());
hydrate_keys(&mut records);
Ok(records)
}
/// Load the key-less agent *definitions* (former personas) from the unified
/// store. The persona compatibility shim (`load_personas`) presents these in
/// the legacy shape via `to_definition_view`.
@@ -275,6 +300,16 @@ pub(crate) fn load_agent_definitions(app: &AppHandle) -> Result<Vec<ManagedAgent
Ok(records)
}
/// Scoped variant: load key-less agent definitions from the given definitions dir.
pub(crate) fn load_agent_definitions_at(
definitions_dir: &Path,
) -> Result<Vec<ManagedAgentRecord>, String> {
let path = managed_agents_store_path_at(definitions_dir);
let mut records = load_agent_store_at(&path)?;
records.retain(|record| record.pubkey.is_empty());
Ok(records)
}
/// Preserve a malformed store file as `<name>.invalid` before the error path
/// unwinds. Copy, not rename: the original stays in place so repeated boots
/// keep failing loudly (rename would make the next launch look like a fresh
@@ -381,6 +416,25 @@ pub fn save_managed_agents(app: &AppHandle, records: &[ManagedAgentRecord]) -> R
write_agent_store(app, definitions, sorted)
}
/// Scoped variant: save keyed agent instances into the given definitions dir.
/// For callers that have captured a `WorkspaceAgentScope`.
pub(crate) fn save_managed_agents_at(
definitions_dir: &Path,
records: &[ManagedAgentRecord],
) -> Result<(), String> {
let definitions = load_agent_definitions_at(definitions_dir).unwrap_or_default();
let mut sorted = records.to_vec();
sorted.retain(|record| !record.pubkey.is_empty());
sorted.sort_by(|left, right| {
left.name
.to_lowercase()
.cmp(&right.name.to_lowercase())
.then_with(|| left.pubkey.cmp(&right.pubkey))
});
persist_agent_keys(&mut sorted);
write_agent_store_at(definitions_dir, definitions, sorted)
}
/// Save the key-less agent *definitions*, preserving the keyed instances —
/// the definition-side mirror of [`save_managed_agents`].
pub(crate) fn save_agent_definitions(
@@ -394,6 +448,19 @@ pub(crate) fn save_agent_definitions(
write_agent_store(app, definitions, instances)
}
/// Scoped variant: save key-less agent definitions into the given definitions dir.
pub(crate) fn save_agent_definitions_at(
definitions_dir: &Path,
definitions: &[ManagedAgentRecord],
) -> Result<(), String> {
let path = managed_agents_store_path_at(definitions_dir);
let mut instances = load_agent_store_at(&path)?;
instances.retain(|record| !record.pubkey.is_empty());
let mut definitions = definitions.to_vec();
definitions.retain(|record| record.pubkey.is_empty());
write_agent_store_at(definitions_dir, definitions, instances)
}
/// Serialize definitions + instances into the single unified store file.
/// Definitions sort first (by slug) for stable diffs; instances keep the
/// name/pubkey order their save path established.
@@ -401,12 +468,36 @@ fn write_agent_store(
app: &AppHandle,
mut definitions: Vec<ManagedAgentRecord>,
instances: Vec<ManagedAgentRecord>,
) -> Result<(), String> {
let path = managed_agents_store_path(app)?;
write_agent_store_to_path(&path, definitions, instances)
}
/// Path-based variant of [`write_agent_store`]. Used by scoped callers.
fn write_agent_store_at(
definitions_dir: &Path,
definitions: Vec<ManagedAgentRecord>,
instances: Vec<ManagedAgentRecord>,
) -> Result<(), String> {
let path = managed_agents_store_path_at(definitions_dir);
// Ensure the directory exists (scoped dirs are created lazily).
if let Some(parent) = path.parent() {
fs::create_dir_all(parent)
.map_err(|e| format!("failed to create scoped store dir: {e}"))?;
}
write_agent_store_to_path(&path, definitions, instances)
}
/// Write definitions + instances to a specific path.
fn write_agent_store_to_path(
path: &Path,
mut definitions: Vec<ManagedAgentRecord>,
instances: Vec<ManagedAgentRecord>,
) -> Result<(), String> {
definitions.sort_by(|left, right| left.slug.cmp(&right.slug));
let mut all = definitions;
all.extend(instances);
let path = managed_agents_store_path(app)?;
let payload = serde_json::to_vec_pretty(&all)
.map_err(|error| format!("failed to serialize agent store: {error}"))?;
@@ -414,7 +505,7 @@ fn write_agent_store(
// fallback. Write it owner-only (`0o600`) unconditionally — harmless for the
// keyring-backed case (it is the user's own agent store) and closes the
// umask window a post-write `chmod` would leave open.
atomic_write_json_restricted(&path, &payload)
atomic_write_json_restricted(path, &payload)
}
/// Write each record's in-memory key to the keyring and blank the inline copy
@@ -13,6 +13,11 @@ pub(crate) fn teams_store_path(app: &AppHandle) -> Result<PathBuf, String> {
Ok(managed_agents_base_dir(app)?.join("teams.json"))
}
/// Scoped variant: resolve `teams.json` under a workspace scope's definitions dir.
pub(crate) fn teams_store_path_at(definitions_dir: &std::path::Path) -> PathBuf {
definitions_dir.join("teams.json")
}
fn sort_teams(records: &mut [TeamRecord]) {
records.sort_by(|left, right| {
let left_builtin = if left.is_builtin { 0 } else { 1 };
@@ -206,6 +211,49 @@ pub fn save_teams(app: &AppHandle, records: &[TeamRecord]) -> Result<(), String>
crate::managed_agents::storage::atomic_write_json(&path, &payload)
}
/// Scoped variant: load teams from the given definitions dir.
pub(crate) fn load_teams_at(definitions_dir: &std::path::Path) -> Result<Vec<TeamRecord>, String> {
let path = teams_store_path_at(definitions_dir);
let now = now_iso();
let records = if path.exists() {
let content = fs::read_to_string(&path)
.map_err(|error| format!("failed to read teams store: {error}"))?;
serde_json::from_str::<Vec<TeamRecord>>(&content)
.map_err(|error| format!("failed to parse teams store: {error}"))?
} else {
Vec::new()
};
let (mut records, changed) = merge_teams(records, &now);
sort_teams(&mut records);
if changed || !path.exists() {
save_teams_at(definitions_dir, &records)?;
}
Ok(records)
}
/// Scoped variant: save teams into the given definitions dir.
pub(crate) fn save_teams_at(
definitions_dir: &std::path::Path,
records: &[TeamRecord],
) -> Result<(), String> {
let mut sorted = records.to_vec();
sort_teams(&mut sorted);
let path = teams_store_path_at(definitions_dir);
// Ensure the directory exists (scoped dirs are created lazily).
if let Some(parent) = path.parent() {
fs::create_dir_all(parent)
.map_err(|e| format!("failed to create scoped store dir: {e}"))?;
}
let payload = serde_json::to_vec_pretty(&sorted)
.map_err(|error| format!("failed to serialize teams store: {error}"))?;
crate::managed_agents::storage::atomic_write_json(&path, &payload)
}
/// Names of managed agents that still reference `team` — either via the
/// legacy `persona_team_dir` link (directory-backed teams only) or the
/// `team_id` field (every team kind, all agents created after the team_id