use std::collections::BTreeMap;
#[derive(Debug, Clone, PartialEq)]
pub struct Sample {
pub name: String,
pub labels: BTreeMap<String, String>,
pub value: f64,
}
pub fn parse_samples(text: &str) -> Vec<Sample> {
text.lines().filter_map(parse_line).collect()
}
fn parse_line(line: &str) -> Option<Sample> {
let line = line.trim();
if line.is_empty() || line.starts_with('#') {
return None;
}
let (head, value_str) = match line.find('}') {
Some(close) => {
let (h, rest) = line.split_at(close + 1);
(h, rest.trim())
}
None => {
let mut parts = line.split_whitespace();
let h = parts.next()?;
(h, line[h.len()..].trim())
}
};
let value: f64 = value_str.split_whitespace().next()?.parse().ok()?;
let (name, labels) = match head.find('{') {
None => (head.trim().to_string(), BTreeMap::new()),
Some(open) => {
let name = head[..open].trim().to_string();
let body = head[open + 1..head.len() - 1].trim_end_matches(',');
(name, parse_labels(body)?)
}
};
if name.is_empty() {
return None;
}
Some(Sample {
name,
labels,
value,
})
}
fn parse_labels(body: &str) -> Option<BTreeMap<String, String>> {
let mut labels = BTreeMap::new();
let mut chars = body.chars().peekable();
loop {
while matches!(chars.peek(), Some(',') | Some(' ')) {
chars.next();
}
if chars.peek().is_none() {
return Some(labels);
}
let mut key = String::new();
for c in chars.by_ref() {
if c == '=' {
break;
}
key.push(c);
}
if chars.next()? != '"' {
return None;
}
let mut value = String::new();
loop {
match chars.next()? {
'\\' => match chars.next()? {
'n' => value.push('\n'),
'\\' => value.push('\\'),
'"' => value.push('"'),
other => {
value.push('\\');
value.push(other);
}
},
'"' => break,
c => value.push(c),
}
}
labels.insert(key.trim().to_string(), value);
}
}
#[derive(Debug, Clone, Default, PartialEq)]
pub struct RowStats {
pub source: String,
pub sink: String,
pub records_in: u64,
pub records_out: u64,
pub rate: f64,
pub source_errors: u64,
pub sink_errors: u64,
pub dlq_records: u64,
pub last_bookmark_unix: f64,
pub finished: Option<bool>,
pub in_flight: bool,
}
#[derive(Debug, Clone, Default, PartialEq)]
pub struct TuiModel {
pub rows: BTreeMap<String, RowStats>,
pub total_in: u64,
pub total_out: u64,
pub total_rate: f64,
}
#[derive(Debug, Default)]
pub struct Sampler {
pipeline: String,
prev_out: BTreeMap<String, (u64, std::time::Duration)>,
}
impl Sampler {
pub fn new(pipeline: impl Into<String>) -> Self {
Self {
pipeline: pipeline.into(),
prev_out: BTreeMap::new(),
}
}
pub fn observe(&mut self, text: &str, elapsed: std::time::Duration) -> TuiModel {
let mut model = TuiModel::default();
for s in parse_samples(text) {
if s.labels.get("pipeline").map(String::as_str) != Some(self.pipeline.as_str()) {
continue;
}
let row_id = s.labels.get("row").cloned().unwrap_or_default();
let row = model.rows.entry(row_id).or_default();
let connector = s.labels.get("connector").cloned().unwrap_or_default();
match s.name.as_str() {
"faucet_source_records_total" => {
row.records_in += s.value as u64;
if !connector.is_empty() {
row.source = connector;
}
}
"faucet_sink_records_total" => {
row.records_out += s.value as u64;
if !connector.is_empty() {
row.sink = connector;
}
}
"faucet_source_errors_total" => row.source_errors += s.value as u64,
"faucet_sink_errors_total" => row.sink_errors += s.value as u64,
"faucet_sink_dlq_records_total" => row.dlq_records += s.value as u64,
"faucet_pipeline_last_bookmark_unix_seconds" => {
row.last_bookmark_unix = s.value;
}
"faucet_pipeline_in_flight" => {
if s.value > 0.0 {
row.in_flight = true;
}
}
"faucet_pipeline_runs_total" => {
if row.source.is_empty()
&& let Some(src) = s.labels.get("source")
{
row.source = src.clone();
}
if row.sink.is_empty()
&& let Some(dst) = s.labels.get("sink")
{
row.sink = dst.clone();
}
if s.value > 0.0 {
match s.labels.get("status").map(String::as_str) {
Some("ok") => row.finished = Some(row.finished.unwrap_or(true)),
Some("err") => row.finished = Some(false),
_ => {}
}
}
}
_ => {}
}
}
let mut prev = std::mem::take(&mut self.prev_out);
for (row_id, row) in &mut model.rows {
if let Some((prev_count, prev_at)) = prev.remove(row_id) {
let dt = elapsed.saturating_sub(prev_at).as_secs_f64();
if dt > 0.0 && row.records_out >= prev_count {
row.rate = (row.records_out - prev_count) as f64 / dt;
}
}
self.prev_out
.insert(row_id.clone(), (row.records_out, elapsed));
model.total_in += row.records_in;
model.total_out += row.records_out;
model.total_rate += row.rate;
}
model
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::time::Duration;
#[test]
fn parses_bare_and_labelled_samples() {
let text = "\
# HELP faucet_x helper
# TYPE faucet_x counter
faucet_up 1
faucet_source_records_total{pipeline=\"p\",row=\"r1\",connector=\"csv\"} 42
";
let samples = parse_samples(text);
assert_eq!(samples.len(), 2);
assert_eq!(samples[0].name, "faucet_up");
assert_eq!(samples[0].value, 1.0);
assert!(samples[0].labels.is_empty());
assert_eq!(samples[1].labels["connector"], "csv");
assert_eq!(samples[1].value, 42.0);
}
#[test]
fn parses_escaped_label_values() {
let text = r#"m{a="quo\"te",b="back\\slash",c="new\nline"} 7"#;
let s = &parse_samples(text)[0];
assert_eq!(s.labels["a"], "quo\"te");
assert_eq!(s.labels["b"], "back\\slash");
assert_eq!(s.labels["c"], "new\nline");
}
#[test]
fn garbage_lines_are_skipped_not_fatal() {
let text = "not a metric\nname_only\n{}} 3\nok_metric 5\n";
let samples = parse_samples(text);
assert_eq!(samples.len(), 1);
assert_eq!(samples[0].name, "ok_metric");
}
fn render(pipeline: &str, row: &str, out: u64) -> String {
format!(
"faucet_sink_records_total{{pipeline=\"{pipeline}\",row=\"{row}\",connector=\"jsonl\"}} {out}\n"
)
}
#[test]
fn sampler_filters_by_pipeline_and_computes_rates() {
let mut s = Sampler::new("mine");
let t0 = Duration::from_secs(1);
let m = s.observe(&(render("mine", "a", 100) + &render("other", "a", 999)), t0);
assert_eq!(m.rows.len(), 1);
assert_eq!(m.rows["a"].records_out, 100);
assert_eq!(m.rows["a"].rate, 0.0, "no window yet");
let m = s.observe(&render("mine", "a", 300), Duration::from_secs(3));
assert!((m.rows["a"].rate - 100.0).abs() < 1e-9, "200 records / 2s");
assert_eq!(m.total_out, 300);
}
#[test]
fn sampler_counter_reset_yields_zero_rate() {
let mut s = Sampler::new("p");
s.observe(&render("p", "a", 500), Duration::from_secs(1));
let m = s.observe(&render("p", "a", 10), Duration::from_secs(2));
assert_eq!(m.rows["a"].rate, 0.0);
}
#[test]
fn sampler_aggregates_row_fields() {
let text = "\
faucet_source_records_total{pipeline=\"p\",row=\"r\",connector=\"spanner\"} 10
faucet_sink_records_total{pipeline=\"p\",row=\"r\",connector=\"jsonl\"} 8
faucet_source_errors_total{pipeline=\"p\",row=\"r\",connector=\"spanner\",kind=\"http\"} 1
faucet_sink_errors_total{pipeline=\"p\",row=\"r\",connector=\"jsonl\",kind=\"io\"} 2
faucet_sink_dlq_records_total{pipeline=\"p\",row=\"r\",connector=\"jsonl\"} 3
faucet_pipeline_last_bookmark_unix_seconds{pipeline=\"p\",row=\"r\"} 1700000000
faucet_pipeline_in_flight{pipeline=\"p\",row=\"r\"} 1
";
let mut s = Sampler::new("p");
let m = s.observe(text, Duration::from_secs(1));
let row = &m.rows["r"];
assert_eq!(row.source, "spanner");
assert_eq!(row.sink, "jsonl");
assert_eq!(row.records_in, 10);
assert_eq!(row.records_out, 8);
assert_eq!(row.source_errors, 1);
assert_eq!(row.sink_errors, 2);
assert_eq!(row.dlq_records, 3);
assert_eq!(row.last_bookmark_unix, 1_700_000_000.0);
assert!(row.in_flight);
assert_eq!(row.finished, None);
}
#[test]
fn run_status_ok_and_err_map_to_finished() {
let ok = "faucet_pipeline_runs_total{pipeline=\"p\",row=\"r\",source=\"csv\",sink=\"stdout\",status=\"ok\"} 1\n";
let err = "faucet_pipeline_runs_total{pipeline=\"p\",row=\"r\",source=\"csv\",sink=\"stdout\",status=\"err\",kind=\"sink\"} 1\n";
let mut s = Sampler::new("p");
let m = s.observe(ok, Duration::from_secs(1));
assert_eq!(m.rows["r"].finished, Some(true));
assert_eq!(m.rows["r"].source, "csv");
assert_eq!(m.rows["r"].sink, "stdout");
let m = s.observe(&(ok.to_string() + err), Duration::from_secs(2));
assert_eq!(m.rows["r"].finished, Some(false));
}
}