Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
48 changes: 39 additions & 9 deletions src/agent_ws.rs
Original file line number Diff line number Diff line change
Expand Up @@ -130,7 +130,7 @@ impl Agent {
/// `tick` is an `Instant` rather than the wall clock, because the rate divides by
/// a duration. NTP stepping the clock backwards -- a fresh boot correcting
/// itself, a restored snapshot -- makes a wall-clock difference negative, and the
/// `.max(1)` guarding the division would then divide a whole minute of bytes by
/// `.max(1.0)` guarding the division would then divide a whole minute of bytes by
/// one second. The agent computes its own rate against `Instant` for the same
/// reason.
#[derive(Debug)]
Expand All @@ -150,18 +150,34 @@ struct Mark {
const MEAN_FLOAT: [&str; 1] = ["cpu"];
const MEAN_INT: [&str; 6] = ["mem_used", "swap_used", "disk_used", "tcp", "udp", "procs"];

/// Running sums for the minute in progress, one slot per averaged field.
/// Rates a history row also carries at their highest over the minute, each as
/// the agent measured it across one report interval, under the column it is
/// stored in. The row's own rate is the minute's mean, which integrates to the
/// traffic totals and therefore stores a 15-second burst at 286 Mbps as 72 Mbps
/// (measured).
///
/// Taken from the agent rather than derived here from the arrival of two
/// frames: the network bunches frames, and a second of bytes divided by the
/// half second between two arrivals would record twice the rate that ran.
const PEAK: [(&str, &str); 2] = [("net_rx", "net_rx_max"), ("net_tx", "net_tx_max")];

/// Running sums for the minute in progress, one slot per averaged field, and
/// the highest of each [`PEAK`] rate.
#[derive(Debug, Default)]
struct Minute {
sums: [f64; MEAN_FLOAT.len() + MEAN_INT.len()],
reports: f64,
peaks: [i64; PEAK.len()],
}

impl Minute {
fn add(&mut self, metrics: &serde_json::Value) {
for (slot, key) in MEAN_FLOAT.iter().chain(&MEAN_INT).enumerate() {
self.sums[slot] += metrics.get(key).and_then(|v| v.as_f64()).unwrap_or(0.0);
}
for (peak, (key, _)) in self.peaks.iter_mut().zip(PEAK) {
*peak = (*peak).max(metrics.get(key).and_then(|v| v.as_i64()).unwrap_or(0));
}
self.reports += 1.0;
}

Expand All @@ -181,6 +197,11 @@ impl Minute {
let mean = if slot < MEAN_FLOAT.len() { json!(mean) } else { json!(mean.round() as i64) };
obj.insert((*key).to_owned(), mean);
}
// Written whatever the report carried, so an agent sending these keys
// itself cannot choose the stored value.
for (peak, (_, column)) in self.peaks.iter().zip(PEAK) {
obj.insert(column.to_owned(), json!(peak));
}
}
}

Expand Down Expand Up @@ -757,9 +778,11 @@ fn report(app: &App, node_id: i64, metrics: serde_json::Value, arrival: Arrival)
entry.minute.write_into(&mut row);
if let (Some((epoch, (rx, tx))), Some(mark), Some(obj)) = (&span, &entry.mark, row.as_object_mut()) {
if mark.epoch == *epoch {
let elapsed = arrival.tick.saturating_duration_since(mark.tick).as_secs().max(1) as i64;
obj.insert("net_rx".into(), json!((rx - mark.counters.0).max(0) / elapsed));
obj.insert("net_tx".into(), json!((tx - mark.counters.1).max(0) / elapsed));
// Fractional seconds: whole ones would drop up to 0.99 s of the
// minute and overstate its rate by up to 1.7%.
let elapsed = arrival.tick.saturating_duration_since(mark.tick).as_secs_f64().max(1.0);
obj.insert("net_rx".into(), json!(((rx - mark.counters.0).max(0) as f64 / elapsed) as i64));
obj.insert("net_tx".into(), json!(((tx - mark.counters.1).max(0) as f64 / elapsed) as i64));
}
}
entry.last_minute = Some(minute);
Expand Down Expand Up @@ -1079,15 +1102,22 @@ mod tests {
};

// Busy for half the minute, then idle; 60 MB arrive in between, and by
// the next sample both have ended.
send(&app, id, &mut session, 0, &burst(1_000, 0, 100.0, 100)).unwrap();
// the next sample both have ended. The sample halfway caught the busiest
// second of the burst. The first lands half a second in, so the row spans
// 59.5 s.
dispatch(&app, id, "ip", &burst(1_000, 0, 100.0, 100), &mut session, at_ms(500)).unwrap();
send(&app, id, &mut session, 30, &burst(1_000 + 45_000_000, 3_000_000, 50.0, 151)).unwrap();
send(&app, id, &mut session, 60, &burst(1_000 + 60_000_000, 0, 0.0, 201)).unwrap();

let row = &app.db.metrics(id, 0, 60).unwrap()[0];
assert_eq!(row["net_rx"], 1_000_000, "60 MB over 60 s is 1 MB/s, not the agent's 0");
assert_eq!(
row["net_rx"], 1_008_403,
"60 MB over 59.5 s, not the agent's 0 nor over 59 whole seconds"
);
assert_eq!(row["net_rx_max"], 3_000_000, "the busiest second survives the mean");
assert_eq!(row["cpu"], 50.0, "the mean of the minute, not the idle second it ended on");
// Integers remain integral: the column is read with as_i64, which returns
// nothing for the 150.5 the raw mean would produce.
// nothing for the 150.67 the raw mean would produce.
assert_eq!(row["mem_used"], 151);
// The live view still shows the instantaneous reading, which is its
// purpose.
Expand Down
13 changes: 11 additions & 2 deletions src/api.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2377,8 +2377,15 @@ mod tests {
let base = Utc::now().timestamp() / 120 * 120 - 120;
// One bucket: a quiet minute and a busy one, then a probe that answered
// once and timed out three times.
app.db.insert_metric(id, base + 10, &json!({"cpu": 0.0, "net_rx": 0})).unwrap();
app.db.insert_metric(id, base + 70, &json!({"cpu": 40.0, "net_rx": 1_000})).unwrap();
// The quiet minute predates the peak column, whose default is 0.
app.db.insert_metric(id, base + 10, &json!({"cpu": 0.0, "net_rx": 0, "net_tx": 3_000})).unwrap();
app.db
.insert_metric(
id,
base + 70,
&json!({"cpu": 40.0, "net_rx": 1_000, "net_rx_max": 4_000, "net_tx": 1_000, "net_tx_max": 2_000}),
)
.unwrap();
for _ in 0..3 {
task(&app, vec![id]);
}
Expand All @@ -2392,6 +2399,8 @@ mod tests {
let m = &app.db.metrics(id, base, 120).unwrap()[0];
assert_eq!(m["cpu"], 20.0, "the bucket is its mean, not one row of it");
assert_eq!(m["net_rx"], 500);
assert_eq!(m["net_rx_max"], 4_000, "the bucket peaks where its busiest minute did");
assert_eq!(m["net_tx_max"], 3_000, "a row without a peak counts as its own mean");
assert_eq!(m["ts"], base, "stamped with the bucket, so every series shares a grid");

// Keyed by task rather than index: the order is the panel's, which
Expand Down
33 changes: 28 additions & 5 deletions src/db.rs
Original file line number Diff line number Diff line change
Expand Up @@ -115,6 +115,7 @@ CREATE TABLE IF NOT EXISTS metric (
mem_used INTEGER NOT NULL, swap_used INTEGER NOT NULL, disk_used INTEGER NOT NULL,
net_rx INTEGER NOT NULL, net_tx INTEGER NOT NULL,
tcp INTEGER NOT NULL, udp INTEGER NOT NULL, procs INTEGER NOT NULL,
net_rx_max INTEGER NOT NULL DEFAULT 0, net_tx_max INTEGER NOT NULL DEFAULT 0,
PRIMARY KEY (node_id, ts)
) WITHOUT ROWID;

Expand Down Expand Up @@ -162,7 +163,7 @@ CREATE TABLE IF NOT EXISTS session (
/// cannot: `open` runs `SCHEMA` before migrating, and on an older file the
/// column is not there yet. `an_upgraded_release_matches_a_fresh_database`
/// holds every migration to these rules, starting from v1.0.0's schema.
const SCHEMA_VERSION: i64 = 9;
const SCHEMA_VERSION: i64 = 10;

/// Adds a column older databases lack. A duplicate column indicates the
/// migration has already run; every other error must propagate.
Expand Down Expand Up @@ -294,6 +295,13 @@ fn migrate_to_9(conn: &Connection) -> Result<()> {
add_column(conn, "ping_task", "sort INTEGER NOT NULL DEFAULT 0")
}

/// Rows written before the peak existed hold 0, which `Db::metrics` reads as
/// the row's mean rather than rewriting every row of history here.
fn migrate_to_10(conn: &Connection) -> Result<()> {
add_column(conn, "metric", "net_rx_max INTEGER NOT NULL DEFAULT 0")?;
add_column(conn, "metric", "net_tx_max INTEGER NOT NULL DEFAULT 0")
}

/// Brings a database already in service up to `SCHEMA_VERSION` and stamps it.
/// `from` is its current version, so a fresh file passes `SCHEMA_VERSION` and
/// receives only the stamp.
Expand Down Expand Up @@ -334,6 +342,9 @@ fn migrate(conn: &Connection, from: i64) -> Result<()> {
if from < 9 {
migrate_to_9(&tx)?;
}
if from < 10 {
migrate_to_10(&tx)?;
}
tx.execute_batch(&format!("PRAGMA user_version = {SCHEMA_VERSION}"))?;
tx.commit()?;
Ok(())
Expand Down Expand Up @@ -1176,8 +1187,9 @@ impl Db {
self.conn()
.prepare_cached(
"INSERT OR REPLACE INTO metric
(node_id, ts, cpu, mem_used, swap_used, disk_used, net_rx, net_tx, tcp, udp, procs)
VALUES (?1,?2,?3,?4,?5,?6,?7,?8,?9,?10,?11)",
(node_id, ts, cpu, mem_used, swap_used, disk_used, net_rx, net_tx, tcp, udp, procs,
net_rx_max, net_tx_max)
VALUES (?1,?2,?3,?4,?5,?6,?7,?8,?9,?10,?11,?12,?13)",
)?
.execute(params![
node_id,
Expand All @@ -1190,7 +1202,9 @@ impl Db {
n("net_tx"),
n("tcp"),
n("udp"),
n("procs")
n("procs"),
n("net_rx_max"),
n("net_tx_max")
])?;
Ok(())
}
Expand All @@ -1212,19 +1226,28 @@ impl Db {
///
/// The stamp is the bucket's start rather than a row inside it, so every
/// series lands on one grid and the probe rows below can be shared.
///
/// `net_rx_max` and `net_tx_max` are the bucket's highest rather than its
/// mean, since a maximum of maxima loses nothing: a week's window peaks at
/// the same rate as the minute that reached it. Each row counts as at least
/// its own mean: rows predating the column hold 0, and the mean, timed by
/// the hub's arrivals rather than the agent's clock, can edge past the
/// agent's own rates by the network's jitter.
pub fn metrics(&self, node_id: i64, since: i64, step: i64) -> Result<Vec<serde_json::Value>> {
let conn = self.conn();
let mut stmt = conn.prepare_cached(
"SELECT (MIN(ts)/?3)*?3, AVG(cpu), CAST(AVG(mem_used) AS INTEGER),
CAST(AVG(disk_used) AS INTEGER),
CAST(AVG(net_rx) AS INTEGER), CAST(AVG(net_tx) AS INTEGER)
CAST(AVG(net_rx) AS INTEGER), CAST(AVG(net_tx) AS INTEGER),
MAX(MAX(net_rx, net_rx_max)), MAX(MAX(net_tx, net_tx_max))
FROM metric WHERE node_id=?1 AND ts>=?2 GROUP BY ts/?3 ORDER BY ts/?3",
)?;
let rows = stmt.query_map(params![node_id, since, step], |r| {
Ok(serde_json::json!({
"ts": r.get::<_, i64>(0)?, "cpu": r.get::<_, f64>(1)?,
"mem_used": r.get::<_, i64>(2)?, "disk_used": r.get::<_, i64>(3)?,
"net_rx": r.get::<_, i64>(4)?, "net_tx": r.get::<_, i64>(5)?,
"net_rx_max": r.get::<_, i64>(6)?, "net_tx_max": r.get::<_, i64>(7)?,
}))
})?;
Ok(rows.collect::<Result<_, _>>()?)
Expand Down
Loading