cairn_mod/storage.rs
1//! SQLite storage layer (§10 architecture).
2//!
3//! Owns the connection pool, migration runner, and the small surface of
4//! cross-cutting queries (currently just a health probe). Domain-specific
5//! query modules (labels, reports, audit, etc.) will land alongside the
6//! writer task in subsequent issues.
7
8use std::path::Path;
9use std::time::Duration;
10
11use serde::Serialize;
12use sqlx::{
13 Pool, Sqlite,
14 sqlite::{SqliteConnectOptions, SqliteJournalMode, SqlitePoolOptions},
15};
16
17use crate::error::Result;
18
19/// Embedded migration bundle. Sources from `./migrations/` relative to the
20/// crate root at compile time.
21pub static MIGRATOR: sqlx::migrate::Migrator = sqlx::migrate!("./migrations");
22
23/// Open a SQLite connection pool at `path`, creating the file if it doesn't
24/// exist, and run all pending migrations.
25///
26/// WAL mode and a 5s busy timeout are set per §F5 "Writer architecture."
27/// Foreign keys are enabled (off by default in SQLite).
28pub async fn open<P: AsRef<Path>>(path: P) -> Result<Pool<Sqlite>> {
29 let opts = SqliteConnectOptions::new()
30 .filename(path.as_ref())
31 .create_if_missing(true)
32 .journal_mode(SqliteJournalMode::Wal)
33 .busy_timeout(Duration::from_millis(5000))
34 .foreign_keys(true);
35
36 let pool = SqlitePoolOptions::new()
37 .max_connections(8)
38 .connect_with(opts)
39 .await?;
40
41 MIGRATOR.run(&pool).await?;
42
43 Ok(pool)
44}
45
46/// Point-in-time storage health snapshot.
47#[derive(Debug, Serialize)]
48pub struct Health {
49 /// Total rows in the `labels` table (emitted labels, including negations).
50 pub labels: i64,
51 /// The current single-instance lease row, if one is held.
52 pub lease: Option<Lease>,
53}
54
55/// The singleton row from `server_instance_lease`.
56#[derive(Debug, Serialize, sqlx::FromRow)]
57pub struct Lease {
58 /// Unique instance identifier of the process holding the
59 /// lease.
60 pub instance_id: String,
61 /// Unix epoch milliseconds when this instance first acquired
62 /// the lease.
63 pub acquired_at: i64,
64 /// Unix epoch milliseconds of the most recent heartbeat —
65 /// staleness > `LEASE_STALE_MS` allows takeover.
66 pub last_heartbeat: i64,
67}
68
69/// Query the current storage health: label count and live instance lease.
70///
71/// Intentionally cheap — exists both as an operator-facing sanity check and
72/// as the smoke query that exercises the compile-time-checked `sqlx` offline
73/// cache in CI.
74pub async fn health(pool: &Pool<Sqlite>) -> Result<Health> {
75 let labels: i64 = sqlx::query_scalar!("SELECT COUNT(*) FROM labels")
76 .fetch_one(pool)
77 .await?;
78
79 let lease = sqlx::query_as!(
80 Lease,
81 "SELECT instance_id, acquired_at, last_heartbeat \
82 FROM server_instance_lease WHERE id = 1"
83 )
84 .fetch_optional(pool)
85 .await?;
86
87 Ok(Health { labels, lease })
88}
89
90#[cfg(test)]
91mod tests {
92 use super::*;
93
94 // `:memory:` + a pooled connection is unsound — each pool connection gets
95 // its own independent in-memory DB, so migrations applied on one wouldn't
96 // be visible to the next. Use a tempfile so the pool shares real storage.
97 #[tokio::test]
98 async fn migrations_apply_to_fresh_db() {
99 let dir = tempfile::tempdir().expect("tempdir");
100 let pool = open(dir.path().join("cairn.db")).await.expect("open pool");
101 let h = health(&pool).await.expect("health query");
102 assert_eq!(h.labels, 0);
103 assert!(h.lease.is_none());
104 }
105}