1use std::path::Path;
21
22use serde::Serialize;
23use sqlx::QueryBuilder;
24use sqlx::{Pool, Sqlite};
25
26use crate::xrpc_gateway::handlers::projections::audit_event::{
27 AuditEventRow, OZONE_ELIGIBLE_ACTIONS, ProjectedModEvent, project_audit_event,
28};
29
30use super::error::CliError;
31
32#[derive(Debug, Clone, Default)]
36pub struct EventsInput {
37 pub subject: Option<String>,
40 pub actor: Option<String>,
43 pub action_type: Option<String>,
49 pub from: Option<String>,
51 pub to: Option<String>,
53 pub limit: Option<u32>,
55 pub cursor: Option<String>,
60 pub ozone_only: bool,
65}
66
67const MAX_LIMIT: u32 = 250;
68const DEFAULT_LIMIT: u32 = 50;
69
70#[derive(Debug, Clone, Serialize)]
74#[serde(tag = "shape")]
75pub enum EventRow {
76 #[serde(rename = "ozone")]
80 Ozone(OzoneEventRow),
81 #[serde(rename = "internal")]
84 Internal(InternalEventRow),
85}
86
87#[derive(Debug, Clone, Serialize)]
92pub struct OzoneEventRow {
93 pub id: i64,
95 pub event: serde_json::Value,
98 pub subject: serde_json::Value,
100 pub created_by: String,
102 pub created_at: String,
104}
105
106#[derive(Debug, Clone, Serialize)]
110pub struct InternalEventRow {
111 pub id: i64,
113 pub action: String,
115 pub actor_did: String,
117 #[serde(skip_serializing_if = "Option::is_none")]
120 pub target: Option<String>,
121 pub outcome: String,
123 pub created_at: String,
125}
126
127#[derive(Debug, Clone, Serialize)]
130pub struct EventsResponse {
131 #[serde(skip_serializing_if = "Option::is_none")]
133 pub cursor: Option<String>,
134 pub events: Vec<EventRow>,
136}
137
138pub async fn list(pool: &Pool<Sqlite>, input: EventsInput) -> Result<EventsResponse, CliError> {
141 let parsed = parse_input(input)?;
142 let raw = fetch_page(pool, &parsed).await?;
143 let has_more = raw.len() > parsed.limit as usize;
144 let surfaced: Vec<RawAuditRow> = raw.into_iter().take(parsed.limit as usize).collect();
145
146 let mut events = Vec::with_capacity(surfaced.len());
147 let mut last_id: Option<i64> = None;
148 for r in &surfaced {
149 last_id = Some(r.audit_id);
150
151 if parsed.ozone_only {
152 let aer = r.to_audit_event_row();
156 if let Some(projected) = project_audit_event(&aer) {
157 events.push(EventRow::Ozone(projected_to_row(projected)));
158 }
159 continue;
160 }
161
162 if OZONE_ELIGIBLE_ACTIONS.contains(&r.audit_action.as_str()) {
165 let aer = r.to_audit_event_row();
166 if let Some(projected) = project_audit_event(&aer) {
167 events.push(EventRow::Ozone(projected_to_row(projected)));
168 continue;
169 }
170 }
176 events.push(EventRow::Internal(InternalEventRow {
177 id: r.audit_id,
178 action: r.audit_action.clone(),
179 actor_did: r.actor_did.clone(),
180 target: r.target.clone(),
181 outcome: r.outcome.clone(),
182 created_at: epoch_ms_to_rfc3339(r.created_at),
183 }));
184 }
185
186 let cursor = if has_more {
187 last_id.map(encode_cursor)
188 } else {
189 None
190 };
191
192 Ok(EventsResponse { cursor, events })
193}
194
195#[derive(Debug)]
200struct ParsedInput {
201 subject_did: Option<String>,
202 subject_uri: Option<String>,
203 actor: Option<String>,
204 action_type: Option<String>,
205 from_ms: Option<i64>,
206 to_ms: Option<i64>,
207 limit: u32,
208 cursor_id: Option<i64>,
209 ozone_only: bool,
210}
211
212fn parse_input(i: EventsInput) -> Result<ParsedInput, CliError> {
213 let limit = match i.limit {
214 None => DEFAULT_LIMIT,
215 Some(0) => return Err(CliError::Config("limit must be at least 1".into())),
216 Some(n) if n > MAX_LIMIT => {
217 return Err(CliError::Config(format!(
218 "limit {n} exceeds maximum {MAX_LIMIT}"
219 )));
220 }
221 Some(n) => n,
222 };
223
224 let cursor_id = match i.cursor.as_deref() {
225 None => None,
226 Some(s) => Some(decode_cursor(s)?),
227 };
228
229 let (subject_did, subject_uri) = match i.subject.as_deref() {
230 None => (None, None),
231 Some(s) if s.starts_with("at://") => {
232 let did = s
233 .strip_prefix("at://")
234 .and_then(|rest| rest.split('/').next())
235 .filter(|d| d.starts_with("did:"))
236 .ok_or_else(|| {
237 CliError::Config(format!("subject AT-URI {s:?} missing DID authority"))
238 })?
239 .to_string();
240 (Some(did), Some(s.to_string()))
241 }
242 Some(s) if s.starts_with("did:") => (Some(s.to_string()), None),
243 Some(s) => {
244 return Err(CliError::Config(format!(
245 "subject {s:?} is neither a DID nor an AT-URI"
246 )));
247 }
248 };
249
250 let from_ms = match i.from.as_deref() {
251 None => None,
252 Some(s) => Some(parse_rfc3339_to_ms(s)?),
253 };
254 let to_ms = match i.to.as_deref() {
255 None => None,
256 Some(s) => Some(parse_rfc3339_to_ms(s)?),
257 };
258
259 Ok(ParsedInput {
260 subject_did,
261 subject_uri,
262 actor: i.actor,
263 action_type: i.action_type,
264 from_ms,
265 to_ms,
266 limit,
267 cursor_id,
268 ozone_only: i.ozone_only,
269 })
270}
271
272fn parse_rfc3339_to_ms(s: &str) -> Result<i64, CliError> {
273 use time::OffsetDateTime;
274 use time::format_description::well_known::Rfc3339;
275 let dt = OffsetDateTime::parse(s, &Rfc3339)
276 .map_err(|_| CliError::Config(format!("malformed RFC-3339 timestamp: {s:?}")))?;
277 Ok((dt.unix_timestamp_nanos() / 1_000_000) as i64)
278}
279
280#[derive(serde::Serialize, serde::Deserialize)]
285struct CursorBody {
286 cursor_id: i64,
287}
288
289fn encode_cursor(id: i64) -> String {
290 use base64::Engine as _;
291 let json = serde_json::to_vec(&CursorBody { cursor_id: id }).expect("cursor serializes");
292 base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(json)
293}
294
295fn decode_cursor(s: &str) -> Result<i64, CliError> {
296 use base64::Engine as _;
297 let bytes = base64::engine::general_purpose::URL_SAFE_NO_PAD
298 .decode(s)
299 .map_err(|e| CliError::Config(format!("malformed cursor (base64): {e}")))?;
300 let body: CursorBody = serde_json::from_slice(&bytes)
301 .map_err(|e| CliError::Config(format!("malformed cursor (json): {e}")))?;
302 Ok(body.cursor_id)
303}
304
305#[derive(Debug, Clone)]
313struct RawAuditRow {
314 audit_id: i64,
315 audit_action: String,
316 actor_did: String,
317 target: Option<String>,
318 outcome: String,
319 created_at: i64,
320 subject_action_type: Option<String>,
321 subject_did: Option<String>,
322 subject_uri: Option<String>,
323 subject_notes: Option<String>,
324 subject_duration: Option<String>,
325 subject_reason_codes: Option<String>,
326 audit_reason: Option<String>,
327}
328
329impl RawAuditRow {
330 fn to_audit_event_row(&self) -> AuditEventRow {
331 AuditEventRow {
332 audit_id: self.audit_id,
333 audit_action: self.audit_action.clone(),
334 actor_did: self.actor_did.clone(),
335 created_at: self.created_at,
336 subject_action_type: self.subject_action_type.clone(),
337 subject_did: self.subject_did.clone(),
338 subject_uri: self.subject_uri.clone(),
339 subject_notes: self.subject_notes.clone(),
340 subject_duration: self.subject_duration.clone(),
341 subject_reason_codes_json: self.subject_reason_codes.clone(),
342 audit_reason_json: self.audit_reason.clone(),
343 }
344 }
345}
346
347async fn fetch_page(
348 pool: &Pool<Sqlite>,
349 parsed: &ParsedInput,
350) -> Result<Vec<RawAuditRow>, CliError> {
351 let mut qb: QueryBuilder<Sqlite> = QueryBuilder::new(
359 r#"SELECT
360 a.id AS audit_id,
361 a.action AS audit_action,
362 a.actor_did AS actor_did,
363 a.target AS target,
364 a.outcome AS outcome,
365 a.created_at AS created_at,
366 sa.action_type AS subject_action_type,
367 sa.subject_did AS subject_did,
368 sa.subject_uri AS subject_uri,
369 sa.notes AS subject_notes,
370 sa.duration AS subject_duration,
371 sa.reason_codes AS subject_reason_codes,
372 a.reason AS audit_reason
373 FROM audit_log a
374 LEFT JOIN subject_actions sa ON
375 (a.action = 'subject_action_recorded' AND sa.audit_log_id = a.id)
376 OR (a.action = 'subject_action_revoked' AND sa.id = CAST(a.target AS INTEGER))
377 WHERE 1=1"#,
378 );
379
380 if parsed.ozone_only {
381 qb.push(" AND a.action IN (");
382 let mut sep = qb.separated(", ");
383 for action in OZONE_ELIGIBLE_ACTIONS {
384 sep.push_bind(*action);
385 }
386 qb.push(")");
387 }
388
389 if let Some(action) = &parsed.action_type {
390 qb.push(" AND a.action = ").push_bind(action.clone());
391 }
392 if let Some(actor) = &parsed.actor {
393 qb.push(" AND a.actor_did = ").push_bind(actor.clone());
394 }
395 if let Some(did) = &parsed.subject_did {
396 qb.push(" AND sa.subject_did = ").push_bind(did.clone());
397 if parsed.subject_uri.is_none() {
398 qb.push(" AND sa.subject_uri IS NULL");
399 }
400 }
401 if let Some(uri) = &parsed.subject_uri {
402 qb.push(" AND sa.subject_uri = ").push_bind(uri.clone());
403 }
404 if let Some(from) = parsed.from_ms {
405 qb.push(" AND a.created_at >= ").push_bind(from);
406 }
407 if let Some(to) = parsed.to_ms {
408 qb.push(" AND a.created_at <= ").push_bind(to);
409 }
410 if let Some(c) = parsed.cursor_id {
411 qb.push(" AND a.id < ").push_bind(c);
414 }
415
416 qb.push(" ORDER BY a.id DESC LIMIT ")
417 .push_bind((parsed.limit + 1) as i64);
418
419 qb.build_query_as::<(
420 i64,
421 String,
422 String,
423 Option<String>,
424 String,
425 i64,
426 Option<String>,
427 Option<String>,
428 Option<String>,
429 Option<String>,
430 Option<String>,
431 Option<String>,
432 Option<String>,
433 )>()
434 .fetch_all(pool)
435 .await
436 .map(|rows| {
437 rows.into_iter()
438 .map(
439 |(
440 audit_id,
441 audit_action,
442 actor_did,
443 target,
444 outcome,
445 created_at,
446 subject_action_type,
447 subject_did,
448 subject_uri,
449 subject_notes,
450 subject_duration,
451 subject_reason_codes,
452 audit_reason,
453 )| RawAuditRow {
454 audit_id,
455 audit_action,
456 actor_did,
457 target,
458 outcome,
459 created_at,
460 subject_action_type,
461 subject_did,
462 subject_uri,
463 subject_notes,
464 subject_duration,
465 subject_reason_codes,
466 audit_reason,
467 },
468 )
469 .collect()
470 })
471 .map_err(|e| CliError::Startup(format!("audit_log query: {e}")))
472}
473
474fn projected_to_row(p: ProjectedModEvent) -> OzoneEventRow {
479 OzoneEventRow {
480 id: p.id,
481 event: p.event,
482 subject: p.subject,
483 created_by: p.created_by,
484 created_at: epoch_ms_to_rfc3339(p.created_at),
485 }
486}
487
488fn epoch_ms_to_rfc3339(ms: i64) -> String {
489 crate::writer::rfc3339_from_epoch_ms(ms)
490 .unwrap_or_else(|_| String::from("1970-01-01T00:00:00.000Z"))
491}
492
493pub fn format_human(resp: &EventsResponse) -> String {
500 use std::fmt::Write;
501 if resp.events.is_empty() {
502 let mut s = String::from("(no events)");
503 if let Some(c) = &resp.cursor {
504 let _ = write!(s, "\nnext cursor: {c}");
505 }
506 return s;
507 }
508
509 let mut out = String::new();
510 for ev in &resp.events {
511 match ev {
512 EventRow::Ozone(o) => {
513 let ty = o.event["$type"].as_str().unwrap_or("?");
514 let subj = o.subject["did"]
515 .as_str()
516 .or_else(|| o.subject["uri"].as_str())
517 .unwrap_or("?");
518 let _ = writeln!(
519 out,
520 "ozone\t{}\t{}\t{}\t{}\tby={}",
521 o.id, o.created_at, ty, subj, o.created_by,
522 );
523 }
524 EventRow::Internal(i) => {
525 let target = i.target.as_deref().unwrap_or("-");
526 let _ = writeln!(
527 out,
528 "internal\t{}\t{}\t{}\t{}\tby={}\toutcome={}",
529 i.id, i.created_at, i.action, target, i.actor_did, i.outcome,
530 );
531 }
532 }
533 }
534 if let Some(c) = &resp.cursor {
535 let _ = write!(out, "next cursor: {c}");
536 } else if out.ends_with('\n') {
537 out.pop();
538 }
539 out
540}
541
542pub fn format_json(resp: &EventsResponse) -> String {
544 serde_json::to_string(resp).expect("EventsResponse serializes")
545}
546
547#[allow(dead_code)]
552fn _imports_used(_p: &Path) {}
553
554#[cfg(test)]
555mod tests {
556 use super::*;
557
558 fn empty_input() -> EventsInput {
559 EventsInput::default()
560 }
561
562 #[test]
565 fn default_limit_is_50() {
566 let p = parse_input(empty_input()).unwrap();
567 assert_eq!(p.limit, 50);
568 }
569
570 #[test]
571 fn limit_capped_at_250() {
572 let mut i = empty_input();
573 i.limit = Some(500);
574 let err = parse_input(i).unwrap_err();
575 assert!(err.to_string().contains("250"));
576 }
577
578 #[test]
579 fn limit_zero_rejected() {
580 let mut i = empty_input();
581 i.limit = Some(0);
582 let err = parse_input(i).unwrap_err();
583 assert!(err.to_string().contains("at least 1"));
584 }
585
586 #[test]
587 fn account_level_subject_parses_to_did() {
588 let mut i = empty_input();
589 i.subject = Some("did:plc:abc".into());
590 let p = parse_input(i).unwrap();
591 assert_eq!(p.subject_did.as_deref(), Some("did:plc:abc"));
592 assert!(p.subject_uri.is_none());
593 }
594
595 #[test]
596 fn record_level_subject_parses_to_did_plus_uri() {
597 let mut i = empty_input();
598 i.subject = Some("at://did:plc:author/c/r".into());
599 let p = parse_input(i).unwrap();
600 assert_eq!(p.subject_did.as_deref(), Some("did:plc:author"));
601 assert_eq!(p.subject_uri.as_deref(), Some("at://did:plc:author/c/r"));
602 }
603
604 #[test]
605 fn malformed_subject_rejected() {
606 let mut i = empty_input();
607 i.subject = Some("not-a-did".into());
608 assert!(parse_input(i).is_err());
609 }
610
611 #[test]
612 fn malformed_rfc3339_rejected() {
613 let mut i = empty_input();
614 i.from = Some("not-a-date".into());
615 assert!(parse_input(i).is_err());
616 }
617
618 #[test]
619 fn cursor_round_trip() {
620 let s = encode_cursor(42);
621 let id = decode_cursor(&s).unwrap();
622 assert_eq!(id, 42);
623 }
624
625 #[test]
626 fn malformed_cursor_rejected() {
627 assert!(decode_cursor("not-base64-!!!").is_err());
628 }
629
630 #[test]
633 fn format_human_empty() {
634 let resp = EventsResponse {
635 cursor: None,
636 events: Vec::new(),
637 };
638 assert_eq!(format_human(&resp), "(no events)");
639 }
640
641 #[test]
642 fn format_human_with_cursor_appends_cursor_line() {
643 let resp = EventsResponse {
644 cursor: Some("c".into()),
645 events: Vec::new(),
646 };
647 let out = format_human(&resp);
648 assert!(out.contains("next cursor: c"));
649 }
650
651 #[test]
652 fn format_human_renders_internal_row() {
653 let resp = EventsResponse {
654 cursor: None,
655 events: vec![EventRow::Internal(InternalEventRow {
656 id: 5,
657 action: "retention_sweep".into(),
658 actor_did: "did:plc:m".into(),
659 target: None,
660 outcome: "success".into(),
661 created_at: "2026-04-29T00:00:00.000Z".into(),
662 })],
663 };
664 let out = format_human(&resp);
665 assert!(out.contains("internal"));
666 assert!(out.contains("retention_sweep"));
667 assert!(out.contains("did:plc:m"));
668 }
669
670 #[test]
671 fn format_human_renders_ozone_row() {
672 let resp = EventsResponse {
673 cursor: None,
674 events: vec![EventRow::Ozone(OzoneEventRow {
675 id: 7,
676 event: serde_json::json!({
677 "$type": "tools.ozone.moderation.defs#modEventLabel",
678 "createLabelVals": ["spam"],
679 }),
680 subject: serde_json::json!({
681 "$type": "com.atproto.admin.defs#repoRef",
682 "did": "did:plc:t",
683 }),
684 created_by: "did:plc:m".into(),
685 created_at: "2026-04-29T00:00:00.000Z".into(),
686 })],
687 };
688 let out = format_human(&resp);
689 assert!(out.contains("ozone"));
690 assert!(out.contains("modEventLabel"));
691 assert!(out.contains("did:plc:t"));
692 }
693
694 #[test]
695 fn format_json_round_trip() {
696 let resp = EventsResponse {
697 cursor: Some("c".into()),
698 events: vec![EventRow::Internal(InternalEventRow {
699 id: 1,
700 action: "report_resolved".into(),
701 actor_did: "did:plc:m".into(),
702 target: Some("42".into()),
703 outcome: "success".into(),
704 created_at: "2026-04-29T00:00:00.000Z".into(),
705 })],
706 };
707 let s = format_json(&resp);
708 let v: serde_json::Value = serde_json::from_str(&s).unwrap();
709 assert_eq!(v["cursor"].as_str(), Some("c"));
710 assert_eq!(v["events"][0]["shape"].as_str(), Some("internal"));
711 assert_eq!(v["events"][0]["action"].as_str(), Some("report_resolved"));
712 }
713}