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