Skip to main content

oximedia_distributed/
audit_log.rs

1#![allow(dead_code)]
2//! Audit logging for coordinator state changes.
3//!
4//! Records every significant state mutation in the distributed coordinator:
5//! job submissions, cancellations, worker joins/leaves, leader elections,
6//! configuration changes, and snapshot operations.
7//!
8//! The audit log is append-only, indexed by monotonic sequence number, and
9//! supports querying by time range, event type, and actor.
10
11use std::collections::VecDeque;
12use std::fmt;
13use std::sync::atomic::{AtomicU64, Ordering};
14use std::time::{Duration, SystemTime, UNIX_EPOCH};
15
16use uuid::Uuid;
17
18// ---------------------------------------------------------------------------
19// AuditEventKind
20// ---------------------------------------------------------------------------
21
22/// Categories of auditable events.
23#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
24pub enum AuditEventKind {
25    /// A new job was submitted.
26    JobSubmitted,
27    /// A job was cancelled.
28    JobCancelled,
29    /// A job completed successfully.
30    JobCompleted,
31    /// A job failed.
32    JobFailed,
33    /// A job was preempted by a higher-priority job.
34    JobPreempted,
35    /// A job was reassigned to a different worker.
36    JobReassigned,
37    /// A worker joined the cluster.
38    WorkerJoined,
39    /// A worker left (or was removed from) the cluster.
40    WorkerLeft,
41    /// A leader election completed.
42    LeaderElected,
43    /// A snapshot was created.
44    SnapshotCreated,
45    /// A snapshot was restored.
46    SnapshotRestored,
47    /// Configuration was changed.
48    ConfigChanged,
49    /// A segment merge was completed.
50    SegmentMerged,
51    /// A custom/application-defined event.
52    Custom,
53}
54
55impl fmt::Display for AuditEventKind {
56    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
57        match self {
58            Self::JobSubmitted => write!(f, "JobSubmitted"),
59            Self::JobCancelled => write!(f, "JobCancelled"),
60            Self::JobCompleted => write!(f, "JobCompleted"),
61            Self::JobFailed => write!(f, "JobFailed"),
62            Self::JobPreempted => write!(f, "JobPreempted"),
63            Self::JobReassigned => write!(f, "JobReassigned"),
64            Self::WorkerJoined => write!(f, "WorkerJoined"),
65            Self::WorkerLeft => write!(f, "WorkerLeft"),
66            Self::LeaderElected => write!(f, "LeaderElected"),
67            Self::SnapshotCreated => write!(f, "SnapshotCreated"),
68            Self::SnapshotRestored => write!(f, "SnapshotRestored"),
69            Self::ConfigChanged => write!(f, "ConfigChanged"),
70            Self::SegmentMerged => write!(f, "SegmentMerged"),
71            Self::Custom => write!(f, "Custom"),
72        }
73    }
74}
75
76// ---------------------------------------------------------------------------
77// AuditSeverity
78// ---------------------------------------------------------------------------
79
80/// Severity level for audit events.
81#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)]
82pub enum AuditSeverity {
83    /// Informational — normal operations.
84    Info,
85    /// Warning — unexpected but non-critical.
86    Warning,
87    /// Error — something failed.
88    Error,
89    /// Critical — system integrity may be at risk.
90    Critical,
91}
92
93impl fmt::Display for AuditSeverity {
94    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
95        match self {
96            Self::Info => write!(f, "INFO"),
97            Self::Warning => write!(f, "WARN"),
98            Self::Error => write!(f, "ERROR"),
99            Self::Critical => write!(f, "CRITICAL"),
100        }
101    }
102}
103
104// ---------------------------------------------------------------------------
105// AuditEntry
106// ---------------------------------------------------------------------------
107
108/// A single entry in the audit log.
109#[derive(Debug, Clone)]
110pub struct AuditEntry {
111    /// Monotonically increasing sequence number.
112    pub sequence: u64,
113    /// Wall-clock timestamp (microseconds since Unix epoch).
114    pub timestamp_us: u64,
115    /// Kind of event.
116    pub kind: AuditEventKind,
117    /// Severity level.
118    pub severity: AuditSeverity,
119    /// Actor that caused the event (e.g., worker ID, coordinator ID, user).
120    pub actor: String,
121    /// Optional target entity (e.g., job ID, worker ID).
122    pub target: Option<String>,
123    /// Human-readable description of the event.
124    pub message: String,
125    /// Optional structured metadata as key-value pairs.
126    pub metadata: Vec<(String, String)>,
127}
128
129impl fmt::Display for AuditEntry {
130    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
131        write!(
132            f,
133            "[#{} {} {} {}] actor={} {}",
134            self.sequence, self.timestamp_us, self.severity, self.kind, self.actor, self.message
135        )
136    }
137}
138
139// ---------------------------------------------------------------------------
140// AuditEntryBuilder
141// ---------------------------------------------------------------------------
142
143/// Builder for constructing [`AuditEntry`] instances.
144pub struct AuditEntryBuilder {
145    kind: AuditEventKind,
146    severity: AuditSeverity,
147    actor: String,
148    target: Option<String>,
149    message: String,
150    metadata: Vec<(String, String)>,
151}
152
153impl AuditEntryBuilder {
154    /// Create a new builder with the given event kind and actor.
155    pub fn new(kind: AuditEventKind, actor: impl Into<String>) -> Self {
156        Self {
157            kind,
158            severity: AuditSeverity::Info,
159            actor: actor.into(),
160            target: None,
161            message: String::new(),
162            metadata: Vec::new(),
163        }
164    }
165
166    /// Set the severity level.
167    pub fn severity(mut self, severity: AuditSeverity) -> Self {
168        self.severity = severity;
169        self
170    }
171
172    /// Set the target entity.
173    pub fn target(mut self, target: impl Into<String>) -> Self {
174        self.target = Some(target.into());
175        self
176    }
177
178    /// Set the message.
179    pub fn message(mut self, message: impl Into<String>) -> Self {
180        self.message = message.into();
181        self
182    }
183
184    /// Add a metadata key-value pair.
185    pub fn meta(mut self, key: impl Into<String>, value: impl Into<String>) -> Self {
186        self.metadata.push((key.into(), value.into()));
187        self
188    }
189
190    /// Build the entry. The sequence number and timestamp will be set by the
191    /// [`AuditLog`] when the entry is appended.
192    fn build(self, sequence: u64, timestamp_us: u64) -> AuditEntry {
193        AuditEntry {
194            sequence,
195            timestamp_us,
196            kind: self.kind,
197            severity: self.severity,
198            actor: self.actor,
199            target: self.target,
200            message: self.message,
201            metadata: self.metadata,
202        }
203    }
204}
205
206// ---------------------------------------------------------------------------
207// AuditLogConfig
208// ---------------------------------------------------------------------------
209
210/// Configuration for the audit log.
211#[derive(Debug, Clone)]
212pub struct AuditLogConfig {
213    /// Maximum number of entries to retain in memory. When exceeded, the
214    /// oldest entries are discarded (ring buffer behaviour).
215    pub max_entries: usize,
216    /// Minimum severity level to record. Events below this level are dropped.
217    pub min_severity: AuditSeverity,
218}
219
220impl Default for AuditLogConfig {
221    fn default() -> Self {
222        Self {
223            max_entries: 100_000,
224            min_severity: AuditSeverity::Info,
225        }
226    }
227}
228
229// ---------------------------------------------------------------------------
230// AuditLog
231// ---------------------------------------------------------------------------
232
233/// An append-only, bounded audit log for coordinator state changes.
234///
235/// Thread-safety note: this struct is *not* internally synchronised. In a
236/// multi-threaded context, wrap it in `Arc<Mutex<AuditLog>>` or similar.
237pub struct AuditLog {
238    config: AuditLogConfig,
239    /// Monotonic sequence counter.
240    next_sequence: AtomicU64,
241    /// Ring buffer of entries.
242    entries: VecDeque<AuditEntry>,
243    /// Total number of entries ever appended (including evicted ones).
244    total_appended: u64,
245}
246
247impl AuditLog {
248    /// Create a new audit log with the given configuration.
249    pub fn new(config: AuditLogConfig) -> Self {
250        Self {
251            next_sequence: AtomicU64::new(1),
252            entries: VecDeque::with_capacity(config.max_entries.min(4096)),
253            total_appended: 0,
254            config,
255        }
256    }
257
258    /// Create a new audit log with default configuration.
259    pub fn with_defaults() -> Self {
260        Self::new(AuditLogConfig::default())
261    }
262
263    /// Append an entry built from a builder.
264    ///
265    /// Returns the assigned sequence number, or `None` if the event was
266    /// filtered out by the minimum severity setting.
267    pub fn append(&mut self, builder: AuditEntryBuilder) -> Option<u64> {
268        if builder.severity < self.config.min_severity {
269            return None;
270        }
271
272        let seq = self.next_sequence.fetch_add(1, Ordering::Relaxed);
273        let timestamp_us = SystemTime::now()
274            .duration_since(UNIX_EPOCH)
275            .unwrap_or(Duration::ZERO)
276            .as_micros() as u64;
277
278        let entry = builder.build(seq, timestamp_us);
279
280        // Evict oldest if at capacity.
281        if self.entries.len() >= self.config.max_entries {
282            self.entries.pop_front();
283        }
284
285        self.entries.push_back(entry);
286        self.total_appended += 1;
287        Some(seq)
288    }
289
290    /// Convenience: log a job submission.
291    pub fn log_job_submitted(&mut self, job_id: Uuid, actor: &str, codec: &str) -> Option<u64> {
292        self.append(
293            AuditEntryBuilder::new(AuditEventKind::JobSubmitted, actor)
294                .target(job_id.to_string())
295                .message(format!("Job {job_id} submitted"))
296                .meta("codec", codec),
297        )
298    }
299
300    /// Convenience: log a job cancellation.
301    pub fn log_job_cancelled(&mut self, job_id: Uuid, actor: &str, reason: &str) -> Option<u64> {
302        self.append(
303            AuditEntryBuilder::new(AuditEventKind::JobCancelled, actor)
304                .severity(AuditSeverity::Warning)
305                .target(job_id.to_string())
306                .message(format!("Job {job_id} cancelled: {reason}")),
307        )
308    }
309
310    /// Convenience: log a job failure.
311    pub fn log_job_failed(&mut self, job_id: Uuid, actor: &str, error: &str) -> Option<u64> {
312        self.append(
313            AuditEntryBuilder::new(AuditEventKind::JobFailed, actor)
314                .severity(AuditSeverity::Error)
315                .target(job_id.to_string())
316                .message(format!("Job {job_id} failed: {error}")),
317        )
318    }
319
320    /// Convenience: log a worker joining the cluster.
321    pub fn log_worker_joined(&mut self, worker_id: &str, addr: &str) -> Option<u64> {
322        self.append(
323            AuditEntryBuilder::new(AuditEventKind::WorkerJoined, "coordinator")
324                .target(worker_id)
325                .message(format!("Worker {worker_id} joined from {addr}"))
326                .meta("address", addr),
327        )
328    }
329
330    /// Convenience: log a worker leaving the cluster.
331    pub fn log_worker_left(&mut self, worker_id: &str, reason: &str) -> Option<u64> {
332        self.append(
333            AuditEntryBuilder::new(AuditEventKind::WorkerLeft, "coordinator")
334                .severity(AuditSeverity::Warning)
335                .target(worker_id)
336                .message(format!("Worker {worker_id} left: {reason}")),
337        )
338    }
339
340    /// Convenience: log a leader election.
341    pub fn log_leader_elected(&mut self, leader_id: &str, term: u64) -> Option<u64> {
342        self.append(
343            AuditEntryBuilder::new(AuditEventKind::LeaderElected, leader_id)
344                .message(format!("Node {leader_id} elected leader for term {term}"))
345                .meta("term", term.to_string()),
346        )
347    }
348
349    /// Number of entries currently in the log.
350    pub fn len(&self) -> usize {
351        self.entries.len()
352    }
353
354    /// Whether the log is empty.
355    pub fn is_empty(&self) -> bool {
356        self.entries.is_empty()
357    }
358
359    /// Total entries ever appended (including evicted).
360    pub fn total_appended(&self) -> u64 {
361        self.total_appended
362    }
363
364    /// Get an entry by sequence number.
365    pub fn get_by_sequence(&self, seq: u64) -> Option<&AuditEntry> {
366        self.entries.iter().find(|e| e.sequence == seq)
367    }
368
369    /// Return the most recent `n` entries (newest last).
370    pub fn recent(&self, n: usize) -> Vec<&AuditEntry> {
371        let start = self.entries.len().saturating_sub(n);
372        self.entries.iter().skip(start).collect()
373    }
374
375    /// Query entries by event kind.
376    pub fn query_by_kind(&self, kind: AuditEventKind) -> Vec<&AuditEntry> {
377        self.entries.iter().filter(|e| e.kind == kind).collect()
378    }
379
380    /// Query entries by actor.
381    pub fn query_by_actor(&self, actor: &str) -> Vec<&AuditEntry> {
382        self.entries.iter().filter(|e| e.actor == actor).collect()
383    }
384
385    /// Query entries by target.
386    pub fn query_by_target(&self, target: &str) -> Vec<&AuditEntry> {
387        self.entries
388            .iter()
389            .filter(|e| e.target.as_deref() == Some(target))
390            .collect()
391    }
392
393    /// Query entries by severity at or above the given level.
394    pub fn query_by_min_severity(&self, min: AuditSeverity) -> Vec<&AuditEntry> {
395        self.entries.iter().filter(|e| e.severity >= min).collect()
396    }
397
398    /// Query entries within a time range (microseconds since Unix epoch).
399    pub fn query_by_time_range(&self, start_us: u64, end_us: u64) -> Vec<&AuditEntry> {
400        self.entries
401            .iter()
402            .filter(|e| e.timestamp_us >= start_us && e.timestamp_us <= end_us)
403            .collect()
404    }
405
406    /// Clear all entries.
407    pub fn clear(&mut self) {
408        self.entries.clear();
409    }
410
411    /// Export all entries as a vector (for serialization or transfer).
412    pub fn export_all(&self) -> Vec<AuditEntry> {
413        self.entries.iter().cloned().collect()
414    }
415}
416
417// ---------------------------------------------------------------------------
418// Tests
419// ---------------------------------------------------------------------------
420
421#[cfg(test)]
422mod tests {
423    use super::*;
424
425    #[test]
426    fn test_append_and_len() {
427        let mut log = AuditLog::with_defaults();
428        assert!(log.is_empty());
429
430        let seq = log.append(
431            AuditEntryBuilder::new(AuditEventKind::JobSubmitted, "user1").message("test job"),
432        );
433        assert!(seq.is_some());
434        assert_eq!(log.len(), 1);
435        assert!(!log.is_empty());
436    }
437
438    #[test]
439    fn test_sequence_numbers_are_monotonic() {
440        let mut log = AuditLog::with_defaults();
441        let s1 = log
442            .append(AuditEntryBuilder::new(AuditEventKind::JobSubmitted, "a").message("1"))
443            .expect("append");
444        let s2 = log
445            .append(AuditEntryBuilder::new(AuditEventKind::JobCancelled, "a").message("2"))
446            .expect("append");
447        let s3 = log
448            .append(AuditEntryBuilder::new(AuditEventKind::JobCompleted, "a").message("3"))
449            .expect("append");
450        assert!(s1 < s2);
451        assert!(s2 < s3);
452    }
453
454    #[test]
455    fn test_max_entries_eviction() {
456        let mut log = AuditLog::new(AuditLogConfig {
457            max_entries: 3,
458            min_severity: AuditSeverity::Info,
459        });
460
461        for i in 0..5 {
462            log.append(
463                AuditEntryBuilder::new(AuditEventKind::Custom, "actor")
464                    .message(format!("event {i}")),
465            );
466        }
467
468        assert_eq!(log.len(), 3);
469        assert_eq!(log.total_appended(), 5);
470        // Oldest two should have been evicted
471        let entries = log.export_all();
472        assert!(entries[0].message.contains("event 2"));
473    }
474
475    #[test]
476    fn test_min_severity_filter() {
477        let mut log = AuditLog::new(AuditLogConfig {
478            max_entries: 100,
479            min_severity: AuditSeverity::Warning,
480        });
481
482        // Info event should be filtered out
483        let seq = log.append(
484            AuditEntryBuilder::new(AuditEventKind::JobSubmitted, "user")
485                .severity(AuditSeverity::Info)
486                .message("should be dropped"),
487        );
488        assert!(seq.is_none());
489        assert_eq!(log.len(), 0);
490
491        // Warning event should be accepted
492        let seq = log.append(
493            AuditEntryBuilder::new(AuditEventKind::JobCancelled, "user")
494                .severity(AuditSeverity::Warning)
495                .message("should be kept"),
496        );
497        assert!(seq.is_some());
498        assert_eq!(log.len(), 1);
499    }
500
501    #[test]
502    fn test_query_by_kind() {
503        let mut log = AuditLog::with_defaults();
504        log.log_job_submitted(Uuid::new_v4(), "user", "av1");
505        log.log_job_cancelled(Uuid::new_v4(), "user", "timeout");
506        log.log_job_submitted(Uuid::new_v4(), "user", "vp9");
507
508        let submitted = log.query_by_kind(AuditEventKind::JobSubmitted);
509        assert_eq!(submitted.len(), 2);
510
511        let cancelled = log.query_by_kind(AuditEventKind::JobCancelled);
512        assert_eq!(cancelled.len(), 1);
513    }
514
515    #[test]
516    fn test_query_by_actor() {
517        let mut log = AuditLog::with_defaults();
518        log.log_job_submitted(Uuid::new_v4(), "alice", "av1");
519        log.log_job_submitted(Uuid::new_v4(), "bob", "vp9");
520        log.log_job_submitted(Uuid::new_v4(), "alice", "opus");
521
522        let alice_events = log.query_by_actor("alice");
523        assert_eq!(alice_events.len(), 2);
524    }
525
526    #[test]
527    fn test_query_by_target() {
528        let mut log = AuditLog::with_defaults();
529        let job_id = Uuid::new_v4();
530        log.log_job_submitted(job_id, "user", "av1");
531        log.log_job_cancelled(job_id, "user", "user request");
532        log.log_job_submitted(Uuid::new_v4(), "user", "vp9");
533
534        let target_events = log.query_by_target(&job_id.to_string());
535        assert_eq!(target_events.len(), 2);
536    }
537
538    #[test]
539    fn test_query_by_min_severity() {
540        let mut log = AuditLog::with_defaults();
541        log.log_job_submitted(Uuid::new_v4(), "user", "av1"); // Info
542        log.log_job_cancelled(Uuid::new_v4(), "user", "timeout"); // Warning
543        log.log_job_failed(Uuid::new_v4(), "user", "crash"); // Error
544
545        let warnings_plus = log.query_by_min_severity(AuditSeverity::Warning);
546        assert_eq!(warnings_plus.len(), 2);
547
548        let errors_only = log.query_by_min_severity(AuditSeverity::Error);
549        assert_eq!(errors_only.len(), 1);
550    }
551
552    #[test]
553    fn test_recent_entries() {
554        let mut log = AuditLog::with_defaults();
555        for i in 0..10 {
556            log.append(
557                AuditEntryBuilder::new(AuditEventKind::Custom, "actor")
558                    .message(format!("event {i}")),
559            );
560        }
561
562        let recent = log.recent(3);
563        assert_eq!(recent.len(), 3);
564        assert!(recent[0].message.contains("event 7"));
565        assert!(recent[2].message.contains("event 9"));
566    }
567
568    #[test]
569    fn test_get_by_sequence() {
570        let mut log = AuditLog::with_defaults();
571        let seq = log
572            .append(
573                AuditEntryBuilder::new(AuditEventKind::LeaderElected, "node-1").message("elected"),
574            )
575            .expect("append");
576
577        let entry = log.get_by_sequence(seq).expect("found");
578        assert_eq!(entry.kind, AuditEventKind::LeaderElected);
579        assert!(log.get_by_sequence(99999).is_none());
580    }
581
582    #[test]
583    fn test_clear() {
584        let mut log = AuditLog::with_defaults();
585        log.log_job_submitted(Uuid::new_v4(), "user", "av1");
586        log.log_job_submitted(Uuid::new_v4(), "user", "vp9");
587        assert_eq!(log.len(), 2);
588
589        log.clear();
590        assert!(log.is_empty());
591        assert_eq!(log.total_appended(), 2); // total is preserved
592    }
593
594    #[test]
595    fn test_convenience_worker_and_leader_logs() {
596        let mut log = AuditLog::with_defaults();
597        log.log_worker_joined("worker-1", "192.168.1.10:50052");
598        log.log_worker_left("worker-1", "heartbeat timeout");
599        log.log_leader_elected("node-3", 42);
600
601        assert_eq!(log.len(), 3);
602
603        let worker_events = log.query_by_kind(AuditEventKind::WorkerJoined);
604        assert_eq!(worker_events.len(), 1);
605
606        let leader_events = log.query_by_kind(AuditEventKind::LeaderElected);
607        assert_eq!(leader_events.len(), 1);
608        assert!(leader_events[0].message.contains("term 42"));
609    }
610
611    #[test]
612    fn test_builder_with_metadata() {
613        let mut log = AuditLog::with_defaults();
614        let seq = log
615            .append(
616                AuditEntryBuilder::new(AuditEventKind::ConfigChanged, "admin")
617                    .severity(AuditSeverity::Warning)
618                    .target("cluster-config")
619                    .message("max_retries changed")
620                    .meta("old_value", "3")
621                    .meta("new_value", "5"),
622            )
623            .expect("append");
624
625        let entry = log.get_by_sequence(seq).expect("found");
626        assert_eq!(entry.metadata.len(), 2);
627        assert_eq!(
628            entry.metadata[0],
629            ("old_value".to_string(), "3".to_string())
630        );
631        assert_eq!(entry.severity, AuditSeverity::Warning);
632    }
633}