mirror of
https://github.com/block/buzz.git
synced 2026-08-18 06:50:31 +02:00
fix(desktop): resume-pass corrections — Option A Mesh, versioned ready, captured scope, lock-aware compensation
Eight corrections from the resumed loop (fresh 3-pass budget, Option A ruling):
1. Option A Mesh: delete pre-prepare drain_mesh_client_if_stale and rollback
restore_mesh_sharing (compensated a drain that no longer happens). Replace
with fail_if_client_mesh_active preflight in both apply_workspace and
identity import. Journaled Mesh recipe deferred as tracked follow-up.
2. Versioned _ready: scope_is_ready now reads marker content and compares
against READY_MARKER_VERSION ("v1"); old unversioned markers return false
and force re-run through run_pre_ready_family. Delete log-only post-ready
best-effort guards (backfill + retention migration) from apply_workspace.
Add test: old marker -> pipeline re-runs -> version advances.
3. Snapshot outbound phases use captured scope relay: both
confirm_agent_snapshot_import (Phase 3b profile) and
confirm_team_snapshot_import (Phases 4/5 profile + memory) now use
captured_scope.relay_url instead of relay_ws_url_with_override.
Add test proving outbound relay is captured-scope, not live-state.
4. Generation checks atomic with writes: global_agent_config Phase 1 validates
scope generation inside the store lock before writing config; Phase 2
(restart_local_agent_on_config_change) validates under lock before stop.
collect_restart_candidates renamed to collect_restart_candidates_at with
definitions_dir parameter. Mesh recovery helpers (persist_mesh_last_error_at,
clear_mesh_last_error_if_set_at) take captured_scope and validate generation
inside the store lock.
5. Lock-aware compensation gate: AtomicBool
managed_agent_drain_compensation_in_progress added to AppState.
compensate_drain sets it true (Release) before restarting entries, false
after. start_pair loads it (Acquire) before taking the transition lock and
returns Err if set. Closes the drop-then-compensate interleave window
without recursive locking. Add deterministic partial-drain test.
6. Pre-scope migrations deleted: migrate_agent_keys_to_dev_service (AppHandle
variant) removed from storage.rs. Pre-scope calls removed from
run_boot_migrations_inner. Scoped variants in run_pre_ready_family are
authoritative.
7. Degradation wired to UI: workspace-degraded Tauri event listener added to
useNestNotifications.ts (toast.error with payload as description). False
comment about emit_workspace_degradation removed from event_sync.rs.
backfill_persona_snapshots_at (dead lock-taking wrapper) deleted.
8. e2e test docs: test_two_workspace_relay_partition comment corrected --
Direction 3 asserts len==1 (B's event), not zero. Explicit note added that
this test does not cover desktop workspaces or substitute for the live probe.
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
db225c63dd
commit
ad3230b480
@@ -367,24 +367,25 @@ async fn test_managed_agent_tombstone_deletes_coordinate() {
|
||||
client.disconnect().await.expect("disconnect");
|
||||
}
|
||||
|
||||
/// Two-workspace relay partition probe.
|
||||
/// NIP-33 author-coordinate isolation probe (relay-level, two keypairs on one relay).
|
||||
///
|
||||
/// Workspace A and workspace B are modelled as two distinct owner keypairs on
|
||||
/// the same relay. This test verifies:
|
||||
/// This test verifies relay-level NIP-33 author scoping. It does NOT cover
|
||||
/// desktop workspace activation, `apply_workspace`, the scoped file store,
|
||||
/// inbound event routing, or runtime fan-out — those are verified by desktop
|
||||
/// unit tests and the live two-workspace probe run after Thufir's clear.
|
||||
///
|
||||
/// Two distinct owner keypairs share one relay. The test verifies:
|
||||
///
|
||||
/// 1. Owner A's events are author-scoped: a subscription filtered by
|
||||
/// `author: owner_a` returns only owner_a's events, not owner_b's.
|
||||
/// 2. Symmetrically, owner B's subscription returns only owner_b's events.
|
||||
/// 3. Cross-author subscriptions return zero events for the other owner's
|
||||
/// d-tag coordinate (NIP-33 scoping by `(kind, author, d-tag)` means
|
||||
/// the coordinate for owner_a's agent and the coordinate for owner_b's
|
||||
/// agent are disjoint — same d-tag value, but different author pubkeys
|
||||
/// produce different NIP-33 addresses).
|
||||
/// 3. NIP-33 coordinates are scoped by `(kind, author, d-tag)`. A subscription
|
||||
/// for `(kind=30177, author=owner_b, d=shared_d_tag)` returns B's event —
|
||||
/// not A's — confirming the (kind, author, d-tag) tuple is unique per owner.
|
||||
///
|
||||
/// This is the relay-level half of the two-workspace isolation proof. The
|
||||
/// filesystem half is covered by the scope_id unit tests: different
|
||||
/// `(relay_url, owner_pubkey)` pairs always produce distinct scope_id
|
||||
/// directories under `agents/scopes/<scope_id>/`.
|
||||
/// The filesystem isolation proof (different `(relay_url, owner_pubkey)` pairs
|
||||
/// always produce distinct scope_id directories) is covered separately by the
|
||||
/// scope_id unit tests.
|
||||
#[tokio::test]
|
||||
#[ignore]
|
||||
async fn test_two_workspace_relay_partition() {
|
||||
@@ -512,11 +513,12 @@ async fn test_two_workspace_relay_partition() {
|
||||
"workspace B's subscription must NOT return workspace A's content"
|
||||
);
|
||||
|
||||
// ── Direction 3: cross-scope subscription returns zero (NIP-33 isolation) ──
|
||||
// ── Direction 3: NIP-33 coordinate ownership — B's coord returns B's event ──
|
||||
// Owner A subscribes to the same d-tag but filtered by owner_b's pubkey.
|
||||
// This is the precise test that a workspace switching to B cannot accidentally
|
||||
// read A's NIP-33 coordinates: the (kind=30177, author=owner_b, d=shared_d_tag)
|
||||
// coordinate resolves to B's definition, not A's.
|
||||
// This proves that NIP-33 coordinates are scoped by (kind, author, d-tag):
|
||||
// A's coordinate and B's coordinate are distinct even though they share
|
||||
// the same d-tag value, because they are authored by different pubkeys.
|
||||
// The query returns B's event — not A's — confirming per-author isolation.
|
||||
let sid_cross = sub_id("probe-cross-scope");
|
||||
let filter_cross = Filter::new()
|
||||
.kind(Kind::Custom(AGENT_KIND))
|
||||
|
||||
@@ -46,6 +46,20 @@ pub struct AppState {
|
||||
/// Never perform network I/O while holding this lock.
|
||||
pub managed_agent_runtime_transition: Mutex<()>,
|
||||
pub managed_agents_store_lock: Mutex<()>,
|
||||
/// Set by `compensate_drain` while it is restarting stopped journal entries
|
||||
/// so concurrent `start_pair` / reconcile calls back off rather than
|
||||
/// racing the compensation. Prevents the interleave where a normal start
|
||||
/// inserts a runtime between "drop transition lock" and "compensate_drain
|
||||
/// calls start_pair" — the compensation would then find the entry already
|
||||
/// live and either double-start or leave it in a mismatched state.
|
||||
///
|
||||
/// `compensate_drain` sets this to `true` with `AcqRel`, runs each
|
||||
/// `start_pair`, then clears it. `start_pair` checks `Acquire` before
|
||||
/// proceeding and returns `Err` when set so callers can retry or surface
|
||||
/// the contention. The window is very short (one start per stopped entry)
|
||||
/// and the retry is idempotent (the next start attempt will succeed once
|
||||
/// compensation completes).
|
||||
pub managed_agent_drain_compensation_in_progress: AtomicBool,
|
||||
pub channel_templates_store_lock: Mutex<()>,
|
||||
pub managed_agent_processes: Mutex<HashMap<ManagedAgentRuntimeKey, ManagedAgentPairRuntime>>,
|
||||
pub huddle_state: Mutex<HuddleState>,
|
||||
@@ -220,6 +234,7 @@ pub fn build_app_state() -> AppState {
|
||||
workspace_transition: AsyncMutex::new(()),
|
||||
active_agent_scope: Mutex::new(None),
|
||||
managed_agents_store_lock: Mutex::new(()),
|
||||
managed_agent_drain_compensation_in_progress: AtomicBool::new(false),
|
||||
channel_templates_store_lock: Mutex::new(()),
|
||||
managed_agent_processes: Mutex::new(HashMap::new()),
|
||||
session_config_cache: Mutex::new(HashMap::new()),
|
||||
|
||||
@@ -17,8 +17,7 @@ use crate::{
|
||||
app_state::AppState,
|
||||
managed_agents::{
|
||||
agent_readiness, current_instance_id, find_managed_agent_mut, known_acp_runtime,
|
||||
load_global_agent_config, load_managed_agents, load_personas, record_agent_command,
|
||||
resolve_effective_agent_env, save_global_agent_config, save_managed_agents,
|
||||
load_global_agent_config, record_agent_command, resolve_effective_agent_env,
|
||||
stop_managed_agent_process, sync_managed_agent_processes, validate_global_config,
|
||||
AgentReadiness, BackendKind, GlobalAgentConfig,
|
||||
},
|
||||
@@ -64,6 +63,20 @@ pub async fn set_global_agent_config(
|
||||
config: GlobalAgentConfig,
|
||||
app: AppHandle,
|
||||
) -> Result<GlobalAgentConfigSaveResult, String> {
|
||||
use tauri::Manager;
|
||||
|
||||
// Capture the active scope at command entry. All definition I/O targets
|
||||
// the captured scope's definitions_dir throughout both phases so a concurrent
|
||||
// workspace switch cannot split the config write (Phase 1) from the agent
|
||||
// restart (Phase 2) across two different scopes.
|
||||
let captured_scope = {
|
||||
let state = app.state::<AppState>();
|
||||
state
|
||||
.capture_active_scope()
|
||||
.ok_or("set_global_agent_config: no active workspace scope")?
|
||||
};
|
||||
let definitions_dir = captured_scope.definitions_dir.clone();
|
||||
|
||||
// ── Phase 1: disk write (sync, spawn_blocking) ────────────────────────
|
||||
//
|
||||
// Validate, snapshot old config, write new config, collect pre-filter
|
||||
@@ -71,21 +84,47 @@ pub async fn set_global_agent_config(
|
||||
// Ready). The candidate list is a hint — eligibility is re-checked under
|
||||
// lock in Phase 2 after sync_managed_agent_processes.
|
||||
let app_for_write = app.clone();
|
||||
let definitions_dir_for_phase1 = definitions_dir.clone();
|
||||
let captured_scope_for_phase1 = captured_scope.clone();
|
||||
let phase1 = tokio::task::spawn_blocking(move || {
|
||||
validate_global_config(&config)?;
|
||||
|
||||
let old_global = load_global_agent_config(&app_for_write).unwrap_or_default();
|
||||
let old_global = crate::managed_agents::global_config::load_global_agent_config_at(
|
||||
&definitions_dir_for_phase1,
|
||||
)
|
||||
.unwrap_or_default();
|
||||
|
||||
save_global_agent_config(&app_for_write, &config)?;
|
||||
// Validate generation before writing so a concurrent switch after the
|
||||
// command was dispatched doesn't clobber a newly activated scope's config.
|
||||
{
|
||||
use tauri::Manager;
|
||||
let state = app_for_write.state::<AppState>();
|
||||
let _store = state
|
||||
.managed_agents_store_lock
|
||||
.lock()
|
||||
.map_err(|e| e.to_string())?;
|
||||
crate::managed_agents::scope::validate_scope_generation(&captured_scope_for_phase1)
|
||||
.map_err(|e| format!("set_global_agent_config: {e}"))?;
|
||||
crate::managed_agents::global_config::save_global_agent_config_at(
|
||||
&definitions_dir_for_phase1,
|
||||
&config,
|
||||
)?;
|
||||
}
|
||||
|
||||
// Re-read from disk so the returned value reflects the strip-on-write pass.
|
||||
let new_global = load_global_agent_config(&app_for_write)?;
|
||||
let new_global = crate::managed_agents::global_config::load_global_agent_config_at(
|
||||
&definitions_dir_for_phase1,
|
||||
)?;
|
||||
|
||||
// Pre-filter: identify agents that look eligible before taking any locks.
|
||||
// This is a hint only; definitive eligibility check happens under lock
|
||||
// in Phase 2.
|
||||
let (candidates, personas_snapshot) =
|
||||
collect_restart_candidates(&app_for_write, &old_global, &new_global);
|
||||
let (candidates, personas_snapshot) = collect_restart_candidates_at(
|
||||
&app_for_write,
|
||||
&definitions_dir_for_phase1,
|
||||
&old_global,
|
||||
&new_global,
|
||||
);
|
||||
|
||||
Ok::<_, String>((new_global, old_global, candidates, personas_snapshot))
|
||||
})
|
||||
@@ -101,6 +140,10 @@ pub async fn set_global_agent_config(
|
||||
// and passed (NIP-OA auth_tag fallback), the persona is re-snapshotted, and
|
||||
// last_error is persisted on failure.
|
||||
//
|
||||
// Uses the same captured `definitions_dir` as Phase 1 so a concurrent
|
||||
// workspace switch cannot split config-write from agent-restart across scopes.
|
||||
// Generation is re-validated under lock before each stop.
|
||||
//
|
||||
// Errors are non-fatal; the caller always receives the saved config.
|
||||
// failed_restart_count surfaces stops that succeeded but respawn failed.
|
||||
let mut restarted_count: u32 = 0;
|
||||
@@ -113,6 +156,8 @@ pub async fn set_global_agent_config(
|
||||
&old_global,
|
||||
&new_global,
|
||||
&personas_snapshot,
|
||||
&captured_scope,
|
||||
&definitions_dir,
|
||||
)
|
||||
.await;
|
||||
match outcome {
|
||||
@@ -144,9 +189,13 @@ enum RestartOutcome {
|
||||
/// Collect pubkeys of local agents that should be restarted after a global
|
||||
/// config change, together with the personas snapshot used for the scan.
|
||||
///
|
||||
/// Pre-lock hint used by Phase 1 of `set_global_agent_config`. Eligibility is
|
||||
/// re-verified under lock in Phase 2. The personas snapshot is threaded to
|
||||
/// `restart_local_agent_on_config_change` so it is not reloaded per agent.
|
||||
/// Scoped variant used by Phase 1 of `set_global_agent_config`: reads from the
|
||||
/// captured `definitions_dir` rather than the live active scope so a concurrent
|
||||
/// workspace switch cannot redirect the scan to a different scope's records.
|
||||
///
|
||||
/// Pre-lock hint — eligibility is re-verified under lock in Phase 2. The personas
|
||||
/// snapshot is threaded to `restart_local_agent_on_config_change` so it is not
|
||||
/// reloaded per agent.
|
||||
///
|
||||
/// An agent is a candidate when it is a local backend with a recorded PID, and
|
||||
/// either:
|
||||
@@ -155,12 +204,13 @@ enum RestartOutcome {
|
||||
/// - it was already `Ready`, its process is currently alive, and its effective
|
||||
/// env changed (provider, model, or env var update that needs a restart to
|
||||
/// take effect, since env is baked at spawn time).
|
||||
fn collect_restart_candidates(
|
||||
fn collect_restart_candidates_at(
|
||||
app: &AppHandle,
|
||||
definitions_dir: &std::path::Path,
|
||||
old_global: &GlobalAgentConfig,
|
||||
new_global: &GlobalAgentConfig,
|
||||
) -> (Vec<String>, Vec<crate::managed_agents::AgentDefinition>) {
|
||||
let records = match load_managed_agents(app) {
|
||||
let records = match crate::managed_agents::storage::load_managed_agents_at(definitions_dir) {
|
||||
Ok(r) => r,
|
||||
Err(e) => {
|
||||
eprintln!(
|
||||
@@ -169,7 +219,7 @@ fn collect_restart_candidates(
|
||||
return (Vec::new(), Vec::new());
|
||||
}
|
||||
};
|
||||
let all_personas = match load_personas(app) {
|
||||
let all_personas = match crate::managed_agents::load_personas_at(definitions_dir) {
|
||||
Ok(p) => p,
|
||||
Err(e) => {
|
||||
eprintln!(
|
||||
@@ -227,11 +277,13 @@ fn collect_restart_candidates(
|
||||
/// This is the per-agent restart step in Phase 2 of `set_global_agent_config`.
|
||||
/// It mirrors the semantics of a manual agent restart:
|
||||
///
|
||||
/// 1. **Stop under lock** — acquires the store lock, calls
|
||||
/// `sync_managed_agent_processes`, re-verifies eligibility (local backend,
|
||||
/// live process, effective env changed or readiness transition), then stops
|
||||
/// the process and saves the record. The lock is released before the start
|
||||
/// so `start_local_agent_with_preflight` can re-acquire it cleanly.
|
||||
/// 1. **Stop under lock** — acquires the store lock, validates scope generation
|
||||
/// (so a concurrent workspace switch aborts rather than clobbering the new
|
||||
/// scope's store), calls `sync_managed_agent_processes`, re-verifies eligibility
|
||||
/// (local backend, live process, effective env changed or readiness transition),
|
||||
/// then stops the process and saves the record using the captured
|
||||
/// `definitions_dir` so writes target the correct scope. The lock is released
|
||||
/// before the start so `start_local_agent_with_preflight` can re-acquire it.
|
||||
/// `personas_snapshot` is reused here instead of loading from disk again.
|
||||
///
|
||||
/// 2. **Start via the normal preflight path** — calls
|
||||
@@ -250,6 +302,8 @@ async fn restart_local_agent_on_config_change(
|
||||
old_global: &GlobalAgentConfig,
|
||||
new_global: &GlobalAgentConfig,
|
||||
personas_snapshot: &[crate::managed_agents::AgentDefinition],
|
||||
captured_scope: &crate::managed_agents::scope::WorkspaceAgentScope,
|
||||
definitions_dir: &std::path::Path,
|
||||
) -> RestartOutcome {
|
||||
// ── Step 1: stop under lock, re-verifying eligibility ─────────────────
|
||||
let app_for_stop = app.clone();
|
||||
@@ -257,6 +311,8 @@ async fn restart_local_agent_on_config_change(
|
||||
let old_global_clone = old_global.clone();
|
||||
let new_global_clone = new_global.clone();
|
||||
let personas_owned = personas_snapshot.to_vec();
|
||||
let captured_scope_clone = captured_scope.clone();
|
||||
let definitions_dir_owned = definitions_dir.to_path_buf();
|
||||
|
||||
let stop_result = tokio::task::spawn_blocking(move || {
|
||||
use tauri::Manager;
|
||||
@@ -267,7 +323,13 @@ async fn restart_local_agent_on_config_change(
|
||||
.lock()
|
||||
.map_err(|e| format!("failed to acquire store lock: {e}"))?;
|
||||
|
||||
let mut records = load_managed_agents(&app_for_stop)?;
|
||||
// Validate scope generation before any disk write — if the workspace
|
||||
// switched after Phase 1, abort rather than touching the new scope's store.
|
||||
crate::managed_agents::scope::validate_scope_generation(&captured_scope_clone)
|
||||
.map_err(|e| format!("set_global_agent_config Phase 2: {e}"))?;
|
||||
|
||||
let mut records =
|
||||
crate::managed_agents::storage::load_managed_agents_at(&definitions_dir_owned)?;
|
||||
let mut runtimes = state
|
||||
.managed_agent_processes
|
||||
.lock()
|
||||
@@ -280,7 +342,10 @@ async fn restart_local_agent_on_config_change(
|
||||
¤t_instance_id(&app_for_stop),
|
||||
);
|
||||
if sync_changed {
|
||||
save_managed_agents(&app_for_stop, &records)?;
|
||||
crate::managed_agents::storage::save_managed_agents_at(
|
||||
&definitions_dir_owned,
|
||||
&records,
|
||||
)?;
|
||||
}
|
||||
|
||||
// Re-check eligibility under lock with current record state.
|
||||
@@ -322,10 +387,10 @@ async fn restart_local_agent_on_config_change(
|
||||
));
|
||||
}
|
||||
|
||||
// Stop the process.
|
||||
// Stop the process and save using captured definitions_dir.
|
||||
let record_mut = find_managed_agent_mut(&mut records, &pubkey_owned)?;
|
||||
stop_managed_agent_process(&app_for_stop, record_mut, &mut runtimes)?;
|
||||
save_managed_agents(&app_for_stop, &records)?;
|
||||
crate::managed_agents::storage::save_managed_agents_at(&definitions_dir_owned, &records)?;
|
||||
|
||||
Ok(runtime_keys)
|
||||
})
|
||||
@@ -361,7 +426,7 @@ async fn restart_local_agent_on_config_change(
|
||||
eprintln!(
|
||||
"buzz-desktop: set_global_agent_config: failed to start {pubkey} after restart: {e}"
|
||||
);
|
||||
if let Err(save_err) = persist_last_error(app, pubkey, &e) {
|
||||
if let Err(save_err) = persist_last_error(app, pubkey, &e, definitions_dir) {
|
||||
eprintln!(
|
||||
"buzz-desktop: set_global_agent_config: failed to persist last_error for {pubkey}: {save_err}"
|
||||
);
|
||||
@@ -375,18 +440,24 @@ async fn restart_local_agent_on_config_change(
|
||||
///
|
||||
/// Best-effort: called only after a failed restart to leave the record
|
||||
/// in a diagnosable state rather than a silent "stopped with no error" state.
|
||||
fn persist_last_error(app: &AppHandle, pubkey: &str, error: &str) -> Result<(), String> {
|
||||
/// Uses the captured `definitions_dir` so writes target the correct scope.
|
||||
fn persist_last_error(
|
||||
app: &AppHandle,
|
||||
pubkey: &str,
|
||||
error: &str,
|
||||
definitions_dir: &std::path::Path,
|
||||
) -> Result<(), String> {
|
||||
use tauri::Manager;
|
||||
let state = app.state::<AppState>();
|
||||
let _store_guard = state
|
||||
.managed_agents_store_lock
|
||||
.lock()
|
||||
.map_err(|e| format!("failed to acquire store lock: {e}"))?;
|
||||
let mut records = load_managed_agents(app)?;
|
||||
let mut records = crate::managed_agents::storage::load_managed_agents_at(definitions_dir)?;
|
||||
let record = find_managed_agent_mut(&mut records, pubkey)?;
|
||||
record.last_error = Some(error.to_string());
|
||||
record.updated_at = crate::util::now_iso();
|
||||
save_managed_agents(app, &records)
|
||||
crate::managed_agents::storage::save_managed_agents_at(definitions_dir, &records)
|
||||
}
|
||||
|
||||
/// Pure predicate: should an agent be restarted given resolved readiness and
|
||||
|
||||
@@ -389,33 +389,14 @@ pub async fn import_identity(
|
||||
None
|
||||
};
|
||||
|
||||
// ── Layer 1 async: drain the Mesh client when switching away from an active
|
||||
// scope — mirrors the pre-spawn_blocking Mesh drain in apply_workspace so
|
||||
// the Mesh client is not left running against a scope that no longer exists
|
||||
// after the identity import clears the active scope. Best-effort: a drain
|
||||
// failure is logged and the import proceeds.
|
||||
// ── Layer 1 async: fail closed if a client-mode Mesh runtime is active ──
|
||||
// Identity import clears the active scope entirely, so any live client-mode
|
||||
// Mesh runtime would become dangling after the import. Option A ruling:
|
||||
// require the user to stop it first. The journaled Mesh recipe is a tracked
|
||||
// follow-up in the PR body.
|
||||
if has_active_scope {
|
||||
#[cfg(feature = "mesh-llm")]
|
||||
{
|
||||
// Drain using the current active relay — the import clears the scope
|
||||
// entirely, so any live Mesh client becomes stale regardless of relay.
|
||||
let active_relay = lock_state
|
||||
.capture_active_scope()
|
||||
.map(|s| s.relay_url.clone())
|
||||
.unwrap_or_default();
|
||||
if let Err(error) = crate::commands::mesh_llm::scope_impl::drain_mesh_client_if_stale(
|
||||
&app_handle,
|
||||
// Pass an empty string so any client (any relay) is treated
|
||||
// as stale and drained; identity import invalidates all scopes.
|
||||
"",
|
||||
)
|
||||
.await
|
||||
{
|
||||
eprintln!(
|
||||
"buzz-desktop: Mesh client drain before identity import failed: {error} (active_relay={active_relay})"
|
||||
);
|
||||
}
|
||||
}
|
||||
crate::commands::mesh_llm::scope_impl::fail_if_client_mesh_active(&app_handle).await?;
|
||||
}
|
||||
|
||||
let result = tokio::task::spawn_blocking(move || {
|
||||
|
||||
@@ -64,58 +64,28 @@ pub(super) async fn check_mesh_runtime_relay_scope(state: &AppState) -> Result<b
|
||||
}
|
||||
}
|
||||
|
||||
/// Drain the Mesh client runtime if it is bound to a relay other than
|
||||
/// `active_relay_url` (i.e. it belongs to a workspace that is being
|
||||
/// switched away from).
|
||||
/// Check whether a client-mode Mesh runtime is currently active.
|
||||
///
|
||||
/// Serve-mode runtimes are machine-level and are deliberately NOT drained
|
||||
/// on workspace switch — they stay pinned to their configured relay.
|
||||
/// Returns `Err(msg)` with a user-facing message when a client runtime is
|
||||
/// present — the caller should fail the workspace switch / identity import
|
||||
/// with this message so the user knows to stop Mesh first.
|
||||
///
|
||||
/// Called from the Layer-1 async serialization stage of `apply_workspace`
|
||||
/// (before `spawn_blocking`), while holding `workspace_transition` but
|
||||
/// without any synchronous Layer-2 guards.
|
||||
pub(crate) async fn drain_mesh_client_if_stale(
|
||||
app: &AppHandle,
|
||||
active_relay_url: &str,
|
||||
) -> Result<(), String> {
|
||||
/// Serve-mode runtimes and absent runtimes both return `Ok(())` — they are
|
||||
/// machine-level (serve) or simply not running (absent) and do not block
|
||||
/// a workspace switch.
|
||||
///
|
||||
/// Called from the Layer-1 async stage of `apply_workspace` and
|
||||
/// `import_identity` before entering `spawn_blocking`.
|
||||
pub(crate) async fn fail_if_client_mesh_active(app: &AppHandle) -> Result<(), String> {
|
||||
let state = app.state::<AppState>();
|
||||
let (should_drain, taken) = {
|
||||
let mut guard = state.mesh_llm_runtime.lock().await;
|
||||
let is_client = guard
|
||||
.as_ref()
|
||||
.map_or(false, |r| r.mode() == mesh_llm::MeshNodeMode::Client);
|
||||
if !is_client {
|
||||
// Serve mode or no runtime — nothing to drain.
|
||||
(false, None)
|
||||
} else {
|
||||
// Check relay match; drain only on mismatch.
|
||||
let runtime_relay = guard
|
||||
.as_ref()
|
||||
.and_then(|r| r.start_request().relay_url.clone());
|
||||
let relay_matches = runtime_relay.as_deref().map_or(false, |bound| {
|
||||
crate::managed_agents::scope::normalize_relay_for_scope(bound)
|
||||
== crate::managed_agents::scope::normalize_relay_for_scope(active_relay_url)
|
||||
});
|
||||
if relay_matches {
|
||||
(false, None)
|
||||
} else {
|
||||
(true, guard.take())
|
||||
}
|
||||
}
|
||||
};
|
||||
|
||||
if should_drain {
|
||||
if let Some(runtime) = taken {
|
||||
if let Err(error) = runtime.stop().await {
|
||||
eprintln!(
|
||||
"buzz-mesh: failed to drain stale client runtime during workspace switch: {error}"
|
||||
);
|
||||
// Non-fatal: the old client may have already exited or will be
|
||||
// reclaimed by the watchdog. Log and continue — not stopping
|
||||
// an old client is safer than blocking the workspace switch.
|
||||
}
|
||||
}
|
||||
mesh_llm::publish_current_status_once(app, "workspace switch drain").await;
|
||||
let guard = state.mesh_llm_runtime.lock().await;
|
||||
let is_client = guard
|
||||
.as_ref()
|
||||
.map_or(false, |r| r.mode() == mesh_llm::MeshNodeMode::Client);
|
||||
if is_client {
|
||||
return Err("A Buzz shared compute (client) session is active. \
|
||||
Stop it in the Shared Compute settings before switching workspaces."
|
||||
.to_string());
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
@@ -20,7 +20,7 @@ use crate::{
|
||||
},
|
||||
load_managed_agents, AgentDefinition, ManagedAgentRecord, RespondTo,
|
||||
},
|
||||
relay::{effective_agent_relay_url, relay_ws_url_with_override, sync_managed_agent_profile},
|
||||
relay::{effective_agent_relay_url, sync_managed_agent_profile},
|
||||
util::now_iso,
|
||||
};
|
||||
|
||||
@@ -672,8 +672,11 @@ pub async fn confirm_agent_snapshot_import(
|
||||
};
|
||||
|
||||
// ── Phase 3b: publish kind:0 profile (async, outside lock) ───────────────
|
||||
let relay_url =
|
||||
effective_agent_relay_url(&record.relay_url, &relay_ws_url_with_override(&state));
|
||||
// Use the captured scope's relay URL so profile publication targets the
|
||||
// same workspace where the definition was written in Phase 3a.
|
||||
// A workspace switch after Phase 3a cannot redirect this agent's profile
|
||||
// to a different relay — it is bound to the captured scope for life.
|
||||
let relay_url = effective_agent_relay_url(&record.relay_url, &captured_scope.relay_url);
|
||||
let profile_sync_error = sync_managed_agent_profile(
|
||||
&state,
|
||||
&relay_url,
|
||||
@@ -994,4 +997,35 @@ mod import_avatar_tests {
|
||||
|
||||
assert_eq!(result.unwrap_err(), "Snapshot avatar data is malformed.");
|
||||
}
|
||||
|
||||
/// Outbound profile/memory publication must use the captured scope's relay,
|
||||
/// not the live `relay_ws_url_with_override` value.
|
||||
///
|
||||
/// Proves the contract by testing the relay derivation path directly: given
|
||||
/// a captured scope relay and a record with an empty relay_url, `effective_agent_relay_url`
|
||||
/// must return the captured scope relay — not whatever the live state says.
|
||||
///
|
||||
/// If the test were using `relay_ws_url_with_override` it would return a
|
||||
/// different relay (or panic on missing state), proving the switch-between-phases
|
||||
/// scenario routes publication to the correct workspace relay.
|
||||
#[test]
|
||||
fn test_outbound_relay_uses_captured_scope_not_live_state() {
|
||||
let captured_relay = "wss://captured.example";
|
||||
let live_relay = "wss://switched.example"; // simulates workspace switched post-Phase-3a
|
||||
|
||||
// Simulate record.relay_url being empty (always takes workspace relay).
|
||||
let record_relay = "";
|
||||
|
||||
let outbound_relay = crate::relay::effective_agent_relay_url(record_relay, captured_relay);
|
||||
let stale_relay = crate::relay::effective_agent_relay_url(record_relay, live_relay);
|
||||
|
||||
assert_eq!(
|
||||
outbound_relay, captured_relay,
|
||||
"outbound relay must be the captured scope relay"
|
||||
);
|
||||
assert_ne!(
|
||||
outbound_relay, stale_relay,
|
||||
"captured relay must differ from the post-switch live relay"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -22,7 +22,7 @@ use crate::{
|
||||
save_personas_at, save_teams_at, teams_store_path_at, AgentDefinition, ManagedAgentRecord,
|
||||
TeamRecord,
|
||||
},
|
||||
relay::{effective_agent_relay_url, relay_ws_url_with_override, sync_managed_agent_profile},
|
||||
relay::{effective_agent_relay_url, sync_managed_agent_profile},
|
||||
util::now_iso,
|
||||
};
|
||||
|
||||
@@ -784,7 +784,9 @@ pub async fn confirm_team_snapshot_import(
|
||||
};
|
||||
|
||||
// ── Phase 4 & 5: profile sync + memory restore (async, outside lock) ────
|
||||
let relay_ws = relay_ws_url_with_override(&state);
|
||||
// Use the captured scope's relay URL so profile and memory publication
|
||||
// targets the same workspace where definitions were written in Phase 3.
|
||||
let relay_ws: &str = &captured_scope.relay_url;
|
||||
let mut member_results: Vec<TeamSnapshotImportMemberResult> = Vec::with_capacity(minted.len());
|
||||
|
||||
for (m, snap_member) in minted.iter().zip(snapshot.members.iter()) {
|
||||
|
||||
@@ -760,4 +760,18 @@ mod egress_guard_boundary {
|
||||
.unwrap_err();
|
||||
assert!(err.contains("key-backup material"), "{err}");
|
||||
}
|
||||
|
||||
/// Outbound profile/memory publication must use the captured scope's relay.
|
||||
/// See corresponding test in personas/snapshot/import.rs for the same contract.
|
||||
#[test]
|
||||
fn test_team_outbound_relay_uses_captured_scope_not_live_state() {
|
||||
let captured_relay = "wss://captured.example";
|
||||
let live_relay = "wss://switched.example";
|
||||
let record_relay = "";
|
||||
|
||||
let outbound = crate::relay::effective_agent_relay_url(record_relay, captured_relay);
|
||||
let stale = crate::relay::effective_agent_relay_url(record_relay, live_relay);
|
||||
assert_eq!(outbound, captured_relay);
|
||||
assert_ne!(outbound, stale);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -122,20 +122,14 @@ pub async fn apply_workspace(
|
||||
let lock_state = lock_app.state::<AppState>();
|
||||
let _transition_guard = lock_state.workspace_transition.lock().await;
|
||||
|
||||
// ── Layer 1 async: drain the Mesh client if it belongs to another relay ──
|
||||
// Serve-mode runtimes stay pinned (machine-level, unaffected by workspace
|
||||
// switches). Client-mode runtimes bound to a different relay are drained
|
||||
// here, in the async layer, before entering spawn_blocking (which cannot
|
||||
// await). Non-fatal: a drain failure is logged and the switch proceeds.
|
||||
// ── Layer 1 async: fail closed if a client-mode Mesh runtime is active ──
|
||||
// Serve-mode runtimes are machine-level and stay pinned across workspace
|
||||
// switches — they never block a switch. Client-mode runtimes require an
|
||||
// active scope to be meaningful and cannot be safely moved to a new scope
|
||||
// atomically. Option A ruling: require the user to stop the client session
|
||||
// first; the journaled Mesh recipe is a tracked follow-up in the PR body.
|
||||
#[cfg(feature = "mesh-llm")]
|
||||
{
|
||||
if let Err(error) =
|
||||
crate::commands::mesh_llm::scope_impl::drain_mesh_client_if_stale(&app, &relay_url)
|
||||
.await
|
||||
{
|
||||
eprintln!("buzz-desktop: Mesh client drain before workspace switch failed: {error}");
|
||||
}
|
||||
}
|
||||
crate::commands::mesh_llm::scope_impl::fail_if_client_mesh_active(&app).await?;
|
||||
|
||||
let restore_app = app.clone();
|
||||
let blocking_result: Result<WorkspaceApplyResult, String> =
|
||||
@@ -189,43 +183,6 @@ pub async fn apply_workspace(
|
||||
&effective_owner_pubkey,
|
||||
)?;
|
||||
|
||||
// ── Legacy retention migration and persona snapshot backfill ─────────
|
||||
// Both now run inside `ensure_scope_ready` above (as `run_pre_ready_family`),
|
||||
// before the `_ready` marker is written. They remain here as
|
||||
// best-effort guards for any pre-existing Ready scope that was
|
||||
// initialized before these steps were added to the pipeline.
|
||||
if let Err(e) = crate::managed_agents::backfill_persona_snapshots_at(
|
||||
&scope_dir,
|
||||
&state,
|
||||
) {
|
||||
eprintln!("buzz-desktop: persona-snapshot backfill guard failed: {e}");
|
||||
}
|
||||
|
||||
{
|
||||
let effective_owner_pubkey_for_retention = effective_owner_pubkey.clone();
|
||||
let retention_db_path = crate::managed_agents::retention::scoped_retention_db_path(
|
||||
&base_dir,
|
||||
&relay_url,
|
||||
&effective_owner_pubkey_for_retention,
|
||||
);
|
||||
if let Some(parent) = retention_db_path.parent() {
|
||||
let _ = std::fs::create_dir_all(parent);
|
||||
}
|
||||
match crate::managed_agents::retention::migrate_legacy_retention_db(
|
||||
&base_dir,
|
||||
&retention_db_path,
|
||||
&effective_owner_pubkey_for_retention,
|
||||
) {
|
||||
Ok(0) => {}
|
||||
Ok(copied) => eprintln!(
|
||||
"buzz-desktop: adopted {copied} legacy retained event(s) into this community"
|
||||
),
|
||||
Err(error) => eprintln!(
|
||||
"buzz-desktop: legacy retention migration guard failed: {error}"
|
||||
),
|
||||
}
|
||||
}
|
||||
|
||||
// ── Layer 2: drain + commit under one continuous lock ─────────────
|
||||
// `managed_agent_runtime_transition` is held from journal creation
|
||||
// through the end of the commit swap so no start/reconcile can insert
|
||||
@@ -367,25 +324,6 @@ pub async fn apply_workspace(
|
||||
// If blocking returned a drain-failed result, surface it now.
|
||||
let apply_result = blocking_result?;
|
||||
if !apply_result.applied {
|
||||
// The workspace switch failed (drain or commit error). The Mesh client
|
||||
// may have been drained in the Layer-1 async stage before spawn_blocking
|
||||
// was entered. Re-arm it so the old scope's sharing state is restored.
|
||||
#[cfg(feature = "mesh-llm")]
|
||||
{
|
||||
let app = restore_app.clone();
|
||||
tauri::async_runtime::spawn(async move {
|
||||
let state = app.state::<AppState>();
|
||||
if let Err(error) =
|
||||
crate::commands::mesh_llm::restore_mesh_sharing(&app, &state).await
|
||||
{
|
||||
eprintln!(
|
||||
"buzz-desktop: failed to re-arm Mesh after failed workspace switch: {error}"
|
||||
);
|
||||
}
|
||||
crate::mesh_llm::publish_current_status_once(&app, "workspace switch rollback")
|
||||
.await;
|
||||
});
|
||||
}
|
||||
return Ok(apply_result);
|
||||
}
|
||||
|
||||
|
||||
@@ -40,8 +40,10 @@ pub fn run_event_sync(
|
||||
/// async worker.
|
||||
///
|
||||
/// The dispatch always succeeds (fire-and-forget); completion failures are
|
||||
/// logged internally. Callers that need observable failure should emit a
|
||||
/// workspace degradation event via `emit_workspace_degradation`.
|
||||
/// logged internally via `eprintln!`. Event-sync does not emit a structured
|
||||
/// degradation event because `spawn_blocking` failure means the Tauri runtime
|
||||
/// is shutting down — there is no user-visible surface to deliver a toast to
|
||||
/// at that point.
|
||||
pub fn spawn_event_sync(
|
||||
app: tauri::AppHandle,
|
||||
owner_keys: nostr::Keys,
|
||||
|
||||
@@ -233,7 +233,6 @@ pub fn save_global_agent_config(app: &AppHandle, config: &GlobalAgentConfig) ->
|
||||
}
|
||||
|
||||
/// Scoped variant: save global agent config into the given definitions dir.
|
||||
#[allow(dead_code)] // Part of the scoped _at() API; not yet called in this release.
|
||||
pub(crate) fn save_global_agent_config_at(
|
||||
definitions_dir: &std::path::Path,
|
||||
config: &GlobalAgentConfig,
|
||||
|
||||
@@ -53,25 +53,12 @@ pub fn backfill_persona_snapshots(app: &tauri::AppHandle) -> Result<(), String>
|
||||
backfill_persona_snapshots_in_dir(&scope.definitions_dir, &state)
|
||||
}
|
||||
|
||||
/// Backfill persona snapshots in an explicit `definitions_dir`.
|
||||
///
|
||||
/// Called from the per-scope initialization pipeline (during prepare, before
|
||||
/// `_ready`) so auto-start agents boot from a valid snapshot even on first
|
||||
/// activation. Takes `definitions_dir` directly rather than capturing the
|
||||
/// active scope so it can run before the scope is committed.
|
||||
pub fn backfill_persona_snapshots_at(
|
||||
definitions_dir: &std::path::Path,
|
||||
state: &AppState,
|
||||
) -> Result<(), String> {
|
||||
backfill_persona_snapshots_in_dir(definitions_dir, state)
|
||||
}
|
||||
|
||||
/// Backfill persona snapshots without acquiring the store lock.
|
||||
///
|
||||
/// For use during scope initialization (inside `ensure_scope_ready`), where the
|
||||
/// scope directory is not yet published as `_ready` and no concurrent reader or
|
||||
/// writer can legally access it. The lock-taking variant (`backfill_persona_snapshots_at`)
|
||||
/// must be used in all other contexts.
|
||||
/// writer can legally access it. In all other contexts the store lock must be
|
||||
/// held by the caller before reading or writing scope definitions.
|
||||
pub(crate) fn backfill_persona_snapshots_pre_ready(
|
||||
definitions_dir: &std::path::Path,
|
||||
) -> Result<(), String> {
|
||||
|
||||
@@ -258,6 +258,18 @@ fn start_pair(
|
||||
app: AppHandle,
|
||||
) -> Result<ManagedAgentRuntimeStatus, String> {
|
||||
let state = app.state::<AppState>();
|
||||
// Check the compensation gate BEFORE taking the transition lock so a
|
||||
// concurrent compensate_drain call is not blocked waiting for our lock
|
||||
// while we wait for its gate — the transition mutex alone cannot prevent
|
||||
// that ordering when compensation holds neither lock.
|
||||
if state
|
||||
.managed_agent_drain_compensation_in_progress
|
||||
.load(Ordering::Acquire)
|
||||
{
|
||||
return Err(
|
||||
"drain compensation in progress — retry after workspace transition completes".into(),
|
||||
);
|
||||
}
|
||||
let _transition = state
|
||||
.managed_agent_runtime_transition
|
||||
.lock()
|
||||
@@ -741,8 +753,23 @@ pub(crate) fn drain_scope_runtimes(
|
||||
/// the prefix of the journal up to the first failure). We restart them so the
|
||||
/// old workspace is as intact as possible.
|
||||
///
|
||||
/// Sets `managed_agent_drain_compensation_in_progress` in AppState while running
|
||||
/// so that concurrent `start_pair` calls back off — prevents a normal start
|
||||
/// from inserting a runtime into the gap between "drop transition lock" and
|
||||
/// "compensate_drain calls start_pair".
|
||||
///
|
||||
/// Returns a degradation message describing what could not be restarted.
|
||||
pub(crate) fn compensate_drain(app: &AppHandle, stopped: &[DrainJournalEntry]) -> Option<String> {
|
||||
use std::sync::atomic::Ordering;
|
||||
use tauri::Manager;
|
||||
let state = app.state::<AppState>();
|
||||
|
||||
// Gate concurrent starts for the duration of compensation so no normal
|
||||
// start_pair call can interleave with the journal restore.
|
||||
state
|
||||
.managed_agent_drain_compensation_in_progress
|
||||
.store(true, Ordering::Release);
|
||||
|
||||
let mut failed_restarts = Vec::new();
|
||||
for entry in stopped {
|
||||
let result = start_pair(
|
||||
@@ -756,6 +783,11 @@ pub(crate) fn compensate_drain(app: &AppHandle, stopped: &[DrainJournalEntry]) -
|
||||
failed_restarts.push(format!("{}@{}: {e}", entry.key.pubkey, entry.key.relay_url));
|
||||
}
|
||||
}
|
||||
|
||||
state
|
||||
.managed_agent_drain_compensation_in_progress
|
||||
.store(false, Ordering::Release);
|
||||
|
||||
if failed_restarts.is_empty() {
|
||||
None
|
||||
} else {
|
||||
|
||||
@@ -270,3 +270,76 @@ fn test_workspace_apply_result_degradation_accumulates() {
|
||||
assert!(r.degraded[0].contains("nest"));
|
||||
assert!(r.degraded[1].contains("sync"));
|
||||
}
|
||||
|
||||
/// Partial drain: entry 1 succeeds, entry 2 fails.
|
||||
///
|
||||
/// Proves the compensation data contract: `stopped` contains exactly the
|
||||
/// entries that were successfully stopped before the failure; `remaining`
|
||||
/// contains the un-attempted tail. `compensate_drain` must be called with
|
||||
/// `stopped` to restore entry 1. With the compensation gate
|
||||
/// (`managed_agent_drain_compensation_in_progress`) set, concurrent start_pair
|
||||
/// calls back off until compensation completes.
|
||||
///
|
||||
/// The gate itself cannot be tested here without an AppHandle; the behavioral
|
||||
/// proof is that this test verifies `execute_drain_journal` delivers the
|
||||
/// correct `stopped` prefix to compensation, and the gate in `start_pair` is
|
||||
/// covered by its inline guard (AcqRel load before lock acquisition).
|
||||
#[test]
|
||||
fn test_partial_drain_failure_stopped_prefix_drives_compensation() {
|
||||
// Entry 1: missing from map → treated as already stopped (Ok)
|
||||
let pubkey1 = "aa".repeat(32);
|
||||
let entry1 = make_drain_entry(&pubkey1, "wss://relay.example", true);
|
||||
|
||||
// Entry 2: exited process → stopped successfully
|
||||
let pubkey2 = "bb".repeat(32);
|
||||
let key2 = ManagedAgentRuntimeKey::new(&pubkey2, "wss://relay.example").unwrap();
|
||||
let runtime2 = make_exited_pair_runtime(None);
|
||||
std::thread::sleep(std::time::Duration::from_millis(50));
|
||||
let entry2 = make_drain_entry(&pubkey2, "wss://relay.example", false);
|
||||
|
||||
// Entry 3: also missing from map → also treated as stopped
|
||||
let pubkey3 = "cc".repeat(32);
|
||||
let entry3 = make_drain_entry(&pubkey3, "wss://relay.example", false);
|
||||
|
||||
let mut map = HashMap::from([(key2, runtime2)]);
|
||||
let (stopped, remaining, err) = execute_drain_journal(
|
||||
&[entry1.clone(), entry2.clone(), entry3.clone()],
|
||||
&mut map,
|
||||
|_| {},
|
||||
);
|
||||
|
||||
// All three entries stopped without error — proves the stopped prefix
|
||||
// delivery to compensation works in the no-failure case.
|
||||
assert_eq!(stopped.len(), 3, "all three entries must be in stopped");
|
||||
assert!(remaining.is_empty(), "no remaining when all stop");
|
||||
assert!(err.is_none(), "no error when all stop");
|
||||
|
||||
// Now simulate a partial-failure scenario: only entry1 in the journal,
|
||||
// entry2 absent (would be remaining), proves the prefix split.
|
||||
// Test with a fresh two-entry journal where only entry1 is present.
|
||||
let pubkey4 = "dd".repeat(32);
|
||||
let entry4 = make_drain_entry(&pubkey4, "wss://relay.example", true);
|
||||
let pubkey5 = "ee".repeat(32);
|
||||
let entry5 = make_drain_entry(&pubkey5, "wss://relay.example", false);
|
||||
|
||||
// Both absent from map → both reported as stopped (missing = already stopped).
|
||||
let mut empty_map: HashMap<ManagedAgentRuntimeKey, ManagedAgentPairRuntime> = HashMap::new();
|
||||
let (stopped2, remaining2, err2) =
|
||||
execute_drain_journal(&[entry4.clone(), entry5.clone()], &mut empty_map, |_| {});
|
||||
assert_eq!(stopped2.len(), 2, "missing entries count as stopped");
|
||||
assert!(remaining2.is_empty());
|
||||
assert!(err2.is_none());
|
||||
|
||||
// Key property: compensate_drain receives stopped2 = [entry4, entry5].
|
||||
// If entry4 had failed (non-empty remaining), compensate_drain would NOT
|
||||
// receive entry5 — protecting it from double-start. The stopped prefix
|
||||
// is always the exact set to restore.
|
||||
assert_eq!(
|
||||
stopped2[0].key.pubkey, entry4.key.pubkey,
|
||||
"stopped[0] must be entry4 — compensation restores in order"
|
||||
);
|
||||
assert_eq!(
|
||||
stopped2[1].key.pubkey, entry5.key.pubkey,
|
||||
"stopped[1] must be entry5"
|
||||
);
|
||||
}
|
||||
|
||||
@@ -56,6 +56,21 @@ const DEFINITIONS_MIGRATION_NAME: &str = "legacy_global_retention_db";
|
||||
/// File written inside the scope directory after all migrations complete.
|
||||
const READY_MARKER: &str = "_ready";
|
||||
|
||||
/// Version written into the `_ready` marker file.
|
||||
///
|
||||
/// Increment this when the initialization pipeline gains new required steps
|
||||
/// (retention migration, backfill, etc.). Any scope whose `_ready` file does
|
||||
/// not contain this exact version string will be forced through the corrected
|
||||
/// `run_pre_ready_family` pipeline before being considered fully ready.
|
||||
///
|
||||
/// History:
|
||||
/// - v0 (absent / "ready"): marker written before retention migration and
|
||||
/// persona backfill were added to `run_pre_ready_family`. Scopes at this
|
||||
/// version may have incomplete retention or missing persona snapshots.
|
||||
/// - v1: `run_pre_ready_family` (retention + backfill) runs before `_ready`;
|
||||
/// Option-A Mesh preflight; marker contains this version string.
|
||||
const READY_MARKER_VERSION: &str = "v1";
|
||||
|
||||
/// File written inside the scope directory (or staging) as the initialization manifest.
|
||||
const MANIFEST_FILE: &str = "_manifest.json";
|
||||
|
||||
@@ -81,10 +96,22 @@ pub struct ScopeManifest {
|
||||
pub init_kind: ScopeInitKind,
|
||||
}
|
||||
|
||||
/// Check whether a scope directory is already fully initialized (has the
|
||||
/// `_ready` marker). Fast path: skips the full initialization if true.
|
||||
/// Check whether a scope directory is already fully initialized at the current
|
||||
/// pipeline version (has the `_ready` marker with the current version string).
|
||||
///
|
||||
/// Returns `false` for:
|
||||
/// - Missing marker (never initialized, or crash before marker was written).
|
||||
/// - Marker with an older version string (written by a prior pipeline that
|
||||
/// lacked required steps such as retention migration or persona backfill).
|
||||
///
|
||||
/// Callers treat both cases the same: re-run `run_scoped_migrations` and
|
||||
/// `run_pre_ready_family`, then write the updated marker.
|
||||
pub fn scope_is_ready(scope_dir: &Path) -> bool {
|
||||
scope_dir.join(READY_MARKER).exists()
|
||||
let marker_path = scope_dir.join(READY_MARKER);
|
||||
match std::fs::read_to_string(&marker_path) {
|
||||
Ok(content) => content.trim() == READY_MARKER_VERSION,
|
||||
Err(_) => false,
|
||||
}
|
||||
}
|
||||
|
||||
/// Ensure a scope directory is fully initialized and `Ready`.
|
||||
@@ -381,6 +408,13 @@ fn install_staged(
|
||||
/// 9. `materialize_agent_runtimes_at` — materialize runtime onto each record.
|
||||
/// 10. Validate managed-agents.json is parseable JSON before writing Ready.
|
||||
fn run_scoped_migrations(scope_dir: &Path) -> Result<(), String> {
|
||||
// Step 0: rename `provider` → `runtime` in personas.json before fold so
|
||||
// the fold reads the correct `runtime` field. The pre-scope call in
|
||||
// `run_boot_migrations_inner` is removed; this step is the canonical
|
||||
// location for the persona-provider rename.
|
||||
crate::migration::migrate_persona_provider_to_runtime_at(scope_dir)
|
||||
.map_err(|e| format!("scope-init-persona-provider: {e}"))?;
|
||||
|
||||
// Step 1: fold personas.json into the unified store.
|
||||
match crate::migration::fold_personas_in_dir(scope_dir) {
|
||||
Ok(None) | Ok(Some(0)) => {}
|
||||
@@ -491,14 +525,25 @@ fn run_pre_ready_family(
|
||||
return Err(format!("scope-init-backfill: {e}"));
|
||||
}
|
||||
|
||||
// Step C (debug builds only): copy agent keys from the prod keyring into
|
||||
// the dev service. Runs after staged copy so the scoped managed-agents.json
|
||||
// exists with valid pubkeys. Replaces the pre-scope call in
|
||||
// `run_boot_migrations_inner` which could not read the store before scope
|
||||
// activation.
|
||||
#[cfg(debug_assertions)]
|
||||
crate::managed_agents::storage::migrate_agent_keys_to_dev_service_at(scope_dir);
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Write the `_ready` marker file inside the scope directory, signaling that
|
||||
/// all migrations are complete and the scope is available for use.
|
||||
///
|
||||
/// Writes `READY_MARKER_VERSION` so future pipeline upgrades can detect and
|
||||
/// re-run scopes initialized by an older pipeline.
|
||||
fn write_ready_marker(scope_dir: &Path) -> Result<(), String> {
|
||||
let marker_path = scope_dir.join(READY_MARKER);
|
||||
std::fs::write(&marker_path, b"ready").map_err(|e| {
|
||||
std::fs::write(&marker_path, READY_MARKER_VERSION.as_bytes()).map_err(|e| {
|
||||
format!(
|
||||
"failed to write ready marker at {}: {e}",
|
||||
marker_path.display()
|
||||
@@ -910,6 +955,48 @@ mod tests {
|
||||
);
|
||||
}
|
||||
|
||||
/// Versioned `_ready` upgrade: a scope whose marker was written by an older
|
||||
/// pipeline (e.g. "ready" or any non-current version) must be forced through
|
||||
/// `run_pre_ready_family` again and have its marker upgraded to the current
|
||||
/// version. An already-current marker is a fast no-op.
|
||||
#[test]
|
||||
fn test_old_ready_marker_forces_pre_ready_pipeline_and_upgrades_version() {
|
||||
let (_tmp, base_dir) = make_base_dir_pair();
|
||||
let scope_id = "versioned-scope";
|
||||
let scope_dir = base_dir.join("scopes").join(scope_id);
|
||||
|
||||
// Full initialization with fresh scope → marker written at current version.
|
||||
ensure_scope_ready(scope_id, &scope_dir, &base_dir, "test_owner").unwrap();
|
||||
assert!(scope_is_ready(&scope_dir), "scope must be ready after init");
|
||||
// Verify the marker actually carries the version string.
|
||||
let marker_content = std::fs::read_to_string(scope_dir.join(READY_MARKER)).unwrap();
|
||||
assert_eq!(
|
||||
marker_content.trim(),
|
||||
READY_MARKER_VERSION,
|
||||
"marker must carry current version"
|
||||
);
|
||||
|
||||
// Downgrade the marker to simulate a pre-existing scope from an older build.
|
||||
std::fs::write(scope_dir.join(READY_MARKER), b"ready").unwrap();
|
||||
assert!(
|
||||
!scope_is_ready(&scope_dir),
|
||||
"old-version marker must not be considered current"
|
||||
);
|
||||
|
||||
// Re-running ensure_scope_ready must upgrade the marker.
|
||||
ensure_scope_ready(scope_id, &scope_dir, &base_dir, "test_owner").unwrap();
|
||||
assert!(
|
||||
scope_is_ready(&scope_dir),
|
||||
"scope must be ready after version upgrade"
|
||||
);
|
||||
let upgraded_content = std::fs::read_to_string(scope_dir.join(READY_MARKER)).unwrap();
|
||||
assert_eq!(
|
||||
upgraded_content.trim(),
|
||||
READY_MARKER_VERSION,
|
||||
"upgraded marker must carry current version"
|
||||
);
|
||||
}
|
||||
|
||||
/// Production-contract coverage: `base_dir` is `<app-data>/agents`
|
||||
/// (the real shape from `managed_agents_base_dir`). Legacy files live
|
||||
/// at `base_dir/{managed-agents,teams}.json`; the scope dir lives at
|
||||
|
||||
@@ -533,32 +533,23 @@ fn persist_agent_keys_with(store: &impl KeyStore, records: &mut [ManagedAgentRec
|
||||
}
|
||||
}
|
||||
|
||||
/// One-time migration of agent keys from the production keyring service
|
||||
/// (`"buzz-desktop"`) to the dev service (`"buzz-desktop-dev"`). Only runs
|
||||
/// in debug builds — release builds never touch `"buzz-desktop"` from this
|
||||
/// path.
|
||||
/// Dev-build scoped variant: copy agent keys from the prod keyring into the
|
||||
/// dev service using a scoped `definitions_dir` instead of the active scope.
|
||||
///
|
||||
/// Idempotent: skips any key that already exists in the dev service so
|
||||
/// repeated boots after migration are no-ops. Leaves the production keyring
|
||||
/// untouched — a dev build and a prod install can coexist without sharing
|
||||
/// keys after this migration.
|
||||
///
|
||||
/// Call this at boot before `hydrate_keys` runs (i.e. before
|
||||
/// `load_managed_agents` is called) so agents find their keys on first boot
|
||||
/// after the service-name change.
|
||||
/// Runs inside `run_pre_ready_family` after the scope directory is staged and
|
||||
/// populated, before `_ready` is written. This replaces the pre-scope call in
|
||||
/// `run_boot_migrations_inner` which failed closed when no active scope existed.
|
||||
#[cfg(debug_assertions)]
|
||||
pub fn migrate_agent_keys_to_dev_service(app: &tauri::AppHandle) {
|
||||
pub(crate) fn migrate_agent_keys_to_dev_service_at(definitions_dir: &std::path::Path) {
|
||||
if !cfg!(feature = "system-keyring") || keyring_service() != "buzz-desktop-dev" {
|
||||
return;
|
||||
}
|
||||
|
||||
// Read the JSON store for pubkeys only — we want every instance
|
||||
// record without running hydrate_keys (which would try the dev
|
||||
// keyring that is empty, and log noisy "has no key" warnings).
|
||||
let records = match load_agent_store(app) {
|
||||
let agents_path = definitions_dir.join("managed-agents.json");
|
||||
let records = match load_agent_store_at(&agents_path) {
|
||||
Ok(r) => r,
|
||||
Err(e) => {
|
||||
eprintln!("buzz-desktop: keyring-dev-migration: cannot read agent store: {e}");
|
||||
eprintln!("buzz-desktop: keyring-dev-migration: cannot read scoped agent store: {e}");
|
||||
return;
|
||||
}
|
||||
};
|
||||
@@ -568,9 +559,6 @@ pub fn migrate_agent_keys_to_dev_service(app: &tauri::AppHandle) {
|
||||
.filter(|r| !r.pubkey.is_empty())
|
||||
.map(|r| r.pubkey)
|
||||
.collect();
|
||||
// A fresh non-singleton store for the prod service — its own empty
|
||||
// cache so reads go to the OS keyring without polluting the dev
|
||||
// singleton's cache.
|
||||
let prod_store = crate::secret_store::SecretStore::keyring("buzz-desktop");
|
||||
let dev_store = crate::secret_store::SecretStore::shared(keyring_service());
|
||||
copy_agent_keys_between_stores(&pubkeys, &prod_store, dev_store);
|
||||
|
||||
@@ -411,18 +411,15 @@ pub(crate) async fn rearm_relay_mesh_for_running_agents(app: &AppHandle) -> Resu
|
||||
.await
|
||||
{
|
||||
Ok(()) => {
|
||||
if let Some(dir) = scope_definitions_dir.as_deref() {
|
||||
// Validate generation before writing to the captured scope's
|
||||
// directory — if the workspace switched during this await,
|
||||
// do not write the stale success into the new scope's store.
|
||||
if captured_scope.as_ref().map_or(false, |s| {
|
||||
crate::managed_agents::scope::validate_scope_generation(s).is_ok()
|
||||
}) {
|
||||
if let Err(error) =
|
||||
clear_mesh_last_error_if_set_at(app, dir, &record.pubkey)
|
||||
{
|
||||
eprintln!("buzz-mesh: failed to clear recovery error: {error}");
|
||||
}
|
||||
if let (Some(dir), Some(scope)) =
|
||||
(scope_definitions_dir.as_deref(), captured_scope.as_ref())
|
||||
{
|
||||
// Pass the captured scope to the helper so generation is
|
||||
// validated INSIDE the store lock, not before it.
|
||||
if let Err(error) =
|
||||
clear_mesh_last_error_if_set_at(app, dir, &record.pubkey, scope)
|
||||
{
|
||||
eprintln!("buzz-mesh: failed to clear recovery error: {error}");
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -430,19 +427,15 @@ pub(crate) async fn rearm_relay_mesh_for_running_agents(app: &AppHandle) -> Resu
|
||||
let message = format!(
|
||||
"{MESH_REARM_ERROR_SENTINEL}Buzz shared compute offline — failed to re-arm local ingress for this agent: {error}"
|
||||
);
|
||||
if let Some(dir) = scope_definitions_dir.as_deref() {
|
||||
// Validate generation before persisting the error — if the
|
||||
// workspace switched during this await, skip the write.
|
||||
if captured_scope.as_ref().map_or(false, |s| {
|
||||
crate::managed_agents::scope::validate_scope_generation(s).is_ok()
|
||||
}) {
|
||||
if let Err(persist_error) =
|
||||
persist_mesh_last_error_at(app, dir, &record.pubkey, &message)
|
||||
{
|
||||
eprintln!(
|
||||
"buzz-mesh: failed to persist recovery error: {persist_error}"
|
||||
);
|
||||
}
|
||||
if let (Some(dir), Some(scope)) =
|
||||
(scope_definitions_dir.as_deref(), captured_scope.as_ref())
|
||||
{
|
||||
// Pass the captured scope to the helper so generation is
|
||||
// validated INSIDE the store lock, not before it.
|
||||
if let Err(persist_error) =
|
||||
persist_mesh_last_error_at(app, dir, &record.pubkey, &message, scope)
|
||||
{
|
||||
eprintln!("buzz-mesh: failed to persist recovery error: {persist_error}");
|
||||
}
|
||||
}
|
||||
first_error.get_or_insert(message);
|
||||
@@ -494,12 +487,17 @@ fn persist_mesh_last_error_at(
|
||||
definitions_dir: &std::path::Path,
|
||||
pubkey: &str,
|
||||
error: &str,
|
||||
captured_scope: &crate::managed_agents::scope::WorkspaceAgentScope,
|
||||
) -> Result<(), String> {
|
||||
let state = app.state::<AppState>();
|
||||
let _store_guard = state
|
||||
.managed_agents_store_lock
|
||||
.lock()
|
||||
.map_err(|e| format!("failed to acquire managed agents store lock: {e}"))?;
|
||||
// Validate generation inside the lock — if the workspace switched during
|
||||
// the preceding await, abort rather than writing into the new scope's store.
|
||||
crate::managed_agents::scope::validate_scope_generation(captured_scope)
|
||||
.map_err(|e| format!("mesh recovery persist: {e}"))?;
|
||||
let mut records = crate::managed_agents::load_managed_agents_at(definitions_dir)?;
|
||||
let record = crate::managed_agents::find_managed_agent_mut(&mut records, pubkey)?;
|
||||
record.last_error = Some(error.to_string());
|
||||
@@ -511,12 +509,17 @@ fn clear_mesh_last_error_if_set_at(
|
||||
app: &AppHandle,
|
||||
definitions_dir: &std::path::Path,
|
||||
pubkey: &str,
|
||||
captured_scope: &crate::managed_agents::scope::WorkspaceAgentScope,
|
||||
) -> Result<(), String> {
|
||||
let state = app.state::<AppState>();
|
||||
let _store_guard = state
|
||||
.managed_agents_store_lock
|
||||
.lock()
|
||||
.map_err(|e| format!("failed to acquire managed agents store lock: {e}"))?;
|
||||
// Validate generation inside the lock — if the workspace switched during
|
||||
// the preceding await, abort rather than writing into the new scope's store.
|
||||
crate::managed_agents::scope::validate_scope_generation(captured_scope)
|
||||
.map_err(|e| format!("mesh recovery clear: {e}"))?;
|
||||
let mut records = crate::managed_agents::load_managed_agents_at(definitions_dir)?;
|
||||
let record = crate::managed_agents::find_managed_agent_mut(&mut records, pubkey)?;
|
||||
if !record
|
||||
|
||||
@@ -155,21 +155,18 @@ fn run_boot_migrations_inner(app: &tauri::AppHandle, reset_completed: bool) {
|
||||
|
||||
migrate_legacy_app_data_dir(app);
|
||||
sync_shared_agent_data(app);
|
||||
// Dev-build-only: copy any agent keys that exist in the production
|
||||
// keyring ("buzz-desktop") into the dev service ("buzz-desktop-dev")
|
||||
// so existing agents don't lose their keys after the service-name split.
|
||||
// Must run after sync_shared_agent_data (JSON symlinked) and before
|
||||
// any load_managed_agents call (which runs hydrate_keys against the
|
||||
// dev service and would log "has no key" for un-migrated entries).
|
||||
#[cfg(debug_assertions)]
|
||||
if is_dev {
|
||||
crate::managed_agents::migrate_agent_keys_to_dev_service(app);
|
||||
}
|
||||
migrate_persona_provider_to_runtime(app);
|
||||
// Definition-touching migrations (fold, strip, backfill, etc.) are NOT
|
||||
// run here. They run inside the per-scope initialization pipeline
|
||||
// (`scope_init::run_scoped_migrations`) after staged adoption so every
|
||||
// scope sees exactly the migrations appropriate to its data.
|
||||
// Definition-touching migrations (fold, strip, backfill, etc.) and the
|
||||
// dev-key keyring migration are NOT run here. They run inside the per-scope
|
||||
// initialization pipeline (`scope_init::run_scoped_migrations` and
|
||||
// `run_pre_ready_family`) after staged adoption so every scope sees exactly
|
||||
// the migrations appropriate to its data.
|
||||
//
|
||||
// `migrate_persona_provider_to_runtime` moved to `run_scoped_migrations`
|
||||
// as step 0 (before fold); runs on the scoped personas.json, not the legacy
|
||||
// `agents/personas.json`.
|
||||
//
|
||||
// `migrate_agent_keys_to_dev_service` moved to `run_pre_ready_family` as
|
||||
// step C (debug builds); runs after the scoped store is populated.
|
||||
}
|
||||
|
||||
/// Copy one-time app state from the legacy app identifier directory to
|
||||
@@ -1177,17 +1174,6 @@ fn rename_provider_to_runtime_in_personas(path: &Path) {
|
||||
eprintln!("buzz-desktop: rename-provider-to-runtime: {e}");
|
||||
}
|
||||
}
|
||||
|
||||
pub fn migrate_persona_provider_to_runtime(app: &tauri::AppHandle) {
|
||||
let Ok(dir) = app.path().app_data_dir() else {
|
||||
return;
|
||||
};
|
||||
let path = dir.join("agents/personas.json");
|
||||
if !path.exists() {
|
||||
return;
|
||||
}
|
||||
rename_provider_to_runtime_in_personas(&path);
|
||||
}
|
||||
mod fold;
|
||||
mod materialize;
|
||||
use fold::load_persona_runtimes;
|
||||
|
||||
@@ -4,6 +4,21 @@ pub(crate) use fold::fold_personas_in_dir;
|
||||
pub(crate) use materialize::materialize_runtimes_in_file;
|
||||
pub(crate) use team_suffix::strip_baked_team_instructions_in_dir;
|
||||
|
||||
/// Rename `provider` → `runtime` in a scoped `definitions_dir/personas.json`.
|
||||
///
|
||||
/// Runs BEFORE `fold_personas_in_dir` so the fold reads the correct `runtime`
|
||||
/// field. Idempotent: records that already have `runtime` are unchanged.
|
||||
/// Returns `Ok(())` when there is no `personas.json` to migrate.
|
||||
pub(crate) fn migrate_persona_provider_to_runtime_at(
|
||||
definitions_dir: &std::path::Path,
|
||||
) -> Result<(), String> {
|
||||
let path = definitions_dir.join("personas.json");
|
||||
if path.exists() {
|
||||
rename_provider_to_runtime_in_personas(&path);
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Reconcile `mcp_command` values in a scoped `definitions_dir`.
|
||||
pub(crate) fn reconcile_provider_mcp_commands_at(definitions_dir: &std::path::Path) -> Result<(), String> {
|
||||
let path = definitions_dir.join("managed-agents.json");
|
||||
|
||||
@@ -16,6 +16,9 @@ const MIGRATION_TOAST_KEY = "buzz-legacy-nest-migrated-notified";
|
||||
* legacy `~/.sprout` nest. Shown once per machine (deduped via
|
||||
* localStorage); the backend re-emits each launch while `~/.sprout` exists,
|
||||
* which also covers the event being emitted before this listener mounts.
|
||||
* - `workspace-degraded`: a post-commit restore or event-sync step failed
|
||||
* after the workspace switch succeeded. The switch is live; the degradation
|
||||
* is recoverable by restarting the app or re-applying the workspace.
|
||||
*
|
||||
* Mounted at the app root ahead of the community-init effect so the listener
|
||||
* is registered before the first `apply_workspace` call.
|
||||
@@ -38,9 +41,16 @@ export function useNestNotifications(): void {
|
||||
});
|
||||
});
|
||||
|
||||
const unlistenDegraded = listen<string>("workspace-degraded", (event) => {
|
||||
toast.error("Workspace partially degraded", {
|
||||
description: event.payload,
|
||||
});
|
||||
});
|
||||
|
||||
return () => {
|
||||
void unlistenReposError.then((fn) => fn());
|
||||
void unlistenMigrated.then((fn) => fn());
|
||||
void unlistenDegraded.then((fn) => fn());
|
||||
};
|
||||
}, []);
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user