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