Skip to main content

synapse/ledger/
mod.rs

1//! Pluggable cost ledger. The hot path enqueues onto a bounded channel drained
2//! by a background writer; on a full channel we drop + count, never block.
3
4use std::sync::Arc;
5
6use async_trait::async_trait;
7use chrono::{DateTime, Utc};
8use futures::future::join_all;
9use metrics::counter;
10use parking_lot::Mutex;
11use tokio::sync::mpsc;
12
13pub mod connect;
14pub mod event;
15#[cfg(feature = "ledger-postgres")]
16pub mod postgres;
17#[cfg(feature = "ledger-pubsub")]
18pub mod pubsub;
19#[cfg(feature = "ledger-sns")]
20pub mod sns;
21#[cfg(feature = "ledger-sqlite")]
22pub mod sqlite;
23
24#[derive(Debug, Clone)]
25pub struct UsageEntry {
26    pub ts: DateTime<Utc>,
27    pub tenant: String,
28    pub workspace: Option<String>,
29    pub route: String,
30    pub provider: String,
31    pub model: String,
32    pub lane: String,
33    pub input_tokens: u64,
34    pub output_tokens: u64,
35    pub cost_usd: f64,
36    pub request_id: String,
37    pub status: String,
38    /// Lane discriminator for the ledger: "chat" or "embedding".
39    pub op: String,
40}
41
42#[derive(Debug, thiserror::Error)]
43pub enum LedgerError {
44    #[error("ledger backend error: {0}")]
45    Backend(String),
46}
47
48#[async_trait]
49pub trait LedgerStore: Send + Sync {
50    async fn record(&self, entry: &UsageEntry) -> Result<(), LedgerError>;
51}
52
53/// Discards all usage events. Used when no ledger sink could be connected.
54#[derive(Default)]
55pub struct NoopLedger;
56
57#[async_trait]
58impl LedgerStore for NoopLedger {
59    async fn record(&self, _entry: &UsageEntry) -> Result<(), LedgerError> {
60        Ok(())
61    }
62}
63
64/// Fire-and-forget handle. Cloneable; the hot path calls `enqueue`.
65#[derive(Clone)]
66pub struct LedgerHandle {
67    tx: mpsc::Sender<UsageEntry>,
68}
69
70impl LedgerHandle {
71    /// Spawn the background writer draining into `store`. `capacity` bounds the
72    /// channel; a full channel drops the entry and bumps `ledger_dropped_total`.
73    pub fn spawn(store: Arc<dyn LedgerStore>, capacity: usize) -> Self {
74        let (tx, mut rx) = mpsc::channel::<UsageEntry>(capacity);
75        tokio::spawn(async move {
76            while let Some(entry) = rx.recv().await {
77                if let Err(e) = store.record(&entry).await {
78                    tracing::warn!(
79                        error = %e,
80                        tenant = %entry.tenant,
81                        request_id = %entry.request_id,
82                        "ledger write failed"
83                    );
84                    counter!("synapse_ledger_errors_total", "backend" => "writer").increment(1);
85                }
86            }
87            tracing::warn!("ledger background writer stopped");
88        });
89        Self { tx }
90    }
91
92    /// Non-blocking enqueue. Never awaits the write; drops + counts on full.
93    pub fn enqueue(&self, entry: UsageEntry) {
94        if self.tx.try_send(entry).is_err() {
95            counter!("synapse_ledger_dropped_total").increment(1);
96        }
97    }
98}
99
100/// In-memory store for tests.
101#[derive(Default)]
102pub struct InMemoryLedger {
103    pub entries: Mutex<Vec<UsageEntry>>,
104}
105
106impl InMemoryLedger {
107    /// Snapshot the recorded entries (test-only read accessor).
108    #[cfg(test)]
109    pub fn entries(&self) -> Vec<UsageEntry> {
110        self.entries.lock().clone()
111    }
112}
113
114#[async_trait]
115impl LedgerStore for InMemoryLedger {
116    async fn record(&self, entry: &UsageEntry) -> Result<(), LedgerError> {
117        self.entries.lock().push(entry.clone());
118        Ok(())
119    }
120}
121
122/// Records each entry to every configured sink, concurrently and independently.
123/// A sink failing never blocks the others; per-sink failures are logged and
124/// counted on `synapse_ledger_errors_total{backend=<label>}`. Always returns
125/// `Ok` — the ledger is fire-and-forget; the fan-out owns error reporting.
126pub struct FanoutLedger {
127    sinks: Vec<(&'static str, Arc<dyn LedgerStore>)>,
128}
129
130impl FanoutLedger {
131    pub fn new(sinks: Vec<(&'static str, Arc<dyn LedgerStore>)>) -> Self {
132        Self { sinks }
133    }
134}
135
136#[async_trait]
137impl LedgerStore for FanoutLedger {
138    async fn record(&self, entry: &UsageEntry) -> Result<(), LedgerError> {
139        let futs = self.sinks.iter().map(|(label, sink)| async move {
140            if let Err(e) = sink.record(entry).await {
141                tracing::warn!(backend = label, error = %e, tenant = %entry.tenant, "ledger sink write failed");
142                counter!("synapse_ledger_errors_total", "backend" => *label).increment(1);
143            }
144        });
145        join_all(futs).await;
146        Ok(())
147    }
148}
149
150#[cfg(test)]
151mod tests {
152    use super::*;
153
154    fn entry() -> UsageEntry {
155        UsageEntry {
156            ts: Utc::now(),
157            tenant: "acme".into(),
158            workspace: None,
159            route: "fast".into(),
160            provider: "vertex".into(),
161            model: "gemini-3-flash".into(),
162            lane: "standard".into(),
163            input_tokens: 3,
164            output_tokens: 5,
165            cost_usd: 0.001,
166            request_id: "r1".into(),
167            status: "ok".into(),
168            op: "chat".into(),
169        }
170    }
171
172    #[tokio::test]
173    async fn in_memory_records_directly() {
174        let store = InMemoryLedger::default();
175        store.record(&entry()).await.unwrap();
176        assert_eq!(store.entries.lock().len(), 1);
177    }
178
179    #[tokio::test]
180    async fn handle_drains_into_store() {
181        let store = Arc::new(InMemoryLedger::default());
182        let handle = LedgerHandle::spawn(store.clone(), 16);
183        handle.enqueue(entry());
184        // give the writer task a tick to drain
185        for _ in 0..50 {
186            if store.entries.lock().len() == 1 {
187                break;
188            }
189            tokio::time::sleep(std::time::Duration::from_millis(5)).await;
190        }
191        assert_eq!(store.entries.lock().len(), 1);
192    }
193
194    struct FailingLedger;
195    #[async_trait]
196    impl LedgerStore for FailingLedger {
197        async fn record(&self, _e: &UsageEntry) -> Result<(), LedgerError> {
198            Err(LedgerError::Backend("boom".into()))
199        }
200    }
201
202    #[tokio::test]
203    async fn fanout_records_to_all_sinks() {
204        let a = Arc::new(InMemoryLedger::default());
205        let b = Arc::new(InMemoryLedger::default());
206        let fanout = FanoutLedger::new(vec![
207            ("a", a.clone() as Arc<dyn LedgerStore>),
208            ("b", b.clone() as Arc<dyn LedgerStore>),
209        ]);
210        fanout.record(&entry()).await.unwrap();
211        assert_eq!(a.entries.lock().len(), 1);
212        assert_eq!(b.entries.lock().len(), 1);
213    }
214
215    #[tokio::test]
216    async fn handle_keeps_accepting_after_write_failures() {
217        let handle = LedgerHandle::spawn(Arc::new(FailingLedger), 16);
218        handle.enqueue(entry());
219        handle.enqueue(entry());
220        tokio::time::sleep(std::time::Duration::from_millis(50)).await;
221        handle.enqueue(entry());
222    }
223
224    #[tokio::test]
225    async fn fanout_survives_a_failing_sink_and_returns_ok() {
226        let healthy = Arc::new(InMemoryLedger::default());
227        let fanout = FanoutLedger::new(vec![
228            ("fail", Arc::new(FailingLedger) as Arc<dyn LedgerStore>),
229            ("mem", healthy.clone() as Arc<dyn LedgerStore>),
230        ]);
231        let r = fanout.record(&entry()).await;
232        assert!(r.is_ok());
233        assert_eq!(healthy.entries.lock().len(), 1);
234    }
235}