diff --git a/desktop/src-tauri/src/managed_agents/restore.rs b/desktop/src-tauri/src/managed_agents/restore.rs index 191062015..e27ee4ea7 100644 --- a/desktop/src-tauri/src/managed_agents/restore.rs +++ b/desktop/src-tauri/src/managed_agents/restore.rs @@ -1,9 +1,12 @@ use super::{ - find_managed_agent_mut, kill_stale_tracked_processes, load_managed_agents, load_personas, - save_managed_agents, spawn_agent_child, sync_managed_agent_processes, BackendKind, - ManagedAgentProcess, + find_managed_agent_mut, kill_stale_tracked_processes, spawn_agent_child, + sync_managed_agent_processes, BackendKind, ManagedAgentProcess, }; use crate::app_state::AppState; +#[cfg(feature = "mesh-llm")] +use crate::managed_agents::global_config::load_global_agent_config_at; +use crate::managed_agents::personas::load_personas_at; +use crate::managed_agents::storage::{load_managed_agents_at, save_managed_agents_at}; use crate::util; use std::sync::atomic::{AtomicBool, Ordering}; use tauri::Manager; @@ -39,12 +42,21 @@ type AgentSpawnResult = (String, SpawnOutcome); /// `effective_config::resolve_effective_config`'s `OrphanedInstance` arm). pub fn backfill_persona_snapshots(app: &tauri::AppHandle) -> Result<(), String> { let state = app.state::(); + + // Capture scope at function entry — all reads/writes in this function use + // this single captured scope so a concurrent workspace switch cannot split + // records and personas across two different definition directories. + let scope = state.capture_active_scope().ok_or_else(|| { + "backfill_persona_snapshots: no active workspace scope — skipping".to_string() + })?; + let definitions_dir = &scope.definitions_dir; + let _store_guard = state .managed_agents_store_lock .lock() .map_err(|error| error.to_string())?; - let mut records = load_managed_agents(app)?; + let mut records = load_managed_agents_at(definitions_dir)?; let needs_backfill = records .iter() .any(|r| r.persona_id.is_some() && r.persona_source_version.is_none()); @@ -52,7 +64,7 @@ pub fn backfill_persona_snapshots(app: &tauri::AppHandle) -> Result<(), String> return Ok(()); } - let personas = load_personas(app)?; + let personas = load_personas_at(definitions_dir)?; let mut changed = false; for record in records.iter_mut() { let Some(persona_id) = record.persona_id.clone() else { @@ -78,7 +90,7 @@ pub fn backfill_persona_snapshots(app: &tauri::AppHandle) -> Result<(), String> } if changed { - save_managed_agents(app, &records)?; + save_managed_agents_at(definitions_dir, &records)?; } Ok(()) } @@ -99,6 +111,14 @@ pub async fn restore_managed_agents_on_launch( let state = app.state::(); + // Capture scope at function entry — all three phases (A, B, C) use this + // single captured definitions_dir so a concurrent workspace switch cannot + // write Phase C's results into a different scope's store than Phase A read. + let scope = state + .capture_active_scope() + .ok_or_else(|| "restore_managed_agents_on_launch: no active workspace scope".to_string())?; + let definitions_dir = scope.definitions_dir.clone(); + // ── Phase A (under lock): housekeeping + collect agents to restore ── let mut agents_to_start: Vec; { @@ -111,7 +131,7 @@ pub async fn restore_managed_agents_on_launch( return Ok(()); } - let mut records = load_managed_agents(app)?; + let mut records = load_managed_agents_at(&definitions_dir)?; let mut runtimes = state .managed_agent_processes .lock() @@ -194,7 +214,7 @@ pub async fn restore_managed_agents_on_launch( // Re-snapshot persona config for agents about to be restored, matching // the interactive spawn path so auto-start agents also pick up the // current persona on app launch. - let personas_for_snapshot = super::load_personas(app).unwrap_or_default(); + let personas_for_snapshot = load_personas_at(&definitions_dir).unwrap_or_default(); for record in records.iter_mut() { if !agents_to_start.iter().any(|r| r.pubkey == record.pubkey) { continue; @@ -220,7 +240,7 @@ pub async fn restore_managed_agents_on_launch( .collect(); if changed { - save_managed_agents(app, &records)?; + save_managed_agents_at(&definitions_dir, &records)?; } } @@ -244,8 +264,10 @@ pub async fn restore_managed_agents_on_launch( // (definition → global fallback). A linked instance's own `provider`/`model`/ // `relay_mesh` bytes never contribute. See `start_local_agent_with_preflight` // in `commands/agents.rs` for the identical rationale on the interactive path. - let personas = load_personas(app).unwrap_or_default(); - let global = super::load_global_agent_config(app).unwrap_or_default(); + // Use the captured scope's definitions_dir for both loads so they read from + // the same scope as Phase A. + let personas = load_personas_at(&definitions_dir).unwrap_or_default(); + let global = load_global_agent_config_at(&definitions_dir).unwrap_or_default(); let mut mesh_preflight_failures = std::collections::HashSet::new(); for record in &agents_to_start { let mesh_model_id = super::effective_config::resolve_effective_relay_mesh_model_id( @@ -261,7 +283,7 @@ pub async fn restore_managed_agents_on_launch( crate::commands::ensure_relay_mesh_for_record(app, mesh_model_id.as_deref(), false) .await { - persist_restore_error(app, &state, &record.pubkey, error)?; + persist_restore_error(app, &state, &record.pubkey, &definitions_dir, error)?; mesh_preflight_failures.insert(record.pubkey.clone()); } } @@ -287,13 +309,13 @@ pub async fn restore_managed_agents_on_launch( } // ── Phase B (transition lock held): resolve commands and spawn in parallel ── - let spawn_results: Vec = std::thread::scope(|scope| { + let spawn_results: Vec = std::thread::scope(|scope_s| { let owner_hex_ref = owner_hex.as_deref(); let handles: Vec<_> = agents_to_start .iter() .filter(|_| !shutdown_started.load(Ordering::SeqCst)) .map(|record| { - let handle = scope.spawn(move || { + let handle = scope_s.spawn(move || { let workspace_relay = crate::relay::relay_ws_url_with_override(&app.state::()); let relay_url = crate::relay::effective_agent_relay_url( @@ -359,11 +381,14 @@ pub async fn restore_managed_agents_on_launch( } // ── Phase C (re-acquire lock): write back PIDs and status to records ── + // Use the same captured definitions_dir from function entry so Phase C + // writes to the same scope Phase A read from, even if a workspace switch + // occurred during Phase B. let _store_guard = state .managed_agents_store_lock .lock() .map_err(|error| error.to_string())?; - let mut records = load_managed_agents(app)?; + let mut records = load_managed_agents_at(&definitions_dir)?; let mut runtimes = state .managed_agent_processes .lock() @@ -417,7 +442,7 @@ pub async fn restore_managed_agents_on_launch( // releasing the lock. This mirrors the fire-and-forget pattern in // start_managed_agent — ensuring boot-restored agents get the same profile // self-healing as UI-started agents. - let reconcile_personas = super::load_personas(app).unwrap_or_default(); + let reconcile_personas = load_personas_at(&definitions_dir).unwrap_or_default(); let reconcile_items: Vec<(String, crate::commands::ProfileReconcileData)> = successfully_spawned .iter() @@ -444,7 +469,7 @@ pub async fn restore_managed_agents_on_launch( }) .collect(); - save_managed_agents(app, &records)?; + save_managed_agents_at(&definitions_dir, &records)?; drop(runtimes); drop(_store_guard); drop(restore_transition); @@ -473,15 +498,16 @@ fn persist_restore_error( app: &tauri::AppHandle, state: &AppState, pubkey: &str, + definitions_dir: &std::path::Path, error: String, ) -> Result<(), String> { let _store_guard = state .managed_agents_store_lock .lock() .map_err(|error| error.to_string())?; - let mut records = load_managed_agents(app)?; + let mut records = load_managed_agents_at(definitions_dir)?; let record = find_managed_agent_mut(&mut records, pubkey)?; record.updated_at = util::now_iso(); record.last_error = Some(error); - save_managed_agents(app, &records) + save_managed_agents_at(definitions_dir, &records) } diff --git a/desktop/src-tauri/src/managed_agents/runtime_commands.rs b/desktop/src-tauri/src/managed_agents/runtime_commands.rs index c0e55184b..58aa899a2 100644 --- a/desktop/src-tauri/src/managed_agents/runtime_commands.rs +++ b/desktop/src-tauri/src/managed_agents/runtime_commands.rs @@ -12,6 +12,9 @@ use super::{ ManagedAgentRuntimeStatus, }; use crate::app_state::AppState; +use crate::managed_agents::global_config::load_global_agent_config_at; +use crate::managed_agents::personas::load_personas_at; +use crate::managed_agents::storage::{load_managed_agents_at, save_managed_agents_at}; const STATUS_EVENT: &str = "managed-agent-runtime-status"; @@ -141,12 +144,22 @@ pub fn put_managed_agent_runtime_lifecycle( pub fn list_managed_agent_runtimes( app: AppHandle, ) -> Result, String> { + // Capture scope at function entry — all reads in this function (personas, + // global config, managed agents) must use the same scope so a concurrent + // workspace switch cannot assemble mixed-scope inputs. + let state = app.state::(); + let scope = state + .capture_active_scope() + .ok_or_else(|| "list_managed_agent_runtimes: no active workspace scope".to_string())?; + let definitions_dir = scope.definitions_dir.clone(); + // This command is polled whenever the members sidebar opens and refetched // on every status event — load the per-row status inputs once, outside // the locks, instead of hitting disk per row while holding them. - let personas = load_personas(&app).unwrap_or_default(); - let global = load_global_agent_config(&app).unwrap_or_default(); - let state = app.state::(); + // Both loads use the captured scope so they are consistent with the + // load_managed_agents below. + let personas = load_personas_at(&definitions_dir).unwrap_or_default(); + let global = load_global_agent_config_at(&definitions_dir).unwrap_or_default(); let _transition = state .managed_agent_runtime_transition .lock() @@ -155,7 +168,7 @@ pub fn list_managed_agent_runtimes( .managed_agents_store_lock .lock() .map_err(|e| e.to_string())?; - let mut records = load_managed_agents(&app)?; + let mut records = load_managed_agents_at(&definitions_dir)?; let mut runtimes = state .managed_agent_processes .lock() @@ -214,7 +227,7 @@ pub fn list_managed_agent_runtimes( // Records are only mutated above when a runtime exited — skip the store // rewrite on the common nothing-changed poll. if records_changed { - save_managed_agents(&app, &records)?; + save_managed_agents_at(&definitions_dir, &records)?; } Ok(statuses) } diff --git a/desktop/src-tauri/src/shutdown.rs b/desktop/src-tauri/src/shutdown.rs index 95f9efc3c..23605961e 100644 --- a/desktop/src-tauri/src/shutdown.rs +++ b/desktop/src-tauri/src/shutdown.rs @@ -128,50 +128,86 @@ pub(crate) fn shutdown_managed_agents(app: &tauri::AppHandle) -> Result<(), Stri .managed_agents_store_lock .lock() .map_err(|error| error.to_string())?; - let mut records = load_managed_agents(app)?; + + // When no workspace scope is active (boot before first apply_workspace, or + // after import_identity cleared the scope) we cannot load the definitions + // store — it fails closed by design. In that case, skip the record-based + // cleanup and drain only from the in-memory runtime map, which may still + // hold processes that were running before the scope was cleared. + let has_scope = state.capture_active_scope().is_some(); + let mut records = if has_scope { + load_managed_agents(app).unwrap_or_default() + } else { + Vec::new() + }; + let mut runtimes = state .managed_agent_processes .lock() .map_err(|error| error.to_string())?; - let (mut changed, _exited) = sync_managed_agent_processes( - &mut records, - &mut runtimes, - &managed_agents::current_instance_id(app), - ); - changed |= kill_stale_tracked_processes( - &mut records, - &runtimes, - &managed_agents::current_instance_id(app), - ); + + let mut changed = false; + if !records.is_empty() { + let (rec_changed, _exited) = sync_managed_agent_processes( + &mut records, + &mut runtimes, + &managed_agents::current_instance_id(app), + ); + changed |= rec_changed; + changed |= kill_stale_tracked_processes( + &mut records, + &runtimes, + &managed_agents::current_instance_id(app), + ); + } // Stop all tracked agents. Send SIGTERM to all process // groups first, then wait for exits in parallel to avoid serial 1s waits. struct AgentToStop { - idx: usize, + /// Index into `records`; `None` when the runtime has no matching record + /// (no-scope path or an orphaned runtime after a scope clear). + record_idx: Option, pid: u32, runtime: Option, } let mut to_stop: Vec = Vec::new(); - for (idx, record) in records.iter().enumerate() { - if record.backend != BackendKind::Local { - continue; - } - // Drain every tracked pair for this record, not just the first — an - // agent can run one harness per community, and each pair gets the - // graceful SIGTERM → 2s wait → SIGKILL fan-out with a stop log - // marker, instead of falling through to the orphan sweep's 200ms - // grace below. - for key in managed_agents::managed_agent_runtime_keys(&runtimes, &record.pubkey) { - let runtime = runtimes.remove(&key); - let Some(pid) = runtime - .as_ref() - .map(|rt| rt.child.id()) - .or(record.runtime_pid) - else { + if !records.is_empty() { + for (idx, record) in records.iter().enumerate() { + if record.backend != BackendKind::Local { continue; - }; - to_stop.push(AgentToStop { idx, pid, runtime }); + } + // Drain every tracked pair for this record, not just the first — an + // agent can run one harness per community, and each pair gets the + // graceful SIGTERM → 2s wait → SIGKILL fan-out with a stop log + // marker, instead of falling through to the orphan sweep's 200ms + // grace below. + for key in managed_agents::managed_agent_runtime_keys(&runtimes, &record.pubkey) { + let runtime = runtimes.remove(&key); + let Some(pid) = runtime + .as_ref() + .map(|rt| rt.child.id()) + .or(record.runtime_pid) + else { + continue; + }; + to_stop.push(AgentToStop { + record_idx: Some(idx), + pid, + runtime, + }); + } + } + } + // No-scope path: drain any runtimes that are still tracked in memory even + // though we have no record store to update. Kill every remaining entry. + if records.is_empty() { + for (_, runtime) in runtimes.drain() { + to_stop.push(AgentToStop { + record_idx: None, + pid: runtime.child.id(), + runtime: Some(runtime), + }); } } @@ -213,31 +249,35 @@ pub(crate) fn shutdown_managed_agents(app: &tauri::AppHandle) -> Result<(), Stri } } - // Reap children and update records. + // Reap children and update records where available. for mut agent in to_stop { if let Some(ref mut rt) = agent.runtime { - // Best-effort reap — don’t block shutdown if the child is stuck + // Best-effort reap — don't block shutdown if the child is stuck // in uninterruptible sleep. The zombie will be cleaned up when // our process exits and launchd reaps it. let _ = rt.child.try_wait(); - // Write log marker (best-effort). - let record = &records[agent.idx]; - let _ = managed_agents::append_log_marker( - &rt.log_path, - &format!( - "=== stopped {} ({}) at {} ===", - record.name, - record.pubkey, - util::now_iso() - ), - ); + // Write log marker (best-effort) only when we have a matching record. + if let Some(idx) = agent.record_idx { + let record = &records[idx]; + let _ = managed_agents::append_log_marker( + &rt.log_path, + &format!( + "=== stopped {} ({}) at {} ===", + record.name, + record.pubkey, + util::now_iso() + ), + ); + } + } + if let Some(idx) = agent.record_idx { + let record = &mut records[idx]; + record.runtime_pid = None; + record.last_stopped_at = Some(util::now_iso()); + record.updated_at = util::now_iso(); + record.last_exit_code = None; + record.last_error = None; } - let record = &mut records[agent.idx]; - record.runtime_pid = None; - record.last_stopped_at = Some(util::now_iso()); - record.updated_at = util::now_iso(); - record.last_exit_code = None; - record.last_error = None; } } @@ -256,7 +296,7 @@ pub(crate) fn shutdown_managed_agents(app: &tauri::AppHandle) -> Result<(), Stri // whose desktop process is no longer running and reap them. managed_agents::reap_dead_instance_agents(&managed_agents::current_instance_id(app), &[]); - if changed { + if changed && !records.is_empty() { save_managed_agents(app, &records)?; }