mirror of
https://github.com/block/buzz.git
synced 2026-08-18 06:50:31 +02:00
fix(mesh): keep shared compute serving through roster changes
Signed-off-by: Michael Neale <michael.neale@gmail.com>
This commit is contained in:
@@ -8,10 +8,10 @@
|
||||
//! P2 auto-model chat — agent sends `model: "auto"`; mesh router picks.
|
||||
//! P3 context-fit regression — an oversized output budget (150k tokens)
|
||||
//! must FAIL with the router's context error (proves the router's fit
|
||||
//! gate — the failure mode the 1024 preset cap protects against).
|
||||
//! gate — the failure mode the 4096 preset cap protects against).
|
||||
//! P4 agentic tool use — agent + buzz-dev-mcp writes a file on disk.
|
||||
//! P5 Buzz reply command — agent runs the exact `buzz messages send`
|
||||
//! command used to publish a reply and receives an accepted result.
|
||||
//! P5 Buzz reply command — agent asks the real developer MCP to run the
|
||||
//! exact `buzz messages send` command used to publish a reply.
|
||||
//!
|
||||
//! The serve node is the same `mesh_llm_sdk::serve` path Share-compute uses
|
||||
//! (publish off, mdns, loopback). The agent legs spawn the real
|
||||
@@ -35,6 +35,36 @@ const DEFAULT_MODEL: &str = "unsloth/gemma-4-E4B-it-GGUF:Q4_K_M";
|
||||
const API_PORT: u16 = 19437;
|
||||
const CONSOLE_PORT: u16 = 13231;
|
||||
|
||||
#[derive(Clone)]
|
||||
struct McpServerEnv {
|
||||
name: String,
|
||||
value: String,
|
||||
}
|
||||
|
||||
#[derive(Clone)]
|
||||
struct McpServerSpec {
|
||||
name: String,
|
||||
command: String,
|
||||
env: Vec<McpServerEnv>,
|
||||
}
|
||||
|
||||
impl McpServerSpec {
|
||||
fn new(name: impl Into<String>, command: impl Into<String>) -> Self {
|
||||
Self {
|
||||
name: name.into(),
|
||||
command: command.into(),
|
||||
env: Vec::new(),
|
||||
}
|
||||
}
|
||||
|
||||
fn set_env(&mut self, name: impl Into<String>, value: impl Into<String>) {
|
||||
self.env.push(McpServerEnv {
|
||||
name: name.into(),
|
||||
value: value.into(),
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
fn main() -> anyhow::Result<()> {
|
||||
tokio::runtime::Builder::new_multi_thread()
|
||||
.enable_all()
|
||||
@@ -168,7 +198,7 @@ async fn run() -> anyhow::Result<()> {
|
||||
let prompt = format!(
|
||||
"Use your developer tools to create {marker_name} in the current working directory containing exactly the text BUZZ_OK (no quotes, no newline commentary). Then confirm."
|
||||
);
|
||||
let mcp = vec![("dev".to_string(), repo_bin("buzz-dev-mcp")?)];
|
||||
let mcp = vec![McpServerSpec::new("dev", repo_bin("buzz-dev-mcp")?)];
|
||||
let (r, marker) =
|
||||
agent_chat_with_marker(&base, &served_id, None, &prompt, &mcp, &marker_name).await;
|
||||
let file_ok = std::fs::read_to_string(&marker)
|
||||
@@ -188,16 +218,19 @@ async fn run() -> anyhow::Result<()> {
|
||||
}
|
||||
let _ = std::fs::remove_file(&marker);
|
||||
|
||||
// P5: prove the model follows Buzz's real reply path. The isolated PATH
|
||||
// contains a fake `buzz` executable so this remains deterministic and
|
||||
// cannot publish to a real community. The fake records the exact argv and
|
||||
// returns the same accepted-write shape as buzz-cli.
|
||||
// P5: prove the model follows Buzz's real reply path. buzz-dev-mcp always
|
||||
// prepends its own multicall `buzz` shim to PATH, so a fake `buzz` later on
|
||||
// PATH cannot isolate this test. Instead, inject a fake shell into the real
|
||||
// MCP server. It records the exact command at the execution boundary and
|
||||
// returns the same accepted-write shape as buzz-cli, without contacting a
|
||||
// community.
|
||||
let send_marker = format!("mesh-e2e-buzz-send-{}.txt", std::process::id());
|
||||
let prompt = "Use the developer shell tool to run exactly this command: buzz messages send --channel test-channel --content BUZZ_REPLY_OK. Do not simulate or substitute another command. After it succeeds, confirm briefly.";
|
||||
let (r, marker) =
|
||||
agent_chat_with_fake_buzz(&base, &served_id, prompt, &mcp, &send_marker).await;
|
||||
agent_chat_with_fake_shell(&base, &served_id, prompt, &mcp, &send_marker).await;
|
||||
let recorded = std::fs::read_to_string(&marker).unwrap_or_default();
|
||||
let send_ok = recorded.trim() == "messages send --channel test-channel --content BUZZ_REPLY_OK";
|
||||
let send_ok =
|
||||
recorded.trim() == "buzz messages send --channel test-channel --content BUZZ_REPLY_OK";
|
||||
match r {
|
||||
Ok(text) => record(
|
||||
"P5 Buzz reply command",
|
||||
@@ -260,7 +293,7 @@ async fn agent_chat(
|
||||
model: &str,
|
||||
max_output_tokens: Option<&str>,
|
||||
prompt: &str,
|
||||
mcp_servers: &[(String, String)],
|
||||
mcp_servers: &[McpServerSpec],
|
||||
) -> anyhow::Result<String> {
|
||||
let (result, _) =
|
||||
agent_chat_in_isolated_home(base, model, max_output_tokens, prompt, mcp_servers, None)
|
||||
@@ -273,7 +306,7 @@ async fn agent_chat_with_marker(
|
||||
model: &str,
|
||||
max_output_tokens: Option<&str>,
|
||||
prompt: &str,
|
||||
mcp_servers: &[(String, String)],
|
||||
mcp_servers: &[McpServerSpec],
|
||||
marker_name: &str,
|
||||
) -> (anyhow::Result<String>, std::path::PathBuf) {
|
||||
let (result, home) =
|
||||
@@ -282,11 +315,11 @@ async fn agent_chat_with_marker(
|
||||
(result, home.join(marker_name))
|
||||
}
|
||||
|
||||
async fn agent_chat_with_fake_buzz(
|
||||
async fn agent_chat_with_fake_shell(
|
||||
base: &str,
|
||||
model: &str,
|
||||
prompt: &str,
|
||||
mcp_servers: &[(String, String)],
|
||||
mcp_servers: &[McpServerSpec],
|
||||
marker_name: &str,
|
||||
) -> (anyhow::Result<String>, std::path::PathBuf) {
|
||||
let (result, home) =
|
||||
@@ -300,8 +333,8 @@ async fn agent_chat_in_isolated_home(
|
||||
model: &str,
|
||||
max_output_tokens: Option<&str>,
|
||||
prompt: &str,
|
||||
mcp_servers: &[(String, String)],
|
||||
fake_buzz_marker_name: Option<&str>,
|
||||
mcp_servers: &[McpServerSpec],
|
||||
fake_shell_marker_name: Option<&str>,
|
||||
) -> (anyhow::Result<String>, std::path::PathBuf) {
|
||||
let agent = match repo_bin("buzz-agent") {
|
||||
Ok(agent) => agent,
|
||||
@@ -313,17 +346,22 @@ async fn agent_chat_in_isolated_home(
|
||||
return (Err(error.into()), home);
|
||||
}
|
||||
|
||||
let inherited_path = std::env::var("PATH").unwrap_or_default();
|
||||
let mut child_path = inherited_path.clone();
|
||||
let mut fake_buzz_marker = None;
|
||||
if let Some(marker_name) = fake_buzz_marker_name {
|
||||
match install_fake_buzz(&home) {
|
||||
Ok(bin_dir) => {
|
||||
child_path = format!("{}:{inherited_path}", bin_dir.to_string_lossy());
|
||||
fake_buzz_marker = Some(home.join(marker_name));
|
||||
}
|
||||
let child_path = std::env::var("PATH").unwrap_or_default();
|
||||
let mut effective_mcp_servers = mcp_servers.to_vec();
|
||||
if let Some(marker_name) = fake_shell_marker_name {
|
||||
let fake_shell = match install_fake_shell(&home) {
|
||||
Ok(fake_shell) => fake_shell,
|
||||
Err(error) => return (Err(error), home),
|
||||
}
|
||||
};
|
||||
let marker = home.join(marker_name);
|
||||
let Some(dev_server) = effective_mcp_servers
|
||||
.iter_mut()
|
||||
.find(|server| server.name == "dev")
|
||||
else {
|
||||
return (Err(anyhow::anyhow!("P5 requires the dev MCP server")), home);
|
||||
};
|
||||
dev_server.set_env("BUZZ_SHELL", fake_shell.to_string_lossy());
|
||||
dev_server.set_env("BUZZ_E2E_SEND_MARKER", marker.to_string_lossy());
|
||||
}
|
||||
|
||||
let mut command = Command::new(&agent);
|
||||
@@ -343,9 +381,6 @@ async fn agent_chat_in_isolated_home(
|
||||
.stdin(Stdio::piped())
|
||||
.stdout(Stdio::piped())
|
||||
.stderr(Stdio::null());
|
||||
if let Some(marker) = &fake_buzz_marker {
|
||||
command.env("BUZZ_E2E_SEND_MARKER", marker);
|
||||
}
|
||||
// P3 deliberately overrides the production default to exercise the
|
||||
// router's context-fit rejection. Normal and tool turns leave it unset,
|
||||
// matching the desktop provider path.
|
||||
@@ -357,27 +392,29 @@ async fn agent_chat_in_isolated_home(
|
||||
Err(error) => return (Err(error.into()), home),
|
||||
};
|
||||
|
||||
let result = drive_acp(&mut child, prompt, mcp_servers, &home).await;
|
||||
let result = drive_acp(&mut child, prompt, &effective_mcp_servers, &home).await;
|
||||
let _ = child.kill().await;
|
||||
(result, home)
|
||||
}
|
||||
|
||||
#[cfg(unix)]
|
||||
fn install_fake_buzz(home: &std::path::Path) -> anyhow::Result<std::path::PathBuf> {
|
||||
fn install_fake_shell(home: &std::path::Path) -> anyhow::Result<std::path::PathBuf> {
|
||||
use std::os::unix::fs::PermissionsExt;
|
||||
|
||||
let bin_dir = home.join("bin");
|
||||
std::fs::create_dir_all(&bin_dir)?;
|
||||
let executable = bin_dir.join("buzz");
|
||||
let executable = bin_dir.join("mesh-e2e-shell");
|
||||
std::fs::write(
|
||||
&executable,
|
||||
r#"#!/bin/sh
|
||||
if [ "$1" = "messages" ] && [ "$2" = "send" ]; then
|
||||
printf '%s\n' "$*" > "$BUZZ_E2E_SEND_MARKER"
|
||||
if [ "$1" = "-lc" ]; then
|
||||
printf '%s\n' "$2" > "$BUZZ_E2E_SEND_MARKER"
|
||||
fi
|
||||
if [ "$1" = "-lc" ] && [ "$2" = "buzz messages send --channel test-channel --content BUZZ_REPLY_OK" ]; then
|
||||
printf '%s\n' '{"event_id":"mesh-e2e-event","accepted":true,"message":"ok"}'
|
||||
exit 0
|
||||
fi
|
||||
printf 'unsupported fake buzz invocation: %s\n' "$*" >&2
|
||||
printf 'unsupported fake shell invocation: %s\n' "$*" >&2
|
||||
exit 2
|
||||
"#,
|
||||
)?;
|
||||
@@ -388,14 +425,14 @@ exit 2
|
||||
}
|
||||
|
||||
#[cfg(not(unix))]
|
||||
fn install_fake_buzz(_home: &std::path::Path) -> anyhow::Result<std::path::PathBuf> {
|
||||
fn install_fake_shell(_home: &std::path::Path) -> anyhow::Result<std::path::PathBuf> {
|
||||
anyhow::bail!("the mesh agent hardware harness currently requires Unix")
|
||||
}
|
||||
|
||||
async fn drive_acp(
|
||||
child: &mut Child,
|
||||
prompt: &str,
|
||||
mcp_servers: &[(String, String)],
|
||||
mcp_servers: &[McpServerSpec],
|
||||
cwd: &std::path::Path,
|
||||
) -> anyhow::Result<String> {
|
||||
let mut stdin = child
|
||||
@@ -410,8 +447,18 @@ async fn drive_acp(
|
||||
|
||||
let mcp_json: Vec<serde_json::Value> = mcp_servers
|
||||
.iter()
|
||||
.map(|(name, command)| {
|
||||
serde_json::json!({ "name": name, "command": command, "args": [], "env": [] })
|
||||
.map(|server| {
|
||||
let env: Vec<serde_json::Value> = server
|
||||
.env
|
||||
.iter()
|
||||
.map(|entry| serde_json::json!({ "name": &entry.name, "value": &entry.value }))
|
||||
.collect();
|
||||
serde_json::json!({
|
||||
"name": &server.name,
|
||||
"command": &server.command,
|
||||
"args": [],
|
||||
"env": env,
|
||||
})
|
||||
})
|
||||
.collect();
|
||||
|
||||
|
||||
@@ -5,12 +5,22 @@ use tauri::{AppHandle, Manager, State};
|
||||
|
||||
use crate::{app_state::AppState, mesh_llm, relay};
|
||||
|
||||
#[derive(Debug, serde::Deserialize, serde::Serialize)]
|
||||
#[derive(Clone, Debug, serde::Deserialize, serde::Serialize)]
|
||||
#[serde(rename_all = "camelCase")]
|
||||
struct MeshSharingConfig {
|
||||
enabled: bool,
|
||||
model_id: String,
|
||||
max_vram_gb: Option<u64>,
|
||||
/// Community relay where Share Compute was explicitly enabled. Older
|
||||
/// configs predate community binding and restore against the active relay.
|
||||
#[serde(default)]
|
||||
relay_url: Option<String>,
|
||||
}
|
||||
|
||||
fn disabled_sharing_checkpoint(config: &MeshSharingConfig) -> MeshSharingConfig {
|
||||
let mut checkpoint = config.clone();
|
||||
checkpoint.enabled = false;
|
||||
checkpoint
|
||||
}
|
||||
|
||||
fn mesh_sharing_config_path(app: &AppHandle) -> Result<PathBuf, String> {
|
||||
@@ -91,6 +101,7 @@ fn sharing_config_from_request(
|
||||
enabled: true,
|
||||
model_id: model_id.to_string(),
|
||||
max_vram_gb: request.max_vram_gb,
|
||||
relay_url: request.relay_url.clone(),
|
||||
})
|
||||
}
|
||||
|
||||
@@ -150,8 +161,13 @@ fn advance_mesh_status_cursor(
|
||||
Ok(cursor)
|
||||
}
|
||||
|
||||
async fn query_mesh_discovery_events(state: &AppState) -> Result<Vec<nostr::Event>, String> {
|
||||
let mut events = relay::query_relay(state, &[mesh_llm::relay_membership_filter()]).await?;
|
||||
async fn query_mesh_discovery_events_at(
|
||||
state: &AppState,
|
||||
relay_url: &str,
|
||||
) -> Result<Vec<nostr::Event>, String> {
|
||||
let api_base_url = relay::relay_http_base_url(relay_url);
|
||||
let mut events =
|
||||
relay::query_relay_at(state, &api_base_url, &[mesh_llm::relay_membership_filter()]).await?;
|
||||
let member_pubkeys = mesh_llm::current_member_pubkeys(&events);
|
||||
if member_pubkeys.is_empty() {
|
||||
// Distinguish "relay returned a membership snapshot listing zero
|
||||
@@ -172,7 +188,7 @@ async fn query_mesh_discovery_events(state: &AppState) -> Result<Vec<nostr::Even
|
||||
let mut previous_cursor: Option<(u64, String)> = None;
|
||||
|
||||
loop {
|
||||
let page = relay::query_relay(state, &[status_filter.clone()]).await?;
|
||||
let page = relay::query_relay_at(state, &api_base_url, &[status_filter.clone()]).await?;
|
||||
let done = page.len() < mesh_llm::MESH_STATUS_PAGE_SIZE;
|
||||
if !done {
|
||||
let cursor = advance_mesh_status_cursor(&mut status_filter, &page)?;
|
||||
@@ -188,6 +204,10 @@ async fn query_mesh_discovery_events(state: &AppState) -> Result<Vec<nostr::Even
|
||||
}
|
||||
}
|
||||
|
||||
async fn query_mesh_discovery_events(state: &AppState) -> Result<Vec<nostr::Event>, String> {
|
||||
query_mesh_discovery_events_at(state, &relay::relay_ws_url_with_override(state)).await
|
||||
}
|
||||
|
||||
/// Resolve the admission roster by intersecting member-signed mesh status
|
||||
/// reporters with the current NIP-43 direct-member list.
|
||||
///
|
||||
@@ -201,6 +221,14 @@ pub(crate) async fn resolve_trusted_owner_ids(state: &AppState) -> Result<Vec<St
|
||||
Ok(mesh_llm::owner_ids_from_events(&events))
|
||||
}
|
||||
|
||||
pub(crate) async fn resolve_trusted_owner_ids_at(
|
||||
state: &AppState,
|
||||
relay_url: &str,
|
||||
) -> Result<Vec<String>, String> {
|
||||
let events = query_mesh_discovery_events_at(state, relay_url).await?;
|
||||
Ok(mesh_llm::owner_ids_from_events(&events))
|
||||
}
|
||||
|
||||
/// Resolve the roster for an initial node *start*, failing closed to self-only
|
||||
/// (an empty roster) when the relay query fails. This is safe only at start:
|
||||
/// there is no established allowlist to preserve yet. The periodic
|
||||
@@ -244,10 +272,11 @@ fn buzz_mesh_join_targets(
|
||||
/// Resolve the validated member endpoint this runtime should join to enter the
|
||||
/// existing Buzz community mesh. `Ok(None)` means this machine is the first
|
||||
/// live serving member (or is itself the shared bootstrap contact).
|
||||
pub(crate) async fn resolve_buzz_mesh_join_targets(
|
||||
pub(crate) async fn resolve_buzz_mesh_join_targets_at(
|
||||
state: &AppState,
|
||||
relay_url: &str,
|
||||
) -> Result<Vec<mesh_llm::MeshServeTarget>, String> {
|
||||
let events = query_mesh_discovery_events(state).await?;
|
||||
let events = query_mesh_discovery_events_at(state, relay_url).await?;
|
||||
let self_owner_id = mesh_llm::ensure_owner_identity()
|
||||
.map_err(|error| format!("failed to load mesh owner identity: {error}"))?
|
||||
.owner_id;
|
||||
@@ -261,8 +290,11 @@ pub(crate) async fn resolve_buzz_mesh_join_targets(
|
||||
/// snapshot. A node start used to repeat the full membership + status query
|
||||
/// for each value, making Share Compute startup both slower and more exposed
|
||||
/// to inconsistent snapshots.
|
||||
async fn resolve_buzz_mesh_startup(state: &AppState) -> (Vec<String>, Option<String>) {
|
||||
match query_mesh_discovery_events(state).await {
|
||||
async fn resolve_buzz_mesh_startup_at(
|
||||
state: &AppState,
|
||||
relay_url: &str,
|
||||
) -> (Vec<String>, Option<String>) {
|
||||
match query_mesh_discovery_events_at(state, relay_url).await {
|
||||
Ok(events) => {
|
||||
let trusted_owner_ids = mesh_llm::owner_ids_from_events(&events);
|
||||
let join_token = mesh_llm::ensure_owner_identity()
|
||||
@@ -291,32 +323,52 @@ async fn resolve_buzz_mesh_startup(state: &AppState) -> (Vec<String>, Option<Str
|
||||
}
|
||||
|
||||
pub(crate) async fn restore_mesh_sharing(app: &AppHandle, state: &AppState) -> CmdResult<()> {
|
||||
let Some(config) = load_mesh_sharing_config(app)? else {
|
||||
let Some(mut config) = load_mesh_sharing_config(app)? else {
|
||||
return Ok(());
|
||||
};
|
||||
if !config.enabled || config.model_id.trim().is_empty() {
|
||||
return Ok(());
|
||||
}
|
||||
config.model_id = mesh_llm::canonical_curated_model_id(&config.model_id).to_string();
|
||||
if state.mesh_llm_runtime.lock().await.is_some() {
|
||||
return Ok(());
|
||||
}
|
||||
let (trusted_owner_ids, join_token) = resolve_buzz_mesh_startup(state).await;
|
||||
let relay_url = config
|
||||
.relay_url
|
||||
.clone()
|
||||
.unwrap_or_else(|| relay::relay_ws_url_with_override(state));
|
||||
let (trusted_owner_ids, join_token) = resolve_buzz_mesh_startup_at(state, &relay_url).await;
|
||||
let mut runtime = state.mesh_llm_runtime.lock().await;
|
||||
if runtime.is_some() {
|
||||
return Ok(());
|
||||
}
|
||||
// The enabled config above has already authorized this restore. Disarm the
|
||||
// next launch while startup is in progress; only proven inference below
|
||||
// re-arms it. This covers both primary weights and lazy package layers.
|
||||
save_mesh_sharing_config(app, &disabled_sharing_checkpoint(&config))?;
|
||||
let request = mesh_llm::StartMeshNodeRequest {
|
||||
mode: mesh_llm::MeshNodeMode::Serve,
|
||||
model_id: Some(config.model_id),
|
||||
model_id: Some(config.model_id.clone()),
|
||||
max_vram_gb: config.max_vram_gb,
|
||||
join_token,
|
||||
mesh_name: Some(buzz_mesh_name(state)),
|
||||
mesh_name: Some(buzz_mesh_name_for_relay(&relay_url)),
|
||||
relay_url: Some(relay_url),
|
||||
trusted_owner_ids: Some(trusted_owner_ids),
|
||||
};
|
||||
let started = mesh_llm::DesktopMeshRuntime::start(request)
|
||||
.await
|
||||
.map_err(|error| format!("failed to restore Share Compute: {error:#}"))?;
|
||||
if let Err(error) = wait_for_mesh_inference(&config.model_id).await {
|
||||
let cleanup = started.stop().await;
|
||||
if let Err(cleanup_error) = cleanup {
|
||||
eprintln!(
|
||||
"buzz-mesh: restored node failed inference readiness and cleanup was incomplete: {cleanup_error:#}"
|
||||
);
|
||||
}
|
||||
return Err(format!("failed to restore Share Compute: {error}"));
|
||||
}
|
||||
*runtime = Some(started);
|
||||
save_mesh_sharing_config(app, &config)?;
|
||||
drop(runtime);
|
||||
mesh_llm::publish_current_status_once(app, "restore").await;
|
||||
Ok(())
|
||||
@@ -328,6 +380,11 @@ pub async fn mesh_start_node(
|
||||
state: State<'_, AppState>,
|
||||
mut request: mesh_llm::StartMeshNodeRequest,
|
||||
) -> CmdResult<mesh_llm::MeshNodeStatus> {
|
||||
let relay_url = relay::relay_ws_url_with_override(&state);
|
||||
request.relay_url = Some(relay_url.clone());
|
||||
if let Some(model_id) = request.model_id.as_mut() {
|
||||
*model_id = mesh_llm::canonical_curated_model_id(model_id).to_string();
|
||||
}
|
||||
let sharing_config = if request.mode == mesh_llm::MeshNodeMode::Serve {
|
||||
Some(sharing_config_from_request(&request)?)
|
||||
} else {
|
||||
@@ -362,13 +419,14 @@ pub async fn mesh_start_node(
|
||||
// Frontend requests never carry a roster. Resolve it and the bootstrap
|
||||
// endpoint from one snapshot so UI startup does not repeat relay probes.
|
||||
if request.trusted_owner_ids.is_none() || request.join_token.is_none() {
|
||||
let (trusted_owner_ids, join_token) = resolve_buzz_mesh_startup(&state).await;
|
||||
let (trusted_owner_ids, join_token) =
|
||||
resolve_buzz_mesh_startup_at(&state, &relay_url).await;
|
||||
request.trusted_owner_ids.get_or_insert(trusted_owner_ids);
|
||||
if request.join_token.is_none() {
|
||||
request.join_token = join_token;
|
||||
}
|
||||
}
|
||||
request.mesh_name = Some(buzz_mesh_name(&state));
|
||||
request.mesh_name = Some(buzz_mesh_name_for_relay(&relay_url));
|
||||
let mut runtime = state.mesh_llm_runtime.lock().await;
|
||||
|
||||
let plan = match runtime.as_ref() {
|
||||
@@ -386,6 +444,13 @@ pub async fn mesh_start_node(
|
||||
return Err("mesh node is already running".to_string());
|
||||
}
|
||||
|
||||
if let Some(config) = sharing_config.as_ref() {
|
||||
// Do not arm launch restoration until the exact inference path used by
|
||||
// agents succeeds. Mesh may bind its ports after primary weights load
|
||||
// while package layers are still downloading.
|
||||
save_mesh_sharing_config(&app, &disabled_sharing_checkpoint(config))?;
|
||||
}
|
||||
|
||||
let started = mesh_llm::DesktopMeshRuntime::start(request)
|
||||
.await
|
||||
.map_err(|error| format!("{error:#}"))?;
|
||||
@@ -409,6 +474,21 @@ pub async fn mesh_start_node(
|
||||
));
|
||||
}
|
||||
};
|
||||
if let Some(config) = sharing_config.as_ref() {
|
||||
if let Err(error) = wait_for_mesh_inference(&config.model_id).await {
|
||||
let cleanup = started.stop().await;
|
||||
if let Err(cleanup_error) = &cleanup {
|
||||
eprintln!(
|
||||
"buzz-mesh: started node failed inference readiness and cleanup was incomplete: {cleanup_error:#}"
|
||||
);
|
||||
}
|
||||
drop(runtime);
|
||||
app.request_restart();
|
||||
return Err(format!(
|
||||
"mesh node started but inference never became ready: {error}; Buzz is restarting to guarantee cleanup"
|
||||
));
|
||||
}
|
||||
}
|
||||
*runtime = Some(started);
|
||||
drop(runtime);
|
||||
if let Some(config) = sharing_config.as_ref() {
|
||||
@@ -612,6 +692,7 @@ pub(crate) async fn ensure_client_node_for_model(
|
||||
max_vram_gb: None,
|
||||
join_token: Some(join_token.clone()),
|
||||
mesh_name: Some(buzz_mesh_name(state)),
|
||||
relay_url: Some(relay::relay_ws_url_with_override(state)),
|
||||
trusted_owner_ids: Some(resolve_trusted_owner_ids_or_self_only(state).await),
|
||||
};
|
||||
let mut runtime = state.mesh_llm_runtime.lock().await;
|
||||
@@ -753,6 +834,18 @@ pub(crate) async fn ensure_relay_mesh_for_record(
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// A persisted Share Compute configuration is authoritative about this
|
||||
// machine's role. If no runtime is currently tracked (for example after a
|
||||
// clean process restart), restore the serving node instead of treating an
|
||||
// agent request as permission to replace it with a client node.
|
||||
if load_mesh_sharing_config(app)?
|
||||
.is_some_and(|config| config.enabled && !config.model_id.trim().is_empty())
|
||||
{
|
||||
restore_mesh_sharing(app, &state).await?;
|
||||
return wait_for_mesh_inference(model_id).await;
|
||||
}
|
||||
|
||||
let target = match resolve_mesh_bootstrap_target(&state, model_id).await {
|
||||
Ok(Some(target)) => target,
|
||||
Ok(None) => {
|
||||
@@ -768,15 +861,9 @@ pub(crate) async fn ensure_relay_mesh_for_record(
|
||||
}
|
||||
};
|
||||
|
||||
// Serve→Client re-arm transition (micspiral review #3, intentional-by-design):
|
||||
// if the dead ingress belonged to a *serve* node with running consumer
|
||||
// agents, this re-arms it as a Client (`MeshNodeMode::Client`). That is the
|
||||
// correct/safe recovery here — config-backed serve restoration is
|
||||
// `restore_mesh_sharing`'s job (`MeshNodeMode::Serve`), and
|
||||
// `ensure_client_node_for_model` reuses any live runtime of *either* mode
|
||||
// (the router resolves per-request), so it only cold-starts a Client when
|
||||
// there is genuinely no runtime. Falling back to Client if a serve node
|
||||
// crashed under local pressure is a desirable fail-safe, not a regression.
|
||||
// No serving configuration exists, so this is a genuine consumer-only
|
||||
// start. A configured serving machine is restored above and never reaches
|
||||
// this client fallback.
|
||||
ensure_client_node_for_model(&state, model_id, Some(target.endpoint_addr)).await?;
|
||||
wait_for_mesh_inference(model_id).await
|
||||
}
|
||||
@@ -792,14 +879,17 @@ pub async fn mesh_stop_node(
|
||||
// role under the lock and, when it's a consume session, leave it running
|
||||
// and return its live status unchanged. The frontend also guards this, but
|
||||
// status can be stale between polls, so the backend is authoritative.
|
||||
let taken = {
|
||||
let (taken, bound_relay_url) = {
|
||||
let mut guard = state.mesh_llm_runtime.lock().await;
|
||||
if let Some(runtime) = guard.as_ref() {
|
||||
if !share_stop_should_teardown(runtime.mode()) {
|
||||
return runtime.status().await.map_err(|error| error.to_string());
|
||||
}
|
||||
}
|
||||
guard.take()
|
||||
let bound_relay_url = guard
|
||||
.as_ref()
|
||||
.and_then(|runtime| runtime.start_request().relay_url.clone());
|
||||
(guard.take(), bound_relay_url)
|
||||
};
|
||||
if let Some(runtime) = taken {
|
||||
runtime.stop().await.map_err(|error| error.to_string())?;
|
||||
@@ -810,9 +900,10 @@ pub async fn mesh_stop_node(
|
||||
enabled: false,
|
||||
model_id: String::new(),
|
||||
max_vram_gb: None,
|
||||
relay_url: None,
|
||||
},
|
||||
)?;
|
||||
mesh_llm::publish_stopped_status_once(&app, "stop").await;
|
||||
mesh_llm::publish_stopped_status_once_at(&app, bound_relay_url.as_deref(), "stop").await;
|
||||
Ok(mesh_llm::stopped_status())
|
||||
}
|
||||
|
||||
|
||||
@@ -110,6 +110,50 @@ fn buzz_mesh_name_is_stable_and_does_not_expose_the_relay() {
|
||||
assert!(!first.contains("example"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn sharing_config_keeps_the_community_where_sharing_was_enabled() {
|
||||
let request = mesh_llm::StartMeshNodeRequest {
|
||||
mode: mesh_llm::MeshNodeMode::Serve,
|
||||
model_id: Some("test-model".to_string()),
|
||||
max_vram_gb: Some(24),
|
||||
join_token: None,
|
||||
mesh_name: Some("buzz-community-test".to_string()),
|
||||
relay_url: Some("wss://community.example".to_string()),
|
||||
trusted_owner_ids: Some(Vec::new()),
|
||||
};
|
||||
|
||||
let config = sharing_config_from_request(&request).expect("valid sharing config");
|
||||
assert_eq!(config.relay_url.as_deref(), Some("wss://community.example"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn legacy_sharing_config_without_community_binding_still_loads() {
|
||||
let config: MeshSharingConfig = serde_json::from_value(serde_json::json!({
|
||||
"enabled": true,
|
||||
"modelId": "test-model",
|
||||
"maxVramGb": null
|
||||
}))
|
||||
.expect("legacy sharing config");
|
||||
|
||||
assert_eq!(config.relay_url, None);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn download_checkpoint_disarms_launch_restore_without_losing_selection() {
|
||||
let config = MeshSharingConfig {
|
||||
enabled: true,
|
||||
model_id: "test-model".to_string(),
|
||||
max_vram_gb: Some(24),
|
||||
relay_url: Some("wss://community.example".to_string()),
|
||||
};
|
||||
|
||||
let checkpoint = disabled_sharing_checkpoint(&config);
|
||||
assert!(!checkpoint.enabled);
|
||||
assert_eq!(checkpoint.model_id, config.model_id);
|
||||
assert_eq!(checkpoint.max_vram_gb, config.max_vram_gb);
|
||||
assert_eq!(checkpoint.relay_url, config.relay_url);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn readiness_failure_is_catalog_sync_when_model_never_visible() {
|
||||
assert_eq!(
|
||||
@@ -345,6 +389,7 @@ fn ensure_serve_runtime_serves_other_model() {
|
||||
max_vram_gb: None,
|
||||
join_token: None,
|
||||
mesh_name: None,
|
||||
relay_url: None,
|
||||
trusted_owner_ids: None,
|
||||
})
|
||||
.await
|
||||
|
||||
@@ -10,32 +10,6 @@ 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
|
||||
@@ -239,31 +213,6 @@ pub async fn apply_workspace(
|
||||
|
||||
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.
|
||||
@@ -334,54 +283,3 @@ 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
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -19,12 +19,14 @@ use mesh_llm_system::vram::{format_rated_capacity, rated_capacity_gb};
|
||||
/// The large pick is resolved through mesh-llm's remote catalog
|
||||
/// (huggingface.co/datasets/meshllm/catalog), so it does not need to exist in
|
||||
/// the compiled `MODEL_CATALOG`; the entry is synthesized below.
|
||||
const CURATED_LARGE: &str = "gemma-4-26B-A4B-it-UD-Q4_K_M";
|
||||
const CURATED_LARGE: &str = "unsloth/gemma-4-26B-A4B-it-GGUF:UD-Q4_K_M";
|
||||
const CURATED_LARGE_ALIAS: &str = "gemma-4-26B-A4B-it-UD-Q4_K_M";
|
||||
const CURATED_LARGE_SIZE: &str = "17GB";
|
||||
const CURATED_LARGE_FILE: &str = "gemma-4-26B-A4B-it-UD-Q4_K_M.gguf";
|
||||
const CURATED_LARGE_DESCRIPTION: &str =
|
||||
"Gemma 4 26B MoE (4B active) — Buzz default for 64GB+ machines";
|
||||
const CURATED_SMALL: &str = "Gemma-4-E4B-it-Q4_K_M";
|
||||
const CURATED_SMALL: &str = "unsloth/gemma-4-E4B-it-GGUF:Q4_K_M";
|
||||
const CURATED_SMALL_ALIAS: &str = "Gemma-4-E4B-it-Q4_K_M";
|
||||
/// Rated-capacity boundary between the two curated tiers, in GB (marketing
|
||||
/// capacity — a "64GB" Mac rates as 64 even though usable AI memory is less).
|
||||
const CURATED_LARGE_MIN_RATED_GB: u64 = 64;
|
||||
@@ -37,6 +39,16 @@ fn buzz_recommended_model(rated_gb: Option<u64>) -> &'static str {
|
||||
}
|
||||
}
|
||||
|
||||
/// Convert Buzz's pre-0.74 curated package aliases into the canonical model
|
||||
/// ids advertised and accepted by Mesh's OpenAI ingress.
|
||||
pub(crate) fn canonical_curated_model_id(model_id: &str) -> &str {
|
||||
match model_id.trim() {
|
||||
CURATED_SMALL_ALIAS => CURATED_SMALL,
|
||||
CURATED_LARGE_ALIAS => CURATED_LARGE,
|
||||
other => other,
|
||||
}
|
||||
}
|
||||
|
||||
/// How a model sits inside this machine's usable AI memory.
|
||||
/// Mirrors mesh-llm's private `fit_code_for_size_label` thresholds.
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
|
||||
@@ -146,12 +158,13 @@ fn build_catalog(
|
||||
.filter(|m| !is_draft_only(&m.name))
|
||||
.map(|m| {
|
||||
let size_gb = parse_size_gb(&m.size);
|
||||
let name = canonical_curated_model_id(&m.name).to_string();
|
||||
MeshCatalogEntry {
|
||||
fit: fit_code(size_gb, vram_gb),
|
||||
installed: is_installed(&m.file, &m.name),
|
||||
installed: is_installed(&m.file, &name) || is_installed(&m.file, &m.name),
|
||||
recommended: false,
|
||||
curated: false,
|
||||
name: m.name.clone(),
|
||||
name,
|
||||
size: m.size.clone(),
|
||||
size_gb,
|
||||
description: m.description.clone(),
|
||||
@@ -166,7 +179,8 @@ fn build_catalog(
|
||||
let size_gb = parse_size_gb(CURATED_LARGE_SIZE);
|
||||
entries.push(MeshCatalogEntry {
|
||||
fit: fit_code(size_gb, vram_gb),
|
||||
installed: is_installed(CURATED_LARGE_FILE, CURATED_LARGE),
|
||||
installed: is_installed(CURATED_LARGE_FILE, CURATED_LARGE)
|
||||
|| is_installed(CURATED_LARGE_FILE, CURATED_LARGE_ALIAS),
|
||||
recommended: false,
|
||||
curated: false,
|
||||
name: CURATED_LARGE.to_string(),
|
||||
@@ -256,8 +270,8 @@ mod tests {
|
||||
|
||||
#[test]
|
||||
fn recommendation_follows_buzz_curated_tiers() {
|
||||
assert_eq!(CURATED_SMALL, "Gemma-4-E4B-it-Q4_K_M");
|
||||
assert_eq!(CURATED_LARGE, "gemma-4-26B-A4B-it-UD-Q4_K_M");
|
||||
assert_eq!(CURATED_SMALL, "unsloth/gemma-4-E4B-it-GGUF:Q4_K_M");
|
||||
assert_eq!(CURATED_LARGE, "unsloth/gemma-4-26B-A4B-it-GGUF:UD-Q4_K_M");
|
||||
// 64GB+ rated machines get the large curated pick.
|
||||
let large = build_catalog(None, 64_000_000_000, 64.0, &[]);
|
||||
assert_eq!(large.recommended.as_deref(), Some(CURATED_LARGE));
|
||||
@@ -271,6 +285,22 @@ mod tests {
|
||||
assert_eq!(tiny.recommended.as_deref(), Some(CURATED_SMALL));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn curated_package_aliases_migrate_to_openai_model_ids() {
|
||||
assert_eq!(
|
||||
canonical_curated_model_id(CURATED_SMALL_ALIAS),
|
||||
CURATED_SMALL
|
||||
);
|
||||
assert_eq!(
|
||||
canonical_curated_model_id(CURATED_LARGE_ALIAS),
|
||||
CURATED_LARGE
|
||||
);
|
||||
assert_eq!(
|
||||
canonical_curated_model_id("other/model:Q4"),
|
||||
"other/model:Q4"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn curated_picks_lead_the_catalog() {
|
||||
let catalog = build_catalog(None, 96_000_000_000, 96.0, &[]);
|
||||
|
||||
@@ -132,7 +132,7 @@ pub async fn start_coordinator(app: AppHandle) {
|
||||
/// MeshLLM establishes the encrypted peer transport itself.
|
||||
async fn reconcile_buzz_mesh_join(app: &AppHandle) -> Result<(), String> {
|
||||
let state = app.state::<AppState>();
|
||||
let peer_ids = {
|
||||
let (peer_ids, relay_url) = {
|
||||
let runtime = state.mesh_llm_runtime.lock().await;
|
||||
let Some(runtime) = runtime.as_ref() else {
|
||||
return Ok(());
|
||||
@@ -141,10 +141,16 @@ async fn reconcile_buzz_mesh_join(app: &AppHandle) -> Result<(), String> {
|
||||
.status_report_payload()
|
||||
.await
|
||||
.map_err(|error| error.to_string())?;
|
||||
visible_peer_ids(&payload)
|
||||
let relay_url = runtime
|
||||
.start_request()
|
||||
.relay_url
|
||||
.clone()
|
||||
.unwrap_or_else(|| crate::relay::relay_ws_url_with_override(&state));
|
||||
(visible_peer_ids(&payload), relay_url)
|
||||
};
|
||||
|
||||
let targets = crate::commands::mesh_llm::resolve_buzz_mesh_join_targets(&state).await?;
|
||||
let targets =
|
||||
crate::commands::mesh_llm::resolve_buzz_mesh_join_targets_at(&state, &relay_url).await?;
|
||||
let Some(target) = targets
|
||||
.into_iter()
|
||||
.find(|target| !target_is_visible(target, &peer_ids))
|
||||
@@ -288,7 +294,12 @@ async fn reconcile_roster(
|
||||
// other member on a transient relay blip (the flapping restart loop). Keep
|
||||
// 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 relay_url = current_request
|
||||
.relay_url
|
||||
.as_deref()
|
||||
.map(str::to_owned)
|
||||
.unwrap_or_else(|| crate::relay::relay_ws_url_with_override(&state));
|
||||
let query = crate::commands::mesh_llm::resolve_trusted_owner_ids_at(&state, &relay_url).await;
|
||||
match roster_reconcile_action(current_owners, pending_shrink.as_deref(), query) {
|
||||
RosterReconcileAction::Keep => {
|
||||
*pending_shrink = None;
|
||||
@@ -350,11 +361,15 @@ pub(crate) async fn publish_current_status_once(app: &AppHandle, reason: &str) {
|
||||
}
|
||||
}
|
||||
|
||||
pub(crate) async fn publish_stopped_status_once(app: &AppHandle, reason: &str) {
|
||||
pub(crate) async fn publish_stopped_status_once_at(
|
||||
app: &AppHandle,
|
||||
relay_url: Option<&str>,
|
||||
reason: &str,
|
||||
) {
|
||||
let state = app.state::<AppState>();
|
||||
match tokio::time::timeout(
|
||||
STATUS_PUBLISH_TIMEOUT,
|
||||
publish_stopped_status_for_state(&state),
|
||||
publish_stopped_status_for_state(&state, relay_url),
|
||||
)
|
||||
.await
|
||||
{
|
||||
@@ -369,26 +384,43 @@ pub(crate) async fn publish_stopped_status_once(app: &AppHandle, reason: &str) {
|
||||
async fn publish_current_status_for_state(state: &AppState) -> Result<(), String> {
|
||||
let identity = super::ensure_owner_identity()
|
||||
.map_err(|error| format!("failed to load mesh owner identity: {error}"))?;
|
||||
let mut payload = {
|
||||
let (mut payload, relay_url) = {
|
||||
let runtime = state.mesh_llm_runtime.lock().await;
|
||||
match runtime.as_ref() {
|
||||
Some(runtime) => runtime
|
||||
.status_report_payload()
|
||||
.await
|
||||
.map_err(|error| error.to_string())?,
|
||||
None => stopped_status_payload(&identity),
|
||||
Some(runtime) => {
|
||||
let payload = runtime
|
||||
.status_report_payload()
|
||||
.await
|
||||
.map_err(|error| error.to_string())?;
|
||||
let relay_url = runtime
|
||||
.start_request()
|
||||
.relay_url
|
||||
.clone()
|
||||
.unwrap_or_else(|| crate::relay::relay_ws_url_with_override(state));
|
||||
(payload, relay_url)
|
||||
}
|
||||
None => (
|
||||
stopped_status_payload(&identity),
|
||||
crate::relay::relay_ws_url_with_override(state),
|
||||
),
|
||||
}
|
||||
};
|
||||
bind_payload_to_member(state, &identity, &mut payload)?;
|
||||
publish_status_report(state, payload).await
|
||||
publish_status_report_at(state, &relay_url, payload).await
|
||||
}
|
||||
|
||||
async fn publish_stopped_status_for_state(state: &AppState) -> Result<(), String> {
|
||||
async fn publish_stopped_status_for_state(
|
||||
state: &AppState,
|
||||
relay_url: Option<&str>,
|
||||
) -> Result<(), String> {
|
||||
let identity = super::ensure_owner_identity()
|
||||
.map_err(|error| format!("failed to load mesh owner identity: {error}"))?;
|
||||
let mut payload = stopped_status_payload(&identity);
|
||||
bind_payload_to_member(state, &identity, &mut payload)?;
|
||||
publish_status_report(state, payload).await
|
||||
let relay_url = relay_url
|
||||
.map(str::to_owned)
|
||||
.unwrap_or_else(|| crate::relay::relay_ws_url_with_override(state));
|
||||
publish_status_report_at(state, &relay_url, payload).await
|
||||
}
|
||||
|
||||
fn stopped_status_payload(identity: &super::identity::OwnerIdentity) -> serde_json::Value {
|
||||
@@ -442,13 +474,21 @@ pub(crate) fn build_status_report_event(
|
||||
.tags([d, k]))
|
||||
}
|
||||
|
||||
pub(crate) async fn publish_status_report(
|
||||
async fn publish_status_report_at(
|
||||
state: &AppState,
|
||||
relay_url: &str,
|
||||
payload: serde_json::Value,
|
||||
) -> Result<(), String> {
|
||||
crate::relay::submit_event(build_status_report_event(payload)?, state)
|
||||
.await
|
||||
.map(|_| ())
|
||||
let api_base_url = crate::relay::relay_http_base_url(relay_url);
|
||||
let keys = state.signing_keys()?;
|
||||
crate::relay::submit_event_at_with_keys(
|
||||
build_status_report_event(payload)?,
|
||||
state,
|
||||
&api_base_url,
|
||||
&keys,
|
||||
)
|
||||
.await
|
||||
.map(|_| ())
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
|
||||
@@ -1,7 +1,7 @@
|
||||
use std::collections::BTreeMap;
|
||||
|
||||
mod coordinator;
|
||||
pub(crate) use coordinator::{publish_current_status_once, publish_stopped_status_once};
|
||||
pub(crate) use coordinator::{publish_current_status_once, publish_stopped_status_once_at};
|
||||
pub use coordinator::{start_coordinator, MeshCoordinator, KIND_BUZZ_MESH_MEMBER_STATUS};
|
||||
|
||||
mod discovery;
|
||||
@@ -14,6 +14,7 @@ pub(crate) use discovery::{
|
||||
use discovery::{device_name_from_status, endpoint_id_from_status, enrich_status_payload_identity};
|
||||
|
||||
mod catalog;
|
||||
pub(crate) use catalog::canonical_curated_model_id;
|
||||
pub use catalog::{model_catalog, MeshModelCatalog};
|
||||
|
||||
mod identity;
|
||||
@@ -200,6 +201,11 @@ pub struct StartMeshNodeRequest {
|
||||
/// accepted from the frontend and contains no relay address.
|
||||
#[serde(default, skip_deserializing)]
|
||||
pub mesh_name: Option<String>,
|
||||
/// Relay this runtime's community membership and discovery are bound to.
|
||||
/// Injected by the backend when sharing starts and retained across UI
|
||||
/// workspace switches; moving a share requires an explicit stop/start.
|
||||
#[serde(default, skip_deserializing)]
|
||||
pub relay_url: Option<String>,
|
||||
/// Mesh owner ids admitted to this node (the member roster from
|
||||
/// member-signed discovery notes). `None` = caller did not resolve a roster
|
||||
/// (tests, direct invocations): the node runs without allowlist
|
||||
@@ -308,17 +314,20 @@ pub const MESH_WORKER_STACK_SIZE: usize = 8 * 1024 * 1024;
|
||||
/// before the node starts. Without this the download happens *inside*
|
||||
/// `serve::start()` where the UI can only show a frozen "starting…" state.
|
||||
/// Already-installed models return immediately from the cache scan.
|
||||
async fn ensure_model_downloaded(model: &str) -> anyhow::Result<()> {
|
||||
let model_owned = model.to_string();
|
||||
let installed = tokio::task::spawn_blocking(move || {
|
||||
async fn model_is_installed(model: &str) -> bool {
|
||||
let model_owned = model.replace("@main", "");
|
||||
tokio::task::spawn_blocking(move || {
|
||||
let cache = mesh_llm_node::models::default_huggingface_cache_dir();
|
||||
mesh_llm_node::models::scan_installed_models(cache)
|
||||
.iter()
|
||||
.any(|m| m.model_ref.contains(&model_owned))
|
||||
.any(|m| m.model_ref.replace("@main", "").contains(&model_owned))
|
||||
})
|
||||
.await
|
||||
.unwrap_or(false);
|
||||
if installed {
|
||||
.unwrap_or(false)
|
||||
}
|
||||
|
||||
async fn ensure_model_downloaded(model: &str) -> anyhow::Result<()> {
|
||||
if model_is_installed(model).await {
|
||||
return Ok(());
|
||||
}
|
||||
mesh_llm_host_runtime::models::download_model_ref_with_progress_details(model, true)
|
||||
|
||||
@@ -12,6 +12,7 @@ fn pending_client_runtime(
|
||||
max_vram_gb: None,
|
||||
join_token: Some("initial-token".to_string()),
|
||||
mesh_name: None,
|
||||
relay_url: None,
|
||||
trusted_owner_ids: None,
|
||||
};
|
||||
super::DesktopMeshRuntime {
|
||||
|
||||
@@ -153,6 +153,13 @@ fn should_evict_after_probe(
|
||||
|| consecutive >= DEAD_PROBE_EVICT_THRESHOLD
|
||||
}
|
||||
|
||||
fn requires_process_restart(
|
||||
mode: crate::mesh_llm::MeshNodeMode,
|
||||
startup_in_progress: bool,
|
||||
) -> bool {
|
||||
startup_in_progress || mode == crate::mesh_llm::MeshNodeMode::Serve
|
||||
}
|
||||
|
||||
/// Probe and, when justified, remove one stale runtime. A closed port is
|
||||
/// decisive for a foreground agent start; watchdog and ambiguous/unhealthy
|
||||
/// ports require consecutive failures to avoid restarting on a transient load
|
||||
@@ -161,22 +168,23 @@ pub(crate) async fn recover_stale_mesh_runtime(
|
||||
state: &AppState,
|
||||
urgency: MeshRecoveryUrgency,
|
||||
) -> MeshRuntimeRecovery {
|
||||
let (candidate_id, startup_in_progress) = match state.mesh_llm_runtime.lock().await.as_ref() {
|
||||
Some(runtime) => (runtime.id(), runtime.is_starting().await),
|
||||
None => {
|
||||
state.mesh_recovery.reset_probe_streak();
|
||||
// A cancelled SDK startup can outlive its Buzz-side task briefly
|
||||
// because the embedded runtime runs on its own thread. Never start
|
||||
// a replacement merely because the tracked handle is gone: first
|
||||
// prove the old ingress is either still useful or has released the
|
||||
// port. This closes the port-conflict loop in #2304.
|
||||
return match probe_mesh_ingress().await {
|
||||
MeshIngressProbe::Live => MeshRuntimeRecovery::Live,
|
||||
MeshIngressProbe::PortClosed => MeshRuntimeRecovery::Absent,
|
||||
MeshIngressProbe::Unhealthy => MeshRuntimeRecovery::ReleasePending,
|
||||
};
|
||||
}
|
||||
};
|
||||
let (candidate_id, startup_in_progress, candidate_mode) =
|
||||
match state.mesh_llm_runtime.lock().await.as_ref() {
|
||||
Some(runtime) => (runtime.id(), runtime.is_starting().await, runtime.mode()),
|
||||
None => {
|
||||
state.mesh_recovery.reset_probe_streak();
|
||||
// A cancelled SDK startup can outlive its Buzz-side task briefly
|
||||
// because the embedded runtime runs on its own thread. Never start
|
||||
// a replacement merely because the tracked handle is gone: first
|
||||
// prove the old ingress is either still useful or has released the
|
||||
// port. This closes the port-conflict loop in #2304.
|
||||
return match probe_mesh_ingress().await {
|
||||
MeshIngressProbe::Live => MeshRuntimeRecovery::Live,
|
||||
MeshIngressProbe::PortClosed => MeshRuntimeRecovery::Absent,
|
||||
MeshIngressProbe::Unhealthy => MeshRuntimeRecovery::ReleasePending,
|
||||
};
|
||||
}
|
||||
};
|
||||
let probe = probe_mesh_ingress().await;
|
||||
if probe == MeshIngressProbe::Live {
|
||||
state.mesh_recovery.reset_probe_streak();
|
||||
@@ -196,12 +204,13 @@ pub(crate) async fn recover_stale_mesh_runtime(
|
||||
return MeshRuntimeRecovery::Debouncing;
|
||||
}
|
||||
|
||||
// The pinned SDK does not yield its control handle until the management
|
||||
// API is ready. Dropping its still-pending start future would detach the
|
||||
// embedded runtime thread without sending a shutdown request, so Buzz must
|
||||
// not evict it and race a replacement onto the same ports. A controlled
|
||||
// app relaunch is the only process-owned cleanup boundary in this state.
|
||||
if startup_in_progress {
|
||||
// Never replace a serving runtime in-process. Its native listeners and
|
||||
// model host are process-owned; stopping it here and then cold-starting a
|
||||
// client silently disables Share Compute and can race ports 9337/3131.
|
||||
// Pending client startups have the same ownership problem because the SDK
|
||||
// has not yielded a shutdown handle yet. In both cases, process restart is
|
||||
// the only boundary that preserves the configured role safely.
|
||||
if requires_process_restart(candidate_mode, startup_in_progress) {
|
||||
state.mesh_recovery.reset_probe_streak();
|
||||
return MeshRuntimeRecovery::RestartRequired;
|
||||
}
|
||||
@@ -241,6 +250,12 @@ pub(crate) async fn recover_stale_mesh_runtime(
|
||||
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());
|
||||
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
|
||||
@@ -254,6 +269,13 @@ pub(crate) async fn rearm_relay_mesh_for_running_agents(app: &AppHandle) -> Resu
|
||||
| MeshRuntimeRecovery::Debouncing
|
||||
| MeshRuntimeRecovery::Replaced => return Ok(()),
|
||||
MeshRuntimeRecovery::RestartRequired => {
|
||||
if runtime_mode == Some(crate::mesh_llm::MeshNodeMode::Serve) {
|
||||
eprintln!(
|
||||
"buzz-mesh: serving ingress failed; restarting Buzz to restore Share Compute without changing roles"
|
||||
);
|
||||
app.request_restart();
|
||||
return Ok(());
|
||||
}
|
||||
let records = crate::managed_agents::load_managed_agents(app).unwrap_or_default();
|
||||
if !records.iter().any(|record| {
|
||||
running_relay_mesh_model_id(record, &active_pubkeys, &personas, &global).is_some()
|
||||
@@ -482,6 +504,22 @@ mod tests {
|
||||
assert!(STALE_STOP_TIMEOUT <= Duration::from_secs(15));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn failed_serving_runtime_requires_process_restart_instead_of_client_fallback() {
|
||||
assert!(requires_process_restart(
|
||||
crate::mesh_llm::MeshNodeMode::Serve,
|
||||
false
|
||||
));
|
||||
assert!(requires_process_restart(
|
||||
crate::mesh_llm::MeshNodeMode::Client,
|
||||
true
|
||||
));
|
||||
assert!(!requires_process_restart(
|
||||
crate::mesh_llm::MeshNodeMode::Client,
|
||||
false
|
||||
));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn only_running_relay_mesh_agents_trigger_rearm() {
|
||||
let personas: Vec<crate::managed_agents::AgentDefinition> = Vec::new();
|
||||
|
||||
Reference in New Issue
Block a user