mirror of
https://github.com/block/buzz.git
synced 2026-08-18 06:50:31 +02:00
fix(desktop): extract mesh readiness helpers to satisfy file-size ratchet
mesh_llm.rs was 1016 lines (split-count 1017, gate limit 1000). Extract MeshReadinessFailure, classify_mesh_readiness_failure, mesh_readiness_failure_message, and wait_for_mesh_inference into a new mesh_llm_readiness.rs submodule (140 lines). mesh_llm.rs is now 887 lines (split-count 888). All other files remain under the 1000-line limit. Tests that used the readiness items via use super::* add explicit use super::readiness::* imports since the items moved out of the parent namespace. No behavior change. Co-authored-by: Will Pfleger <pfleger.will@gmail.com> Signed-off-by: Will Pfleger <pfleger.will@gmail.com>
This commit is contained in:
co-authored by
Will Pfleger
parent
8b046bc62b
commit
c29cd651db
@@ -533,138 +533,9 @@ pub async fn mesh_start_node(
|
||||
Ok(status)
|
||||
}
|
||||
|
||||
/// Which startup stage a mesh client is stuck at when it never becomes
|
||||
/// inference-ready. The two live-observed failure modes are physically
|
||||
/// distinct and want different user copy:
|
||||
///
|
||||
/// * `CatalogNeverSynced` — the local client node came up and connected to
|
||||
/// the host at the control level (ping/RTT fine), but the served model
|
||||
/// never appeared in the local `/v1/models` catalog. That catalog is
|
||||
/// populated by the peer gossip exchange; when the gossip bi-stream can't
|
||||
/// establish across the network (observed as iroh
|
||||
/// `MultipathNotNegotiated` / unreachable direct path), the catalog stays
|
||||
/// empty forever and every request is rejected "model not available".
|
||||
/// Root cause is the network path between this machine and the host.
|
||||
/// * `RoutingNeverCompleted` — the model *did* sync into the catalog, but
|
||||
/// inference requests never completed (routing/transport to the host
|
||||
/// failing per-request). The host is discoverable and advertised but not
|
||||
/// actually serving us.
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||
enum MeshReadinessFailure {
|
||||
CatalogNeverSynced,
|
||||
RoutingNeverCompleted,
|
||||
}
|
||||
|
||||
/// Pure classifier: given whether the served model was ever observed in the
|
||||
/// local `/v1/models` catalog during the wait, decide which stage failed.
|
||||
/// Split out so the diagnosis is unit-testable without a live mesh.
|
||||
fn classify_mesh_readiness_failure(model_ever_visible: bool) -> MeshReadinessFailure {
|
||||
if model_ever_visible {
|
||||
MeshReadinessFailure::RoutingNeverCompleted
|
||||
} else {
|
||||
MeshReadinessFailure::CatalogNeverSynced
|
||||
}
|
||||
}
|
||||
|
||||
/// Actionable, non-technical copy for a readiness failure. `last_detail` is the
|
||||
/// last raw transport/HTTP error, appended for support triage.
|
||||
fn mesh_readiness_failure_message(
|
||||
failure: MeshReadinessFailure,
|
||||
model_id: &str,
|
||||
last_detail: &str,
|
||||
) -> String {
|
||||
match failure {
|
||||
MeshReadinessFailure::CatalogNeverSynced => format!(
|
||||
"Buzz shared compute connected to the serving member but could not sync \
|
||||
the model list for \"{model_id}\" — this is a network path problem \
|
||||
between this machine and the host (the compute node is reachable for \
|
||||
pings but the model-sync stream did not establish). Try again, or have \
|
||||
the host and this machine on a more direct network. (last: {last_detail})"
|
||||
),
|
||||
MeshReadinessFailure::RoutingNeverCompleted => format!(
|
||||
"Buzz shared compute found \"{model_id}\" on a serving member but inference \
|
||||
requests did not complete — the host is discoverable but not currently \
|
||||
reachable for requests. Try again shortly. (last: {last_detail})"
|
||||
),
|
||||
}
|
||||
}
|
||||
|
||||
/// Poll the local mesh OpenAI ingress until a real inference for `model_id`
|
||||
/// succeeds, or a deadline elapses. On failure, returns a stage-specific,
|
||||
/// actionable message (see [`MeshReadinessFailure`]) rather than a raw
|
||||
/// `HTTP 429`, so the UI can tell "still warming up" apart from "can't reach
|
||||
/// the host".
|
||||
async fn wait_for_mesh_inference(model_id: &str) -> CmdResult<()> {
|
||||
let client = reqwest::Client::builder()
|
||||
.timeout(std::time::Duration::from_secs(30))
|
||||
.build()
|
||||
.map_err(|error| format!("failed to build mesh readiness client: {error}"))?;
|
||||
let deadline = tokio::time::Instant::now() + std::time::Duration::from_secs(120);
|
||||
let models_url = format!("{}/models", crate::managed_agents::RELAY_MESH_API_BASE_URL);
|
||||
let chat_url = format!(
|
||||
"{}/chat/completions",
|
||||
crate::managed_agents::RELAY_MESH_API_BASE_URL
|
||||
);
|
||||
let mut last_error = "mesh inference is not ready".to_string();
|
||||
// Track whether the served model ever reached the local catalog — the
|
||||
// signal that splits "catalog never synced" from "routing never completed".
|
||||
let mut model_ever_visible = false;
|
||||
|
||||
while tokio::time::Instant::now() < deadline {
|
||||
// Refresh catalog visibility. "auto" delegates model choice to the
|
||||
// router, so any advertised model counts as the catalog having synced.
|
||||
if let Ok(response) = client
|
||||
.get(&models_url)
|
||||
.bearer_auth(crate::managed_agents::RELAY_MESH_API_KEY_PLACEHOLDER)
|
||||
.send()
|
||||
.await
|
||||
{
|
||||
if let Ok(body) = response.json::<serde_json::Value>().await {
|
||||
if let Some(data) = body.get("data").and_then(|d| d.as_array()) {
|
||||
let wanted = model_id.trim().replace("@main", "");
|
||||
let visible = !data.is_empty()
|
||||
&& (model_id == crate::mesh_llm::AUTO_MODEL_ID
|
||||
|| data.iter().any(|m| {
|
||||
m.get("id")
|
||||
.and_then(|id| id.as_str())
|
||||
.map(|id| id.replace("@main", "") == wanted)
|
||||
.unwrap_or(false)
|
||||
}));
|
||||
model_ever_visible |= visible;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
match client
|
||||
.post(&chat_url)
|
||||
.bearer_auth(crate::managed_agents::RELAY_MESH_API_KEY_PLACEHOLDER)
|
||||
.json(&serde_json::json!({
|
||||
"model": model_id,
|
||||
"messages": [{"role": "user", "content": "Reply OK"}],
|
||||
"max_tokens": 1,
|
||||
"stream": false
|
||||
}))
|
||||
.send()
|
||||
.await
|
||||
{
|
||||
Ok(response) if response.status().is_success() => return Ok(()),
|
||||
Ok(response) => {
|
||||
let status = response.status();
|
||||
let body = response.text().await.unwrap_or_default();
|
||||
last_error = format!("HTTP {status}: {body}");
|
||||
}
|
||||
Err(error) => last_error = error.to_string(),
|
||||
}
|
||||
tokio::time::sleep(std::time::Duration::from_secs(2)).await;
|
||||
}
|
||||
|
||||
let failure = classify_mesh_readiness_failure(model_ever_visible);
|
||||
Err(mesh_readiness_failure_message(
|
||||
failure,
|
||||
model_id,
|
||||
&last_error,
|
||||
))
|
||||
}
|
||||
#[path = "mesh_llm_readiness.rs"]
|
||||
mod readiness;
|
||||
use readiness::wait_for_mesh_inference;
|
||||
|
||||
pub(crate) async fn ensure_client_node_for_model(
|
||||
state: &AppState,
|
||||
|
||||
@@ -0,0 +1,138 @@
|
||||
//! Mesh inference readiness polling helpers.
|
||||
//!
|
||||
//! Extracted from `mesh_llm.rs` to stay within the file-size ratchet.
|
||||
//! All items are `pub(super)` — visible to `mesh_llm` and its test submodule
|
||||
//! via `use super::readiness::*`.
|
||||
|
||||
/// Which startup stage a mesh client is stuck at when it never becomes
|
||||
/// inference-ready. The two live-observed failure modes are physically
|
||||
/// distinct and want different user copy:
|
||||
///
|
||||
/// * `CatalogNeverSynced` — the local client node came up and connected to
|
||||
/// the host at the control level (ping/RTT fine), but the served model
|
||||
/// never appeared in the local `/v1/models` catalog. That catalog is
|
||||
/// populated by the peer gossip exchange; when the gossip bi-stream can't
|
||||
/// establish across the network (observed as iroh
|
||||
/// `MultipathNotNegotiated` / unreachable direct path), the catalog stays
|
||||
/// empty forever and every request is rejected "model not available".
|
||||
/// Root cause is the network path between this machine and the host.
|
||||
/// * `RoutingNeverCompleted` — the model *did* sync into the catalog, but
|
||||
/// inference requests never completed (routing/transport to the host
|
||||
/// failing per-request). The host is discoverable and advertised but not
|
||||
/// actually serving us.
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||
pub(super) enum MeshReadinessFailure {
|
||||
CatalogNeverSynced,
|
||||
RoutingNeverCompleted,
|
||||
}
|
||||
|
||||
/// Pure classifier: given whether the served model was ever observed in the
|
||||
/// local `/v1/models` catalog during the wait, decide which stage failed.
|
||||
/// Split out so the diagnosis is unit-testable without a live mesh.
|
||||
pub(super) fn classify_mesh_readiness_failure(model_ever_visible: bool) -> MeshReadinessFailure {
|
||||
if model_ever_visible {
|
||||
MeshReadinessFailure::RoutingNeverCompleted
|
||||
} else {
|
||||
MeshReadinessFailure::CatalogNeverSynced
|
||||
}
|
||||
}
|
||||
|
||||
/// Actionable, non-technical copy for a readiness failure. `last_detail` is the
|
||||
/// last raw transport/HTTP error, appended for support triage.
|
||||
pub(super) fn mesh_readiness_failure_message(
|
||||
failure: MeshReadinessFailure,
|
||||
model_id: &str,
|
||||
last_detail: &str,
|
||||
) -> String {
|
||||
match failure {
|
||||
MeshReadinessFailure::CatalogNeverSynced => format!(
|
||||
"Buzz shared compute connected to the serving member but could not sync \
|
||||
the model list for \"{model_id}\" — this is a network path problem \
|
||||
between this machine and the host (the compute node is reachable for \
|
||||
pings but the model-sync stream did not establish). Try again, or have \
|
||||
the host and this machine on a more direct network. (last: {last_detail})"
|
||||
),
|
||||
MeshReadinessFailure::RoutingNeverCompleted => format!(
|
||||
"Buzz shared compute found \"{model_id}\" on a serving member but inference \
|
||||
requests did not complete — the host is discoverable but not currently \
|
||||
reachable for requests. Try again shortly. (last: {last_detail})"
|
||||
),
|
||||
}
|
||||
}
|
||||
|
||||
/// Poll the local mesh OpenAI ingress until a real inference for `model_id`
|
||||
/// succeeds, or a deadline elapses. On failure, returns a stage-specific,
|
||||
/// actionable message (see [`MeshReadinessFailure`]) rather than a raw
|
||||
/// `HTTP 429`, so the UI can tell "still warming up" apart from "can't reach
|
||||
/// the host".
|
||||
pub(super) async fn wait_for_mesh_inference(model_id: &str) -> Result<(), String> {
|
||||
let client = reqwest::Client::builder()
|
||||
.timeout(std::time::Duration::from_secs(30))
|
||||
.build()
|
||||
.map_err(|error| format!("failed to build mesh readiness client: {error}"))?;
|
||||
let deadline = tokio::time::Instant::now() + std::time::Duration::from_secs(120);
|
||||
let models_url = format!("{}/models", crate::managed_agents::RELAY_MESH_API_BASE_URL);
|
||||
let chat_url = format!(
|
||||
"{}/chat/completions",
|
||||
crate::managed_agents::RELAY_MESH_API_BASE_URL
|
||||
);
|
||||
let mut last_error = "mesh inference is not ready".to_string();
|
||||
// Track whether the served model ever reached the local catalog — the
|
||||
// signal that splits "catalog never synced" from "routing never completed".
|
||||
let mut model_ever_visible = false;
|
||||
|
||||
while tokio::time::Instant::now() < deadline {
|
||||
// Refresh catalog visibility. "auto" delegates model choice to the
|
||||
// router, so any advertised model counts as the catalog having synced.
|
||||
if let Ok(response) = client
|
||||
.get(&models_url)
|
||||
.bearer_auth(crate::managed_agents::RELAY_MESH_API_KEY_PLACEHOLDER)
|
||||
.send()
|
||||
.await
|
||||
{
|
||||
if let Ok(body) = response.json::<serde_json::Value>().await {
|
||||
if let Some(data) = body.get("data").and_then(|d| d.as_array()) {
|
||||
let wanted = model_id.trim().replace("@main", "");
|
||||
let visible = !data.is_empty()
|
||||
&& (model_id == crate::mesh_llm::AUTO_MODEL_ID
|
||||
|| data.iter().any(|m| {
|
||||
m.get("id")
|
||||
.and_then(|id| id.as_str())
|
||||
.map(|id| id.replace("@main", "") == wanted)
|
||||
.unwrap_or(false)
|
||||
}));
|
||||
model_ever_visible |= visible;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
match client
|
||||
.post(&chat_url)
|
||||
.bearer_auth(crate::managed_agents::RELAY_MESH_API_KEY_PLACEHOLDER)
|
||||
.json(&serde_json::json!({
|
||||
"model": model_id,
|
||||
"messages": [{"role": "user", "content": "Reply OK"}],
|
||||
"max_tokens": 1,
|
||||
"stream": false
|
||||
}))
|
||||
.send()
|
||||
.await
|
||||
{
|
||||
Ok(response) if response.status().is_success() => return Ok(()),
|
||||
Ok(response) => {
|
||||
let status = response.status();
|
||||
let body = response.text().await.unwrap_or_default();
|
||||
last_error = format!("HTTP {status}: {body}");
|
||||
}
|
||||
Err(error) => last_error = error.to_string(),
|
||||
}
|
||||
tokio::time::sleep(std::time::Duration::from_secs(2)).await;
|
||||
}
|
||||
|
||||
let failure = classify_mesh_readiness_failure(model_ever_visible);
|
||||
Err(mesh_readiness_failure_message(
|
||||
failure,
|
||||
model_id,
|
||||
&last_error,
|
||||
))
|
||||
}
|
||||
@@ -116,6 +116,7 @@ pub(crate) async fn fail_if_client_mesh_active<R: tauri::Runtime>(
|
||||
/// Alternatively, callers that cannot fit their body into a sync `FnOnce` call
|
||||
/// `run_mesh_transition_preflight` (which they invoke after acquiring the lock
|
||||
/// themselves) to share the preflight logic without lifetime constraints.
|
||||
#[cfg_attr(not(test), allow(dead_code))]
|
||||
pub(crate) async fn with_workspace_transition_preflight<R, F, T>(
|
||||
app: &AppHandle<R>,
|
||||
transition_body: F,
|
||||
|
||||
@@ -1,3 +1,6 @@
|
||||
use super::readiness::{
|
||||
classify_mesh_readiness_failure, mesh_readiness_failure_message, MeshReadinessFailure,
|
||||
};
|
||||
use super::*;
|
||||
use crate::app_state::build_app_state;
|
||||
|
||||
|
||||
Reference in New Issue
Block a user