1use crate::textdiff;
31use rusqlite::{params, Connection, OptionalExtension};
32use serde::Serialize;
33use sha2::{Digest, Sha256};
34
35pub const DEFAULT_KEEP: i64 = 50;
36pub const DEFAULT_DAYS: i64 = 90;
37pub const MAX_BYTES: usize = 2 * 1024 * 1024;
40
41pub fn init_schema(conn: &Connection) -> Result<(), String> {
42 conn.execute_batch(
43 "CREATE TABLE IF NOT EXISTS config_snapshots (
44 seq INTEGER PRIMARY KEY AUTOINCREMENT,
45 project TEXT NOT NULL,
46 project_name TEXT NOT NULL,
47 exporter TEXT NOT NULL,
48 at TEXT NOT NULL,
49 sha256 TEXT NOT NULL,
50 bytes INTEGER NOT NULL,
51 content TEXT NOT NULL,
52 masked TEXT NOT NULL,
53 exposed TEXT NOT NULL DEFAULT '[]',
54 cause TEXT NOT NULL,
55 actor TEXT,
56 chain TEXT NOT NULL
57 );
58 CREATE INDEX IF NOT EXISTS config_snapshots_stream
59 ON config_snapshots(project, exporter, seq);
60 CREATE INDEX IF NOT EXISTS config_snapshots_sha ON config_snapshots(sha256);
61 CREATE TABLE IF NOT EXISTS config_checkpoints (
62 id INTEGER PRIMARY KEY AUTOINCREMENT,
63 project TEXT NOT NULL,
64 exporter TEXT NOT NULL,
65 pruned_through_seq INTEGER NOT NULL,
66 chain_at TEXT NOT NULL,
67 pruned_count INTEGER NOT NULL,
68 at TEXT NOT NULL
69 );",
70 )
71 .map_err(|e| e.to_string())?;
72 let _ = conn.execute_batch(
74 "ALTER TABLE config_snapshots ADD COLUMN exposed TEXT NOT NULL DEFAULT '[]';",
75 );
76 Ok(())
77}
78
79fn sha(s: &str) -> String {
80 hex::encode(Sha256::digest(s.as_bytes()))
81}
82
83fn chain_of(
84 prev: &str,
85 project: &str,
86 exporter: &str,
87 content: &str,
88 masked: &str,
89 exposed: &str,
90 at: &str,
91) -> String {
92 sha(&format!(
95 "envv-history-v2|{prev}|{project}|{exporter}|{}|{}|{}|{at}",
96 sha(content),
97 sha(masked),
98 sha(exposed)
99 ))
100}
101
102#[derive(Debug, Clone, Copy, PartialEq, Serialize)]
105pub struct Policy {
106 pub enabled: bool,
107 pub keep: i64,
108 pub days: i64,
109}
110
111fn meta_get(conn: &Connection, key: &str) -> Option<String> {
112 conn.query_row("SELECT value FROM vault_meta WHERE key = ?1", [key], |r| {
113 r.get(0)
114 })
115 .optional()
116 .ok()
117 .flatten()
118}
119
120fn meta_set(conn: &Connection, key: &str, value: &str) -> Result<(), String> {
121 conn.execute(
122 "INSERT OR REPLACE INTO vault_meta (key, value) VALUES (?1, ?2)",
123 params![key, value],
124 )
125 .map(|_| ())
126 .map_err(|e| e.to_string())
127}
128
129pub fn policy(conn: &Connection) -> Policy {
130 let num = |k: &str, d: i64| {
131 meta_get(conn, k)
132 .and_then(|v| v.parse::<i64>().ok())
133 .filter(|n| *n >= 0)
134 .unwrap_or(d)
135 };
136 Policy {
137 enabled: meta_get(conn, "history_enabled").as_deref() != Some("0"),
138 keep: num("history_keep", DEFAULT_KEEP),
139 days: num("history_days", DEFAULT_DAYS),
140 }
141}
142
143pub fn set_policy(
144 conn: &Connection,
145 enabled: Option<bool>,
146 keep: Option<i64>,
147 days: Option<i64>,
148) -> Result<Policy, String> {
149 if let Some(k) = keep {
150 if !(1..=100_000).contains(&k) {
151 return Err(
152 "keep must be between 1 and 100000 (a stream always keeps its newest snapshot)"
153 .into(),
154 );
155 }
156 meta_set(conn, "history_keep", &k.to_string())?;
157 }
158 if let Some(d) = days {
159 if !(0..=36_500).contains(&d) {
160 return Err("days must be between 0 and 36500".into());
161 }
162 meta_set(conn, "history_days", &d.to_string())?;
163 }
164 if let Some(e) = enabled {
165 meta_set(conn, "history_enabled", if e { "1" } else { "0" })?;
166 }
167 Ok(policy(conn))
168}
169
170pub struct Stream<'a> {
174 pub project: &'a str,
175 pub project_name: &'a str,
176 pub exporter: &'a str,
177}
178
179#[derive(Debug, Clone, Serialize, PartialEq)]
180pub struct SnapMeta {
181 pub seq: i64,
182 pub project: String,
183 pub project_name: String,
184 pub exporter: String,
185 pub at: String,
186 pub sha256: String,
187 pub bytes: i64,
188 pub cause: String,
189 pub actor: Option<String>,
190}
191
192const META_COLS: &str = "seq, project, project_name, exporter, at, sha256, bytes, cause, actor";
193
194fn meta_row(r: &rusqlite::Row) -> rusqlite::Result<SnapMeta> {
195 Ok(SnapMeta {
196 seq: r.get(0)?,
197 project: r.get(1)?,
198 project_name: r.get(2)?,
199 exporter: r.get(3)?,
200 at: r.get(4)?,
201 sha256: r.get(5)?,
202 bytes: r.get(6)?,
203 cause: r.get(7)?,
204 actor: r.get(8)?,
205 })
206}
207
208pub fn record(
211 conn: &Connection,
212 s: &Stream,
213 content: &str,
214 masked: &str,
215 exposed: &str,
216 cause: &str,
217 actor: Option<&str>,
218) -> Result<Option<i64>, String> {
219 if !policy(conn).enabled {
220 return Ok(None);
221 }
222 if content.len() > MAX_BYTES {
223 return Err(format!(
224 "{} {} renders to {} bytes; history keeps files up to {} bytes",
225 s.project_name,
226 s.exporter,
227 content.len(),
228 MAX_BYTES
229 ));
230 }
231 let tx = rusqlite::Transaction::new_unchecked(conn, rusqlite::TransactionBehavior::Immediate)
235 .map_err(|e| e.to_string())?;
236 let last: Option<(String, String)> = tx
237 .query_row(
238 "SELECT sha256, chain FROM config_snapshots WHERE project = ?1 AND exporter = ?2 \
239 ORDER BY seq DESC LIMIT 1",
240 params![s.project, s.exporter],
241 |r| Ok((r.get(0)?, r.get(1)?)),
242 )
243 .optional()
244 .map_err(|e| e.to_string())?;
245 let new_sha = sha(content);
246 if last.as_ref().is_some_and(|(h, _)| *h == new_sha) {
247 return Ok(None);
248 }
249 let prev = last.map_or_else(|| genesis_for(&tx, s.project, s.exporter), |(_, c)| c);
252 let at = crate::iso_now();
253 let chain = chain_of(&prev, s.project, s.exporter, content, masked, exposed, &at);
254 tx.execute(
255 "INSERT INTO config_snapshots \
256 (project, project_name, exporter, at, sha256, bytes, content, masked, exposed, cause, actor, chain) \
257 VALUES (?1,?2,?3,?4,?5,?6,?7,?8,?9,?10,?11,?12)",
258 params![
259 s.project,
260 s.project_name,
261 s.exporter,
262 at,
263 new_sha,
264 content.len() as i64,
265 content,
266 masked,
267 exposed,
268 cause,
269 actor,
270 chain
271 ],
272 )
273 .map_err(|e| e.to_string())?;
274 let seq = tx.last_insert_rowid();
275 tx.commit().map_err(|e| e.to_string())?;
276 Ok(Some(seq))
277}
278
279fn genesis_for(conn: &Connection, project: &str, exporter: &str) -> String {
282 conn.query_row(
283 "SELECT chain_at FROM config_checkpoints WHERE project = ?1 AND exporter = ?2 \
284 ORDER BY pruned_through_seq DESC LIMIT 1",
285 params![project, exporter],
286 |r| r.get(0),
287 )
288 .optional()
289 .ok()
290 .flatten()
291 .unwrap_or_else(|| "genesis".to_string())
292}
293
294pub fn list(
297 conn: &Connection,
298 project: Option<&str>,
299 exporter: Option<&str>,
300 since: Option<&str>,
301 limit: i64,
302) -> Result<Vec<SnapMeta>, String> {
303 let mut stmt = conn
304 .prepare(&format!(
305 "SELECT {META_COLS} FROM config_snapshots \
306 WHERE (?1 IS NULL OR project = ?1 OR project_name = ?1) \
307 AND (?2 IS NULL OR exporter = ?2) AND (?3 IS NULL OR at >= ?3) \
308 ORDER BY seq DESC LIMIT ?4"
309 ))
310 .map_err(|e| e.to_string())?;
311 let rows = stmt
312 .query_map(
313 params![project, exporter, since, limit.clamp(1, 10_000)],
314 meta_row,
315 )
316 .map_err(|e| e.to_string())?;
317 rows.collect::<Result<Vec<_>, _>>()
318 .map_err(|e| e.to_string())
319}
320
321pub struct Snapshot {
322 pub meta: SnapMeta,
323 pub content: String,
324 pub masked: String,
325 pub exposed: String,
327}
328
329pub fn get(conn: &Connection, seq: i64) -> Result<Option<Snapshot>, String> {
330 conn.query_row(
331 &format!(
332 "SELECT {META_COLS}, content, masked, exposed FROM config_snapshots WHERE seq = ?1"
333 ),
334 [seq],
335 |r| {
336 Ok(Snapshot {
337 meta: meta_row(r)?,
338 content: r.get(9)?,
339 masked: r.get(10)?,
340 exposed: r.get(11)?,
341 })
342 },
343 )
344 .optional()
345 .map_err(|e| e.to_string())
346}
347
348pub fn find_by_sha(conn: &Connection, sha256: &str, limit: i64) -> Result<Vec<SnapMeta>, String> {
351 let mut stmt = conn
352 .prepare(&format!(
353 "SELECT {META_COLS} FROM config_snapshots WHERE sha256 = ?1 ORDER BY seq DESC LIMIT ?2"
354 ))
355 .map_err(|e| e.to_string())?;
356 let rows = stmt
357 .query_map(params![sha256, limit.clamp(1, 100)], meta_row)
358 .map_err(|e| e.to_string())?;
359 rows.collect::<Result<Vec<_>, _>>()
360 .map_err(|e| e.to_string())
361}
362
363pub fn exposed_by_sha(
367 conn: &Connection,
368 sha256: &str,
369) -> Result<Option<Vec<crate::blast::Exposed>>, String> {
370 let mut stmt = conn
371 .prepare("SELECT exposed FROM config_snapshots WHERE sha256 = ?1")
372 .map_err(|e| e.to_string())?;
373 let rows: Vec<String> = stmt
374 .query_map([sha256], |r| r.get(0))
375 .map_err(|e| e.to_string())?
376 .collect::<Result<_, _>>()
377 .map_err(|e| e.to_string())?;
378 if rows.is_empty() {
379 return Ok(None);
380 }
381 let mut all: Vec<crate::blast::Exposed> = Vec::new();
382 for json in rows {
383 let list: Vec<crate::blast::Exposed> =
384 serde_json::from_str(&json).map_err(|e| e.to_string())?;
385 for e in list {
386 if !all.contains(&e) {
387 all.push(e);
388 }
389 }
390 }
391 Ok(Some(all))
392}
393
394pub fn latest(
396 conn: &Connection,
397 project: &str,
398 exporter: &str,
399) -> Result<Option<Snapshot>, String> {
400 let seq: Option<i64> = conn
401 .query_row(
402 "SELECT seq FROM config_snapshots WHERE project = ?1 AND exporter = ?2 \
403 ORDER BY seq DESC LIMIT 1",
404 params![project, exporter],
405 |r| r.get(0),
406 )
407 .optional()
408 .map_err(|e| e.to_string())?;
409 match seq {
410 Some(s) => get(conn, s),
411 None => Ok(None),
412 }
413}
414
415pub fn diff(
418 conn: &Connection,
419 a: i64,
420 b: i64,
421 masked: bool,
422 context: usize,
423) -> Result<String, String> {
424 let sa = get(conn, a)?.ok_or_else(|| format!("No snapshot #{a}"))?;
425 let sb = get(conn, b)?.ok_or_else(|| format!("No snapshot #{b}"))?;
426 if (&sa.meta.project, &sa.meta.exporter) != (&sb.meta.project, &sb.meta.exporter) {
427 return Err(
428 "Those snapshots belong to different configs; compare two of one stream".into(),
429 );
430 }
431 let pick = |s: &Snapshot| {
432 if masked {
433 s.masked.clone()
434 } else {
435 s.content.clone()
436 }
437 };
438 Ok(textdiff::unified(
439 &format!("#{a} {}", sa.meta.at),
440 &format!("#{b} {}", sb.meta.at),
441 &pick(&sa),
442 &pick(&sb),
443 context,
444 ))
445}
446
447pub fn resolve_pair(
451 conn: &Connection,
452 project: Option<&str>,
453 exporter: Option<&str>,
454 from: Option<i64>,
455 to: Option<i64>,
456) -> Result<(i64, i64), String> {
457 if let (Some(a), Some(b)) = (from, to) {
458 return Ok((a, b));
459 }
460 let project = project.ok_or("Name a project, or give both --from and --to")?;
461 let rows = list(conn, Some(project), exporter, None, 1000)?;
462 let mut exporters: Vec<&str> = rows.iter().map(|m| m.exporter.as_str()).collect();
463 exporters.sort_unstable();
464 exporters.dedup();
465 if exporters.len() > 1 {
466 return Err(format!(
467 "'{project}' has several configs ({}); say which with --exporter",
468 exporters.join(", ")
469 ));
470 }
471 let newest = rows
473 .first()
474 .ok_or_else(|| format!("No snapshots of '{project}' yet"))?;
475 let b = to.unwrap_or(newest.seq);
476 let a = match from {
477 Some(a) => a,
478 None => rows
479 .iter()
480 .find(|m| m.seq < b)
481 .map(|m| m.seq)
482 .ok_or_else(|| {
483 format!("'{project}' has only one snapshot, so there is nothing to compare")
484 })?,
485 };
486 Ok((a, b))
487}
488
489#[derive(Debug, Serialize, PartialEq)]
492pub struct PruneReport {
493 pub would_delete: i64,
494 pub deleted: i64,
495 pub streams_touched: i64,
496 pub bytes_freed: i64,
497}
498
499pub fn prune(
503 conn: &Connection,
504 keep: i64,
505 days: i64,
506 dry_run: bool,
507 actor: Option<&str>,
508) -> Result<PruneReport, String> {
509 let cutoff = cutoff_iso(days)?;
510 let streams: Vec<(String, String)> = {
511 let mut st = conn
512 .prepare("SELECT DISTINCT project, exporter FROM config_snapshots")
513 .map_err(|e| e.to_string())?;
514 let rows = st
515 .query_map([], |r| Ok((r.get(0)?, r.get(1)?)))
516 .map_err(|e| e.to_string())?;
517 rows.collect::<Result<_, _>>().map_err(|e| e.to_string())?
518 };
519 let mut rep = PruneReport {
520 would_delete: 0,
521 deleted: 0,
522 streams_touched: 0,
523 bytes_freed: 0,
524 };
525 let tx = rusqlite::Transaction::new_unchecked(conn, rusqlite::TransactionBehavior::Immediate)
529 .map_err(|e| e.to_string())?;
530 for (project, exporter) in streams {
531 let doomed: Vec<(i64, String, i64)> = {
534 let mut st = tx
535 .prepare(
536 "SELECT seq, chain, bytes FROM config_snapshots \
537 WHERE project = ?1 AND exporter = ?2 AND at < ?3 \
538 AND seq NOT IN (SELECT seq FROM config_snapshots \
539 WHERE project = ?1 AND exporter = ?2 \
540 ORDER BY seq DESC LIMIT ?4) \
541 ORDER BY seq ASC",
542 )
543 .map_err(|e| e.to_string())?;
544 let rows = st
545 .query_map(params![project, exporter, cutoff, keep.max(1)], |r| {
546 Ok((r.get(0)?, r.get(1)?, r.get(2)?))
547 })
548 .map_err(|e| e.to_string())?;
549 rows.collect::<Result<_, _>>().map_err(|e| e.to_string())?
550 };
551 let Some((last_seq, last_chain, _)) = doomed.last().cloned() else {
554 continue;
555 };
556 rep.would_delete += doomed.len() as i64;
557 rep.bytes_freed += doomed.iter().map(|d| d.2).sum::<i64>();
558 rep.streams_touched += 1;
559 if dry_run {
560 continue;
561 }
562 tx.execute(
563 "DELETE FROM config_snapshots WHERE project = ?1 AND exporter = ?2 AND seq <= ?3",
564 params![project, exporter, last_seq],
565 )
566 .map_err(|e| e.to_string())?;
567 let at = crate::iso_now();
568 tx.execute(
569 "INSERT INTO config_checkpoints \
570 (project, exporter, pruned_through_seq, chain_at, pruned_count, at) \
571 VALUES (?1,?2,?3,?4,?5,?6)",
572 params![
573 project,
574 exporter,
575 last_seq,
576 last_chain,
577 doomed.len() as i64,
578 at
579 ],
580 )
581 .map_err(|e| e.to_string())?;
582 rep.deleted += doomed.len() as i64;
583 crate::record_event(
586 &tx,
587 "config.prune",
588 &project,
589 Some(
590 &serde_json::json!({
591 "exporter": exporter, "through_seq": last_seq,
592 "chain_at": last_chain, "count": doomed.len()
593 })
594 .to_string(),
595 ),
596 actor,
597 )?;
598 }
599 tx.commit().map_err(|e| e.to_string())?;
600 Ok(rep)
601}
602
603fn cutoff_iso(days: i64) -> Result<String, String> {
604 let now = time::OffsetDateTime::now_utc();
605 let then = now
606 .checked_sub(time::Duration::days(days))
607 .ok_or("days is out of range")?;
608 Ok(format!(
609 "{:04}-{:02}-{:02}T{:02}:{:02}:{:02}Z",
610 then.year(),
611 u8::from(then.month()),
612 then.day(),
613 then.hour(),
614 then.minute(),
615 then.second()
616 ))
617}
618
619#[derive(Debug, Serialize, PartialEq)]
622pub struct VerifyReport {
623 pub streams: i64,
624 pub snapshots: i64,
625 pub checkpoints: i64,
626 pub problems: Vec<String>,
627}
628
629pub fn verify(conn: &Connection) -> Result<VerifyReport, String> {
632 let mut rep = VerifyReport {
633 streams: 0,
634 snapshots: 0,
635 checkpoints: 0,
636 problems: vec![],
637 };
638 let streams: Vec<(String, String)> = {
639 let mut st = conn
640 .prepare(
641 "SELECT project, exporter FROM config_snapshots \
642 UNION SELECT project, exporter FROM config_checkpoints",
643 )
644 .map_err(|e| e.to_string())?;
645 let rows = st
646 .query_map([], |r| Ok((r.get(0)?, r.get(1)?)))
647 .map_err(|e| e.to_string())?;
648 rows.collect::<Result<_, _>>().map_err(|e| e.to_string())?
649 };
650 let audit: Vec<String> = {
651 let mut st = conn
652 .prepare("SELECT COALESCE(details,'') FROM vault_audit WHERE action = 'config.prune'")
653 .map_err(|e| e.to_string())?;
654 let rows = st.query_map([], |r| r.get(0)).map_err(|e| e.to_string())?;
655 rows.collect::<Result<_, _>>().map_err(|e| e.to_string())?
656 };
657 for (project, exporter) in streams {
658 rep.streams += 1;
659 let label = format!("{project}/{exporter}");
660 let cps: Vec<(i64, String)> = {
661 let mut st = conn
662 .prepare(
663 "SELECT pruned_through_seq, chain_at FROM config_checkpoints \
664 WHERE project = ?1 AND exporter = ?2 ORDER BY pruned_through_seq ASC",
665 )
666 .map_err(|e| e.to_string())?;
667 let rows = st
668 .query_map(params![project, exporter], |r| Ok((r.get(0)?, r.get(1)?)))
669 .map_err(|e| e.to_string())?;
670 rows.collect::<Result<_, _>>().map_err(|e| e.to_string())?
671 };
672 for (through, chain_at) in &cps {
673 rep.checkpoints += 1;
674 if !audit.iter().any(|d| d.contains(chain_at.as_str())) {
675 rep.problems.push(format!(
676 "{label}: checkpoint through #{through} has no matching row in the audit chain"
677 ));
678 }
679 }
680 let mut prev = cps
681 .last()
682 .map_or_else(|| "genesis".to_string(), |c| c.1.clone());
683 let mut st = conn
684 .prepare(
685 "SELECT seq, at, content, masked, sha256, bytes, chain, exposed FROM config_snapshots \
686 WHERE project = ?1 AND exporter = ?2 ORDER BY seq ASC",
687 )
688 .map_err(|e| e.to_string())?;
689 let rows = st
690 .query_map(params![project, exporter], |r| {
691 Ok((
692 r.get::<_, i64>(0)?,
693 r.get::<_, String>(1)?,
694 r.get::<_, String>(2)?,
695 r.get::<_, String>(3)?,
696 r.get::<_, String>(4)?,
697 r.get::<_, i64>(5)?,
698 r.get::<_, String>(6)?,
699 r.get::<_, String>(7)?,
700 ))
701 })
702 .map_err(|e| e.to_string())?;
703 for row in rows {
704 let (seq, at, content, masked, stored_sha, bytes, chain, exposed) =
705 row.map_err(|e| e.to_string())?;
706 rep.snapshots += 1;
707 if sha(&content) != stored_sha || content.len() as i64 != bytes {
708 rep.problems
709 .push(format!("{label}: #{seq} content does not match its hash"));
710 }
711 let want = chain_of(&prev, &project, &exporter, &content, &masked, &exposed, &at);
712 if want != chain {
713 rep.problems
714 .push(format!("{label}: #{seq} breaks the chain"));
715 }
716 prev = chain;
717 }
718 }
719 Ok(rep)
720}
721
722#[derive(Debug, Serialize)]
723pub struct Stats {
724 pub snapshots: i64,
725 pub streams: i64,
726 pub bytes: i64,
727 pub oldest: Option<String>,
728}
729
730pub fn stats(conn: &Connection) -> Result<Stats, String> {
731 conn.query_row(
732 "SELECT COUNT(*), COUNT(DISTINCT project || '/' || exporter), COALESCE(SUM(bytes),0), MIN(at) \
733 FROM config_snapshots",
734 [],
735 |r| {
736 Ok(Stats {
737 snapshots: r.get(0)?,
738 streams: r.get(1)?,
739 bytes: r.get(2)?,
740 oldest: r.get(3)?,
741 })
742 },
743 )
744 .map_err(|e| e.to_string())
745}
746
747#[cfg(test)]
748mod tests {
749 use super::*;
750
751 fn db() -> Connection {
752 let c = Connection::open_in_memory().unwrap();
753 crate::init_schema(&c).unwrap();
754 c
755 }
756
757 const S: Stream = Stream {
758 project: "p1",
759 project_name: "edge",
760 exporter: "nginx",
761 };
762
763 fn rec(c: &Connection, body: &str) -> Option<i64> {
764 record(c, &S, body, &masked_of(body), "[]", "save", Some("owner")).unwrap()
765 }
766
767 fn masked_of(body: &str) -> String {
770 format!("m:{}", &sha(body)[..8])
771 }
772
773 fn age(c: &Connection, seq: i64, at: &str) {
776 c.execute(
777 "UPDATE config_snapshots SET at = ?1 WHERE seq = ?2",
778 params![at, seq],
779 )
780 .unwrap();
781 rechain(c);
782 }
783
784 fn rechain(c: &Connection) {
785 let rows: Vec<(i64, String, String, String, String)> = {
786 let mut st = c
787 .prepare(
788 "SELECT seq, at, content, masked, exposed FROM config_snapshots ORDER BY seq",
789 )
790 .unwrap();
791 st.query_map([], |r| {
792 Ok((r.get(0)?, r.get(1)?, r.get(2)?, r.get(3)?, r.get(4)?))
793 })
794 .unwrap()
795 .map(Result::unwrap)
796 .collect()
797 };
798 let mut prev = genesis_for(c, "p1", "nginx");
799 for (seq, at, content, masked, exposed) in rows {
800 prev = chain_of(&prev, "p1", "nginx", &content, &masked, &exposed, &at);
801 c.execute(
802 "UPDATE config_snapshots SET chain = ?1 WHERE seq = ?2",
803 params![prev, seq],
804 )
805 .unwrap();
806 }
807 }
808
809 #[test]
810 fn concurrent_writers_never_fork_the_chain() {
811 let n = std::time::SystemTime::now()
812 .duration_since(std::time::UNIX_EPOCH)
813 .unwrap()
814 .as_nanos();
815 let path = std::env::temp_dir().join(format!("unv-hist-conc-{n}.db"));
816 {
817 let c = Connection::open(&path).unwrap();
818 crate::init_schema(&c).unwrap();
819 }
820 let handles: Vec<_> = (0..4)
821 .map(|t| {
822 let path = path.clone();
823 std::thread::spawn(move || {
824 let c = Connection::open(&path).unwrap();
825 for i in 0..25 {
826 let body = format!("thread {t} version {i}");
827 record(&c, &S, &body, &masked_of(&body), "[]", "save", None).unwrap();
828 }
829 })
830 })
831 .collect();
832 for h in handles {
833 h.join().unwrap();
834 }
835 let c = Connection::open(&path).unwrap();
836 let v = verify(&c).unwrap();
837 assert_eq!(v.snapshots, 100);
838 assert!(v.problems.is_empty(), "{:?}", v.problems);
839 }
840
841 #[test]
842 fn a_change_is_recorded_and_an_identical_render_is_not() {
843 let c = db();
844 assert!(rec(&c, "a").is_some());
845 assert!(rec(&c, "a").is_none());
846 assert!(rec(&c, "b").is_some());
847 assert_eq!(list(&c, Some("edge"), None, None, 10).unwrap().len(), 2);
848 }
849
850 #[test]
851 fn going_back_to_an_earlier_text_is_a_new_snapshot() {
852 let c = db();
853 rec(&c, "a");
854 rec(&c, "b");
855 assert!(rec(&c, "a").is_some(), "reverting a change is a change");
856 }
857
858 #[test]
859 fn a_stream_is_per_project_and_exporter() {
860 let c = db();
861 rec(&c, "a");
862 let other = Stream {
863 project: "p1",
864 project_name: "edge",
865 exporter: "compose",
866 };
867 assert!(record(&c, &other, "a", "m", "[]", "save", None)
868 .unwrap()
869 .is_some());
870 }
871
872 #[test]
873 fn disabled_history_records_nothing() {
874 let c = db();
875 set_policy(&c, Some(false), None, None).unwrap();
876 assert!(rec(&c, "a").is_none());
877 set_policy(&c, Some(true), None, None).unwrap();
878 assert!(rec(&c, "a").is_some());
879 }
880
881 #[test]
882 fn default_views_read_the_masked_text_and_only_get_returns_the_real_one() {
883 let c = db();
884 let s1 = rec(&c, "key=REAL1").unwrap();
885 let s2 = rec(&c, "key=REAL2").unwrap();
886 let d = diff(&c, s1, s2, true, 3).unwrap();
887 assert!(d.contains(&format!("-{}", masked_of("key=REAL1"))));
888 assert!(d.contains(&format!("+{}", masked_of("key=REAL2"))));
889 assert!(
890 !d.contains("REAL"),
891 "a masked diff carried the real value: {d}"
892 );
893 let real = diff(&c, s1, s2, false, 3).unwrap();
894 assert!(real.contains("-key=REAL1") && real.contains("+key=REAL2"));
895 }
896
897 #[test]
898 fn the_default_diff_compares_the_two_newest_and_asks_when_a_project_has_two_configs() {
899 let c = db();
900 let a = rec(&c, "a").unwrap();
901 let b = rec(&c, "b").unwrap();
902 assert_eq!(
903 resolve_pair(&c, Some("edge"), None, None, None).unwrap(),
904 (a, b)
905 );
906 assert_eq!(
907 resolve_pair(&c, Some("edge"), None, Some(a), None).unwrap(),
908 (a, b)
909 );
910 assert_eq!(
911 resolve_pair(&c, None, None, Some(1), Some(2)).unwrap(),
912 (1, 2)
913 );
914 assert!(resolve_pair(&c, None, None, None, None).is_err());
915 let other = Stream {
916 project: "p1",
917 project_name: "edge",
918 exporter: "compose",
919 };
920 record(&c, &other, "x", "m", "[]", "save", None).unwrap();
921 let e = resolve_pair(&c, Some("edge"), None, None, None).unwrap_err();
922 assert!(e.contains("--exporter"), "{e}");
923 assert!(resolve_pair(&c, Some("edge"), Some("nginx"), None, None).is_ok());
924 }
925
926 #[test]
927 fn one_snapshot_has_nothing_to_compare_with() {
928 let c = db();
929 rec(&c, "a");
930 assert!(resolve_pair(&c, Some("edge"), None, None, None)
931 .unwrap_err()
932 .contains("only one"));
933 }
934
935 #[test]
936 fn a_diff_of_two_different_streams_is_refused() {
937 let c = db();
938 let a = rec(&c, "a").unwrap();
939 let other = Stream {
940 project: "p2",
941 project_name: "x",
942 exporter: "nginx",
943 };
944 let b = record(&c, &other, "b", "m", "[]", "save", None)
945 .unwrap()
946 .unwrap();
947 assert!(diff(&c, a, b, true, 3)
948 .unwrap_err()
949 .contains("different configs"));
950 }
951
952 fn ex(id: &str, fp: &str) -> String {
953 serde_json::to_string(&vec![crate::blast::Exposed {
954 entry_id: id.into(),
955 provider: "P".into(),
956 key_id: String::new(),
957 field: "api_key".into(),
958 fp: fp.into(),
959 short: false,
960 }])
961 .unwrap()
962 }
963
964 #[test]
965 fn what_a_file_contained_is_the_union_over_every_snapshot_with_that_hash_and_none_when_unknown()
966 {
967 let c = db();
968 record(&c, &S, "same", "m1", &ex("e1", "f1"), "save", None).unwrap();
969 record(&c, &S, "other", "m2", "[]", "save", None).unwrap();
970 record(&c, &S, "same", "m1", &ex("e2", "f2"), "save", None).unwrap();
972 let got = exposed_by_sha(&c, &sha("same")).unwrap().unwrap();
973 let ids: Vec<&str> = got.iter().map(|e| e.entry_id.as_str()).collect();
974 assert_eq!(ids, vec!["e1", "e2"]);
975 assert_eq!(exposed_by_sha(&c, &sha("other")).unwrap().unwrap().len(), 0);
976 assert!(exposed_by_sha(&c, &sha("never recorded"))
977 .unwrap()
978 .is_none());
979 }
980
981 #[test]
982 fn hiding_an_entry_from_the_exposed_list_breaks_the_chain() {
983 let c = db();
984 record(&c, &S, "t", "m", &ex("e1", "f1"), "save", None).unwrap();
985 assert!(verify(&c).unwrap().problems.is_empty());
986 c.execute(
987 "UPDATE config_snapshots SET exposed = '[]' WHERE seq = 1",
988 [],
989 )
990 .unwrap();
991 assert!(
992 verify(&c)
993 .unwrap()
994 .problems
995 .iter()
996 .any(|p| p.contains("#1 breaks the chain")),
997 "an edited exposed list went unnoticed"
998 );
999 }
1000
1001 #[test]
1002 fn the_file_on_a_host_is_matched_to_the_snapshot_it_came_from() {
1003 let c = db();
1004 let first = rec(&c, "one").unwrap();
1005 rec(&c, "two");
1006 let hits = find_by_sha(&c, &sha("one"), 5).unwrap();
1007 assert_eq!(hits.len(), 1);
1008 assert_eq!(hits[0].seq, first);
1009 assert!(find_by_sha(&c, &sha("never rendered"), 5)
1010 .unwrap()
1011 .is_empty());
1012 }
1013
1014 #[test]
1015 fn verify_passes_on_an_untouched_history_and_catches_each_kind_of_tampering() {
1016 let c = db();
1017 rec(&c, "a");
1018 rec(&c, "b");
1019 rec(&c, "c");
1020 let ok = verify(&c).unwrap();
1021 assert_eq!((ok.snapshots, ok.problems.len()), (3, 0));
1022
1023 c.execute(
1025 "UPDATE config_snapshots SET content = 'evil' WHERE seq = 2",
1026 [],
1027 )
1028 .unwrap();
1029 assert!(verify(&c)
1030 .unwrap()
1031 .problems
1032 .iter()
1033 .any(|p| p.contains("#2 content")));
1034 c.execute(
1035 "UPDATE config_snapshots SET content = 'b' WHERE seq = 2",
1036 [],
1037 )
1038 .unwrap();
1039 assert!(verify(&c).unwrap().problems.is_empty());
1040
1041 c.execute("UPDATE config_snapshots SET masked = 'x' WHERE seq = 2", [])
1043 .unwrap();
1044 assert!(verify(&c)
1045 .unwrap()
1046 .problems
1047 .iter()
1048 .any(|p| p.contains("#2 breaks the chain")));
1049 c.execute(
1050 "UPDATE config_snapshots SET masked = ?1 WHERE seq = 2",
1051 [masked_of("b")],
1052 )
1053 .unwrap();
1054
1055 c.execute("DELETE FROM config_snapshots WHERE seq = 2", [])
1057 .unwrap();
1058 assert!(!verify(&c).unwrap().problems.is_empty());
1059 }
1060
1061 #[test]
1062 fn deleting_the_oldest_snapshots_without_a_checkpoint_is_detected() {
1063 let c = db();
1064 rec(&c, "a");
1065 rec(&c, "b");
1066 c.execute("DELETE FROM config_snapshots WHERE seq = 1", [])
1067 .unwrap();
1068 assert!(verify(&c)
1069 .unwrap()
1070 .problems
1071 .iter()
1072 .any(|p| p.contains("breaks the chain")));
1073 }
1074
1075 #[test]
1076 fn prune_keeps_the_newest_and_the_recent_and_leaves_a_verifiable_chain() {
1077 let c = db();
1078 for i in 0..6 {
1079 rec(&c, &format!("v{i}"));
1080 }
1081 for seq in 1..=4 {
1083 age(&c, seq, "2020-01-01T00:00:00Z");
1084 }
1085 let dry = prune(&c, 3, 30, true, None).unwrap();
1086 assert_eq!(
1087 (dry.would_delete, dry.deleted),
1088 (3, 0),
1089 "keep 3 newest: only #1..#3 are old and beyond"
1090 );
1091 assert_eq!(
1092 stats(&c).unwrap().snapshots,
1093 6,
1094 "a dry run deleted something"
1095 );
1096 let real = prune(&c, 3, 30, false, Some("owner")).unwrap();
1097 assert_eq!((real.deleted, real.streams_touched), (3, 1));
1098 let left: Vec<i64> = list(&c, None, None, None, 10)
1099 .unwrap()
1100 .iter()
1101 .map(|m| m.seq)
1102 .collect();
1103 assert_eq!(left, vec![6, 5, 4]);
1104 let v = verify(&c).unwrap();
1105 assert!(v.problems.is_empty(), "{:?}", v.problems);
1106 assert_eq!(v.checkpoints, 1);
1107
1108 rec(&c, "v6");
1110 assert!(verify(&c).unwrap().problems.is_empty());
1111 for seq in 4..=7 {
1113 age(&c, seq, "2020-02-01T00:00:00Z");
1114 }
1115 prune(&c, 1, 30, false, None).unwrap();
1116 let v = verify(&c).unwrap();
1117 assert!(v.problems.is_empty(), "{:?}", v.problems);
1118 assert_eq!(list(&c, None, None, None, 10).unwrap().len(), 1);
1119 }
1120
1121 #[test]
1122 fn a_stream_that_lost_every_row_after_a_prune_restarts_from_its_checkpoint() {
1123 let c = db();
1127 rec(&c, "a");
1128 rec(&c, "b");
1129 age(&c, 1, "2000-01-01T00:00:00Z");
1130 prune(&c, 1, 30, false, Some("owner")).unwrap();
1131 c.execute("DELETE FROM config_snapshots", []).unwrap();
1132 rec(&c, "c");
1133 assert!(verify(&c).unwrap().problems.is_empty());
1134 }
1135
1136 #[test]
1137 fn prune_never_empties_a_stream() {
1138 let c = db();
1139 rec(&c, "a");
1140 age(&c, 1, "2000-01-01T00:00:00Z");
1141 let r = prune(&c, 1, 0, false, None).unwrap();
1142 assert_eq!(r.deleted, 0);
1143 assert_eq!(stats(&c).unwrap().snapshots, 1);
1144 }
1145
1146 #[test]
1147 fn a_forged_checkpoint_with_no_audit_row_is_caught() {
1148 let c = db();
1149 rec(&c, "a");
1150 rec(&c, "b");
1151 rec(&c, "c");
1152 let chain1: String = c
1154 .query_row(
1155 "SELECT chain FROM config_snapshots WHERE seq = 1",
1156 [],
1157 |r| r.get(0),
1158 )
1159 .unwrap();
1160 c.execute("DELETE FROM config_snapshots WHERE seq = 1", [])
1161 .unwrap();
1162 c.execute(
1163 "INSERT INTO config_checkpoints (project, exporter, pruned_through_seq, chain_at, pruned_count, at) \
1164 VALUES ('p1','nginx',1,?1,1,'t')",
1165 [chain1],
1166 )
1167 .unwrap();
1168 let v = verify(&c).unwrap();
1169 assert!(
1170 v.problems
1171 .iter()
1172 .any(|p| p.contains("no matching row in the audit chain")),
1173 "{:?}",
1174 v.problems
1175 );
1176 }
1177
1178 #[test]
1179 fn a_real_prune_writes_the_audit_row_the_checkpoint_is_checked_against() {
1180 let c = db();
1181 rec(&c, "a");
1182 rec(&c, "b");
1183 age(&c, 1, "2000-01-01T00:00:00Z");
1184 prune(&c, 1, 30, false, Some("owner")).unwrap();
1185 let n: i64 = c
1186 .query_row(
1187 "SELECT COUNT(*) FROM vault_audit WHERE action = 'config.prune'",
1188 [],
1189 |r| r.get(0),
1190 )
1191 .unwrap();
1192 assert_eq!(n, 1);
1193 }
1194
1195 #[test]
1196 fn policy_has_defaults_and_refuses_nonsense() {
1197 let c = db();
1198 assert_eq!(
1199 policy(&c),
1200 Policy {
1201 enabled: true,
1202 keep: DEFAULT_KEEP,
1203 days: DEFAULT_DAYS
1204 }
1205 );
1206 assert!(set_policy(&c, None, Some(0), None).is_err());
1207 assert!(set_policy(&c, None, None, Some(-1)).is_err());
1208 let p = set_policy(&c, None, Some(5), Some(7)).unwrap();
1209 assert_eq!((p.keep, p.days), (5, 7));
1210 }
1211
1212 #[test]
1213 fn an_oversized_render_is_refused_not_truncated() {
1214 let c = db();
1215 let big = "x".repeat(MAX_BYTES + 1);
1216 assert!(record(&c, &S, &big, "m", "[]", "save", None)
1217 .unwrap_err()
1218 .contains("history keeps files up to"));
1219 }
1220}