Skip to main content

bezant_server/events/
persistence.rs

1//! Optional sqlite persistence for captured events.
2//!
3//! The ring buffer in [`super::ring::TopicRing`] is the primary read
4//! path — fast, in-memory, bounded. Sqlite is the **historical** store:
5//! every event we push to a ring is also appended to a single
6//! `events` table so the `GET /events/{topic}/history?since_ts=…`
7//! endpoint can reach back beyond ring capacity (and beyond container
8//! lifetime).
9//!
10//! Schema:
11//!
12//! ```sql
13//! CREATE TABLE events (
14//!   id           INTEGER PRIMARY KEY AUTOINCREMENT,
15//!   cursor       INTEGER NOT NULL,
16//!   topic        TEXT    NOT NULL,
17//!   received_at  TEXT    NOT NULL,    -- ISO 8601
18//!   reset_epoch  INTEGER NOT NULL,
19//!   payload      TEXT    NOT NULL     -- JSON
20//! );
21//! CREATE INDEX events_topic_received_at ON events(topic, received_at);
22//! ```
23//!
24//! Retention is configurable per topic:
25//! - `orders`     — 90 days (default)
26//! - `pnl`        — 90 days
27//! - `marketdata:*` — 14 days (high volume; bigger window doesn't pay)
28//! - `gap`        — 365 days (low volume, useful forever)
29//!
30//! [`EventLog::prune`] is intended to be called from a
31//! once-an-hour task; nothing in the read path waits on it.
32
33use std::path::{Path, PathBuf};
34use std::sync::Mutex;
35
36use rusqlite::{params, Connection, OptionalExtension};
37
38use super::ObservedEvent;
39
40/// Sqlite-backed historical event store. Wraps a [`Connection`] in a
41/// [`Mutex`] so it can be shared across the connector + axum tasks
42/// without unsafe juggling. Reads are fast; writes serialise.
43pub struct EventLog {
44    conn: Mutex<Connection>,
45    /// Path the connection was opened against (for diagnostics).
46    path: PathBuf,
47}
48
49impl std::fmt::Debug for EventLog {
50    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
51        f.debug_struct("EventLog")
52            .field("path", &self.path)
53            .finish()
54    }
55}
56
57/// Per-topic retention policy. Topics not in the map default to 30 days.
58#[derive(Clone, Debug)]
59pub struct RetentionPolicy {
60    /// Default retention for topics not explicitly listed.
61    pub default_days: i64,
62    /// Days for `orders`.
63    pub orders_days: i64,
64    /// Days for `pnl`.
65    pub pnl_days: i64,
66    /// Days for `marketdata:*` topics.
67    pub marketdata_days: i64,
68    /// Days for `gap`.
69    pub gap_days: i64,
70}
71
72impl Default for RetentionPolicy {
73    fn default() -> Self {
74        Self {
75            default_days: 30,
76            orders_days: 90,
77            pnl_days: 90,
78            marketdata_days: 14,
79            gap_days: 365,
80        }
81    }
82}
83
84impl RetentionPolicy {
85    /// How many days to keep events for `topic`.
86    #[must_use]
87    pub fn days_for(&self, topic: &str) -> i64 {
88        if topic == "orders" {
89            self.orders_days
90        } else if topic == "pnl" {
91            self.pnl_days
92        } else if topic == "gap" {
93            self.gap_days
94        } else if topic.starts_with("marketdata:") {
95            self.marketdata_days
96        } else {
97            self.default_days
98        }
99    }
100}
101
102impl EventLog {
103    /// Open or create a sqlite database at `path`. Runs the schema
104    /// migration on first call (idempotent).
105    pub fn open(path: impl AsRef<Path>) -> rusqlite::Result<Self> {
106        let conn = Connection::open(path.as_ref())?;
107        conn.pragma_update(None, "journal_mode", "WAL")?;
108        conn.pragma_update(None, "synchronous", "NORMAL")?;
109        conn.execute_batch(
110            r"
111            CREATE TABLE IF NOT EXISTS events (
112                id           INTEGER PRIMARY KEY AUTOINCREMENT,
113                cursor       INTEGER NOT NULL,
114                topic        TEXT    NOT NULL,
115                received_at  TEXT    NOT NULL,
116                reset_epoch  INTEGER NOT NULL,
117                payload      TEXT    NOT NULL
118            );
119            CREATE INDEX IF NOT EXISTS events_topic_received_at
120                ON events(topic, received_at);
121            ",
122        )?;
123        Ok(Self {
124            conn: Mutex::new(conn),
125            path: path.as_ref().to_path_buf(),
126        })
127    }
128
129    /// Open an in-memory database — used by tests.
130    pub fn open_in_memory() -> rusqlite::Result<Self> {
131        let conn = Connection::open_in_memory()?;
132        conn.execute_batch(
133            r"
134            CREATE TABLE events (
135                id           INTEGER PRIMARY KEY AUTOINCREMENT,
136                cursor       INTEGER NOT NULL,
137                topic        TEXT    NOT NULL,
138                received_at  TEXT    NOT NULL,
139                reset_epoch  INTEGER NOT NULL,
140                payload      TEXT    NOT NULL
141            );
142            CREATE INDEX events_topic_received_at
143                ON events(topic, received_at);
144            ",
145        )?;
146        Ok(Self {
147            conn: Mutex::new(conn),
148            path: PathBuf::from(":memory:"),
149        })
150    }
151
152    /// Path the log was opened against.
153    #[must_use]
154    pub fn path(&self) -> &Path {
155        &self.path
156    }
157
158    /// Append an event. Returns the row id.
159    pub fn append(&self, event: &ObservedEvent) -> rusqlite::Result<i64> {
160        let payload_str = serde_json::to_string(&event.payload)
161            .map_err(|e| rusqlite::Error::ToSqlConversionFailure(Box::new(e)))?;
162        let conn = self.conn.lock().unwrap();
163        conn.execute(
164            "INSERT INTO events (cursor, topic, received_at, reset_epoch, payload) \
165             VALUES (?1, ?2, ?3, ?4, ?5)",
166            params![
167                event.cursor as i64,
168                event.topic,
169                event.received_at,
170                event.reset_epoch as i64,
171                payload_str,
172            ],
173        )?;
174        Ok(conn.last_insert_rowid())
175    }
176
177    /// Read events for a topic since `since_ts` (ISO 8601). Newest-last.
178    pub fn query_since(
179        &self,
180        topic: &str,
181        since_ts: &str,
182        limit: usize,
183    ) -> rusqlite::Result<Vec<ObservedEvent>> {
184        let conn = self.conn.lock().unwrap();
185        let mut stmt = conn.prepare(
186            "SELECT cursor, topic, received_at, reset_epoch, payload \
187             FROM events WHERE topic = ?1 AND received_at > ?2 \
188             ORDER BY received_at ASC LIMIT ?3",
189        )?;
190        let rows = stmt.query_map(params![topic, since_ts, limit as i64], |row| {
191            let cursor: i64 = row.get(0)?;
192            let topic: String = row.get(1)?;
193            let received_at: String = row.get(2)?;
194            let reset_epoch: i64 = row.get(3)?;
195            let payload_str: String = row.get(4)?;
196            let payload: serde_json::Value =
197                serde_json::from_str(&payload_str).unwrap_or(serde_json::Value::Null);
198            Ok(ObservedEvent {
199                cursor: cursor as u64,
200                topic,
201                received_at,
202                reset_epoch: reset_epoch as u64,
203                payload,
204            })
205        })?;
206        let mut out = Vec::new();
207        for r in rows {
208            out.push(r?);
209        }
210        Ok(out)
211    }
212
213    /// Drop events older than the topic's retention cutoff. Returns
214    /// total rows deleted across all topics.
215    pub fn prune(&self, policy: &RetentionPolicy) -> rusqlite::Result<usize> {
216        let conn = self.conn.lock().unwrap();
217        let now: chrono::DateTime<chrono::Utc> = std::time::SystemTime::now().into();
218        let mut total = 0usize;
219        let topics: Vec<String> = {
220            let mut stmt = conn.prepare("SELECT DISTINCT topic FROM events")?;
221            let mapped = stmt.query_map([], |r| r.get::<_, String>(0))?;
222            let mut v = Vec::new();
223            for m in mapped {
224                v.push(m?);
225            }
226            v
227        };
228        for topic in topics {
229            let days = policy.days_for(&topic);
230            let cutoff = now - chrono::Duration::days(days);
231            let cutoff_str = cutoff.to_rfc3339_opts(chrono::SecondsFormat::Millis, true);
232            let n = conn.execute(
233                "DELETE FROM events WHERE topic = ?1 AND received_at < ?2",
234                params![topic, cutoff_str],
235            )?;
236            total += n;
237        }
238        Ok(total)
239    }
240
241    /// Total row count. Used by tests + diagnostics.
242    pub fn count(&self) -> rusqlite::Result<usize> {
243        let conn = self.conn.lock().unwrap();
244        let n: i64 = conn
245            .query_row("SELECT COUNT(*) FROM events", [], |r| r.get(0))
246            .optional()?
247            .unwrap_or(0);
248        Ok(n as usize)
249    }
250}
251
252#[cfg(test)]
253mod tests {
254    use super::*;
255    use serde_json::json;
256
257    fn evt(cursor: u64, topic: &str, ts: &str) -> ObservedEvent {
258        ObservedEvent {
259            cursor,
260            topic: topic.to_string(),
261            received_at: ts.to_string(),
262            reset_epoch: 1,
263            payload: json!({"id": cursor}),
264        }
265    }
266
267    #[test]
268    fn append_and_query_round_trip() {
269        let log = EventLog::open_in_memory().unwrap();
270        log.append(&evt(1, "orders", "2026-05-06T13:30:00Z"))
271            .unwrap();
272        log.append(&evt(2, "orders", "2026-05-06T13:31:00Z"))
273            .unwrap();
274        log.append(&evt(3, "pnl", "2026-05-06T13:32:00Z")).unwrap();
275
276        let orders = log
277            .query_since("orders", "1970-01-01T00:00:00Z", 100)
278            .unwrap();
279        assert_eq!(orders.len(), 2);
280        assert_eq!(orders[0].cursor, 1);
281        assert_eq!(orders[1].cursor, 2);
282
283        let pnl = log.query_since("pnl", "1970-01-01T00:00:00Z", 100).unwrap();
284        assert_eq!(pnl.len(), 1);
285        assert_eq!(pnl[0].topic, "pnl");
286    }
287
288    #[test]
289    fn query_filters_by_since_ts() {
290        let log = EventLog::open_in_memory().unwrap();
291        log.append(&evt(1, "orders", "2026-05-06T13:30:00Z"))
292            .unwrap();
293        log.append(&evt(2, "orders", "2026-05-06T13:31:00Z"))
294            .unwrap();
295        log.append(&evt(3, "orders", "2026-05-06T13:32:00Z"))
296            .unwrap();
297
298        let recent = log
299            .query_since("orders", "2026-05-06T13:30:30Z", 100)
300            .unwrap();
301        assert_eq!(recent.len(), 2);
302        assert_eq!(recent[0].cursor, 2);
303    }
304
305    #[test]
306    fn query_respects_limit() {
307        let log = EventLog::open_in_memory().unwrap();
308        for i in 1..=10 {
309            let ts = format!("2026-05-06T13:{:02}:00Z", 30 + i);
310            log.append(&evt(i, "orders", &ts)).unwrap();
311        }
312        let result = log
313            .query_since("orders", "1970-01-01T00:00:00Z", 3)
314            .unwrap();
315        assert_eq!(result.len(), 3);
316        assert_eq!(result[0].cursor, 1);
317    }
318
319    #[test]
320    fn prune_drops_old_events_per_topic_policy() {
321        let log = EventLog::open_in_memory().unwrap();
322        // Old events (>90 days for orders / >14 days for marketdata)
323        let ancient = (chrono::Utc::now() - chrono::Duration::days(120))
324            .to_rfc3339_opts(chrono::SecondsFormat::Millis, true);
325        let mid = (chrono::Utc::now() - chrono::Duration::days(20))
326            .to_rfc3339_opts(chrono::SecondsFormat::Millis, true);
327        let recent = (chrono::Utc::now() - chrono::Duration::days(1))
328            .to_rfc3339_opts(chrono::SecondsFormat::Millis, true);
329
330        log.append(&evt(1, "orders", &ancient)).unwrap();
331        log.append(&evt(2, "orders", &mid)).unwrap();
332        log.append(&evt(3, "orders", &recent)).unwrap();
333        log.append(&evt(10, "marketdata:1", &mid)).unwrap();
334        log.append(&evt(11, "marketdata:1", &recent)).unwrap();
335
336        let policy = RetentionPolicy::default();
337        let dropped = log.prune(&policy).unwrap();
338        // orders: 120-day-old dropped; marketdata: 20-day-old dropped
339        assert_eq!(dropped, 2);
340
341        assert_eq!(log.count().unwrap(), 3);
342    }
343
344    #[test]
345    fn retention_policy_picks_per_topic_days() {
346        let p = RetentionPolicy::default();
347        assert_eq!(p.days_for("orders"), 90);
348        assert_eq!(p.days_for("pnl"), 90);
349        assert_eq!(p.days_for("gap"), 365);
350        assert_eq!(p.days_for("marketdata:265598"), 14);
351        assert_eq!(p.days_for("unknown"), 30);
352    }
353}