mirror of
https://github.com/block/buzz.git
synced 2026-08-18 06:50:31 +02:00
fix(mesh): restart safely across community boundaries
Signed-off-by: Michael Neale <michael.neale@gmail.com>
This commit is contained in:
@@ -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))
|
||||
}
|
||||
|
||||
|
||||
@@ -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::<AppState>();
|
||||
|
||||
// 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
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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<String>),
|
||||
/// 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<String> = 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.
|
||||
|
||||
Reference in New Issue
Block a user