use std::sync::Arc;
use crate::observability::{LlmCallRecord, ToolCallRecord};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum RecordLevel {
Debug,
Info,
Warn,
Error,
}
#[async_trait::async_trait]
pub trait LogExporter: Send + Sync + 'static {
fn accepts(&self, _record_level: RecordLevel) -> bool {
true
}
async fn export_llm(&self, record: &LlmCallRecord);
async fn export_tool(&self, record: &ToolCallRecord);
fn validate(&self) -> Result<(), String> {
Ok(())
}
}
pub struct ExporterRouter {
exporters: Vec<Arc<dyn LogExporter>>,
}
impl ExporterRouter {
pub fn new() -> Self {
Self {
exporters: Vec::new(),
}
}
pub fn with_capacity(n: usize) -> Self {
Self {
exporters: Vec::with_capacity(n),
}
}
pub fn register(&mut self, exporter: Arc<dyn LogExporter>) -> Result<(), String> {
if self
.exporters
.iter()
.any(|existing| Arc::ptr_eq(existing, &exporter))
{
return Ok(());
}
exporter.validate()?;
self.exporters.push(exporter);
Ok(())
}
pub fn len(&self) -> usize {
self.exporters.len()
}
pub fn is_empty(&self) -> bool {
self.exporters.is_empty()
}
pub async fn log_llm(&self, level: RecordLevel, record: &LlmCallRecord) {
for exporter in &self.exporters {
if !exporter.accepts(level) {
continue;
}
exporter.export_llm(record).await;
}
}
pub async fn log_tool(&self, level: RecordLevel, record: &ToolCallRecord) {
for exporter in &self.exporters {
if !exporter.accepts(level) {
continue;
}
exporter.export_tool(record).await;
}
}
pub fn log_llm_spawned(&self, level: RecordLevel, record: LlmCallRecord) {
let router = self.clone();
tokio::spawn(async move { router.log_llm(level, &record).await });
}
}
impl Clone for ExporterRouter {
fn clone(&self) -> Self {
Self {
exporters: self.exporters.clone(),
}
}
}
impl Default for ExporterRouter {
fn default() -> Self {
Self::new()
}
}
struct ClosureExporter<L, T> {
fn_llm: L,
fn_tool: T,
}
#[async_trait::async_trait]
impl<L, T> LogExporter for ClosureExporter<L, T>
where
L: Fn(&LlmCallRecord) + Send + Sync + 'static,
T: Fn(&ToolCallRecord) + Send + Sync + 'static,
{
async fn export_llm(&self, record: &LlmCallRecord) {
(self.fn_llm)(record);
}
async fn export_tool(&self, record: &ToolCallRecord) {
(self.fn_tool)(record);
}
}
pub fn closure<L, T>(fn_llm: L, fn_tool: T) -> Arc<dyn LogExporter>
where
L: Fn(&LlmCallRecord) + Send + Sync + 'static,
T: Fn(&ToolCallRecord) + Send + Sync + 'static,
{
Arc::new(ClosureExporter { fn_llm, fn_tool })
}
pub struct LogRing {
inner: std::sync::Mutex<RingState>,
}
struct RingState {
records: std::collections::VecDeque<Arc<LlmCallRecord>>,
cap: usize,
}
impl LogRing {
pub fn new(cap: usize) -> Self {
Self {
inner: std::sync::Mutex::new(RingState {
records: std::collections::VecDeque::new(),
cap: cap.max(1),
}),
}
}
pub fn new_exporter(self: &Arc<Self>) -> Arc<dyn LogExporter> {
self.clone()
}
pub fn push(&self, record: LlmCallRecord) {
let mut state = match self.inner.lock() {
Ok(state) => state,
Err(poisoned) => poisoned.into_inner(),
};
state.records.push_back(Arc::new(record));
while state.records.len() > state.cap {
state.records.pop_front();
}
}
pub fn snapshot(&self) -> Vec<Arc<LlmCallRecord>> {
let state = match self.inner.lock() {
Ok(state) => state,
Err(poisoned) => poisoned.into_inner(),
};
state.records.iter().cloned().collect()
}
pub fn len(&self) -> usize {
let state = match self.inner.lock() {
Ok(state) => state,
Err(poisoned) => poisoned.into_inner(),
};
state.records.len()
}
pub fn is_empty(&self) -> bool {
self.len() == 0
}
pub fn set_cap(&self, cap: usize) {
let mut state = match self.inner.lock() {
Ok(state) => state,
Err(poisoned) => poisoned.into_inner(),
};
state.cap = cap.max(1);
while state.records.len() > state.cap {
state.records.pop_front();
}
}
}
#[async_trait::async_trait]
impl LogExporter for LogRing {
async fn export_llm(&self, record: &LlmCallRecord) {
self.push(record.clone());
}
async fn export_tool(&self, _record: &ToolCallRecord) {}
}
#[derive(Debug, Clone, Copy, Default)]
pub struct TracingExporter;
impl TracingExporter {
pub fn new() -> Self {
Self
}
}
#[async_trait::async_trait]
impl LogExporter for TracingExporter {
async fn export_llm(&self, record: &LlmCallRecord) {
let latency_ms = record.latency_ms;
let model = record.model.as_str();
let provider = record.provider.as_str();
let status = record.status.as_str();
let cached_tokens = record.cached_tokens;
let total_time_ms = record.total_time_ms;
if status == "success" {
tracing::info!(
latency_ms,
model,
provider,
status,
cached_tokens,
total_time_ms,
"LLM call completed"
);
} else {
tracing::warn!(
latency_ms,
model,
provider,
status,
cached_tokens,
total_time_ms,
"LLM call failed"
);
}
}
async fn export_tool(&self, record: &ToolCallRecord) {
let latency_ms = record.latency_ms;
let status = record.status.as_str();
let tool_name = record.tool_name.as_str();
let tool_type = record.tool_type.as_str();
if status == "success" {
tracing::info!(
latency_ms,
status,
tool_name,
tool_type,
"Tool call completed"
);
} else {
tracing::warn!(latency_ms, status, tool_name, tool_type, "Tool call failed");
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::Mutex;
type Capture = Arc<(Mutex<usize>, Mutex<Vec<String>>)>;
struct MockExporter {
tag: &'static str,
errors_only: bool,
capture: Capture,
shared_log: Option<Arc<Mutex<Vec<String>>>>,
}
impl MockExporter {
fn new(tag: &'static str) -> Self {
Self {
tag,
errors_only: false,
capture: Arc::new((Mutex::new(0), Mutex::new(Vec::new()))),
shared_log: None,
}
}
fn errors_only(tag: &'static str) -> Self {
Self {
errors_only: true,
..Self::new(tag)
}
}
fn ordered(tag: &'static str, log: Arc<Mutex<Vec<String>>>) -> Self {
Self {
shared_log: Some(log),
..Self::new(tag)
}
}
fn record(&self, kind: &str) {
*self.capture.0.lock().unwrap() += 1;
self.capture.1.lock().unwrap().push(kind.to_string());
if let Some(log) = &self.shared_log {
log.lock().unwrap().push(self.tag.to_string());
}
}
}
#[async_trait::async_trait]
impl LogExporter for MockExporter {
fn accepts(&self, level: RecordLevel) -> bool {
!self.errors_only || level == RecordLevel::Error
}
async fn export_llm(&self, _record: &LlmCallRecord) {
self.record("llm");
}
async fn export_tool(&self, _record: &ToolCallRecord) {
self.record("tool");
}
}
struct BrokenExporter;
#[async_trait::async_trait]
impl LogExporter for BrokenExporter {
async fn export_llm(&self, _record: &LlmCallRecord) {}
async fn export_tool(&self, _record: &ToolCallRecord) {}
fn validate(&self) -> Result<(), String> {
Err(String::from("broken exporter refuses validation"))
}
}
fn sample_llm() -> LlmCallRecord {
LlmCallRecord {
step_index: 0,
provider: String::from("openai"),
model: String::from("gpt-4o"),
prompt_tokens: 10,
completion_tokens: 5,
latency_ms: 42,
status: String::from("success"),
cached_tokens: Some(7),
total_time_ms: Some(42),
}
}
#[test]
fn sample_record_carries_optional_usage_fields() {
let record = sample_llm();
assert_eq!(record.cached_tokens, Some(7));
assert_eq!(record.total_time_ms, Some(42));
}
async fn router_register_and_export(exporter: &Arc<dyn LogExporter>, record: &LlmCallRecord) {
exporter.export_llm(record).await;
}
fn sample_tool() -> ToolCallRecord {
ToolCallRecord {
step_index: 0,
tool_name: String::from("calculator"),
tool_type: String::from("builtin"),
arguments: serde_json::json!({ "expression": "1 + 1" }),
result: None,
latency_ms: 7,
status: String::from("success"),
}
}
#[tokio::test]
async fn router_fans_out_to_all_exporters() {
let mut router = ExporterRouter::new();
let first = MockExporter::new("first");
let second = MockExporter::new("second");
let first_capture = first.capture.clone();
let second_capture = second.capture.clone();
router.register(Arc::new(first)).unwrap();
router.register(Arc::new(second)).unwrap();
assert_eq!(router.len(), 2);
router.log_llm(RecordLevel::Info, &sample_llm()).await;
router.log_tool(RecordLevel::Info, &sample_tool()).await;
for capture in [first_capture, second_capture] {
let counts = capture.0.lock().unwrap();
assert_eq!(*counts, 2, "each exporter sees both records");
let kinds = capture.1.lock().unwrap();
assert_eq!(*kinds, vec![String::from("llm"), String::from("tool")]);
}
}
#[tokio::test]
async fn accepts_gate_filters_records() {
let mut router = ExporterRouter::new();
let exporter = MockExporter::errors_only("gate");
let capture = exporter.capture.clone();
router.register(Arc::new(exporter)).unwrap();
router.log_llm(RecordLevel::Info, &sample_llm()).await;
router.log_tool(RecordLevel::Warn, &sample_tool()).await;
assert_eq!(
*capture.0.lock().unwrap(),
0,
"Info and Warn are filtered out"
);
router.log_llm(RecordLevel::Error, &sample_llm()).await;
assert_eq!(*capture.0.lock().unwrap(), 1, "Error passes the gate");
}
#[tokio::test]
async fn duplicate_registration_is_skipped() {
let mut router = ExporterRouter::new();
let exporter: Arc<dyn LogExporter> = Arc::new(MockExporter::new("dup"));
assert!(router.register(exporter.clone()).is_ok());
assert!(
router.register(exporter.clone()).is_ok(),
"duplicate register is a silent no-op, not an error"
);
assert_eq!(router.len(), 1);
}
#[tokio::test]
async fn validate_rejects_broken_exporter() {
let mut router = ExporterRouter::new();
let error = router
.register(Arc::new(BrokenExporter))
.expect_err("validation failure must reject registration");
assert_eq!(error, "broken exporter refuses validation");
assert!(router.is_empty(), "rejected exporter is not stored");
}
#[tokio::test]
async fn ring_trims_at_capacity() {
let ring = Arc::new(LogRing::new(3));
let exporter = ring.new_exporter();
for i in 0..5u32 {
let mut record = sample_llm();
record.prompt_tokens = i as i64;
router_register_and_export(&exporter, &record).await;
}
assert_eq!(ring.len(), 3, "capacity bounds the retained records");
let snap = ring.snapshot();
let prompts: Vec<i64> = snap.iter().map(|r| r.prompt_tokens).collect();
assert_eq!(prompts, vec![2, 3, 4], "oldest entries evicted first");
ring.set_cap(1);
let snap = ring.snapshot();
assert_eq!(snap.len(), 1);
assert_eq!(snap[0].prompt_tokens, 4, "newest survivor retained");
ring.set_cap(4);
assert_eq!(ring.len(), 1);
}
#[tokio::test]
async fn ring_ignores_tool_records() {
let ring = Arc::new(LogRing::new(4));
let exporter = ring.new_exporter();
exporter.export_tool(&sample_tool()).await;
assert!(ring.is_empty(), "tool records are deliberately not stored");
}
#[tokio::test]
async fn slow_exporter_does_not_block_others() {
let mut router = ExporterRouter::new();
let shared: Arc<Mutex<Vec<String>>> = Arc::new(Mutex::new(Vec::new()));
router
.register(Arc::new(MockExporter::ordered("a", shared.clone())))
.unwrap();
router
.register(Arc::new(MockExporter::ordered("b", shared.clone())))
.unwrap();
router.log_llm(RecordLevel::Info, &sample_llm()).await;
assert_eq!(
*shared.lock().unwrap(),
vec![String::from("a"), String::from("b")],
"fan-out is sequential in registration order"
);
}
}