From 38009f1f2eec5f01ccdfd47b5fb1edbeb9bfdb4d Mon Sep 17 00:00:00 2001 From: stqfdyr <89149493+stqfdyr@users.noreply.github.com> Date: Thu, 1 Oct 2026 12:47:44 +0800 Subject: [PATCH 1/2] =?UTF-8?q?perf:=20=E5=8E=86=E5=8F=B2=E5=9B=BE?= =?UTF-8?q?=E8=A1=A8=E7=9A=84=E6=9F=A5=E8=AF=A2=E8=B5=B0=E7=8B=AC=E7=AB=8B?= =?UTF-8?q?=E7=9A=84=E5=8F=AA=E8=AF=BB=E8=BF=9E=E6=8E=A5=EF=BC=8C=E4=B8=8D?= =?UTF-8?q?=E5=86=8D=E6=8C=A1=E4=BD=8F=20agent=20=E4=B8=8A=E6=8A=A5?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 打开数据库时另开一条只读连接,`/api/nodes/{id}/metrics` 的资源、延迟、探测名称三个查询都走它。WAL 下它读的同时写连接照常提交,agent 上报和面板请求不再排在历史扫描后面 - 每次查询包在一个读事务里,水位线和它两侧的行来自同一个快照 - 只读连接先于写连接关闭,hub 正常停机后库仍是单个文件,不留 -wal - `:memory:` 开不了第二条连接,测试库仍走写连接 - 历史闸门仍是 4 个,扫描改在只读连接上排队,闸门决定的是图表最长等多久 --- src/api.rs | 41 +++++++++--------- src/db.rs | 123 +++++++++++++++++++++++++++++++++++++++++++++++------ 2 files changed, 130 insertions(+), 34 deletions(-) diff --git a/src/api.rs b/src/api.rs index d88606b..48ed81b 100644 --- a/src/api.rs +++ b/src/api.rs @@ -341,19 +341,19 @@ fn default_hours() -> i64 { /// `span` bounds what one request costs; this bounds how many may run, /// closing the same gap `main::RELAY_GATE` and `auth::PASSWORD_GATE` close on /// the other two paths an anonymous caller can make expensive. This is the most -/// expensive of the three: every request holds the single connection the agents -/// report through for its entire scan, measured at 118 ms for a week of four -/// probes, the most minute rows any window reads. Without a gate, 120 requests -/// from one machine took the panel's own node list from 1 ms to 2.8 s. +/// expensive of the three: a week of probe results is a scan of 98 ms at four +/// 60-second probes and 486 ms at eight 10-second ones, growing with the probes +/// the admin configured rather than with anything the caller sends. /// -/// Four, because the requests serialise on that one connection regardless: a -/// fifth in flight buys no throughput and merely places another scan ahead of -/// the next agent report. What the number actually sets is how long that wait -/// can become -- four at roughly 120 ms is half a second -- while leaving room -/// for several people opening charts simultaneously. +/// The scans run on the database's read-only connection, so the agents' writes +/// do not wait on them, and they serialise on that one connection instead. Four, +/// because a fifth in flight buys no throughput and only lengthens the wait +/// behind the others: what the number sets is how long a chart can wait, four +/// scans of the slower kind being 2 s, while leaving room for several people +/// opening charts simultaneously. /// /// Refused rather than queued, as in `auth`: a queue admits the same flood, -/// merely later. +/// merely later, and each request waiting in it holds a blocking thread. /// /// **This gate is ineffective without the `spawn_blocking` below.** The body of /// this handler never awaits, so a permit taken and dropped within it is held @@ -387,10 +387,9 @@ pub async fn metrics( let span = span(hours, w.points, Utc::now().timestamp()); let wants = |name: &str| w.series.as_deref().is_none_or(|s| s == name); let (want_metrics, want_ping) = (wants("metrics"), wants("ping")); - // Off the runtime, for the reason given in `db_stats` below: this reads every - // probe result the node has retained within the window and holds the - // connection the agents report through throughout. That route is behind - // `Admin` and cheaper than this one, which anyone can reach. + // Off the runtime: this reads every probe result the node has retained + // within the window, and a worker thread blocked on that scan, or on the + // reader while another request holds it, serves nothing else. // // It is also what makes the gate above effective: the permit is held across // an await, so exactly four callers are inside it at once rather than however @@ -399,7 +398,7 @@ pub async fn metrics( // Probe names accompany the samples they label, so the page needs no // second request. Names only: targets and assignments remain behind // `Admin`. Skipped when probes were not requested, since the resources tab - // has nothing to label and this costs a turn at the write connection. + // has nothing to label and this costs a turn at the reader. let probes = if want_ping { app.db.ping_task_names(id).unwrap_or_else(|_| json!({})) } else { json!({}) }; let metrics = if want_metrics { app.db.metrics(id, span)? } else { vec![] }; @@ -2406,8 +2405,7 @@ mod tests { /// A chart request costs roughly the same whatever it spans. This path /// requires no session, so an unbounded window would be megabytes of JSON any - /// caller could have the hub build on the connection the agents report - /// through. + /// caller could have the hub build. #[test] fn a_history_window_costs_the_same_however_wide_it_is() { let app = app(); @@ -3298,11 +3296,10 @@ mod tests { assert_eq!(app.db.node(id).unwrap().unwrap().disk_total, 30i64 << 30, "and no extra write to get it"); } - /// `span` bounds one window; this bounds how many are built - /// concurrently. Each holds the connection the agents report through for its - /// entire scan, and the path takes no credentials. `PASSWORD_GATE` refuses the - /// same way; `RELAY_GATE` queues briefly instead, as a batch install is one - /// burst of legitimate requests. + /// `span` bounds one window; this bounds how many are built concurrently. + /// Each holds the reader for its entire scan, and the path takes no + /// credentials. `PASSWORD_GATE` refuses the same way; `RELAY_GATE` queues + /// briefly instead, as a batch install is one burst of legitimate requests. #[tokio::test] async fn history_queries_past_the_gate_are_refused_rather_than_queued() { let _serial = HISTORY_TESTS.lock().await; diff --git a/src/db.rs b/src/db.rs index 8642239..0e12a41 100644 --- a/src/db.rs +++ b/src/db.rs @@ -1,18 +1,34 @@ //! SQLite storage. A single writer connection behind a mutex: at a handful of -//! nodes reporting every few seconds, every statement here is sub-millisecond. -// ponytail: single global connection; move to a read pool if the dashboard ever -// blocks behind ingest. +//! nodes reporting every few seconds, every statement on it is sub-millisecond. +//! The history charts, whose scans are not, read through a second connection. use std::collections::{HashMap, HashSet}; use std::sync::Mutex; use anyhow::{Context, Result}; use chrono::{DateTime, Datelike, Local, NaiveDate, Utc}; -use rusqlite::{params, Connection, OptionalExtension}; +use rusqlite::{params, Connection, OpenFlags, OptionalExtension}; use serde::{Deserialize, Serialize}; use tracing::info; -pub struct Db(Mutex); +pub struct Db { + /// A read-only connection to the same file, for the history charts. A week + /// of one node's probe results is a scan of 98 ms at four 60-second probes + /// and 486 ms at eight 10-second ones. With four of the latter in flight on + /// `conn`, every agent report and panel request behind them waited 1.1 s at + /// the median, and 1.8 ms with the scans here: WAL lets this connection + /// read while `conn` commits. `None` for `:memory:`, which a second + /// connection cannot open. + /// + /// Declared before `conn` so that it closes first. The last connection to + /// close folds the WAL into the database file and deletes it, which a + /// read-only one cannot do, and a stopped hub would otherwise leave rows in + /// a -wal that a copy of the database file alone misses. + // ponytail: one reader, so chart requests queue behind one another, at most + // `api::HISTORY_SLOTS` deep; a pool if that wait becomes visible. + reader: Option>, + conn: Mutex, +} const SCHEMA: &str = r#" PRAGMA journal_mode = WAL; @@ -677,6 +693,18 @@ fn own_only(file: &str) { } } +/// Opens [`Db`]'s reader. Read-only from the open, so no statement reaching it +/// can write. It inherits none of the PRAGMAs in `SCHEMA`, so the busy timeout +/// is set again, as in [`Db::backup_into`]. +fn read_only(file: &str) -> Result { + let conn = Connection::open_with_flags( + file, + OpenFlags::SQLITE_OPEN_READ_ONLY | OpenFlags::SQLITE_OPEN_NO_MUTEX, + )?; + conn.busy_timeout(std::time::Duration::from_secs(5))?; + Ok(conn) +} + /// The `main` database's path as SQLite reports it, empty for `:memory:`. /// Queried rather than cached so there is a single answer to which file is /// open. @@ -898,11 +926,21 @@ impl Db { let version: i64 = conn.query_row("PRAGMA user_version", [], |r| r.get(0))?; migrate(&conn, if fresh { SCHEMA_VERSION } else { version })?; - Ok(Self(Mutex::new(conn))) + let reader = match main_file(&conn) { + file if file.is_empty() => None, + file => Some(Mutex::new(read_only(&file)?)), + }; + Ok(Self { reader, conn: Mutex::new(conn) }) } fn conn(&self) -> std::sync::MutexGuard<'_, Connection> { - self.0.lock().unwrap_or_else(|e| e.into_inner()) + self.conn.lock().unwrap_or_else(|e| e.into_inner()) + } + + /// The connection the history charts read through: the read-only one, or + /// the writer for `:memory:`. + fn reader(&self) -> std::sync::MutexGuard<'_, Connection> { + self.reader.as_ref().unwrap_or(&self.conn).lock().unwrap_or_else(|e| e.into_inner()) } // ---- settings ---- @@ -1539,7 +1577,11 @@ impl Db { /// hours between the watermark and that week are missing from the chart /// until they are folded. pub fn metrics(&self, node_id: i64, span: Span) -> Result> { - let conn = self.conn(); + let mut reader = self.reader(); + // One snapshot for every statement below. On the reader, a rollup can + // otherwise commit between reading the watermark and reading the rows + // on either side of it. + let conn = reader.transaction()?; let row = |r: &rusqlite::Row<'_>| { Ok(serde_json::json!({ "ts": r.get::<_, i64>(0)?, "cpu": r.get::<_, f64>(1)?, @@ -1937,7 +1979,7 @@ impl Db { /// routinely carries a hostname or a customer, and the rest of the table /// belongs to nodes this caller may not be able to see. pub fn ping_task_names(&self, node_id: i64) -> Result { - let conn = self.conn(); + let conn = self.reader(); let mut stmt = conn.prepare( "SELECT id, name FROM ping_task WHERE id IN (SELECT task_id FROM ping_node WHERE node_id=?1)", )?; @@ -2018,7 +2060,9 @@ impl Db { node_id: i64, span: Span, ) -> Result<(Vec, serde_json::Value)> { - let conn = self.conn(); + let mut reader = self.reader(); + // One snapshot for every statement below, as in `metrics`. + let conn = reader.transaction()?; let step = span.step; let mut out = Vec::new(); // Per probe in the bucket being filled. @@ -2084,13 +2128,14 @@ impl Db { .enumerate() .map(|(i, id)| id.map(|id| (id, i))) .collect::>()?; - // Sorted after releasing the connection the agents write through. A - // probe missing from the rank, which the assignment filter in + // Sorted after releasing the reader, which the next chart request waits + // on. A probe missing from the rank, which the assignment filter in // `PING_ROWS` rules out today, goes last rather than taking the first // colour. drop(rows); drop(stmt); drop(conn); + drop(reader); out.sort_by_cached_key(|row| { row["task_id"].as_i64().and_then(|id| rank.get(&id).copied()).unwrap_or(usize::MAX) }); @@ -2521,6 +2566,60 @@ mod tests { } } + /// The history charts read through their own connection, so a scan never + /// holds up the writer: a week of eight 10-second probes is 486 ms, which + /// every agent report would otherwise wait out. + #[test] + fn history_is_read_while_the_writer_is_held() { + let scratch = Scratch::new(); + let db = std::sync::Arc::new(Db::open(&scratch.0).unwrap()); + let id = db.create_node(&Node { name: "n".into(), ..Default::default() }, "token").unwrap(); + let task = db + .save_ping_task(&PingTask { + name: "probe".into(), + target: "1.1.1.1:443".into(), + interval: 60, + nodes: vec![id], + ..Default::default() + }) + .unwrap(); + db.insert_metric(id, 60, &serde_json::json!({"cpu": 1.0})).unwrap(); + db.insert_pings(id, &[(task, 60, 42)]).unwrap(); + + let writer = db.conn(); + let (send, receive) = std::sync::mpsc::channel(); + let reading = std::thread::spawn({ + let db = db.clone(); + move || { + let span = Span::minutes(0, 60); + let (metrics, pings) = (db.metrics(id, span).unwrap(), db.ping_records(id, span).unwrap().0); + send.send((metrics, pings, db.ping_task_names(id).unwrap())).unwrap(); + } + }); + // Bounded, so a read that does wait fails here rather than hanging the + // test; released below, it then finishes. + let read = receive.recv_timeout(std::time::Duration::from_secs(5)); + drop(writer); + reading.join().unwrap(); + let (metrics, pings, names) = read.expect("history is read while the writer is held"); + assert_eq!((metrics.len(), pings.len()), (1, 1), "and it sees what the writer committed"); + assert_eq!(names[&task.to_string()], "probe"); + } + + /// A closed database is one file again, as before the reader existed, so a + /// copy of the file taken with the hub stopped holds every row. + #[test] + fn a_closed_database_leaves_no_wal_behind() { + let scratch = Scratch::new(); + let db = Db::open(&scratch.0).unwrap(); + let id = db.create_node(&Node { name: "n".into(), ..Default::default() }, "token").unwrap(); + db.insert_metric(id, 60, &serde_json::json!({"cpu": 1.0})).unwrap(); + // The reader joins the WAL on its first read, not when it opens. + db.metrics(id, Span::minutes(0, 60)).unwrap(); + drop(db); + assert!(!std::path::Path::new(&format!("{}-wal", scratch.0)).exists()); + } + /// Backup and restore are the two operations that can lose every row in the /// database, so this exercises the whole path: take a copy, modify the live /// database, restore the copy, and confirm the change is gone. From ae59c7ac6576c746af217197fea636132d506d55 Mon Sep 17 00:00:00 2001 From: stqfdyr <89149493+stqfdyr@users.noreply.github.com> Date: Thu, 1 Oct 2026 15:21:19 +0800 Subject: [PATCH 2/2] =?UTF-8?q?fix:=20=E5=9B=BE=E8=A1=A8=E8=AF=B7=E6=B1=82?= =?UTF-8?q?=E6=8E=A5=E8=BF=9E=E4=B8=8D=E6=96=AD=E6=97=B6=20WAL=20=E4=B8=8D?= =?UTF-8?q?=E5=86=8D=E6=97=A0=E9=99=90=E5=A2=9E=E9=95=BF?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 只读连接被接连的图表请求占着快照时,自动 checkpoint 重置不了 WAL。取读连接时若 WAL 超过 4 MiB,先在写连接上做一次 TRUNCATE checkpoint,不等待其他读者,导出备份时不会占住写连接 - 导出备份改用同一个只读连接的打开方式 --- src/db.rs | 93 +++++++++++++++++++++++++++++++++++++++++++++---------- 1 file changed, 77 insertions(+), 16 deletions(-) diff --git a/src/db.rs b/src/db.rs index 0e12a41..f00cce4 100644 --- a/src/db.rs +++ b/src/db.rs @@ -15,9 +15,9 @@ pub struct Db { /// A read-only connection to the same file, for the history charts. A week /// of one node's probe results is a scan of 98 ms at four 60-second probes /// and 486 ms at eight 10-second ones. With four of the latter in flight on - /// `conn`, every agent report and panel request behind them waited 1.1 s at - /// the median, and 1.8 ms with the scans here: WAL lets this connection - /// read while `conn` commits. `None` for `:memory:`, which a second + /// `conn`, every agent report and panel request behind them would wait 1.1 s + /// at the median; with the scans here they wait 1.8 ms, as WAL lets this + /// connection read while `conn` commits. `None` for `:memory:`, which a second /// connection cannot open. /// /// Declared before `conn` so that it closes first. The last connection to @@ -693,9 +693,11 @@ fn own_only(file: &str) { } } -/// Opens [`Db`]'s reader. Read-only from the open, so no statement reaching it -/// can write. It inherits none of the PRAGMAs in `SCHEMA`, so the busy timeout -/// is set again, as in [`Db::backup_into`]. +/// Opens a read-only connection to the database: [`Db`]'s reader, and the one +/// [`Db::backup_into`] exports through. Read-only from the open, so no statement +/// reaching it can write. It inherits none of the PRAGMAs in `SCHEMA`, so the +/// busy timeout is set again; without it, a checkpoint racing a read would +/// return SQLITE_BUSY immediately. fn read_only(file: &str) -> Result { let conn = Connection::open_with_flags( file, @@ -939,8 +941,31 @@ impl Db { /// The connection the history charts read through: the read-only one, or /// the writer for `:memory:`. + /// + /// A scan holds its snapshot throughout, and with chart requests queued the + /// next begins as the last ends, so no commit's checkpoint would find the + /// reader idle and the WAL would keep every commit for as long as anyone + /// sent them: 20 MB in 20 s at 100 commits a second, from four requests kept + /// in flight. Before the scan the reader is idle, and a checkpoint taken + /// here resets the WAL itself. Past 4 MiB only, as without readers the WAL + /// stays within the 1 MiB `journal_size_limit` plus one transaction. + /// + /// Without waiting: an export is the one other reader, and waiting on it + /// would hold the writer, and every agent report behind it, for the busy + /// timeout. + /// + /// The writer is taken after the reader here and must never be held while + /// taking the reader. fn reader(&self) -> std::sync::MutexGuard<'_, Connection> { - self.reader.as_ref().unwrap_or(&self.conn).lock().unwrap_or_else(|e| e.into_inner()) + let Some(reader) = &self.reader else { return self.conn() }; + let reader = reader.lock().unwrap_or_else(|e| e.into_inner()); + if bytes_of(&format!("{}-wal", reader.path().unwrap_or_default())) > 4 << 20 { + let conn = self.conn(); + let _ = conn.busy_timeout(std::time::Duration::ZERO); + let _ = conn.query_row("PRAGMA wal_checkpoint(TRUNCATE)", [], |_| Ok(())); + let _ = conn.busy_timeout(std::time::Duration::from_secs(5)); + } + reader } // ---- settings ---- @@ -2209,15 +2234,11 @@ impl Db { /// reads the whole file, so the caller runs it off the runtime -- every other /// statement here is sub-millisecond, this one is not. pub fn backup_into(&self, dest: &str) -> Result<()> { - // A second connection to the same file. `VACUUM INTO` only reads, and WAL - // allows it to read a consistent snapshot while the agents continue - // writing through the first -- exporting is the one heavy operation here - // that need not block them. A fresh connection inherits none of the - // PRAGMAs in SCHEMA, so the busy timeout must be set again or a - // checkpoint racing this read returns SQLITE_BUSY immediately. - let reader = Connection::open(self.file())?; - reader.busy_timeout(std::time::Duration::from_secs(5))?; - reader.execute("VACUUM INTO ?1", [dest])?; + // A connection of its own. `VACUUM INTO` only reads, and WAL allows it to + // read a consistent snapshot while the agents continue writing through + // the first -- exporting is the one heavy operation here that need not + // block them. Not the charts' reader, which it would hold throughout. + read_only(&self.file())?.execute("VACUUM INTO ?1", [dest])?; // The copy is the credential store in one portable file: node tokens in // the clear, the GitHub secret, the password hash. SQLite creates it // under the umask, which at the usual 022 is world-readable. @@ -2606,6 +2627,46 @@ mod tests { assert_eq!(names[&task.to_string()], "probe"); } + /// Chart reads kept back to back leave no commit an idle reader to + /// checkpoint past, so the reader checkpoints for them and the WAL stays + /// bounded however long they continue. Without that, these writes would + /// grow it to 85 MB. + #[test] + fn chart_reads_back_to_back_do_not_grow_the_wal() { + let scratch = Scratch::new(); + let db = std::sync::Arc::new(Db::open(&scratch.0).unwrap()); + let id = db.create_node(&Node { name: "n".into(), ..Default::default() }, "token").unwrap(); + let task = db + .save_ping_task(&PingTask { + name: "probe".into(), + target: "1.1.1.1:443".into(), + interval: 10, + nodes: vec![id], + ..Default::default() + }) + .unwrap(); + let stop = std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false)); + let readers: Vec<_> = (0..4) + .map(|_| { + let (db, stop) = (db.clone(), stop.clone()); + std::thread::spawn(move || { + while !stop.load(std::sync::atomic::Ordering::Relaxed) { + db.ping_records(id, Span::minutes(0, 60)).unwrap(); + } + }) + }) + .collect(); + let wal = format!("{}-wal", scratch.0); + let mut largest = 0; + for ts in 0..20_000 { + db.insert_pings(id, &[(task, ts * 10, 40)]).unwrap(); + largest = largest.max(bytes_of(&wal)); + } + stop.store(true, std::sync::atomic::Ordering::Relaxed); + readers.into_iter().for_each(|r| r.join().unwrap()); + assert!(largest < 16 << 20, "the WAL reached {largest} bytes"); + } + /// A closed database is one file again, as before the reader existed, so a /// copy of the file taken with the hub stopped holds every row. #[test]