mirror of
https://github.com/block/buzz.git
synced 2026-08-18 06:50:31 +02:00
feat(desktop): scope-stable restore, shutdown, and runtime list
Convert restore.rs, shutdown.rs, and list_managed_agent_runtimes to use a single captured workspace scope per logical operation rather than re-resolving the active scope on every load/save call. restore.rs: - backfill_persona_snapshots captures scope at entry; uses load_managed_agents_at / save_managed_agents_at / load_personas_at throughout the single store-lock epoch. - restore_managed_agents_on_launch captures scope at function entry and clones definitions_dir; all three phases (A: collect, B: spawn, C: write-back) use the same captured path, preventing a concurrent workspace switch from writing Phase C results into the wrong scope. - persist_restore_error receives definitions_dir explicitly. - Both functions return Err (with a clear message) when no scope is active, keeping the fail-closed invariant. shutdown.rs: - When no workspace scope is active (boot before apply_workspace, or after import_identity cleared the scope) skip load_managed_agents and drain only from the in-memory runtime map. Prevents the shutdown path from panicking with 'no active workspace scope'. - record_idx: Option<usize> on AgentToStop distinguishes runtimes with a backing record from those drained without one. - save_managed_agents only called when records were actually loaded. runtime_commands.rs: - list_managed_agent_runtimes captures scope at function entry and uses load_personas_at / load_global_agent_config_at / load_managed_agents_at / save_managed_agents_at so all reads in one poll see the same scope, even if a workspace switch races between the pre-lock and in-lock loads. 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
80a68775d5
commit
6cb0a4aac7
@@ -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::<AppState>();
|
||||
|
||||
// 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::<AppState>();
|
||||
|
||||
// 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<super::ManagedAgentRecord>;
|
||||
{
|
||||
@@ -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<AgentSpawnResult> = std::thread::scope(|scope| {
|
||||
let spawn_results: Vec<AgentSpawnResult> = 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::<AppState>());
|
||||
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)
|
||||
}
|
||||
|
||||
@@ -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<Vec<ManagedAgentRuntimeStatus>, 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::<AppState>();
|
||||
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::<AppState>();
|
||||
// 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)
|
||||
}
|
||||
|
||||
@@ -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<usize>,
|
||||
pid: u32,
|
||||
runtime: Option<managed_agents::ManagedAgentPairRuntime>,
|
||||
}
|
||||
|
||||
let mut to_stop: Vec<AgentToStop> = 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)?;
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user