mirror of
https://github.com/block/buzz.git
synced 2026-08-18 06:50:31 +02:00
feat(mesh): download mesh-llm node on demand instead of compiling it in
The mesh-llm feature previously statically linked the whole mesh node (~52 MB binary weight, +17.5 MB compressed) which kept it off by default behind 'just mesh=1' builds. Replace the in-process embedding with the pattern Buzz already uses for every other capability binary: - mesh_llm/node_install.rs: first use downloads the official pinned mesh-llm release (~30 MB), verifies the published sha256, extracts to a versioned cache dir, and emits progress app events for the UI. - mesh_llm/node_process.rs: spawns 'mesh-llm serve/client --headless' as a supervised child process (own process group, SIGTERM-then-kill teardown) and drives it over its local management API — the same HTTP surface the SDK's embedded handle used internally. - DesktopMeshRuntime keeps its exact API; dial_endpoint_addr respawns with the extra --join token (the standalone node has no join-at- runtime endpoint yet; upstream candidate: POST /api/join). - mesh_node_install_status command so the UI can say 'first start will download the mesh runtime' up front. - mesh-llm-sdk/host-runtime git deps removed entirely; the mesh-llm cargo feature is now a few KB of client code (57.9 MB binary with the feature on vs 57.3 without — was 109.7 MB). Admission/coordination is unchanged and mesh-native: nostr discovery, invite-token admission, iroh transport — private by default (no --publish), same flags as the embedded config used. Ignored e2e test covers the full first-use flow (download, checksum, cache hit, spawn, status + attestation verification, /v1/models, stop): cargo test --features mesh-llm -- --ignored mesh_node_e2e
This commit is contained in:
Generated
+91
-4564
File diff suppressed because it is too large
Load Diff
@@ -18,7 +18,12 @@ crate-type = ["staticlib", "cdylib", "rlib"]
|
||||
|
||||
[features]
|
||||
default = ["system-keyring"]
|
||||
mesh-llm = ["dep:mesh-llm-sdk", "dep:mesh-llm-host-runtime"]
|
||||
# Mesh compute UI + node management. The mesh-llm node itself is NOT
|
||||
# compiled in (that cost ~52 MB) — it is downloaded on first use as a
|
||||
# checksum-verified release binary and run as a supervised child process
|
||||
# (see mesh_llm/node_install.rs). This feature is therefore only a few KB
|
||||
# of client code and stays cheap to enable everywhere.
|
||||
mesh-llm = []
|
||||
# OS keyring backing for desktop secret storage (nsec private keys). When
|
||||
# disabled, secrets fall back to 0o600 files. On by default for real builds.
|
||||
system-keyring = ["dep:keyring"]
|
||||
@@ -65,7 +70,7 @@ tauri-plugin-updater = "2"
|
||||
tauri-plugin-process = "2"
|
||||
infer = "0.19"
|
||||
hex = "0.4"
|
||||
tokio = { version = "1", features = ["fs", "sync", "rt", "macros", "time"] }
|
||||
tokio = { version = "1", features = ["fs", "sync", "rt", "macros", "time", "process"] }
|
||||
tokio-tungstenite = { version = "0.29", features = ["rustls-tls-webpki-roots"] }
|
||||
tokio-util = { version = "0.7", features = ["rt"] }
|
||||
bytes = "1"
|
||||
@@ -84,11 +89,10 @@ buzz_core_pkg = { package = "buzz-core", path = "../../crates/buzz-core" }
|
||||
buzz_persona_pkg = { package = "buzz-persona", path = "../../crates/buzz-persona" }
|
||||
buzz_sdk_pkg = { package = "buzz-sdk", path = "../../crates/buzz-sdk" }
|
||||
buzz_agent_pkg = { package = "buzz-agent", path = "../../crates/buzz-agent" }
|
||||
mesh-llm-sdk = { git = "https://github.com/Mesh-LLM/mesh-llm.git", rev = "49cf03427ad86becf8dc0596b738c055d166e7e2", package = "mesh-llm-sdk", default-features = false, features = ["client", "serving"], optional = true }
|
||||
mesh-llm-host-runtime = { git = "https://github.com/Mesh-LLM/mesh-llm.git", rev = "49cf03427ad86becf8dc0596b738c055d166e7e2", package = "mesh-llm-host-runtime", default-features = false, features = ["dynamic-native-runtime"], optional = true }
|
||||
base64 = "0.22"
|
||||
sha2 = "0.11"
|
||||
tar = "0.4"
|
||||
flate2 = "1"
|
||||
bzip2 = "0.6"
|
||||
chrono = { version = "0.4", features = ["serde"] }
|
||||
tauri-plugin-global-shortcut = "2"
|
||||
|
||||
@@ -28,7 +28,9 @@ pub async fn mesh_start_node(
|
||||
return Err("mesh node is already running".to_string());
|
||||
}
|
||||
|
||||
let started = mesh_llm::DesktopMeshRuntime::start(request)
|
||||
// First use downloads the mesh-llm node binary; pass the app handle so
|
||||
// install progress reaches the UI as `mesh-node-install-progress` events.
|
||||
let started = mesh_llm::DesktopMeshRuntime::start_with_app(request, Some(&app))
|
||||
.await
|
||||
.map_err(|error| error.to_string())?;
|
||||
let status = started
|
||||
@@ -438,6 +440,17 @@ pub fn mesh_agent_preset(
|
||||
mesh_llm::agent_preset(request)
|
||||
}
|
||||
|
||||
/// Whether the mesh-llm node binary is installed, and at which pinned
|
||||
/// version. Lets the UI say "first start will download the mesh runtime
|
||||
/// (~30 MB)" instead of surprising the user.
|
||||
#[tauri::command]
|
||||
pub fn mesh_node_install_status() -> CmdResult<serde_json::Value> {
|
||||
Ok(serde_json::json!({
|
||||
"installed": mesh_llm::node_installed(),
|
||||
"version": mesh_llm::MESH_NODE_VERSION,
|
||||
}))
|
||||
}
|
||||
|
||||
#[cfg(all(test, feature = "mesh-llm"))]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
@@ -540,6 +540,7 @@ pub fn run() {
|
||||
mesh_node_status,
|
||||
mesh_installed_models,
|
||||
mesh_agent_preset,
|
||||
mesh_node_install_status,
|
||||
update_managed_agent,
|
||||
discover_backend_providers,
|
||||
probe_backend_provider,
|
||||
|
||||
@@ -13,7 +13,12 @@ use discovery::{device_name_from_status, endpoint_id_from_status, enrich_status_
|
||||
mod preset;
|
||||
pub use preset::{agent_preset, MeshAgentPreset, MeshAgentPresetRequest};
|
||||
|
||||
use mesh_llm_sdk::{client, serve, EmbeddedNodeHandle, MeshDiscoveryMode};
|
||||
mod node_install;
|
||||
pub use node_install::{ensure_node_installed, node_installed, MESH_NODE_VERSION};
|
||||
|
||||
pub(crate) mod node_process;
|
||||
use node_process::{NodeProcess, NodeSpawnConfig, NodeStatus};
|
||||
|
||||
use serde::{Deserialize, Serialize};
|
||||
|
||||
const DEFAULT_MESH_API_PORT: u16 = 9337;
|
||||
@@ -197,83 +202,52 @@ pub fn stopped_status() -> MeshNodeStatus {
|
||||
}
|
||||
|
||||
pub struct DesktopMeshRuntime {
|
||||
handle: EmbeddedNodeHandle,
|
||||
/// The spawned node plus the config it was spawned with. Interior
|
||||
/// mutability because `dial_endpoint_addr` joins a mesh by respawning
|
||||
/// the node with an extra `--join` token (the standalone node has no
|
||||
/// join-at-runtime management endpoint yet) while callers hold `&self`.
|
||||
node: tokio::sync::Mutex<(NodeProcess, NodeSpawnConfig)>,
|
||||
mode: MeshNodeMode,
|
||||
model_id: Option<String>,
|
||||
model_name: Option<String>,
|
||||
}
|
||||
|
||||
async fn initialize_mesh_native_runtime() -> anyhow::Result<()> {
|
||||
let cache = mesh_llm_sdk::native_runtime::native_runtime_cache(None)?;
|
||||
let installed = cache.installed()?;
|
||||
let current = mesh_llm_sdk::native_runtime::CURRENT_MESH_VERSION;
|
||||
if !installed
|
||||
.iter()
|
||||
.any(|runtime| runtime.mesh_version == current)
|
||||
{
|
||||
anyhow::bail!(
|
||||
"mesh native runtime for MeshLLM {current} is not installed; run `just mesh=1 staging` or `just mesh-e2e-hardware` to prepare it"
|
||||
);
|
||||
}
|
||||
// initialize_host_runtime became async upstream (mesh-llm #900s series).
|
||||
mesh_llm_host_runtime::initialize_host_runtime()
|
||||
.await
|
||||
.map_err(|error| {
|
||||
anyhow::anyhow!(
|
||||
"mesh native runtime failed to load; run `just mesh=1 staging` or `just mesh-e2e-hardware` to repair it: {error}"
|
||||
)
|
||||
})
|
||||
}
|
||||
|
||||
impl DesktopMeshRuntime {
|
||||
pub async fn start(request: StartMeshNodeRequest) -> anyhow::Result<Self> {
|
||||
Self::start_with_app(request, None).await
|
||||
}
|
||||
|
||||
/// Start the node, downloading the mesh-llm binary first if this is the
|
||||
/// first use (progress emitted as app events when `app` is provided).
|
||||
pub async fn start_with_app(
|
||||
request: StartMeshNodeRequest,
|
||||
app: Option<&tauri::AppHandle>,
|
||||
) -> anyhow::Result<Self> {
|
||||
validate_no_leak_request(&request)?;
|
||||
initialize_mesh_native_runtime().await?;
|
||||
let binary = ensure_node_installed(app)
|
||||
.await
|
||||
.map_err(|error| anyhow::anyhow!("mesh node install failed: {error}"))?;
|
||||
let model_id = request
|
||||
.model_id
|
||||
.clone()
|
||||
.filter(|value| !value.trim().is_empty());
|
||||
let model_name = model_id.clone();
|
||||
let handle = match request.mode {
|
||||
MeshNodeMode::Serve => {
|
||||
let model = model_id
|
||||
.clone()
|
||||
.ok_or_else(|| anyhow::anyhow!("modelId is required for serve mode"))?;
|
||||
let mut builder = serve::EmbeddedServeConfig::builder()
|
||||
.model(model)
|
||||
.api_port(mesh_api_port()?)
|
||||
.console_port(mesh_console_port()?)
|
||||
.publish(false)
|
||||
.auto_join(false)
|
||||
.disable_iroh_relays(true)
|
||||
.discovery_mode(MeshDiscoveryMode::Nostr)
|
||||
.console_ui(true);
|
||||
if let Some(max_vram_gb) = request.max_vram_gb {
|
||||
builder = builder.max_vram_gb(max_vram_gb as f64);
|
||||
}
|
||||
if let Some(join_token) = request.join_token.as_deref() {
|
||||
builder = builder.join_token(join_token);
|
||||
}
|
||||
serve::start(builder.build()).await?
|
||||
}
|
||||
MeshNodeMode::Client => {
|
||||
let mut builder = client::EmbeddedClientConfig::builder()
|
||||
.api_port(mesh_api_port()?)
|
||||
.console_port(mesh_console_port()?)
|
||||
.publish(false)
|
||||
.auto_join(false)
|
||||
.disable_iroh_relays(true)
|
||||
.discovery_mode(MeshDiscoveryMode::Nostr)
|
||||
.console_ui(true);
|
||||
if let Some(join_token) = request.join_token.as_deref() {
|
||||
builder = builder.join_token(join_token);
|
||||
}
|
||||
client::start(builder.build()).await?
|
||||
}
|
||||
if matches!(request.mode, MeshNodeMode::Serve) && model_id.is_none() {
|
||||
anyhow::bail!("modelId is required for serve mode");
|
||||
}
|
||||
let config = NodeSpawnConfig {
|
||||
binary,
|
||||
serve: matches!(request.mode, MeshNodeMode::Serve),
|
||||
model: model_id.clone(),
|
||||
api_port: mesh_api_port()?,
|
||||
console_port: mesh_console_port()?,
|
||||
max_vram_gb: request.max_vram_gb.map(|v| v as f64),
|
||||
join_tokens: request.join_token.clone().into_iter().collect(),
|
||||
};
|
||||
let node = NodeProcess::spawn(config.clone()).await?;
|
||||
|
||||
Ok(Self {
|
||||
handle,
|
||||
node: tokio::sync::Mutex::new((node, config)),
|
||||
mode: request.mode,
|
||||
model_id,
|
||||
model_name,
|
||||
@@ -281,30 +255,48 @@ impl DesktopMeshRuntime {
|
||||
}
|
||||
|
||||
pub async fn status(&self) -> anyhow::Result<MeshNodeStatus> {
|
||||
let status = self.handle.status().await?;
|
||||
self.status_from_sdk(status)
|
||||
let status = self.node.lock().await.0.status().await?;
|
||||
self.status_from_node(status)
|
||||
}
|
||||
|
||||
pub async fn status_report_payload(&self) -> anyhow::Result<serde_json::Value> {
|
||||
let status = self.handle.status().await?;
|
||||
let status = self.node.lock().await.0.status().await?;
|
||||
let mut payload = status.payload;
|
||||
enrich_status_payload_identity(&mut payload, status.invite_token.as_deref());
|
||||
Ok(payload)
|
||||
}
|
||||
|
||||
/// Join a mesh at `endpoint_addr` (invite token).
|
||||
///
|
||||
/// The standalone node accepts join tokens only at startup (`--join`);
|
||||
/// there is no join-at-runtime management endpoint yet (upstream
|
||||
/// candidate: `POST /api/join`). A repeated token is a no-op fast path;
|
||||
/// a new token respawns the node with the token appended — model state
|
||||
/// is preserved because serve mode reloads its configured model and
|
||||
/// client mode has none.
|
||||
pub async fn dial_endpoint_addr(&self, endpoint_addr: impl Into<String>) -> anyhow::Result<()> {
|
||||
self.handle.join_token(endpoint_addr).await
|
||||
let token = endpoint_addr.into();
|
||||
let mut guard = self.node.lock().await;
|
||||
if guard.1.join_tokens.iter().any(|t| t == &token) {
|
||||
return Ok(());
|
||||
}
|
||||
let mut config = guard.1.clone();
|
||||
config.join_tokens.push(token);
|
||||
// Stop the old node first — same ports — then spawn with the new
|
||||
// token. On spawn failure the lock is released with the node gone;
|
||||
// callers see the error and the next start recreates it.
|
||||
guard.0.take_stop().await.ok();
|
||||
let node = NodeProcess::spawn(config.clone()).await?;
|
||||
*guard = (node, config);
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub async fn installed_models(&self) -> anyhow::Result<Vec<MeshModelOption>> {
|
||||
let status = self.handle.status().await?;
|
||||
let status = self.node.lock().await.0.status().await?;
|
||||
Ok(models_from_status_payload(Some(&status.payload)))
|
||||
}
|
||||
|
||||
fn status_from_sdk(
|
||||
&self,
|
||||
status: mesh_llm_sdk::EmbeddedNodeStatus,
|
||||
) -> anyhow::Result<MeshNodeStatus> {
|
||||
fn status_from_node(&self, status: NodeStatus) -> anyhow::Result<MeshNodeStatus> {
|
||||
let health = health_from_payload(&status.payload);
|
||||
let endpoint_id = endpoint_id_from_status(&status.payload, status.invite_token.as_deref());
|
||||
let device_name = device_name_from_status(&status.payload, endpoint_id.as_deref());
|
||||
@@ -329,7 +321,7 @@ impl DesktopMeshRuntime {
|
||||
}
|
||||
|
||||
pub async fn stop(self) -> anyhow::Result<()> {
|
||||
self.handle.stop().await
|
||||
self.node.into_inner().0.stop().await
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -0,0 +1,336 @@
|
||||
//! Download-on-first-use install of the `mesh-llm` node binary.
|
||||
//!
|
||||
//! Buzz does not compile the mesh node in (that costs ~52 MB of binary);
|
||||
//! instead the official release artifact is fetched the first time a user
|
||||
//! starts a mesh node, sha256-verified against the published sidecar
|
||||
//! checksum, and cached under Buzz's data dir. This mirrors how mesh-llm
|
||||
//! itself treats its native inference runtime (versioned tarball + checksum
|
||||
//! + cache) and how Buzz treats models: pay for the capability when you use
|
||||
//! it, not in the app download.
|
||||
//!
|
||||
//! The binary additionally carries an ed25519 release attestation that the
|
||||
//! node self-verifies and reports via `/api/status` (`release_attestation`),
|
||||
//! which the runtime layer surfaces after spawn.
|
||||
|
||||
use std::path::PathBuf;
|
||||
|
||||
use serde::Serialize;
|
||||
use sha2::Digest as _;
|
||||
use tauri::Emitter as _;
|
||||
|
||||
/// Pinned mesh-llm release. Bump deliberately (with the flag/API surface
|
||||
/// re-checked) rather than tracking "latest" — the spawn flags and
|
||||
/// management API below are validated against this version.
|
||||
pub const MESH_NODE_VERSION: &str = "v0.72.2";
|
||||
|
||||
const RELEASE_BASE: &str = "https://github.com/Mesh-LLM/mesh-llm/releases/download";
|
||||
|
||||
/// Tauri event emitted with download progress so the UI can show
|
||||
/// "downloading mesh runtime…" the same way model downloads do.
|
||||
pub const INSTALL_PROGRESS_EVENT: &str = "mesh-node-install-progress";
|
||||
|
||||
#[derive(Debug, Clone, Serialize)]
|
||||
#[serde(rename_all = "camelCase")]
|
||||
pub struct NodeInstallProgress {
|
||||
pub version: String,
|
||||
pub downloaded_bytes: u64,
|
||||
pub total_bytes: Option<u64>,
|
||||
pub phase: &'static str,
|
||||
}
|
||||
|
||||
/// Release asset name for the current platform, or an explanation of why the
|
||||
/// platform is unsupported.
|
||||
fn release_asset() -> Result<&'static str, String> {
|
||||
match (std::env::consts::OS, std::env::consts::ARCH) {
|
||||
("macos", "aarch64") => Ok("mesh-llm-aarch64-apple-darwin.tar.gz"),
|
||||
("linux", "x86_64") => Ok("mesh-llm-x86_64-unknown-linux-gnu.tar.gz"),
|
||||
("linux", "aarch64") => Ok("mesh-llm-aarch64-unknown-linux-gnu.tar.gz"),
|
||||
("windows", "x86_64") => Ok("mesh-llm-x86_64-pc-windows-msvc.zip"),
|
||||
(os, arch) => Err(format!(
|
||||
"mesh-llm has no release build for {os}/{arch}; see https://github.com/Mesh-LLM/mesh-llm/releases"
|
||||
)),
|
||||
}
|
||||
}
|
||||
|
||||
fn node_binary_name() -> &'static str {
|
||||
if cfg!(windows) {
|
||||
"mesh-llm.exe"
|
||||
} else {
|
||||
"mesh-llm"
|
||||
}
|
||||
}
|
||||
|
||||
/// Versioned cache directory for the node binary.
|
||||
pub fn node_install_dir() -> Result<PathBuf, String> {
|
||||
let base = dirs::data_dir().ok_or("no platform data dir available")?;
|
||||
Ok(base.join("buzz").join("mesh-node").join(MESH_NODE_VERSION))
|
||||
}
|
||||
|
||||
/// Path the installed node binary is expected at (may not exist yet).
|
||||
pub fn installed_node_path() -> Result<PathBuf, String> {
|
||||
Ok(node_install_dir()?.join(node_binary_name()))
|
||||
}
|
||||
|
||||
/// True if the pinned node version is already installed.
|
||||
pub fn node_installed() -> bool {
|
||||
installed_node_path().map(|p| p.is_file()).unwrap_or(false)
|
||||
}
|
||||
|
||||
/// Ensure the pinned mesh-llm node binary is installed, downloading and
|
||||
/// verifying it if needed. Returns the binary path. Progress is emitted as
|
||||
/// [`INSTALL_PROGRESS_EVENT`] app events when an `AppHandle` is provided.
|
||||
pub async fn ensure_node_installed(app: Option<&tauri::AppHandle>) -> Result<PathBuf, String> {
|
||||
let path = installed_node_path()?;
|
||||
if path.is_file() {
|
||||
return Ok(path);
|
||||
}
|
||||
|
||||
let asset = release_asset()?;
|
||||
let url = format!("{RELEASE_BASE}/{MESH_NODE_VERSION}/{asset}");
|
||||
let checksum_url = format!("{url}.sha256");
|
||||
|
||||
let client = reqwest::Client::new();
|
||||
|
||||
// Published checksum first — fail before the big download if missing.
|
||||
let checksum_body = client
|
||||
.get(&checksum_url)
|
||||
.send()
|
||||
.await
|
||||
.map_err(|e| format!("mesh node checksum fetch failed: {e}"))?
|
||||
.error_for_status()
|
||||
.map_err(|e| format!("mesh node checksum fetch failed: {e}"))?
|
||||
.text()
|
||||
.await
|
||||
.map_err(|e| format!("mesh node checksum read failed: {e}"))?;
|
||||
let expected_sha = checksum_body
|
||||
.split_whitespace()
|
||||
.next()
|
||||
.filter(|s| s.len() == 64)
|
||||
.ok_or("mesh node checksum sidecar is malformed")?
|
||||
.to_ascii_lowercase();
|
||||
|
||||
// Stream the archive to a temp file next to the final location (same
|
||||
// filesystem → atomic-ish rename), hashing as we go.
|
||||
let dir = node_install_dir()?;
|
||||
tokio::fs::create_dir_all(&dir)
|
||||
.await
|
||||
.map_err(|e| format!("mesh node cache dir create failed: {e}"))?;
|
||||
let archive_path = dir.join(format!("{asset}.partial"));
|
||||
|
||||
let response = client
|
||||
.get(&url)
|
||||
.send()
|
||||
.await
|
||||
.map_err(|e| format!("mesh node download failed: {e}"))?
|
||||
.error_for_status()
|
||||
.map_err(|e| format!("mesh node download failed: {e}"))?;
|
||||
let total_bytes = response.content_length();
|
||||
|
||||
let mut hasher = sha2::Sha256::new();
|
||||
let mut downloaded: u64 = 0;
|
||||
let mut last_emit: u64 = 0;
|
||||
{
|
||||
let mut file = tokio::fs::File::create(&archive_path)
|
||||
.await
|
||||
.map_err(|e| format!("mesh node archive create failed: {e}"))?;
|
||||
let mut stream = response.bytes_stream();
|
||||
use futures_util::StreamExt as _;
|
||||
use tokio::io::AsyncWriteExt as _;
|
||||
while let Some(chunk) = stream.next().await {
|
||||
let chunk = chunk.map_err(|e| format!("mesh node download failed: {e}"))?;
|
||||
hasher.update(&chunk);
|
||||
downloaded += chunk.len() as u64;
|
||||
file.write_all(&chunk)
|
||||
.await
|
||||
.map_err(|e| format!("mesh node archive write failed: {e}"))?;
|
||||
// Emit at most every 2 MB so the event stream stays light.
|
||||
if downloaded - last_emit >= 2 * 1024 * 1024 {
|
||||
last_emit = downloaded;
|
||||
emit_progress(app, downloaded, total_bytes, "downloading");
|
||||
}
|
||||
}
|
||||
file.flush()
|
||||
.await
|
||||
.map_err(|e| format!("mesh node archive flush failed: {e}"))?;
|
||||
}
|
||||
|
||||
let actual_sha = hex::encode(hasher.finalize());
|
||||
if actual_sha != expected_sha {
|
||||
let _ = tokio::fs::remove_file(&archive_path).await;
|
||||
return Err(format!(
|
||||
"mesh node download checksum mismatch (expected {expected_sha}, got {actual_sha})"
|
||||
));
|
||||
}
|
||||
emit_progress(app, downloaded, total_bytes, "extracting");
|
||||
|
||||
// Extract on a blocking thread (tar/zip are sync APIs).
|
||||
let archive_for_extract = archive_path.clone();
|
||||
let dir_for_extract = dir.clone();
|
||||
let is_zip = asset.ends_with(".zip");
|
||||
tokio::task::spawn_blocking(move || {
|
||||
extract_node_binary(&archive_for_extract, &dir_for_extract, is_zip)
|
||||
})
|
||||
.await
|
||||
.map_err(|e| format!("mesh node extract task failed: {e}"))??;
|
||||
|
||||
let _ = tokio::fs::remove_file(&archive_path).await;
|
||||
if !path.is_file() {
|
||||
return Err("mesh node archive did not contain the mesh-llm binary".to_string());
|
||||
}
|
||||
emit_progress(app, downloaded, total_bytes, "installed");
|
||||
Ok(path)
|
||||
}
|
||||
|
||||
fn emit_progress(
|
||||
app: Option<&tauri::AppHandle>,
|
||||
downloaded_bytes: u64,
|
||||
total_bytes: Option<u64>,
|
||||
phase: &'static str,
|
||||
) {
|
||||
if let Some(app) = app {
|
||||
let _ = app.emit(
|
||||
INSTALL_PROGRESS_EVENT,
|
||||
NodeInstallProgress {
|
||||
version: MESH_NODE_VERSION.to_string(),
|
||||
downloaded_bytes,
|
||||
total_bytes,
|
||||
phase,
|
||||
},
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
/// Pull the `mesh-llm` binary (searched by basename, whatever directory
|
||||
/// layout the archive uses) out of the downloaded archive into `dir`.
|
||||
fn extract_node_binary(
|
||||
archive: &std::path::Path,
|
||||
dir: &std::path::Path,
|
||||
is_zip: bool,
|
||||
) -> Result<(), String> {
|
||||
let wanted = node_binary_name();
|
||||
let dest = dir.join(wanted);
|
||||
if is_zip {
|
||||
let file = std::fs::File::open(archive).map_err(|e| format!("archive open failed: {e}"))?;
|
||||
let mut zip =
|
||||
zip::ZipArchive::new(file).map_err(|e| format!("zip archive read failed: {e}"))?;
|
||||
for index in 0..zip.len() {
|
||||
let mut entry = zip
|
||||
.by_index(index)
|
||||
.map_err(|e| format!("zip entry read failed: {e}"))?;
|
||||
let matches = entry
|
||||
.enclosed_name()
|
||||
.and_then(|p| p.file_name().map(|n| n == wanted))
|
||||
.unwrap_or(false);
|
||||
if matches {
|
||||
let mut out = std::fs::File::create(&dest)
|
||||
.map_err(|e| format!("binary write failed: {e}"))?;
|
||||
std::io::copy(&mut entry, &mut out)
|
||||
.map_err(|e| format!("binary write failed: {e}"))?;
|
||||
return Ok(());
|
||||
}
|
||||
}
|
||||
Err("mesh-llm binary not found in zip archive".to_string())
|
||||
} else {
|
||||
let file = std::fs::File::open(archive).map_err(|e| format!("archive open failed: {e}"))?;
|
||||
let gz = flate2::read::GzDecoder::new(file);
|
||||
let mut tar = tar::Archive::new(gz);
|
||||
for entry in tar
|
||||
.entries()
|
||||
.map_err(|e| format!("tar archive read failed: {e}"))?
|
||||
{
|
||||
let mut entry = entry.map_err(|e| format!("tar entry read failed: {e}"))?;
|
||||
let matches = entry
|
||||
.path()
|
||||
.ok()
|
||||
.and_then(|p| p.file_name().map(|n| n == wanted))
|
||||
.unwrap_or(false);
|
||||
if matches {
|
||||
entry
|
||||
.unpack(&dest)
|
||||
.map_err(|e| format!("binary unpack failed: {e}"))?;
|
||||
#[cfg(unix)]
|
||||
{
|
||||
use std::os::unix::fs::PermissionsExt as _;
|
||||
std::fs::set_permissions(&dest, std::fs::Permissions::from_mode(0o755))
|
||||
.map_err(|e| format!("binary chmod failed: {e}"))?;
|
||||
}
|
||||
return Ok(());
|
||||
}
|
||||
}
|
||||
Err("mesh-llm binary not found in tar archive".to_string())
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn release_asset_known_for_this_platform() {
|
||||
// The dev platforms Buzz builds on must all be mapped.
|
||||
release_asset().expect("current platform should have a mesh-llm release asset");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn install_dir_is_versioned() {
|
||||
let dir = node_install_dir().expect("data dir");
|
||||
assert!(dir.ends_with(std::path::Path::new("mesh-node").join(MESH_NODE_VERSION)));
|
||||
}
|
||||
|
||||
/// Full first-use flow: download the pinned release, verify the
|
||||
/// checksum, extract, spawn a client node, drive its management and
|
||||
/// OpenAI APIs, and stop it. Network + ~30 MB download; run with
|
||||
/// `cargo test --features mesh-llm -- --ignored mesh_node_e2e`.
|
||||
#[tokio::test(flavor = "multi_thread")]
|
||||
#[ignore = "downloads the mesh-llm release binary; network-dependent"]
|
||||
async fn mesh_node_e2e_download_spawn_openai_surface() {
|
||||
let binary = ensure_node_installed(None).await.expect("install");
|
||||
assert!(binary.is_file());
|
||||
assert!(node_installed());
|
||||
|
||||
// Second call is a cache hit (no re-download): must return fast.
|
||||
let started = std::time::Instant::now();
|
||||
ensure_node_installed(None).await.expect("cache hit");
|
||||
assert!(started.elapsed() < std::time::Duration::from_secs(2));
|
||||
|
||||
let node = crate::mesh_llm::node_process::NodeProcess::spawn(
|
||||
crate::mesh_llm::node_process::NodeSpawnConfig {
|
||||
binary,
|
||||
serve: false,
|
||||
model: None,
|
||||
api_port: 29337,
|
||||
console_port: 23131,
|
||||
max_vram_gb: None,
|
||||
join_tokens: Vec::new(),
|
||||
},
|
||||
)
|
||||
.await
|
||||
.expect("spawn client node");
|
||||
|
||||
let status = node.status().await.expect("status");
|
||||
assert_eq!(
|
||||
status.payload.get("node_state").and_then(|v| v.as_str()),
|
||||
Some("client")
|
||||
);
|
||||
assert!(status.invite_token.is_some(), "invite token published");
|
||||
// Release attestation is self-verified by the node.
|
||||
assert_eq!(
|
||||
status
|
||||
.payload
|
||||
.pointer("/release_attestation/verified")
|
||||
.and_then(serde_json::Value::as_bool),
|
||||
Some(true)
|
||||
);
|
||||
|
||||
// OpenAI-compatible surface answers.
|
||||
let models: serde_json::Value = reqwest::get(format!("{}/models", node.api_base_url()))
|
||||
.await
|
||||
.expect("GET /v1/models")
|
||||
.json()
|
||||
.await
|
||||
.expect("models json");
|
||||
assert_eq!(models.get("object").and_then(|v| v.as_str()), Some("list"));
|
||||
|
||||
node.stop().await.expect("stop");
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,171 @@
|
||||
//! Spawned mesh-llm node process + management-API client.
|
||||
//!
|
||||
//! Replaces the in-process `mesh_llm_sdk::serve/client::start()` embedding
|
||||
//! (which statically linked the whole node, ~52 MB) with the same pattern
|
||||
//! Buzz uses for every other capability binary: a supervised child process
|
||||
//! driven over its local HTTP surface. The SDK's own embedded handle already
|
||||
//! spoke HTTP to itself (`GET {console}/api/status`), so this changes where
|
||||
//! the node lives, not how it is driven:
|
||||
//!
|
||||
//! - inference: http://127.0.0.1:{api_port}/v1 (OpenAI-compatible)
|
||||
//! - management: http://127.0.0.1:{console_port}/api/status
|
||||
//!
|
||||
//! Admission/coordination is mesh-native and unchanged: the node does its
|
||||
//! own Nostr discovery, invite-token admission, and iroh transport exactly
|
||||
//! as `mesh-llm serve/client` does standalone.
|
||||
|
||||
use std::path::PathBuf;
|
||||
use std::process::Stdio;
|
||||
use std::time::Duration;
|
||||
|
||||
use tokio::process::{Child, Command};
|
||||
|
||||
/// One spawned mesh-llm node.
|
||||
pub struct NodeProcess {
|
||||
child: Child,
|
||||
api_port: u16,
|
||||
console_port: u16,
|
||||
}
|
||||
|
||||
/// Mirror of the SDK's `EmbeddedNodeStatus`: raw status payload plus the
|
||||
/// derived fields callers use.
|
||||
pub struct NodeStatus {
|
||||
pub api_base_url: String,
|
||||
pub console_url: String,
|
||||
pub invite_token: Option<String>,
|
||||
pub payload: serde_json::Value,
|
||||
}
|
||||
|
||||
#[derive(Clone)]
|
||||
pub struct NodeSpawnConfig {
|
||||
pub binary: PathBuf,
|
||||
/// `serve` (share a model) or `client` (consume only).
|
||||
pub serve: bool,
|
||||
pub model: Option<String>,
|
||||
pub api_port: u16,
|
||||
pub console_port: u16,
|
||||
pub max_vram_gb: Option<f64>,
|
||||
pub join_tokens: Vec<String>,
|
||||
}
|
||||
|
||||
impl NodeProcess {
|
||||
/// Spawn the node and wait for its management API to come up.
|
||||
pub async fn spawn(config: NodeSpawnConfig) -> anyhow::Result<Self> {
|
||||
let mut cmd = Command::new(&config.binary);
|
||||
cmd.arg(if config.serve { "serve" } else { "client" });
|
||||
if let Some(model) = config.model.as_deref() {
|
||||
cmd.arg("--model").arg(model);
|
||||
}
|
||||
cmd.arg("--port").arg(config.api_port.to_string());
|
||||
cmd.arg("--console").arg(config.console_port.to_string());
|
||||
// Headless: management API without the web console UI. Private by
|
||||
// default (no --publish): joinable only via invite token, matching
|
||||
// the embedded config (`publish(false)`, `auto_join(false)`).
|
||||
cmd.arg("--headless");
|
||||
cmd.arg("--disable-iroh-relays");
|
||||
cmd.arg("--mesh-discovery-mode").arg("nostr");
|
||||
cmd.arg("--log-format").arg("json");
|
||||
if let Some(max_vram) = config.max_vram_gb {
|
||||
cmd.arg("--max-vram").arg(max_vram.to_string());
|
||||
}
|
||||
for token in &config.join_tokens {
|
||||
cmd.arg("--join").arg(token);
|
||||
}
|
||||
cmd.stdin(Stdio::null())
|
||||
.stdout(Stdio::null())
|
||||
.stderr(Stdio::null())
|
||||
.kill_on_drop(true);
|
||||
#[cfg(unix)]
|
||||
{
|
||||
// Own process group so teardown can signal the whole tree.
|
||||
cmd.process_group(0);
|
||||
}
|
||||
|
||||
let child = cmd
|
||||
.spawn()
|
||||
.map_err(|e| anyhow::anyhow!("mesh-llm node spawn failed: {e}"))?;
|
||||
|
||||
let node = Self {
|
||||
child,
|
||||
api_port: config.api_port,
|
||||
console_port: config.console_port,
|
||||
};
|
||||
node.wait_ready(Duration::from_secs(60)).await?;
|
||||
Ok(node)
|
||||
}
|
||||
|
||||
pub fn api_base_url(&self) -> String {
|
||||
format!("http://127.0.0.1:{}/v1", self.api_port)
|
||||
}
|
||||
|
||||
pub fn console_url(&self) -> String {
|
||||
format!("http://127.0.0.1:{}", self.console_port)
|
||||
}
|
||||
|
||||
async fn wait_ready(&self, timeout: Duration) -> anyhow::Result<()> {
|
||||
let deadline = tokio::time::Instant::now() + timeout;
|
||||
let url = format!("{}/api/status", self.console_url());
|
||||
loop {
|
||||
if tokio::time::Instant::now() >= deadline {
|
||||
anyhow::bail!("mesh-llm node did not become ready within {timeout:?}");
|
||||
}
|
||||
if let Ok(response) = reqwest::get(&url).await {
|
||||
if response.status().is_success() {
|
||||
return Ok(());
|
||||
}
|
||||
}
|
||||
tokio::time::sleep(Duration::from_millis(500)).await;
|
||||
}
|
||||
}
|
||||
|
||||
/// Fetch node status from the management API.
|
||||
pub async fn status(&self) -> anyhow::Result<NodeStatus> {
|
||||
let url = format!("{}/api/status", self.console_url());
|
||||
let payload: serde_json::Value = reqwest::get(&url)
|
||||
.await
|
||||
.map_err(|e| anyhow::anyhow!("mesh node status fetch failed: {e}"))?
|
||||
.error_for_status()
|
||||
.map_err(|e| anyhow::anyhow!("mesh node status fetch failed: {e}"))?
|
||||
.json()
|
||||
.await
|
||||
.map_err(|e| anyhow::anyhow!("mesh node status parse failed: {e}"))?;
|
||||
let invite_token = payload
|
||||
.get("token")
|
||||
.and_then(serde_json::Value::as_str)
|
||||
.map(ToString::to_string);
|
||||
Ok(NodeStatus {
|
||||
api_base_url: self.api_base_url(),
|
||||
console_url: self.console_url(),
|
||||
invite_token,
|
||||
payload,
|
||||
})
|
||||
}
|
||||
|
||||
/// Stop the node process (consuming form).
|
||||
pub async fn stop(mut self) -> anyhow::Result<()> {
|
||||
self.take_stop().await
|
||||
}
|
||||
|
||||
/// Stop the node process in place (for callers that only have `&mut`,
|
||||
/// e.g. replacing the node behind a mutex). Safe to call twice.
|
||||
pub async fn take_stop(&mut self) -> anyhow::Result<()> {
|
||||
#[cfg(unix)]
|
||||
{
|
||||
// Graceful first: SIGTERM the process group, then escalate.
|
||||
if let Some(pid) = self.child.id() {
|
||||
unsafe {
|
||||
libc::kill(-(pid as i32), libc::SIGTERM);
|
||||
}
|
||||
let grace = tokio::time::timeout(Duration::from_secs(5), self.child.wait()).await;
|
||||
if grace.is_ok() {
|
||||
return Ok(());
|
||||
}
|
||||
}
|
||||
}
|
||||
self.child
|
||||
.kill()
|
||||
.await
|
||||
.map_err(|e| anyhow::anyhow!("mesh node kill failed: {e}"))?;
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
@@ -55,6 +55,11 @@ pub async fn mesh_installed_models(
|
||||
Err("mesh-llm feature not enabled".to_string())
|
||||
}
|
||||
|
||||
#[tauri::command]
|
||||
pub fn mesh_node_install_status() -> CmdResult<serde_json::Value> {
|
||||
Err("mesh-llm feature not enabled".to_string())
|
||||
}
|
||||
|
||||
#[tauri::command]
|
||||
pub fn mesh_agent_preset(_request: serde_json::Value) -> CmdResult<serde_json::Value> {
|
||||
Err("mesh-llm feature not enabled".to_string())
|
||||
|
||||
Reference in New Issue
Block a user