From f1ee9f8cb1cd37050e5b54937e569fe0f42b0638 Mon Sep 17 00:00:00 2001 From: Michael Neale Date: Wed, 29 Jul 2026 17:01:30 +1000 Subject: [PATCH] fix(mesh): restart safely across community boundaries Signed-off-by: Michael Neale --- desktop/src-tauri/src/commands/mesh_llm.rs | 2 +- desktop/src-tauri/src/commands/workspace.rs | 103 ++++++++++++++++++ desktop/src-tauri/src/mesh_llm/coordinator.rs | 81 +++++--------- 3 files changed, 131 insertions(+), 55 deletions(-) diff --git a/desktop/src-tauri/src/commands/mesh_llm.rs b/desktop/src-tauri/src/commands/mesh_llm.rs index 305c54a20..31a4f2b10 100644 --- a/desktop/src-tauri/src/commands/mesh_llm.rs +++ b/desktop/src-tauri/src/commands/mesh_llm.rs @@ -133,7 +133,7 @@ fn buzz_mesh_name_for_relay(relay_url: &str) -> String { format!("buzz-community-{}", &digest[..32]) } -fn buzz_mesh_name(state: &AppState) -> String { +pub(super) fn buzz_mesh_name(state: &AppState) -> String { buzz_mesh_name_for_relay(&relay::relay_ws_url_with_override(state)) } diff --git a/desktop/src-tauri/src/commands/workspace.rs b/desktop/src-tauri/src/commands/workspace.rs index 731a99d9d..9ffbae103 100644 --- a/desktop/src-tauri/src/commands/workspace.rs +++ b/desktop/src-tauri/src/commands/workspace.rs @@ -10,6 +10,32 @@ use crate::managed_agents::{ }; use crate::relay; +#[cfg(feature = "mesh-llm")] +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +enum MeshWorkspaceApplyAction { + Continue, + RestartProcess, +} + +/// Decide whether an embedded MeshLLM runtime belongs to the workspace that +/// has just been applied. A running node is relay-scoped through `mesh_name`; +/// carrying it across a workspace boundary would publish old-community state +/// on the new relay. Rebuilding it also requires a process boundary because +/// MeshLLM's native listeners are not safe to stop and recreate in-process. +#[cfg(feature = "mesh-llm")] +fn mesh_workspace_apply_action( + request: Option<&crate::mesh_llm::StartMeshNodeRequest>, + expected_mesh_name: &str, +) -> MeshWorkspaceApplyAction { + match request { + None => MeshWorkspaceApplyAction::Continue, + Some(request) if request.mesh_name.as_deref() == Some(expected_mesh_name) => { + MeshWorkspaceApplyAction::Continue + } + Some(_) => MeshWorkspaceApplyAction::RestartProcess, + } +} + /// Adopt the pre-scoping global retention database's pending rows into `scope`. /// /// Best-effort: a failure is logged and the boot proceeds. The migration's own @@ -212,6 +238,32 @@ pub async fn apply_workspace( .map_err(|e| format!("spawn_blocking failed: {e}"))??; let state = restore_app.state::(); + + // A running MeshLLM node is bound to the relay-derived community mesh + // name it started with. Never let the old runtime survive a workspace + // switch: besides advertising the wrong community, trying to repair it by + // stopping and starting the embedded native runtime can terminate Buzz. + // The persisted Share Compute config is intentionally left intact, so the + // normal launch restore recreates the same serving model for this workspace. + #[cfg(feature = "mesh-llm")] + { + let expected_mesh_name = crate::commands::mesh_llm::buzz_mesh_name(&state); + let action = { + let runtime = state.mesh_llm_runtime.lock().await; + mesh_workspace_apply_action( + runtime.as_ref().map(|runtime| runtime.start_request()), + &expected_mesh_name, + ) + }; + if action == MeshWorkspaceApplyAction::RestartProcess { + eprintln!( + "buzz-mesh: workspace community changed; restarting Buzz to bind MeshLLM to the selected relay" + ); + restore_app.request_restart(); + return Ok(()); + } + } + // 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. @@ -282,3 +334,54 @@ pub async fn apply_workspace( Ok(()) } + +#[cfg(all(test, feature = "mesh-llm"))] +mod tests { + use super::*; + + fn mesh_request(mesh_name: Option<&str>) -> crate::mesh_llm::StartMeshNodeRequest { + crate::mesh_llm::StartMeshNodeRequest { + mode: crate::mesh_llm::MeshNodeMode::Serve, + model_id: Some("test-model".to_string()), + max_vram_gb: None, + join_token: None, + mesh_name: mesh_name.map(str::to_string), + trusted_owner_ids: Some(Vec::new()), + } + } + + #[test] + fn workspace_apply_continues_without_a_mesh_runtime() { + assert_eq!( + mesh_workspace_apply_action(None, "buzz-community-new"), + MeshWorkspaceApplyAction::Continue + ); + } + + #[test] + fn workspace_apply_keeps_runtime_bound_to_selected_community() { + let request = mesh_request(Some("buzz-community-current")); + assert_eq!( + mesh_workspace_apply_action(Some(&request), "buzz-community-current"), + MeshWorkspaceApplyAction::Continue + ); + } + + #[test] + fn workspace_apply_restarts_for_runtime_bound_to_another_community() { + let request = mesh_request(Some("buzz-community-old")); + assert_eq!( + mesh_workspace_apply_action(Some(&request), "buzz-community-new"), + MeshWorkspaceApplyAction::RestartProcess + ); + } + + #[test] + fn workspace_apply_restarts_when_runtime_binding_is_unknown() { + let request = mesh_request(None); + assert_eq!( + mesh_workspace_apply_action(Some(&request), "buzz-community-new"), + MeshWorkspaceApplyAction::RestartProcess + ); + } +} diff --git a/desktop/src-tauri/src/mesh_llm/coordinator.rs b/desktop/src-tauri/src/mesh_llm/coordinator.rs index 1e279353b..6bc1aedc2 100644 --- a/desktop/src-tauri/src/mesh_llm/coordinator.rs +++ b/desktop/src-tauri/src/mesh_llm/coordinator.rs @@ -201,8 +201,13 @@ fn target_is_visible(target: &crate::mesh_llm::MeshServeTarget, peer_ids: &[Stri enum RosterReconcileAction { /// Keep the running allowlist untouched (no-op, or a failure we ride out). Keep, - /// Restart the node with a freshly resolved roster. - Restart(Vec), + /// Restart Buzz so MeshLLM is rebuilt with a freshly resolved roster. + /// + /// MeshLLM's native listeners are process-owned in practice: stopping and + /// starting the embedded runtime in one process can terminate Buzz or race + /// ports 9337/3131. The process boundary is therefore part of the safety + /// contract, not an implementation detail. + RestartProcess, /// Observed a *shrink* (or empty) once. Hold the current allowlist and /// require the same reduced roster on the next poll before tearing down, /// so a single transient short-read never drops a member mid-inference. @@ -224,9 +229,9 @@ fn roster_shrinks(current: &[String], fresh: &[String]) -> bool { /// Rules: /// - query failed (`Err`) → `Keep` (never de-admit on a relay blip) /// - resolved roster == current → `Keep` (no-op) -/// - grows (only additions) → `Restart` immediately (fast admission) +/// - grows (only additions) → `RestartProcess` immediately (fast admission) /// - shrinks/empties, first observation → `AwaitConfirm` (hold, re-check next poll) -/// - shrinks/empties, confirmed → `Restart` (same reduced roster twice) +/// - shrinks/empties, confirmed → `RestartProcess` (same reduced roster twice) fn roster_reconcile_action( current_owners: &[String], pending_shrink: Option<&[String]>, @@ -248,13 +253,13 @@ fn roster_reconcile_action( // Growth (pure additions) is safe to apply immediately. if !roster_shrinks(current_owners, &fresh) { - return RosterReconcileAction::Restart(fresh); + return RosterReconcileAction::RestartProcess; } // A shrink (including down to empty) must be confirmed across two // consecutive polls with the *same* reduced roster before we tear down. match pending_shrink { - Some(pending) if pending == fresh => RosterReconcileAction::Restart(fresh), + Some(pending) if pending == fresh => RosterReconcileAction::RestartProcess, _ => RosterReconcileAction::AwaitConfirm(fresh), } } @@ -284,7 +289,7 @@ async fn reconcile_roster( // the current allowlist and try again on the next poll. A shrink is held // for one extra poll (hysteresis) so a single short-read never tears down. let query = crate::commands::mesh_llm::resolve_trusted_owner_ids(&state).await; - let fresh = match roster_reconcile_action(current_owners, pending_shrink.as_deref(), query) { + match roster_reconcile_action(current_owners, pending_shrink.as_deref(), query) { RosterReconcileAction::Keep => { *pending_shrink = None; return Ok(()); @@ -294,34 +299,12 @@ async fn reconcile_roster( *pending_shrink = Some(reduced); return Ok(()); } - RosterReconcileAction::Restart(fresh) => { + RosterReconcileAction::RestartProcess => { *pending_shrink = None; - fresh } - }; + } - let mut request = current_request.clone(); - request.trusted_owner_ids = Some(fresh); - // Bootstrap endpoints are live device state, not configuration. The - // endpoint used at the previous start may belong to the member that just - // left or to a device whose iroh identity rotated while offline. Resolve a - // fresh validated peer for this restart; starting isolated is safe because - // the join watcher will converge it when a member next publishes. - request.join_token = match crate::commands::mesh_llm::resolve_buzz_mesh_join_targets(&state) - .await - { - Ok(targets) => targets - .into_iter() - .next() - .map(|target| target.endpoint_addr), - Err(error) => { - eprintln!( - "buzz-mesh: could not refresh bootstrap endpoint for roster restart; starting isolated: {error}" - ); - None - } - }; - let mut guard = state.mesh_llm_runtime.lock().await; + let guard = state.mesh_llm_runtime.lock().await; let startup_pending = match guard.as_ref() { Some(runtime) => runtime.is_starting().await, None => false, @@ -342,24 +325,14 @@ async fn reconcile_roster( // snapshot. return Ok(()); } - let Some(running) = guard.take() else { + if guard.is_none() { return Ok(()); - }; - eprintln!("buzz-mesh: membership roster changed; restarting mesh node with fresh allowlist"); - if let Err(error) = running.stop().await { - drop(guard); - eprintln!( - "buzz-mesh: stopping mesh node for roster restart failed; restarting Buzz instead of racing the occupied ingress: {error}" - ); - app.request_restart(); - return Err(format!( - "mesh node shutdown failed during roster change: {error}" - )); } - let replacement = crate::mesh_llm::DesktopMeshRuntime::start(request) - .await - .map_err(|error| format!("mesh node restart after roster change failed: {error:#}"))?; - *guard = Some(replacement); + drop(guard); + eprintln!( + "buzz-mesh: membership roster changed; restarting Buzz to rebuild MeshLLM with the fresh community allowlist" + ); + app.request_restart(); Ok(()) } @@ -554,11 +527,11 @@ mod tests { // Growth (pure additions) applies immediately — fast admission is fine. #[test] - fn roster_growth_restarts_immediately() { + fn roster_growth_requests_process_restart_immediately() { let current = vec!["owner-a".to_string()]; let fresh = vec!["owner-a".to_string(), "owner-c".to_string()]; - let action = roster_reconcile_action(¤t, None, Ok(fresh.clone())); - assert_eq!(action, RosterReconcileAction::Restart(fresh)); + let action = roster_reconcile_action(¤t, None, Ok(fresh)); + assert_eq!(action, RosterReconcileAction::RestartProcess); } // A shrink is NOT applied on first observation — it must be confirmed. @@ -572,11 +545,11 @@ mod tests { // The same reduced roster on two consecutive polls confirms the shrink. #[test] - fn roster_shrink_restarts_once_confirmed() { + fn roster_shrink_requests_process_restart_once_confirmed() { let current = vec!["owner-a".to_string(), "owner-b".to_string()]; let reduced = vec!["owner-a".to_string()]; let action = roster_reconcile_action(¤t, Some(&reduced), Ok(reduced.clone())); - assert_eq!(action, RosterReconcileAction::Restart(reduced)); + assert_eq!(action, RosterReconcileAction::RestartProcess); } // A shrink that changes between polls is not confirmed — it re-holds with @@ -600,7 +573,7 @@ mod tests { assert_eq!(first, RosterReconcileAction::AwaitConfirm(Vec::new())); let empty: Vec = Vec::new(); let confirmed = roster_reconcile_action(¤t, Some(&empty), Ok(Vec::new())); - assert_eq!(confirmed, RosterReconcileAction::Restart(Vec::new())); + assert_eq!(confirmed, RosterReconcileAction::RestartProcess); } // A shrink followed by recovery to the full roster cancels the teardown.