From a432e76655215dfb7c10b33841a262258bf2eed9 Mon Sep 17 00:00:00 2001 From: James Dumay Date: Wed, 25 Feb 2026 14:05:55 +1100 Subject: [PATCH] Use SSE for telemetry throughput --- .idea/.gitignore | 10 + .idea/.name | 1 + .idea/decentralized-inference.iml | 11 + .idea/inspectionProfiles/Project_Default.xml | 6 + .idea/modules.xml | 8 + .idea/prettier.xml | 6 + .idea/vcs.xml | 7 + mesh-llm/Cargo.lock | 75 +- mesh-llm/Cargo.toml | 1 + mesh-llm/src/api.rs | 179 +++- mesh-llm/src/benchmark.rs | 548 ++++++++++++ mesh-llm/src/console.html | 472 ++++++++++- mesh-llm/src/main.rs | 72 +- mesh-llm/src/mesh.rs | 37 +- mesh-llm/src/proxy.rs | 35 +- mesh-llm/src/telemetry.rs | 840 +++++++++++++++++++ 16 files changed, 2261 insertions(+), 47 deletions(-) create mode 100644 .idea/.gitignore create mode 100644 .idea/.name create mode 100644 .idea/decentralized-inference.iml create mode 100644 .idea/inspectionProfiles/Project_Default.xml create mode 100644 .idea/modules.xml create mode 100644 .idea/prettier.xml create mode 100644 .idea/vcs.xml create mode 100644 mesh-llm/src/benchmark.rs create mode 100644 mesh-llm/src/telemetry.rs diff --git a/.idea/.gitignore b/.idea/.gitignore new file mode 100644 index 000000000..ab1f4164e --- /dev/null +++ b/.idea/.gitignore @@ -0,0 +1,10 @@ +# Default ignored files +/shelf/ +/workspace.xml +# Ignored default folder with query files +/queries/ +# Datasource local storage ignored files +/dataSources/ +/dataSources.local.xml +# Editor-based HTTP Client requests +/httpRequests/ diff --git a/.idea/.name b/.idea/.name new file mode 100644 index 000000000..83fc3eac7 --- /dev/null +++ b/.idea/.name @@ -0,0 +1 @@ +telemetry.rs \ No newline at end of file diff --git a/.idea/decentralized-inference.iml b/.idea/decentralized-inference.iml new file mode 100644 index 000000000..ea135a485 --- /dev/null +++ b/.idea/decentralized-inference.iml @@ -0,0 +1,11 @@ + + + + + + + + + + + \ No newline at end of file diff --git a/.idea/inspectionProfiles/Project_Default.xml b/.idea/inspectionProfiles/Project_Default.xml new file mode 100644 index 000000000..03d9549ea --- /dev/null +++ b/.idea/inspectionProfiles/Project_Default.xml @@ -0,0 +1,6 @@ + + + + \ No newline at end of file diff --git a/.idea/modules.xml b/.idea/modules.xml new file mode 100644 index 000000000..4b4daf0c6 --- /dev/null +++ b/.idea/modules.xml @@ -0,0 +1,8 @@ + + + + + + + + \ No newline at end of file diff --git a/.idea/prettier.xml b/.idea/prettier.xml new file mode 100644 index 000000000..b0c1c68fb --- /dev/null +++ b/.idea/prettier.xml @@ -0,0 +1,6 @@ + + + + + \ No newline at end of file diff --git a/.idea/vcs.xml b/.idea/vcs.xml new file mode 100644 index 000000000..2376cc2c4 --- /dev/null +++ b/.idea/vcs.xml @@ -0,0 +1,7 @@ + + + + + + + \ No newline at end of file diff --git a/mesh-llm/Cargo.lock b/mesh-llm/Cargo.lock index 4952fca51..7aa9251f2 100644 --- a/mesh-llm/Cargo.lock +++ b/mesh-llm/Cargo.lock @@ -12,6 +12,18 @@ dependencies = [ "generic-array", ] +[[package]] +name = "ahash" +version = "0.8.12" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5a15f179cd60c4584b8a8c596927aadc462e27f2ca70c04e0071964a73ba7a75" +dependencies = [ + "cfg-if", + "once_cell", + "version_check", + "zerocopy", +] + [[package]] name = "aho-corasick" version = "1.1.4" @@ -959,6 +971,18 @@ dependencies = [ "windows-sys 0.61.2", ] +[[package]] +name = "fallible-iterator" +version = "0.3.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2acce4a10f12dc2fb14a218589d4f1f62ef011b2d0cc4b3cb1bba8e94da14649" + +[[package]] +name = "fallible-streaming-iterator" +version = "0.1.9" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7360491ce676a36bf9bb3c56c1aa791658183a54d2744120f27285738d90465a" + [[package]] name = "fastbloom" version = "0.14.1" @@ -1238,6 +1262,15 @@ dependencies = [ "byteorder", ] +[[package]] +name = "hashbrown" +version = "0.14.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e5274423e17b7c9fc20b6e7e208532f9b19825d82dfd615708b70edd83df41f1" +dependencies = [ + "ahash", +] + [[package]] name = "hashbrown" version = "0.16.1" @@ -1249,6 +1282,15 @@ dependencies = [ "foldhash", ] +[[package]] +name = "hashlink" +version = "0.9.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6ba4ff7128dee98c7dc9794b6a411377e1404dba1c97deb8d1a55297bd25d8af" +dependencies = [ + "hashbrown 0.14.5", +] + [[package]] name = "heapless" version = "0.7.17" @@ -1647,7 +1689,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "7714e70437a7dc3ac8eb7e6f8df75fd8eb422675fc7678aff7364301092b1017" dependencies = [ "equivalent", - "hashbrown", + "hashbrown 0.16.1", ] [[package]] @@ -1967,6 +2009,16 @@ dependencies = [ "libc", ] +[[package]] +name = "libsqlite3-sys" +version = "0.30.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2e99fb7a497b1e3339bc746195567ed8d3e24945ecd636e3619d20b9de9e9149" +dependencies = [ + "pkg-config", + "vcpkg", +] + [[package]] name = "linux-raw-sys" version = "0.11.0" @@ -2019,7 +2071,7 @@ version = "0.16.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "a1dc47f592c06f33f8e3aea9591776ec7c9f9e4124778ff8a3c3b87159f7e593" dependencies = [ - "hashbrown", + "hashbrown 0.16.1", ] [[package]] @@ -2062,6 +2114,7 @@ dependencies = [ "nostr-sdk", "rand 0.9.2", "reqwest", + "rusqlite", "rustls", "serde", "serde_json", @@ -3110,6 +3163,20 @@ dependencies = [ "windows-sys 0.52.0", ] +[[package]] +name = "rusqlite" +version = "0.32.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7753b721174eb8ff87a9a0e799e2d7bc3749323e773db92e0984debb00019d6e" +dependencies = [ + "bitflags", + "fallible-iterator", + "fallible-streaming-iterator", + "hashlink", + "libsqlite3-sys", + "smallvec", +] + [[package]] name = "rustc-hash" version = "2.1.1" @@ -3135,7 +3202,7 @@ dependencies = [ "errno", "libc", "linux-raw-sys", - "windows-sys 0.52.0", + "windows-sys 0.61.2", ] [[package]] @@ -3638,7 +3705,7 @@ dependencies = [ "getrandom 0.3.4", "once_cell", "rustix", - "windows-sys 0.52.0", + "windows-sys 0.61.2", ] [[package]] diff --git a/mesh-llm/Cargo.toml b/mesh-llm/Cargo.toml index c2ffe77d4..baea6ec34 100644 --- a/mesh-llm/Cargo.toml +++ b/mesh-llm/Cargo.toml @@ -20,3 +20,4 @@ nostr-sdk = { version = "0.44.1", default-features = false } rustls = "0.23.36" reqwest = { version = "0.12", features = ["stream", "json"] } tokio-stream = "0.1" +rusqlite = "0.32" diff --git a/mesh-llm/src/api.rs b/mesh-llm/src/api.rs index ccdb52d6a..994393726 100644 --- a/mesh-llm/src/api.rs +++ b/mesh-llm/src/api.rs @@ -10,7 +10,7 @@ //! The console is read-only — shows status, topology, models. //! All mutations happen via CLI flags (--join, --model, --auto). -use crate::{download, election, mesh, nostr}; +use crate::{download, election, mesh, nostr, telemetry}; use serde::Serialize; use std::sync::Arc; use tokio::io::{AsyncReadExt, AsyncWriteExt}; @@ -29,6 +29,7 @@ pub struct MeshApi { struct ApiInner { node: mesh::Node, + telemetry: telemetry::Telemetry, is_host: bool, is_client: bool, llama_ready: bool, @@ -82,11 +83,27 @@ struct MeshModelPayload { size_gb: f64, } +#[derive(Serialize)] +struct TelemetryEventsPayload { + live: telemetry::LiveSnapshot, + nodes: Vec, + rollup: Vec, + node_history: Vec, + benchmarks: Vec, +} + impl MeshApi { - pub fn new(node: mesh::Node, model_name: String, api_port: u16, model_size_bytes: u64) -> Self { + pub fn new( + node: mesh::Node, + telemetry: telemetry::Telemetry, + model_name: String, + api_port: u16, + model_size_bytes: u64, + ) -> Self { MeshApi { inner: Arc::new(Mutex::new(ApiInner { node, + telemetry, is_host: false, is_client: false, llama_ready: false, @@ -372,6 +389,135 @@ async fn handle_request(mut stream: TcpStream, state: &MeshApi) -> anyhow::Resul stream.write_all(resp.as_bytes()).await?; } + // ── Telemetry: local live snapshot ── + "/api/metrics/live" => { + let telemetry = state.inner.lock().await.telemetry.clone(); + let json = serde_json::to_string(&telemetry.snapshot())?; + let resp = format!( + "HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\n\r\n{}", + json.len(), json + ); + stream.write_all(resp.as_bytes()).await?; + } + + // ── Telemetry SSE stream (full package for UI) ── + p if p.starts_with("/api/metrics/events") => { + let header = "HTTP/1.1 200 OK\r\nContent-Type: text/event-stream\r\nCache-Control: no-cache\r\nConnection: keep-alive\r\n\r\n"; + stream.write_all(header.as_bytes()).await?; + let telemetry = state.inner.lock().await.telemetry.clone(); + let minutes = query_param_u32(p, "minutes").unwrap_or(180).clamp(1, 7 * 24 * 60); + let node_id = query_param(p, "id").map(|s| s.to_string()); + let limit = query_param_u32(p, "limit").unwrap_or(300).clamp(1, 2000); + loop { + let node_history = if let Some(id) = node_id.clone() { + telemetry.node_history_for(id, minutes).await + } else { + telemetry.node_history(minutes).await + }; + let payload = match ( + telemetry.all_nodes_latest(minutes).await, + telemetry.rollup_history(minutes).await, + node_history, + telemetry.benchmark_history(minutes, limit).await, + ) { + (Ok(nodes), Ok(rollup), Ok(node_history), Ok(benchmarks)) => TelemetryEventsPayload { + live: telemetry.snapshot(), + nodes, + rollup, + node_history, + benchmarks, + }, + _ => { + if stream.write_all(b"event: error\ndata: {\"error\":\"telemetry query failed\"}\n\n").await.is_err() { + break; + } + tokio::time::sleep(std::time::Duration::from_secs(5)).await; + continue; + } + }; + let json = serde_json::to_string(&payload)?; + if stream.write_all(format!("data: {json}\n\n").as_bytes()).await.is_err() { + break; + } + tokio::time::sleep(std::time::Duration::from_secs(5)).await; + } + } + + // ── Telemetry: local node history (SQLite) ── + p if p.starts_with("/api/metrics/node") => { + let telemetry = state.inner.lock().await.telemetry.clone(); + let minutes = query_param_u32(p, "minutes").unwrap_or(60).clamp(1, 24 * 60); + let node_id = query_param(p, "id").map(|s| s.to_string()); + let result = if let Some(id) = node_id { + telemetry.node_history_for(id, minutes).await + } else { + telemetry.node_history(minutes).await + }; + match result { + Ok(rows) => { + let json = serde_json::to_string(&rows)?; + let resp = format!( + "HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\n\r\n{}", + json.len(), json + ); + stream.write_all(resp.as_bytes()).await?; + } + Err(e) => respond_error(&mut stream, 500, &format!("Telemetry query failed: {e}")).await?, + } + } + + // ── Telemetry: latest row per known node (SQLite, local DB) ── + p if p.starts_with("/api/metrics/nodes") => { + let telemetry = state.inner.lock().await.telemetry.clone(); + let minutes = query_param_u32(p, "minutes").unwrap_or(180).clamp(1, 7 * 24 * 60); + match telemetry.all_nodes_latest(minutes).await { + Ok(rows) => { + let json = serde_json::to_string(&rows)?; + let resp = format!( + "HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\n\r\n{}", + json.len(), json + ); + stream.write_all(resp.as_bytes()).await?; + } + Err(e) => respond_error(&mut stream, 500, &format!("Telemetry query failed: {e}")).await?, + } + } + + // ── Telemetry: rollup history across rows in local DB ── + p if p.starts_with("/api/metrics/rollup") => { + let telemetry = state.inner.lock().await.telemetry.clone(); + let minutes = query_param_u32(p, "minutes").unwrap_or(60).clamp(1, 7 * 24 * 60); + match telemetry.rollup_history(minutes).await { + Ok(rows) => { + let json = serde_json::to_string(&rows)?; + let resp = format!( + "HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\n\r\n{}", + json.len(), json + ); + stream.write_all(resp.as_bytes()).await?; + } + Err(e) => respond_error(&mut stream, 500, &format!("Telemetry query failed: {e}")).await?, + } + } + + // ── Telemetry: benchmark run history (raw local rows, recent only) ── + p if p.starts_with("/api/metrics/benchmarks") => { + let telemetry = state.inner.lock().await.telemetry.clone(); + let minutes = query_param_u32(p, "minutes").unwrap_or(60).clamp(1, 7 * 24 * 60); + let limit = query_param_u32(p, "limit").unwrap_or(200).clamp(1, 2000); + match telemetry.benchmark_history(minutes, limit).await { + Ok(rows) => { + let json = serde_json::to_string(&rows)?; + let resp = format!( + "HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\n\r\n{}", + json.len(), json + ); + stream.write_all(resp.as_bytes()).await?; + } + Err(e) => respond_error(&mut stream, 500, &format!("Telemetry query failed: {e}")).await?, + } + } + // ── SSE event stream ── "/api/events" => { let header = "HTTP/1.1 200 OK\r\nContent-Type: text/event-stream\r\nCache-Control: no-cache\r\nConnection: keep-alive\r\n\r\n"; @@ -394,10 +540,14 @@ async fn handle_request(mut stream: TcpStream, state: &MeshApi) -> anyhow::Resul // ── Chat proxy (routes through inference API port) ── p if p.starts_with("/api/chat") => { + let telemetry = state.inner.lock().await.telemetry.clone(); + let span = telemetry.start_request(telemetry::RouteKind::Local); let inner = state.inner.lock().await; if !inner.llama_ready && !inner.is_client { drop(inner); - return respond_error(&mut stream, 503, "LLM not ready").await; + let res = respond_error(&mut stream, 503, "LLM not ready").await; + span.finish(false); + return res; } let port = inner.api_port; drop(inner); @@ -405,9 +555,13 @@ async fn handle_request(mut stream: TcpStream, state: &MeshApi) -> anyhow::Resul if let Ok(mut upstream) = TcpStream::connect(&target).await { let rewritten = req.replacen("/api/chat", "/v1/chat/completions", 1); upstream.write_all(rewritten.as_bytes()).await?; - tokio::io::copy_bidirectional(&mut stream, &mut upstream).await?; + let result = tokio::io::copy_bidirectional(&mut stream, &mut upstream).await; + span.finish(result.is_ok()); + result?; } else { - respond_error(&mut stream, 502, "Cannot reach LLM server").await?; + let res = respond_error(&mut stream, 502, "Cannot reach LLM server").await; + span.finish(false); + res?; } } @@ -418,6 +572,21 @@ async fn handle_request(mut stream: TcpStream, state: &MeshApi) -> anyhow::Resul Ok(()) } +fn query_param_u32(path: &str, key: &str) -> Option { + query_param(path, key)?.parse::().ok() +} + +fn query_param<'a>(path: &'a str, key: &str) -> Option<&'a str> { + let (_, query) = path.split_once('?')?; + for pair in query.split('&') { + let (k, v) = pair.split_once('=')?; + if k == key { + return Some(v); + } + } + None +} + async fn respond_error(stream: &mut TcpStream, code: u16, msg: &str) -> anyhow::Result<()> { let body = format!("{{\"error\":\"{msg}\"}}"); let status = match code { diff --git a/mesh-llm/src/benchmark.rs b/mesh-llm/src/benchmark.rs new file mode 100644 index 000000000..ff7bca771 --- /dev/null +++ b/mesh-llm/src/benchmark.rs @@ -0,0 +1,548 @@ +use crate::{mesh, telemetry}; +use serde::Deserialize; +use serde::Serialize; +use serde_json::json; +use std::time::Instant; + +#[derive(Debug, Clone, Serialize)] +pub struct BenchmarkProbe { + pub name: &'static str, + pub prompt: &'static str, + pub settings: BenchmarkSettings, +} + +#[derive(Debug, Clone, Serialize)] +pub struct BenchmarkSettings { + pub temperature: f64, + pub top_p: f64, + pub max_tokens: u32, + pub stream: bool, +} + +#[allow(dead_code)] +pub fn analysis_report_v1() -> BenchmarkProbe { + BenchmarkProbe { + name: "analysis_report_v1", + prompt: ANALYSIS_REPORT_V1_PROMPT, + settings: BenchmarkSettings { + temperature: 0.0, + top_p: 1.0, + max_tokens: 300, + stream: false, + }, + } +} + +#[allow(dead_code)] +pub fn analysis_report_v1_stream() -> BenchmarkProbe { + BenchmarkProbe { + name: "analysis_report_v1_stream", + prompt: ANALYSIS_REPORT_V1_PROMPT, + settings: BenchmarkSettings { + temperature: 0.0, + top_p: 1.0, + max_tokens: 300, + stream: true, + }, + } +} + +pub const ANALYSIS_REPORT_V1_PROMPT: &str = r#"You are a senior AI systems analyst. + +Read the following internal engineering report and produce: + +1. A concise executive summary (5 bullet points) +2. A list of the top 5 technical risks +3. Three concrete performance improvement recommendations +4. A final 2–3 sentence conclusion + +Be precise, structured, and avoid generic language. + +--- + +INTERNAL REPORT: + +System Overview: +The LLM inference service currently runs a 13B parameter transformer model behind an HTTP API. The deployment uses tensor parallelism across 2 GPUs with dynamic batching enabled. Requests are routed through a queue that aggregates requests for up to 15ms before dispatch. The average request context length is 1,200 input tokens with a mean generation length of 220 output tokens. + +Observed Performance: +Over the past 7 days, median latency has remained stable at 820ms. However, p95 latency increased from 1.4s to 2.1s during peak traffic windows (14:00–18:00 UTC). GPU utilization averages 72% during off-peak hours and reaches 96–98% during peak hours. KV cache usage frequently exceeds 85% during sustained load. Occasional batch size reductions have been observed when memory pressure increases. + +Infrastructure Notes: +Each node has 2x A100 40GB GPUs. CPU utilization remains below 55% under load. Network I/O does not appear saturated. No significant disk I/O is observed. Horizontal scaling is currently manual. Autoscaling is under consideration but not implemented. + +Recent Changes: +A new feature was deployed enabling longer context windows (up to 8k tokens). After deployment, memory fragmentation increased and average prefill time rose by 18%. Decode token throughput per GPU decreased from 145 tokens/sec to 123 tokens/sec under peak concurrency. + +Error Metrics: +HTTP 5xx errors remain below 0.3%. However, request timeouts increased by 1.1% during peak hours. No GPU OOM crashes have been recorded, but soft memory allocation retries increased. + +Operational Constraints: +Latency SLO: p95 < 1.5s +Target throughput: 30 requests/sec sustained +Cost sensitivity: moderate — GPU overprovisioning should be avoided if possible. + +--- + +Produce your response now."#; + +pub fn spawn_probe_loop( + node: mesh::Node, + telemetry: telemetry::Telemetry, + api_port: u16, + interval_secs: u64, + model_override: Option, +) { + let interval_secs = interval_secs.max(10); + tokio::spawn(async move { + let client = match reqwest::Client::builder() + .timeout(std::time::Duration::from_secs(55)) + .build() + { + Ok(c) => c, + Err(e) => { + tracing::warn!("Benchmark runner disabled: failed to build HTTP client: {e}"); + return; + } + }; + let probe_nonstream = analysis_report_v1(); + let probe_stream = analysis_report_v1_stream(); + let mut run_stream_next = false; + loop { + let probe = if run_stream_next { &probe_stream } else { &probe_nonstream }; + if let Err(e) = run_probe_once(&client, &node, &telemetry, api_port, probe, model_override.as_deref()).await { + tracing::debug!("Benchmark probe failed to record: {e}"); + } + run_stream_next = !run_stream_next; + tokio::time::sleep(std::time::Duration::from_secs(interval_secs)).await; + } + }); +} + +async fn run_probe_once( + client: &reqwest::Client, + node: &mesh::Node, + telemetry: &telemetry::Telemetry, + api_port: u16, + probe: &BenchmarkProbe, + model_override: Option<&str>, +) -> anyhow::Result<()> { + let ts = unix_ts_now(); + let mesh_id = node.mesh_id().await; + let model = match model_override { + Some(m) => m.to_string(), + None => match discover_model(client, api_port).await { + Some(m) => m, + None => { + tracing::debug!("Benchmark probe {} skipped: no model available from /v1/models yet", probe.name); + return Ok(()); + } + }, + }; + + let url = format!("http://127.0.0.1:{api_port}/v1/chat/completions"); + let body = if probe.settings.stream { + json!({ + "model": model, + "messages": [ + {"role": "user", "content": probe.prompt} + ], + "temperature": probe.settings.temperature, + "top_p": probe.settings.top_p, + "max_tokens": probe.settings.max_tokens, + "stream": true, + "stream_options": {"include_usage": true}, + }) + } else { + json!({ + "model": model, + "messages": [ + {"role": "user", "content": probe.prompt} + ], + "temperature": probe.settings.temperature, + "top_p": probe.settings.top_p, + "max_tokens": probe.settings.max_tokens, + "stream": false, + }) + }; + + if probe.settings.stream { + return run_stream_probe_once(client, telemetry, probe, ts, mesh_id, model, url, body).await; + } + + let started = Instant::now(); + let resp = client.post(&url).json(&body).send().await; + let latency_ms = started.elapsed().as_millis().min(u128::from(u32::MAX)) as u32; + + match resp { + Ok(resp) => { + let status = resp.status(); + let status_code = Some(status.as_u16()); + let text = resp.text().await.unwrap_or_default(); + if !status.is_success() { + telemetry + .insert_benchmark_run(telemetry::BenchmarkRunInput { + ts, + mesh_id, + target_node_id: None, + model, + probe_name: probe.name.into(), + probe_type: "nonstream".into(), + stream: probe.settings.stream, + temperature: probe.settings.temperature, + top_p: probe.settings.top_p, + max_tokens: probe.settings.max_tokens, + prompt_hash: Some(fnv1a_hex(probe.prompt.as_bytes())), + route_kind: None, + success: false, + status_code, + latency_ms: Some(latency_ms), + ttft_ms: None, + prompt_tokens: None, + completion_tokens: None, + tokens_per_sec: None, + error_kind: Some("http_error".into()), + error_message: Some(truncate(&text, 300)), + response_shape_ok: None, + }) + .await?; + return Ok(()); + } + + let parsed: Option = serde_json::from_str(&text).ok(); + let (prompt_tokens, completion_tokens, response_shape_ok) = if let Some(ref p) = parsed { + let content = p.choices.first().and_then(|c| c.message.content.as_deref()).unwrap_or(""); + ( + p.usage.as_ref().map(|u| u.prompt_tokens), + p.usage.as_ref().map(|u| u.completion_tokens), + Some(validate_analysis_report_shape(content)), + ) + } else { + (None, None, None) + }; + let tps = completion_tokens.and_then(|n| { + if latency_ms == 0 { + None + } else { + Some(n as f64 / (latency_ms as f64 / 1000.0)) + } + }); + + telemetry + .insert_benchmark_run(telemetry::BenchmarkRunInput { + ts, + mesh_id, + target_node_id: None, + model, + probe_name: probe.name.into(), + probe_type: "nonstream".into(), + stream: probe.settings.stream, + temperature: probe.settings.temperature, + top_p: probe.settings.top_p, + max_tokens: probe.settings.max_tokens, + prompt_hash: Some(fnv1a_hex(probe.prompt.as_bytes())), + route_kind: None, + success: true, + status_code, + latency_ms: Some(latency_ms), + ttft_ms: None, + prompt_tokens, + completion_tokens, + tokens_per_sec: tps, + error_kind: None, + error_message: None, + response_shape_ok, + }) + .await?; + } + Err(e) => { + telemetry + .insert_benchmark_run(telemetry::BenchmarkRunInput { + ts, + mesh_id, + target_node_id: None, + model, + probe_name: probe.name.into(), + probe_type: "nonstream".into(), + stream: probe.settings.stream, + temperature: probe.settings.temperature, + top_p: probe.settings.top_p, + max_tokens: probe.settings.max_tokens, + prompt_hash: Some(fnv1a_hex(probe.prompt.as_bytes())), + route_kind: None, + success: false, + status_code: None, + latency_ms: Some(latency_ms), + ttft_ms: None, + prompt_tokens: None, + completion_tokens: None, + tokens_per_sec: None, + error_kind: Some(if e.is_timeout() { "timeout" } else { "request_error" }.into()), + error_message: Some(truncate(&e.to_string(), 300)), + response_shape_ok: None, + }) + .await?; + } + } + + Ok(()) +} + +async fn run_stream_probe_once( + client: &reqwest::Client, + telemetry: &telemetry::Telemetry, + probe: &BenchmarkProbe, + ts: i64, + mesh_id: Option, + model: String, + url: String, + body: serde_json::Value, +) -> anyhow::Result<()> { + let started = Instant::now(); + let resp = client.post(&url).json(&body).send().await; + + let mut resp = match resp { + Ok(r) => r, + Err(e) => { + let latency_ms = started.elapsed().as_millis().min(u128::from(u32::MAX)) as u32; + telemetry.insert_benchmark_run(telemetry::BenchmarkRunInput { + ts, + mesh_id, + target_node_id: None, + model, + probe_name: probe.name.into(), + probe_type: "stream".into(), + stream: true, + temperature: probe.settings.temperature, + top_p: probe.settings.top_p, + max_tokens: probe.settings.max_tokens, + prompt_hash: Some(fnv1a_hex(probe.prompt.as_bytes())), + route_kind: None, + success: false, + status_code: None, + latency_ms: Some(latency_ms), + ttft_ms: None, + prompt_tokens: None, + completion_tokens: None, + tokens_per_sec: None, + error_kind: Some(if e.is_timeout() { "timeout" } else { "request_error" }.into()), + error_message: Some(truncate(&e.to_string(), 300)), + response_shape_ok: None, + }).await?; + return Ok(()); + } + }; + + let status = resp.status(); + let status_code = Some(status.as_u16()); + if !status.is_success() { + let text = resp.text().await.unwrap_or_default(); + let latency_ms = started.elapsed().as_millis().min(u128::from(u32::MAX)) as u32; + telemetry.insert_benchmark_run(telemetry::BenchmarkRunInput { + ts, + mesh_id, + target_node_id: None, + model, + probe_name: probe.name.into(), + probe_type: "stream".into(), + stream: true, + temperature: probe.settings.temperature, + top_p: probe.settings.top_p, + max_tokens: probe.settings.max_tokens, + prompt_hash: Some(fnv1a_hex(probe.prompt.as_bytes())), + route_kind: None, + success: false, + status_code, + latency_ms: Some(latency_ms), + ttft_ms: None, + prompt_tokens: None, + completion_tokens: None, + tokens_per_sec: None, + error_kind: Some("http_error".into()), + error_message: Some(truncate(&text, 300)), + response_shape_ok: None, + }).await?; + return Ok(()); + } + + let mut pending = String::new(); + let mut saw_done = false; + let mut parse_error: Option = None; + let mut content = String::new(); + let mut ttft_ms: Option = None; + let mut prompt_tokens: Option = None; + let mut completion_tokens: Option = None; + + while let Some(chunk) = resp.chunk().await? { + pending.push_str(&String::from_utf8_lossy(&chunk)); + while let Some(idx) = pending.find("\n\n") { + let event = pending[..idx].to_string(); + pending.drain(..idx + 2); + + for line in event.lines() { + let Some(data) = line.strip_prefix("data: ") else { continue }; + let payload = data.trim(); + if payload == "[DONE]" { + saw_done = true; + continue; + } + match serde_json::from_str::(payload) { + Ok(chunk) => { + if let Some(usage) = chunk.usage { + prompt_tokens = Some(usage.prompt_tokens); + completion_tokens = Some(usage.completion_tokens); + } + if let Some(choice) = chunk.choices.first() { + if let Some(delta) = &choice.delta { + if let Some(piece) = &delta.content { + if !piece.is_empty() { + if ttft_ms.is_none() { + ttft_ms = Some(started.elapsed().as_millis().min(u128::from(u32::MAX)) as u32); + } + content.push_str(piece); + } + } + } + } + } + Err(e) => { + if parse_error.is_none() { + parse_error = Some(truncate(&e.to_string(), 200)); + } + } + } + } + } + } + + let latency_ms = started.elapsed().as_millis().min(u128::from(u32::MAX)) as u32; + let success = parse_error.is_none() && saw_done; + let response_shape_ok = if content.is_empty() { None } else { Some(validate_analysis_report_shape(&content)) }; + let tokens_per_sec = completion_tokens.and_then(|n| { + if latency_ms == 0 { None } else { Some(n as f64 / (latency_ms as f64 / 1000.0)) } + }); + + telemetry.insert_benchmark_run(telemetry::BenchmarkRunInput { + ts, + mesh_id, + target_node_id: None, + model, + probe_name: probe.name.into(), + probe_type: "stream".into(), + stream: true, + temperature: probe.settings.temperature, + top_p: probe.settings.top_p, + max_tokens: probe.settings.max_tokens, + prompt_hash: Some(fnv1a_hex(probe.prompt.as_bytes())), + route_kind: None, + success, + status_code, + latency_ms: Some(latency_ms), + ttft_ms, + prompt_tokens, + completion_tokens, + tokens_per_sec, + error_kind: if success { None } else if parse_error.is_some() { Some("stream_parse_error".into()) } else { Some("stream_incomplete".into()) }, + error_message: parse_error, + response_shape_ok, + }).await?; + + Ok(()) +} + +async fn discover_model(client: &reqwest::Client, api_port: u16) -> Option { + let url = format!("http://127.0.0.1:{api_port}/v1/models"); + let resp = client.get(url).send().await.ok()?; + if !resp.status().is_success() { + return None; + } + let text = resp.text().await.ok()?; + let parsed: ModelsListResponse = serde_json::from_str(&text).ok()?; + parsed.data.into_iter().map(|m| m.id).next() +} + +fn validate_analysis_report_shape(content: &str) -> bool { + let c = content.to_lowercase(); + (c.contains("executive summary") || c.contains("summary")) + && c.contains("risk") + && c.contains("recommend") + && c.contains("conclusion") +} + +fn truncate(s: &str, max: usize) -> String { + if s.len() <= max { + s.to_string() + } else { + s[..max].to_string() + } +} + +fn fnv1a_hex(bytes: &[u8]) -> String { + let mut hash: u64 = 0xcbf29ce484222325; + for b in bytes { + hash ^= u64::from(*b); + hash = hash.wrapping_mul(0x100000001b3); + } + format!("{hash:016x}") +} + +fn unix_ts_now() -> i64 { + std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .unwrap_or_default() + .as_secs() as i64 +} + +#[derive(Debug, Deserialize)] +struct ModelsListResponse { + data: Vec, +} + +#[derive(Debug, Deserialize)] +struct ModelListItem { + id: String, +} + +#[derive(Debug, Deserialize)] +struct ChatCompletionResponse { + choices: Vec, + #[serde(default)] + usage: Option, +} + +#[derive(Debug, Deserialize)] +struct ChatChoice { + message: ChatMessage, +} + +#[derive(Debug, Deserialize)] +struct ChatMessage { + #[serde(default)] + content: Option, +} + +#[derive(Debug, Deserialize)] +struct ChatUsage { + prompt_tokens: u32, + completion_tokens: u32, +} + +#[derive(Debug, Deserialize)] +struct ChatCompletionStreamChunk { + #[serde(default)] + choices: Vec, + #[serde(default)] + usage: Option, +} + +#[derive(Debug, Deserialize)] +struct ChatStreamChoice { + #[serde(default)] + delta: Option, +} + +#[derive(Debug, Deserialize)] +struct ChatStreamDelta { + #[serde(default)] + content: Option, +} diff --git a/mesh-llm/src/console.html b/mesh-llm/src/console.html index 2f75a3a9d..e974b2027 100644 --- a/mesh-llm/src/console.html +++ b/mesh-llm/src/console.html @@ -162,7 +162,7 @@ /* Status line */ .status-line { - font-size: 0.7rem; + font-size: 0.78rem; color: #666; margin-bottom: 0.6rem; } @@ -173,7 +173,7 @@ /* Section labels */ .section-label { - font-size: 0.6rem; + font-size: 0.68rem; text-transform: uppercase; letter-spacing: 0.05em; color: #444; @@ -190,7 +190,7 @@ display: flex; align-items: center; gap: 0.4rem; - font-size: 0.65rem; + font-size: 0.72rem; padding: 0.15rem 0; } @@ -220,7 +220,7 @@ .model-info { color: #444; - font-size: 0.55rem; + font-size: 0.62rem; margin-left: auto; } @@ -232,7 +232,7 @@ } .agent-name { - font-size: 0.6rem; + font-size: 0.68rem; font-weight: 600; color: #555; min-width: 36px; @@ -241,7 +241,7 @@ .cmd-sm { font-family: monospace; - font-size: 0.55rem; + font-size: 0.62rem; color: #aaa; background: #080808; border: 1px solid #1a1a1a; @@ -258,7 +258,7 @@ } .join-details>summary { - font-size: 0.6rem; + font-size: 0.68rem; color: #444; cursor: pointer; } @@ -285,13 +285,13 @@ } .discover-name { - font-size: 0.6rem; + font-size: 0.68rem; font-weight: 600; color: #6cf; } .discover-region { - font-size: 0.5rem; + font-size: 0.58rem; color: #666; background: #111; padding: 0.1rem 0.25rem; @@ -299,12 +299,12 @@ } .discover-stats { - font-size: 0.55rem; + font-size: 0.62rem; color: #888; } .discover-models { - font-size: 0.55rem; + font-size: 0.62rem; color: #aaa; font-family: monospace; white-space: nowrap; @@ -321,7 +321,7 @@ } details.panel-section > summary { - font-size: 0.6rem; + font-size: 0.7rem; font-weight: 600; letter-spacing: 0.04em; color: #888; @@ -352,7 +352,7 @@ } .panel-desc { - font-size: 0.6rem; + font-size: 0.7rem; color: #777; line-height: 1.5; margin-bottom: 0.5rem; @@ -363,7 +363,7 @@ } .panel-label { - font-size: 0.55rem; + font-size: 0.62rem; font-weight: 600; color: #666; text-transform: uppercase; @@ -373,7 +373,7 @@ .panel-val { font-family: monospace; - font-size: 0.6rem; + font-size: 0.68rem; color: #ccc; background: #0a0a0a; border: 1px solid #1a1a1a; @@ -386,7 +386,7 @@ .panel-tok { font-family: monospace; - font-size: 0.5rem; + font-size: 0.6rem; color: #999; background: #0a0a0a; border: 1px solid #1a1a1a; @@ -404,7 +404,7 @@ padding: 0.25rem 0.6rem; color: #e0e0e0; cursor: pointer; - font-size: 0.55rem; + font-size: 0.62rem; white-space: nowrap; flex-shrink: 0; } @@ -418,7 +418,7 @@ padding: 0.2rem 0.5rem; color: #aaa; cursor: pointer; - font-size: 0.55rem; + font-size: 0.62rem; white-space: nowrap; flex-shrink: 0; } @@ -544,10 +544,125 @@ margin-bottom: 0.6rem; } + .mmeta { + color: #666; + font-size: 0.72rem; + margin-top: -0.35rem; + margin-bottom: 0.7rem; + font-family: monospace; + } + .hidden { display: none !important; } + .metrics-grid { + display: grid; + grid-template-columns: repeat(2, minmax(0, 1fr)); + gap: 0.35rem; + margin: 0.35rem 0 0.55rem; + } + + .metric-card { + border: 1px solid #1a1a1a; + border-radius: 4px; + background: #090909; + padding: 0.35rem 0.45rem; + } + + .metric-k { + color: #555; + font-size: 0.58rem; + text-transform: uppercase; + letter-spacing: 0.04em; + margin-bottom: 0.1rem; + } + + .metric-v { + color: #ccc; + font-family: monospace; + font-size: 0.84rem; + } + + .metric-sub { + color: #666; + font-size: 0.58rem; + margin-top: 0.1rem; + } + + .telemetry-chart { + border: 1px solid #1a1a1a; + border-radius: 4px; + background: #090909; + padding: 0.3rem 0.35rem; + margin: 0.35rem 0; + } + + .telemetry-chart .label { + color: #666; + font-size: 0.6rem; + margin-bottom: 0.2rem; + display: flex; + justify-content: space-between; + } + + .telemetry-chart svg { + width: 100%; + height: 36px; + display: block; + } + + .telemetry-table { + width: 100%; + border-collapse: collapse; + font-size: 0.64rem; + margin-top: 0.3rem; + } + + .telemetry-table th { + color: #555; + text-align: left; + font-weight: 600; + border-bottom: 1px solid #1a1a1a; + padding: 0.2rem 0.18rem; + font-size: 0.56rem; + text-transform: uppercase; + letter-spacing: 0.04em; + } + + .telemetry-table td { + color: #aaa; + border-bottom: 1px solid #111; + padding: 0.22rem 0.18rem; + font-family: monospace; + vertical-align: middle; + } + + .telemetry-row { + cursor: pointer; + } + + .telemetry-row:hover td { + background: #0d0d0d; + } + + .telemetry-row.sel td { + background: #0d141a; + color: #cce9ff; + } + + .fresh-dot { + display: inline-block; + width: 5px; + height: 5px; + border-radius: 50%; + margin-right: 0.28rem; + background: #444; + } + .fresh-dot.ok { background: #4c4; } + .fresh-dot.warn { background: #cc4; } + .fresh-dot.stale { background: #c66; } + @media (max-width: 720px) { .main { flex-direction: column; @@ -569,6 +684,10 @@ width: 100%; min-height: 50vh; } + + .metrics-grid { + grid-template-columns: 1fr; + } } @@ -592,6 +711,7 @@

