Skip to main content

tokio_prompt_orchestrator/
audit.rs

1//! # Request Replay / Audit Log
2//!
3//! An append-only, capacity-bounded audit log for LLM inference requests and
4//! responses.  Entries can be filtered, queried, and exported as JSONL.
5//!
6//! ## Example
7//!
8//! ```rust
9//! use tokio_prompt_orchestrator::audit::{AuditEntry, AuditFilter, AuditLog};
10//! use chrono::Utc;
11//!
12//! let mut log = AuditLog::new(1000);
13//!
14//! let entry = AuditEntry {
15//!     id: 1,
16//!     timestamp: Utc::now(),
17//!     model_id: "claude-3-5-sonnet".to_string(),
18//!     input: "Hello!".to_string(),
19//!     output: vec!["Hi there!".to_string()],
20//!     latency_ms: 320,
21//!     cache_hit: false,
22//!     tags: vec!["prod".to_string()],
23//! };
24//!
25//! log.append(entry);
26//!
27//! let results = log.query(&AuditFilter::default());
28//! assert_eq!(results.len(), 1);
29//! ```
30
31use chrono::{DateTime, Utc};
32use serde::{Deserialize, Serialize};
33use std::collections::{HashMap, VecDeque};
34use std::io::Write;
35
36// ── Entry ──────────────────────────────────────────────────────────────────────
37
38/// A single audit log entry recording one inference request/response pair.
39#[derive(Debug, Clone, Serialize, Deserialize)]
40pub struct AuditEntry {
41    /// Monotonically increasing entry identifier within this log instance.
42    pub id: u64,
43    /// UTC timestamp at which the request was received.
44    pub timestamp: DateTime<Utc>,
45    /// Model identifier (e.g. `"claude-3-5-sonnet"`, `"gpt-4o"`).
46    pub model_id: String,
47    /// The raw prompt / input text sent to the model.
48    pub input: String,
49    /// One or more output tokens/chunks returned by the model.
50    pub output: Vec<String>,
51    /// End-to-end latency in milliseconds.
52    pub latency_ms: u64,
53    /// Whether the response was served from a local cache.
54    pub cache_hit: bool,
55    /// Arbitrary string tags for grouping / filtering (e.g. `"prod"`, `"user:abc"`).
56    pub tags: Vec<String>,
57}
58
59// ── Filter ─────────────────────────────────────────────────────────────────────
60
61/// Predicate for [`AuditLog::query`].
62///
63/// All provided fields are AND-combined.  A `None` field matches any value.
64#[derive(Debug, Clone, Default)]
65pub struct AuditFilter {
66    /// Only return entries with `timestamp >= since`.
67    pub since: Option<DateTime<Utc>>,
68    /// Only return entries whose `model_id` equals this value.
69    pub model_id: Option<String>,
70    /// Only return entries matching this cache-hit status.
71    pub cache_hit: Option<bool>,
72    /// Only return entries with `latency_ms >= min_latency_ms`.
73    pub min_latency_ms: Option<u64>,
74}
75
76impl AuditFilter {
77    /// Return `true` if `entry` satisfies all filter predicates.
78    pub fn matches(&self, entry: &AuditEntry) -> bool {
79        if let Some(since) = self.since {
80            if entry.timestamp < since {
81                return false;
82            }
83        }
84        if let Some(ref model) = self.model_id {
85            if &entry.model_id != model {
86                return false;
87            }
88        }
89        if let Some(hit) = self.cache_hit {
90            if entry.cache_hit != hit {
91                return false;
92            }
93        }
94        if let Some(min_lat) = self.min_latency_ms {
95            if entry.latency_ms < min_lat {
96                return false;
97            }
98        }
99        true
100    }
101}
102
103// ── Stats ──────────────────────────────────────────────────────────────────────
104
105/// Aggregate statistics computed over all entries currently in the log.
106#[derive(Debug, Clone)]
107pub struct AuditStats {
108    /// Total number of entries in the log.
109    pub total_entries: usize,
110    /// Fraction of entries that were cache hits (0.0–1.0).
111    pub cache_hit_rate: f64,
112    /// Mean latency across all entries in milliseconds.
113    pub avg_latency_ms: f64,
114    /// Per-model entry counts.
115    pub by_model: HashMap<String, u64>,
116}
117
118// ── Log ───────────────────────────────────────────────────────────────────────
119
120/// Append-only audit log backed by a [`VecDeque`].
121///
122/// When the log reaches `capacity` entries the oldest entry is evicted before
123/// the new one is inserted, keeping memory usage bounded.
124pub struct AuditLog {
125    entries: VecDeque<AuditEntry>,
126    capacity: usize,
127    next_id: u64,
128}
129
130impl AuditLog {
131    /// Create a new log with the given maximum `capacity`.
132    ///
133    /// A capacity of `0` is legal but means every appended entry is immediately
134    /// evicted.
135    pub fn new(capacity: usize) -> Self {
136        Self {
137            entries: VecDeque::with_capacity(capacity.min(4096)),
138            capacity,
139            next_id: 1,
140        }
141    }
142
143    /// Append `entry` to the log, assigning a fresh `id` to it.
144    ///
145    /// If the log is already at capacity the oldest entry is removed first.
146    pub fn append(&mut self, mut entry: AuditEntry) {
147        if self.capacity == 0 {
148            return;
149        }
150        if self.entries.len() >= self.capacity {
151            self.entries.pop_front();
152        }
153        entry.id = self.next_id;
154        self.next_id += 1;
155        self.entries.push_back(entry);
156    }
157
158    /// Return all entries matching `filter`, in insertion order (oldest first).
159    pub fn query(&self, filter: &AuditFilter) -> Vec<&AuditEntry> {
160        self.entries.iter().filter(|e| filter.matches(e)).collect()
161    }
162
163    /// Compute aggregate statistics over all entries currently in the log.
164    pub fn stats(&self) -> AuditStats {
165        let total = self.entries.len();
166        if total == 0 {
167            return AuditStats {
168                total_entries: 0,
169                cache_hit_rate: 0.0,
170                avg_latency_ms: 0.0,
171                by_model: HashMap::new(),
172            };
173        }
174        let hits = self.entries.iter().filter(|e| e.cache_hit).count();
175        let total_latency: u64 = self.entries.iter().map(|e| e.latency_ms).sum();
176        let mut by_model: HashMap<String, u64> = HashMap::new();
177        for e in &self.entries {
178            *by_model.entry(e.model_id.clone()).or_insert(0) += 1;
179        }
180        AuditStats {
181            total_entries: total,
182            cache_hit_rate: hits as f64 / total as f64,
183            avg_latency_ms: total_latency as f64 / total as f64,
184            by_model,
185        }
186    }
187
188    /// Serialize all entries to `writer` as newline-delimited JSON (JSONL).
189    ///
190    /// Each line is a JSON object corresponding to one [`AuditEntry`].
191    pub fn export_jsonl(&self, writer: &mut dyn Write) -> std::io::Result<()> {
192        for entry in &self.entries {
193            let line = serde_json::to_string(entry)
194                .map_err(std::io::Error::other)?;
195            writeln!(writer, "{}", line)?;
196        }
197        Ok(())
198    }
199
200    /// Return the number of entries currently stored.
201    pub fn len(&self) -> usize {
202        self.entries.len()
203    }
204
205    /// Return `true` if the log contains no entries.
206    pub fn is_empty(&self) -> bool {
207        self.entries.is_empty()
208    }
209
210    /// Return the configured capacity.
211    pub fn capacity(&self) -> usize {
212        self.capacity
213    }
214}
215
216// ── HTTP handler helpers ───────────────────────────────────────────────────────
217// These are standalone functions rather than axum routes so the audit module
218// compiles without the `web-api` feature.  The web_api module can re-export
219// them when the feature is active.
220
221/// Serialisable response body for `GET /api/v1/audit`.
222#[derive(Debug, Serialize)]
223pub struct AuditQueryResponse {
224    /// Entries returned by the query.
225    pub entries: Vec<AuditEntry>,
226    /// Total number of entries returned.
227    pub count: usize,
228}
229
230/// Serialisable response body for `GET /api/v1/audit/stats`.
231#[derive(Debug, Serialize)]
232pub struct AuditStatsResponse {
233    /// Total entries in the log.
234    pub total_entries: usize,
235    /// Cache hit rate (0.0–1.0).
236    pub cache_hit_rate: f64,
237    /// Average latency in milliseconds.
238    pub avg_latency_ms: f64,
239    /// Per-model entry counts.
240    pub by_model: HashMap<String, u64>,
241}
242
243impl From<AuditStats> for AuditStatsResponse {
244    fn from(s: AuditStats) -> Self {
245        Self {
246            total_entries: s.total_entries,
247            cache_hit_rate: s.cache_hit_rate,
248            avg_latency_ms: s.avg_latency_ms,
249            by_model: s.by_model,
250        }
251    }
252}
253
254// ── Tests ──────────────────────────────────────────────────────────────────────
255
256#[cfg(test)]
257mod tests {
258    use super::*;
259    use chrono::{Duration, Utc};
260
261    fn make_entry(model: &str, latency_ms: u64, cache_hit: bool) -> AuditEntry {
262        AuditEntry {
263            id: 0, // will be assigned by AuditLog::append
264            timestamp: Utc::now(),
265            model_id: model.to_string(),
266            input: "test input".to_string(),
267            output: vec!["test output".to_string()],
268            latency_ms,
269            cache_hit,
270            tags: vec![],
271        }
272    }
273
274    // ── Basic append / query ──────────────────────────────────────────────────
275
276    #[test]
277    fn append_and_query_all() {
278        let mut log = AuditLog::new(100);
279        log.append(make_entry("gpt-4o", 200, false));
280        log.append(make_entry("claude-3", 150, true));
281        let all = log.query(&AuditFilter::default());
282        assert_eq!(all.len(), 2);
283    }
284
285    #[test]
286    fn ids_are_assigned_sequentially() {
287        let mut log = AuditLog::new(100);
288        log.append(make_entry("m", 100, false));
289        log.append(make_entry("m", 100, false));
290        log.append(make_entry("m", 100, false));
291        let all = log.query(&AuditFilter::default());
292        let ids: Vec<u64> = all.iter().map(|e| e.id).collect();
293        assert_eq!(ids, vec![1, 2, 3]);
294    }
295
296    #[test]
297    fn capacity_evicts_oldest() {
298        let mut log = AuditLog::new(2);
299        log.append(make_entry("m", 100, false));
300        log.append(make_entry("m", 200, false));
301        log.append(make_entry("m", 300, false)); // evicts first
302        assert_eq!(log.len(), 2);
303        let latencies: Vec<u64> = log.query(&AuditFilter::default())
304            .iter().map(|e| e.latency_ms).collect();
305        assert_eq!(latencies, vec![200, 300]);
306    }
307
308    #[test]
309    fn capacity_zero_stores_nothing() {
310        let mut log = AuditLog::new(0);
311        log.append(make_entry("m", 100, false));
312        assert!(log.is_empty());
313    }
314
315    // ── Filter ────────────────────────────────────────────────────────────────
316
317    #[test]
318    fn filter_by_model_id() {
319        let mut log = AuditLog::new(100);
320        log.append(make_entry("gpt-4o", 100, false));
321        log.append(make_entry("claude-3", 200, false));
322        let filter = AuditFilter { model_id: Some("gpt-4o".to_string()), ..Default::default() };
323        let results = log.query(&filter);
324        assert_eq!(results.len(), 1);
325        assert_eq!(results[0].model_id, "gpt-4o");
326    }
327
328    #[test]
329    fn filter_by_cache_hit() {
330        let mut log = AuditLog::new(100);
331        log.append(make_entry("m", 100, false));
332        log.append(make_entry("m", 200, true));
333        log.append(make_entry("m", 300, true));
334        let filter = AuditFilter { cache_hit: Some(true), ..Default::default() };
335        let results = log.query(&filter);
336        assert_eq!(results.len(), 2);
337    }
338
339    #[test]
340    fn filter_by_min_latency() {
341        let mut log = AuditLog::new(100);
342        log.append(make_entry("m", 50, false));
343        log.append(make_entry("m", 100, false));
344        log.append(make_entry("m", 500, false));
345        let filter = AuditFilter { min_latency_ms: Some(100), ..Default::default() };
346        let results = log.query(&filter);
347        assert_eq!(results.len(), 2);
348    }
349
350    #[test]
351    fn filter_by_since() {
352        let mut log = AuditLog::new(100);
353        let old_entry = AuditEntry {
354            id: 0,
355            timestamp: Utc::now() - Duration::hours(2),
356            model_id: "m".to_string(),
357            input: "old".to_string(),
358            output: vec![],
359            latency_ms: 100,
360            cache_hit: false,
361            tags: vec![],
362        };
363        let new_entry = make_entry("m", 100, false);
364        log.append(old_entry);
365        log.append(new_entry);
366        let cutoff = Utc::now() - Duration::minutes(30);
367        let filter = AuditFilter { since: Some(cutoff), ..Default::default() };
368        let results = log.query(&filter);
369        assert_eq!(results.len(), 1);
370    }
371
372    #[test]
373    fn filter_combined_predicates() {
374        let mut log = AuditLog::new(100);
375        log.append(make_entry("gpt-4o", 100, false));
376        log.append(make_entry("gpt-4o", 500, true));
377        log.append(make_entry("claude-3", 500, true));
378        let filter = AuditFilter {
379            model_id: Some("gpt-4o".to_string()),
380            cache_hit: Some(true),
381            ..Default::default()
382        };
383        let results = log.query(&filter);
384        assert_eq!(results.len(), 1);
385        assert_eq!(results[0].latency_ms, 500);
386    }
387
388    // ── Stats ─────────────────────────────────────────────────────────────────
389
390    #[test]
391    fn stats_empty_log() {
392        let log = AuditLog::new(100);
393        let stats = log.stats();
394        assert_eq!(stats.total_entries, 0);
395        assert_eq!(stats.cache_hit_rate, 0.0);
396        assert_eq!(stats.avg_latency_ms, 0.0);
397        assert!(stats.by_model.is_empty());
398    }
399
400    #[test]
401    fn stats_cache_hit_rate() {
402        let mut log = AuditLog::new(100);
403        log.append(make_entry("m", 100, true));
404        log.append(make_entry("m", 100, false));
405        let stats = log.stats();
406        assert!((stats.cache_hit_rate - 0.5).abs() < 1e-9);
407    }
408
409    #[test]
410    fn stats_avg_latency() {
411        let mut log = AuditLog::new(100);
412        log.append(make_entry("m", 100, false));
413        log.append(make_entry("m", 200, false));
414        let stats = log.stats();
415        assert!((stats.avg_latency_ms - 150.0).abs() < 1e-9);
416    }
417
418    #[test]
419    fn stats_by_model_counts() {
420        let mut log = AuditLog::new(100);
421        log.append(make_entry("gpt-4o", 100, false));
422        log.append(make_entry("gpt-4o", 100, false));
423        log.append(make_entry("claude-3", 100, false));
424        let stats = log.stats();
425        assert_eq!(stats.by_model["gpt-4o"], 2);
426        assert_eq!(stats.by_model["claude-3"], 1);
427    }
428
429    // ── JSONL export ──────────────────────────────────────────────────────────
430
431    #[test]
432    fn export_jsonl_produces_valid_lines() {
433        let mut log = AuditLog::new(100);
434        log.append(make_entry("gpt-4o", 100, false));
435        log.append(make_entry("claude-3", 200, true));
436
437        let mut buf: Vec<u8> = Vec::new();
438        log.export_jsonl(&mut buf).unwrap();
439
440        let text = String::from_utf8(buf).unwrap();
441        let lines: Vec<&str> = text.lines().collect();
442        assert_eq!(lines.len(), 2);
443
444        // Each line should parse as valid JSON with expected fields.
445        for line in &lines {
446            let val: serde_json::Value = serde_json::from_str(line).unwrap();
447            assert!(val.get("model_id").is_some());
448            assert!(val.get("latency_ms").is_some());
449        }
450    }
451
452    #[test]
453    fn export_jsonl_empty_log_empty_output() {
454        let log = AuditLog::new(100);
455        let mut buf: Vec<u8> = Vec::new();
456        log.export_jsonl(&mut buf).unwrap();
457        assert!(buf.is_empty());
458    }
459
460    // ── AuditStatsResponse conversion ─────────────────────────────────────────
461
462    #[test]
463    fn audit_stats_response_from_stats() {
464        let mut log = AuditLog::new(100);
465        log.append(make_entry("m", 100, true));
466        let stats = log.stats();
467        let resp = AuditStatsResponse::from(stats.clone());
468        assert_eq!(resp.total_entries, stats.total_entries);
469        assert!((resp.cache_hit_rate - stats.cache_hit_rate).abs() < 1e-9);
470    }
471}