1use std::path::{Path, PathBuf};
19use std::sync::{Condvar, Mutex, MutexGuard};
20use std::time::{Duration, Instant};
21
22use rusqlite::{Connection, OpenFlags, OptionalExtension, params};
23use serde::{Deserialize, Serialize};
24use serde_json::{Value, json};
25
26use crate::error::{Error, Result};
27
28pub const DEFAULT_RETENTION: Duration = Duration::from_secs(90 * 86400);
30
31pub(crate) const GENESIS: &str = "0000000000000000000000000000000000000000000000000000000000000000";
33
34const PRUNE_EVERY: Duration = Duration::from_secs(3600);
36
37const MIGRATIONS: &[&str] = &[
38 "
39 CREATE TABLE audit (
40 id INTEGER PRIMARY KEY,
41 time INTEGER NOT NULL,
42 org TEXT,
43 actor TEXT NOT NULL,
44 actor_kind TEXT NOT NULL,
45 user_id INTEGER,
46 user_email TEXT,
47 token_id INTEGER,
48 token_name TEXT,
49 surface TEXT NOT NULL,
50 action TEXT NOT NULL,
51 target TEXT,
52 details TEXT NOT NULL DEFAULT '{}',
53 outcome TEXT NOT NULL,
54 ip TEXT,
55 user_agent TEXT,
56 request_id TEXT,
57 prev_hash TEXT NOT NULL,
58 hash TEXT NOT NULL
59 );
60 CREATE INDEX audit_org ON audit(org, id);
61 CREATE INDEX audit_time ON audit(time);
62 CREATE TABLE audit_meta (k TEXT PRIMARY KEY, v TEXT NOT NULL);
63 CREATE TRIGGER audit_no_update BEFORE UPDATE ON audit
64 BEGIN SELECT RAISE(ABORT, 'the audit log is append-only'); END;
65 CREATE TRIGGER audit_no_delete BEFORE DELETE ON audit
66 WHEN (SELECT v FROM audit_meta WHERE k = 'pruning') IS NOT '1'
67 BEGIN SELECT RAISE(ABORT, 'the audit log is append-only; rows leave only by retention'); END;
68 ",
69 "
72 CREATE TABLE history (
73 id INTEGER PRIMARY KEY,
74 time INTEGER NOT NULL,
75 source TEXT NOT NULL,
76 org TEXT,
77 project TEXT,
78 kind TEXT NOT NULL,
79 object_type TEXT,
80 object TEXT,
81 objects TEXT NOT NULL DEFAULT '',
82 actor TEXT,
83 level TEXT,
84 message TEXT,
85 details TEXT NOT NULL DEFAULT '{}',
86 prev_hash TEXT NOT NULL,
87 hash TEXT NOT NULL
88 );
89 CREATE INDEX history_org ON history(org, id);
90 CREATE INDEX history_time ON history(time);
91 CREATE TRIGGER history_no_update BEFORE UPDATE ON history
92 BEGIN SELECT RAISE(ABORT, 'the history is append-only'); END;
93 CREATE TRIGGER history_no_delete BEFORE DELETE ON history
94 WHEN (SELECT v FROM audit_meta WHERE k = 'pruning') IS NOT '1'
95 BEGIN SELECT RAISE(ABORT, 'the history is append-only; rows leave only by retention'); END;
96 ",
97];
98
99#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
101#[serde(rename_all = "snake_case")]
102pub enum ActorKind {
103 Person,
106 Agent,
108 Local,
110 Webhook,
112 Anonymous,
115 Superadmin,
119}
120
121impl ActorKind {
122 pub fn as_str(self) -> &'static str {
123 match self {
124 ActorKind::Person => "person",
125 ActorKind::Agent => "agent",
126 ActorKind::Local => "local",
127 ActorKind::Webhook => "webhook",
128 ActorKind::Anonymous => "anonymous",
129 ActorKind::Superadmin => "superadmin",
130 }
131 }
132}
133
134#[derive(Debug, Clone, PartialEq, Eq, Default)]
136pub struct Actor {
137 pub name: String,
139 pub kind: Option<ActorKind>,
140 pub user_id: Option<i64>,
141 pub email: Option<String>,
142 pub token_id: Option<i64>,
143 pub token_name: Option<String>,
144}
145
146impl Actor {
147 pub fn from_principal(p: &crate::auth::Principal) -> Actor {
148 let mut a = Actor {
149 name: p.user.email.clone(),
150 kind: Some(ActorKind::Person),
151 user_id: Some(p.user.id),
152 email: Some(p.user.email.clone()),
153 ..Default::default()
154 };
155 match &p.kind {
156 crate::auth::PrincipalKind::ApiToken { id, name, .. } => {
157 a.kind = Some(ActorKind::Agent);
158 a.token_id = Some(*id);
159 a.token_name = Some(name.clone());
160 }
161 crate::auth::PrincipalKind::Workspace { name, .. } => a.as_workspace(name),
162 crate::auth::PrincipalKind::Agent { label } => a.as_agent(label),
163 crate::auth::PrincipalKind::Superadmin { source } => {
164 a.name = source.label();
165 a.kind = Some(ActorKind::Superadmin);
166 if p.user.id <= 0 {
167 a.user_id = None;
168 a.email = None;
169 }
170 if let crate::auth::SuperadminSource::Token { id, name } = source {
171 a.token_id = Some(*id);
172 a.token_name = Some(name.clone());
173 }
174 }
175 _ => {}
176 }
177 a
178 }
179
180 pub fn local(uid: Option<u32>) -> Actor {
181 Actor {
182 name: match uid {
183 Some(u) => format!("local(uid {u})"),
184 None => "local".into(),
185 },
186 kind: Some(ActorKind::Local),
187 ..Default::default()
188 }
189 }
190
191 pub fn cli() -> Actor {
193 Actor::local(Some(rustix::process::getuid().as_raw()))
194 }
195
196 pub fn webhook(provider: &str) -> Actor {
197 Actor {
198 name: format!("webhook:{provider}"),
199 kind: Some(ActorKind::Webhook),
200 ..Default::default()
201 }
202 }
203
204 pub fn anonymous(name: impl Into<String>) -> Actor {
205 Actor {
206 name: name.into(),
207 kind: Some(ActorKind::Anonymous),
208 ..Default::default()
209 }
210 }
211
212 pub fn claimed(email: &str) -> Actor {
214 let e = clip(email.trim().to_lowercase(), 254);
215 Actor {
216 name: e.clone(),
217 kind: Some(ActorKind::Anonymous),
218 email: Some(e),
219 ..Default::default()
220 }
221 }
222}
223
224#[derive(Debug, Clone, Default, PartialEq, Eq)]
226pub struct Origin {
227 pub surface: String,
229 pub ip: Option<String>,
230 pub user_agent: Option<String>,
231 pub request_id: Option<String>,
232}
233
234#[derive(Debug, Clone, Default)]
237pub struct NewEntry {
238 pub org: Option<String>,
240 pub actor: Actor,
241 pub origin: Origin,
242 pub action: String,
244 pub target: Option<String>,
245 pub details: Value,
247 pub outcome: String,
249}
250
251#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
253pub struct Entry {
254 pub id: i64,
255 pub time: i64,
257 pub org: Option<String>,
258 pub actor: String,
259 pub actor_kind: String,
260 pub user_id: Option<i64>,
261 pub user_email: Option<String>,
262 pub token_id: Option<i64>,
263 pub token_name: Option<String>,
264 pub surface: String,
265 pub action: String,
266 pub target: Option<String>,
267 pub details: Value,
268 pub outcome: String,
269 pub ip: Option<String>,
270 pub user_agent: Option<String>,
271 pub request_id: Option<String>,
272 pub prev_hash: String,
273 pub hash: String,
274}
275
276#[derive(Debug, Clone, PartialEq, Eq)]
278pub enum Visibility {
279 All,
281 Orgs(Vec<String>),
283}
284
285#[derive(Debug, Clone, Default, Deserialize)]
287#[serde(default, deny_unknown_fields)]
288pub struct Query {
289 pub org: Option<String>,
291 pub platform: bool,
292 pub actor: Option<String>,
294 pub user_id: Option<i64>,
295 pub token_id: Option<i64>,
296 pub action: Option<String>,
298 pub target: Option<String>,
300 pub outcome: Option<String>,
302 pub surface: Option<String>,
303 pub since: Option<i64>,
305 pub until: Option<i64>,
307 pub before: Option<i64>,
309 pub after: Option<i64>,
311 pub limit: Option<usize>,
313 pub object: Option<String>,
316 pub object_exact: bool,
317 pub ascending: bool,
319}
320
321#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
323pub struct Verified {
324 pub ok: bool,
325 pub rows: u64,
326 pub head: Option<(i64, String)>,
328 pub pruned_through: Option<i64>,
330 pub broken: Option<(i64, String)>,
332}
333
334pub struct AuditLog {
335 conn: Mutex<Connection>,
336 path: Option<PathBuf>,
337 retention: Duration,
338 pub(crate) history_retention: Duration,
340 pub(crate) history_max_rows: i64,
341 generation: Mutex<u64>,
343 appended: Condvar,
344 last_prune: Mutex<Option<Instant>>,
345 pub(crate) clock: Box<dyn Fn() -> i64 + Send + Sync>,
346}
347
348impl std::fmt::Debug for AuditLog {
349 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
350 f.debug_struct("AuditLog")
351 .field("path", &self.path)
352 .field("retention", &self.retention)
353 .finish()
354 }
355}
356
357pub fn db_path(state_dir: &Path) -> PathBuf {
358 state_dir.join("audit.db")
359}
360
361#[doc(hidden)]
362pub fn now_ms() -> i64 {
363 std::time::SystemTime::now()
364 .duration_since(std::time::UNIX_EPOCH)
365 .map(|d| d.as_millis() as i64)
366 .unwrap_or(0)
367}
368
369pub(crate) fn clip(mut s: String, n: usize) -> String {
370 if s.len() > n {
371 let mut i = n;
372 while !s.is_char_boundary(i) {
373 i -= 1;
374 }
375 s.truncate(i);
376 }
377 s.retain(|c| !c.is_control());
378 s
379}
380
381pub(crate) fn db_err(step: &str, e: rusqlite::Error) -> Error {
382 Error::Protocol(format!("audit log: {step}: {e}"))
383}
384
385pub(crate) fn canonical(v: &Value) -> String {
387 match v {
388 Value::Object(m) => {
389 let mut keys: Vec<&String> = m.keys().collect();
390 keys.sort();
391 let parts: Vec<String> = keys
392 .into_iter()
393 .map(|k| format!("{}:{}", Value::String(k.clone()), canonical(&m[k])))
394 .collect();
395 format!("{{{}}}", parts.join(","))
396 }
397 Value::Array(a) => format!(
398 "[{}]",
399 a.iter().map(canonical).collect::<Vec<_>>().join(",")
400 ),
401 other => other.to_string(),
402 }
403}
404
405pub(crate) fn hex(b: &[u8]) -> String {
406 b.iter().map(|x| format!("{x:02x}")).collect()
407}
408
409fn row_hash(e: &Entry) -> String {
411 let body = json!({
412 "id": e.id, "time": e.time, "org": e.org, "actor": e.actor,
413 "actor_kind": e.actor_kind, "user_id": e.user_id, "user_email": e.user_email,
414 "token_id": e.token_id, "token_name": e.token_name, "surface": e.surface,
415 "action": e.action, "target": e.target, "details": e.details,
416 "outcome": e.outcome, "ip": e.ip, "user_agent": e.user_agent,
417 "request_id": e.request_id,
418 });
419 let mut ctx = ring::digest::Context::new(&ring::digest::SHA256);
420 ctx.update(e.prev_hash.as_bytes());
421 ctx.update(b"\n");
422 ctx.update(canonical(&body).as_bytes());
423 hex(ctx.finish().as_ref())
424}
425
426const COLS: &str = "id, time, org, actor, actor_kind, user_id, user_email, token_id, token_name, \
427 surface, action, target, details, outcome, ip, user_agent, request_id, prev_hash, hash";
428
429fn row(r: &rusqlite::Row) -> rusqlite::Result<Entry> {
430 let details: String = r.get(12)?;
431 Ok(Entry {
432 id: r.get(0)?,
433 time: r.get(1)?,
434 org: r.get(2)?,
435 actor: r.get(3)?,
436 actor_kind: r.get(4)?,
437 user_id: r.get(5)?,
438 user_email: r.get(6)?,
439 token_id: r.get(7)?,
440 token_name: r.get(8)?,
441 surface: r.get(9)?,
442 action: r.get(10)?,
443 target: r.get(11)?,
444 details: serde_json::from_str(&details).unwrap_or(Value::Null),
445 outcome: r.get(13)?,
446 ip: r.get(14)?,
447 user_agent: r.get(15)?,
448 request_id: r.get(16)?,
449 prev_hash: r.get(17)?,
450 hash: r.get(18)?,
451 })
452}
453
454pub(crate) fn sql_glob(p: &str) -> String {
456 p.replace("[!", "[^")
457}
458
459pub(crate) fn meta(conn: &Connection, k: &str) -> rusqlite::Result<Option<String>> {
460 conn.query_row("SELECT v FROM audit_meta WHERE k = ?1", [k], |r| r.get(0))
461 .optional()
462}
463
464pub(crate) fn set_meta(conn: &Connection, k: &str, v: &str) -> rusqlite::Result<()> {
465 conn.execute(
466 "INSERT INTO audit_meta (k, v) VALUES (?1, ?2) ON CONFLICT(k) DO UPDATE SET v = ?2",
467 params![k, v],
468 )?;
469 Ok(())
470}
471
472impl AuditLog {
473 pub fn open(path: &Path, retention: Duration) -> Result<AuditLog> {
475 use std::os::unix::fs::{DirBuilderExt, OpenOptionsExt, PermissionsExt};
476 if let Some(dir) = path.parent().filter(|d| !d.as_os_str().is_empty()) {
477 if !dir.exists() {
478 std::fs::DirBuilder::new()
479 .recursive(true)
480 .mode(0o700)
481 .create(dir)?;
482 }
483 }
484 std::fs::OpenOptions::new()
485 .create(true)
486 .append(true)
487 .mode(0o600)
488 .open(path)?;
489 std::fs::set_permissions(path, std::fs::Permissions::from_mode(0o600))?;
490 let conn = Connection::open_with_flags(
491 path,
492 OpenFlags::SQLITE_OPEN_READ_WRITE
493 | OpenFlags::SQLITE_OPEN_CREATE
494 | OpenFlags::SQLITE_OPEN_NO_MUTEX,
495 )
496 .map_err(|e| db_err(&format!("open {}", path.display()), e))?;
497 Self::build(conn, Some(path.to_path_buf()), retention)
498 }
499
500 pub fn in_memory() -> Result<AuditLog> {
502 let conn = Connection::open_in_memory().map_err(|e| db_err("open", e))?;
503 Self::build(conn, None, DEFAULT_RETENTION)
504 }
505
506 fn build(conn: Connection, path: Option<PathBuf>, retention: Duration) -> Result<AuditLog> {
507 conn.busy_timeout(Duration::from_secs(5))
508 .map_err(|e| db_err("configure", e))?;
509 if path.is_some() {
510 let _: String = conn
511 .query_row("PRAGMA journal_mode = WAL", [], |r| r.get(0))
512 .map_err(|e| db_err("configure", e))?;
513 }
514 conn.execute_batch("PRAGMA synchronous = NORMAL;")
515 .map_err(|e| db_err("configure", e))?;
516 migrate(&conn)?;
517 let log = AuditLog {
518 conn: Mutex::new(conn),
519 path,
520 retention,
521 history_retention: crate::history::DEFAULT_RETENTION,
522 history_max_rows: crate::history::DEFAULT_MAX_ROWS,
523 generation: Mutex::new(0),
524 appended: Condvar::new(),
525 last_prune: Mutex::new(None),
526 clock: Box::new(now_ms),
527 };
528 log.prune()?;
529 Ok(log)
530 }
531
532 pub fn with_clock(mut self, clock: impl Fn() -> i64 + Send + Sync + 'static) -> Self {
534 self.clock = Box::new(clock);
535 self
536 }
537
538 pub fn with_history_limits(mut self, retention: Duration, max_rows: i64) -> Self {
540 self.history_retention = retention;
541 self.history_max_rows = max_rows.max(1000);
542 self
543 }
544
545 pub fn path(&self) -> Option<&Path> {
546 self.path.as_deref()
547 }
548
549 pub fn retention(&self) -> Duration {
550 self.retention
551 }
552
553 pub(crate) fn db(&self) -> MutexGuard<'_, Connection> {
554 self.conn.lock().unwrap_or_else(|e| e.into_inner())
555 }
556
557 pub fn append(&self, n: NewEntry) -> Result<Entry> {
559 let mut e = Entry {
560 id: 0,
561 time: (self.clock)(),
562 org: n.org.map(|o| clip(o, 64)),
563 actor: clip(n.actor.name, 300),
564 actor_kind: n
565 .actor
566 .kind
567 .unwrap_or(ActorKind::Anonymous)
568 .as_str()
569 .to_string(),
570 user_id: n.actor.user_id,
571 user_email: n.actor.email.map(|s| clip(s, 254)),
572 token_id: n.actor.token_id,
573 token_name: n.actor.token_name.map(|s| clip(s, 100)),
574 surface: clip(n.origin.surface, 16),
575 action: clip(n.action, 128),
576 target: n.target.map(|s| clip(s, 256)).filter(|s| !s.is_empty()),
577 details: match n.details {
578 Value::Null => json!({}),
579 v => v,
580 },
581 outcome: clip(n.outcome, 64),
582 ip: n.origin.ip.map(|s| clip(s, 64)),
583 user_agent: n.origin.user_agent.map(|s| clip(s, 256)),
584 request_id: n.origin.request_id.map(|s| clip(s, 64)),
585 prev_hash: String::new(),
586 hash: String::new(),
587 };
588 let db = self.db();
589 db.execute_batch("BEGIN IMMEDIATE")
590 .map_err(|e| db_err("append", e))?;
591 let r = (|| -> rusqlite::Result<()> {
592 let head_id: i64 = meta(&db, "head_id")?
593 .and_then(|v| v.parse().ok())
594 .unwrap_or(0);
595 e.prev_hash = meta(&db, "head_hash")?.unwrap_or_else(|| GENESIS.into());
596 e.id = head_id + 1;
597 e.hash = row_hash(&e);
598 db.execute(
599 &format!(
600 "INSERT INTO audit ({COLS}) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, \
601 ?11, ?12, ?13, ?14, ?15, ?16, ?17, ?18, ?19)"
602 ),
603 params![
604 e.id,
605 e.time,
606 e.org,
607 e.actor,
608 e.actor_kind,
609 e.user_id,
610 e.user_email,
611 e.token_id,
612 e.token_name,
613 e.surface,
614 e.action,
615 e.target,
616 e.details.to_string(),
617 e.outcome,
618 e.ip,
619 e.user_agent,
620 e.request_id,
621 e.prev_hash,
622 e.hash
623 ],
624 )?;
625 set_meta(&db, "head_id", &e.id.to_string())?;
626 set_meta(&db, "head_hash", &e.hash)?;
627 Ok(())
628 })();
629 match r {
630 Ok(()) => db
631 .execute_batch("COMMIT")
632 .map_err(|e| db_err("append", e))?,
633 Err(err) => {
634 let _ = db.execute_batch("ROLLBACK");
635 return Err(db_err("append", err));
636 }
637 }
638 drop(db);
639 self.appended_one();
640 Ok(e)
641 }
642
643 pub(crate) fn appended_one(&self) {
645 *self.generation.lock().unwrap_or_else(|e| e.into_inner()) += 1;
646 self.appended.notify_all();
647 let due = self
648 .last_prune
649 .lock()
650 .unwrap_or_else(|e| e.into_inner())
651 .is_none_or(|t| t.elapsed() >= PRUNE_EVERY);
652 if due {
653 if let Err(err) = self.prune() {
654 eprintln!("isb: {err}");
655 }
656 }
657 }
658
659 pub fn prune(&self) -> Result<usize> {
662 *self.last_prune.lock().unwrap_or_else(|e| e.into_inner()) = Some(Instant::now());
663 let n = self.prune_history()?;
664 Ok(n + self.prune_audit()?)
665 }
666
667 fn prune_audit(&self) -> Result<usize> {
668 let cutoff = (self.clock)() - self.retention.as_millis() as i64;
669 let db = self.db();
670 db.execute_batch("BEGIN IMMEDIATE")
671 .map_err(|e| db_err("prune", e))?;
672 let r = (|| -> rusqlite::Result<usize> {
673 let last: Option<(i64, String)> = db
674 .query_row(
675 "SELECT id, hash FROM audit WHERE time < ?1 ORDER BY id DESC LIMIT 1",
676 [cutoff],
677 |r| Ok((r.get(0)?, r.get(1)?)),
678 )
679 .optional()?;
680 let Some((id, hash)) = last else {
681 return Ok(0);
682 };
683 set_meta(&db, "pruning", "1")?;
684 let n = db.execute("DELETE FROM audit WHERE id <= ?1", [id])?;
685 set_meta(&db, "pruning", "0")?;
686 set_meta(&db, "pruned_through", &id.to_string())?;
687 set_meta(&db, "pruned_hash", &hash)?;
688 Ok(n)
689 })();
690 match r {
691 Ok(n) => {
692 db.execute_batch("COMMIT").map_err(|e| db_err("prune", e))?;
693 Ok(n)
694 }
695 Err(err) => {
696 let _ = db.execute_batch("ROLLBACK");
697 Err(db_err("prune", err))
698 }
699 }
700 }
701
702 #[expect(
705 clippy::too_many_lines,
706 reason = "predates the lint ratchet; split it when next changed"
707 )]
708 pub fn list(&self, q: &Query, vis: &Visibility) -> Result<Vec<Entry>> {
709 let mut sql = format!("SELECT {COLS} FROM audit WHERE 1=1");
710 let mut args: Vec<rusqlite::types::Value> = Vec::new();
711 fn push(
712 sql: &mut String,
713 args: &mut Vec<rusqlite::types::Value>,
714 cond: &str,
715 v: rusqlite::types::Value,
716 ) {
717 args.push(v);
718 sql.push_str(&cond.replace('?', &format!("?{}", args.len())));
719 }
720 use rusqlite::types::Value as V;
721 if let Visibility::Orgs(orgs) = vis {
722 if orgs.is_empty() {
723 return Ok(Vec::new());
724 }
725 sql.push_str(" AND org IN (");
726 for (i, o) in orgs.iter().enumerate() {
727 if i > 0 {
728 sql.push(',');
729 }
730 push(&mut sql, &mut args, "?", V::Text(o.clone()));
731 }
732 sql.push(')');
733 }
734 if q.platform {
735 sql.push_str(" AND org IS NULL");
736 } else if let Some(o) = &q.org {
737 push(&mut sql, &mut args, " AND org = ?", V::Text(o.clone()));
738 }
739 if let Some(a) = &q.actor {
740 let g = sql_glob(a);
741 args.push(V::Text(g));
742 let n = args.len();
743 sql.push_str(&format!(
744 " AND (actor GLOB ?{n} OR IFNULL(user_email, '') GLOB ?{n})"
745 ));
746 }
747 if let Some(u) = q.user_id {
748 push(&mut sql, &mut args, " AND user_id = ?", V::Integer(u));
749 }
750 if let Some(t) = q.token_id {
751 push(&mut sql, &mut args, " AND token_id = ?", V::Integer(t));
752 }
753 if let Some(a) = &q.action {
754 push(
755 &mut sql,
756 &mut args,
757 " AND action GLOB ?",
758 V::Text(sql_glob(a)),
759 );
760 }
761 if let Some(t) = &q.target {
762 push(
763 &mut sql,
764 &mut args,
765 " AND IFNULL(target, '') GLOB ?",
766 V::Text(sql_glob(t)),
767 );
768 }
769 match q.outcome.as_deref() {
770 None | Some("") => {}
771 Some("error") => sql.push_str(" AND outcome != 'ok'"),
772 Some(o) => push(&mut sql, &mut args, " AND outcome = ?", V::Text(o.into())),
773 }
774 if let Some(s) = &q.surface {
775 push(&mut sql, &mut args, " AND surface = ?", V::Text(s.clone()));
776 }
777 if let Some(t) = q.since {
778 push(&mut sql, &mut args, " AND time >= ?", V::Integer(t));
779 }
780 if let Some(t) = q.until {
781 push(&mut sql, &mut args, " AND time < ?", V::Integer(t));
782 }
783 if let Some(b) = q.before {
784 push(&mut sql, &mut args, " AND id < ?", V::Integer(b));
785 }
786 if let Some(a) = q.after {
787 push(&mut sql, &mut args, " AND id > ?", V::Integer(a));
788 }
789 if let Some(o) = q.object.as_deref().map(str::trim).filter(|o| !o.is_empty()) {
790 let esc = o
791 .replace('\\', "\\\\")
792 .replace('%', "\\%")
793 .replace('_', "\\_");
794 if q.object_exact {
795 args.push(V::Text(o.to_string()));
796 args.push(V::Text(format!("%\"{esc}\"%")));
797 let n = args.len();
798 sql.push_str(&format!(
799 " AND (target = ?{} OR details LIKE ?{n} ESCAPE '\\')",
800 n - 1
801 ));
802 } else {
803 args.push(V::Text(format!("%{esc}%")));
804 let n = args.len();
805 sql.push_str(&format!(
806 " AND (IFNULL(target, '') LIKE ?{n} ESCAPE '\\' OR details LIKE ?{n} ESCAPE '\\')"
807 ));
808 }
809 }
810 let limit = q.limit.unwrap_or(100).clamp(1, 1000);
811 sql.push_str(if q.after.is_some() || q.ascending {
812 " ORDER BY id ASC"
813 } else {
814 " ORDER BY id DESC"
815 });
816 sql.push_str(&format!(" LIMIT {limit}"));
817 let db = self.db();
818 let mut st = db.prepare(&sql).map_err(|e| db_err("query", e))?;
819 let rows = st
820 .query_map(rusqlite::params_from_iter(args), row)
821 .map_err(|e| db_err("query", e))?;
822 rows.collect::<rusqlite::Result<_>>()
823 .map_err(|e| db_err("query", e))
824 }
825
826 pub fn generation(&self) -> u64 {
828 *self.generation.lock().unwrap_or_else(|e| e.into_inner())
829 }
830
831 pub fn wait_change(&self, seen: u64, timeout: Duration) -> u64 {
835 let g = self.generation.lock().unwrap_or_else(|e| e.into_inner());
836 let g = self
837 .appended
838 .wait_timeout_while(g, timeout, |s| *s <= seen)
839 .unwrap_or_else(|e| e.into_inner())
840 .0;
841 *g
842 }
843
844 pub fn head(&self) -> Result<i64> {
846 let db = self.db();
847 Ok(meta(&db, "head_id")
848 .map_err(|e| db_err("head", e))?
849 .and_then(|v| v.parse().ok())
850 .unwrap_or(0))
851 }
852
853 pub fn verify(&self) -> Result<Verified> {
855 let db = self.db();
856 let m = |k: &str| meta(&db, k).map_err(|e| db_err("verify", e));
857 let pruned_through: Option<i64> = m("pruned_through")?.and_then(|v| v.parse().ok());
858 let mut expect = m("pruned_hash")?.unwrap_or_else(|| GENESIS.into());
859 let head_id: Option<i64> = m("head_id")?.and_then(|v| v.parse().ok());
860 let head_hash = m("head_hash")?;
861 let mut st = db
862 .prepare(&format!("SELECT {COLS} FROM audit ORDER BY id ASC"))
863 .map_err(|e| db_err("verify", e))?;
864 let rows = st.query_map([], row).map_err(|e| db_err("verify", e))?;
865 let mut out = Verified {
866 ok: true,
867 rows: 0,
868 head: None,
869 pruned_through,
870 broken: None,
871 };
872 let mut last_id = pruned_through.unwrap_or(0);
873 for r in rows {
874 let e = r.map_err(|e| db_err("verify", e))?;
875 out.rows += 1;
876 let why = if e.id != last_id + 1 {
877 Some(format!("expected row {} next, found {}", last_id + 1, e.id))
878 } else if e.prev_hash != expect {
879 Some("prev_hash does not match the row before it".to_string())
880 } else if row_hash(&e) != e.hash {
881 Some("the row's contents do not match its hash".to_string())
882 } else {
883 None
884 };
885 if let Some(w) = why {
886 out.ok = false;
887 out.broken = Some((e.id, w));
888 return Ok(out);
889 }
890 expect = e.hash.clone();
891 last_id = e.id;
892 out.head = Some((e.id, e.hash));
893 }
894 let tail_ok = match (&out.head, head_id, &head_hash) {
895 (Some((id, h)), Some(hid), Some(hh)) => *id == hid && h == hh,
896 (None, Some(hid), _) => Some(hid) == pruned_through,
898 (None, None, None) => true,
899 _ => false,
900 };
901 if !tail_ok {
902 out.ok = false;
903 out.broken = Some((
904 head_id.unwrap_or(0),
905 "the newest rows are missing (the head does not match)".into(),
906 ));
907 }
908 Ok(out)
909 }
910}
911
912fn migrate(conn: &Connection) -> Result<()> {
913 conn.execute_batch("CREATE TABLE IF NOT EXISTS audit_version (version INTEGER NOT NULL)")
914 .map_err(|e| db_err("migrate", e))?;
915 let mut v: i64 = conn
916 .query_row(
917 "SELECT IFNULL(MAX(version), 0) FROM audit_version",
918 [],
919 |r| r.get(0),
920 )
921 .map_err(|e| db_err("migrate", e))?;
922 let want = MIGRATIONS.len() as i64;
923 if v > want {
924 return Err(Error::invalid(format!(
925 "the audit log is schema version {v}, newer than this isb understands ({want}); upgrade isb"
926 )));
927 }
928 while v < want {
929 conn.execute_batch("BEGIN IMMEDIATE")
930 .map_err(|e| db_err("migrate", e))?;
931 let r = (|| -> rusqlite::Result<()> {
932 let now: i64 = conn.query_row(
933 "SELECT IFNULL(MAX(version), 0) FROM audit_version",
934 [],
935 |r| r.get(0),
936 )?;
937 if now != v {
938 return Ok(());
939 }
940 conn.execute_batch(MIGRATIONS[v as usize])?;
941 conn.execute("INSERT INTO audit_version (version) VALUES (?1)", [v + 1])?;
942 Ok(())
943 })();
944 match r {
945 Ok(()) => conn
946 .execute_batch("COMMIT")
947 .map_err(|e| db_err("migrate", e))?,
948 Err(e) => {
949 let _ = conn.execute_batch("ROLLBACK");
950 return Err(db_err(&format!("migration {}", v + 1), e));
951 }
952 }
953 v = conn
954 .query_row(
955 "SELECT IFNULL(MAX(version), 0) FROM audit_version",
956 [],
957 |r| r.get(0),
958 )
959 .map_err(|e| db_err("migrate", e))?;
960 }
961 Ok(())
962}
963
964pub fn ago_ms(d: Duration) -> i64 {
966 now_ms() - d.as_millis() as i64
967}
968
969#[cfg(test)]
970mod tests {
971 use super::*;
972 use std::sync::Arc;
973 use std::sync::atomic::{AtomicI64, Ordering};
974
975 fn entry(org: Option<&str>, action: &str, outcome: &str) -> NewEntry {
976 NewEntry {
977 org: org.map(String::from),
978 actor: Actor {
979 name: "a@x.io".into(),
980 kind: Some(ActorKind::Person),
981 user_id: Some(1),
982 email: Some("a@x.io".into()),
983 ..Default::default()
984 },
985 origin: Origin {
986 surface: "rest".into(),
987 ..Default::default()
988 },
989 action: action.into(),
990 target: Some("web".into()),
991 details: json!({"app": "web", "b": 1}),
992 outcome: outcome.into(),
993 }
994 }
995
996 #[test]
997 fn chain_verifies_and_detects_edits() {
998 let dir = tempfile::tempdir().unwrap();
999 let p = dir.path().join("audit.db");
1000 let log = AuditLog::open(&p, DEFAULT_RETENTION).unwrap();
1001 for i in 0..5 {
1002 log.append(entry(Some("acme"), &format!("stack_deploy{i}"), "ok"))
1003 .unwrap();
1004 }
1005 let v = log.verify().unwrap();
1006 assert!(v.ok, "{v:?}");
1007 assert_eq!(v.rows, 5);
1008 assert_eq!(v.head.as_ref().unwrap().0, 5);
1009 drop(log);
1010 let c = Connection::open(&p).unwrap();
1012 c.execute_batch(
1013 "DROP TRIGGER audit_no_update;
1014 UPDATE audit SET outcome = 'forbidden' WHERE id = 3;",
1015 )
1016 .unwrap();
1017 drop(c);
1018 let log = AuditLog::open(&p, DEFAULT_RETENTION).unwrap();
1019 let v = log.verify().unwrap();
1020 assert!(!v.ok);
1021 assert_eq!(v.broken.as_ref().unwrap().0, 3, "{v:?}");
1022 }
1023
1024 #[test]
1025 fn deleting_rows_is_refused_and_detected() {
1026 let dir = tempfile::tempdir().unwrap();
1027 let p = dir.path().join("audit.db");
1028 let log = AuditLog::open(&p, DEFAULT_RETENTION).unwrap();
1029 for _ in 0..4 {
1030 log.append(entry(None, "auth.login", "ok")).unwrap();
1031 }
1032 {
1033 let db = log.db();
1034 assert!(db.execute("DELETE FROM audit WHERE id = 2", []).is_err());
1035 assert!(
1036 db.execute("UPDATE audit SET actor = 'x' WHERE id = 2", [])
1037 .is_err()
1038 );
1039 }
1040 drop(log);
1041 let c = Connection::open(&p).unwrap();
1042 c.execute_batch("DROP TRIGGER audit_no_delete; DELETE FROM audit WHERE id = 2;")
1043 .unwrap();
1044 drop(c);
1045 let log = AuditLog::open(&p, DEFAULT_RETENTION).unwrap();
1046 let v = log.verify().unwrap();
1047 assert_eq!(v.broken.as_ref().unwrap().0, 3, "{v:?}");
1048 drop(log);
1050 let c = Connection::open(&p).unwrap();
1051 c.execute_batch("DELETE FROM audit WHERE id >= 2;").unwrap();
1052 drop(c);
1053 let log = AuditLog::open(&p, DEFAULT_RETENTION).unwrap();
1054 assert!(!log.verify().unwrap().ok);
1055 }
1056
1057 #[test]
1058 fn retention_prunes_and_keeps_the_chain_anchored() {
1059 let now = Arc::new(AtomicI64::new(1_000_000_000_000));
1060 let n = now.clone();
1061 let log = AuditLog::in_memory()
1062 .unwrap()
1063 .with_clock(move || n.load(Ordering::SeqCst));
1064 for _ in 0..3 {
1065 log.append(entry(Some("acme"), "a", "ok")).unwrap();
1066 }
1067 now.fetch_add(91 * 86_400_000, Ordering::SeqCst);
1068 log.append(entry(Some("acme"), "b", "ok")).unwrap();
1069 assert_eq!(log.prune().unwrap(), 3);
1070 let v = log.verify().unwrap();
1071 assert!(v.ok, "{v:?}");
1072 assert_eq!((v.rows, v.pruned_through), (1, Some(3)));
1073 log.append(entry(Some("acme"), "c", "ok")).unwrap();
1074 assert!(log.verify().unwrap().ok);
1075 now.fetch_add(91 * 86_400_000, Ordering::SeqCst);
1077 log.prune().unwrap();
1078 let v = log.verify().unwrap();
1079 assert!(v.ok && v.rows == 0, "{v:?}");
1080 let e = log.append(entry(None, "d", "ok")).unwrap();
1081 assert_eq!(e.id, 6);
1082 assert!(log.verify().unwrap().ok);
1083 }
1084
1085 #[test]
1086 fn queries_filter_and_respect_visibility() {
1087 let log = AuditLog::in_memory().unwrap();
1088 log.append(entry(Some("acme"), "stack_deploy", "ok"))
1089 .unwrap();
1090 log.append(entry(Some("acme"), "secret_get", "forbidden"))
1091 .unwrap();
1092 log.append(entry(Some("beta"), "secret_set", "ok")).unwrap();
1093 log.append(entry(None, "auth.login", "ok")).unwrap();
1094 let all = Visibility::All;
1095 let acme = Visibility::Orgs(vec!["acme".into()]);
1096 let q = |f: &dyn Fn(&mut Query)| {
1097 let mut q = Query::default();
1098 f(&mut q);
1099 q
1100 };
1101 assert_eq!(log.list(&Query::default(), &all).unwrap().len(), 4);
1102 let mine = log.list(&Query::default(), &acme).unwrap();
1103 assert_eq!(mine.len(), 2);
1104 assert!(mine[0].id > mine[1].id, "newest first");
1105 assert_eq!(
1106 log.list(&q(&|q| q.action = Some("secret_*".into())), &all)
1107 .unwrap()
1108 .len(),
1109 2
1110 );
1111 assert_eq!(
1112 log.list(&q(&|q| q.outcome = Some("error".into())), &all)
1113 .unwrap()
1114 .len(),
1115 1
1116 );
1117 assert_eq!(log.list(&q(&|q| q.platform = true), &all).unwrap().len(), 1);
1118 assert!(
1120 log.list(&q(&|q| q.platform = true), &acme)
1121 .unwrap()
1122 .is_empty()
1123 );
1124 assert!(
1125 log.list(&q(&|q| q.org = Some("beta".into())), &acme)
1126 .unwrap()
1127 .is_empty()
1128 );
1129 let tail = log.list(&q(&|q| q.after = Some(1)), &all).unwrap();
1130 assert_eq!(tail.iter().map(|e| e.id).collect::<Vec<_>>(), [2, 3, 4]);
1131 let page = log
1132 .list(
1133 &q(&|q| {
1134 q.before = Some(3);
1135 q.limit = Some(1)
1136 }),
1137 &all,
1138 )
1139 .unwrap();
1140 assert_eq!(page[0].id, 2);
1141 assert_eq!(
1142 log.list(&q(&|q| q.actor = Some("*@x.io".into())), &all)
1143 .unwrap()
1144 .len(),
1145 4
1146 );
1147 }
1148
1149 #[test]
1150 fn canonical_is_order_free() {
1151 let a: Value = serde_json::from_str(r#"{"b":1,"a":{"y":2,"x":[1,"s"]}}"#).unwrap();
1152 assert_eq!(canonical(&a), r#"{"a":{"x":[1,"s"],"y":2},"b":1}"#);
1153 }
1154}