use crate::ids::{AppId, CommandName, NodeId};
use crate::operational_journal::FileOperationalJournal;
use crate::redact_text;
use crate::trace::TraceContext;
use parking_lot::Mutex;
use std::collections::VecDeque;
use std::sync::Arc;
const MAX_AUDIT_RECORDS: usize = 10_000;
#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum AuditOutcome {
Accepted,
Rejected,
Error,
}
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
pub struct AuditRecord {
pub command_id: String,
pub command_name: CommandName,
pub app_id: AppId,
pub node_id: NodeId,
pub timestamp_ms: u64,
pub outcome: AuditOutcome,
pub message: Option<String>,
pub trace: Option<TraceContext>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum AuditCategory {
Command,
Query,
Event,
Scheduler,
ControlPlane,
PeerRpc,
Runtime,
}
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
pub struct AuditEntry {
pub category: AuditCategory,
pub operation_id: String,
pub operation_name: String,
pub app_id: Option<String>,
pub node_id: Option<String>,
pub started_at_ms: u64,
pub completed_at_ms: u64,
pub latency_ms: u64,
pub outcome: AuditOutcome,
pub message: Option<String>,
pub trace: Option<TraceContext>,
}
impl AuditEntry {
pub fn new(
category: AuditCategory,
operation_id: impl Into<String>,
operation_name: impl Into<String>,
started_at_ms: u64,
completed_at_ms: u64,
outcome: AuditOutcome,
) -> Self {
Self {
category,
operation_id: operation_id.into(),
operation_name: operation_name.into(),
app_id: None,
node_id: None,
started_at_ms,
completed_at_ms,
latency_ms: completed_at_ms.saturating_sub(started_at_ms),
outcome,
message: None,
trace: None,
}
}
pub fn with_runtime_scope(mut self, app_id: &AppId, node_id: &NodeId) -> Self {
self.app_id = Some(app_id.as_str().to_string());
self.node_id = Some(node_id.as_str().to_string());
self
}
pub fn with_message(mut self, message: Option<String>) -> Self {
self.message = message.map(|value| redact_text(&value));
self
}
pub fn with_trace(mut self, trace: Option<TraceContext>) -> Self {
self.trace = trace;
self
}
}
#[derive(Debug, Default)]
pub struct AuditLog {
records: Mutex<VecDeque<AuditRecord>>,
entries: Mutex<VecDeque<AuditEntry>>,
journal: Mutex<Option<Arc<FileOperationalJournal>>>,
journal_error: Mutex<Option<String>>,
}
impl Clone for AuditLog {
fn clone(&self) -> Self {
let guard = self.records.lock();
Self {
records: Mutex::new(guard.clone()),
entries: Mutex::new(self.entries.lock().clone()),
journal: Mutex::new(self.journal.lock().clone()),
journal_error: Mutex::new(self.journal_error.lock().clone()),
}
}
}
impl AuditLog {
pub fn new() -> Self {
Self::default()
}
pub fn attach_journal(&self, journal: Arc<FileOperationalJournal>) {
let mut entries = journal.audit_entries();
if entries.len() > MAX_AUDIT_RECORDS {
entries.drain(..entries.len() - MAX_AUDIT_RECORDS);
}
*self.entries.lock() = entries.into();
*self.journal.lock() = Some(journal);
*self.journal_error.lock() = None;
}
pub fn durability_error(&self) -> Option<String> {
self.journal_error.lock().clone()
}
pub fn push(&self, mut record: AuditRecord) {
record.message = record.message.map(|message| redact_text(&message));
let completed_at_ms = now_ms();
let entry = AuditEntry::new(
AuditCategory::Command,
record.command_id.clone(),
record.command_name.as_str(),
record.timestamp_ms,
completed_at_ms,
record.outcome,
)
.with_runtime_scope(&record.app_id, &record.node_id)
.with_message(record.message.clone())
.with_trace(record.trace.clone());
let mut guard = self.records.lock();
while guard.len() >= MAX_AUDIT_RECORDS {
guard.pop_front();
}
guard.push_back(record);
drop(guard);
self.push_entry(entry);
}
pub fn push_entry(&self, mut entry: AuditEntry) {
entry.message = entry.message.map(|message| redact_text(&message));
if let Some(journal) = self.journal.lock().clone() {
if let Err(error) = journal.append_audit(entry.clone()) {
*self.journal_error.lock() = Some(redact_text(&format!("{error:?}")));
}
}
let mut entries = self.entries.lock();
while entries.len() >= MAX_AUDIT_RECORDS {
entries.pop_front();
}
entries.push_back(entry);
}
pub fn len(&self) -> usize {
self.records.lock().len()
}
pub fn is_empty(&self) -> bool {
self.records.lock().is_empty()
}
pub fn records(&self) -> Vec<AuditRecord> {
self.records.lock().iter().cloned().collect()
}
pub fn entries(&self) -> Vec<AuditEntry> {
self.entries.lock().iter().cloned().collect()
}
pub fn export_jsonl(&self) -> Result<String, serde_json::Error> {
let entries = self.entries.lock();
let mut output = String::new();
for entry in entries.iter() {
output.push_str(&serde_json::to_string(entry)?);
output.push('\n');
}
Ok(output)
}
}
fn now_ms() -> u64 {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|duration| duration.as_millis() as u64)
.unwrap_or(0)
}
#[cfg(test)]
#[path = "audit_tests.rs"]
mod tests;