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 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 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#[derive(Clone)]
54pub struct LedgerHandle {
55 tx: mpsc::Sender<UsageEntry>,
56}
57
58impl LedgerHandle {
59 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 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#[derive(Default)]
84pub struct InMemoryLedger {
85 pub entries: Mutex<Vec<UsageEntry>>,
86}
87
88impl InMemoryLedger {
89 #[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
104pub 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 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}