g1t/services/events/src/audit.rs

592 lines19,966 bytesCodeBlame

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 API1//! 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
7use g1t_contracts::audit::{
8 AuditEntry, AuditPage, AuditVisibility, ListAuditArgs, MAX_AUDIT_PAGE, NewAuditEntry,
9 RecordAuditArgs,
10};
Audit logs are kept by plan: a week on free, 90 days on the plan, and what staff set for an account in sudo11use g1t_contracts::billing::{AuditRetention, AuditRetentionArgs};
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API12use g1t_contracts::events::{Event, WorkspaceRenamed};
13use g1t_contracts::new_id;
14use g1t_contracts::time::rfc3339;
15use g1t_kit::now_ms;
16use serde::Deserialize;
17use worker::wasm_bindgen::JsValue;
Audit logs are kept by plan: a week on free, 90 days on the plan, and what staff set for an account in sudo18use worker::{D1Database, Fetcher, Result};
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API19
20const DEFAULT_PAGE: u32 = 100;
21/// More than one request ever records.
22const MAX_BATCH: usize = 50;
23const MAX_TEXT: usize = 500;
24
25fn text(value: &Option<String>) -> JsValue {
26 value
27 .as_deref()
28 .map_or(JsValue::NULL, |value| JsValue::from(clip(value)))
29}
30
31fn clip(value: &str) -> String {
32 value.chars().take(MAX_TEXT).collect()
33}
34
35#[derive(Deserialize)]
36struct Row {
37 id: String,
38 time: String,
39 workspace: String,
40 actor_kind: String,
41 actor: String,
42 actor_id: String,
43 agent: Option<String>,
44 on_behalf_of: Option<String>,
45 run_id: Option<String>,
46 run_kind: Option<String>,
47 credential_id: Option<String>,
48 action: String,
49 surface: String,
50 repo: Option<String>,
51 number: Option<f64>,
52 git_ref: Option<String>,
53 path: Option<String>,
54 outcome: String,
55 rule: String,
56 result: Option<String>,
57 message: Option<String>,
58 request_id: String,
59}
60
61impl Row {
62 /// Read back through the contract's own names, so the two cannot drift.
63 fn into_entry(self) -> Option<AuditEntry> {
64 let entry: NewAuditEntry = serde_json::from_value(serde_json::json!({
65 "actorKind": self.actor_kind,
66 "actor": self.actor,
67 "actorId": self.actor_id,
68 "agent": self.agent,
69 "onBehalfOf": self.on_behalf_of,
70 "runId": self.run_id,
71 "runKind": self.run_kind,
72 "credentialId": self.credential_id,
73 "action": self.action,
74 "surface": self.surface,
75 "workspace": self.workspace,
76 "repo": self.repo,
77 "number": self.number.map(|n| n as u32),
78 "gitRef": self.git_ref,
79 "path": self.path,
80 "outcome": self.outcome,
81 "rule": self.rule,
82 "result": self.result,
83 "message": self.message,
84 "requestId": self.request_id,
85 }))
86 .ok()?;
87 Some(AuditEntry {
88 id: self.id,
89 time: self.time,
90 entry,
91 })
92 }
93}
94
95/// Appends entries. Returns how many were kept: one without a workspace
96/// or an actor belongs to nobody's log and is dropped.
97pub async fn record(db: &D1Database, a: RecordAuditArgs) -> Result<u32> {
98 let now = now_ms();
99 let time = rfc3339(now);
100 let mut statements = Vec::new();
101 for entry in a.entries.into_iter().take(MAX_BATCH) {
102 let Some(kind) = entry.actor.actor_kind else {
103 continue;
104 };
105 if entry.target.workspace.is_empty() || entry.actor.actor.is_empty() {
106 continue;
107 }
108 statements.push(
109 db.prepare(
110 "INSERT INTO audit_entries (id, time, workspace, actor_kind, actor, actor_id,
111 agent, on_behalf_of, run_id, run_kind, credential_id, action, surface, repo,
112 number, git_ref, path, outcome, rule, result, message, request_id)
113 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
114 )
115 .bind(&[
116 new_id("aud", now).into(),
117 time.as_str().into(),
118 entry.target.workspace.to_lowercase().into(),
119 kind.as_str().into(),
120 clip(&entry.actor.actor).into(),
121 clip(&entry.actor.actor_id).into(),
122 text(&entry.actor.agent),
123 text(&entry.actor.on_behalf_of),
124 text(&entry.actor.run_id),
125 text(&entry.actor.run_kind),
126 text(&entry.actor.credential_id),
127 clip(&entry.action).into(),
128 serde_json::to_value(entry.surface)?
129 .as_str()
130 .unwrap_or("rest")
131 .into(),
132 text(&entry.target.repo),
133 entry.target.number.map_or(JsValue::NULL, JsValue::from),
134 text(&entry.target.git_ref),
135 text(&entry.target.path),
136 entry.outcome.as_str().into(),
137 clip(&entry.rule).into(),
138 text(&entry.result),
139 text(&entry.message),
140 clip(&entry.request_id).into(),
141 ])?,
142 );
143 }
144 let kept = statements.len() as u32;
145 if kept > 0 {
146 db.batch(statements).await?;
147 }
148 Ok(kept)
149}
150
151/// The conditions and values of a query, built together so they stay in
152/// step.
153#[derive(Default)]
154struct Filter {
155 conditions: Vec<String>,
156 values: Vec<Param>,
157}
158
159/// A bound value, kept apart from `JsValue` so filters can be tested.
160#[derive(Debug, PartialEq)]
161enum Param {
162 Text(String),
163 Number(u32),
164}
165
166impl From<&str> for Param {
167 fn from(value: &str) -> Self {
168 Param::Text(value.to_owned())
169 }
170}
171
172impl From<String> for Param {
173 fn from(value: String) -> Self {
174 Param::Text(value)
175 }
176}
177
178impl From<u32> for Param {
179 fn from(value: u32) -> Self {
180 Param::Number(value)
181 }
182}
183
184impl From<&Param> for JsValue {
185 fn from(param: &Param) -> Self {
186 match param {
187 Param::Text(text) => JsValue::from(text.as_str()),
188 Param::Number(number) => JsValue::from(*number),
189 }
190 }
191}
192
193impl Filter {
194 fn add(&mut self, condition: &str, values: impl IntoIterator<Item = Param>) {
195 self.conditions.push(condition.to_owned());
196 self.values.extend(values);
197 }
198}
199
200fn filter(a: &ListAuditArgs) -> Filter {
201 let mut f = Filter::default();
202 f.add("workspace = ?", [a.workspace.to_lowercase().into()]);
203 if let AuditVisibility::Projects { username } = &a.visibility {
204 f.add(
205 "(repo IS NOT NULL OR actor = ? OR on_behalf_of = ?)",
206 [username.as_str().into(), username.as_str().into()],
207 );
208 }
209 if let Some(actor) = a.actor.as_deref().filter(|v| !v.is_empty()) {
210 f.add(
211 "(actor = ? OR on_behalf_of = ?)",
212 [actor.into(), actor.into()],
213 );
214 }
215 if let Some(agent) = a.agent.as_deref().filter(|v| !v.is_empty()) {
216 f.add("agent = ?", [agent.into()]);
217 }
218 if let Some(action) = a.action.as_deref().filter(|v| !v.is_empty()) {
219 f.add("action = ?", [action.into()]);
220 }
221 if let Some(repo) = a.repo.as_deref().filter(|v| !v.is_empty()) {
222 f.add("repo = ? COLLATE NOCASE", [repo.into()]);
223 }
224 if let Some(number) = a.number {
225 f.add("number = ?", [number.into()]);
226 }
227 if let Some(outcome) = a.outcome {
228 f.add("outcome = ?", [outcome.as_str().into()]);
229 }
230 if let Some(kind) = a.actor_kind {
231 f.add("actor_kind = ?", [kind.as_str().into()]);
232 }
233 if !a.run_ids.is_empty() {
234 let ids: Vec<&String> = a.run_ids.iter().take(50).collect();
235 let marks = vec!["?"; ids.len()].join(", ");
236 f.add(
237 &format!("run_id IN ({marks})"),
238 ids.into_iter().map(|id| Param::from(id.as_str())),
239 );
240 }
241 if let Some(since) = a.since.as_deref().filter(|v| !v.is_empty()) {
242 f.add("time >= ?", [since.into()]);
243 }
244 if let Some(until) = a.until.as_deref().filter(|v| !v.is_empty()) {
245 f.add("time < ?", [until.into()]);
246 }
247 if let Some(before) = a.before.as_deref().filter(|v| !v.is_empty()) {
248 f.add("id < ?", [before.into()]);
249 }
250 f
251}
252
253/// Newest first. The caller has checked that the viewer may see the
254/// workspace's log, and says how much of it in `visibility`.
255pub async fn list(db: &D1Database, a: ListAuditArgs) -> Result<AuditPage> {
256 let limit = a.limit.unwrap_or(DEFAULT_PAGE).clamp(1, MAX_AUDIT_PAGE);
257 let Filter { conditions, values } = filter(&a);
258 let mut values: Vec<JsValue> = values.iter().map(JsValue::from).collect();
259 values.push((limit + 1).into());
260 let rows = db
261 .prepare(format!(
262 "SELECT * FROM audit_entries WHERE {} ORDER BY id DESC LIMIT ?",
263 conditions.join(" AND ")
264 ))
265 .bind(&values)?
266 .all()
267 .await?
268 .results::<Row>()?;
269 let more = rows.len() > limit as usize;
270 let entries: Vec<AuditEntry> = rows
271 .into_iter()
272 .take(limit as usize)
273 .filter_map(Row::into_entry)
274 .collect();
275 let next = more
276 .then(|| entries.last().map(|entry| entry.id.clone()))
277 .flatten();
278 Ok(AuditPage { entries, next })
279}
280
Audit logs are kept by plan: a week on free, 90 days on the plan, and what staff set for an account in sudo281/// No entry is kept longer than this, whatever its workspace's plan, unless
282/// `AUDIT_MAX_DAYS` says otherwise. Keep it at least billing's
283/// `AUDIT_MAX_DAYS`, the most staff can set for an account.
284pub const DEFAULT_MAX_DAYS: u32 = 400;
285/// The shortest any workspace keeps (`AUDIT_MIN_DAYS`, a free workspace's
286/// 7 days): only workspaces with entries older than this are asked about.
287pub const DEFAULT_MIN_DAYS: u32 = 7;
288/// Workspaces looked at in one daily run, so one run never runs long; the
289/// next run carries on after the last one.
290pub const WORKSPACES_PER_RUN: u32 = 200;
291/// Workspaces asked about in one call to billing, which reads each one's
292/// account and plan.
293const ASK_AT_ONCE: usize = 50;
Team plan, an open-source pool, monthly trials and honest metering; the sidebar for everyone; a workspace that stays put294/// Rows removed per statement, so one purge never runs long.
295const PURGE_BATCH: u32 = 5_000;
296
297/// The oldest time an entry is kept from, `keep_days` before `now_ms`.
298pub fn keep_from(now_ms: u64, keep_days: u32) -> String {
299 g1t_contracts::time::rfc3339(now_ms.saturating_sub(u64::from(keep_days) * 24 * 60 * 60 * 1000))
300}
301
Audit logs are kept by plan: a week on free, 90 days on the plan, and what staff set for an account in sudo302/// Each workspace with the time its entries are kept from. Its days are
303/// held between the shortest and the longest any workspace keeps: fewer
304/// than the shortest would not be looked for, and more than the longest
305/// are deleted anyway.
306pub fn cutoffs(
307 now_ms: u64,
308 retention: &[AuditRetention],
309 min_days: u32,
310 max_days: u32,
311) -> Vec<(String, String)> {
312 retention
313 .iter()
314 .map(|r| {
315 let days = r.days.max(min_days).min(max_days);
316 (r.workspace.clone(), keep_from(now_ms, days))
317 })
318 .collect()
319}
320
Team plan, an open-source pool, monthly trials and honest metering; the sidebar for everyone; a workspace that stays put321/// Removes entries older than every plan keeps, a batch at a time, up to
322/// `rounds` batches. Returns how many went.
323pub async fn purge(db: &D1Database, before: &str, rounds: u32) -> Result<u32> {
324 let mut removed = 0;
325 for _ in 0..rounds {
326 let result = db
327 .prepare(
328 "DELETE FROM audit_entries WHERE id IN
329 (SELECT id FROM audit_entries WHERE time < ? ORDER BY time LIMIT ?)",
330 )
331 .bind(&[before.into(), PURGE_BATCH.into()])?
332 .run()
333 .await?;
334 let changed = result.meta()?.and_then(|meta| meta.changes).unwrap_or(0) as u32;
335 removed += changed;
336 if changed < PURGE_BATCH {
337 break;
338 }
339 }
340 Ok(removed)
341}
342
Audit logs are kept by plan: a week on free, 90 days on the plan, and what staff set for an account in sudo343/// Removes one workspace's entries older than `before`, the same way.
344async fn purge_workspace(
345 db: &D1Database,
346 workspace: &str,
347 before: &str,
348 rounds: u32,
349) -> Result<u32> {
350 let mut removed = 0;
351 for _ in 0..rounds {
352 let result = db
353 .prepare(
354 "DELETE FROM audit_entries WHERE id IN
355 (SELECT id FROM audit_entries WHERE workspace = ? AND time < ? ORDER BY time LIMIT ?)",
356 )
357 .bind(&[workspace.into(), before.into(), PURGE_BATCH.into()])?
358 .run()
359 .await?;
360 let changed = result.meta()?.and_then(|meta| meta.changes).unwrap_or(0) as u32;
361 removed += changed;
362 if changed < PURGE_BATCH {
363 break;
364 }
365 }
366 Ok(removed)
367}
368
369/// Up to `limit` workspaces after `after`, in order, that have entries
370/// older than `before`. Each step seeks the next workspace in the index on
371/// (workspace, time) rather than reading every row, which a DISTINCT over
372/// the rows older than a week would: a workspace on the plan always has
373/// weeks of them.
374async fn workspaces_past(
375 db: &D1Database,
376 after: &str,
377 before: &str,
378 limit: u32,
379) -> Result<Vec<String>> {
380 #[derive(Deserialize)]
381 struct Found {
382 workspace: String,
383 }
384 Ok(db
385 .prepare(
386 "WITH RECURSIVE w(workspace) AS (
387 SELECT (SELECT MIN(workspace) FROM audit_entries WHERE workspace > ?1)
388 UNION ALL
389 SELECT (SELECT MIN(e.workspace) FROM audit_entries e WHERE e.workspace > w.workspace)
390 FROM w WHERE w.workspace IS NOT NULL
391 )
392 SELECT workspace FROM w
393 WHERE workspace IS NOT NULL
394 AND EXISTS (SELECT 1 FROM audit_entries a WHERE a.workspace = w.workspace AND a.time < ?2)
395 LIMIT ?3",
396 )
397 .bind(&[after.into(), before.into(), limit.into()])?
398 .all()
399 .await?
400 .results::<Found>()?
401 .into_iter()
402 .map(|found| found.workspace)
403 .collect())
404}
405
406/// Removes each workspace's entries older than its plan keeps, for up to
407/// `WORKSPACES_PER_RUN` workspaces after where the last run stopped. Their
408/// days come from billing; if it cannot be reached, nothing is removed and
409/// the next run tries the same workspaces again. Returns how many went.
410pub async fn purge_by_plan(
411 db: &D1Database,
412 billing: &Fetcher,
413 now_ms: u64,
414 min_days: u32,
415 max_days: u32,
416) -> Result<u32> {
417 #[derive(Deserialize)]
418 struct Cursor {
419 after: String,
420 }
421 let after = db
422 .prepare("SELECT after FROM audit_purge_cursor WHERE id = 1")
423 .first::<Cursor>(None)
424 .await?
425 .map_or_else(String::new, |cursor| cursor.after);
426 let found =
427 workspaces_past(db, &after, &keep_from(now_ms, min_days), WORKSPACES_PER_RUN).await?;
428 // Every workspace's days are asked for before anything is removed, so a
429 // billing that cannot be reached removes nothing at all.
430 let mut retention: Vec<AuditRetention> = Vec::with_capacity(found.len());
431 for chunk in found.chunks(ASK_AT_ONCE) {
432 let args = AuditRetentionArgs {
433 workspaces: chunk.to_vec(),
434 };
435 let answered: Vec<AuditRetention> =
436 g1t_kit::call(billing, "audit_retention", &args).await?;
437 retention.extend(answered);
438 }
439 let mut removed = 0;
440 for (workspace, before) in cutoffs(now_ms, &retention, min_days, max_days) {
441 removed += purge_workspace(db, &workspace, &before, 4).await?;
442 }
443 // A short page means the end was reached: the next run starts over.
444 let next = if found.len() < WORKSPACES_PER_RUN as usize {
445 String::new()
446 } else {
447 found.last().cloned().unwrap_or_default()
448 };
449 db.prepare(
450 "INSERT INTO audit_purge_cursor (id, after) VALUES (1, ?1)
451 ON CONFLICT (id) DO UPDATE SET after = ?1",
452 )
453 .bind(&[next.into()])?
454 .run()
455 .await?;
456 Ok(removed)
457}
458
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API459/// Moves a renamed workspace's rows to its new slug.
460pub async fn follow_renames(db: &D1Database, events: &[Event]) -> Result<()> {
461 for event in events
462 .iter()
463 .filter(|event| event.kind == "workspace.renamed")
464 {
465 let Ok(renamed) = serde_json::from_value::<WorkspaceRenamed>(event.data.clone()) else {
466 continue;
467 };
468 let (from, to) = (renamed.from.to_lowercase(), renamed.to.to_lowercase());
469 if from == to {
470 continue;
471 }
472 db.batch(vec![
473 db.prepare(
474 "UPDATE audit_entries SET repo = ? || substr(repo, length(?) + 1)
475 WHERE workspace = ? AND repo LIKE ? || '/%'",
476 )
477 .bind(&[
478 to.as_str().into(),
479 from.as_str().into(),
480 from.as_str().into(),
481 from.as_str().into(),
482 ])?,
483 db.prepare("UPDATE audit_entries SET workspace = ? WHERE workspace = ?")
484 .bind(&[to.as_str().into(), from.as_str().into()])?,
485 ])
486 .await?;
487 }
488 Ok(())
489}
490
491#[cfg(test)]
492mod tests {
493 use super::*;
494 use g1t_contracts::audit::{ActorKind, AuditOutcome};
495
Team plan, an open-source pool, monthly trials and honest metering; the sidebar for everyone; a workspace that stays put496 #[test]
497 fn entries_are_kept_as_long_as_the_longest_plan_reads_back() {
498 // 2026-10-05T00:00:00Z, a year back.
499 let now = 1_791_158_400_000;
500 assert_eq!(keep_from(now, 365), "2025-10-05T00:00:00.000Z");
Audit logs are kept by plan: a week on free, 90 days on the plan, and what staff set for an account in sudo501 assert_eq!(DEFAULT_MAX_DAYS, 400);
502 assert_eq!(DEFAULT_MIN_DAYS, 7);
503 }
504
505 #[test]
506 fn each_workspace_is_cut_off_at_its_own_days() {
507 // 2026-10-05T00:00:00Z.
508 let now = 1_791_158_400_000;
509 let kept = |workspace: &str, days: u32| AuditRetention {
510 workspace: workspace.into(),
511 days,
512 };
513 let retention = [
514 kept("free", 7),
515 kept("plan", 90),
516 kept("longer", 365),
517 kept("shorter", 1),
518 kept("past-the-most", 1_000),
519 ];
520 let expected = [
521 ("free", "2026-09-28T00:00:00.000Z"),
522 ("plan", "2026-07-07T00:00:00.000Z"),
523 ("longer", "2025-10-05T00:00:00.000Z"),
524 // Fewer days than the shortest are never looked for.
525 ("shorter", "2026-09-28T00:00:00.000Z"),
526 // More than the most are deleted by the ceiling anyway.
527 ("past-the-most", "2025-08-31T00:00:00.000Z"),
528 ];
529 let expected: Vec<(String, String)> = expected
530 .iter()
531 .map(|(w, t)| ((*w).to_owned(), (*t).to_owned()))
532 .collect();
533 assert_eq!(cutoffs(now, &retention, 7, 400), expected);
Team plan, an open-source pool, monthly trials and honest metering; the sidebar for everyone; a workspace that stays put534 }
535
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API536 fn args() -> ListAuditArgs {
537 ListAuditArgs {
538 workspace: "Acme".to_owned(),
539 visibility: AuditVisibility::All,
540 actor: None,
541 agent: None,
542 action: None,
543 repo: None,
544 number: None,
545 outcome: None,
546 actor_kind: None,
547 run_ids: vec![],
548 since: None,
549 until: None,
550 before: None,
551 limit: None,
552 }
553 }
554
555 #[test]
556 fn an_owner_sees_the_whole_workspace() {
557 let f = filter(&args());
558 assert_eq!(f.conditions, ["workspace = ?"]);
559 assert_eq!(f.values.len(), 1);
560 }
561
562 #[test]
563 fn a_member_sees_projects_and_their_own() {
564 let mut a = args();
565 a.visibility = AuditVisibility::Projects {
566 username: "ana".to_owned(),
567 };
568 let f = filter(&a);
569 assert!(f.conditions[1].contains("repo IS NOT NULL"));
570 assert_eq!(f.values.len(), 3);
571 }
572
573 #[test]
574 fn every_filter_binds_its_values() {
575 let mut a = args();
576 a.actor = Some("syntaqx".to_owned());
g1t is one name: its agent's work, commits and comments show as @g1t, and nobody can claim g1t or g1t-agent577 a.agent = Some("g1t".to_owned());
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API578 a.action = Some("git.push".to_owned());
579 a.repo = Some("acme/rocket".to_owned());
580 a.number = Some(4);
581 a.outcome = Some(AuditOutcome::Denied);
582 a.actor_kind = Some(ActorKind::Agent);
583 a.run_ids = vec!["run_1".to_owned(), "run_2".to_owned()];
584 a.since = Some("2026-10-01T00:00:00Z".to_owned());
585 a.until = Some("2026-10-05T00:00:00Z".to_owned());
586 a.before = Some("aud_9".to_owned());
587 let f = filter(&a);
588 let marks: usize = f.conditions.iter().map(|c| c.matches('?').count()).sum();
589 assert_eq!(marks, f.values.len());
590 assert!(f.conditions.iter().any(|c| c == "run_id IN (?, ?)"));
591 }
592}