1use 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 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#[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#[derive(Clone)]
66pub struct LedgerHandle {
67 tx: mpsc::Sender<UsageEntry>,
68}
69
70impl LedgerHandle {
71 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 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#[derive(Default)]
102pub struct InMemoryLedger {
103 pub entries: Mutex<Vec<UsageEntry>>,
104}
105
106impl InMemoryLedger {
107 #[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
122pub 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 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}