use serde_json::{json, Value};
use std::io::Write as _;
use std::path::PathBuf;
use std::time::Instant;
#[derive(Clone, Copy, PartialEq, Eq, Hash, Debug)]
pub enum Phase {
Signaling,
Presence,
Establishing,
Ready,
L2Open,
Up,
}
impl Phase {
fn name(self) -> &'static str {
match self {
Phase::Signaling => "signaling",
Phase::Presence => "presence",
Phase::Establishing => "establishing",
Phase::Ready => "ready",
Phase::L2Open => "l2open",
Phase::Up => "up",
}
}
}
pub fn budget_ms(p: Phase) -> u64 {
match p {
Phase::Signaling => 1200,
Phase::Presence => 2500,
Phase::Establishing => 3000,
Phase::Ready => 500,
Phase::L2Open => 800,
Phase::Up => 0,
}
}
pub fn over_budget(p: Phase, elapsed_ms: u64) -> bool {
let b = budget_ms(p);
b > 0 && elapsed_ms > b
}
fn correlation_id() -> String {
crate::fresh_secret()[..8].to_string()
}
pub fn peer_hash_from_secret(secret: &str) -> String {
crate::channel_of(secret)[..10].to_string()
}
pub struct Attempt {
id: String,
server: String,
peer: String,
role: &'static str,
start: Instant,
phase: Phase,
phase_at: Instant,
timings: Vec<PhaseTiming>,
}
#[derive(Clone, Copy, Debug)]
pub struct PhaseTiming {
pub phase: Phase,
pub dur_ms: u64,
pub over_budget: bool,
}
impl Phase {
pub fn label(self) -> &'static str {
self.name()
}
}
impl Attempt {
pub fn new(server: &str, peer_hash: &str, role: &'static str) -> Attempt {
let now = Instant::now();
let a = Attempt {
id: correlation_id(),
server: server.to_string(),
peer: peer_hash.to_string(),
role,
start: now,
phase: Phase::Signaling,
phase_at: now,
timings: Vec::new(),
};
a.emit("start", json!({ "phase": a.phase.name() }));
a
}
pub fn enter(&mut self, next: Phase) {
let now = Instant::now();
let dur_ms = now.duration_since(self.phase_at).as_millis() as u64;
let prev = self.phase;
let ob = over_budget(prev, dur_ms);
self.timings.push(PhaseTiming { phase: prev, dur_ms, over_budget: ob });
self.emit(
"phase",
json!({
"phase": next.name(),
"prev": prev.name(),
"dur_ms": dur_ms,
"over_budget": ob,
}),
);
self.phase = next;
self.phase_at = now;
}
pub fn up(&mut self, route: &str, transport: &str) {
let now = Instant::now();
let dur_ms = now.duration_since(self.phase_at).as_millis() as u64;
let prev = self.phase;
if prev != Phase::Up {
self.timings.push(PhaseTiming { phase: prev, dur_ms, over_budget: over_budget(prev, dur_ms) });
}
let total_ms = self.start.elapsed().as_millis() as u64;
self.phase = Phase::Up;
self.phase_at = now;
self.emit(
"up",
json!({
"phase": Phase::Up.name(),
"total_ms": total_ms,
"route": route,
"transport": transport,
}),
);
}
pub fn fail(&mut self, reason: &str) {
let total_ms = self.start.elapsed().as_millis() as u64;
self.emit(
"fail",
json!({
"last_phase": self.phase.name(),
"total_ms": total_ms,
"reason": reason,
}),
);
}
pub fn stall(&mut self, phase: Phase, elapsed_ms: u64) {
self.emit(
"stall",
json!({
"phase": phase.name(),
"elapsed_ms": elapsed_ms,
}),
);
}
pub fn timings(&self) -> &[PhaseTiming] {
&self.timings
}
pub fn total_ms(&self) -> u64 {
self.start.elapsed().as_millis() as u64
}
fn emit(&self, ev: &str, mut fields: Value) {
let obj = fields.as_object_mut().expect("diag fields are an object");
obj.insert("ev".into(), json!(ev));
obj.insert("src".into(), json!("cli"));
obj.insert("role".into(), json!(self.role));
obj.insert("peer".into(), json!(self.peer));
obj.insert("id".into(), json!(self.id));
write_jsonl(&fields);
beacon(&self.server, fields);
}
}
const MAX_JSONL_BYTES: u64 = 512 * 1024;
fn diag_path() -> PathBuf {
crate::platform::Paths::config_path("diag.jsonl")
}
fn write_jsonl(v: &Value) {
let path = diag_path();
if let Some(dir) = path.parent() {
let _ = std::fs::create_dir_all(dir);
}
let over = std::fs::metadata(&path).map(|m| m.len() > MAX_JSONL_BYTES).unwrap_or(false);
let mut opts = std::fs::OpenOptions::new();
opts.create(true).write(true);
if over {
opts.truncate(true);
} else {
opts.append(true);
}
if let Ok(mut f) = opts.open(&path) {
let mut line = v.to_string();
line.push('\n');
let _ = f.write_all(line.as_bytes());
}
}
fn beacon(server: &str, body: Value) {
let url = format!("{server}/api/telemetry");
tokio::spawn(async move {
let Ok(client) = reqwest::Client::builder()
.timeout(std::time::Duration::from_secs(2))
.build()
else {
return;
};
let _ = client.post(&url).json(&body).send().await;
});
}
#[derive(Debug, Default, Clone)]
pub struct Summary {
pub considered: usize,
pub ups: usize,
pub fails: usize,
pub median_total_ms: Option<u64>,
pub worst_phase: Option<(Phase, usize)>,
pub spans_with_stall: usize,
}
fn phase_from_label(s: &str) -> Option<Phase> {
match s {
"signaling" => Some(Phase::Signaling),
"presence" => Some(Phase::Presence),
"establishing" => Some(Phase::Establishing),
"ready" => Some(Phase::Ready),
"l2open" => Some(Phase::L2Open),
"up" => Some(Phase::Up),
_ => None,
}
}
pub fn summarize(limit: usize) -> Summary {
let path = diag_path();
let raw = std::fs::read_to_string(&path).unwrap_or_default();
summarize_lines(&raw, limit)
}
pub fn summarize_lines(text: &str, limit: usize) -> Summary {
let mut order: Vec<String> = Vec::new();
let mut spans: std::collections::HashMap<String, SpanAcc> = std::collections::HashMap::new();
for line in text.lines() {
let line = line.trim();
if line.is_empty() {
continue;
}
let Ok(v) = serde_json::from_str::<Value>(line) else { continue };
let id = match v["id"].as_str() {
Some(i) => i.to_string(),
None => continue,
};
let acc = spans.entry(id.clone()).or_insert_with(|| {
order.push(id.clone());
SpanAcc::default()
});
match v["ev"].as_str() {
Some("phase") => {
if v["over_budget"].as_bool().unwrap_or(false) {
if let Some(p) = v["prev"].as_str().and_then(phase_from_label) {
acc.over_budget_phases.push(p);
}
}
}
Some("stall") => acc.had_stall = true,
Some("up") => {
acc.terminal = Some(Terminal::Up);
acc.total_ms = v["total_ms"].as_u64();
}
Some("fail") => {
acc.terminal = Some(Terminal::Fail);
acc.total_ms = v["total_ms"].as_u64();
}
_ => {}
}
}
let mut terminal: Vec<&SpanAcc> = order
.iter()
.rev()
.filter_map(|id| spans.get(id))
.filter(|s| s.terminal.is_some())
.take(limit)
.collect();
let considered = terminal.len();
let ups = terminal.iter().filter(|s| matches!(s.terminal, Some(Terminal::Up))).count();
let fails = terminal.iter().filter(|s| matches!(s.terminal, Some(Terminal::Fail))).count();
let spans_with_stall = terminal.iter().filter(|s| s.had_stall).count();
let mut totals: Vec<u64> = terminal.iter().filter_map(|s| s.total_ms).collect();
totals.sort_unstable();
let median_total_ms = median(&totals);
let mut tally: std::collections::HashMap<Phase, usize> = std::collections::HashMap::new();
for s in terminal.iter_mut() {
for p in &s.over_budget_phases {
*tally.entry(*p).or_default() += 1;
}
}
let worst_phase = tally
.into_iter()
.max_by_key(|(_, c)| *c)
.map(|(p, c)| (p, c));
Summary { considered, ups, fails, median_total_ms, worst_phase, spans_with_stall }
}
pub fn latest_span_ladder() -> Option<(Vec<PhaseTiming>, Phase)> {
let raw = std::fs::read_to_string(diag_path()).ok()?;
latest_span_ladder_lines(&raw)
}
pub fn latest_span_ladder_lines(text: &str) -> Option<(Vec<PhaseTiming>, Phase)> {
let mut last_id: Option<String> = None;
for line in text.lines() {
if let Ok(v) = serde_json::from_str::<Value>(line.trim()) {
if let Some(id) = v["id"].as_str() {
if last_id.as_deref() != Some(id) {
last_id = Some(id.to_string());
}
}
}
}
let id = last_id?;
let mut timings = Vec::new();
let mut last_phase = Phase::Signaling;
for line in text.lines() {
let Ok(v) = serde_json::from_str::<Value>(line.trim()) else { continue };
if v["id"].as_str() != Some(id.as_str()) {
continue;
}
match v["ev"].as_str() {
Some("phase") => {
if let Some(prev) = v["prev"].as_str().and_then(phase_from_label) {
let dur_ms = v["dur_ms"].as_u64().unwrap_or(0);
timings.push(PhaseTiming { phase: prev, dur_ms, over_budget: over_budget(prev, dur_ms) });
}
if let Some(p) = v["phase"].as_str().and_then(phase_from_label) {
last_phase = p;
}
}
Some("fail") => {
if let Some(p) = v["last_phase"].as_str().and_then(phase_from_label) {
last_phase = p;
}
}
Some("up") => last_phase = Phase::Up,
_ => {}
}
}
Some((timings, last_phase))
}
fn median(sorted: &[u64]) -> Option<u64> {
if sorted.is_empty() {
return None;
}
Some(sorted[(sorted.len() - 1) / 2])
}
#[derive(Default)]
struct SpanAcc {
terminal: Option<Terminal>,
total_ms: Option<u64>,
over_budget_phases: Vec<Phase>,
had_stall: bool,
}
enum Terminal {
Up,
Fail,
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn budgets_are_positive_except_up() {
for p in [Phase::Signaling, Phase::Presence, Phase::Establishing, Phase::Ready, Phase::L2Open] {
assert!(budget_ms(p) > 0, "{:?} should have a budget", p);
}
assert_eq!(budget_ms(Phase::Up), 0);
}
#[test]
fn critical_path_budget_under_5s() {
let critical: u64 = budget_ms(Phase::Signaling) + budget_ms(Phase::Establishing) + budget_ms(Phase::L2Open);
assert!(critical <= 5000, "critical-path budget {critical}ms exceeds 5s target");
}
#[test]
fn over_budget_is_strict_and_phase_aware() {
assert!(!over_budget(Phase::Signaling, 0));
assert!(!over_budget(Phase::Signaling, budget_ms(Phase::Signaling)));
assert!(over_budget(Phase::Signaling, budget_ms(Phase::Signaling) + 1));
assert!(!over_budget(Phase::Up, 10_000_000));
}
#[test]
fn establishing_budget_covers_slow_but_real_ice() {
assert!(budget_ms(Phase::Establishing) >= 2000);
}
fn line(id: &str, ev: &str, extra: &str) -> String {
if extra.is_empty() {
format!(r#"{{"id":"{id}","ev":"{ev}"}}"#)
} else {
format!(r#"{{"id":"{id}","ev":"{ev}",{extra}}}"#)
}
}
#[test]
fn summarize_empty_is_all_zero() {
let s = summarize_lines("", 10);
assert_eq!(s.considered, 0);
assert_eq!(s.ups, 0);
assert_eq!(s.fails, 0);
assert!(s.median_total_ms.is_none());
assert!(s.worst_phase.is_none());
assert_eq!(s.spans_with_stall, 0);
}
#[test]
fn summarize_counts_ups_fails_and_median() {
let mut text = String::new();
text += &line("a", "start", "");
text.push('\n');
text += &line("a", "up", r#""total_ms":1000"#);
text.push('\n');
text += &line("b", "up", r#""total_ms":3000"#);
text.push('\n');
text += &line("c", "fail", r#""total_ms":5000"#);
text.push('\n');
let s = summarize_lines(&text, 10);
assert_eq!(s.considered, 3);
assert_eq!(s.ups, 2);
assert_eq!(s.fails, 1);
assert_eq!(s.median_total_ms, Some(3000));
}
#[test]
fn summarize_skips_in_flight_spans() {
let text = format!("{}\n", line("x", "start", ""));
let s = summarize_lines(&text, 10);
assert_eq!(s.considered, 0);
}
#[test]
fn summarize_finds_worst_phase_and_stall_rate() {
let mut text = String::new();
text += &line("a", "phase", r#""prev":"establishing","over_budget":true"#);
text.push('\n');
text += &line("a", "stall", r#""phase":"establishing""#);
text.push('\n');
text += &line("a", "up", r#""total_ms":4000"#);
text.push('\n');
text += &line("b", "phase", r#""prev":"establishing","over_budget":true"#);
text.push('\n');
text += &line("b", "up", r#""total_ms":4200"#);
text.push('\n');
text += &line("c", "phase", r#""prev":"presence","over_budget":true"#);
text.push('\n');
text += &line("c", "fail", r#""total_ms":9000"#);
text.push('\n');
let s = summarize_lines(&text, 10);
assert_eq!(s.considered, 3);
assert_eq!(s.spans_with_stall, 1);
let (p, c) = s.worst_phase.expect("a worst phase");
assert_eq!(p, Phase::Establishing);
assert_eq!(c, 2);
}
#[test]
fn summarize_limit_keeps_newest() {
let mut text = String::new();
for (id, t) in [("a", 100u64), ("b", 200), ("c", 300), ("d", 400)] {
text += &line(id, "up", &format!(r#""total_ms":{t}"#));
text.push('\n');
}
let s = summarize_lines(&text, 2);
assert_eq!(s.considered, 2);
assert_eq!(s.median_total_ms, Some(300));
}
#[test]
fn median_handles_odd_even_empty() {
assert_eq!(median(&[]), None);
assert_eq!(median(&[5]), Some(5));
assert_eq!(median(&[1, 2, 3]), Some(2));
assert_eq!(median(&[1, 2, 3, 4]), Some(2)); }
#[test]
fn peer_hash_is_short_and_not_the_input() {
let secret = "abcdef0123456789";
let h = peer_hash_from_secret(secret);
assert_eq!(h.len(), 10);
assert!(!secret.contains(&h));
assert_eq!(h, peer_hash_from_secret(secret));
}
}