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