flagon-io/g1t

public

Git for AI scale: a forge for thousands of agents working on the same code at once.

g1t/services/events/src/audit.rs

415 lines13,650 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};
11use g1t_contracts::events::{Event, WorkspaceRenamed};
12use g1t_contracts::new_id;
13use g1t_contracts::time::rfc3339;
14use g1t_kit::now_ms;
15use serde::Deserialize;
16use worker::wasm_bindgen::JsValue;
17use worker::{D1Database, Result};
18
19const DEFAULT_PAGE: u32 = 100;
20/// More than one request ever records.
21const MAX_BATCH: usize = 50;
22const MAX_TEXT: usize = 500;
23
24fn text(value: &Option<String>) -> JsValue {
25 value
26 .as_deref()
27 .map_or(JsValue::NULL, |value| JsValue::from(clip(value)))
28}
29
30fn clip(value: &str) -> String {
31 value.chars().take(MAX_TEXT).collect()
32}
33
34#[derive(Deserialize)]
35struct 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
60impl 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.
96pub 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)]
153struct 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)]
160enum Param {
161 Text(String),
162 Number(u32),
163}
164
165impl From<&str> for Param {
166 fn from(value: &str) -> Self {
167 Param::Text(value.to_owned())
168 }
169}
170
171impl From<String> for Param {
172 fn from(value: String) -> Self {
173 Param::Text(value)
174 }
175}
176
177impl From<u32> for Param {
178 fn from(value: u32) -> Self {
179 Param::Number(value)
180 }
181}
182
183impl 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
192impl 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
199fn 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`.
254pub 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 put280/// Entries are kept this many days unless `AUDIT_KEEP_DAYS` says otherwise:
281/// the longest any plan reads back (a year, on Team). Shorter windows,
282/// such as 30 days without Team, are applied where the log is read.
283pub const DEFAULT_KEEP_DAYS: u32 = 365;
284/// Rows removed per statement, so one purge never runs long.
285const PURGE_BATCH: u32 = 5_000;
286
287/// The oldest time an entry is kept from, `keep_days` before `now_ms`.
288pub fn keep_from(now_ms: u64, keep_days: u32) -> String {
289 g1t_contracts::time::rfc3339(now_ms.saturating_sub(u64::from(keep_days) * 24 * 60 * 60 * 1000))
290}
291
292/// Removes entries older than every plan keeps, a batch at a time, up to
293/// `rounds` batches. Returns how many went.
294pub async fn purge(db: &D1Database, before: &str, rounds: u32) -> Result<u32> {
295 let mut removed = 0;
296 for _ in 0..rounds {
297 let result = db
298 .prepare(
299 "DELETE FROM audit_entries WHERE id IN
300 (SELECT id FROM audit_entries WHERE time < ? ORDER BY time LIMIT ?)",
301 )
302 .bind(&[before.into(), PURGE_BATCH.into()])?
303 .run()
304 .await?;
305 let changed = result.meta()?.and_then(|meta| meta.changes).unwrap_or(0) as u32;
306 removed += changed;
307 if changed < PURGE_BATCH {
308 break;
309 }
310 }
311 Ok(removed)
312}
313
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API314/// Moves a renamed workspace's rows to its new slug.
315pub async fn follow_renames(db: &D1Database, events: &[Event]) -> Result<()> {
316 for event in events
317 .iter()
318 .filter(|event| event.kind == "workspace.renamed")
319 {
320 let Ok(renamed) = serde_json::from_value::<WorkspaceRenamed>(event.data.clone()) else {
321 continue;
322 };
323 let (from, to) = (renamed.from.to_lowercase(), renamed.to.to_lowercase());
324 if from == to {
325 continue;
326 }
327 db.batch(vec![
328 db.prepare(
329 "UPDATE audit_entries SET repo = ? || substr(repo, length(?) + 1)
330 WHERE workspace = ? AND repo LIKE ? || '/%'",
331 )
332 .bind(&[
333 to.as_str().into(),
334 from.as_str().into(),
335 from.as_str().into(),
336 from.as_str().into(),
337 ])?,
338 db.prepare("UPDATE audit_entries SET workspace = ? WHERE workspace = ?")
339 .bind(&[to.as_str().into(), from.as_str().into()])?,
340 ])
341 .await?;
342 }
343 Ok(())
344}
345
346#[cfg(test)]
347mod tests {
348 use super::*;
349 use g1t_contracts::audit::{ActorKind, AuditOutcome};
350
Team plan, an open-source pool, monthly trials and honest metering; the sidebar for everyone; a workspace that stays put351 #[test]
352 fn entries_are_kept_as_long_as_the_longest_plan_reads_back() {
353 // 2026-10-05T00:00:00Z, a year back.
354 let now = 1_791_158_400_000;
355 assert_eq!(keep_from(now, 365), "2025-10-05T00:00:00.000Z");
356 assert_eq!(DEFAULT_KEEP_DAYS, 365);
357 }
358
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API359 fn args() -> ListAuditArgs {
360 ListAuditArgs {
361 workspace: "Acme".to_owned(),
362 visibility: AuditVisibility::All,
363 actor: None,
364 agent: None,
365 action: None,
366 repo: None,
367 number: None,
368 outcome: None,
369 actor_kind: None,
370 run_ids: vec![],
371 since: None,
372 until: None,
373 before: None,
374 limit: None,
375 }
376 }
377
378 #[test]
379 fn an_owner_sees_the_whole_workspace() {
380 let f = filter(&args());
381 assert_eq!(f.conditions, ["workspace = ?"]);
382 assert_eq!(f.values.len(), 1);
383 }
384
385 #[test]
386 fn a_member_sees_projects_and_their_own() {
387 let mut a = args();
388 a.visibility = AuditVisibility::Projects {
389 username: "ana".to_owned(),
390 };
391 let f = filter(&a);
392 assert!(f.conditions[1].contains("repo IS NOT NULL"));
393 assert_eq!(f.values.len(), 3);
394 }
395
396 #[test]
397 fn every_filter_binds_its_values() {
398 let mut a = args();
399 a.actor = Some("syntaqx".to_owned());
400 a.agent = Some("g1t-agent".to_owned());
401 a.action = Some("git.push".to_owned());
402 a.repo = Some("acme/rocket".to_owned());
403 a.number = Some(4);
404 a.outcome = Some(AuditOutcome::Denied);
405 a.actor_kind = Some(ActorKind::Agent);
406 a.run_ids = vec!["run_1".to_owned(), "run_2".to_owned()];
407 a.since = Some("2026-10-01T00:00:00Z".to_owned());
408 a.until = Some("2026-10-05T00:00:00Z".to_owned());
409 a.before = Some("aud_9".to_owned());
410 let f = filter(&a);
411 let marks: usize = f.conditions.iter().map(|c| c.matches('?').count()).sum();
412 assert_eq!(marks, f.values.len());
413 assert!(f.conditions.iter().any(|c| c == "run_id IN (?, ?)"));
414 }
415}