1use std::str::FromStr;
4
5use async_trait::async_trait;
6use sqlx::sqlite::{SqliteConnectOptions, SqlitePoolOptions};
7use sqlx::SqlitePool;
8
9use crate::ledger::{LedgerError, LedgerStore, UsageEntry};
10
11pub struct SqliteLedger {
12 pool: SqlitePool,
13}
14
15impl SqliteLedger {
16 pub async fn connect(dsn: &str) -> Result<Self, LedgerError> {
24 let opts = SqliteConnectOptions::from_str(dsn)
25 .map_err(|e| LedgerError::Backend(e.to_string()))?
26 .create_if_missing(true);
27
28 let pool = SqlitePoolOptions::new()
29 .max_connections(1)
30 .connect_with(opts)
31 .await
32 .map_err(|e| LedgerError::Backend(e.to_string()))?;
33
34 sqlx::raw_sql(include_str!("../../migrations/0001_usage_events.sql"))
37 .execute(&pool)
38 .await
39 .map_err(|e| LedgerError::Backend(e.to_string()))?;
40
41 let _ = sqlx::query("ALTER TABLE usage_events ADD COLUMN user_id TEXT")
44 .execute(&pool)
45 .await;
46 let _ = sqlx::query("ALTER TABLE usage_events ADD COLUMN thread_id TEXT")
47 .execute(&pool)
48 .await;
49 let _ = sqlx::query("ALTER TABLE usage_events ADD COLUMN message_id TEXT")
50 .execute(&pool)
51 .await;
52
53 Ok(Self { pool })
54 }
55}
56
57#[async_trait]
58impl LedgerStore for SqliteLedger {
59 async fn record(&self, e: &UsageEntry) -> Result<(), LedgerError> {
60 sqlx::query(
61 "INSERT INTO usage_events \
62 (ts, tenant, workspace, user_id, thread_id, message_id, route, provider, model, lane, \
63 input_tokens, output_tokens, cost_usd, request_id, status) \
64 VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)",
65 )
66 .bind(e.ts.to_rfc3339())
67 .bind(&e.tenant)
68 .bind(&e.workspace)
69 .bind(&e.user)
70 .bind(&e.thread)
71 .bind(&e.message)
72 .bind(&e.route)
73 .bind(&e.provider)
74 .bind(&e.model)
75 .bind(&e.lane)
76 .bind(e.input_tokens as i64)
77 .bind(e.output_tokens as i64)
78 .bind(e.cost_usd)
79 .bind(&e.request_id)
80 .bind(&e.status)
81 .execute(&self.pool)
82 .await
83 .map_err(|e| LedgerError::Backend(e.to_string()))?;
84 Ok(())
85 }
86}