Skip to main content

synapse/ledger/
sqlite.rs

1//! SQLite ledger backend (feature `ledger-sqlite`).
2
3use 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    /// Connect (DSN like `sqlite://synapse.db?mode=rwc` or `sqlite::memory:`)
17    /// and create the table if absent.
18    ///
19    /// Uses `max_connections(1)` so that both file-backed and in-memory databases
20    /// work correctly: with `sqlite::memory:` every connection gets its own
21    /// isolated database, so a single connection ensures the migration and all
22    /// subsequent writes share the same in-memory DB.
23    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        // Run the multi-statement migration via raw_sql which supports
35        // multiple `;`-separated statements in a single call.
36        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        // Best-effort for databases created before the user_id column existed;
42        // SQLite has no ADD COLUMN IF NOT EXISTS, so ignore "duplicate column".
43        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}