mirror of
https://github.com/block/buzz.git
synced 2026-08-18 06:50:31 +02:00
feat(desktop): implement Phase 3 runtime ownership + Mesh scope rules
Phase 3 of workspace-scoped agent store:
3a: reconcile_managed_agent_runtimes loses its `communities` parameter.
The backend derives the sole target relay from the captured active scope;
cross-scope fan-out is no longer representable at the API level.
- runtime_commands.rs: capture active scope relay; remove communities Vec
- runtime_types.rs: remove ManagedAgentCommunityTarget struct
- tauriManagedAgents.ts: reconcileManagedAgentRuntimes() takes no args
- managedAgentRuntimeHooks.ts: bootstrapManagedAgentRuntimePairs calls
parameterless reconcile; drop communities list construction
- useManagedAgentRuntimeReconciliation.ts: rewritten to track a single
activeCommunityKey instead of per-relay state; simplified retry logic
- AppShell.tsx: pass `${activeCommunity?.id}-${reinitKey}` as the key
3b: Mesh relay-match reuse rule + fail-closed serve preflight + watchdog.
- mesh_llm.rs: ensure_relay_mesh_for_record captures scope relay at entry;
a live runtime is only reused when its relay matches the scope relay;
serve-mode mismatch fails closed with a precise 'Share Compute is
currently pinned to <relay>' error; client-mode mismatch falls through
to re-arm; drain_mesh_client_if_stale drains a client whose relay
differs from the incoming workspace relay (Layer-1 async, non-fatal).
- recovery.rs: rearm_relay_mesh_for_running_agents captures one scope per
pass; Live early-return only taken on relay match; serve-mode Live
mismatch skips the pass (machine-level pinning).
- personas.rs: add scoped load_personas_at / save_personas_at variants.
3c: Drain journal + compensation + apply_workspace rewrite.
- runtime_commands.rs: DrainJournalEntry struct, drain_scope_runtimes
(snapshot journal + stop all live runtimes, returns stopped/remaining/
first_error), compensate_drain (restart exactly the stopped entries).
- workspace.rs: apply_workspace return type changed from () to
WorkspaceApplyResult. Layer-1 async drains the Mesh client before
spawn_blocking. Drain stage acquires managed_agent_runtime_transition,
calls drain_scope_runtimes; on failure calls compensate_drain and
returns applied:false. Per-transition restore replaces the launch-only
managed_agent_restore_pending one-shot. Post-commit failures (event
sync, restore) surface as degraded entries on WorkspaceApplyResult.
3d: Scope-tagged runtime map entries.
- runtime_types.rs: ManagedAgentPairRuntime gains scope_id: Option<String>
- starting() constructor takes scope_id; captured from active scope at
spawn time in runtime_commands.rs, restore.rs, and runtime.rs.
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
b795ed6c79
commit
00a3381902
@@ -813,6 +813,16 @@ fn pick_serve_target_for_model(
|
||||
/// a relay query failure ("could not refresh targets") is not the same as a
|
||||
/// relay that answered with no live target for this model ("peer offline").
|
||||
/// Non relay-mesh records are a no-op.
|
||||
///
|
||||
/// **Scope ownership rule (v4):** a live runtime is only reused when its bound
|
||||
/// relay matches the active workspace scope's relay. A relay mismatch means the
|
||||
/// runtime belongs to a different scope:
|
||||
/// - Serve mode (Share Compute): fail closed with a precise error — the
|
||||
/// process has one singleton runtime slot and one `:9337` ingress; no client
|
||||
/// can start while serve occupies it; the user must stop sharing first.
|
||||
/// - Client mode: treat as absent and fall through to re-arm (the old client
|
||||
/// should have been drained on the workspace switch, but this is a safety
|
||||
/// net for any edge where drain did not reach it).
|
||||
pub(crate) async fn ensure_relay_mesh_for_record(
|
||||
app: &AppHandle,
|
||||
model_id: Option<&str>,
|
||||
@@ -822,6 +832,14 @@ pub(crate) async fn ensure_relay_mesh_for_record(
|
||||
let Some(model_id) = model_id else {
|
||||
return Ok(());
|
||||
};
|
||||
|
||||
// Capture the active scope relay once. None means no workspace applied yet
|
||||
// — fail closed rather than starting a mesh client with an unknown relay.
|
||||
let scope_relay = state
|
||||
.capture_active_scope()
|
||||
.map(|scope| scope.relay_url.clone())
|
||||
.ok_or_else(|| "Buzz shared compute cannot start: no active workspace scope".to_string())?;
|
||||
|
||||
// A local serve/client runtime already owns the OpenAI ingress and its
|
||||
// router can resolve both `auto` and explicit remote models. Do not require
|
||||
// a separate relay-advertised target in that case — BUT only trust it when
|
||||
@@ -832,37 +850,75 @@ pub(crate) async fn ensure_relay_mesh_for_record(
|
||||
// runtime and fall through to re-arm it. The mesh coordinator watchdog also
|
||||
// calls this path after eviction so recovery is not start-only (Brad #2304).
|
||||
if state.mesh_llm_runtime.lock().await.is_some() {
|
||||
match mesh_llm::recover_stale_mesh_runtime(
|
||||
&state,
|
||||
mesh_llm::MeshRecoveryUrgency::Foreground,
|
||||
)
|
||||
.await
|
||||
{
|
||||
mesh_llm::MeshRuntimeRecovery::Live => {
|
||||
return wait_for_mesh_inference(model_id).await;
|
||||
// Before probing liveness, verify the runtime's relay matches the
|
||||
// active scope. A mismatch means it belongs to a switched-away scope.
|
||||
let (runtime_relay, runtime_mode) = {
|
||||
let guard = state.mesh_llm_runtime.lock().await;
|
||||
let relay = guard
|
||||
.as_ref()
|
||||
.and_then(|r| r.start_request().relay_url.clone());
|
||||
let mode = guard.as_ref().map(|r| r.mode());
|
||||
(relay, mode)
|
||||
};
|
||||
|
||||
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(&scope_relay)
|
||||
});
|
||||
|
||||
if !relay_matches {
|
||||
match runtime_mode {
|
||||
Some(mesh_llm::MeshNodeMode::Serve) => {
|
||||
// Fail closed: Share Compute is pinned to another relay.
|
||||
// The process has one runtime slot and one :9337 ingress.
|
||||
// No client can start while serve occupies it.
|
||||
let pinned_relay = runtime_relay.as_deref().unwrap_or("another relay");
|
||||
return Err(format!(
|
||||
"Share Compute is currently pinned to {pinned_relay}. \
|
||||
Stop sharing first, then switch workspaces to use \
|
||||
Buzz shared compute on this workspace."
|
||||
));
|
||||
}
|
||||
Some(mesh_llm::MeshNodeMode::Client) | None => {
|
||||
// Stale client from a prior workspace. Treat as absent —
|
||||
// fall through to re-arm a new client for the active scope.
|
||||
// The drain stage in apply_workspace should have cleared
|
||||
// this; this is a safety net for missed drains.
|
||||
}
|
||||
}
|
||||
mesh_llm::MeshRuntimeRecovery::Evicted | mesh_llm::MeshRuntimeRecovery::Absent => {}
|
||||
mesh_llm::MeshRuntimeRecovery::Debouncing => {
|
||||
return Err(
|
||||
"Buzz shared compute ingress is temporarily unresponsive; recovery is already scheduled. Try again shortly."
|
||||
.to_string(),
|
||||
);
|
||||
}
|
||||
mesh_llm::MeshRuntimeRecovery::ReleasePending => {
|
||||
return Err(
|
||||
"Buzz shared compute is still shutting down its previous local ingress. Try again shortly."
|
||||
.to_string(),
|
||||
);
|
||||
}
|
||||
mesh_llm::MeshRuntimeRecovery::Replaced => {
|
||||
return wait_for_mesh_inference(model_id).await;
|
||||
}
|
||||
mesh_llm::MeshRuntimeRecovery::RestartRequired => {
|
||||
app.request_restart();
|
||||
return Err(
|
||||
"Buzz shared compute startup lost its local ingress before shutdown control became available. Buzz is restarting to recover it."
|
||||
.to_string(),
|
||||
);
|
||||
} else {
|
||||
match mesh_llm::recover_stale_mesh_runtime(
|
||||
&state,
|
||||
mesh_llm::MeshRecoveryUrgency::Foreground,
|
||||
)
|
||||
.await
|
||||
{
|
||||
mesh_llm::MeshRuntimeRecovery::Live => {
|
||||
return wait_for_mesh_inference(model_id).await;
|
||||
}
|
||||
mesh_llm::MeshRuntimeRecovery::Evicted | mesh_llm::MeshRuntimeRecovery::Absent => {}
|
||||
mesh_llm::MeshRuntimeRecovery::Debouncing => {
|
||||
return Err(
|
||||
"Buzz shared compute ingress is temporarily unresponsive; recovery is already scheduled. Try again shortly."
|
||||
.to_string(),
|
||||
);
|
||||
}
|
||||
mesh_llm::MeshRuntimeRecovery::ReleasePending => {
|
||||
return Err(
|
||||
"Buzz shared compute is still shutting down its previous local ingress. Try again shortly."
|
||||
.to_string(),
|
||||
);
|
||||
}
|
||||
mesh_llm::MeshRuntimeRecovery::Replaced => {
|
||||
return wait_for_mesh_inference(model_id).await;
|
||||
}
|
||||
mesh_llm::MeshRuntimeRecovery::RestartRequired => {
|
||||
app.request_restart();
|
||||
return Err(
|
||||
"Buzz shared compute startup lost its local ingress before shutdown control became available. Buzz is restarting to recover it."
|
||||
.to_string(),
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -900,6 +956,63 @@ pub(crate) async fn ensure_relay_mesh_for_record(
|
||||
wait_for_mesh_inference(model_id).await
|
||||
}
|
||||
|
||||
/// 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).
|
||||
///
|
||||
/// Serve-mode runtimes are machine-level and are deliberately NOT drained
|
||||
/// on workspace switch — they stay pinned to their configured relay.
|
||||
///
|
||||
/// 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.
|
||||
#[cfg(feature = "mesh-llm")]
|
||||
pub(crate) async fn drain_mesh_client_if_stale(
|
||||
app: &AppHandle,
|
||||
active_relay_url: &str,
|
||||
) -> 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;
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
#[tauri::command]
|
||||
pub async fn mesh_stop_node(
|
||||
app: AppHandle,
|
||||
|
||||
@@ -116,13 +116,19 @@ pub async fn validate_repos_dir(dir: String) -> Result<(), String> {
|
||||
/// Tauri backend with the selected workspace's relay URL, keys, and repos
|
||||
/// directory.
|
||||
///
|
||||
/// Returns `WorkspaceApplyResult`:
|
||||
/// - `applied: true` → new scope committed; post-commit failures surface as
|
||||
/// `degraded` entries (informational — workspace IS active).
|
||||
/// - `applied: false` → drain failed; old scope still active; `degraded`
|
||||
/// names what could not be stopped or restored by compensation.
|
||||
///
|
||||
/// A bad `repos_dir` is non-fatal: relay/keys always apply (the relay is the
|
||||
/// active workspace's own choice — orthogonal to the filesystem repos dir),
|
||||
/// the bad value is NOT persisted (so the next boot starts clean), the
|
||||
/// `REPOS` symlink is skipped (REPOS stays a real dir), a `repos-dir-error`
|
||||
/// event surfaces the reason, and the command returns `Ok`. The dialogs
|
||||
/// already block a bad path at Save (`validate_repos_dir`); this fallback only
|
||||
/// catches a value that went bad after save (deleted dir, unmounted volume).
|
||||
/// event surfaces the reason. The dialogs already block a bad path at Save
|
||||
/// (`validate_repos_dir`); this fallback only catches a value that went bad
|
||||
/// after save (deleted dir, unmounted volume).
|
||||
#[tauri::command]
|
||||
pub async fn apply_workspace(
|
||||
relay_url: String,
|
||||
@@ -130,7 +136,9 @@ pub async fn apply_workspace(
|
||||
repos_dir: Option<String>,
|
||||
agent_managed_profiles: Option<bool>,
|
||||
app: AppHandle,
|
||||
) -> Result<(), String> {
|
||||
) -> Result<crate::managed_agents::scope::WorkspaceApplyResult, String> {
|
||||
use crate::managed_agents::scope::WorkspaceApplyResult;
|
||||
|
||||
// ── Layer 1: async serialization lock ────────────────────────────────────
|
||||
// workspace_transition serializes apply_workspace and live identity import
|
||||
// so scope transitions are never concurrent. We acquire via a clone so the
|
||||
@@ -139,142 +147,162 @@ 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.
|
||||
#[cfg(feature = "mesh-llm")]
|
||||
{
|
||||
if let Err(error) =
|
||||
crate::commands::mesh_llm::drain_mesh_client_if_stale(&app, &relay_url).await
|
||||
{
|
||||
eprintln!("buzz-desktop: Mesh client drain before workspace switch failed: {error}");
|
||||
}
|
||||
}
|
||||
|
||||
let restore_app = app.clone();
|
||||
tokio::task::spawn_blocking(move || {
|
||||
let state = app.state::<AppState>();
|
||||
let blocking_result: Result<WorkspaceApplyResult, String> =
|
||||
tokio::task::spawn_blocking(move || {
|
||||
let state = app.state::<AppState>();
|
||||
|
||||
// ── Validate before mutating ──────────────────────────────────────────
|
||||
let parsed_keys = match nsec.as_deref().map(str::trim).filter(|s| !s.is_empty()) {
|
||||
Some(nsec_trimmed) => {
|
||||
Some(Keys::parse(nsec_trimmed).map_err(|e| format!("invalid nsec: {e}"))?)
|
||||
}
|
||||
None => None,
|
||||
};
|
||||
|
||||
// Decide the effective repos_dir from the candidate. A bad path does NOT
|
||||
// reject — it is treated as if no override were set: relay/keys still
|
||||
// apply, the bad value is not persisted, and a `repos-dir-error` surfaces
|
||||
// the reason. Persisting a bad path would make every later boot read it,
|
||||
// fail to resolve the symlink, and silently skip agent restore. One
|
||||
// validate (inside `effective_repos_dir`) drives both the emit and the
|
||||
// persisted value. `nest` is resolved softly: when absent there is nothing
|
||||
// to persist or symlink, and relay/keys must still apply unconditionally.
|
||||
let nest = nest_dir();
|
||||
let effective_repos_dir = match nest.as_deref() {
|
||||
Some(nest) => match effective_repos_dir(nest, repos_dir.as_deref()) {
|
||||
Ok(value) => value,
|
||||
Err(error) => {
|
||||
let _ = app.emit("repos-dir-error", error);
|
||||
None
|
||||
// ── Validate before mutating ──────────────────────────────────────
|
||||
let parsed_keys = match nsec.as_deref().map(str::trim).filter(|s| !s.is_empty()) {
|
||||
Some(nsec_trimmed) => {
|
||||
Some(Keys::parse(nsec_trimmed).map_err(|e| format!("invalid nsec: {e}"))?)
|
||||
}
|
||||
},
|
||||
None => None,
|
||||
};
|
||||
None => None,
|
||||
};
|
||||
|
||||
// ── Prepare: derive target scope and run staged initialization ────────
|
||||
// This is the reversible prepare stage: the old scope remains active
|
||||
// throughout. We derive the effective owner pubkey (candidate keys win
|
||||
// over existing, mirroring the commit below) and call ensure_scope_ready
|
||||
// which handles the canonical claim ledger, staged install, idempotent
|
||||
// migrations, and the Ready marker. Any error here leaves the old scope
|
||||
// untouched and returns Err before any state mutation.
|
||||
let base_dir = crate::managed_agents::managed_agents_base_dir(&app).unwrap_or_default();
|
||||
let effective_owner_pubkey = match &parsed_keys {
|
||||
Some(keys) => keys.public_key().to_hex(),
|
||||
None => state
|
||||
.keys
|
||||
.lock()
|
||||
.map_err(|e| e.to_string())?
|
||||
.public_key()
|
||||
.to_hex(),
|
||||
};
|
||||
let target_scope_id =
|
||||
crate::managed_agents::scope::derive_scope_id(&relay_url, &effective_owner_pubkey);
|
||||
let scope_dir =
|
||||
crate::managed_agents::scope::scoped_definitions_dir(&base_dir, &target_scope_id);
|
||||
crate::managed_agents::scope_init::ensure_scope_ready(
|
||||
&target_scope_id,
|
||||
&scope_dir,
|
||||
&base_dir,
|
||||
)?;
|
||||
// Decide the effective repos_dir from the candidate. A bad path does NOT
|
||||
// reject — it is treated as if no override were set: relay/keys still
|
||||
// apply, the bad value is not persisted, and a `repos-dir-error` surfaces
|
||||
// the reason.
|
||||
let nest = nest_dir();
|
||||
let effective_repos_dir = match nest.as_deref() {
|
||||
Some(nest) => match effective_repos_dir(nest, repos_dir.as_deref()) {
|
||||
Ok(value) => value,
|
||||
Err(error) => {
|
||||
let _ = app.emit("repos-dir-error", error);
|
||||
None
|
||||
}
|
||||
},
|
||||
None => None,
|
||||
};
|
||||
|
||||
// ── Layer 2: synchronous commit epoch ────────────────────────────────
|
||||
// No .await may be held while any Layer-2 guard is live. Relay override,
|
||||
// keys, and the active scope are all committed in this critical section.
|
||||
{
|
||||
let mut override_guard = state.relay_url_override.lock().map_err(|e| e.to_string())?;
|
||||
*override_guard = Some(relay_url.clone());
|
||||
}
|
||||
// Reset the Rust-side admission gate when switching workspace/community,
|
||||
// matching `resetRateLimitGate()` on the TS side (useCommunityInit.ts:38).
|
||||
crate::relay_admission::reset_gate_for_workspace_change();
|
||||
|
||||
if let Some(keys) = parsed_keys {
|
||||
let mut keys_guard = state.keys.lock().map_err(|e| e.to_string())?;
|
||||
*keys_guard = keys;
|
||||
}
|
||||
|
||||
// Keep the backend-side reconcile guard aligned with the frontend
|
||||
// experiment before launch-time restore can spawn any agents. Missing
|
||||
// means the stable behavior: desktop remains authoritative.
|
||||
state
|
||||
.managed_agent_profile_reconcile_enabled
|
||||
.store(!agent_managed_profiles.unwrap_or(false), Ordering::Release);
|
||||
|
||||
// ── Commit the active workspace agent scope ───────────────────────────
|
||||
// Derive the scope from the now-applied relay + owner keys and commit it
|
||||
// as the active scope. All subsequent store reads/writes (via
|
||||
// load_managed_agents, save_managed_agents, etc.) resolve through
|
||||
// capture_active_scope() → scoped definitions directory. There is NO
|
||||
// fallback to the legacy unscoped root.
|
||||
{
|
||||
let owner_pubkey = state
|
||||
.keys
|
||||
.lock()
|
||||
.map_err(|e| e.to_string())?
|
||||
.public_key()
|
||||
.to_hex();
|
||||
let generation = crate::managed_agents::scope::next_scope_generation();
|
||||
let scope = crate::managed_agents::scope::WorkspaceAgentScope::new(
|
||||
relay_url,
|
||||
owner_pubkey,
|
||||
// ── Prepare: derive target scope and run staged initialization ────
|
||||
// Reversible prepare stage: the old scope remains active throughout.
|
||||
let base_dir = crate::managed_agents::managed_agents_base_dir(&app).unwrap_or_default();
|
||||
let effective_owner_pubkey = match &parsed_keys {
|
||||
Some(keys) => keys.public_key().to_hex(),
|
||||
None => state
|
||||
.keys
|
||||
.lock()
|
||||
.map_err(|e| e.to_string())?
|
||||
.public_key()
|
||||
.to_hex(),
|
||||
};
|
||||
let target_scope_id =
|
||||
crate::managed_agents::scope::derive_scope_id(&relay_url, &effective_owner_pubkey);
|
||||
let scope_dir =
|
||||
crate::managed_agents::scope::scoped_definitions_dir(&base_dir, &target_scope_id);
|
||||
crate::managed_agents::scope_init::ensure_scope_ready(
|
||||
&target_scope_id,
|
||||
&scope_dir,
|
||||
&base_dir,
|
||||
generation,
|
||||
);
|
||||
state.commit_active_scope(scope);
|
||||
}
|
||||
)?;
|
||||
|
||||
// ── Filesystem side-effect (non-fatal) ────────────────────────────────
|
||||
// Persist the *effective* repos_dir (None when the candidate failed
|
||||
// validation) for the backend to read at boot, then re-point REPOS to
|
||||
// match. Persisting first makes the dotfile authoritative even if the
|
||||
// symlink apply fails here (e.g. a non-empty real REPOS): the next boot
|
||||
// reads the persisted value and resolves the symlink before any agent can
|
||||
// clone into REPOS. A bad candidate persists `None`, so the next boot is
|
||||
// clean and agent restore proceeds. Failure of either must NOT fail the
|
||||
// command — relay/keys are already applied. Surface symlink errors via
|
||||
// `repos-dir-error`.
|
||||
if let Some(nest) = nest.as_deref() {
|
||||
if let Err(error) = write_persisted_repos_dir(nest, effective_repos_dir.as_deref()) {
|
||||
eprintln!("buzz-desktop: persist repos dir failed: {error}");
|
||||
// ── Drain: journal + stop all old-scope runtimes ──────────────────
|
||||
// Before committing the new scope, drain all live managed-agent
|
||||
// processes. Take managed_agent_runtime_transition (Layer 2) now;
|
||||
// the commit below also holds it.
|
||||
let (stopped_entries, _remaining, drain_error) = {
|
||||
let _rt_transition = state
|
||||
.managed_agent_runtime_transition
|
||||
.lock()
|
||||
.map_err(|e| e.to_string())?;
|
||||
crate::managed_agents::drain_scope_runtimes(&app, &state)
|
||||
};
|
||||
|
||||
if let Some(drain_err) = drain_error {
|
||||
// Drain failed — compensate by restarting what we stopped.
|
||||
let comp_err = crate::managed_agents::compensate_drain(&app, &stopped_entries);
|
||||
let degraded_msg = match comp_err {
|
||||
Some(comp) => {
|
||||
format!("drain failed ({drain_err}); compensation also failed: {comp}")
|
||||
}
|
||||
None => format!("drain failed ({drain_err}); old runtimes restored"),
|
||||
};
|
||||
return Ok(WorkspaceApplyResult::drain_failed(degraded_msg));
|
||||
}
|
||||
if let Err(error) = ensure_repos_symlink(nest, effective_repos_dir.as_deref()) {
|
||||
eprintln!("buzz-desktop: repos dir setup failed: {error}");
|
||||
let _ = app.emit("repos-dir-error", error);
|
||||
|
||||
// ── Layer 2: synchronous commit epoch ────────────────────────────
|
||||
// No .await may be held while any Layer-2 guard is live.
|
||||
{
|
||||
let mut override_guard =
|
||||
state.relay_url_override.lock().map_err(|e| e.to_string())?;
|
||||
*override_guard = Some(relay_url.clone());
|
||||
}
|
||||
}
|
||||
crate::relay_admission::reset_gate_for_workspace_change();
|
||||
|
||||
try_regenerate_nest(&app);
|
||||
if let Some(keys) = parsed_keys {
|
||||
let mut keys_guard = state.keys.lock().map_err(|e| e.to_string())?;
|
||||
*keys_guard = keys;
|
||||
}
|
||||
|
||||
Ok::<(), String>(())
|
||||
})
|
||||
.await
|
||||
.map_err(|e| format!("spawn_blocking failed: {e}"))??;
|
||||
state
|
||||
.managed_agent_profile_reconcile_enabled
|
||||
.store(!agent_managed_profiles.unwrap_or(false), Ordering::Release);
|
||||
|
||||
// ── Commit the active workspace agent scope ───────────────────────
|
||||
{
|
||||
let owner_pubkey = state
|
||||
.keys
|
||||
.lock()
|
||||
.map_err(|e| e.to_string())?
|
||||
.public_key()
|
||||
.to_hex();
|
||||
let generation = crate::managed_agents::scope::next_scope_generation();
|
||||
let scope = crate::managed_agents::scope::WorkspaceAgentScope::new(
|
||||
relay_url,
|
||||
owner_pubkey,
|
||||
&base_dir,
|
||||
generation,
|
||||
);
|
||||
state.commit_active_scope(scope);
|
||||
}
|
||||
|
||||
// ── Filesystem side-effects (non-fatal) ───────────────────────────
|
||||
if let Some(nest) = nest.as_deref() {
|
||||
if let Err(error) = write_persisted_repos_dir(nest, effective_repos_dir.as_deref())
|
||||
{
|
||||
eprintln!("buzz-desktop: persist repos dir failed: {error}");
|
||||
}
|
||||
if let Err(error) = ensure_repos_symlink(nest, effective_repos_dir.as_deref()) {
|
||||
eprintln!("buzz-desktop: repos dir setup failed: {error}");
|
||||
let _ = app.emit("repos-dir-error", error);
|
||||
}
|
||||
}
|
||||
|
||||
try_regenerate_nest(&app);
|
||||
|
||||
Ok::<WorkspaceApplyResult, String>(WorkspaceApplyResult::success())
|
||||
})
|
||||
.await
|
||||
.map_err(|e| format!("spawn_blocking failed: {e}"))?;
|
||||
|
||||
// If blocking returned a drain-failed result, surface it now.
|
||||
let apply_result = blocking_result?;
|
||||
if !apply_result.applied {
|
||||
return Ok(apply_result);
|
||||
}
|
||||
|
||||
// ── Post-commit (non-rollback) ────────────────────────────────────────────
|
||||
// The workspace HAS switched. Post-commit failures surface as degradation
|
||||
// on the applied result — we never pretend the old scope survived.
|
||||
let mut degraded: Vec<String> = Vec::new();
|
||||
|
||||
let state = restore_app.state::<AppState>();
|
||||
// Backfill this exact relay+owner scope only after the workspace has been
|
||||
// applied. Running at process boot would target the fallback relay and
|
||||
// collapse every community into one pending-event store.
|
||||
match crate::managed_agents::retention::active_retention_scope(&restore_app, &state) {
|
||||
Ok(scope) => {
|
||||
// Adopt whatever the pre-scoping release left queued in the global
|
||||
@@ -283,10 +311,6 @@ pub async fn apply_workspace(
|
||||
// instead of being abandoned by the storage cutover.
|
||||
migrate_legacy_retention_into(&restore_app, &scope);
|
||||
|
||||
// The active scope was committed in the spawn_blocking above.
|
||||
// If it is somehow None here, event sync is skipped rather than
|
||||
// falling back to the legacy unscoped root (which would recreate
|
||||
// split-brain storage).
|
||||
if let Some(agent_scope) = state.capture_active_scope() {
|
||||
crate::event_sync::spawn_event_sync(
|
||||
restore_app.clone(),
|
||||
@@ -295,53 +319,44 @@ pub async fn apply_workspace(
|
||||
agent_scope.definitions_dir,
|
||||
);
|
||||
} else {
|
||||
eprintln!(
|
||||
"buzz-desktop: active agent scope unavailable after workspace apply — \
|
||||
event sync skipped"
|
||||
degraded.push(
|
||||
"active agent scope unavailable after workspace apply — event sync skipped"
|
||||
.to_string(),
|
||||
);
|
||||
}
|
||||
}
|
||||
Err(error) => {
|
||||
eprintln!("buzz-desktop: scoped event-sync unavailable after workspace apply: {error}");
|
||||
degraded.push(format!(
|
||||
"scoped event-sync unavailable after workspace apply: {error}"
|
||||
));
|
||||
}
|
||||
}
|
||||
|
||||
let restore_pending = state
|
||||
.managed_agent_restore_pending
|
||||
.swap(false, Ordering::AcqRel);
|
||||
|
||||
// The coordinator starts before React applies the selected workspace, so
|
||||
// its startup publication may have used the fallback relay and placeholder
|
||||
// identity. Correct it off the command path so an unavailable relay cannot
|
||||
// hold the frontend on its loading gate. On initial launch, restore MeshLLM
|
||||
// first so a slow stopped-status request cannot overwrite a newly restored
|
||||
// serving status, then restore managed agents after the admission identity
|
||||
// has been published (or the bounded publication attempt has timed out).
|
||||
// Per-transition restore: always restore the new scope's auto-start agents
|
||||
// (replaces the launch-only `managed_agent_restore_pending.swap` one-shot).
|
||||
// Fire-and-forget spawn so the command returns promptly; failures are logged.
|
||||
#[cfg(feature = "mesh-llm")]
|
||||
{
|
||||
let app = restore_app.clone();
|
||||
tauri::async_runtime::spawn(async move {
|
||||
let state = app.state::<AppState>();
|
||||
if restore_pending {
|
||||
if let Err(error) =
|
||||
crate::commands::mesh_llm::restore_mesh_sharing(&app, &state).await
|
||||
{
|
||||
eprintln!("buzz-desktop: failed to restore Share Compute: {error}");
|
||||
}
|
||||
// Restore mesh sharing first so a slow stopped-status request cannot
|
||||
// overwrite a newly restored serving status.
|
||||
if let Err(error) = crate::commands::mesh_llm::restore_mesh_sharing(&app, &state).await
|
||||
{
|
||||
eprintln!("buzz-desktop: failed to restore Share Compute: {error}");
|
||||
}
|
||||
crate::mesh_llm::publish_current_status_once(&app, "workspace apply").await;
|
||||
if restore_pending {
|
||||
if let Err(error) =
|
||||
restore_managed_agents_on_launch(&app, &state.shutdown_started).await
|
||||
{
|
||||
eprintln!("buzz-desktop: failed to restore managed agents: {error}");
|
||||
}
|
||||
if let Err(error) =
|
||||
restore_managed_agents_on_launch(&app, &state.shutdown_started).await
|
||||
{
|
||||
eprintln!("buzz-desktop: failed to restore managed agents: {error}");
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
#[cfg(not(feature = "mesh-llm"))]
|
||||
if restore_pending {
|
||||
{
|
||||
let app = restore_app.clone();
|
||||
tauri::async_runtime::spawn(async move {
|
||||
let state = app.state::<AppState>();
|
||||
@@ -353,5 +368,13 @@ pub async fn apply_workspace(
|
||||
});
|
||||
}
|
||||
|
||||
Ok(())
|
||||
if degraded.is_empty() {
|
||||
Ok(WorkspaceApplyResult::success())
|
||||
} else {
|
||||
Ok(degraded
|
||||
.into_iter()
|
||||
.fold(WorkspaceApplyResult::success(), |r, msg| {
|
||||
r.with_degradation(msg)
|
||||
}))
|
||||
}
|
||||
}
|
||||
|
||||
@@ -345,6 +345,28 @@ pub fn load_personas(app: &AppHandle) -> Result<Vec<AgentDefinition>, String> {
|
||||
Ok(records)
|
||||
}
|
||||
|
||||
/// Scoped variant of [`load_personas`]: load from an explicit definitions
|
||||
/// directory instead of resolving through the active scope. Used by
|
||||
/// operations that captured a [`WorkspaceAgentScope`] at entry to guarantee
|
||||
/// scope stability across awaits.
|
||||
pub(crate) fn load_personas_at(
|
||||
definitions_dir: &std::path::Path,
|
||||
) -> Result<Vec<AgentDefinition>, String> {
|
||||
let now = now_iso();
|
||||
|
||||
let records = crate::managed_agents::storage::load_agent_definitions_at(definitions_dir)?
|
||||
.iter()
|
||||
.filter_map(|record| record.to_definition_view())
|
||||
.collect();
|
||||
|
||||
let (records, changed) = merge_personas(records, &now);
|
||||
if changed {
|
||||
save_personas_at(definitions_dir, &records)?;
|
||||
}
|
||||
|
||||
Ok(records)
|
||||
}
|
||||
|
||||
/// Read the raw persona records at `path` — no built-in merge, no write-back.
|
||||
/// The single disk-read seam for persona definitions: `load_personas` layers
|
||||
/// the built-in merge on top, and the boot-time readers that need raw records
|
||||
@@ -376,5 +398,22 @@ pub fn save_personas(app: &AppHandle, records: &[AgentDefinition]) -> Result<(),
|
||||
crate::managed_agents::storage::save_agent_definitions(app, &definitions)
|
||||
}
|
||||
|
||||
/// Scoped variant of [`save_personas`]: write to an explicit definitions
|
||||
/// directory. Used by [`load_personas_at`] write-back and any operation that
|
||||
/// captured a [`WorkspaceAgentScope`] at entry.
|
||||
pub(crate) fn save_personas_at(
|
||||
definitions_dir: &std::path::Path,
|
||||
records: &[AgentDefinition],
|
||||
) -> Result<(), String> {
|
||||
let mut sorted = records.to_vec();
|
||||
sort_personas(&mut sorted);
|
||||
|
||||
let definitions: Vec<_> = sorted
|
||||
.into_iter()
|
||||
.map(|persona| persona.into_agent_record())
|
||||
.collect();
|
||||
crate::managed_agents::storage::save_agent_definitions_at(definitions_dir, &definitions)
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests;
|
||||
|
||||
@@ -425,7 +425,10 @@ pub async fn restore_managed_agents_on_launch(
|
||||
record.last_stopped_at = None;
|
||||
record.last_exit_code = None;
|
||||
record.last_error = None;
|
||||
runtimes.insert(key, super::ManagedAgentPairRuntime::starting(process));
|
||||
runtimes.insert(
|
||||
key,
|
||||
super::ManagedAgentPairRuntime::starting(process, Some(scope.scope_id.clone())),
|
||||
);
|
||||
successfully_spawned.push(pubkey);
|
||||
}
|
||||
SpawnOutcome::Failed(error) => {
|
||||
@@ -495,7 +498,7 @@ pub async fn restore_managed_agents_on_launch(
|
||||
|
||||
#[cfg(feature = "mesh-llm")]
|
||||
fn persist_restore_error(
|
||||
app: &tauri::AppHandle,
|
||||
_app: &tauri::AppHandle,
|
||||
state: &AppState,
|
||||
pubkey: &str,
|
||||
definitions_dir: &std::path::Path,
|
||||
|
||||
@@ -1019,7 +1019,12 @@ pub fn start_managed_agent_process(
|
||||
record.last_error = None;
|
||||
record.last_error_code = None;
|
||||
|
||||
runtimes.insert(key, ManagedAgentPairRuntime::starting(process));
|
||||
let scope_id = {
|
||||
use tauri::Manager;
|
||||
let state = app.state::<crate::app_state::AppState>();
|
||||
state.capture_active_scope().map(|s| s.scope_id.clone())
|
||||
};
|
||||
runtimes.insert(key, ManagedAgentPairRuntime::starting(process, scope_id));
|
||||
Ok(())
|
||||
}
|
||||
|
||||
|
||||
@@ -296,6 +296,9 @@ fn start_pair(
|
||||
.lock()
|
||||
.ok()
|
||||
.map(|keys| keys.public_key().to_hex());
|
||||
let scope_id = state
|
||||
.capture_active_scope()
|
||||
.map(|scope| scope.scope_id.clone());
|
||||
let mut process = spawn_agent_child(&app, record, &key.relay_url, lazy, owner.as_deref())?;
|
||||
let now = crate::util::now_iso();
|
||||
let receipt = ManagedAgentRuntimeReceipt {
|
||||
@@ -314,7 +317,10 @@ fn start_pair(
|
||||
record.last_started_at = Some(now);
|
||||
record.last_stopped_at = None;
|
||||
record.last_error = None;
|
||||
runtimes.insert(key.clone(), ManagedAgentPairRuntime::starting(process));
|
||||
runtimes.insert(
|
||||
key.clone(),
|
||||
ManagedAgentPairRuntime::starting(process, scope_id),
|
||||
);
|
||||
let status = status_for(&app, record, &key, runtimes.get(&key), None);
|
||||
drop(runtimes);
|
||||
save_managed_agents(&app, &records)?;
|
||||
@@ -459,35 +465,37 @@ fn unkeyable_failed_status(
|
||||
}
|
||||
}
|
||||
|
||||
/// Spawn a lazy harness pair for every eligible (agent, community) pair.
|
||||
/// Spawn a lazy harness pair for every auto-start local agent in the active
|
||||
/// workspace scope.
|
||||
///
|
||||
/// Eligibility is deliberately gated on `start_on_app_launch`: auto-start is
|
||||
/// the *proactive fan-out* policy — "keep this agent warm in every community" —
|
||||
/// not a correctness prerequisite. A manual-start agent still works on demand
|
||||
/// everywhere: attaching it to a channel ensures its pair, an @mention wakes a
|
||||
/// pair, the members sidebar and Settings controls start pairs, and restore
|
||||
/// preserves running pairs across relaunch. Fanning out warm-socket pairs for
|
||||
/// agents the user chose *not* to auto-start would contradict that choice, so
|
||||
/// reconcile leaves them alone until something explicitly asks for them.
|
||||
/// The target relay is derived from the captured active scope — the
|
||||
/// `communities` fan-out parameter has been removed. Under the active-scope-only
|
||||
/// runtime policy, reconcile targets exactly one relay: the relay the current
|
||||
/// workspace is bound to. Cross-scope fan-out is no longer representable at the
|
||||
/// API level.
|
||||
///
|
||||
/// Eligibility is gated on `start_on_app_launch`: auto-start is the proactive
|
||||
/// fan-out policy — agents not set to auto-start are left alone until something
|
||||
/// explicitly asks for them.
|
||||
#[tauri::command]
|
||||
pub async fn reconcile_managed_agent_runtimes(
|
||||
communities: Vec<super::ManagedAgentCommunityTarget>,
|
||||
app: AppHandle,
|
||||
) -> Result<Vec<ManagedAgentRuntimeStatus>, String> {
|
||||
use futures_util::{stream, StreamExt};
|
||||
|
||||
let state = app.state::<AppState>();
|
||||
let scope = state
|
||||
.capture_active_scope()
|
||||
.ok_or_else(|| "reconcile_managed_agent_runtimes: no active workspace scope".to_string())?;
|
||||
let relay_url = scope.relay_url.clone();
|
||||
|
||||
let records = load_managed_agents(&app)?;
|
||||
let mut jobs = Vec::new();
|
||||
for community in communities {
|
||||
for record in records
|
||||
.iter()
|
||||
.filter(|record| record.start_on_app_launch && record.backend == BackendKind::Local)
|
||||
// The legacy per-record relay pin is deliberately ignored here — see
|
||||
// `effective_agent_relay_url`. Every local auto-start agent fans out
|
||||
// to every configured community.
|
||||
{
|
||||
jobs.push((record.clone(), community.relay_url.clone()));
|
||||
}
|
||||
for record in records
|
||||
.iter()
|
||||
.filter(|record| record.start_on_app_launch && record.backend == BackendKind::Local)
|
||||
{
|
||||
jobs.push((record.clone(), relay_url.clone()));
|
||||
}
|
||||
let probes: Vec<_> = stream::iter(jobs)
|
||||
.map(|(record, requested)| {
|
||||
@@ -581,6 +589,159 @@ pub async fn reconcile_managed_agent_runtimes(
|
||||
.map_err(|e| format!("spawn_blocking failed: {e}"))
|
||||
}
|
||||
|
||||
/// A single entry in the drain journal: enough to restart the process if
|
||||
/// compensation is needed after a partial drain failure.
|
||||
#[derive(Debug, Clone)]
|
||||
pub(crate) struct DrainJournalEntry {
|
||||
pub key: ManagedAgentRuntimeKey,
|
||||
/// Whether the agent would auto-start on app launch (used to determine
|
||||
/// whether compensation should restart it as auto-start or lazy).
|
||||
pub start_on_app_launch: bool,
|
||||
}
|
||||
|
||||
/// Drain all live runtimes from the runtime map and return a drain journal
|
||||
/// (keys + restart recipes) for use by compensation.
|
||||
///
|
||||
/// This runs under the `managed_agent_runtime_transition` lock (Layer 2
|
||||
/// synchronous epoch — no `.await`). Callers are responsible for acquiring
|
||||
/// that lock before calling this function.
|
||||
///
|
||||
/// Returns `(stopped, remaining, first_stop_error)`. `stopped` contains the
|
||||
/// entries that were successfully killed (compensation restores these).
|
||||
/// `remaining` contains entries that were NOT attempted (due to early-exit on
|
||||
/// first failure). On success `remaining` is empty.
|
||||
pub(crate) fn drain_scope_runtimes(
|
||||
app: &AppHandle,
|
||||
state: &AppState,
|
||||
) -> (
|
||||
Vec<DrainJournalEntry>,
|
||||
Vec<DrainJournalEntry>,
|
||||
Option<String>,
|
||||
) {
|
||||
// Snapshot the journal from the live runtime map before any stops.
|
||||
let journal: Vec<DrainJournalEntry> = {
|
||||
let runtimes = match state.managed_agent_processes.lock() {
|
||||
Ok(r) => r,
|
||||
Err(e) => {
|
||||
return (
|
||||
vec![],
|
||||
vec![],
|
||||
Some(format!("runtime map lock poisoned: {e}")),
|
||||
)
|
||||
}
|
||||
};
|
||||
runtimes
|
||||
.iter()
|
||||
.map(|(key, _runtime)| {
|
||||
// Look up start_on_app_launch from the current store; if we
|
||||
// can't read it, assume true (safer for compensation — we'd
|
||||
// rather restart too many than too few).
|
||||
let start_on_app_launch = load_managed_agents(app)
|
||||
.ok()
|
||||
.and_then(|records| {
|
||||
records
|
||||
.iter()
|
||||
.find(|r| r.pubkey == key.pubkey)
|
||||
.map(|r| r.start_on_app_launch)
|
||||
})
|
||||
.unwrap_or(true);
|
||||
DrainJournalEntry {
|
||||
key: key.clone(),
|
||||
start_on_app_launch,
|
||||
}
|
||||
})
|
||||
.collect()
|
||||
};
|
||||
|
||||
let mut stopped: Vec<DrainJournalEntry> = Vec::new();
|
||||
let mut first_error: Option<String> = None;
|
||||
|
||||
for (idx, entry) in journal.iter().enumerate() {
|
||||
let key = &entry.key;
|
||||
let stop_result = {
|
||||
let mut runtimes = match state.managed_agent_processes.lock() {
|
||||
Ok(r) => r,
|
||||
Err(e) => {
|
||||
first_error.get_or_insert_with(|| {
|
||||
format!("runtime map lock poisoned during drain: {e}")
|
||||
});
|
||||
// Return remaining as the un-attempted tail.
|
||||
return (stopped, journal[idx..].to_vec(), first_error);
|
||||
}
|
||||
};
|
||||
if let Some(mut runtime) = runtimes.remove(key) {
|
||||
let kill_result = if super::process_is_running(runtime.child.id()) {
|
||||
super::terminate_process(runtime.child.id())
|
||||
} else {
|
||||
Ok(())
|
||||
}
|
||||
.and_then(|()| runtime.child.wait().map_err(|e| e.to_string()));
|
||||
|
||||
match kill_result {
|
||||
Ok(_) => {
|
||||
// Remove the receipt so sweep/restore see a clean slate.
|
||||
super::remove_agent_runtime_receipt(app, key);
|
||||
state.clear_agent_session_cache(key);
|
||||
Ok(())
|
||||
}
|
||||
Err(e) => {
|
||||
// Put it back so the map is consistent.
|
||||
runtimes.insert(key.clone(), runtime);
|
||||
Err(e)
|
||||
}
|
||||
}
|
||||
} else {
|
||||
// Nothing live at this key — treat as already stopped.
|
||||
Ok(())
|
||||
}
|
||||
};
|
||||
|
||||
match stop_result {
|
||||
Ok(()) => stopped.push(entry.clone()),
|
||||
Err(e) => {
|
||||
let msg = format!("failed to stop agent {}@{}: {e}", key.pubkey, key.relay_url);
|
||||
first_error.get_or_insert(msg);
|
||||
// Return the un-attempted tail (idx+1 onward) as remaining.
|
||||
return (stopped, journal[idx + 1..].to_vec(), first_error);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
(stopped, vec![], first_error)
|
||||
}
|
||||
|
||||
/// Compensate a partial drain by restarting the entries that were successfully
|
||||
/// stopped before the failure.
|
||||
///
|
||||
/// `stopped` is the slice of journal entries that were actually stopped (i.e.,
|
||||
/// the prefix of the journal up to the first failure). We restart them so the
|
||||
/// old workspace is as intact as possible.
|
||||
///
|
||||
/// Returns a degradation message describing what could not be restarted.
|
||||
pub(crate) fn compensate_drain(app: &AppHandle, stopped: &[DrainJournalEntry]) -> Option<String> {
|
||||
let mut failed_restarts = Vec::new();
|
||||
for entry in stopped {
|
||||
let result = start_pair(
|
||||
entry.key.pubkey.clone(),
|
||||
entry.key.relay_url.clone(),
|
||||
entry.start_on_app_launch,
|
||||
None,
|
||||
app.clone(),
|
||||
);
|
||||
if let Err(e) = result {
|
||||
failed_restarts.push(format!("{}@{}: {e}", entry.key.pubkey, entry.key.relay_url));
|
||||
}
|
||||
}
|
||||
if failed_restarts.is_empty() {
|
||||
None
|
||||
} else {
|
||||
Some(format!(
|
||||
"workspace drain compensation failed for: {}",
|
||||
failed_restarts.join(", ")
|
||||
))
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
@@ -50,6 +50,12 @@ pub struct ManagedAgentPairRuntime {
|
||||
/// Unpredictable identity for this exact harness generation. Lifecycle
|
||||
/// frames from prior processes are rejected even when the pair is live.
|
||||
pub start_nonce: String,
|
||||
/// Scope ID of the workspace this runtime was spawned into. Used by drain
|
||||
/// filtering and `list_managed_agent_runtimes` to detect cross-scope
|
||||
/// entries (the seam that option 2 background-runtime pinning would build
|
||||
/// on). Under active-scope-only policy, all live entries should always
|
||||
/// match the current scope; this field makes the invariant testable.
|
||||
pub scope_id: Option<String>,
|
||||
}
|
||||
|
||||
impl std::ops::Deref for ManagedAgentPairRuntime {
|
||||
@@ -67,13 +73,14 @@ impl std::ops::DerefMut for ManagedAgentPairRuntime {
|
||||
}
|
||||
|
||||
impl ManagedAgentPairRuntime {
|
||||
pub fn starting(process: ManagedAgentProcess) -> Self {
|
||||
pub fn starting(process: ManagedAgentProcess, scope_id: Option<String>) -> Self {
|
||||
let start_nonce = process.start_nonce.clone();
|
||||
Self {
|
||||
process,
|
||||
lifecycle: ManagedAgentRuntimeLifecycle::Starting,
|
||||
error: None,
|
||||
start_nonce,
|
||||
scope_id,
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -104,12 +111,6 @@ pub struct ManagedAgentRuntimeLifecycleObserverPayload {
|
||||
pub error: Option<String>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Deserialize)]
|
||||
#[serde(rename_all = "camelCase")]
|
||||
pub struct ManagedAgentCommunityTarget {
|
||||
pub relay_url: String,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
|
||||
#[serde(rename_all = "camelCase")]
|
||||
pub struct ManagedAgentRuntimeReceipt {
|
||||
|
||||
@@ -265,15 +265,29 @@ pub(crate) async fn recover_stale_mesh_runtime(
|
||||
}
|
||||
|
||||
/// Post-launch recovery for actively running relay-mesh agents.
|
||||
///
|
||||
/// Captures one active scope at function entry for the entire recovery pass.
|
||||
/// A live runtime is only treated as healthy when its bound relay matches the
|
||||
/// captured scope's relay — a mismatched runtime (from a switched-away scope)
|
||||
/// is treated as absent and re-arming proceeds for the current scope.
|
||||
pub(crate) async fn rearm_relay_mesh_for_running_agents(app: &AppHandle) -> Result<(), String> {
|
||||
let state = app.state::<AppState>();
|
||||
let _rearm_guard = state.mesh_recovery.rearm_lock.lock().await;
|
||||
let runtime_mode = state
|
||||
.mesh_llm_runtime
|
||||
.lock()
|
||||
.await
|
||||
.as_ref()
|
||||
.map(|runtime| runtime.mode());
|
||||
|
||||
// Capture scope once for the entire pass; a concurrent workspace switch
|
||||
// that commits after this point is handled on the next watchdog cycle.
|
||||
let scope_relay = state
|
||||
.capture_active_scope()
|
||||
.map(|scope| scope.relay_url.clone());
|
||||
|
||||
let (runtime_mode, runtime_relay) = {
|
||||
let guard = state.mesh_llm_runtime.lock().await;
|
||||
let mode = guard.as_ref().map(|r| r.mode());
|
||||
let relay = guard
|
||||
.as_ref()
|
||||
.and_then(|r| r.start_request().relay_url.clone());
|
||||
(mode, relay)
|
||||
};
|
||||
let recovery = recover_stale_mesh_runtime(&state, MeshRecoveryUrgency::Watchdog).await;
|
||||
let active_pubkeys = active_managed_agent_pubkeys(&state);
|
||||
// Mesh participation is resolved through the same definition-authoritative
|
||||
@@ -282,10 +296,37 @@ pub(crate) async fn rearm_relay_mesh_for_running_agents(app: &AppHandle) -> Resu
|
||||
let personas = crate::managed_agents::load_personas(app).unwrap_or_default();
|
||||
let global = crate::managed_agents::load_global_agent_config(app).unwrap_or_default();
|
||||
|
||||
// Helper: does the live runtime's relay match the active scope relay?
|
||||
// When either is None we treat it as a mismatch (fail closed).
|
||||
let runtime_relay_matches_scope = || -> bool {
|
||||
let Some(scope_r) = scope_relay.as_deref() else {
|
||||
return false;
|
||||
};
|
||||
let Some(runtime_r) = runtime_relay.as_deref() else {
|
||||
return false;
|
||||
};
|
||||
crate::managed_agents::scope::normalize_relay_for_scope(runtime_r)
|
||||
== crate::managed_agents::scope::normalize_relay_for_scope(scope_r)
|
||||
};
|
||||
|
||||
match recovery {
|
||||
MeshRuntimeRecovery::Live
|
||||
| MeshRuntimeRecovery::Debouncing
|
||||
| MeshRuntimeRecovery::Replaced => return Ok(()),
|
||||
MeshRuntimeRecovery::Live => {
|
||||
// Only trust a live runtime whose relay matches the active scope.
|
||||
// A mismatched live runtime (stale from a switched-away scope) is
|
||||
// not healthy for the current scope — fall through to re-arm.
|
||||
if runtime_relay_matches_scope() {
|
||||
return Ok(());
|
||||
}
|
||||
// Mismatch: Serve-mode runtimes stay pinned (machine-level) and
|
||||
// are never bounced by the watchdog — just skip this pass.
|
||||
// Client-mode mismatch: let the loop below attempt re-arm; it
|
||||
// will find the mismatch via ensure_relay_mesh_for_record and
|
||||
// produce the appropriate error or start a new client.
|
||||
if runtime_mode == Some(crate::mesh_llm::MeshNodeMode::Serve) {
|
||||
return Ok(());
|
||||
}
|
||||
}
|
||||
MeshRuntimeRecovery::Debouncing | MeshRuntimeRecovery::Replaced => return Ok(()),
|
||||
MeshRuntimeRecovery::RestartRequired => {
|
||||
if runtime_mode == Some(crate::mesh_llm::MeshNodeMode::Serve) {
|
||||
eprintln!(
|
||||
|
||||
@@ -123,7 +123,9 @@ export function AppShell() {
|
||||
const mainInsetRef = React.useRef<HTMLElement>(null);
|
||||
const location = useLocation();
|
||||
const queryClient = useQueryClient();
|
||||
useManagedAgentRuntimeReconciliation(communitiesHook.communities); // sync storage snapshot
|
||||
useManagedAgentRuntimeReconciliation(
|
||||
`${communitiesHook.activeCommunity?.id ?? "none"}-${communitiesHook.reinitKey}`,
|
||||
); // re-runs on workspace switch
|
||||
const {
|
||||
goAgents,
|
||||
goChannel,
|
||||
|
||||
@@ -68,13 +68,17 @@ export function cacheReconciledManagedAgentRuntimes(
|
||||
}
|
||||
|
||||
/**
|
||||
* Bootstrap runtime pairs in every configured community (fire-and-forget).
|
||||
* Bootstrap runtime pairs for all auto-start agents in the active workspace
|
||||
* (fire-and-forget).
|
||||
*
|
||||
* Called after an agent create: the create command spawns only the active
|
||||
* community's pair, and the startup reconcile won't run again until the next
|
||||
* launch or community switch, so without this kick a brand-new agent stays
|
||||
* deaf in every other community. Idempotent — live pairs are skipped and
|
||||
* deaf in the active workspace. Idempotent — live pairs are skipped and
|
||||
* missing ones spawn lazily (warm socket, no LLM until first mention).
|
||||
*
|
||||
* The backend derives the sole target relay from the captured active scope;
|
||||
* the frontend no longer passes a communities list.
|
||||
*/
|
||||
export function bootstrapManagedAgentRuntimePairs(
|
||||
queryClient: QueryClient,
|
||||
@@ -82,10 +86,7 @@ export function bootstrapManagedAgentRuntimePairs(
|
||||
const baseline = queryClient.getQueryData<ManagedAgentRuntimeStatus[]>(
|
||||
managedAgentRuntimesQueryKey,
|
||||
);
|
||||
const communities = loadCommunities().map((community) => ({
|
||||
relayUrl: community.relayUrl,
|
||||
}));
|
||||
void reconcileManagedAgentRuntimes(communities)
|
||||
void reconcileManagedAgentRuntimes()
|
||||
.then((runtimes) => {
|
||||
cacheReconciledManagedAgentRuntimes(queryClient, baseline, runtimes);
|
||||
})
|
||||
|
||||
@@ -1,48 +1,42 @@
|
||||
import { useQueryClient } from "@tanstack/react-query";
|
||||
import * as React from "react";
|
||||
|
||||
import {
|
||||
canonicalCommunityRelays,
|
||||
classifyReconcileResult,
|
||||
pendingReconcileRelays,
|
||||
reconcileRetryDelayMs,
|
||||
} from "@/features/agents/managedAgentReconciliationPlan";
|
||||
import { reconcileRetryDelayMs } from "@/features/agents/managedAgentReconciliationPlan";
|
||||
import {
|
||||
cacheReconciledManagedAgentRuntimes,
|
||||
managedAgentRuntimesQueryKey,
|
||||
} from "@/features/agents/managedAgentRuntimeHooks";
|
||||
import { canonicalRelayUrl } from "@/features/agents/managedAgentRuntimeStatus";
|
||||
import type { ManagedAgentRuntimeStatus } from "@/shared/api/types";
|
||||
import { reconcileManagedAgentRuntimes } from "@/shared/api/tauriManagedAgents";
|
||||
|
||||
/**
|
||||
* Bootstrap a lazy harness pair for every auto-start local agent in every
|
||||
* configured community, incrementally and with retry.
|
||||
* Bootstrap a lazy harness pair for every auto-start local agent in the active
|
||||
* workspace, with retry on failure.
|
||||
*
|
||||
* Reconciliation is keyed by canonical relay URL: each configured relay is
|
||||
* reconciled once it appears (so adding a community mid-session spawns pairs
|
||||
* there without needing the add flow to also switch communities), and a relay
|
||||
* whose reconcile fails is retried with a capped backoff (5s / 30s / 2m) rather
|
||||
* than left un-spawned until the next switch or relaunch. Relays that reconcile
|
||||
* cleanly are never re-hit; once nothing is outstanding, no timer is left
|
||||
* running.
|
||||
* Under the active-scope-only runtime policy the backend derives the sole
|
||||
* target relay from the captured active scope — no community list is passed.
|
||||
* Reconciliation runs once on mount; if it fails it is retried with a capped
|
||||
* backoff (5s / 30s / 2m). Once it succeeds, no timer is left running.
|
||||
*
|
||||
* The `activeCommunityKey` parameter is a stable key that changes whenever the
|
||||
* active workspace changes (e.g. `"${communityId}-${reinitKey}"`). A workspace
|
||||
* switch unmounts/remounts the effect, resetting the reconcile state and
|
||||
* re-running for the new scope.
|
||||
*/
|
||||
export function useManagedAgentRuntimeReconciliation(
|
||||
communities: readonly { relayUrl: string }[],
|
||||
activeCommunityKey: string,
|
||||
): void {
|
||||
const queryClient = useQueryClient();
|
||||
// Canonical relay URLs that have reconciled cleanly — never re-hit.
|
||||
const reconciledRef = React.useRef<Set<string>>(new Set());
|
||||
// Canonical relay URLs with a reconcile call in flight — not re-dispatched.
|
||||
const inFlightRef = React.useRef<Set<string>>(new Set());
|
||||
// Consecutive failures per canonical relay URL, driving the retry backoff.
|
||||
const failuresRef = React.useRef<Map<string, number>>(new Map());
|
||||
const failureCountRef = React.useRef(0);
|
||||
const retryTimerRef = React.useRef<ReturnType<typeof setTimeout> | null>(
|
||||
null,
|
||||
);
|
||||
|
||||
// activeCommunityKey changes on workspace switch, resetting the effect.
|
||||
// biome-ignore lint/correctness/useExhaustiveDependencies: activeCommunityKey is an intentional trigger dependency — the effect must re-run on workspace switch to reset reconcile state for the new scope.
|
||||
React.useEffect(() => {
|
||||
let cancelled = false;
|
||||
failureCountRef.current = 0;
|
||||
|
||||
const clearRetryTimer = () => {
|
||||
if (retryTimerRef.current !== null) {
|
||||
@@ -51,77 +45,35 @@ export function useManagedAgentRuntimeReconciliation(
|
||||
}
|
||||
};
|
||||
|
||||
const scheduleRetry = (failed: readonly string[]) => {
|
||||
// One shared timer fires at the soonest per-relay backoff; every failing
|
||||
// relay is retried together (reconcile is idempotent), so re-hitting a
|
||||
// longer-backoff relay early is harmless.
|
||||
let soonest: number | null = null;
|
||||
for (const relay of failed) {
|
||||
const nextCount = (failuresRef.current.get(relay) ?? 0) + 1;
|
||||
failuresRef.current.set(relay, nextCount);
|
||||
const delay = reconcileRetryDelayMs(nextCount);
|
||||
if (delay !== null && (soonest === null || delay < soonest)) {
|
||||
soonest = delay;
|
||||
}
|
||||
}
|
||||
const scheduleRetry = () => {
|
||||
const nextCount = failureCountRef.current + 1;
|
||||
failureCountRef.current = nextCount;
|
||||
const delay = reconcileRetryDelayMs(nextCount);
|
||||
clearRetryTimer();
|
||||
if (soonest === null) return; // all failing relays hit the retry cap
|
||||
if (delay === null) return; // retry cap exhausted
|
||||
retryTimerRef.current = setTimeout(() => {
|
||||
retryTimerRef.current = null;
|
||||
if (!cancelled) runReconcile();
|
||||
}, soonest);
|
||||
}, delay);
|
||||
};
|
||||
|
||||
const runReconcile = () => {
|
||||
const canonicalToRequested = canonicalCommunityRelays(
|
||||
communities,
|
||||
canonicalRelayUrl,
|
||||
);
|
||||
// Forget bookkeeping for relays that are no longer configured so the sets
|
||||
// stay bounded and re-adding a removed community reconciles it afresh.
|
||||
for (const done of [...reconciledRef.current]) {
|
||||
if (!canonicalToRequested.has(done)) reconciledRef.current.delete(done);
|
||||
}
|
||||
for (const failing of [...failuresRef.current.keys()]) {
|
||||
if (!canonicalToRequested.has(failing)) {
|
||||
failuresRef.current.delete(failing);
|
||||
}
|
||||
}
|
||||
|
||||
const pending = pendingReconcileRelays(
|
||||
canonicalToRequested,
|
||||
reconciledRef.current,
|
||||
inFlightRef.current,
|
||||
);
|
||||
if (pending.length === 0) {
|
||||
clearRetryTimer();
|
||||
return;
|
||||
}
|
||||
|
||||
for (const relay of pending) inFlightRef.current.add(relay);
|
||||
const targets = pending.map((relay) => ({
|
||||
relayUrl: canonicalToRequested.get(relay) as string,
|
||||
}));
|
||||
const baseline = queryClient.getQueryData<ManagedAgentRuntimeStatus[]>(
|
||||
managedAgentRuntimesQueryKey,
|
||||
);
|
||||
|
||||
void reconcileManagedAgentRuntimes(targets)
|
||||
void reconcileManagedAgentRuntimes()
|
||||
.then((runtimes) => {
|
||||
cacheReconciledManagedAgentRuntimes(queryClient, baseline, runtimes);
|
||||
return classifyReconcileResult(pending, runtimes, canonicalRelayUrl);
|
||||
if (!cancelled) {
|
||||
cacheReconciledManagedAgentRuntimes(
|
||||
queryClient,
|
||||
baseline,
|
||||
runtimes,
|
||||
);
|
||||
}
|
||||
})
|
||||
.catch((error) => {
|
||||
console.warn("[managed-agent-runtimes] reconcile failed:", error);
|
||||
return classifyReconcileResult(pending, null, canonicalRelayUrl);
|
||||
})
|
||||
.then(({ succeeded, failed }) => {
|
||||
for (const relay of pending) inFlightRef.current.delete(relay);
|
||||
for (const relay of succeeded) {
|
||||
reconciledRef.current.add(relay);
|
||||
failuresRef.current.delete(relay);
|
||||
}
|
||||
if (!cancelled && failed.length > 0) scheduleRetry(failed);
|
||||
if (!cancelled) scheduleRetry();
|
||||
});
|
||||
};
|
||||
|
||||
@@ -131,5 +83,5 @@ export function useManagedAgentRuntimeReconciliation(
|
||||
cancelled = true;
|
||||
clearRetryTimer();
|
||||
};
|
||||
}, [communities, queryClient]);
|
||||
}, [activeCommunityKey, queryClient]);
|
||||
}
|
||||
|
||||
@@ -89,8 +89,8 @@ export async function putManagedAgentRuntimeLifecycle(
|
||||
});
|
||||
}
|
||||
|
||||
export async function reconcileManagedAgentRuntimes(
|
||||
communities: readonly { relayUrl: string }[],
|
||||
): Promise<ManagedAgentRuntimeStatus[]> {
|
||||
return invokeTauri("reconcile_managed_agent_runtimes", { communities });
|
||||
export async function reconcileManagedAgentRuntimes(): Promise<
|
||||
ManagedAgentRuntimeStatus[]
|
||||
> {
|
||||
return invokeTauri("reconcile_managed_agent_runtimes", {});
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user