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 pub user_task_type: Option<String>,
46 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#[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#[derive(Clone)]
76pub struct LedgerHandle {
77 tx: mpsc::Sender<UsageEntry>,
78}
79
80impl LedgerHandle {
81 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 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#[derive(Default)]
112pub struct InMemoryLedger {
113 pub entries: Mutex<Vec<UsageEntry>>,
114}
115
116impl InMemoryLedger {
117 #[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
132pub 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 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}