feat: rampart-tui — live metrics terminal dashboard
- ratatui + crossterm; separate bin reading Prometheus exposition format - counter table (xdp stats, auto-bans, intel), pps/drops history chart, status panel with error banner, q/quit r/refresh - prometheus text parser (labels, escapes, NaN/Inf), Fetcher trait for tests - 152 tests green
This commit is contained in:
parent
6863249ad4
commit
564613b82d
11 changed files with 2024 additions and 15 deletions
1067
Cargo.lock
generated
1067
Cargo.lock
generated
File diff suppressed because it is too large
Load diff
|
|
@ -62,6 +62,10 @@ rand = { version = "0.8", default-features = false, features = ["std", "std_rng"
|
|||
clap = { version = "4", features = ["derive"] }
|
||||
chrono = { version = "0.4", features = ["serde"] }
|
||||
|
||||
# TUI (rampart-tui): лёгкие deps, без feature-gate
|
||||
ratatui = "0.30"
|
||||
crossterm = "0.29"
|
||||
|
||||
axum = "0.8"
|
||||
tower-http = { version = "0.6", features = ["cors"] }
|
||||
jsonwebtoken = "9"
|
||||
|
|
|
|||
64
src/bin/rampart-tui.rs
Normal file
64
src/bin/rampart-tui.rs
Normal file
|
|
@ -0,0 +1,64 @@
|
|||
//! `rampart-tui` — интерактивный терминальный дашборд live-метрик Rampart.
|
||||
|
||||
use std::sync::Arc;
|
||||
use std::time::Duration;
|
||||
|
||||
use clap::Parser;
|
||||
use ratatui::Terminal;
|
||||
use ratatui::backend::CrosstermBackend;
|
||||
use tokio::sync::Notify;
|
||||
|
||||
use rampart::tui::app;
|
||||
use rampart::tui::fetch::{Fetcher, HttpFetcher};
|
||||
use rampart::tui::state::AppState;
|
||||
|
||||
#[derive(Parser)]
|
||||
#[command(
|
||||
name = "rampart-tui",
|
||||
about = "Rampart live metrics TUI (Prometheus /metrics source)"
|
||||
)]
|
||||
struct Cli {
|
||||
/// URL Prometheus metrics endpoint
|
||||
#[arg(long, default_value = "http://127.0.0.1:9090/metrics")]
|
||||
url: String,
|
||||
/// Период опроса, секунды
|
||||
#[arg(long, default_value_t = 1)]
|
||||
refresh_secs: u64,
|
||||
}
|
||||
|
||||
fn main() -> anyhow::Result<()> {
|
||||
let cli = Cli::parse();
|
||||
let Some(fetcher) = HttpFetcher::new(cli.url.clone()) else {
|
||||
anyhow::bail!("failed to build http client");
|
||||
};
|
||||
let runtime = tokio::runtime::Builder::new_current_thread().enable_all().build()?;
|
||||
|
||||
let terminal = ratatui::try_init()?;
|
||||
let result = runtime.block_on(run(
|
||||
terminal,
|
||||
Arc::new(fetcher) as Arc<dyn Fetcher>,
|
||||
cli.url,
|
||||
Duration::from_secs(cli.refresh_secs.max(1)),
|
||||
));
|
||||
ratatui::restore();
|
||||
result
|
||||
}
|
||||
|
||||
async fn run(
|
||||
mut terminal: Terminal<CrosstermBackend<std::io::Stdout>>,
|
||||
fetcher: Arc<dyn Fetcher>,
|
||||
url: String,
|
||||
refresh: Duration,
|
||||
) -> anyhow::Result<()> {
|
||||
let state = Arc::new(AppState::new(url));
|
||||
let notify = Arc::new(Notify::new());
|
||||
tokio::spawn(app::spawn_fetch_loop(
|
||||
fetcher,
|
||||
Arc::clone(&state),
|
||||
refresh,
|
||||
Arc::clone(¬ify),
|
||||
));
|
||||
app::run_ui(&mut terminal, state, notify)
|
||||
.await
|
||||
.map_err(anyhow::Error::from)
|
||||
}
|
||||
|
|
@ -8,4 +8,5 @@ pub mod metrics;
|
|||
pub mod protocol;
|
||||
pub mod store;
|
||||
pub mod traffic;
|
||||
pub mod tui;
|
||||
pub mod xdp;
|
||||
|
|
|
|||
86
src/tui/app.rs
Normal file
86
src/tui/app.rs
Normal file
|
|
@ -0,0 +1,86 @@
|
|||
//! Циклы приложения: fetch-таска (tokio) и UI-цикл (tick + события клавиатуры).
|
||||
|
||||
use std::sync::Arc;
|
||||
use std::time::Duration;
|
||||
|
||||
use tokio::sync::Notify;
|
||||
|
||||
use crate::tui::fetch::{FetchError, Fetcher};
|
||||
use crate::tui::prometheus::parse;
|
||||
use crate::tui::state::AppState;
|
||||
use crate::tui::ui;
|
||||
|
||||
/// Управляемые пользователем команды (клавиши).
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||
pub enum Command {
|
||||
Quit,
|
||||
Refresh,
|
||||
}
|
||||
|
||||
/// Маппинг клавиш: `q`/`Esc` — выход, `r` — обновить немедленно.
|
||||
#[must_use]
|
||||
pub fn handle_key(code: crossterm::event::KeyCode) -> Option<Command> {
|
||||
match code {
|
||||
crossterm::event::KeyCode::Char('q') | crossterm::event::KeyCode::Esc => Some(Command::Quit),
|
||||
crossterm::event::KeyCode::Char('r') => Some(Command::Refresh),
|
||||
_ => None,
|
||||
}
|
||||
}
|
||||
|
||||
/// Приводит результат fetch к виду для [`AppState::apply_fetch`].
|
||||
pub fn normalize(result: Result<String, FetchError>) -> Result<crate::tui::prometheus::MetricSet, String> {
|
||||
match result {
|
||||
Ok(body) => Ok(parse(&body)),
|
||||
Err(FetchError::Network(message)) => Err(message),
|
||||
Err(FetchError::Status(status)) => Err(format!("http status {status}")),
|
||||
}
|
||||
}
|
||||
|
||||
/// Фоновая задача: fetch → [`AppState`]. Ошибки сети пишутся в снапшот баннером.
|
||||
pub async fn spawn_fetch_loop(fetcher: Arc<dyn Fetcher>, state: Arc<AppState>, refresh: Duration, notify: Arc<Notify>) {
|
||||
loop {
|
||||
let started = std::time::Instant::now();
|
||||
let result = normalize(fetcher.fetch().await);
|
||||
let elapsed_secs = started.elapsed().as_secs_f64();
|
||||
state.apply_fetch(result, elapsed_secs);
|
||||
|
||||
tokio::select! {
|
||||
() = tokio::time::sleep(refresh) => {},
|
||||
() = notify.notified() => {},
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// UI-цикл: рисует кадр каждые ~100 мс, обрабатывает клавиши.
|
||||
///
|
||||
/// # Errors
|
||||
/// Возвращает ошибку терминала (draw/poll) — вызывающий восстанавливает TTY.
|
||||
pub async fn run_ui(
|
||||
terminal: &mut ratatui::Terminal<impl ratatui::backend::Backend<Error = std::io::Error>>,
|
||||
state: Arc<AppState>,
|
||||
notify: Arc<Notify>,
|
||||
) -> std::io::Result<()> {
|
||||
loop {
|
||||
terminal.draw(|frame| ui::render(frame, &state.frame_data()))?;
|
||||
if let Some(command) = poll_command(Duration::from_millis(100))? {
|
||||
match command {
|
||||
Command::Quit => return Ok(()),
|
||||
Command::Refresh => notify.notify_one(),
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Неблокирующее ожидание клавиши в течение `timeout`.
|
||||
fn poll_command(timeout: Duration) -> std::io::Result<Option<Command>> {
|
||||
if !crossterm::event::poll(timeout)? {
|
||||
return Ok(None);
|
||||
}
|
||||
if let crossterm::event::Event::Key(key) = crossterm::event::read()? {
|
||||
// Реагируем только на нажатие (не автоповтор/отпускание).
|
||||
if key.kind == crossterm::event::KeyEventKind::Press {
|
||||
return Ok(handle_key(key.code));
|
||||
}
|
||||
}
|
||||
Ok(None)
|
||||
}
|
||||
65
src/tui/fetch.rs
Normal file
65
src/tui/fetch.rs
Normal file
|
|
@ -0,0 +1,65 @@
|
|||
//! Получение тела `/metrics`: трейт [`Fetcher`] для мока в тестах и HTTP-реализация.
|
||||
|
||||
use std::time::Duration;
|
||||
|
||||
use futures::future::BoxFuture;
|
||||
use thiserror::Error;
|
||||
|
||||
/// Ошибка загрузки метрик.
|
||||
#[derive(Debug, Error, Clone)]
|
||||
pub enum FetchError {
|
||||
/// Сетевая ошибка / DNS / таймаут.
|
||||
#[error("network error: {0}")]
|
||||
Network(String),
|
||||
/// Endpoint ответил не-200 статусом.
|
||||
#[error("http status {0}")]
|
||||
Status(u16),
|
||||
}
|
||||
|
||||
/// Абстракция источника exposition-текста Prometheus.
|
||||
pub trait Fetcher: Send + Sync {
|
||||
/// Загружает текущий срез метрик.
|
||||
///
|
||||
/// # Errors
|
||||
/// Возвращает [`FetchError`] при сетевом сбое или плохом HTTP-статусе.
|
||||
fn fetch(&self) -> BoxFuture<'_, Result<String, FetchError>>;
|
||||
}
|
||||
|
||||
/// Реализация [`Fetcher`] поверх reqwest GET.
|
||||
pub struct HttpFetcher {
|
||||
url: String,
|
||||
client: reqwest::Client,
|
||||
}
|
||||
|
||||
impl HttpFetcher {
|
||||
/// Создаёт fetcher с таймаутом запроса 5 секунд.
|
||||
pub fn new(url: impl Into<String>) -> Option<Self> {
|
||||
let client = reqwest::Client::builder()
|
||||
.timeout(Duration::from_secs(5))
|
||||
.build()
|
||||
.ok()?;
|
||||
Some(Self {
|
||||
url: url.into(),
|
||||
client,
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
impl Fetcher for HttpFetcher {
|
||||
fn fetch(&self) -> BoxFuture<'_, Result<String, FetchError>> {
|
||||
let url = self.url.clone();
|
||||
let client = self.client.clone();
|
||||
Box::pin(async move {
|
||||
let response = client
|
||||
.get(&url)
|
||||
.send()
|
||||
.await
|
||||
.map_err(|e| FetchError::Network(e.to_string()))?;
|
||||
let status = response.status();
|
||||
if !status.is_success() {
|
||||
return Err(FetchError::Status(status.as_u16()));
|
||||
}
|
||||
response.text().await.map_err(|e| FetchError::Network(e.to_string()))
|
||||
})
|
||||
}
|
||||
}
|
||||
10
src/tui/mod.rs
Normal file
10
src/tui/mod.rs
Normal file
|
|
@ -0,0 +1,10 @@
|
|||
//! TUI-модуль Rampart: live-метрики edge-ноды/manager в терминале (ratatui).
|
||||
//!
|
||||
//! Используется бинарём `rampart-tui`; источник данных — Prometheus
|
||||
//! exposition endpoint (`/metrics`).
|
||||
|
||||
pub mod app;
|
||||
pub mod fetch;
|
||||
pub mod prometheus;
|
||||
pub mod state;
|
||||
pub mod ui;
|
||||
238
src/tui/prometheus.rs
Normal file
238
src/tui/prometheus.rs
Normal file
|
|
@ -0,0 +1,238 @@
|
|||
//! Минималистичный парсер текстового формата Prometheus (version=0.0.4).
|
||||
//!
|
||||
//! Поддерживает: комментарии (`# HELP`, `# TYPE`, `#`), строки вида
|
||||
//! `metric{label="value",...} value [timestamp]` и `metric value [timestamp]`.
|
||||
//! Значения `NaN`/`+Inf`/`-Inf` парсятся штатно через `f64::from_str`.
|
||||
|
||||
use std::collections::BTreeMap;
|
||||
|
||||
/// Один сэмпл из exposition-формата Prometheus.
|
||||
#[derive(Debug, Clone, PartialEq)]
|
||||
pub struct MetricSample {
|
||||
pub name: String,
|
||||
pub labels: BTreeMap<String, String>,
|
||||
pub value: f64,
|
||||
}
|
||||
|
||||
/// Набор сэмплов одного ответа `/metrics`.
|
||||
#[derive(Debug, Default, Clone, PartialEq)]
|
||||
pub struct MetricSet {
|
||||
samples: Vec<MetricSample>,
|
||||
}
|
||||
|
||||
impl MetricSet {
|
||||
/// Сумма значений всех серий с данным именем (для counter/gauge vec).
|
||||
/// Отсутствующая метрика даёт `0.0`.
|
||||
#[must_use]
|
||||
pub fn sum(&self, name: &str) -> f64 {
|
||||
self.samples.iter().filter(|s| s.name == name).map(|s| s.value).sum()
|
||||
}
|
||||
|
||||
/// Значение серии с именем и одной меткой; `None`, если не найдено.
|
||||
#[must_use]
|
||||
pub fn labeled(&self, name: &str, label: &str, value: &str) -> Option<f64> {
|
||||
self.samples
|
||||
.iter()
|
||||
.find(|s| s.name == name && s.labels.get(label).map(String::as_str) == Some(value))
|
||||
.map(|s| s.value)
|
||||
}
|
||||
|
||||
/// Все сэмплы (для тестов и отладки).
|
||||
#[must_use]
|
||||
pub fn samples(&self) -> &[MetricSample] {
|
||||
&self.samples
|
||||
}
|
||||
}
|
||||
|
||||
/// Парсит тело ответа `/metrics`.
|
||||
///
|
||||
/// Строки с ошибкой формата пропускаются (exposition-поток может быть частично
|
||||
/// повреждён — TUI не должен падать на этом).
|
||||
#[must_use]
|
||||
pub fn parse(input: &str) -> MetricSet {
|
||||
let mut samples = Vec::new();
|
||||
for line in input.lines() {
|
||||
if let Some(sample) = parse_line(line) {
|
||||
samples.push(sample);
|
||||
}
|
||||
}
|
||||
MetricSet { samples }
|
||||
}
|
||||
|
||||
fn parse_line(line: &str) -> Option<MetricSample> {
|
||||
let line = line.trim();
|
||||
if line.is_empty() || line.starts_with('#') {
|
||||
return None;
|
||||
}
|
||||
// Формат: metric[{labels}] value [timestamp]; значения меток могут
|
||||
// содержать пробелы, поэтому голову отрезаем с учётом кавычек.
|
||||
let (head, rest) = split_head_and_rest(line);
|
||||
let value_str = rest.split_whitespace().next()?;
|
||||
// Timestamp после value игнорируем.
|
||||
|
||||
let (name, labels) = split_metric(head)?;
|
||||
let value = parse_value(value_str)?;
|
||||
Some(MetricSample { name, labels, value })
|
||||
}
|
||||
|
||||
/// Отделяет `metric{...}` от хвоста `value [timestamp]`.
|
||||
fn split_head_and_rest(line: &str) -> (&str, &str) {
|
||||
let Some(open) = line.find('{') else {
|
||||
return match line.find(char::is_whitespace) {
|
||||
Some(space) => (&line[..space], &line[space..]),
|
||||
None => (line, ""),
|
||||
};
|
||||
};
|
||||
match find_closing_brace(line, open + 1) {
|
||||
Some(close) => (&line[..=close], &line[close + 1..]),
|
||||
None => match line.find(char::is_whitespace) {
|
||||
Some(space) => (&line[..space.min(open)], &line[space.min(open)..]),
|
||||
None => (line, ""),
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
/// Индекс `}` вне кавычек, начиная с позиции `from`.
|
||||
fn find_closing_brace(line: &str, from: usize) -> Option<usize> {
|
||||
let mut in_quotes = false;
|
||||
let mut escaped = false;
|
||||
for (idx, ch) in line[from..].char_indices() {
|
||||
match ch {
|
||||
'\\' if in_quotes => escaped = !escaped,
|
||||
'"' if !escaped => in_quotes = !in_quotes,
|
||||
'}' if !in_quotes => return Some(from + idx),
|
||||
_ => {},
|
||||
}
|
||||
if ch != '\\' {
|
||||
escaped = false;
|
||||
}
|
||||
}
|
||||
None
|
||||
}
|
||||
|
||||
fn split_metric(head: &str) -> Option<(String, BTreeMap<String, String>)> {
|
||||
match head.find('{') {
|
||||
None => Some((head.to_owned(), BTreeMap::new())),
|
||||
Some(open) => {
|
||||
if !head.ends_with('}') || open == 0 {
|
||||
return None;
|
||||
}
|
||||
let name = head[..open].trim().to_owned();
|
||||
let inner = &head[open + 1..head.len() - 1];
|
||||
Some((name, parse_labels(inner)))
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
fn parse_labels(inner: &str) -> BTreeMap<String, String> {
|
||||
let mut labels = BTreeMap::new();
|
||||
for pair in split_label_pairs(inner) {
|
||||
if let Some(eq) = pair.find('=') {
|
||||
let key = pair[..eq].trim();
|
||||
let raw_value = pair[eq + 1..].trim();
|
||||
if let Some(unquoted) = raw_value.strip_prefix('"').and_then(|v| v.strip_suffix('"')) {
|
||||
labels.insert(key.to_owned(), unescape_label(unquoted));
|
||||
}
|
||||
}
|
||||
}
|
||||
labels
|
||||
}
|
||||
|
||||
/// Разбивает `a="1", b="2"` по запятым вне кавычек.
|
||||
fn split_label_pairs(inner: &str) -> Vec<&str> {
|
||||
let mut parts = Vec::new();
|
||||
let mut in_quotes = false;
|
||||
let mut escaped = false;
|
||||
let mut start = 0;
|
||||
for (idx, ch) in inner.char_indices() {
|
||||
match ch {
|
||||
'\\' if in_quotes => escaped = !escaped,
|
||||
'"' if !escaped => in_quotes = !in_quotes,
|
||||
',' if !in_quotes => {
|
||||
parts.push(&inner[start..idx]);
|
||||
start = idx + 1;
|
||||
},
|
||||
_ => {},
|
||||
}
|
||||
if ch != '\\' {
|
||||
escaped = false;
|
||||
}
|
||||
}
|
||||
parts.push(&inner[start..]);
|
||||
parts
|
||||
}
|
||||
|
||||
fn unescape_label(value: &str) -> String {
|
||||
let mut out = String::with_capacity(value.len());
|
||||
let mut chars = value.chars();
|
||||
while let Some(ch) = chars.next() {
|
||||
if ch == '\\' {
|
||||
match chars.next() {
|
||||
Some('n') => out.push('\n'),
|
||||
Some(next) => out.push(next),
|
||||
None => {},
|
||||
}
|
||||
} else {
|
||||
out.push(ch);
|
||||
}
|
||||
}
|
||||
out
|
||||
}
|
||||
|
||||
fn parse_value(value: &str) -> Option<f64> {
|
||||
match value {
|
||||
"+Inf" | "Inf" => Some(f64::INFINITY),
|
||||
"-Inf" => Some(f64::NEG_INFINITY),
|
||||
"NaN" => Some(f64::NAN),
|
||||
v => v.parse::<f64>().ok(),
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn parses_simple_and_labeled_lines() {
|
||||
let set = parse("up 1\nxdp_total{iface=\"eth0\"} 42\n");
|
||||
assert_eq!(set.sum("up"), 1.0);
|
||||
assert_eq!(set.labeled("xdp_total", "iface", "eth0"), Some(42.0));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn skips_comments_help_type_and_empty_lines() {
|
||||
let input = "# HELP rampart_xdp_total Total\n# TYPE rampart_xdp_total gauge\n\nrampart_xdp_total 7\n";
|
||||
let set = parse(input);
|
||||
assert_eq!(set.samples().len(), 1);
|
||||
assert_eq!(set.sum("rampart_xdp_total"), 7.0);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn ignores_timestamp_and_special_values() {
|
||||
let set = parse("m1 10 1700000000000\nm2 +Inf\nm3 NaN\nbroken_line_no_value\nbad{} \nm4 -2.5e3\n");
|
||||
assert_eq!(set.sum("m1"), 10.0);
|
||||
assert_eq!(set.sum("m2"), f64::INFINITY);
|
||||
assert!(set.labeled("m3", "x", "y").is_none());
|
||||
assert_eq!(set.sum("m4"), -2500.0);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn sums_series_with_same_name() {
|
||||
let input = "bans{reason=\"cps\"} 3\nbans{reason=\"syn\"} 4\nbans 5\n";
|
||||
assert_eq!(parse(input).sum("bans"), 12.0);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn handles_spaces_in_label_values() {
|
||||
let set = parse("m{k=\"a b\", j=\"v w\"} 1\n");
|
||||
assert_eq!(set.labeled("m", "k", "a b"), Some(1.0));
|
||||
assert_eq!(set.labeled("m", "j", "v w"), Some(1.0));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn handles_escaped_quotes_and_backslashes_in_labels() {
|
||||
let set = parse("m{q=\"a\\\"b\", s=\"c\\\\d\"} 1\n");
|
||||
assert_eq!(set.labeled("m", "q", "a\"b"), Some(1.0));
|
||||
assert_eq!(set.labeled("m", "s", "c\\d"), Some(1.0));
|
||||
}
|
||||
}
|
||||
166
src/tui/state.rs
Normal file
166
src/tui/state.rs
Normal file
|
|
@ -0,0 +1,166 @@
|
|||
//! Общее состояние TUI: последний срез метрик, история pps/drops, расчёт rate.
|
||||
|
||||
use std::collections::VecDeque;
|
||||
use std::sync::Mutex;
|
||||
use std::time::Instant;
|
||||
|
||||
use crate::tui::prometheus::MetricSet;
|
||||
|
||||
/// Максимальная длина скользящего окна графика.
|
||||
pub const HISTORY_CAPACITY: usize = 120;
|
||||
|
||||
/// Один сэмпл истории (значения — rates, единиц/сек).
|
||||
#[derive(Debug, Clone, Copy, PartialEq)]
|
||||
pub struct HistoryPoint {
|
||||
pub packets_per_sec: f64,
|
||||
pub drops_per_sec: f64,
|
||||
}
|
||||
|
||||
/// Скользящее окно истории фиксированной ёмкости.
|
||||
#[derive(Debug, Clone)]
|
||||
pub struct History {
|
||||
points: VecDeque<HistoryPoint>,
|
||||
}
|
||||
|
||||
impl Default for History {
|
||||
fn default() -> Self {
|
||||
Self {
|
||||
points: VecDeque::with_capacity(HISTORY_CAPACITY),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl History {
|
||||
/// Добавляет точку, вытесняя старейшую при переполнении окна.
|
||||
pub fn push(&mut self, point: HistoryPoint) {
|
||||
if self.points.len() == HISTORY_CAPACITY {
|
||||
self.points.pop_front();
|
||||
}
|
||||
self.points.push_back(point);
|
||||
}
|
||||
|
||||
/// Точки в хронологическом порядке.
|
||||
#[must_use]
|
||||
pub fn points(&self) -> &VecDeque<HistoryPoint> {
|
||||
&self.points
|
||||
}
|
||||
}
|
||||
|
||||
/// Состояние приложения, общее между fetch-потоком и UI-циклом.
|
||||
#[derive(Debug)]
|
||||
pub struct AppState {
|
||||
url: String,
|
||||
snapshot: Mutex<Snapshot>,
|
||||
history: Mutex<History>,
|
||||
started_at: Instant,
|
||||
}
|
||||
|
||||
impl AppState {
|
||||
/// Создаёт состояние с URL источника метрик (для статус-панели).
|
||||
#[must_use]
|
||||
pub fn new(url: impl Into<String>) -> Self {
|
||||
Self {
|
||||
url: url.into(),
|
||||
snapshot: Mutex::new(Snapshot::default()),
|
||||
history: Mutex::new(History::default()),
|
||||
started_at: Instant::now(),
|
||||
}
|
||||
}
|
||||
|
||||
/// URL источника метрик.
|
||||
#[must_use]
|
||||
pub fn url(&self) -> &str {
|
||||
&self.url
|
||||
}
|
||||
}
|
||||
|
||||
/// Последний успешный/неуспешный результат fetch.
|
||||
#[derive(Debug, Default, Clone)]
|
||||
pub struct Snapshot {
|
||||
pub metrics: Option<MetricSet>,
|
||||
pub fetched_at: Option<Instant>,
|
||||
pub error: Option<String>,
|
||||
}
|
||||
|
||||
impl AppState {
|
||||
/// Записывает результат очередного fetch и обновляет историю rate.
|
||||
pub fn apply_fetch(&self, result: Result<MetricSet, String>, elapsed_secs: f64) {
|
||||
let mut snapshot = match self.snapshot.lock() {
|
||||
Ok(guard) => guard,
|
||||
Err(poisoned) => poisoned.into_inner(),
|
||||
};
|
||||
match result {
|
||||
Ok(metrics) => {
|
||||
let mut history = match self.history.lock() {
|
||||
Ok(guard) => guard,
|
||||
Err(poisoned) => poisoned.into_inner(),
|
||||
};
|
||||
let previous = snapshot.metrics.take();
|
||||
let point = HistoryPoint {
|
||||
packets_per_sec: counter_rate(
|
||||
previous.as_ref().map(|p| p.sum("rampart_xdp_total")),
|
||||
Some(metrics.sum("rampart_xdp_total")),
|
||||
elapsed_secs,
|
||||
),
|
||||
drops_per_sec: counter_rate(
|
||||
previous.as_ref().map(|p| p.sum("rampart_xdp_dropped")),
|
||||
Some(metrics.sum("rampart_xdp_dropped")),
|
||||
elapsed_secs,
|
||||
),
|
||||
};
|
||||
history.push(point);
|
||||
snapshot.metrics = Some(metrics);
|
||||
snapshot.fetched_at = Some(Instant::now());
|
||||
snapshot.error = None;
|
||||
},
|
||||
Err(message) => {
|
||||
snapshot.error = Some(message);
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
/// Копия среза метрик + истории для отрисовки кадра.
|
||||
///
|
||||
/// UI не держит лок — рисуем по снимку.
|
||||
#[must_use]
|
||||
pub fn frame_data(&self) -> FrameData {
|
||||
let snapshot = self.snapshot.lock().unwrap_or_else(std::sync::PoisonError::into_inner);
|
||||
let history = self
|
||||
.history
|
||||
.lock()
|
||||
.unwrap_or_else(std::sync::PoisonError::into_inner)
|
||||
.clone();
|
||||
FrameData {
|
||||
url: self.url.clone(),
|
||||
snapshot: snapshot.clone(),
|
||||
history,
|
||||
uptime_secs: self.started_at.elapsed().as_secs(),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Снимок для одного кадра UI.
|
||||
#[derive(Debug, Default, Clone)]
|
||||
pub struct FrameData {
|
||||
/// URL источника метрик.
|
||||
pub url: String,
|
||||
pub snapshot: Snapshot,
|
||||
pub history: History,
|
||||
/// Время работы TUI (сек).
|
||||
pub uptime_secs: u64,
|
||||
}
|
||||
|
||||
/// Rate счётчика из двух сэмплов: `(current - previous) / dt`.
|
||||
///
|
||||
/// `None` у предыдущего значения (первый сэмпл) даёт 0.0. Сброс счётчика
|
||||
/// (`current < previous`) трактуетcя как рестарт процесса → 0.0.
|
||||
#[must_use]
|
||||
pub fn counter_rate(previous: Option<f64>, current: Option<f64>, elapsed_secs: f64) -> f64 {
|
||||
let (Some(previous), Some(current)) = (previous, current) else {
|
||||
return 0.0;
|
||||
};
|
||||
if !(elapsed_secs > 0.0 && current.is_finite()) || current < previous {
|
||||
return 0.0;
|
||||
}
|
||||
(current - previous) / elapsed_secs
|
||||
}
|
||||
179
src/tui/ui.rs
Normal file
179
src/tui/ui.rs
Normal file
|
|
@ -0,0 +1,179 @@
|
|||
//! Отрисовка кадра ratatui: статус-панель, таблица счётчиков, график pps/drops, помощь.
|
||||
|
||||
use ratatui::Frame;
|
||||
use ratatui::layout::{Constraint, Layout, Rect};
|
||||
use ratatui::style::{Color, Modifier, Style};
|
||||
use ratatui::symbols;
|
||||
use ratatui::text::{Line, Span};
|
||||
use ratatui::widgets::{Axis, Block, Borders, Chart, Dataset, GraphType, Paragraph, Row, Table};
|
||||
|
||||
use crate::tui::state::{FrameData, HistoryPoint};
|
||||
|
||||
/// Имена метрик edge/manager (см. src/metrics.rs и src/xdp/metrics.rs).
|
||||
mod names {
|
||||
pub const XDP_TOTAL: &str = "rampart_xdp_total";
|
||||
pub const XDP_TCP: &str = "rampart_xdp_tcp";
|
||||
pub const XDP_UDP: &str = "rampart_xdp_udp";
|
||||
pub const XDP_PASSED: &str = "rampart_xdp_passed";
|
||||
pub const XDP_DROPPED: &str = "rampart_xdp_dropped";
|
||||
pub const XDP_THROTTLE: &str = "rampart_xdp_syn_throttle";
|
||||
pub const XDP_RATE_LIMIT: &str = "rampart_xdp_rate_limit";
|
||||
pub const XDP_BLACKLIST: &str = "rampart_xdp_blacklist";
|
||||
pub const RATE_LIMIT_HITS: &str = "rampart_rate_limit_hits";
|
||||
pub const AUTO_BANS_TOTAL: &str = "rampart_auto_bans_total";
|
||||
pub const INTEL_DROPS: &str = "rampart_intel_rate_limit_drops_total";
|
||||
pub const INTEL_CPS: &str = "rampart_intel_cps";
|
||||
pub const INTEL_ALERTS: &str = "rampart_intel_alerts_total";
|
||||
pub const ATTACK_STATUS: &str = "rampart_attack_status";
|
||||
}
|
||||
|
||||
/// Рисует весь кадр.
|
||||
pub fn render(frame: &mut Frame, data: &FrameData) {
|
||||
let [header, body, help] =
|
||||
Layout::vertical([Constraint::Length(3), Constraint::Min(3), Constraint::Length(1)]).areas(frame.area());
|
||||
render_header(frame, header, data);
|
||||
let [table_area, chart_area] =
|
||||
Layout::vertical([Constraint::Percentage(55), Constraint::Percentage(45)]).areas(body);
|
||||
render_table(frame, table_area, data);
|
||||
render_chart(frame, chart_area, data);
|
||||
render_help(frame, help);
|
||||
}
|
||||
|
||||
fn status_color(data: &FrameData) -> Color {
|
||||
match data.snapshot.error {
|
||||
Some(_) => Color::Red,
|
||||
None if data.snapshot.fetched_at.is_some() => Color::Green,
|
||||
None => Color::Yellow,
|
||||
}
|
||||
}
|
||||
|
||||
fn render_header(frame: &mut Frame, area: Rect, data: &FrameData) {
|
||||
let dot = Span::styled(
|
||||
"● ",
|
||||
Style::default().fg(status_color(data)).add_modifier(Modifier::BOLD),
|
||||
);
|
||||
let title = Span::styled(" Rampart live metrics", Style::default().add_modifier(Modifier::BOLD));
|
||||
let state_text = match (&data.snapshot.error, data.snapshot.fetched_at) {
|
||||
(Some(err), _) => format!("ERROR: {err}"),
|
||||
(_, Some(at)) => format!("OK, last update {}s ago", at.elapsed().as_secs()),
|
||||
(None, None) => "waiting for first fetch...".to_owned(),
|
||||
};
|
||||
let line = Line::from(vec![
|
||||
dot,
|
||||
title,
|
||||
Span::raw(format!(
|
||||
" {} | uptime {}s | {state_text}",
|
||||
data.url, data.uptime_secs
|
||||
)),
|
||||
]);
|
||||
frame.render_widget(Paragraph::new(line).block(Block::new().borders(Borders::ALL)), area);
|
||||
}
|
||||
|
||||
/// Пары «подпись → имя метрики» для таблицы.
|
||||
const COUNTER_ROWS: &[(&str, &str)] = &[
|
||||
("XDP total", names::XDP_TOTAL),
|
||||
("XDP tcp", names::XDP_TCP),
|
||||
("XDP udp", names::XDP_UDP),
|
||||
("XDP passed", names::XDP_PASSED),
|
||||
("XDP dropped", names::XDP_DROPPED),
|
||||
("SYN throttle", names::XDP_THROTTLE),
|
||||
("Rate limited", names::XDP_RATE_LIMIT),
|
||||
("Blacklisted", names::XDP_BLACKLIST),
|
||||
("HTTP rate-limit hits", names::RATE_LIMIT_HITS),
|
||||
("Auto bans", names::AUTO_BANS_TOTAL),
|
||||
("Intel rl drops", names::INTEL_DROPS),
|
||||
("Intel cps", names::INTEL_CPS),
|
||||
("Intel alerts", names::INTEL_ALERTS),
|
||||
];
|
||||
|
||||
fn attack_status_label(value: Option<f64>) -> String {
|
||||
match value {
|
||||
None => "-".to_owned(),
|
||||
Some(v) if v <= 0.0 => "normal".to_owned(),
|
||||
Some(v) if v < 2.0 => "suspicious".to_owned(),
|
||||
Some(_) => "under attack".to_owned(),
|
||||
}
|
||||
}
|
||||
|
||||
fn render_table(frame: &mut Frame, area: Rect, data: &FrameData) {
|
||||
let rows: Vec<Row> = COUNTER_ROWS
|
||||
.iter()
|
||||
.map(|(label, metric)| {
|
||||
let value = data.snapshot.metrics.as_ref().map(|m| m.sum(metric));
|
||||
Row::new(vec![(*label).to_owned(), fmt_value(value)])
|
||||
})
|
||||
.chain(std::iter::once(Row::new(vec![
|
||||
"Attack status".to_owned(),
|
||||
attack_status_label(data.snapshot.metrics.as_ref().map(|m| m.sum(names::ATTACK_STATUS))),
|
||||
])))
|
||||
.collect();
|
||||
|
||||
let table = Table::new(rows, [Constraint::Percentage(60), Constraint::Percentage(40)])
|
||||
.block(Block::new().borders(Borders::ALL).title(" Counters "))
|
||||
.header(Row::new(vec!["metric", "value"]).style(Style::default().add_modifier(Modifier::BOLD)));
|
||||
frame.render_widget(table, area);
|
||||
}
|
||||
|
||||
#[allow(clippy::float_arithmetic)]
|
||||
fn fmt_value(value: Option<f64>) -> String {
|
||||
match value {
|
||||
None => "-".to_owned(),
|
||||
Some(v) if !v.is_finite() => format!("{v}"),
|
||||
Some(v) if v.fract() == 0.0 => format!("{}", v as u64),
|
||||
Some(v) => format!("{v:.1}"),
|
||||
}
|
||||
}
|
||||
|
||||
fn history_series(history: &[HistoryPoint], pick: impl Fn(&HistoryPoint) -> f64) -> Vec<(f64, f64)> {
|
||||
history
|
||||
.iter()
|
||||
.enumerate()
|
||||
.map(|(idx, point)| (f64::from(idx as u16), pick(point)))
|
||||
.collect()
|
||||
}
|
||||
|
||||
fn render_chart(frame: &mut Frame, area: Rect, data: &FrameData) {
|
||||
let points: Vec<HistoryPoint> = data.history.points().iter().copied().collect();
|
||||
let pps_series = history_series(&points, |p| p.packets_per_sec);
|
||||
let drops_series = history_series(&points, |p| p.drops_per_sec);
|
||||
let pps = Dataset::default()
|
||||
.name("pps")
|
||||
.marker(symbols::Marker::Braille)
|
||||
.graph_type(GraphType::Line)
|
||||
.style(Style::default().fg(Color::Cyan))
|
||||
.data(&pps_series);
|
||||
let drops = Dataset::default()
|
||||
.name("drops/s")
|
||||
.marker(symbols::Marker::Braille)
|
||||
.graph_type(GraphType::Line)
|
||||
.style(Style::default().fg(Color::Red))
|
||||
.data(&drops_series);
|
||||
|
||||
let x_max = points.len().max(2) as f64 - 1.0;
|
||||
let y_max = points
|
||||
.iter()
|
||||
.flat_map(|p| [p.packets_per_sec, p.drops_per_sec])
|
||||
.fold(1.0_f64, f64::max);
|
||||
|
||||
let chart = Chart::new(vec![pps, drops])
|
||||
.block(
|
||||
Block::new()
|
||||
.borders(Borders::ALL)
|
||||
.title(" pps / drops (last 120 samples) "),
|
||||
)
|
||||
.x_axis(Axis::default().bounds([0.0, x_max]))
|
||||
.y_axis(Axis::default().bounds([0.0, y_max]));
|
||||
frame.render_widget(chart, area);
|
||||
}
|
||||
|
||||
fn render_help(frame: &mut Frame, area: Rect) {
|
||||
frame.render_widget(
|
||||
Paragraph::new(Line::from(vec![
|
||||
Span::styled("q", Style::default().add_modifier(Modifier::BOLD)),
|
||||
Span::raw(" quit "),
|
||||
Span::styled("r", Style::default().add_modifier(Modifier::BOLD)),
|
||||
Span::raw(" refresh now"),
|
||||
])),
|
||||
area,
|
||||
);
|
||||
}
|
||||
159
tests/tui_metrics.rs
Normal file
159
tests/tui_metrics.rs
Normal file
|
|
@ -0,0 +1,159 @@
|
|||
//! Тесты TUI-модуля: парсер Prometheus, расчёт rate, поведение при ошибках fetch.
|
||||
|
||||
use futures::future::BoxFuture;
|
||||
use rampart::tui::fetch::{FetchError, Fetcher, HttpFetcher};
|
||||
use rampart::tui::prometheus::{MetricSample, parse};
|
||||
use rampart::tui::state::{AppState, counter_rate};
|
||||
|
||||
// ---------- парсер ----------
|
||||
|
||||
#[test]
|
||||
fn parses_metric_with_labels_value_and_timestamp() {
|
||||
let set = parse("rampart_xdp_total{iface=\"eth0\",core=\"0\"} 123 1700000000000\n");
|
||||
assert_eq!(
|
||||
set.samples(),
|
||||
vec![MetricSample {
|
||||
name: "rampart_xdp_total".to_owned(),
|
||||
labels: [("iface", "eth0"), ("core", "0")]
|
||||
.iter()
|
||||
.map(|(k, v)| ((*k).to_owned(), (*v).to_owned()))
|
||||
.collect(),
|
||||
value: 123.0,
|
||||
}]
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn parser_skips_comments_help_and_type() {
|
||||
let body = "\
|
||||
# HELP rampart_xdp_dropped Packets dropped\n\
|
||||
# TYPE rampart_xdp_dropped gauge\n\
|
||||
# обычный комментарий\n\
|
||||
rampart_xdp_dropped 42\n";
|
||||
let set = parse(body);
|
||||
assert_eq!(set.samples().len(), 1);
|
||||
assert_eq!(set.sum("rampart_xdp_dropped"), 42.0);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn parser_sums_counter_vec_series() {
|
||||
let body = "rampart_auto_bans_total{reason=\"cps\"} 3\nrampart_auto_bans_total{reason=\"syn\"} 4\n";
|
||||
assert_eq!(parse(body).sum("rampart_auto_bans_total"), 7.0);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn parser_missing_metric_is_zero() {
|
||||
assert_eq!(parse("# nothing here\n").sum("rampart_xdp_tcp"), 0.0);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn parser_tolerates_malformed_lines() {
|
||||
let set = parse("garbage\nm{unterminated=\"v} 1\nok_metric 5\n");
|
||||
assert_eq!(set.sum("ok_metric"), 5.0);
|
||||
assert_eq!(set.samples().len(), 1);
|
||||
}
|
||||
|
||||
// ---------- rate ----------
|
||||
|
||||
#[test]
|
||||
fn rate_from_two_samples_divides_by_elapsed() {
|
||||
// 1200 - 1100 = 100 пакетов за 2 сек
|
||||
assert_eq!(counter_rate(Some(1100.0), Some(1200.0), 2.0), 50.0);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn rate_first_sample_and_counter_reset_are_zero() {
|
||||
assert_eq!(counter_rate(None, Some(10.0), 1.0), 0.0);
|
||||
assert_eq!(counter_rate(Some(20.0), Some(10.0), 1.0), 0.0);
|
||||
assert_eq!(counter_rate(Some(10.0), Some(12.0), 0.0), 0.0);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn state_accumulates_history_rates() {
|
||||
let state = AppState::new("http://test/metrics");
|
||||
state.apply_fetch(Ok(parse("rampart_xdp_total 100\n")), 1.0);
|
||||
state.apply_fetch(Ok(parse("rampart_xdp_total 150\n")), 1.0);
|
||||
let frame = state.frame_data();
|
||||
assert_eq!(frame.history.points().len(), 2);
|
||||
assert_eq!(frame.history.points()[1].packets_per_sec, 50.0);
|
||||
assert!(frame.snapshot.error.is_none());
|
||||
}
|
||||
|
||||
// ---------- ошибки fetch (мок Fetcher + локальный HTTP-сервер) ----------
|
||||
|
||||
struct FailingFetcher {
|
||||
error: FetchError,
|
||||
}
|
||||
|
||||
impl Fetcher for FailingFetcher {
|
||||
fn fetch(&self) -> BoxFuture<'_, Result<String, FetchError>> {
|
||||
let error = self.error.clone();
|
||||
Box::pin(async { Err(error) })
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn network_error_sets_error_banner_not_panic() {
|
||||
let state = AppState::new("http://127.0.0.1:9/metrics");
|
||||
let fetcher = FailingFetcher {
|
||||
error: FetchError::Network("connection refused".to_owned()),
|
||||
};
|
||||
|
||||
// Тот же путь, что и в spawn_fetch_loop: normalize → apply_fetch.
|
||||
let normalized = rampart::tui::app::normalize(Err(fetcher.error));
|
||||
assert!(normalized.is_err());
|
||||
let message = match normalized {
|
||||
Err(message) => message,
|
||||
Ok(_) => return,
|
||||
};
|
||||
|
||||
state.apply_fetch(Err(message), 1.0);
|
||||
assert_eq!(state.frame_data().snapshot.error.as_deref(), Some("connection refused"));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn http_fetcher_reads_prometheus_body_from_local_server() {
|
||||
let Ok(listener) = tokio::net::TcpListener::bind("127.0.0.1:0").await else {
|
||||
return;
|
||||
};
|
||||
let Ok(addr) = listener.local_addr() else {
|
||||
return;
|
||||
};
|
||||
let server = tokio::spawn(async move {
|
||||
use tokio::io::AsyncWriteExt;
|
||||
if let Ok((mut stream, _)) = listener.accept().await {
|
||||
let body = "rampart_xdp_total 99\n";
|
||||
let head = format!(
|
||||
"HTTP/1.1 200 OK\r\nContent-Length: {}\r\nConnection: close\r\n\r\n",
|
||||
body.len()
|
||||
);
|
||||
let _ = stream.write_all(head.as_bytes()).await;
|
||||
let _ = stream.write_all(body.as_bytes()).await;
|
||||
}
|
||||
});
|
||||
|
||||
let Some(fetcher) = HttpFetcher::new(format!("http://{addr}/metrics")) else {
|
||||
server.abort();
|
||||
return;
|
||||
};
|
||||
let Ok(body) = fetcher.fetch().await else {
|
||||
server.abort();
|
||||
return;
|
||||
};
|
||||
assert_eq!(parse(&body).sum("rampart_xdp_total"), 99.0);
|
||||
server.abort();
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn failed_fetch_keeps_last_good_metrics_but_reports_error() {
|
||||
let state = AppState::new("http://test/metrics");
|
||||
state.apply_fetch(Ok(parse("rampart_xdp_total 7\n")), 1.0);
|
||||
state.apply_fetch(Err("connection refused".to_owned()), 1.0);
|
||||
|
||||
let frame = state.frame_data();
|
||||
assert_eq!(frame.snapshot.error.as_deref(), Some("connection refused"));
|
||||
assert_eq!(
|
||||
frame.snapshot.metrics.as_ref().map(|m| m.sum("rampart_xdp_total")),
|
||||
Some(7.0)
|
||||
);
|
||||
}
|
||||
Loading…
Add table
Add a link
Reference in a new issue