Skip to main content

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}