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