use super::{Health, RowStatus, StatusReport, bookmark_text};
use chrono::{DateTime, Utc};
pub fn age(secs: i64) -> String {
let secs = secs.max(0);
match secs {
s if s < 60 => format!("{s}s"),
s if s < 3_600 => format!("{}m", s / 60),
s if s < 86_400 => format!("{}h", s / 3_600),
s => format!("{}d", s / 86_400),
}
}
fn clip(s: &str, max: usize) -> String {
if s.chars().count() <= max {
return s.to_string();
}
let mut out: String = s.chars().take(max.saturating_sub(1)).collect();
out.push('…');
out
}
fn when(at: DateTime<Utc>, now: DateTime<Utc>) -> String {
format!(
"{} ({})",
at.format("%Y-%m-%d %H:%M"),
age((now - at).num_seconds())
)
}
fn health_label(h: Health) -> String {
match h {
Health::Failed | Health::Degraded => h.as_str().to_ascii_uppercase(),
other => other.as_str().to_string(),
}
}
fn cells(r: &RowStatus, now: DateTime<Utc>) -> [String; 7] {
let dlq = if !r.dlq.configured {
"—".to_string()
} else if r.dlq.readable {
r.dlq.count.to_string()
} else {
"?".to_string()
};
[
r.row.clone(),
health_label(r.health),
r.last_success
.as_ref()
.map(|s| when(s.at, now))
.unwrap_or_else(|| "never".into()),
r.bookmark
.as_ref()
.map(|b| clip(&bookmark_text(b), 28))
.unwrap_or_else(|| "—".into()),
r.lag
.as_ref()
.map(|l| clip(&l.human, 16))
.unwrap_or_else(|| "—".into()),
dlq,
clip(&r.resume, 60),
]
}
pub fn render(report: &StatusReport) -> String {
let now = report.generated_at;
let header = [
"row",
"status",
"last success",
"bookmark",
"lag",
"dlq",
"next run resumes at",
];
let rows: Vec<[String; 7]> = report.rows.iter().map(|r| cells(r, now)).collect();
let mut widths: Vec<usize> = header.iter().map(|h| h.chars().count()).collect();
for r in &rows {
for (i, c) in r.iter().enumerate() {
widths[i] = widths[i].max(c.chars().count());
}
}
let line = |cols: Vec<&str>| -> String {
let mut s = String::from(" ");
for (i, c) in cols.iter().enumerate() {
if i + 1 == cols.len() {
s.push_str(c);
} else {
s.push_str(&format!("{c:<w$} ", w = widths[i]));
}
}
s.trim_end().to_string() + "\n"
};
let state = if report.state.kinds.is_empty() {
"none".to_string()
} else {
report.state.kinds.join(", ")
};
let mut out = format!(
"pipeline {} ({} {}{}) — {} state: {}\n",
report.pipeline,
report.rows.len(),
if report.topology { "sink node" } else { "row" },
if report.rows.len() == 1 { "" } else { "s" },
health_label(report.health),
state
);
if let Some(note) = &report.state.note {
out.push_str(&format!(" note: {note}\n"));
}
out.push_str(&line(header.to_vec()));
let indent = " ".repeat(2 + widths[0] + 2);
for (r, c) in report.rows.iter().zip(&rows) {
out.push_str(&line(c.iter().map(String::as_str).collect()));
for d in details(r, now) {
out.push_str(&format!("{indent}└ {d}\n"));
}
}
out
}
fn details(r: &RowStatus, now: DateTime<Utc>) -> Vec<String> {
let mut out = Vec::new();
if let Some(l) = &r.running {
out.push(format!(
"running: run {} (pid {}{}) since {}",
l.run_id,
l.pid,
l.host
.as_deref()
.map(|h| format!(" on {h}"))
.unwrap_or_default(),
when(l.acquired_at, now)
));
}
if let Some(f) = &r.last_failure
&& r.health == Health::Failed
{
out.push(format!(
"last error: {}{} ({}{}){}",
f.error_kind
.as_deref()
.map(|k| format!("{k}: "))
.unwrap_or_default(),
clip(f.error.as_deref().unwrap_or("failed"), 160),
f.at.format("%Y-%m-%d %H:%M"),
f.run_id
.as_deref()
.map(|id| format!(", run {}", clip(id, 12)))
.unwrap_or_default(),
if r.consecutive_failures > 1 {
format!(" · {} consecutive failures", r.consecutive_failures)
} else {
String::new()
}
));
}
let mut marks = Vec::new();
for v in &r.sla {
marks.push(format!("SLA {}: {}", v.kind, v.message));
}
if let Some(eo) = &r.exactly_once {
let sink = eo
.sink
.as_ref()
.and_then(|s| s.seq)
.map(|s| s.to_string())
.unwrap_or_else(|| "—".into());
marks.push(format!(
"exactly-once: state seq {} · sink seq {} · {:?} → next run trusts the {}",
eo.state_seq, sink, eo.agreement, eo.trusted
));
if let Some(e) = &eo.probe_error {
marks.push(format!("watermark probe failed: {e}"));
}
}
if r.dlq.count > 0 {
marks.push(format!(
"DLQ: {} record(s){}",
r.dlq.count,
r.dlq
.oldest
.map(|o| format!(", oldest {}", when(o, now)))
.unwrap_or_default()
));
}
if let Some(n) = &r.dlq.note {
marks.push(format!("DLQ: {n}"));
}
if let Some(b) = &r.batches
&& b.unclean() > 0
{
marks.push(super::batch_note(b));
}
if let Some(f) = &r.state_format
&& f.status != "current"
{
marks.push(format!(
"state: {} — {}",
f.status,
f.detail.as_deref().unwrap_or_default()
));
}
if let Some(l) = &r.lag {
marks.push(format!(
"lag: {} ({}, {})",
l.human,
if l.measured == "probe" {
"queried now"
} else {
"at the end of the last run"
},
when(l.at, now)
));
}
if let Some(p) = &r.profiling
&& p.drift > 0
{
marks.push(format!(
"profiling: {} drift finding(s) in the latest run",
p.drift
));
}
if r.rollback.undoable_runs > 0 {
marks.push(format!(
"rollback: {} undoable run(s), newest {}",
r.rollback.undoable_runs,
r.rollback.newest.as_deref().unwrap_or("?")
));
}
if let Some(o) = &r.overwrite_staging {
marks.push(format!("overwrite staging: {} — {}", o.state, o.note));
}
for c in &r.children {
marks.push(format!(
"children '{}': {} bookmark(s), {} failed — worst {}",
c.row,
c.bookmarks,
c.failed,
c.worst.as_str()
));
}
for reason in r.reasons.iter().filter(|x| {
x.contains("unknown between runs")
|| x.contains("without releasing")
|| x.contains("in flight")
}) {
marks.push(reason.clone());
}
for e in &r.errors {
marks.push(format!("unreadable: {e}"));
}
out.extend(marks);
out
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn ages_and_clipping() {
assert_eq!(age(-5), "0s");
assert_eq!(age(59), "59s");
assert_eq!(age(120), "2m");
assert_eq!(age(7_200), "2h");
assert_eq!(age(200_000), "2d");
assert_eq!(clip("abcdef", 4), "abc…");
assert_eq!(clip("abc", 4), "abc");
}
}