g1t/services/events/src/audit.rs
| 1 | //! The audit log, kept beside the event log. See `g1t_contracts::audit`. |
| 2 | //! |
| 3 | //! Rows are only ever appended. The one exception is a workspace rename, |
| 4 | //! which moves the workspace's rows to its new slug, as every service does |
| 5 | //! with what it keeps under a slug. |
| 6 | |
| 7 | use g1t_contracts::audit::{ |
| 8 | AuditEntry, AuditPage, AuditVisibility, ListAuditArgs, MAX_AUDIT_PAGE, NewAuditEntry, |
| 9 | RecordAuditArgs, |
| 10 | }; |
| 11 | use g1t_contracts::events::{Event, WorkspaceRenamed}; |
| 12 | use g1t_contracts::new_id; |
| 13 | use g1t_contracts::time::rfc3339; |
| 14 | use g1t_kit::now_ms; |
| 15 | use serde::Deserialize; |
| 16 | use worker::wasm_bindgen::JsValue; |
| 17 | use worker::{D1Database, Result}; |
| 18 | |
| 19 | const DEFAULT_PAGE: u32 = 100; |
| 20 | /// More than one request ever records. |
| 21 | const MAX_BATCH: usize = 50; |
| 22 | const MAX_TEXT: usize = 500; |
| 23 | |
| 24 | fn text(value: &Option<String>) -> JsValue { |
| 25 | value |
| 26 | .as_deref() |
| 27 | .map_or(JsValue::NULL, |value| JsValue::from(clip(value))) |
| 28 | } |
| 29 | |
| 30 | fn clip(value: &str) -> String { |
| 31 | value.chars().take(MAX_TEXT).collect() |
| 32 | } |
| 33 | |
| 34 | #[derive(Deserialize)] |
| 35 | struct Row { |
| 36 | id: String, |
| 37 | time: String, |
| 38 | workspace: String, |
| 39 | actor_kind: String, |
| 40 | actor: String, |
| 41 | actor_id: String, |
| 42 | agent: Option<String>, |
| 43 | on_behalf_of: Option<String>, |
| 44 | run_id: Option<String>, |
| 45 | run_kind: Option<String>, |
| 46 | credential_id: Option<String>, |
| 47 | action: String, |
| 48 | surface: String, |
| 49 | repo: Option<String>, |
| 50 | number: Option<f64>, |
| 51 | git_ref: Option<String>, |
| 52 | path: Option<String>, |
| 53 | outcome: String, |
| 54 | rule: String, |
| 55 | result: Option<String>, |
| 56 | message: Option<String>, |
| 57 | request_id: String, |
| 58 | } |
| 59 | |
| 60 | impl Row { |
| 61 | /// Read back through the contract's own names, so the two cannot drift. |
| 62 | fn into_entry(self) -> Option<AuditEntry> { |
| 63 | let entry: NewAuditEntry = serde_json::from_value(serde_json::json!({ |
| 64 | "actorKind": self.actor_kind, |
| 65 | "actor": self.actor, |
| 66 | "actorId": self.actor_id, |
| 67 | "agent": self.agent, |
| 68 | "onBehalfOf": self.on_behalf_of, |
| 69 | "runId": self.run_id, |
| 70 | "runKind": self.run_kind, |
| 71 | "credentialId": self.credential_id, |
| 72 | "action": self.action, |
| 73 | "surface": self.surface, |
| 74 | "workspace": self.workspace, |
| 75 | "repo": self.repo, |
| 76 | "number": self.number.map(|n| n as u32), |
| 77 | "gitRef": self.git_ref, |
| 78 | "path": self.path, |
| 79 | "outcome": self.outcome, |
| 80 | "rule": self.rule, |
| 81 | "result": self.result, |
| 82 | "message": self.message, |
| 83 | "requestId": self.request_id, |
| 84 | })) |
| 85 | .ok()?; |
| 86 | Some(AuditEntry { |
| 87 | id: self.id, |
| 88 | time: self.time, |
| 89 | entry, |
| 90 | }) |
| 91 | } |
| 92 | } |
| 93 | |
| 94 | /// Appends entries. Returns how many were kept: one without a workspace |
| 95 | /// or an actor belongs to nobody's log and is dropped. |
| 96 | pub async fn record(db: &D1Database, a: RecordAuditArgs) -> Result<u32> { |
| 97 | let now = now_ms(); |
| 98 | let time = rfc3339(now); |
| 99 | let mut statements = Vec::new(); |
| 100 | for entry in a.entries.into_iter().take(MAX_BATCH) { |
| 101 | let Some(kind) = entry.actor.actor_kind else { |
| 102 | continue; |
| 103 | }; |
| 104 | if entry.target.workspace.is_empty() || entry.actor.actor.is_empty() { |
| 105 | continue; |
| 106 | } |
| 107 | statements.push( |
| 108 | db.prepare( |
| 109 | "INSERT INTO audit_entries (id, time, workspace, actor_kind, actor, actor_id, |
| 110 | agent, on_behalf_of, run_id, run_kind, credential_id, action, surface, repo, |
| 111 | number, git_ref, path, outcome, rule, result, message, request_id) |
| 112 | VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)", |
| 113 | ) |
| 114 | .bind(&[ |
| 115 | new_id("aud", now).into(), |
| 116 | time.as_str().into(), |
| 117 | entry.target.workspace.to_lowercase().into(), |
| 118 | kind.as_str().into(), |
| 119 | clip(&entry.actor.actor).into(), |
| 120 | clip(&entry.actor.actor_id).into(), |
| 121 | text(&entry.actor.agent), |
| 122 | text(&entry.actor.on_behalf_of), |
| 123 | text(&entry.actor.run_id), |
| 124 | text(&entry.actor.run_kind), |
| 125 | text(&entry.actor.credential_id), |
| 126 | clip(&entry.action).into(), |
| 127 | serde_json::to_value(entry.surface)? |
| 128 | .as_str() |
| 129 | .unwrap_or("rest") |
| 130 | .into(), |
| 131 | text(&entry.target.repo), |
| 132 | entry.target.number.map_or(JsValue::NULL, JsValue::from), |
| 133 | text(&entry.target.git_ref), |
| 134 | text(&entry.target.path), |
| 135 | entry.outcome.as_str().into(), |
| 136 | clip(&entry.rule).into(), |
| 137 | text(&entry.result), |
| 138 | text(&entry.message), |
| 139 | clip(&entry.request_id).into(), |
| 140 | ])?, |
| 141 | ); |
| 142 | } |
| 143 | let kept = statements.len() as u32; |
| 144 | if kept > 0 { |
| 145 | db.batch(statements).await?; |
| 146 | } |
| 147 | Ok(kept) |
| 148 | } |
| 149 | |
| 150 | /// The conditions and values of a query, built together so they stay in |
| 151 | /// step. |
| 152 | #[derive(Default)] |
| 153 | struct Filter { |
| 154 | conditions: Vec<String>, |
| 155 | values: Vec<Param>, |
| 156 | } |
| 157 | |
| 158 | /// A bound value, kept apart from `JsValue` so filters can be tested. |
| 159 | #[derive(Debug, PartialEq)] |
| 160 | enum Param { |
| 161 | Text(String), |
| 162 | Number(u32), |
| 163 | } |
| 164 | |
| 165 | impl From<&str> for Param { |
| 166 | fn from(value: &str) -> Self { |
| 167 | Param::Text(value.to_owned()) |
| 168 | } |
| 169 | } |
| 170 | |
| 171 | impl From<String> for Param { |
| 172 | fn from(value: String) -> Self { |
| 173 | Param::Text(value) |
| 174 | } |
| 175 | } |
| 176 | |
| 177 | impl From<u32> for Param { |
| 178 | fn from(value: u32) -> Self { |
| 179 | Param::Number(value) |
| 180 | } |
| 181 | } |
| 182 | |
| 183 | impl From<&Param> for JsValue { |
| 184 | fn from(param: &Param) -> Self { |
| 185 | match param { |
| 186 | Param::Text(text) => JsValue::from(text.as_str()), |
| 187 | Param::Number(number) => JsValue::from(*number), |
| 188 | } |
| 189 | } |
| 190 | } |
| 191 | |
| 192 | impl Filter { |
| 193 | fn add(&mut self, condition: &str, values: impl IntoIterator<Item = Param>) { |
| 194 | self.conditions.push(condition.to_owned()); |
| 195 | self.values.extend(values); |
| 196 | } |
| 197 | } |
| 198 | |
| 199 | fn filter(a: &ListAuditArgs) -> Filter { |
| 200 | let mut f = Filter::default(); |
| 201 | f.add("workspace = ?", [a.workspace.to_lowercase().into()]); |
| 202 | if let AuditVisibility::Projects { username } = &a.visibility { |
| 203 | f.add( |
| 204 | "(repo IS NOT NULL OR actor = ? OR on_behalf_of = ?)", |
| 205 | [username.as_str().into(), username.as_str().into()], |
| 206 | ); |
| 207 | } |
| 208 | if let Some(actor) = a.actor.as_deref().filter(|v| !v.is_empty()) { |
| 209 | f.add( |
| 210 | "(actor = ? OR on_behalf_of = ?)", |
| 211 | [actor.into(), actor.into()], |
| 212 | ); |
| 213 | } |
| 214 | if let Some(agent) = a.agent.as_deref().filter(|v| !v.is_empty()) { |
| 215 | f.add("agent = ?", [agent.into()]); |
| 216 | } |
| 217 | if let Some(action) = a.action.as_deref().filter(|v| !v.is_empty()) { |
| 218 | f.add("action = ?", [action.into()]); |
| 219 | } |
| 220 | if let Some(repo) = a.repo.as_deref().filter(|v| !v.is_empty()) { |
| 221 | f.add("repo = ? COLLATE NOCASE", [repo.into()]); |
| 222 | } |
| 223 | if let Some(number) = a.number { |
| 224 | f.add("number = ?", [number.into()]); |
| 225 | } |
| 226 | if let Some(outcome) = a.outcome { |
| 227 | f.add("outcome = ?", [outcome.as_str().into()]); |
| 228 | } |
| 229 | if let Some(kind) = a.actor_kind { |
| 230 | f.add("actor_kind = ?", [kind.as_str().into()]); |
| 231 | } |
| 232 | if !a.run_ids.is_empty() { |
| 233 | let ids: Vec<&String> = a.run_ids.iter().take(50).collect(); |
| 234 | let marks = vec!["?"; ids.len()].join(", "); |
| 235 | f.add( |
| 236 | &format!("run_id IN ({marks})"), |
| 237 | ids.into_iter().map(|id| Param::from(id.as_str())), |
| 238 | ); |
| 239 | } |
| 240 | if let Some(since) = a.since.as_deref().filter(|v| !v.is_empty()) { |
| 241 | f.add("time >= ?", [since.into()]); |
| 242 | } |
| 243 | if let Some(until) = a.until.as_deref().filter(|v| !v.is_empty()) { |
| 244 | f.add("time < ?", [until.into()]); |
| 245 | } |
| 246 | if let Some(before) = a.before.as_deref().filter(|v| !v.is_empty()) { |
| 247 | f.add("id < ?", [before.into()]); |
| 248 | } |
| 249 | f |
| 250 | } |
| 251 | |
| 252 | /// Newest first. The caller has checked that the viewer may see the |
| 253 | /// workspace's log, and says how much of it in `visibility`. |
| 254 | pub async fn list(db: &D1Database, a: ListAuditArgs) -> Result<AuditPage> { |
| 255 | let limit = a.limit.unwrap_or(DEFAULT_PAGE).clamp(1, MAX_AUDIT_PAGE); |
| 256 | let Filter { conditions, values } = filter(&a); |
| 257 | let mut values: Vec<JsValue> = values.iter().map(JsValue::from).collect(); |
| 258 | values.push((limit + 1).into()); |
| 259 | let rows = db |
| 260 | .prepare(format!( |
| 261 | "SELECT * FROM audit_entries WHERE {} ORDER BY id DESC LIMIT ?", |
| 262 | conditions.join(" AND ") |
| 263 | )) |
| 264 | .bind(&values)? |
| 265 | .all() |
| 266 | .await? |
| 267 | .results::<Row>()?; |
| 268 | let more = rows.len() > limit as usize; |
| 269 | let entries: Vec<AuditEntry> = rows |
| 270 | .into_iter() |
| 271 | .take(limit as usize) |
| 272 | .filter_map(Row::into_entry) |
| 273 | .collect(); |
| 274 | let next = more |
| 275 | .then(|| entries.last().map(|entry| entry.id.clone())) |
| 276 | .flatten(); |
| 277 | Ok(AuditPage { entries, next }) |
| 278 | } |
| 279 | |
| 280 | /// Entries are kept this many days unless `AUDIT_KEEP_DAYS` says otherwise: |
| 281 | /// what the audit log reads back, the same on every plan (90 days). |
| 282 | pub const DEFAULT_KEEP_DAYS: u32 = 90; |
| 283 | /// Rows removed per statement, so one purge never runs long. |
| 284 | const PURGE_BATCH: u32 = 5_000; |
| 285 | |
| 286 | /// The oldest time an entry is kept from, `keep_days` before `now_ms`. |
| 287 | pub fn keep_from(now_ms: u64, keep_days: u32) -> String { |
| 288 | g1t_contracts::time::rfc3339(now_ms.saturating_sub(u64::from(keep_days) * 24 * 60 * 60 * 1000)) |
| 289 | } |
| 290 | |
| 291 | /// Removes entries older than every plan keeps, a batch at a time, up to |
| 292 | /// `rounds` batches. Returns how many went. |
| 293 | pub async fn purge(db: &D1Database, before: &str, rounds: u32) -> Result<u32> { |
| 294 | let mut removed = 0; |
| 295 | for _ in 0..rounds { |
| 296 | let result = db |
| 297 | .prepare( |
| 298 | "DELETE FROM audit_entries WHERE id IN |
| 299 | (SELECT id FROM audit_entries WHERE time < ? ORDER BY time LIMIT ?)", |
| 300 | ) |
| 301 | .bind(&[before.into(), PURGE_BATCH.into()])? |
| 302 | .run() |
| 303 | .await?; |
| 304 | let changed = result.meta()?.and_then(|meta| meta.changes).unwrap_or(0) as u32; |
| 305 | removed += changed; |
| 306 | if changed < PURGE_BATCH { |
| 307 | break; |
| 308 | } |
| 309 | } |
| 310 | Ok(removed) |
| 311 | } |
| 312 | |
| 313 | /// Moves a renamed workspace's rows to its new slug. |
| 314 | pub async fn follow_renames(db: &D1Database, events: &[Event]) -> Result<()> { |
| 315 | for event in events |
| 316 | .iter() |
| 317 | .filter(|event| event.kind == "workspace.renamed") |
| 318 | { |
| 319 | let Ok(renamed) = serde_json::from_value::<WorkspaceRenamed>(event.data.clone()) else { |
| 320 | continue; |
| 321 | }; |
| 322 | let (from, to) = (renamed.from.to_lowercase(), renamed.to.to_lowercase()); |
| 323 | if from == to { |
| 324 | continue; |
| 325 | } |
| 326 | db.batch(vec![ |
| 327 | db.prepare( |
| 328 | "UPDATE audit_entries SET repo = ? || substr(repo, length(?) + 1) |
| 329 | WHERE workspace = ? AND repo LIKE ? || '/%'", |
| 330 | ) |
| 331 | .bind(&[ |
| 332 | to.as_str().into(), |
| 333 | from.as_str().into(), |
| 334 | from.as_str().into(), |
| 335 | from.as_str().into(), |
| 336 | ])?, |
| 337 | db.prepare("UPDATE audit_entries SET workspace = ? WHERE workspace = ?") |
| 338 | .bind(&[to.as_str().into(), from.as_str().into()])?, |
| 339 | ]) |
| 340 | .await?; |
| 341 | } |
| 342 | Ok(()) |
| 343 | } |
| 344 | |
| 345 | #[cfg(test)] |
| 346 | mod tests { |
| 347 | use super::*; |
| 348 | use g1t_contracts::audit::{ActorKind, AuditOutcome}; |
| 349 | |
| 350 | #[test] |
| 351 | fn entries_are_kept_as_long_as_the_longest_plan_reads_back() { |
| 352 | // 2026-10-05T00:00:00Z, a year back. |
| 353 | let now = 1_791_158_400_000; |
| 354 | assert_eq!(keep_from(now, 365), "2025-10-05T00:00:00.000Z"); |
| 355 | assert_eq!(DEFAULT_KEEP_DAYS, 90); |
| 356 | } |
| 357 | |
| 358 | fn args() -> ListAuditArgs { |
| 359 | ListAuditArgs { |
| 360 | workspace: "Acme".to_owned(), |
| 361 | visibility: AuditVisibility::All, |
| 362 | actor: None, |
| 363 | agent: None, |
| 364 | action: None, |
| 365 | repo: None, |
| 366 | number: None, |
| 367 | outcome: None, |
| 368 | actor_kind: None, |
| 369 | run_ids: vec![], |
| 370 | since: None, |
| 371 | until: None, |
| 372 | before: None, |
| 373 | limit: None, |
| 374 | } |
| 375 | } |
| 376 | |
| 377 | #[test] |
| 378 | fn an_owner_sees_the_whole_workspace() { |
| 379 | let f = filter(&args()); |
| 380 | assert_eq!(f.conditions, ["workspace = ?"]); |
| 381 | assert_eq!(f.values.len(), 1); |
| 382 | } |
| 383 | |
| 384 | #[test] |
| 385 | fn a_member_sees_projects_and_their_own() { |
| 386 | let mut a = args(); |
| 387 | a.visibility = AuditVisibility::Projects { |
| 388 | username: "ana".to_owned(), |
| 389 | }; |
| 390 | let f = filter(&a); |
| 391 | assert!(f.conditions[1].contains("repo IS NOT NULL")); |
| 392 | assert_eq!(f.values.len(), 3); |
| 393 | } |
| 394 | |
| 395 | #[test] |
| 396 | fn every_filter_binds_its_values() { |
| 397 | let mut a = args(); |
| 398 | a.actor = Some("syntaqx".to_owned()); |
| 399 | a.agent = Some("g1t-agent".to_owned()); |
| 400 | a.action = Some("git.push".to_owned()); |
| 401 | a.repo = Some("acme/rocket".to_owned()); |
| 402 | a.number = Some(4); |
| 403 | a.outcome = Some(AuditOutcome::Denied); |
| 404 | a.actor_kind = Some(ActorKind::Agent); |
| 405 | a.run_ids = vec!["run_1".to_owned(), "run_2".to_owned()]; |
| 406 | a.since = Some("2026-10-01T00:00:00Z".to_owned()); |
| 407 | a.until = Some("2026-10-05T00:00:00Z".to_owned()); |
| 408 | a.before = Some("aud_9".to_owned()); |
| 409 | let f = filter(&a); |
| 410 | let marks: usize = f.conditions.iter().map(|c| c.matches('?').count()).sum(); |
| 411 | assert_eq!(marks, f.values.len()); |
| 412 | assert!(f.conditions.iter().any(|c| c == "run_id IN (?, ?)")); |
| 413 | } |
| 414 | } |