use std::collections::HashMap;
use std::io::{self, Write};
use crate::output::Theme;
use crate::strutil::truncate_max_bytes;
use crate::types::*;
const HISTORY_SIZE: usize = 16;
const MIN_SAMPLES: usize = 3;
const HIGH_RATIO: f64 = 0.80;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Verdict {
Healthy,
TxStalled,
RxStarved,
}
impl Verdict {
pub fn as_str(self) -> &'static str {
match self {
Self::Healthy => "HEALTHY",
Self::TxStalled => "TX-STALLED",
Self::RxStarved => "RX-STARVED",
}
}
fn is_stalled(self) -> bool {
!matches!(self, Self::Healthy)
}
}
#[derive(Debug, Clone, Copy)]
struct StallSample {
send_q: u64,
recv_q: u64,
send_buf: u64,
recv_buf: u64,
timestamp: i64,
}
struct StallEntry {
pid: i32,
fd: i32,
command: String,
proto: String,
peer: String,
history: Vec<StallSample>,
verdict: Verdict,
seen: bool,
}
pub struct StallDetector {
table: HashMap<(i32, i32), StallEntry>,
iteration: u64,
}
impl StallDetector {
pub fn new() -> Self {
Self {
table: HashMap::new(),
iteration: 0,
}
}
pub fn update(&mut self, procs: &[Process]) {
self.iteration += 1;
let now = chrono::Utc::now().timestamp();
for entry in self.table.values_mut() {
entry.seen = false;
}
for p in procs {
for f in &p.files {
let FdName::Number(fd) = f.fd else {
continue;
};
let Some(si) = &f.socket_info else {
continue;
};
let (Some(send_q), Some(recv_q)) = (si.send_queue, si.recv_queue) else {
continue;
};
let sample = StallSample {
send_q,
recv_q,
send_buf: si.send_buf_size.unwrap_or(0),
recv_buf: si.recv_buf_size.unwrap_or(0),
timestamp: now,
};
let key = (p.pid, fd);
let entry = self.table.entry(key).or_insert_with(|| StallEntry {
pid: p.pid,
fd,
command: p.command.clone(),
proto: si.protocol.clone(),
peer: f.name.clone(),
history: Vec::with_capacity(HISTORY_SIZE),
verdict: Verdict::Healthy,
seen: false,
});
entry.seen = true;
if entry.command != p.command {
entry.command = p.command.clone();
entry.history.clear();
entry.verdict = Verdict::Healthy;
}
entry.proto = si.protocol.clone();
entry.peer = f.name.clone();
if entry.history.len() >= HISTORY_SIZE {
entry.history.remove(0);
}
entry.history.push(sample);
entry.verdict = classify(&entry.history);
}
}
self.table.retain(|_, e| e.seen || e.verdict.is_stalled());
}
pub fn report(&self, theme: &Theme) {
let out = io::stdout();
let mut out = out.lock();
let mut stalled: Vec<&StallEntry> = self
.table
.values()
.filter(|e| e.seen && e.verdict.is_stalled())
.collect();
stalled.sort_by_key(|e| (e.pid, e.fd));
let scanned = self.table.values().filter(|e| e.seen).count();
let _ = writeln!(
out,
"\n{bold}═══ lsofrs socket backpressure ═══{reset}",
bold = theme.bold(),
reset = theme.reset(),
);
let _ = writeln!(
out,
" iteration: {} | sockets: {} | stalled: {red}{}{reset}\n",
self.iteration,
scanned,
stalled.len(),
red = if stalled.is_empty() { "" } else { theme.red() },
reset = theme.reset(),
);
if stalled.is_empty() {
let _ = writeln!(
out,
" {green}No stalled sockets detected.{reset}",
green = theme.green(),
reset = theme.reset(),
);
let _ = writeln!(out);
return;
}
let _ = writeln!(
out,
" {hdr}{bold}{:>7} {:>5} {:<12} {:<11} {:>7} PEER{reset}",
"PID",
"FD",
"COMMAND",
"VERDICT",
"PROTO",
hdr = theme.hdr_bg(),
bold = theme.bold(),
reset = theme.reset(),
);
for e in &stalled {
let cmd = truncate_max_bytes(&e.command, 12);
let peer = truncate_max_bytes(&e.peer, 40);
let color = match e.verdict {
Verdict::TxStalled => theme.red(),
Verdict::RxStarved => theme.yellow(),
Verdict::Healthy => theme.green(),
};
let _ = writeln!(
out,
" {red}{:>7}{reset} {:>5} {cyan}{:<12}{reset} {color}{:<11}{reset} {:>7} {}",
e.pid,
e.fd,
cmd,
e.verdict.as_str(),
e.proto,
peer,
red = theme.red(),
cyan = theme.cyan(),
color = color,
reset = theme.reset(),
);
}
let _ = writeln!(out);
}
}
impl Default for StallDetector {
fn default() -> Self {
Self::new()
}
}
fn tx_severity(w: &[StallSample]) -> Option<f64> {
if w.iter().all(|s| s.send_buf > 0) {
let ratios: Vec<f64> = w
.iter()
.map(|s| s.send_q as f64 / s.send_buf as f64)
.collect();
let min_r = ratios.iter().copied().fold(f64::INFINITY, f64::min);
let first = *ratios.first().unwrap();
let last = *ratios.last().unwrap();
if min_r >= HIGH_RATIO && last >= first - f64::EPSILON {
return Some(min_r);
}
None
} else {
let qs: Vec<u64> = w.iter().map(|s| s.send_q).collect();
let all_pos = qs.iter().all(|&q| q > 0);
let non_decreasing = qs.windows(2).all(|p| p[1] >= p[0]);
if all_pos && non_decreasing {
let max = *qs.iter().max().unwrap() as f64;
return Some(if max > 0.0 {
*qs.last().unwrap() as f64 / max
} else {
0.0
});
}
None
}
}
fn rx_severity(w: &[StallSample]) -> Option<f64> {
if w.iter().all(|s| s.recv_buf > 0) {
let ratios: Vec<f64> = w
.iter()
.map(|s| s.recv_q as f64 / s.recv_buf as f64)
.collect();
let min_r = ratios.iter().copied().fold(f64::INFINITY, f64::min);
let first = *ratios.first().unwrap();
let last = *ratios.last().unwrap();
if min_r >= HIGH_RATIO && last >= first - f64::EPSILON {
return Some(min_r);
}
None
} else {
let qs: Vec<u64> = w.iter().map(|s| s.recv_q).collect();
let all_pos = qs.iter().all(|&q| q > 0);
let non_decreasing = qs.windows(2).all(|p| p[1] >= p[0]);
let grew = *qs.last().unwrap() > *qs.first().unwrap();
if all_pos && non_decreasing && grew {
let max = *qs.iter().max().unwrap() as f64;
return Some(if max > 0.0 {
*qs.last().unwrap() as f64 / max
} else {
0.0
});
}
None
}
}
fn classify(history: &[StallSample]) -> Verdict {
let n = history.len();
if n < MIN_SAMPLES {
return Verdict::Healthy;
}
let w = &history[n - MIN_SAMPLES..];
match (tx_severity(w), rx_severity(w)) {
(Some(tx), Some(rx)) => {
if tx >= rx {
Verdict::TxStalled
} else {
Verdict::RxStarved
}
}
(Some(_), None) => Verdict::TxStalled,
(None, Some(_)) => Verdict::RxStarved,
(None, None) => Verdict::Healthy,
}
}
#[cfg(test)]
mod tests {
use super::*;
fn sample(send_q: u64, recv_q: u64, send_buf: u64, recv_buf: u64) -> StallSample {
StallSample {
send_q,
recv_q,
send_buf,
recv_buf,
timestamp: 0,
}
}
#[test]
fn classifier_directional_verdicts() {
let warmup = vec![sample(98, 2, 100, 100), sample(97, 3, 100, 100)];
assert_eq!(classify(&warmup), Verdict::Healthy);
let tx = vec![
sample(95, 4, 100, 100),
sample(96, 5, 100, 100),
sample(97, 3, 100, 100),
sample(98, 4, 100, 100),
];
assert_eq!(classify(&tx), Verdict::TxStalled);
let rx = vec![
sample(3, 85, 100, 100),
sample(2, 90, 100, 100),
sample(4, 95, 100, 100),
sample(1, 98, 100, 100),
];
assert_eq!(classify(&rx), Verdict::RxStarved);
let draining = vec![
sample(95, 4, 100, 100),
sample(50, 3, 100, 100),
sample(10, 5, 100, 100),
sample(4, 2, 100, 100),
];
assert_eq!(classify(&draining), Verdict::Healthy);
let idle = vec![
sample(10, 8, 100, 100),
sample(12, 6, 100, 100),
sample(9, 7, 100, 100),
];
assert_eq!(classify(&idle), Verdict::Healthy);
let linux_rx = vec![
sample(0, 1000, 0, 0),
sample(0, 4000, 0, 0),
sample(0, 9000, 0, 0),
];
assert_eq!(classify(&linux_rx), Verdict::RxStarved);
}
}