From a57635cd6178261ee7915380c4bb377f5ee39158 Mon Sep 17 00:00:00 2001 From: stqfdyr <89149493+stqfdyr@users.noreply.github.com> Date: Sat, 26 Sep 2026 09:26:42 +0800 Subject: [PATCH 1/2] =?UTF-8?q?feat:=20=E5=8E=86=E5=8F=B2=E8=AE=B0?= =?UTF-8?q?=E5=BD=95=E4=BF=9D=E7=95=99=E6=AF=8F=E5=88=86=E9=92=9F=E7=9A=84?= =?UTF-8?q?=E5=B3=B0=E5=80=BC=E7=BD=91=E9=80=9F?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 历史行的网速是这一分钟的平均值,积分等于累计流量,但一次 15 秒的测速被摊到整分钟。 同一行另存这一分钟 agent 上报过的最高网速,接口每行多 net_rx_max / net_tx_max, 聚合成宽窗口时取桶内最大值。 --- src/agent_ws.rs | 30 +++++++++++++++++++++++++++--- src/api.rs | 13 +++++++++++-- src/db.rs | 33 ++++++++++++++++++++++++++++----- 3 files changed, 66 insertions(+), 10 deletions(-) diff --git a/src/agent_ws.rs b/src/agent_ws.rs index 8b6032b3..3003fbc9 100644 --- a/src/agent_ws.rs +++ b/src/agent_ws.rs @@ -150,11 +150,24 @@ 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 { @@ -162,6 +175,9 @@ impl Minute { 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; } @@ -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)); + } } } @@ -1079,15 +1100,18 @@ mod tests { }; // Busy for half the minute, then idle; 60 MB arrive in between, and by - // the next sample both have ended. + // the next sample both have ended. The sample halfway caught the busiest + // second of the burst. send(&app, id, &mut session, 0, &burst(1_000, 0, 100.0, 100)).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_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. diff --git a/src/api.rs b/src/api.rs index eb2c3845..1d86cf52 100644 --- a/src/api.rs +++ b/src/api.rs @@ -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]); } @@ -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 diff --git a/src/db.rs b/src/db.rs index 27085b83..283ac2d7 100644 --- a/src/db.rs +++ b/src/db.rs @@ -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; @@ -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. @@ -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. @@ -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(()) @@ -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, @@ -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(()) } @@ -1212,12 +1226,20 @@ 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 divides by + /// the hub's clock in whole seconds, which over a minute of steady traffic + /// can put it 1-2% above the agent's per-second rates. pub fn metrics(&self, node_id: i64, since: i64, step: i64) -> Result> { 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| { @@ -1225,6 +1247,7 @@ impl Db { "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::>()?) From 181bbbce8c76e67bc281ef55c5fc026891903aaf Mon Sep 17 00:00:00 2001 From: stqfdyr <89149493+stqfdyr@users.noreply.github.com> Date: Sun, 27 Sep 2026 22:23:51 +0800 Subject: [PATCH 2/2] =?UTF-8?q?fix:=20=E5=88=86=E9=92=9F=E5=B9=B3=E5=9D=87?= =?UTF-8?q?=E7=BD=91=E9=80=9F=E6=8C=89=E5=B0=8F=E6=95=B0=E7=A7=92=E8=AE=A1?= =?UTF-8?q?=E7=AE=97?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- src/agent_ws.rs | 20 +++++++++++++------- src/db.rs | 6 +++--- 2 files changed, 16 insertions(+), 10 deletions(-) diff --git a/src/agent_ws.rs b/src/agent_ws.rs index 3003fbc9..f1f78cd3 100644 --- a/src/agent_ws.rs +++ b/src/agent_ws.rs @@ -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)] @@ -778,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); @@ -1101,13 +1103,17 @@ mod tests { // Busy for half the minute, then idle; 60 MB arrive in between, and by // the next sample both have ended. The sample halfway caught the busiest - // second of the burst. - send(&app, id, &mut session, 0, &burst(1_000, 0, 100.0, 100)).unwrap(); + // 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 diff --git a/src/db.rs b/src/db.rs index 283ac2d7..3a04541d 100644 --- a/src/db.rs +++ b/src/db.rs @@ -1230,9 +1230,9 @@ impl Db { /// `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 divides by - /// the hub's clock in whole seconds, which over a minute of steady traffic - /// can put it 1-2% above the agent's per-second rates. + /// 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> { let conn = self.conn(); let mut stmt = conn.prepare_cached(