1use std::path::{Path, PathBuf};
18use std::sync::Mutex;
19use std::time::{Duration, SystemTime, UNIX_EPOCH};
20
21use redb::{Database, ReadableDatabase, ReadableTable, ReadableTableMetadata, TableDefinition};
22use serde::{Deserialize, Serialize};
23use teksilo_core::telemetry::OwnedEvent;
24
25use super::EventQueue;
26
27const EVENTS: TableDefinition<u64, &[u8]> = TableDefinition::new("events");
28const NEXT_ID: TableDefinition<&str, u64> = TableDefinition::new("next_id");
29const NEXT_ID_KEY: &str = "next";
30
31#[derive(Debug, Clone, Serialize, Deserialize)]
35struct PersistedRecord {
36 #[serde(default = "default_version")]
38 record_version: u32,
39 event: OwnedEvent,
40 enqueued_at_unix_ms: u64,
41 #[serde(default)]
44 attempts: u32,
45 #[serde(default)]
48 next_attempt_at_unix_ms: u64,
49}
50
51fn default_version() -> u32 {
52 1
53}
54
55#[derive(Debug, thiserror::Error)]
56pub enum PersistentQueueError {
57 #[error("i/o error: {0}")]
58 Io(#[from] std::io::Error),
59 #[error("redb error: {0}")]
60 Database(#[from] redb::Error),
61 #[error("serialize error: {0}")]
62 Serialize(#[from] serde_json::Error),
63}
64
65impl From<redb::DatabaseError> for PersistentQueueError {
66 fn from(e: redb::DatabaseError) -> Self {
67 Self::Database(e.into())
68 }
69}
70
71impl From<redb::TransactionError> for PersistentQueueError {
72 fn from(e: redb::TransactionError) -> Self {
73 Self::Database(e.into())
74 }
75}
76
77impl From<redb::TableError> for PersistentQueueError {
78 fn from(e: redb::TableError) -> Self {
79 Self::Database(e.into())
80 }
81}
82
83impl From<redb::StorageError> for PersistentQueueError {
84 fn from(e: redb::StorageError) -> Self {
85 Self::Database(e.into())
86 }
87}
88
89impl From<redb::CommitError> for PersistentQueueError {
90 fn from(e: redb::CommitError) -> Self {
91 Self::Database(e.into())
92 }
93}
94
95pub struct PersistentEventQueue {
99 db: Mutex<Database>,
100 path: PathBuf,
101 capacity: usize,
102 max_age: Duration,
103}
104
105impl PersistentEventQueue {
106 pub fn open(path: impl Into<PathBuf>) -> Result<Self, PersistentQueueError> {
109 Self::open_with(path, 10_000, Duration::from_secs(60 * 60 * 24 * 7))
110 }
111
112 pub fn open_with(
116 path: impl Into<PathBuf>,
117 capacity: usize,
118 max_age: Duration,
119 ) -> Result<Self, PersistentQueueError> {
120 let path = path.into();
121 if let Some(parent) = path.parent() {
122 std::fs::create_dir_all(parent)?;
123 }
124 let db = Database::create(&path)?;
125 let write_txn = db.begin_write()?;
128 {
129 let _ = write_txn.open_table(EVENTS)?;
130 let _ = write_txn.open_table(NEXT_ID)?;
131 }
132 write_txn.commit()?;
133 Ok(Self {
134 db: Mutex::new(db),
135 path,
136 capacity: capacity.max(1),
137 max_age,
138 })
139 }
140
141 pub fn path(&self) -> &Path {
142 &self.path
143 }
144
145 fn now_unix_ms() -> u64 {
146 SystemTime::now()
147 .duration_since(UNIX_EPOCH)
148 .map(|d| d.as_millis() as u64)
149 .unwrap_or(0)
150 }
151
152 fn evict_locked(
155 write_txn: &redb::WriteTransaction,
156 capacity: usize,
157 max_age: Duration,
158 now_ms: u64,
159 ) -> Result<(), PersistentQueueError> {
160 let mut events = write_txn.open_table(EVENTS)?;
161 let max_age_ms = max_age.as_millis() as u64;
162 let mut to_delete: Vec<u64> = Vec::new();
165 for entry in events.iter()? {
166 let (k, v) = entry?;
167 let key = k.value();
168 let bytes = v.value();
169 if let Ok(record) = serde_json::from_slice::<PersistedRecord>(bytes) {
170 let age_ms = now_ms.saturating_sub(record.enqueued_at_unix_ms);
171 if age_ms >= max_age_ms {
172 to_delete.push(key);
173 }
174 } else {
175 to_delete.push(key);
177 }
178 }
179 let len_after_age = events.len()? as usize - to_delete.len();
183 if len_after_age + 1 > capacity {
184 let drop_n = len_after_age + 1 - capacity;
185 let surviving_keys: Vec<u64> = events
186 .iter()?
187 .filter_map(|e| e.ok().map(|(k, _)| k.value()))
188 .filter(|k| !to_delete.contains(k))
189 .collect();
190 for key in surviving_keys.into_iter().take(drop_n) {
191 to_delete.push(key);
192 }
193 }
194 for key in to_delete {
195 events.remove(key)?;
196 }
197 Ok(())
198 }
199
200 fn next_id_locked(write_txn: &redb::WriteTransaction) -> Result<u64, PersistentQueueError> {
201 let mut next_id_table = write_txn.open_table(NEXT_ID)?;
202 let current = next_id_table
203 .get(NEXT_ID_KEY)?
204 .map(|g| g.value())
205 .unwrap_or(0);
206 let next = current.wrapping_add(1);
207 next_id_table.insert(NEXT_ID_KEY, next)?;
208 Ok(next)
209 }
210}
211
212impl EventQueue for PersistentEventQueue {
213 fn push(&self, event: OwnedEvent) {
214 let db = self.db.lock().expect("queue db mutex poisoned");
215 if let Err(e) = (|| -> Result<(), PersistentQueueError> {
216 let now_ms = Self::now_unix_ms();
217 let write_txn = db.begin_write()?;
218 Self::evict_locked(&write_txn, self.capacity, self.max_age, now_ms)?;
220 let id = Self::next_id_locked(&write_txn)?;
221 {
222 let record = PersistedRecord {
223 record_version: 1,
224 event,
225 enqueued_at_unix_ms: now_ms,
226 attempts: 0,
227 next_attempt_at_unix_ms: 0,
228 };
229 let bytes = serde_json::to_vec(&record)?;
230 let mut events = write_txn.open_table(EVENTS)?;
231 events.insert(id, bytes.as_slice())?;
232 }
233 write_txn.commit()?;
234 Ok(())
235 })() {
236 eprintln!("teksilo-telemetry: persistent queue push failed: {e}");
240 }
241 }
242
243 fn drain_batch(&self, n: usize) -> Vec<OwnedEvent> {
244 let db = self.db.lock().expect("queue db mutex poisoned");
245 match (|| -> Result<Vec<OwnedEvent>, PersistentQueueError> {
246 let write_txn = db.begin_write()?;
247 let mut taken: Vec<(u64, OwnedEvent)> = Vec::new();
248 {
249 let events = write_txn.open_table(EVENTS)?;
250 for entry in events.iter()? {
251 if taken.len() >= n {
252 break;
253 }
254 let (k, v) = entry?;
255 let key = k.value();
256 let bytes = v.value();
257 match serde_json::from_slice::<PersistedRecord>(bytes) {
258 Ok(rec) => taken.push((key, rec.event)),
259 Err(_) => taken.push((
260 key,
261 OwnedEvent {
262 name: "telemetry.corrupt".into(),
266 category: teksilo_core::telemetry::EventCategory::Custom,
267 timestamp: SystemTime::UNIX_EPOCH,
268 install_id: None,
269 session_id: String::new(),
270 schema_version: 0,
271 props: vec![],
272 },
273 )),
274 }
275 }
276 }
277 {
278 let mut events = write_txn.open_table(EVENTS)?;
279 for (k, _) in &taken {
280 events.remove(*k)?;
281 }
282 }
283 write_txn.commit()?;
284 Ok(taken.into_iter().map(|(_, e)| e).collect())
285 })() {
286 Ok(events) => events,
287 Err(e) => {
288 eprintln!("teksilo-telemetry: persistent queue drain failed: {e}");
289 Vec::new()
290 }
291 }
292 }
293
294 fn len(&self) -> usize {
295 let db = self.db.lock().expect("queue db mutex poisoned");
296 (|| -> Result<usize, PersistentQueueError> {
297 let read_txn = db.begin_read()?;
298 let events = read_txn.open_table(EVENTS)?;
299 Ok(events.len()? as usize)
300 })()
301 .unwrap_or_default()
302 }
303
304 fn discard_all(&self) {
305 let db = self.db.lock().expect("queue db mutex poisoned");
306 if let Err(e) = (|| -> Result<(), PersistentQueueError> {
307 let write_txn = db.begin_write()?;
308 {
309 let mut events = write_txn.open_table(EVENTS)?;
310 let keys: Vec<u64> = events
311 .iter()?
312 .filter_map(|e| e.ok().map(|(k, _)| k.value()))
313 .collect();
314 for key in keys {
315 events.remove(key)?;
316 }
317 }
318 write_txn.commit()?;
319 Ok(())
320 })() {
321 eprintln!("teksilo-telemetry: persistent queue discard failed: {e}");
322 }
323 }
324
325 fn peek_recent(&self, n: usize) -> Vec<OwnedEvent> {
326 let db = self.db.lock().expect("queue db mutex poisoned");
327 (|| -> Result<Vec<OwnedEvent>, PersistentQueueError> {
328 let read_txn = db.begin_read()?;
329 let events = read_txn.open_table(EVENTS)?;
330 let mut out = Vec::with_capacity(n);
331 let mut all: Vec<(u64, Vec<u8>)> = Vec::new();
335 for entry in events.iter()? {
336 let (k, v) = entry?;
337 all.push((k.value(), v.value().to_vec()));
338 }
339 for (_, bytes) in all.into_iter().rev().take(n) {
340 if let Ok(rec) = serde_json::from_slice::<PersistedRecord>(&bytes) {
341 out.push(rec.event);
342 }
343 }
344 Ok(out)
345 })()
346 .unwrap_or_default()
347 }
348}
349
350impl std::fmt::Debug for PersistentEventQueue {
351 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
352 f.debug_struct("PersistentEventQueue")
353 .field("path", &self.path)
354 .field("capacity", &self.capacity)
355 .field("max_age", &self.max_age)
356 .field("len", &self.len())
357 .finish()
358 }
359}
360
361#[cfg(test)]
362mod tests {
363 use super::*;
364 use teksilo_core::telemetry::EventCategory;
365 use tempfile::tempdir;
366
367 fn ev(name: &str) -> OwnedEvent {
368 OwnedEvent {
369 name: name.to_string(),
370 category: EventCategory::Intent,
371 timestamp: SystemTime::UNIX_EPOCH,
372 install_id: None,
373 session_id: "test".into(),
374 schema_version: 1,
375 props: vec![],
376 }
377 }
378
379 #[test]
380 fn open_creates_file_and_starts_empty() {
381 let dir = tempdir().unwrap();
382 let path = dir.path().join("queue.redb");
383 let q = PersistentEventQueue::open(&path).unwrap();
384 assert!(path.exists());
385 assert_eq!(q.len(), 0);
386 assert!(q.is_empty());
387 }
388
389 #[test]
390 fn push_then_drain_round_trip() {
391 let dir = tempdir().unwrap();
392 let path = dir.path().join("q.redb");
393 let q = PersistentEventQueue::open(&path).unwrap();
394 q.push(ev("a"));
395 q.push(ev("b"));
396 q.push(ev("c"));
397 assert_eq!(q.len(), 3);
398 let batch = q.drain_batch(10);
399 assert_eq!(batch.len(), 3);
400 assert_eq!(batch[0].name, "a");
401 assert_eq!(batch[2].name, "c");
402 assert_eq!(q.len(), 0);
403 }
404
405 #[test]
406 fn events_survive_reopen() {
407 let dir = tempdir().unwrap();
408 let path = dir.path().join("survive.redb");
409 {
410 let q = PersistentEventQueue::open(&path).unwrap();
411 q.push(ev("event_1"));
412 q.push(ev("event_2"));
413 assert_eq!(q.len(), 2);
414 }
416 let q = PersistentEventQueue::open(&path).unwrap();
418 assert_eq!(q.len(), 2);
419 let batch = q.drain_batch(10);
420 assert_eq!(batch.len(), 2);
421 assert_eq!(batch[0].name, "event_1");
422 assert_eq!(batch[1].name, "event_2");
423 }
424
425 #[test]
426 fn capacity_evicts_oldest_on_push() {
427 let dir = tempdir().unwrap();
428 let path = dir.path().join("cap.redb");
429 let q = PersistentEventQueue::open_with(&path, 3, Duration::from_secs(3600)).unwrap();
430 q.push(ev("a"));
431 q.push(ev("b"));
432 q.push(ev("c"));
433 q.push(ev("d"));
434 q.push(ev("e"));
435 assert_eq!(q.len(), 3);
436 let batch = q.drain_batch(10);
437 let names: Vec<String> = batch.iter().map(|e| e.name.clone()).collect();
438 assert_eq!(names, vec!["c".to_string(), "d".into(), "e".into()]);
439 }
440
441 #[test]
442 fn discard_all_empties() {
443 let dir = tempdir().unwrap();
444 let path = dir.path().join("discard.redb");
445 let q = PersistentEventQueue::open(&path).unwrap();
446 q.push(ev("a"));
447 q.push(ev("b"));
448 q.discard_all();
449 assert!(q.is_empty());
450 }
451
452 #[test]
453 fn discard_persists_across_reopen() {
454 let dir = tempdir().unwrap();
455 let path = dir.path().join("discard2.redb");
456 {
457 let q = PersistentEventQueue::open(&path).unwrap();
458 q.push(ev("a"));
459 q.discard_all();
460 }
461 let q = PersistentEventQueue::open(&path).unwrap();
462 assert!(q.is_empty());
463 }
464
465 #[test]
466 fn peek_recent_returns_newest_first() {
467 let dir = tempdir().unwrap();
468 let path = dir.path().join("peek.redb");
469 let q = PersistentEventQueue::open(&path).unwrap();
470 q.push(ev("a"));
471 q.push(ev("b"));
472 q.push(ev("c"));
473 let recent = q.peek_recent(2);
474 assert_eq!(recent[0].name, "c");
475 assert_eq!(recent[1].name, "b");
476 }
477
478 #[test]
479 fn drain_batch_takes_only_up_to_n() {
480 let dir = tempdir().unwrap();
481 let path = dir.path().join("batch.redb");
482 let q = PersistentEventQueue::open(&path).unwrap();
483 for name in ["a", "b", "c", "d", "e"] {
484 q.push(ev(name));
485 }
486 let batch = q.drain_batch(2);
487 assert_eq!(batch.len(), 2);
488 assert_eq!(batch[0].name, "a");
489 assert_eq!(batch[1].name, "b");
490 assert_eq!(q.len(), 3);
491 }
492
493 #[test]
494 fn fifo_order_preserved_across_reopen() {
495 let dir = tempdir().unwrap();
497 let path = dir.path().join("fifo.redb");
498 {
499 let q = PersistentEventQueue::open(&path).unwrap();
500 for name in ["first", "second", "third", "fourth"] {
501 q.push(ev(name));
502 }
503 }
504 let q = PersistentEventQueue::open(&path).unwrap();
505 let batch = q.drain_batch(10);
506 let names: Vec<String> = batch.iter().map(|e| e.name.clone()).collect();
507 assert_eq!(
508 names,
509 vec![
510 "first".to_string(),
511 "second".into(),
512 "third".into(),
513 "fourth".into()
514 ]
515 );
516 }
517}