Skip to main content

tuff_console/
store.rs

1//! The console's SQLite file (RFC-108 D6).
2
3use std::path::{Path, PathBuf};
4use std::sync::Mutex;
5
6use rusqlite::{Connection, OptionalExtension, TransactionBehavior, params};
7use serde::Serialize;
8use sha2::{Digest, Sha256};
9use tuff_core::error::{ErrorKind, Result, TuffError};
10use tuff_core::report::{Report, normalize_remote};
11
12use crate::events::{self, Snapshot};
13
14/// File name of the database inside the data directory.
15pub const DATABASE_FILE: &str = "console.sqlite";
16
17/// Prefix of every publish key, so a leaked one is recognisable.
18pub const KEY_PREFIX: &str = "tuffc_";
19
20/// Each entry moves the schema from version `index` to `index + 1`, tracked
21/// in `PRAGMA user_version`. Entries are never edited once released.
22const MIGRATIONS: &[&str] = &[
23    // 1: the D6 tables. `inventory` and `events` are filled by later
24    // milestones (the audit trail and the read API) and exist from the start
25    // so the first release of the file format holds the whole schema.
26    "
27    CREATE TABLE projects (
28        id INTEGER PRIMARY KEY,
29        repository TEXT NOT NULL,
30        path TEXT NOT NULL,
31        name TEXT NOT NULL,
32        first_report_at TEXT NOT NULL,
33        last_report_at TEXT NOT NULL,
34        UNIQUE (repository, path)
35    );
36    CREATE TABLE reports (
37        id INTEGER PRIMARY KEY,
38        project_id INTEGER NOT NULL REFERENCES projects (id),
39        received_at TEXT NOT NULL,
40        generated_at TEXT NOT NULL,
41        commit_sha TEXT,
42        branch TEXT,
43        tuff_version TEXT NOT NULL,
44        digest TEXT NOT NULL,
45        body TEXT NOT NULL
46    );
47    CREATE INDEX reports_by_project ON reports (project_id, id);
48    CREATE TABLE inventory (
49        project_id INTEGER NOT NULL REFERENCES projects (id),
50        capability_type TEXT NOT NULL,
51        capability_id TEXT NOT NULL,
52        target TEXT NOT NULL,
53        version TEXT,
54        source TEXT,
55        status TEXT NOT NULL,
56        PRIMARY KEY (project_id, capability_type, capability_id, target)
57    );
58    CREATE TABLE events (
59        id INTEGER PRIMARY KEY,
60        project_id INTEGER NOT NULL REFERENCES projects (id),
61        report_id INTEGER NOT NULL REFERENCES reports (id),
62        kind TEXT NOT NULL,
63        capability_type TEXT,
64        capability_id TEXT,
65        target TEXT,
66        detail TEXT,
67        commit_sha TEXT,
68        occurred_at TEXT NOT NULL
69    );
70    CREATE INDEX events_by_project ON events (project_id, id);
71    CREATE TABLE keys (
72        name TEXT PRIMARY KEY,
73        sha256 TEXT NOT NULL UNIQUE,
74        created_at TEXT NOT NULL,
75        last_used_at TEXT
76    );
77    ",
78    // 2: a key may be bound to one repository (D5).
79    "ALTER TABLE keys ADD COLUMN repository TEXT;",
80];
81
82/// The data directory when `--data` is not given:
83/// `$XDG_DATA_HOME/tuff/console`, or `~/.local/share/tuff/console`.
84pub fn default_data_dir(home: &Path) -> PathBuf {
85    tuff_core::paths::user_data(home).join("console")
86}
87
88/// What ingesting a report did.
89#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
90#[serde(rename_all = "camelCase")]
91pub struct IngestOutcome {
92    pub project_id: i64,
93    /// The stored report, or the previous one when `deduplicated`.
94    pub report_id: i64,
95    /// The report equalled the project's previous one, so only the
96    /// project's last report time changed.
97    pub deduplicated: bool,
98    pub project_first_seen: bool,
99}
100
101#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
102#[serde(rename_all = "camelCase")]
103pub struct ProjectRow {
104    pub id: i64,
105    pub repository: String,
106    pub path: String,
107    pub name: String,
108    pub first_report_at: String,
109    pub last_report_at: String,
110    pub report_count: i64,
111}
112
113#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
114#[serde(rename_all = "camelCase")]
115pub struct KeyInfo {
116    pub name: String,
117    /// The one repository the key may publish for, when it is scoped.
118    pub repository: Option<String>,
119    pub created_at: String,
120    pub last_used_at: Option<String>,
121}
122
123/// One stored report, without its body.
124#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
125#[serde(rename_all = "camelCase")]
126pub struct ReportSummary {
127    pub id: i64,
128    pub received_at: String,
129    pub generated_at: String,
130    pub commit: Option<String>,
131    pub branch: Option<String>,
132    pub tuff_version: String,
133    pub digest: String,
134}
135
136/// What a live key allows.
137#[derive(Debug, Clone, PartialEq, Eq)]
138pub struct KeyGrant {
139    pub name: String,
140    /// The repository the key is bound to, normalised; `None` for any.
141    pub repository: Option<String>,
142}
143
144/// One change recorded between two consecutive reports of a project.
145#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
146#[serde(rename_all = "camelCase")]
147pub struct EventRow {
148    pub id: i64,
149    pub project_id: i64,
150    pub report_id: i64,
151    pub kind: String,
152    pub capability_type: Option<String>,
153    pub capability_id: Option<String>,
154    pub target: Option<String>,
155    pub detail: Option<String>,
156    pub commit: Option<String>,
157    pub occurred_at: String,
158}
159
160/// Which events to list. Empty fields match everything.
161#[derive(Debug, Clone, Default)]
162pub struct EventFilter {
163    pub project_id: Option<i64>,
164    pub capability: Option<String>,
165    pub kind: Option<String>,
166    /// RFC 3339 time or a date; events from then on.
167    pub since: Option<String>,
168    /// Only events older than this event id, for paging back.
169    pub before: Option<i64>,
170    pub limit: Option<u32>,
171}
172
173pub struct Store {
174    conn: Mutex<Connection>,
175}
176
177fn db_error(error: rusqlite::Error) -> TuffError {
178    TuffError::of(ErrorKind::Source, format!("console database: {error}")).with_source(error)
179}
180
181fn now() -> String {
182    chrono::Utc::now().to_rfc3339_opts(chrono::SecondsFormat::Secs, true)
183}
184
185fn sha256_hex(bytes: &[u8]) -> String {
186    hex(&Sha256::digest(bytes))
187}
188
189fn hex(bytes: &[u8]) -> String {
190    use std::fmt::Write;
191    bytes.iter().fold(String::new(), |mut out, byte| {
192        let _ = write!(out, "{byte:02x}");
193        out
194    })
195}
196
197/// JSON with object keys in sorted order and no whitespace, so equal
198/// documents produce equal bytes whatever order their keys arrived in.
199fn write_canonical(value: &serde_json::Value, out: &mut String) {
200    match value {
201        serde_json::Value::Object(map) => {
202            let mut keys: Vec<_> = map.keys().collect();
203            keys.sort();
204            out.push('{');
205            for (index, key) in keys.into_iter().enumerate() {
206                if index > 0 {
207                    out.push(',');
208                }
209                out.push_str(&serde_json::Value::String(key.clone()).to_string());
210                out.push(':');
211                write_canonical(&map[key], out);
212            }
213            out.push('}');
214        }
215        serde_json::Value::Array(items) => {
216            out.push('[');
217            for (index, item) in items.iter().enumerate() {
218                if index > 0 {
219                    out.push(',');
220                }
221                write_canonical(item, out);
222            }
223            out.push(']');
224        }
225        other => out.push_str(&other.to_string()),
226    }
227}
228
229/// SHA-256 over the report's canonical JSON without `generatedAt` and
230/// the project's `commit`, `branch`, and `dirty` (D6). Those describe when
231/// and where the report was taken; two reports that differ only there say
232/// the same thing about the project, so a CI job publishing on every push
233/// does not add a row per commit.
234pub fn report_digest(report: &serde_json::Value) -> String {
235    let mut value = report.clone();
236    if let Some(map) = value.as_object_mut() {
237        map.remove("generatedAt");
238        if let Some(project) = map.get_mut("project").and_then(|p| p.as_object_mut()) {
239            for key in ["commit", "branch", "dirty"] {
240                project.remove(key);
241            }
242        }
243    }
244    let mut canonical = String::new();
245    write_canonical(&value, &mut canonical);
246    sha256_hex(canonical.as_bytes())
247}
248
249fn valid_key_name(name: &str) -> bool {
250    !name.is_empty()
251        && name.len() <= 64
252        && name
253            .chars()
254            .all(|c| c.is_ascii_alphanumeric() || matches!(c, '-' | '_' | '.'))
255}
256
257impl Store {
258    /// Open, creating the directory and file when missing, and bring the
259    /// schema up to date.
260    pub fn open(data_dir: &Path) -> Result<Self> {
261        std::fs::create_dir_all(data_dir).map_err(|error| {
262            TuffError::of(
263                ErrorKind::Io,
264                format!("cannot create {}: {error}", data_dir.display()),
265            )
266            .with_hint("pass --data <dir> with a folder you can write to")
267        })?;
268        let path = data_dir.join(DATABASE_FILE);
269        let created = !path.exists();
270        let conn = Connection::open(&path).map_err(|error| {
271            TuffError::of(
272                ErrorKind::Source,
273                format!("cannot open {}: {error}", path.display()),
274            )
275            .with_hint("pass --data <dir> with a folder you can write to")
276        })?;
277        #[cfg(unix)]
278        if created {
279            use std::os::unix::fs::PermissionsExt;
280            let _ = std::fs::set_permissions(&path, std::fs::Permissions::from_mode(0o600));
281        }
282        #[cfg(not(unix))]
283        let _ = created;
284        Self::from_connection(conn, &path.display().to_string())
285    }
286
287    /// An in-memory store, for tests.
288    pub fn open_in_memory() -> Result<Self> {
289        Self::from_connection(Connection::open_in_memory().map_err(db_error)?, ":memory:")
290    }
291
292    fn from_connection(mut conn: Connection, label: &str) -> Result<Self> {
293        conn.busy_timeout(std::time::Duration::from_secs(5))
294            .map_err(db_error)?;
295        conn.pragma_update(None, "foreign_keys", true)
296            .map_err(db_error)?;
297        // Not every database can change mode, such as one in memory.
298        let _ = conn.pragma_update(None, "journal_mode", "WAL");
299        migrate(&mut conn, label)?;
300        Ok(Self {
301            conn: Mutex::new(conn),
302        })
303    }
304
305    fn conn(&self) -> std::sync::MutexGuard<'_, Connection> {
306        // A panic while holding the lock leaves the connection usable.
307        self.conn
308            .lock()
309            .unwrap_or_else(|poisoned| poisoned.into_inner())
310    }
311
312    /// The schema version of the file.
313    pub fn schema_version(&self) -> Result<u32> {
314        self.conn()
315            .pragma_query_value(None, "user_version", |row| row.get(0))
316            .map_err(db_error)
317    }
318
319    /// Store one report. `raw` is the report as received, kept verbatim in
320    /// the `reports` row. A report equal to the project's previous one
321    /// (digest over everything except `generatedAt`, commit, branch, and
322    /// dirty) adds no row: the previous row takes the new report's commit,
323    /// branch, body, and times, so the project shows where it was last
324    /// seen, and the project's last report time moves.
325    pub fn ingest(&self, report: &Report, raw: &serde_json::Value) -> Result<IngestOutcome> {
326        self.ingest_at(report, raw, &now())
327    }
328
329    /// [`Store::ingest`] with the time the report counts as received at,
330    /// an RFC 3339 UTC string. `--demo` uses it to give sample reports a
331    /// history.
332    pub fn ingest_at(
333        &self,
334        report: &Report,
335        raw: &serde_json::Value,
336        received_at: &str,
337    ) -> Result<IngestOutcome> {
338        let digest = report_digest(raw);
339        let body = serde_json::to_string(raw)?;
340        let mut conn = self.conn();
341        let tx = conn
342            .transaction_with_behavior(TransactionBehavior::Immediate)
343            .map_err(db_error)?;
344
345        let project: Option<i64> = tx
346            .query_row(
347                "SELECT id FROM projects WHERE repository = ?1 AND path = ?2",
348                params![report.project.repository, report.project.path],
349                |row| row.get(0),
350            )
351            .optional()
352            .map_err(db_error)?;
353        let (project_id, project_first_seen) = match project {
354            Some(id) => (id, false),
355            None => {
356                tx.execute(
357                    "INSERT INTO projects (repository, path, name, first_report_at, last_report_at)
358                     VALUES (?1, ?2, ?3, ?4, ?4)",
359                    params![
360                        report.project.repository,
361                        report.project.path,
362                        report.project.name,
363                        received_at
364                    ],
365                )
366                .map_err(db_error)?;
367                (tx.last_insert_rowid(), true)
368            }
369        };
370
371        let previous: Option<(i64, String, String)> = tx
372            .query_row(
373                "SELECT id, digest, body FROM reports WHERE project_id = ?1 ORDER BY id DESC LIMIT 1",
374                params![project_id],
375                |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?)),
376            )
377            .optional()
378            .map_err(db_error)?;
379
380        let outcome = match previous {
381            Some((report_id, previous_digest, _)) if previous_digest == digest => {
382                tx.execute(
383                    "UPDATE reports SET received_at = ?2, generated_at = ?3, commit_sha = ?4,
384                                        branch = ?5, body = ?6
385                     WHERE id = ?1",
386                    params![
387                        report_id,
388                        received_at,
389                        report.generated_at,
390                        report.project.commit,
391                        report.project.branch,
392                        body
393                    ],
394                )
395                .map_err(db_error)?;
396                tx.execute(
397                    "UPDATE projects SET last_report_at = ?2 WHERE id = ?1",
398                    params![project_id, received_at],
399                )
400                .map_err(db_error)?;
401                IngestOutcome {
402                    project_id,
403                    report_id,
404                    deduplicated: true,
405                    project_first_seen,
406                }
407            }
408            previous => {
409                tx.execute(
410                    "INSERT INTO reports (project_id, received_at, generated_at, commit_sha, branch,
411                                          tuff_version, digest, body)
412                     VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8)",
413                    params![
414                        project_id,
415                        received_at,
416                        report.generated_at,
417                        report.project.commit,
418                        report.project.branch,
419                        report.tuff_version,
420                        digest,
421                        body
422                    ],
423                )
424                .map_err(db_error)?;
425                let report_id = tx.last_insert_rowid();
426                tx.execute(
427                    "UPDATE projects SET name = ?2, last_report_at = ?3 WHERE id = ?1",
428                    params![project_id, report.project.name, received_at],
429                )
430                .map_err(db_error)?;
431
432                let current = Snapshot::from_report(raw);
433                let before = previous
434                    .and_then(|(_, _, body)| serde_json::from_str::<serde_json::Value>(&body).ok())
435                    .map(|body| Snapshot::from_report(&body));
436                for event in events::diff(before.as_ref(), &current) {
437                    tx.execute(
438                        "INSERT INTO events (project_id, report_id, kind, capability_type,
439                                             capability_id, target, detail, commit_sha, occurred_at)
440                         VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9)",
441                        params![
442                            project_id,
443                            report_id,
444                            event.kind,
445                            event.capability_type,
446                            event.capability_id,
447                            event.target,
448                            event.detail,
449                            report.project.commit,
450                            received_at
451                        ],
452                    )
453                    .map_err(db_error)?;
454                }
455                tx.execute(
456                    "DELETE FROM inventory WHERE project_id = ?1",
457                    params![project_id],
458                )
459                .map_err(db_error)?;
460                for row in &current.rows {
461                    tx.execute(
462                        "INSERT OR REPLACE INTO inventory (project_id, capability_type, capability_id,
463                                                          target, version, source, status)
464                         VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)",
465                        params![
466                            project_id,
467                            row.capability_type,
468                            row.capability_id,
469                            row.target,
470                            row.version,
471                            row.source,
472                            row.status
473                        ],
474                    )
475                    .map_err(db_error)?;
476                }
477                IngestOutcome {
478                    project_id,
479                    report_id,
480                    deduplicated: false,
481                    project_first_seen,
482                }
483            }
484        };
485        tx.commit().map_err(db_error)?;
486        Ok(outcome)
487    }
488
489    /// Every project, in the order they were first seen.
490    pub fn projects(&self) -> Result<Vec<ProjectRow>> {
491        let conn = self.conn();
492        let mut statement = conn
493            .prepare(
494                "SELECT p.id, p.repository, p.path, p.name, p.first_report_at, p.last_report_at,
495                        (SELECT COUNT(*) FROM reports r WHERE r.project_id = p.id)
496                 FROM projects p ORDER BY p.id",
497            )
498            .map_err(db_error)?;
499        statement
500            .query_map([], project_row)
501            .map_err(db_error)?
502            .collect::<std::result::Result<_, _>>()
503            .map_err(db_error)
504    }
505
506    /// One project and its latest report as stored.
507    pub fn project(&self, id: i64) -> Result<Option<(ProjectRow, serde_json::Value)>> {
508        let conn = self.conn();
509        let row = conn
510            .query_row(
511                "SELECT p.id, p.repository, p.path, p.name, p.first_report_at, p.last_report_at,
512                        (SELECT COUNT(*) FROM reports r WHERE r.project_id = p.id)
513                 FROM projects p WHERE p.id = ?1",
514                params![id],
515                project_row,
516            )
517            .optional()
518            .map_err(db_error)?;
519        let Some(row) = row else { return Ok(None) };
520        let body: String = conn
521            .query_row(
522                "SELECT body FROM reports WHERE project_id = ?1 ORDER BY id DESC LIMIT 1",
523                params![id],
524                |row| row.get(0),
525            )
526            .map_err(db_error)?;
527        Ok(Some((row, serde_json::from_str(&body)?)))
528    }
529
530    /// Every project with its latest report as stored, in the order the
531    /// projects were first seen.
532    pub fn latest_reports(&self) -> Result<Vec<(ProjectRow, serde_json::Value)>> {
533        let conn = self.conn();
534        let mut statement = conn
535            .prepare(
536                "SELECT p.id, p.repository, p.path, p.name, p.first_report_at, p.last_report_at,
537                        (SELECT COUNT(*) FROM reports r WHERE r.project_id = p.id),
538                        (SELECT body FROM reports r WHERE r.project_id = p.id
539                         ORDER BY r.id DESC LIMIT 1)
540                 FROM projects p ORDER BY p.id",
541            )
542            .map_err(db_error)?;
543        let rows = statement
544            .query_map([], |row| {
545                Ok((project_row(row)?, row.get::<_, Option<String>>(7)?))
546            })
547            .map_err(db_error)?;
548        let mut out = Vec::new();
549        for row in rows {
550            let (project, body) = row.map_err(db_error)?;
551            if let Some(body) = body {
552                out.push((project, serde_json::from_str(&body)?));
553            }
554        }
555        Ok(out)
556    }
557
558    /// The stored reports of a project, newest first, without their bodies.
559    pub fn report_history(&self, project_id: i64, limit: u32) -> Result<Vec<ReportSummary>> {
560        let conn = self.conn();
561        let mut statement = conn
562            .prepare(
563                "SELECT id, received_at, generated_at, commit_sha, branch, tuff_version, digest
564                 FROM reports WHERE project_id = ?1 ORDER BY id DESC LIMIT ?2",
565            )
566            .map_err(db_error)?;
567        statement
568            .query_map(params![project_id, i64::from(limit)], |row| {
569                Ok(ReportSummary {
570                    id: row.get(0)?,
571                    received_at: row.get(1)?,
572                    generated_at: row.get(2)?,
573                    commit: row.get(3)?,
574                    branch: row.get(4)?,
575                    tuff_version: row.get(5)?,
576                    digest: row.get(6)?,
577                })
578            })
579            .map_err(db_error)?
580            .collect::<std::result::Result<_, _>>()
581            .map_err(db_error)
582    }
583
584    /// Create a publish key and return its secret. Only the SHA-256 of the
585    /// secret is stored, so this is the one time the secret exists in full.
586    /// With `repository`, the key may publish for that repository only.
587    pub fn create_key(&self, name: &str, repository: Option<&str>) -> Result<String> {
588        if !valid_key_name(name) {
589            return Err(
590                TuffError::usage(format!("'{name}' is not a valid key name"))
591                    .with_hint("use 1 to 64 letters, digits, '-', '_' or '.'"),
592            );
593        }
594        let repository = match repository {
595            Some(text) => {
596                let normalized = normalize_remote(text);
597                if !normalized.contains('/') {
598                    return Err(
599                        TuffError::usage(format!("'{text}' does not name a repository"))
600                            .with_hint("use host/owner/name, such as github.com/acme/web"),
601                    );
602                }
603                Some(normalized)
604            }
605            None => None,
606        };
607        let mut bytes = [0u8; 32];
608        getrandom::fill(&mut bytes).map_err(|error| {
609            TuffError::of(
610                ErrorKind::Internal,
611                format!("no random bytes for a key: {error}"),
612            )
613        })?;
614        let secret = format!("{KEY_PREFIX}{}", hex(&bytes));
615        let inserted = self.conn().execute(
616            "INSERT INTO keys (name, sha256, created_at, repository) VALUES (?1, ?2, ?3, ?4)",
617            params![name, sha256_hex(secret.as_bytes()), now(), repository],
618        );
619        match inserted {
620            Ok(_) => Ok(secret),
621            Err(rusqlite::Error::SqliteFailure(failure, _))
622                if failure.code == rusqlite::ErrorCode::ConstraintViolation =>
623            {
624                Err(
625                    TuffError::refused(format!("a key named '{name}' already exists")).with_hint(
626                        format!("run 'tuff console key revoke {name}' first, or pick another name"),
627                    ),
628                )
629            }
630            Err(error) => Err(db_error(error)),
631        }
632    }
633
634    /// Keys in creation order.
635    pub fn keys(&self) -> Result<Vec<KeyInfo>> {
636        let conn = self.conn();
637        let mut statement = conn
638            .prepare(
639                "SELECT name, repository, created_at, last_used_at FROM keys
640                 ORDER BY created_at, name",
641            )
642            .map_err(db_error)?;
643        statement
644            .query_map([], |row| {
645                Ok(KeyInfo {
646                    name: row.get(0)?,
647                    repository: row.get(1)?,
648                    created_at: row.get(2)?,
649                    last_used_at: row.get(3)?,
650                })
651            })
652            .map_err(db_error)?
653            .collect::<std::result::Result<_, _>>()
654            .map_err(db_error)
655    }
656
657    pub fn key_count(&self) -> Result<u64> {
658        self.conn()
659            .query_row("SELECT COUNT(*) FROM keys", [], |row| row.get(0))
660            .map_err(db_error)
661    }
662
663    pub fn revoke_key(&self, name: &str) -> Result<()> {
664        let removed = self
665            .conn()
666            .execute("DELETE FROM keys WHERE name = ?1", params![name])
667            .map_err(db_error)?;
668        if removed == 0 {
669            return Err(TuffError::not_found(format!("no key named '{name}'"))
670                .with_hint("run 'tuff console key list' to see the names"));
671        }
672        Ok(())
673    }
674
675    /// The grant of `secret` when it is a live key. A match records its use.
676    pub fn verify_key(&self, secret: &str) -> Result<Option<KeyGrant>> {
677        let hash = sha256_hex(secret.as_bytes());
678        let conn = self.conn();
679        let grant = conn
680            .query_row(
681                "SELECT name, repository FROM keys WHERE sha256 = ?1",
682                params![hash],
683                |row| {
684                    Ok(KeyGrant {
685                        name: row.get(0)?,
686                        repository: row.get(1)?,
687                    })
688                },
689            )
690            .optional()
691            .map_err(db_error)?;
692        if grant.is_some() {
693            conn.execute(
694                "UPDATE keys SET last_used_at = ?2 WHERE sha256 = ?1",
695                params![hash, now()],
696            )
697            .map_err(db_error)?;
698        }
699        Ok(grant)
700    }
701
702    /// Recorded events, newest first.
703    pub fn events(&self, filter: &EventFilter) -> Result<Vec<EventRow>> {
704        let conn = self.conn();
705        let mut statement = conn
706            .prepare(
707                "SELECT e.id, e.project_id, e.report_id, e.kind, e.capability_type,
708                        e.capability_id, e.target, e.detail, e.commit_sha, e.occurred_at
709                 FROM events e
710                 WHERE (?1 IS NULL OR e.project_id = ?1)
711                   AND (?2 IS NULL OR e.capability_id = ?2)
712                   AND (?3 IS NULL OR e.kind = ?3)
713                   AND (?4 IS NULL OR e.occurred_at >= ?4)
714                   AND (?6 IS NULL OR e.id < ?6)
715                 ORDER BY e.id DESC
716                 LIMIT ?5",
717            )
718            .map_err(db_error)?;
719        statement
720            .query_map(
721                params![
722                    filter.project_id,
723                    filter.capability,
724                    filter.kind,
725                    filter.since,
726                    i64::from(filter.limit.unwrap_or(500)),
727                    filter.before
728                ],
729                |row| {
730                    Ok(EventRow {
731                        id: row.get(0)?,
732                        project_id: row.get(1)?,
733                        report_id: row.get(2)?,
734                        kind: row.get(3)?,
735                        capability_type: row.get(4)?,
736                        capability_id: row.get(5)?,
737                        target: row.get(6)?,
738                        detail: row.get(7)?,
739                        commit: row.get(8)?,
740                        occurred_at: row.get(9)?,
741                    })
742                },
743            )
744            .map_err(db_error)?
745            .collect::<std::result::Result<_, _>>()
746            .map_err(db_error)
747    }
748}
749
750fn project_row(row: &rusqlite::Row<'_>) -> rusqlite::Result<ProjectRow> {
751    Ok(ProjectRow {
752        id: row.get(0)?,
753        repository: row.get(1)?,
754        path: row.get(2)?,
755        name: row.get(3)?,
756        first_report_at: row.get(4)?,
757        last_report_at: row.get(5)?,
758        report_count: row.get(6)?,
759    })
760}
761
762fn migrate(conn: &mut Connection, label: &str) -> Result<()> {
763    let current: u32 = conn
764        .pragma_query_value(None, "user_version", |row| row.get(0))
765        .map_err(db_error)?;
766    let known = MIGRATIONS.len() as u32;
767    if current > known {
768        return Err(TuffError::unsupported(format!(
769            "{label} has schema version {current}, and this tuff reads up to {known}"
770        ))
771        .with_hint("upgrade tuff, or point --data at another folder"));
772    }
773    for (index, sql) in MIGRATIONS.iter().enumerate().skip(current as usize) {
774        let tx = conn.transaction().map_err(db_error)?;
775        tx.execute_batch(sql).map_err(db_error)?;
776        tx.pragma_update(None, "user_version", index as u32 + 1)
777            .map_err(db_error)?;
778        tx.commit().map_err(db_error)?;
779    }
780    Ok(())
781}
782
783#[cfg(test)]
784mod tests {
785    use super::*;
786    use serde_json::json;
787    use tuff_core::report::ProjectIdentity;
788
789    fn ingest(
790        store: &Store,
791        generated_at: &str,
792        commit: &str,
793        lockfile: serde_json::Value,
794    ) -> IngestOutcome {
795        let report = Report {
796            schema: 1,
797            tuff_version: "0.12.0".into(),
798            generated_at: generated_at.into(),
799            project: ProjectIdentity {
800                repository: "github.com/acme/agents".into(),
801                path: "apps/billing-agent".into(),
802                name: "billing-agent".into(),
803                commit: Some(commit.into()),
804                branch: Some("main".into()),
805                dirty: false,
806            },
807            lockfile,
808            check: json!({ "valid": true, "results": [] }),
809            outdated: None,
810        };
811        let raw = serde_json::to_value(&report).unwrap();
812        store.ingest(&report, &raw).unwrap()
813    }
814
815    #[test]
816    fn a_new_database_is_migrated_to_the_latest_schema() {
817        let store = Store::open_in_memory().unwrap();
818        assert_eq!(store.schema_version().unwrap(), MIGRATIONS.len() as u32);
819        assert!(store.projects().unwrap().is_empty());
820    }
821
822    #[test]
823    fn a_database_from_a_newer_tuff_is_refused() {
824        let temp = tempfile::tempdir().unwrap();
825        {
826            let store = Store::open(temp.path()).unwrap();
827            store
828                .conn()
829                .pragma_update(None, "user_version", 99)
830                .unwrap();
831        }
832        let error = Store::open(temp.path()).err().unwrap();
833        assert_eq!(error.kind(), ErrorKind::Unsupported);
834        assert!(error.hint().unwrap().contains("upgrade tuff"));
835    }
836
837    #[test]
838    fn reopening_keeps_the_data() {
839        let temp = tempfile::tempdir().unwrap();
840        Store::open(temp.path())
841            .unwrap()
842            .create_key("ci", None)
843            .unwrap();
844        assert_eq!(Store::open(temp.path()).unwrap().key_count().unwrap(), 1);
845    }
846
847    #[test]
848    fn the_first_report_creates_the_project() {
849        let store = Store::open_in_memory().unwrap();
850        let outcome = ingest(
851            &store,
852            "2026-09-16T18:00:00Z",
853            "aaa",
854            json!({ "version": 3 }),
855        );
856        assert!(outcome.project_first_seen);
857        assert!(!outcome.deduplicated);
858        let projects = store.projects().unwrap();
859        assert_eq!(projects.len(), 1);
860        assert_eq!(projects[0].repository, "github.com/acme/agents");
861        assert_eq!(projects[0].path, "apps/billing-agent");
862        assert_eq!(projects[0].report_count, 1);
863    }
864
865    #[test]
866    fn a_report_equal_but_for_generated_at_is_deduplicated() {
867        let store = Store::open_in_memory().unwrap();
868        let first = ingest(
869            &store,
870            "2026-09-16T18:00:00Z",
871            "aaa",
872            json!({ "version": 3 }),
873        );
874        let second = ingest(
875            &store,
876            "2026-09-16T19:30:00Z",
877            "aaa",
878            json!({ "version": 3 }),
879        );
880        assert!(second.deduplicated);
881        assert!(!second.project_first_seen);
882        assert_eq!(second.report_id, first.report_id);
883        assert_eq!(store.projects().unwrap()[0].report_count, 1);
884    }
885
886    #[test]
887    fn a_changed_report_is_stored_and_a_return_to_an_older_one_too() {
888        let store = Store::open_in_memory().unwrap();
889        let a = ingest(
890            &store,
891            "2026-09-16T18:00:00Z",
892            "aaa",
893            json!({ "version": 3 }),
894        );
895        let b = ingest(
896            &store,
897            "2026-09-16T18:05:00Z",
898            "bbb",
899            json!({ "version": 3, "capabilities": [] }),
900        );
901        let c = ingest(
902            &store,
903            "2026-09-16T18:10:00Z",
904            "ccc",
905            json!({ "version": 3 }),
906        );
907        assert!(!b.deduplicated && !c.deduplicated);
908        assert_ne!(a.report_id, b.report_id);
909        assert_ne!(b.report_id, c.report_id);
910        assert_eq!(store.projects().unwrap()[0].report_count, 3);
911    }
912
913    #[test]
914    fn a_new_commit_with_the_same_content_moves_the_latest_report() {
915        let store = Store::open_in_memory().unwrap();
916        let first = ingest(
917            &store,
918            "2026-09-16T18:00:00Z",
919            "aaa",
920            json!({ "version": 3 }),
921        );
922        let second = ingest(
923            &store,
924            "2026-09-17T09:00:00Z",
925            "bbb",
926            json!({ "version": 3 }),
927        );
928        assert!(second.deduplicated);
929        assert_eq!(second.report_id, first.report_id);
930        let reports = store.report_history(first.project_id, 10).unwrap();
931        assert_eq!(reports.len(), 1);
932        assert_eq!(reports[0].commit.as_deref(), Some("bbb"));
933        assert_eq!(reports[0].generated_at, "2026-09-17T09:00:00Z");
934    }
935
936    #[test]
937    fn digests_ignore_key_order_and_generated_at_only() {
938        let one = json!({ "generatedAt": "x", "b": [1, { "d": 1, "c": 2 }], "a": null });
939        let two = json!({ "a": null, "generatedAt": "y", "b": [1, { "c": 2, "d": 1 }] });
940        assert_eq!(report_digest(&one), report_digest(&two));
941        let three = json!({ "a": null, "b": [1, { "c": 2, "d": 2 }] });
942        assert_ne!(report_digest(&one), report_digest(&three));
943        let at = |commit: &str, dirty: bool| json!({ "project": { "repository": "r", "commit": commit, "branch": "main", "dirty": dirty }, "x": 1 });
944        assert_eq!(
945            report_digest(&at("aaa", false)),
946            report_digest(&at("bbb", true))
947        );
948        let elsewhere = json!({ "project": { "repository": "s", "commit": "aaa" }, "x": 1 });
949        assert_ne!(report_digest(&at("aaa", false)), report_digest(&elsewhere));
950    }
951
952    #[test]
953    fn the_latest_report_is_read_back_as_stored() {
954        let store = Store::open_in_memory().unwrap();
955        ingest(
956            &store,
957            "2026-09-16T18:00:00Z",
958            "aaa",
959            json!({ "version": 3 }),
960        );
961        let outcome = ingest(
962            &store,
963            "2026-09-16T18:05:00Z",
964            "bbb",
965            json!({ "version": 3 }),
966        );
967        let (project, latest) = store.project(outcome.project_id).unwrap().unwrap();
968        assert_eq!(project.report_count, 1, "only the commit changed");
969        assert_eq!(latest["project"]["commit"], "bbb");
970        assert!(store.project(999).unwrap().is_none());
971    }
972
973    #[test]
974    fn a_key_secret_is_shown_once_and_only_its_hash_is_stored() {
975        let store = Store::open_in_memory().unwrap();
976        let secret = store.create_key("ci", None).unwrap();
977        assert!(secret.starts_with(KEY_PREFIX));
978        assert_eq!(secret.len(), KEY_PREFIX.len() + 64);
979
980        let stored: String = store
981            .conn()
982            .query_row("SELECT sha256 FROM keys WHERE name = 'ci'", [], |row| {
983                row.get(0)
984            })
985            .unwrap();
986        assert_eq!(stored, sha256_hex(secret.as_bytes()));
987        assert!(!stored.contains(&secret));
988
989        assert!(store.verify_key(&secret).unwrap().is_some());
990        assert!(store.verify_key("tuffc_wrong").unwrap().is_none());
991        assert!(store.verify_key("").unwrap().is_none());
992        let keys = store.keys().unwrap();
993        assert_eq!(keys.len(), 1);
994        assert_eq!(keys[0].name, "ci");
995        assert!(keys[0].last_used_at.is_some());
996    }
997
998    #[test]
999    fn a_key_is_unused_until_it_authenticates() {
1000        let store = Store::open_in_memory().unwrap();
1001        store.create_key("ci", None).unwrap();
1002        assert!(store.keys().unwrap()[0].last_used_at.is_none());
1003    }
1004
1005    #[test]
1006    fn a_revoked_key_stops_working() {
1007        let store = Store::open_in_memory().unwrap();
1008        let secret = store.create_key("ci", None).unwrap();
1009        store.revoke_key("ci").unwrap();
1010        assert!(store.verify_key(&secret).unwrap().is_none());
1011        assert_eq!(store.key_count().unwrap(), 0);
1012        let error = store.revoke_key("ci").unwrap_err();
1013        assert_eq!(error.kind(), ErrorKind::NotFound);
1014    }
1015
1016    #[test]
1017    fn key_names_are_unique_and_validated() {
1018        let store = Store::open_in_memory().unwrap();
1019        store.create_key("ci", None).unwrap();
1020        assert_eq!(
1021            store.create_key("ci", None).unwrap_err().kind(),
1022            ErrorKind::Refused
1023        );
1024        for bad in ["", "has space", "a/b", &"x".repeat(65)] {
1025            assert_eq!(
1026                store.create_key(bad, None).unwrap_err().kind(),
1027                ErrorKind::Usage,
1028                "{bad:?}"
1029            );
1030        }
1031        let one = store.create_key("one", None).unwrap();
1032        let two = store.create_key("two", None).unwrap();
1033        assert_ne!(one, two);
1034    }
1035
1036    fn carrying(
1037        commit: &str,
1038        version: &str,
1039        status: &str,
1040        gap: bool,
1041    ) -> (Report, serde_json::Value) {
1042        let mut policy = json!({
1043            "name": "guard", "type": "policy", "target": "codex", "version": "1",
1044            "source": { "kind": "local", "path": "guard" },
1045        });
1046        if gap {
1047            policy["unenforced_rules"] =
1048                json!([{ "rule": 2, "description": "deny read", "reason": "none" }]);
1049        }
1050        let report = Report {
1051            schema: 1,
1052            tuff_version: "0.12.0".into(),
1053            generated_at: "2026-09-30T10:00:00Z".into(),
1054            project: ProjectIdentity {
1055                repository: "github.com/acme/agents".into(),
1056                path: ".".into(),
1057                name: "agents".into(),
1058                commit: Some(commit.into()),
1059                branch: None,
1060                dirty: false,
1061            },
1062            lockfile: json!({ "version": 3, "capabilities": [
1063                { "name": "lint", "type": "skill", "target": "claude", "version": version,
1064                  "source": { "kind": "local", "path": "lint" } },
1065                policy,
1066            ] }),
1067            check: json!({ "valid": true, "results": [
1068                { "id": "lint", "type": "skill", "target": "claude", "status": status },
1069            ] }),
1070            outdated: None,
1071        };
1072        let raw = serde_json::to_value(&report).unwrap();
1073        (report, raw)
1074    }
1075
1076    fn kinds_of(store: &Store) -> Vec<String> {
1077        let mut rows = store.events(&EventFilter::default()).unwrap();
1078        rows.reverse();
1079        rows.into_iter().map(|row| row.kind).collect()
1080    }
1081
1082    #[test]
1083    fn two_different_reports_record_the_events_between_them() {
1084        let store = Store::open_in_memory().unwrap();
1085        let (report, raw) = carrying("aaa", "1.0.0", "ok", true);
1086        store.ingest(&report, &raw).unwrap();
1087        assert_eq!(
1088            kinds_of(&store),
1089            [
1090                "project_first_seen",
1091                "capability_added",
1092                "capability_added",
1093                "policy_gap_added"
1094            ]
1095        );
1096
1097        let (report, raw) = carrying("bbb", "1.1.0", "modified", false);
1098        let outcome = store.ingest(&report, &raw).unwrap();
1099        let all = store.events(&EventFilter::default()).unwrap();
1100        let second: Vec<_> = all
1101            .iter()
1102            .filter(|row| row.report_id == outcome.report_id)
1103            .collect();
1104        let mut kinds: Vec<&str> = second.iter().map(|row| row.kind.as_str()).collect();
1105        kinds.sort_unstable();
1106        assert_eq!(
1107            kinds,
1108            ["drift_detected", "policy_gap_closed", "version_changed"]
1109        );
1110        assert!(
1111            second
1112                .iter()
1113                .all(|row| row.commit.as_deref() == Some("bbb"))
1114        );
1115        assert!(
1116            second
1117                .iter()
1118                .all(|row| row.project_id == outcome.project_id)
1119        );
1120
1121        let (report, raw) = carrying("ccc", "1.1.0", "ok", false);
1122        store.ingest(&report, &raw).unwrap();
1123        assert_eq!(
1124            kinds_of(&store).last().map(String::as_str),
1125            Some("drift_cleared")
1126        );
1127    }
1128
1129    #[test]
1130    fn the_same_report_twice_records_nothing_the_second_time() {
1131        let store = Store::open_in_memory().unwrap();
1132        let (report, raw) = carrying("aaa", "1.0.0", "ok", true);
1133        store.ingest(&report, &raw).unwrap();
1134        let before = store.events(&EventFilter::default()).unwrap();
1135        let mut again = raw.clone();
1136        again["generatedAt"] = json!("2026-10-01T00:00:00Z");
1137        let outcome = store.ingest(&report, &again).unwrap();
1138        assert!(outcome.deduplicated);
1139        assert_eq!(store.events(&EventFilter::default()).unwrap(), before);
1140    }
1141
1142    #[test]
1143    fn the_inventory_holds_the_latest_report_as_rows() {
1144        let store = Store::open_in_memory().unwrap();
1145        let (report, raw) = carrying("aaa", "1.0.0", "ok", false);
1146        store.ingest(&report, &raw).unwrap();
1147        let (report, raw) = carrying("bbb", "1.1.0", "modified", false);
1148        store.ingest(&report, &raw).unwrap();
1149        type Row = (String, String, String, String, String);
1150        let rows: Vec<Row> = store
1151            .conn()
1152            .prepare(
1153                "SELECT capability_type, capability_id, target, version, status
1154                 FROM inventory ORDER BY capability_id",
1155            )
1156            .unwrap()
1157            .query_map([], |row| {
1158                Ok((
1159                    row.get(0)?,
1160                    row.get(1)?,
1161                    row.get(2)?,
1162                    row.get(3)?,
1163                    row.get(4)?,
1164                ))
1165            })
1166            .unwrap()
1167            .collect::<std::result::Result<_, _>>()
1168            .unwrap();
1169        assert_eq!(rows.len(), 2);
1170        assert_eq!(
1171            rows[1],
1172            (
1173                "skill".into(),
1174                "lint".into(),
1175                "claude".into(),
1176                "1.1.0".into(),
1177                "modified".into()
1178            )
1179        );
1180    }
1181
1182    #[test]
1183    fn events_filter_by_project_capability_kind_and_time() {
1184        let store = Store::open_in_memory().unwrap();
1185        let (report, raw) = carrying("aaa", "1.0.0", "ok", true);
1186        let outcome = store.ingest(&report, &raw).unwrap();
1187        let count = |filter: EventFilter| store.events(&filter).unwrap().len();
1188        assert_eq!(
1189            count(EventFilter {
1190                project_id: Some(outcome.project_id),
1191                ..Default::default()
1192            }),
1193            4
1194        );
1195        assert_eq!(
1196            count(EventFilter {
1197                project_id: Some(999),
1198                ..Default::default()
1199            }),
1200            0
1201        );
1202        assert_eq!(
1203            count(EventFilter {
1204                capability: Some("lint".into()),
1205                ..Default::default()
1206            }),
1207            1
1208        );
1209        assert_eq!(
1210            count(EventFilter {
1211                kind: Some("capability_added".into()),
1212                ..Default::default()
1213            }),
1214            2
1215        );
1216        assert_eq!(
1217            count(EventFilter {
1218                since: Some("2999-01-01".into()),
1219                ..Default::default()
1220            }),
1221            0
1222        );
1223        assert_eq!(
1224            count(EventFilter {
1225                since: Some("2000-01-01".into()),
1226                limit: Some(1),
1227                ..Default::default()
1228            }),
1229            1
1230        );
1231    }
1232
1233    #[test]
1234    fn a_scoped_key_remembers_its_repository() {
1235        let store = Store::open_in_memory().unwrap();
1236        let secret = store
1237            .create_key("web", Some("git@github.com:Acme/web.git"))
1238            .unwrap();
1239        let grant = store.verify_key(&secret).unwrap().unwrap();
1240        assert_eq!(grant.repository.as_deref(), Some("github.com/Acme/web"));
1241        assert_eq!(
1242            store.keys().unwrap()[0].repository.as_deref(),
1243            Some("github.com/Acme/web")
1244        );
1245        assert_eq!(
1246            store.create_key("bad", Some("web")).unwrap_err().kind(),
1247            ErrorKind::Usage
1248        );
1249    }
1250
1251    #[test]
1252    fn a_version_one_database_gains_the_key_scope_column() {
1253        let temp = tempfile::tempdir().unwrap();
1254        {
1255            let conn = Connection::open(temp.path().join(DATABASE_FILE)).unwrap();
1256            conn.execute_batch(MIGRATIONS[0]).unwrap();
1257            conn.pragma_update(None, "user_version", 1).unwrap();
1258            conn.execute(
1259                "INSERT INTO keys (name, sha256, created_at) VALUES ('old', 'x', 't')",
1260                [],
1261            )
1262            .unwrap();
1263        }
1264        let store = Store::open(temp.path()).unwrap();
1265        assert_eq!(store.schema_version().unwrap(), 2);
1266        let keys = store.keys().unwrap();
1267        assert_eq!(keys[0].name, "old");
1268        assert_eq!(keys[0].repository, None);
1269    }
1270}