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