mesh·llm

+
@@ -610,6 +730,8 @@

mesh·llm

- \ No newline at end of file + diff --git a/mesh-llm/src/main.rs b/mesh-llm/src/main.rs index 3e9b68bc1..94d3e30be 100644 --- a/mesh-llm/src/main.rs +++ b/mesh-llm/src/main.rs @@ -1,4 +1,5 @@ mod api; +mod benchmark; mod download; mod election; mod launch; @@ -6,6 +7,7 @@ mod mesh; mod nostr; mod proxy; mod rewrite; +mod telemetry; mod tunnel; use anyhow::{Context, Result}; @@ -133,6 +135,18 @@ struct Cli { /// Nostr relay URLs for publishing/discovery (default: damus, nos.lol, nostr.band). #[arg(long)] nostr_relay: Vec, + + /// Disable the synthetic benchmark probe runner (enabled by default). + #[arg(long)] + no_benchmark: bool, + + /// Benchmark probe cadence in seconds (default: 60). + #[arg(long, default_value = "60")] + benchmark_interval_secs: u64, + + /// Optional model override for benchmark probes (defaults to local served model or first /v1/models entry). + #[arg(long)] + benchmark_model: Option, } #[derive(Subcommand, Debug)] @@ -732,6 +746,8 @@ async fn run_auto(mut cli: Cli, resolved_models: Vec, requested_model_n // Start mesh node let (node, channels) = mesh::Node::start(NodeRole::Worker, &cli.relay, cli.bind_port, cli.max_vram).await?; + let telemetry = telemetry::Telemetry::new(node.id().fmt_short().to_string()).await?; + node.set_telemetry(telemetry.clone()); node.start_accepting(); let token = node.invite_token(); @@ -809,9 +825,10 @@ async fn run_auto(mut cli: Cli, resolved_models: Vec, requested_model_n let mut bootstrap_listener_tx = if !cli.join.is_empty() { let (stop_tx, stop_rx) = tokio::sync::mpsc::channel::>(1); let boot_node = node.clone(); + let boot_telemetry = telemetry.clone(); let boot_port = api_port; tokio::spawn(async move { - bootstrap_proxy(boot_node, boot_port, stop_rx, cli.listen_all).await; + bootstrap_proxy(boot_node, boot_telemetry, boot_port, stop_rx, cli.listen_all).await; }); Some(stop_tx) } else { @@ -924,16 +941,27 @@ async fn run_auto(mut cli: Cli, resolved_models: Vec, requested_model_n // API proxy: model-aware routing let proxy_node = node.clone(); + let proxy_telemetry = telemetry.clone(); let proxy_rx = target_rx.clone(); tokio::spawn(async move { - api_proxy(proxy_node, api_port, proxy_rx, drop_tx, existing_listener, cli.listen_all).await; + api_proxy(proxy_node, proxy_telemetry, api_port, proxy_rx, drop_tx, existing_listener, cli.listen_all).await; }); + if !cli.no_benchmark { + benchmark::spawn_probe_loop( + node.clone(), + telemetry.clone(), + api_port, + cli.benchmark_interval_secs, + cli.benchmark_model.clone(), + ); + } + // Console (optional) let model_name_for_console = model_name.clone(); let console_state = if let Some(cport) = console_port { let model_size_bytes = election::total_model_bytes(&model); - let cs = api::MeshApi::new(node.clone(), model_name_for_console.clone(), api_port, model_size_bytes); + let cs = api::MeshApi::new(node.clone(), telemetry.clone(), model_name_for_console.clone(), api_port, model_size_bytes); cs.set_nostr_relays(nostr_relays(&cli.nostr_relay)).await; if let Some(draft) = &cli.draft { let dn = draft.file_stem().unwrap_or_default().to_string_lossy().to_string(); @@ -1099,8 +1127,10 @@ async fn run_idle(cli: Cli, _bin_dir: PathBuf) -> Result<()> { // Start a dormant node just for the console let (node, _channels) = mesh::Node::start(NodeRole::Worker, &cli.relay, cli.bind_port, cli.max_vram).await?; node.set_available_models(local_models).await; + let telemetry = telemetry::Telemetry::new(node.id().fmt_short().to_string()).await?; + node.set_telemetry(telemetry.clone()); - let cs = api::MeshApi::new(node.clone(), "(idle)".into(), cli.port, 0); + let cs = api::MeshApi::new(node.clone(), telemetry, "(idle)".into(), cli.port, 0); cs.set_nostr_relays(nostr_relays(&cli.nostr_relay)).await; cs.update(false, false).await; let cs2 = cs.clone(); @@ -1121,6 +1151,8 @@ async fn run_idle(cli: Cli, _bin_dir: PathBuf) -> Result<()> { /// Returns Ok(None) on clean shutdown. async fn run_passive(cli: &Cli, node: mesh::Node, is_client: bool) -> Result> { let local_port = cli.port; + let telemetry = telemetry::Telemetry::new(node.id().fmt_short().to_string()).await?; + node.set_telemetry(telemetry.clone()); // Nostr publishing (if --publish, for idle GPU nodes advertising capacity) if cli.publish && !is_client { @@ -1164,7 +1196,7 @@ async fn run_passive(cli: &Cli, node: mesh::Node, is_client: bool) -> Result Result(1); @@ -1221,7 +1263,8 @@ async fn run_passive(cli: &Cli, node: mesh::Node, is_client: bool) -> Result { eprintln!("⬆️ Standby promoting to serve: {model_name}"); @@ -1239,7 +1282,15 @@ async fn run_passive(cli: &Cli, node: mesh::Node, is_client: bool) -> Result, drop_tx: tokio::sync::mpsc::UnboundedSender, existing_listener: Option, listen_all: bool) { +async fn api_proxy( + node: mesh::Node, + telemetry: telemetry::Telemetry, + port: u16, + target_rx: tokio::sync::watch::Receiver, + drop_tx: tokio::sync::mpsc::UnboundedSender, + existing_listener: Option, + listen_all: bool, +) { let listener = match existing_listener { Some(l) => l, None => { @@ -1265,6 +1316,7 @@ async fn api_proxy(node: mesh::Node, port: u16, target_rx: tokio::sync::watch::R let node = node.clone(); let drop_tx = drop_tx.clone(); + let telemetry = telemetry.clone(); tokio::spawn(async move { // Read the HTTP request to extract the model name let mut buf = vec![0u8; 32768]; @@ -1302,7 +1354,7 @@ async fn api_proxy(node: mesh::Node, port: u16, target_rx: tokio::sync::watch::R first_available_target(&targets) }; - proxy::route_to_target(node, tcp_stream, target).await; + proxy::route_to_target(node, tcp_stream, target, telemetry).await; } Err(_) => return, }; @@ -1314,6 +1366,7 @@ async fn api_proxy(node: mesh::Node, port: u16, target_rx: tokio::sync::watch::R /// Returns the TcpListener when signaled to stop (so api_proxy can take it over). async fn bootstrap_proxy( node: mesh::Node, + telemetry: telemetry::Telemetry, port: u16, mut stop_rx: tokio::sync::mpsc::Receiver>, listen_all: bool, @@ -1338,7 +1391,8 @@ async fn bootstrap_proxy( }; let _ = tcp_stream.set_nodelay(true); let node = node.clone(); - tokio::spawn(proxy::handle_mesh_request(node, tcp_stream, true)); + let telemetry = telemetry.clone(); + tokio::spawn(proxy::handle_mesh_request(node, tcp_stream, true, telemetry)); } resp_tx = stop_rx.recv() => { // Hand over listener to api_proxy diff --git a/mesh-llm/src/mesh.rs b/mesh-llm/src/mesh.rs index 4354756a0..3a3b21c20 100644 --- a/mesh-llm/src/mesh.rs +++ b/mesh-llm/src/mesh.rs @@ -3,6 +3,7 @@ //! Single ALPN, single connection per peer. Bi-streams multiplexed by //! first byte: 0x01 = gossip, 0x02 = tunnel (RPC), 0x03 = tunnel map, 0x04 = tunnel (HTTP). +use crate::telemetry; use anyhow::Result; use base64::Engine; use iroh::{Endpoint, EndpointAddr, EndpointId, SecretKey}; @@ -72,6 +73,9 @@ struct PeerAnnouncement { /// Generated once by the originator, propagated via gossip. #[serde(default)] mesh_id: Option, + /// Compact recent telemetry summaries (local historian rows for this node). + #[serde(default)] + telemetry_summaries: Vec, } #[derive(Debug, Clone)] @@ -344,6 +348,7 @@ pub struct Node { pub peer_change_rx: watch::Receiver, tunnel_tx: tokio::sync::mpsc::Sender<(iroh::endpoint::SendStream, iroh::endpoint::RecvStream)>, tunnel_http_tx: tokio::sync::mpsc::Sender<(iroh::endpoint::SendStream, iroh::endpoint::RecvStream)>, + telemetry: Arc>>, } struct MeshState { @@ -465,6 +470,7 @@ impl Node { peer_change_rx, tunnel_tx, tunnel_http_tx, + telemetry: Arc::new(std::sync::Mutex::new(None)), }; // Accept loop starts but waits for start_accepting() before processing connections. @@ -706,6 +712,10 @@ impl Node { self.requested_models.lock().await.clone() } + pub fn set_telemetry(&self, telemetry: telemetry::Telemetry) { + *self.telemetry.lock().unwrap() = Some(telemetry); + } + /// Get all models the mesh has ever wanted — survives peer removal. pub async fn mesh_wanted_models(&self) -> std::collections::HashSet { self.state.lock().await.mesh_wanted.clone() @@ -1509,6 +1519,7 @@ impl Node { let mut buf = vec![0u8; len]; recv.read_exact(&mut buf).await?; let their_announcements: Vec = serde_json::from_slice(&buf)?; + self.ingest_gossip_telemetry(&their_announcements).await; // Wait for stream to fully close, then small delay for accept_bi to re-arm let _ = recv.read_to_end(0).await; @@ -1555,6 +1566,7 @@ impl Node { let mut buf = vec![0u8; len]; recv.read_exact(&mut buf).await?; let their_announcements: Vec = serde_json::from_slice(&buf)?; + self.ingest_gossip_telemetry(&their_announcements).await; // Send our announcements let our_announcements = self.collect_announcements().await; @@ -1718,6 +1730,12 @@ impl Node { } async fn collect_announcements(&self) -> Vec { + let telemetry_handle = { self.telemetry.lock().unwrap().clone() }; + let telemetry_summaries = if let Some(t) = telemetry_handle { + t.latest_summaries_for_gossip().await.unwrap_or_default() + } else { + Vec::new() + }; let state = self.state.lock().await; let my_role = self.role.lock().await.clone(); let my_models = self.models.lock().await.clone(); @@ -1738,6 +1756,7 @@ impl Node { requested_models: p.requested_models.clone(), request_rates: p.request_rates.clone(), mesh_id: my_mesh_id.clone(), + telemetry_summaries: Vec::new(), }) .collect(); let my_rates = self.snapshot_request_rates(); @@ -1752,9 +1771,25 @@ impl Node { requested_models: my_requested, request_rates: my_rates, mesh_id: my_mesh_id, + telemetry_summaries, }); announcements } + + async fn ingest_gossip_telemetry(&self, anns: &[PeerAnnouncement]) { + let telemetry_handle = { self.telemetry.lock().unwrap().clone() }; + let Some(t) = telemetry_handle else { return; }; + let mut summaries = Vec::new(); + for ann in anns { + summaries.extend(ann.telemetry_summaries.iter().cloned()); + } + if summaries.is_empty() { + return; + } + if let Err(e) = t.upsert_peer_summaries(&summaries).await { + tracing::debug!("Telemetry gossip upsert failed: {e}"); + } + } } /// Generate a mesh ID for a new mesh. @@ -1850,5 +1885,3 @@ async fn load_or_create_key() -> Result { tracing::info!("Generated new key, saved to {}", key_path.display()); Ok(key) } - - diff --git a/mesh-llm/src/proxy.rs b/mesh-llm/src/proxy.rs index 5a3868824..df63c482c 100644 --- a/mesh-llm/src/proxy.rs +++ b/mesh-llm/src/proxy.rs @@ -3,7 +3,7 @@ //! Used by the API proxy (port 9337), bootstrap proxy, and passive mode. //! All inference traffic flows through these functions. -use crate::{election, mesh, tunnel}; +use crate::{election, mesh, telemetry, tunnel}; use anyhow::Result; use tokio::io::AsyncWriteExt; use tokio::net::TcpStream; @@ -54,7 +54,12 @@ pub fn is_drop_request(buf: &[u8]) -> bool { /// by model name (or falls back to any host), and tunnels the request via QUIC. /// /// Set `track_demand` to record requests for demand-based rebalancing. -pub async fn handle_mesh_request(node: mesh::Node, tcp_stream: TcpStream, track_demand: bool) { +pub async fn handle_mesh_request( + node: mesh::Node, + tcp_stream: TcpStream, + track_demand: bool, + telemetry: telemetry::Telemetry, +) { let mut buf = vec![0u8; 32768]; let (n, model_name) = match peek_request(&tcp_stream, &mut buf).await { Ok(v) => v, @@ -86,22 +91,29 @@ pub async fn handle_mesh_request(node: mesh::Node, tcp_stream: TcpStream, track_ None => match node.any_host().await { Some(p) => p.id, None => { + let span = telemetry.start_request(telemetry::RouteKind::Remote); let _ = send_503(tcp_stream).await; + span.finish(false); return; } }, }; // Tunnel via QUIC + let span = telemetry.start_request(telemetry::RouteKind::Remote); match node.open_http_tunnel(target_host).await { Ok((quic_send, quic_recv)) => { if let Err(e) = tunnel::relay_tcp_via_quic(tcp_stream, quic_send, quic_recv).await { tracing::debug!("HTTP tunnel relay ended: {e}"); + span.finish(false); + return; } + span.finish(true); } Err(e) => { tracing::warn!("Failed to tunnel to host {}: {e}", target_host.fmt_short()); let _ = send_503(tcp_stream).await; + span.finish(false); } } } @@ -109,37 +121,54 @@ pub async fn handle_mesh_request(node: mesh::Node, tcp_stream: TcpStream, track_ /// Route a request to a known inference target (local llama-server or remote host). /// /// Used by the API proxy after election has determined the target. -pub async fn route_to_target(node: mesh::Node, tcp_stream: TcpStream, target: election::InferenceTarget) { +pub async fn route_to_target( + node: mesh::Node, + tcp_stream: TcpStream, + target: election::InferenceTarget, + telemetry: telemetry::Telemetry, +) { match target { election::InferenceTarget::Local(llama_port) => { + let span = telemetry.start_request(telemetry::RouteKind::Local); match TcpStream::connect(format!("127.0.0.1:{llama_port}")).await { Ok(upstream) => { let _ = upstream.set_nodelay(true); if let Err(e) = tunnel::relay_tcp_streams(tcp_stream, upstream).await { tracing::debug!("API proxy (local) ended: {e}"); + span.finish(false); + return; } + span.finish(true); } Err(e) => { tracing::warn!("API proxy: can't reach llama-server on {llama_port}: {e}"); let _ = send_503(tcp_stream).await; + span.finish(false); } } } election::InferenceTarget::Remote(host_id) => { + let span = telemetry.start_request(telemetry::RouteKind::Remote); match node.open_http_tunnel(host_id).await { Ok((quic_send, quic_recv)) => { if let Err(e) = tunnel::relay_tcp_via_quic(tcp_stream, quic_send, quic_recv).await { tracing::debug!("API proxy (remote) ended: {e}"); + span.finish(false); + return; } + span.finish(true); } Err(e) => { tracing::warn!("API proxy: can't tunnel to host {}: {e}", host_id.fmt_short()); let _ = send_503(tcp_stream).await; + span.finish(false); } } } election::InferenceTarget::None => { + let span = telemetry.start_request(telemetry::RouteKind::Remote); let _ = send_503(tcp_stream).await; + span.finish(false); } } } diff --git a/mesh-llm/src/telemetry.rs b/mesh-llm/src/telemetry.rs new file mode 100644 index 000000000..77ea8cb06 --- /dev/null +++ b/mesh-llm/src/telemetry.rs @@ -0,0 +1,840 @@ +use crate::{benchmark, tunnel}; +use anyhow::Result; +use rusqlite::{params, Connection}; +use serde::Serialize; +use std::path::PathBuf; +use std::sync::{Arc, Mutex}; +use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH}; + +#[derive(Clone)] +pub struct Telemetry { + inner: Arc, +} + +struct Inner { + node_id: String, + db_path: PathBuf, + state: Mutex, +} + +struct State { + current_minute: i64, + requests: u64, + errors: u64, + requests_local: u64, + requests_remote: u64, + request_time_ms_total: u64, + active_requests: u32, + active_requests_peak: u32, + latency_ms_samples: Vec, + tunnel_bytes_minute_start: u64, +} + +#[derive(Clone, Copy)] +pub enum RouteKind { + Local, + Remote, +} + +pub struct RequestSpan { + telemetry: Telemetry, + route: RouteKind, + started: Instant, + finished: bool, +} + +#[derive(Serialize, Clone)] +pub struct LiveSnapshot { + pub ts_minute: i64, + pub node_id: String, + pub active_requests: u32, + pub active_requests_peak: u32, + pub requests: u64, + pub errors: u64, + pub requests_local: u64, + pub requests_remote: u64, + pub request_time_ms_total: u64, + pub utilization_pct: f64, + pub tunnel_bytes_total: u64, +} + +#[derive(Serialize, Clone)] +pub struct NodeMetricRow { + pub ts_minute: i64, + pub source_node_id: String, + pub requests: u64, + pub errors: u64, + pub requests_local: u64, + pub requests_remote: u64, + pub request_time_ms_total: u64, + pub utilization_pct: f64, + pub active_requests_peak: u32, + pub latency_p50_ms: Option, + pub latency_p95_ms: Option, + pub latency_p99_ms: Option, + pub tunnel_bytes_total: u64, + pub observed_at: i64, +} + +#[derive(Debug, Clone, Serialize, serde::Deserialize)] +pub struct NodeMetricSummary { + pub ts_minute: i64, + pub source_node_id: String, + pub requests: u64, + pub errors: u64, + pub requests_local: u64, + pub requests_remote: u64, + pub request_time_ms_total: u64, + pub utilization_pct: f64, + pub active_requests_peak: u32, + pub latency_p50_ms: Option, + pub latency_p95_ms: Option, + pub latency_p99_ms: Option, + pub tunnel_bytes_total: u64, + pub observed_at: i64, +} + +#[derive(Serialize, Clone)] +pub struct RollupMetricRow { + pub ts_minute: i64, + pub node_count: u64, + pub requests: u64, + pub errors: u64, + pub requests_local: u64, + pub requests_remote: u64, + pub request_time_ms_total: u64, + pub utilization_pct_avg_nodes: f64, + pub active_requests_peak_max: u32, + pub tunnel_bytes_total: u64, + pub latency_p95_ms_avg_nodes: Option, + pub latency_p95_ms_max_nodes: Option, +} + +#[derive(Serialize, Clone)] +pub struct BenchmarkRunRow { + pub ts: i64, + pub source_node_id: String, + pub model: String, + pub probe_name: String, + pub probe_type: String, + pub stream: bool, + pub success: bool, + pub status_code: Option, + pub latency_ms: Option, + pub ttft_ms: Option, + pub completion_tokens: Option, + pub tokens_per_sec: Option, + pub error_kind: Option, + pub observed_at: i64, +} + +#[derive(Debug, Clone)] +#[allow(dead_code)] +pub struct BenchmarkRunInput { + pub ts: i64, + pub mesh_id: Option, + pub target_node_id: Option, + pub model: String, + pub probe_name: String, + pub probe_type: String, + pub stream: bool, + pub temperature: f64, + pub top_p: f64, + pub max_tokens: u32, + pub prompt_hash: Option, + pub route_kind: Option, + pub success: bool, + pub status_code: Option, + pub latency_ms: Option, + pub ttft_ms: Option, + pub prompt_tokens: Option, + pub completion_tokens: Option, + pub tokens_per_sec: Option, + pub error_kind: Option, + pub error_message: Option, + pub response_shape_ok: Option, +} + +struct FlushRow { + ts_minute: i64, + source_node_id: String, + requests: u64, + errors: u64, + requests_local: u64, + requests_remote: u64, + request_time_ms_total: u64, + utilization_pct: f64, + active_requests_peak: u32, + latency_p50_ms: Option, + latency_p95_ms: Option, + latency_p99_ms: Option, + tunnel_bytes_total: u64, + observed_at: i64, +} + +impl Telemetry { + pub async fn new(node_id: String) -> Result { + let db_path = telemetry_db_path(); + let db_path_for_init = db_path.clone(); + tokio::task::spawn_blocking(move || init_db(&db_path_for_init)).await??; + + let now_minute = unix_minute_now(); + let tunnel_bytes = tunnel::bytes_transferred(); + let telemetry = Self { + inner: Arc::new(Inner { + node_id, + db_path, + state: Mutex::new(State { + current_minute: now_minute, + requests: 0, + errors: 0, + requests_local: 0, + requests_remote: 0, + request_time_ms_total: 0, + active_requests: 0, + active_requests_peak: 0, + latency_ms_samples: Vec::new(), + tunnel_bytes_minute_start: tunnel_bytes, + }), + }), + }; + telemetry.start_flush_loop(); + Ok(telemetry) + } + + pub fn start_request(&self, route: RouteKind) -> RequestSpan { + let mut state = self.inner.state.lock().expect("telemetry state mutex poisoned"); + state.active_requests = state.active_requests.saturating_add(1); + state.active_requests_peak = state.active_requests_peak.max(state.active_requests); + RequestSpan { + telemetry: self.clone(), + route, + started: Instant::now(), + finished: false, + } + } + + fn finish_request(&self, route: RouteKind, latency: Duration, success: bool) { + let mut state = self.inner.state.lock().expect("telemetry state mutex poisoned"); + state.requests = state.requests.saturating_add(1); + if !success { + state.errors = state.errors.saturating_add(1); + } + match route { + RouteKind::Local => state.requests_local = state.requests_local.saturating_add(1), + RouteKind::Remote => state.requests_remote = state.requests_remote.saturating_add(1), + } + let ms = latency.as_millis().min(u128::from(u32::MAX)) as u32; + state.request_time_ms_total = state.request_time_ms_total.saturating_add(ms as u64); + state.latency_ms_samples.push(ms); + state.active_requests = state.active_requests.saturating_sub(1); + } + + pub fn snapshot(&self) -> LiveSnapshot { + let tunnel_now = tunnel::bytes_transferred(); + let state = self.inner.state.lock().expect("telemetry state mutex poisoned"); + LiveSnapshot { + ts_minute: state.current_minute, + node_id: self.inner.node_id.clone(), + active_requests: state.active_requests, + active_requests_peak: state.active_requests_peak, + requests: state.requests, + errors: state.errors, + requests_local: state.requests_local, + requests_remote: state.requests_remote, + request_time_ms_total: state.request_time_ms_total, + utilization_pct: state.request_time_ms_total as f64 / 60_000.0 * 100.0, + tunnel_bytes_total: tunnel_now.saturating_sub(state.tunnel_bytes_minute_start), + } + } + + pub fn current_minute_summary(&self) -> NodeMetricSummary { + let tunnel_now = tunnel::bytes_transferred(); + let state = self.inner.state.lock().expect("telemetry state mutex poisoned"); + let row = flush_row_from_state(&self.inner.node_id, &state, tunnel_now); + NodeMetricSummary::from(row) + } + + pub async fn node_history(&self, minutes: u32) -> Result> { + self.node_history_for(self.inner.node_id.clone(), minutes).await + } + + pub async fn node_history_for(&self, node_id: String, minutes: u32) -> Result> { + let db_path = self.inner.db_path.clone(); + tokio::task::spawn_blocking(move || query_node_history(&db_path, &node_id, minutes)) + .await? + } + + pub async fn all_nodes_latest(&self, minutes: u32) -> Result> { + let db_path = self.inner.db_path.clone(); + tokio::task::spawn_blocking(move || query_all_nodes_latest(&db_path, minutes)).await? + } + + pub async fn rollup_history(&self, minutes: u32) -> Result> { + let db_path = self.inner.db_path.clone(); + tokio::task::spawn_blocking(move || query_rollup_history(&db_path, minutes)).await? + } + + pub async fn benchmark_history(&self, minutes: u32, limit: u32) -> Result> { + let db_path = self.inner.db_path.clone(); + tokio::task::spawn_blocking(move || query_benchmark_history(&db_path, minutes, limit)).await? + } + + pub async fn latest_summaries_for_gossip(&self) -> Result> { + let db_path = self.inner.db_path.clone(); + let current = self.current_minute_summary(); + let current_minute = current.ts_minute; + let node_id = self.inner.node_id.clone(); + let mut rows = tokio::task::spawn_blocking(move || { + query_node_history_exact_minutes(&db_path, &node_id, &[current_minute - 1]) + }) + .await??; + rows.push(current); + rows.sort_by_key(|r| r.ts_minute); + rows.dedup_by_key(|r| r.ts_minute); + Ok(rows) + } + + pub async fn upsert_peer_summaries(&self, summaries: &[NodeMetricSummary]) -> Result<()> { + if summaries.is_empty() { + return Ok(()); + } + let db_path = self.inner.db_path.clone(); + let rows = summaries.to_vec(); + tokio::task::spawn_blocking(move || { + for summary in &rows { + insert_node_metric_summary(&db_path, summary)?; + } + Ok::<(), anyhow::Error>(()) + }) + .await??; + Ok(()) + } + + #[allow(dead_code)] + pub async fn insert_benchmark_run(&self, run: BenchmarkRunInput) -> Result<()> { + let db_path = self.inner.db_path.clone(); + let source_node_id = self.inner.node_id.clone(); + tokio::task::spawn_blocking(move || insert_benchmark_run(&db_path, &source_node_id, &run)) + .await??; + Ok(()) + } + + fn start_flush_loop(&self) { + let telemetry = self.clone(); + tokio::spawn(async move { + loop { + tokio::time::sleep(Duration::from_secs(1)).await; + if let Err(e) = telemetry.maybe_flush().await { + tracing::debug!("Telemetry flush failed: {e}"); + } + } + }); + } + + async fn maybe_flush(&self) -> Result<()> { + let now_minute = unix_minute_now(); + let tunnel_now = tunnel::bytes_transferred(); + let row = { + let mut state = self.inner.state.lock().expect("telemetry state mutex poisoned"); + if now_minute <= state.current_minute { + return Ok(()); + } + let row = flush_row_from_state(&self.inner.node_id, &state, tunnel_now); + state.current_minute = now_minute; + state.requests = 0; + state.errors = 0; + state.requests_local = 0; + state.requests_remote = 0; + state.request_time_ms_total = 0; + state.active_requests_peak = state.active_requests; + state.latency_ms_samples.clear(); + state.tunnel_bytes_minute_start = tunnel_now; + row + }; + + let db_path = self.inner.db_path.clone(); + tokio::task::spawn_blocking(move || insert_node_metric(&db_path, &row)).await??; + Ok(()) + } +} + +impl RequestSpan { + pub fn finish(mut self, success: bool) { + if !self.finished { + self.telemetry + .finish_request(self.route, self.started.elapsed(), success); + self.finished = true; + } + } +} + +impl From for NodeMetricSummary { + fn from(row: FlushRow) -> Self { + Self { + ts_minute: row.ts_minute, + source_node_id: row.source_node_id, + requests: row.requests, + errors: row.errors, + requests_local: row.requests_local, + requests_remote: row.requests_remote, + request_time_ms_total: row.request_time_ms_total, + utilization_pct: row.utilization_pct, + active_requests_peak: row.active_requests_peak, + latency_p50_ms: row.latency_p50_ms, + latency_p95_ms: row.latency_p95_ms, + latency_p99_ms: row.latency_p99_ms, + tunnel_bytes_total: row.tunnel_bytes_total, + observed_at: row.observed_at, + } + } +} + +impl Drop for RequestSpan { + fn drop(&mut self) { + if !self.finished { + self.telemetry + .finish_request(self.route, self.started.elapsed(), false); + self.finished = true; + } + } +} + +fn flush_row_from_state(node_id: &str, state: &State, tunnel_now: u64) -> FlushRow { + let mut sorted = state.latency_ms_samples.clone(); + sorted.sort_unstable(); + FlushRow { + ts_minute: state.current_minute, + source_node_id: node_id.to_string(), + requests: state.requests, + errors: state.errors, + requests_local: state.requests_local, + requests_remote: state.requests_remote, + request_time_ms_total: state.request_time_ms_total, + utilization_pct: state.request_time_ms_total as f64 / 60_000.0 * 100.0, + active_requests_peak: state.active_requests_peak, + latency_p50_ms: percentile(&sorted, 0.50), + latency_p95_ms: percentile(&sorted, 0.95), + latency_p99_ms: percentile(&sorted, 0.99), + tunnel_bytes_total: tunnel_now.saturating_sub(state.tunnel_bytes_minute_start), + observed_at: unix_ts_now(), + } +} + +fn percentile(samples: &[u32], q: f64) -> Option { + if samples.is_empty() { + return None; + } + let idx = ((samples.len() - 1) as f64 * q).round() as usize; + samples.get(idx).map(|v| *v as f64) +} + +fn telemetry_db_path() -> PathBuf { + let home = dirs::home_dir().unwrap_or_else(|| PathBuf::from(".")); + let dir = home.join(".meshllm"); + let _ = std::fs::create_dir_all(&dir); + dir.join("metrics.db") +} + +fn unix_ts_now() -> i64 { + SystemTime::now() + .duration_since(UNIX_EPOCH) + .unwrap_or_default() + .as_secs() as i64 +} + +fn unix_minute_now() -> i64 { + unix_ts_now() / 60 +} + +fn init_db(db_path: &PathBuf) -> Result<()> { + let conn = Connection::open(db_path)?; + conn.pragma_update(None, "journal_mode", "WAL")?; + conn.pragma_update(None, "synchronous", "NORMAL")?; + conn.execute_batch( + r#" + CREATE TABLE IF NOT EXISTS node_metrics_1m ( + ts_minute INTEGER NOT NULL, + source_node_id TEXT NOT NULL, + requests INTEGER NOT NULL, + errors INTEGER NOT NULL, + requests_local INTEGER NOT NULL, + requests_remote INTEGER NOT NULL, + request_time_ms_total INTEGER NOT NULL, + utilization_pct REAL NOT NULL, + active_requests_peak INTEGER NOT NULL, + latency_p50_ms REAL, + latency_p95_ms REAL, + latency_p99_ms REAL, + tunnel_bytes_total INTEGER NOT NULL, + observed_at INTEGER NOT NULL, + PRIMARY KEY (ts_minute, source_node_id) + ); + CREATE INDEX IF NOT EXISTS idx_node_metrics_time ON node_metrics_1m (ts_minute); + CREATE INDEX IF NOT EXISTS idx_node_metrics_node_time ON node_metrics_1m (source_node_id, ts_minute); + + CREATE TABLE IF NOT EXISTS benchmark_runs ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + ts INTEGER NOT NULL, + source_node_id TEXT NOT NULL, + mesh_id TEXT, + target_node_id TEXT, + model TEXT NOT NULL, + probe_name TEXT NOT NULL, + probe_type TEXT NOT NULL, + stream INTEGER NOT NULL, + temperature REAL NOT NULL, + top_p REAL NOT NULL, + max_tokens INTEGER NOT NULL, + prompt_hash TEXT, + route_kind TEXT, + success INTEGER NOT NULL, + status_code INTEGER, + latency_ms INTEGER, + ttft_ms INTEGER, + prompt_tokens INTEGER, + completion_tokens INTEGER, + tokens_per_sec REAL, + error_kind TEXT, + error_message TEXT, + response_shape_ok INTEGER, + settings_json TEXT, + observed_at INTEGER NOT NULL + ); + CREATE INDEX IF NOT EXISTS idx_benchmark_runs_time ON benchmark_runs (ts); + CREATE INDEX IF NOT EXISTS idx_benchmark_runs_node_time ON benchmark_runs (source_node_id, ts); + CREATE INDEX IF NOT EXISTS idx_benchmark_runs_probe_time ON benchmark_runs (probe_name, ts); + "#, + )?; + Ok(()) +} + +fn insert_node_metric(db_path: &PathBuf, row: &FlushRow) -> Result<()> { + let conn = Connection::open(db_path)?; + conn.execute( + r#" + INSERT INTO node_metrics_1m ( + ts_minute, source_node_id, requests, errors, requests_local, requests_remote, + request_time_ms_total, utilization_pct, active_requests_peak, + latency_p50_ms, latency_p95_ms, latency_p99_ms, tunnel_bytes_total, observed_at + ) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14) + ON CONFLICT(ts_minute, source_node_id) DO UPDATE SET + requests = excluded.requests, + errors = excluded.errors, + requests_local = excluded.requests_local, + requests_remote = excluded.requests_remote, + request_time_ms_total = excluded.request_time_ms_total, + utilization_pct = excluded.utilization_pct, + active_requests_peak = excluded.active_requests_peak, + latency_p50_ms = excluded.latency_p50_ms, + latency_p95_ms = excluded.latency_p95_ms, + latency_p99_ms = excluded.latency_p99_ms, + tunnel_bytes_total = excluded.tunnel_bytes_total, + observed_at = excluded.observed_at + "#, + params![ + row.ts_minute, + row.source_node_id, + row.requests, + row.errors, + row.requests_local, + row.requests_remote, + row.request_time_ms_total, + row.utilization_pct, + row.active_requests_peak, + row.latency_p50_ms, + row.latency_p95_ms, + row.latency_p99_ms, + row.tunnel_bytes_total, + row.observed_at, + ], + )?; + Ok(()) +} + +fn insert_node_metric_summary(db_path: &PathBuf, row: &NodeMetricSummary) -> Result<()> { + let conn = Connection::open(db_path)?; + conn.execute( + r#" + INSERT INTO node_metrics_1m ( + ts_minute, source_node_id, requests, errors, requests_local, requests_remote, + request_time_ms_total, utilization_pct, active_requests_peak, + latency_p50_ms, latency_p95_ms, latency_p99_ms, tunnel_bytes_total, observed_at + ) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14) + ON CONFLICT(ts_minute, source_node_id) DO UPDATE SET + requests = excluded.requests, + errors = excluded.errors, + requests_local = excluded.requests_local, + requests_remote = excluded.requests_remote, + request_time_ms_total = excluded.request_time_ms_total, + utilization_pct = excluded.utilization_pct, + active_requests_peak = excluded.active_requests_peak, + latency_p50_ms = excluded.latency_p50_ms, + latency_p95_ms = excluded.latency_p95_ms, + latency_p99_ms = excluded.latency_p99_ms, + tunnel_bytes_total = excluded.tunnel_bytes_total, + observed_at = excluded.observed_at + "#, + params![ + row.ts_minute, + row.source_node_id, + row.requests, + row.errors, + row.requests_local, + row.requests_remote, + row.request_time_ms_total, + row.utilization_pct, + row.active_requests_peak, + row.latency_p50_ms, + row.latency_p95_ms, + row.latency_p99_ms, + row.tunnel_bytes_total, + row.observed_at, + ], + )?; + Ok(()) +} + +fn insert_benchmark_run(db_path: &PathBuf, source_node_id: &str, run: &BenchmarkRunInput) -> Result<()> { + let conn = Connection::open(db_path)?; + let settings_json = serde_json::to_string(&benchmark::BenchmarkSettings { + temperature: run.temperature, + top_p: run.top_p, + max_tokens: run.max_tokens, + stream: run.stream, + })?; + conn.execute( + r#" + INSERT INTO benchmark_runs ( + ts, source_node_id, mesh_id, target_node_id, model, probe_name, probe_type, + stream, temperature, top_p, max_tokens, prompt_hash, route_kind, success, status_code, + latency_ms, ttft_ms, prompt_tokens, completion_tokens, tokens_per_sec, error_kind, + error_message, response_shape_ok, settings_json, observed_at + ) VALUES ( + ?1, ?2, ?3, ?4, ?5, ?6, ?7, + ?8, ?9, ?10, ?11, ?12, ?13, ?14, ?15, + ?16, ?17, ?18, ?19, ?20, ?21, + ?22, ?23, ?24, ?25 + ) + "#, + params![ + run.ts, + source_node_id, + run.mesh_id, + run.target_node_id, + run.model, + run.probe_name, + run.probe_type, + if run.stream { 1 } else { 0 }, + run.temperature, + run.top_p, + run.max_tokens, + run.prompt_hash, + run.route_kind, + if run.success { 1 } else { 0 }, + run.status_code, + run.latency_ms, + run.ttft_ms, + run.prompt_tokens, + run.completion_tokens, + run.tokens_per_sec, + run.error_kind, + run.error_message, + run.response_shape_ok.map(|v| if v { 1 } else { 0 }), + settings_json, + unix_ts_now(), + ], + )?; + Ok(()) +} + +fn query_node_history(db_path: &PathBuf, node_id: &str, minutes: u32) -> Result> { + let conn = Connection::open(db_path)?; + let cutoff = unix_minute_now() - minutes as i64; + let mut stmt = conn.prepare( + r#" + SELECT + ts_minute, source_node_id, requests, errors, requests_local, requests_remote, + request_time_ms_total, utilization_pct, active_requests_peak, + latency_p50_ms, latency_p95_ms, latency_p99_ms, tunnel_bytes_total, observed_at + FROM node_metrics_1m + WHERE source_node_id = ?1 AND ts_minute >= ?2 + ORDER BY ts_minute ASC + "#, + )?; + let rows = stmt + .query_map(params![node_id, cutoff], map_node_row)? + .collect::>>()?; + Ok(rows) +} + +fn query_node_history_exact_minutes( + db_path: &PathBuf, + node_id: &str, + minutes: &[i64], +) -> Result> { + if minutes.is_empty() { + return Ok(Vec::new()); + } + let conn = Connection::open(db_path)?; + let mut out = Vec::new(); + for minute in minutes { + let mut stmt = conn.prepare( + r#" + SELECT + ts_minute, source_node_id, requests, errors, requests_local, requests_remote, + request_time_ms_total, utilization_pct, active_requests_peak, + latency_p50_ms, latency_p95_ms, latency_p99_ms, tunnel_bytes_total, observed_at + FROM node_metrics_1m + WHERE source_node_id = ?1 AND ts_minute = ?2 + LIMIT 1 + "#, + )?; + let mut rows = stmt.query(params![node_id, minute])?; + if let Some(row) = rows.next()? { + out.push(NodeMetricSummary { + ts_minute: row.get(0)?, + source_node_id: row.get(1)?, + requests: row.get(2)?, + errors: row.get(3)?, + requests_local: row.get(4)?, + requests_remote: row.get(5)?, + request_time_ms_total: row.get(6)?, + utilization_pct: row.get(7)?, + active_requests_peak: row.get::<_, u32>(8)?, + latency_p50_ms: row.get(9)?, + latency_p95_ms: row.get(10)?, + latency_p99_ms: row.get(11)?, + tunnel_bytes_total: row.get(12)?, + observed_at: row.get(13)?, + }); + } + } + Ok(out) +} + +fn query_all_nodes_latest(db_path: &PathBuf, minutes: u32) -> Result> { + let conn = Connection::open(db_path)?; + let cutoff = unix_minute_now() - minutes as i64; + let mut stmt = conn.prepare( + r#" + SELECT m.ts_minute, m.source_node_id, m.requests, m.errors, m.requests_local, m.requests_remote, + m.request_time_ms_total, m.utilization_pct, m.active_requests_peak, + m.latency_p50_ms, m.latency_p95_ms, m.latency_p99_ms, m.tunnel_bytes_total, m.observed_at + FROM node_metrics_1m m + INNER JOIN ( + SELECT source_node_id, MAX(ts_minute) AS max_ts + FROM node_metrics_1m + WHERE ts_minute >= ?1 + GROUP BY source_node_id + ) latest + ON m.source_node_id = latest.source_node_id AND m.ts_minute = latest.max_ts + ORDER BY m.source_node_id ASC + "#, + )?; + let rows = stmt + .query_map(params![cutoff], map_node_row)? + .collect::>>()?; + Ok(rows) +} + +fn query_rollup_history(db_path: &PathBuf, minutes: u32) -> Result> { + let conn = Connection::open(db_path)?; + let cutoff = unix_minute_now() - minutes as i64; + let mut stmt = conn.prepare( + r#" + SELECT + ts_minute, + COUNT(*) AS node_count, + COALESCE(SUM(requests), 0) AS requests, + COALESCE(SUM(errors), 0) AS errors, + COALESCE(SUM(requests_local), 0) AS requests_local, + COALESCE(SUM(requests_remote), 0) AS requests_remote, + COALESCE(SUM(request_time_ms_total), 0) AS request_time_ms_total, + COALESCE(AVG(utilization_pct), 0.0) AS utilization_pct_avg_nodes, + COALESCE(MAX(active_requests_peak), 0) AS active_requests_peak_max, + COALESCE(SUM(tunnel_bytes_total), 0) AS tunnel_bytes_total, + AVG(latency_p95_ms) AS latency_p95_ms_avg_nodes, + MAX(latency_p95_ms) AS latency_p95_ms_max_nodes + FROM node_metrics_1m + WHERE ts_minute >= ?1 + GROUP BY ts_minute + ORDER BY ts_minute ASC + "#, + )?; + let rows = stmt + .query_map(params![cutoff], |row| { + Ok(RollupMetricRow { + ts_minute: row.get(0)?, + node_count: row.get(1)?, + requests: row.get(2)?, + errors: row.get(3)?, + requests_local: row.get(4)?, + requests_remote: row.get(5)?, + request_time_ms_total: row.get(6)?, + utilization_pct_avg_nodes: row.get(7)?, + active_requests_peak_max: row.get::<_, u32>(8)?, + tunnel_bytes_total: row.get(9)?, + latency_p95_ms_avg_nodes: row.get(10)?, + latency_p95_ms_max_nodes: row.get(11)?, + }) + })? + .collect::>>()?; + Ok(rows) +} + +fn query_benchmark_history(db_path: &PathBuf, minutes: u32, limit: u32) -> Result> { + let conn = Connection::open(db_path)?; + let cutoff = unix_ts_now() - (minutes as i64 * 60); + let mut stmt = conn.prepare( + r#" + SELECT ts, source_node_id, model, probe_name, probe_type, stream, success, status_code, + latency_ms, ttft_ms, completion_tokens, tokens_per_sec, error_kind, observed_at + FROM benchmark_runs + WHERE ts >= ?1 + ORDER BY ts DESC + LIMIT ?2 + "#, + )?; + let rows = stmt.query_map(params![cutoff, limit], |row| { + let status_code_i: Option = row.get(7)?; + let status_code = status_code_i.and_then(|v| u16::try_from(v).ok()); + Ok(BenchmarkRunRow { + ts: row.get(0)?, + source_node_id: row.get(1)?, + model: row.get(2)?, + probe_name: row.get(3)?, + probe_type: row.get(4)?, + stream: row.get::<_, i64>(5)? != 0, + success: row.get::<_, i64>(6)? != 0, + status_code, + latency_ms: row.get::<_, Option>(8)?, + ttft_ms: row.get::<_, Option>(9)?, + completion_tokens: row.get::<_, Option>(10)?, + tokens_per_sec: row.get::<_, Option>(11)?, + error_kind: row.get(12)?, + observed_at: row.get(13)?, + }) + })? + .collect::>>()?; + Ok(rows) +} + +fn map_node_row(row: &rusqlite::Row<'_>) -> rusqlite::Result { + Ok(NodeMetricRow { + ts_minute: row.get(0)?, + source_node_id: row.get(1)?, + requests: row.get(2)?, + errors: row.get(3)?, + requests_local: row.get(4)?, + requests_remote: row.get(5)?, + request_time_ms_total: row.get(6)?, + utilization_pct: row.get(7)?, + active_requests_peak: row.get::<_, u32>(8)?, + latency_p50_ms: row.get(9)?, + latency_p95_ms: row.get(10)?, + latency_p99_ms: row.get(11)?, + tunnel_bytes_total: row.get(12)?, + observed_at: row.get(13)?, + }) +}