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