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
41 changes: 19 additions & 22 deletions src/api.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand All @@ -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![] };
Expand Down Expand Up @@ -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();
Expand Down Expand Up @@ -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;
Expand Down
202 changes: 181 additions & 21 deletions src/db.rs
Original file line number Diff line number Diff line change
@@ -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<Connection>);
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 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
/// 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<Mutex<Connection>>,
conn: Mutex<Connection>,
}

const SCHEMA: &str = r#"
PRAGMA journal_mode = WAL;
Expand Down Expand Up @@ -677,6 +693,20 @@ fn own_only(file: &str) {
}
}

/// 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<Connection> {
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.
Expand Down Expand Up @@ -898,11 +928,44 @@ 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:`.
///
/// 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> {
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 ----
Expand Down Expand Up @@ -1539,7 +1602,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<Vec<serde_json::Value>> {
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)?,
Expand Down Expand Up @@ -1937,7 +2004,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<serde_json::Value> {
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)",
)?;
Expand Down Expand Up @@ -2018,7 +2085,9 @@ impl Db {
node_id: i64,
span: Span,
) -> Result<(Vec<serde_json::Value>, 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.
Expand Down Expand Up @@ -2084,13 +2153,14 @@ impl Db {
.enumerate()
.map(|(i, id)| id.map(|id| (id, i)))
.collect::<Result<_, _>>()?;
// 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)
});
Expand Down Expand Up @@ -2164,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.
Expand Down Expand Up @@ -2521,6 +2587,100 @@ 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");
}

/// 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]
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.
Expand Down
Loading