1use std::path::{Path, PathBuf};
34use std::sync::Mutex;
35
36use rusqlite::{params, Connection, OptionalExtension};
37
38use super::ObservedEvent;
39
40pub struct EventLog {
44 conn: Mutex<Connection>,
45 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#[derive(Clone, Debug)]
59pub struct RetentionPolicy {
60 pub default_days: i64,
62 pub orders_days: i64,
64 pub pnl_days: i64,
66 pub marketdata_days: i64,
68 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 #[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 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 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 #[must_use]
154 pub fn path(&self) -> &Path {
155 &self.path
156 }
157
158 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 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 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 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 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 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}