g1t/services/events/src/audit.rs
Pick any line to see why it is the way it is: the commit, the pull request and issue it came from, and what the agent was thinking.
| Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API | 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 | ||
| Team plan, an open-source pool, monthly trials and honest metering; the sidebar for everyone; a workspace that stays put | 280 | /// Entries are kept this many days unless `AUDIT_KEEP_DAYS` says otherwise: |
| Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look | 281 | /// what the audit log reads back, the same on every plan (90 days). |
| 282 | pub const DEFAULT_KEEP_DAYS: u32 = 90; | |
| Team plan, an open-source pool, monthly trials and honest metering; the sidebar for everyone; a workspace that stays put | 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 | ||
| Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API | 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 | ||
| Team plan, an open-source pool, monthly trials and honest metering; the sidebar for everyone; a workspace that stays put | 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"); | |
| Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look | 355 | assert_eq!(DEFAULT_KEEP_DAYS, 90); |
| Team plan, an open-source pool, monthly trials and honest metering; the sidebar for everyone; a workspace that stays put | 356 | } |
| 357 | ||
| Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API | 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 | } |