Skip to main content

teksilo_telemetry/queue/
persistent.rs

1// SPDX-License-Identifier: MPL-2.0
2// SPDX-FileCopyrightText: 2026 FernTech
3
4//! redb-backed `EventQueue` that survives process restart.
5//!
6//! Single table, key = monotonic `u64` (FIFO order = ascending key
7//! order), value = `serde_json` bytes of [`PersistedRecord`].
8//! Events are bounded by capacity (oldest dropped past the cap) and
9//! by age (events older than `max_age` are dropped during the next
10//! drain or push).
11//!
12//! Per-event retry metadata (`attempts`, `next_attempt_at_unix_ms`)
13//! is reserved in the schema for a future scheduled-retry feature
14//! but isn't currently consulted on drain — adapters retry the same
15//! batch using their own backoff state.
16
17use 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/// Wire shape persisted on disk. Versioned by tag for future schema
32/// migrations (just bump the variant; serde_json's untagged fallback
33/// makes additive changes painless).
34#[derive(Debug, Clone, Serialize, Deserialize)]
35struct PersistedRecord {
36    /// Schema version of *this record*. Currently always 1.
37    #[serde(default = "default_version")]
38    record_version: u32,
39    event: OwnedEvent,
40    enqueued_at_unix_ms: u64,
41    /// Reserved for future per-event scheduled-retry — currently
42    /// always 0. Adapters drive backoff from their own state today.
43    #[serde(default)]
44    attempts: u32,
45    /// Reserved for future per-event scheduled-retry — currently
46    /// always 0.
47    #[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
95/// redb-backed event queue. `Send + Sync` — internal `Mutex` wraps
96/// the `Database` handle so adapter UI-thread pushes and worker-
97/// thread drains don't conflict.
98pub struct PersistentEventQueue {
99    db: Mutex<Database>,
100    path: PathBuf,
101    capacity: usize,
102    max_age: Duration,
103}
104
105impl PersistentEventQueue {
106    /// Open or create the queue file at `path`. Parent directory is
107    /// created if missing.
108    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    /// Like [`open`](Self::open) but with explicit `capacity` (events
113    /// past which the oldest is dropped) and `max_age` (events older
114    /// than this are dropped on next push or drain).
115    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        // Initialize the tables so first-read transactions don't
126        // observe TableDoesNotExist.
127        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    /// Drop entries older than `max_age` and the oldest entries past
153    /// `capacity`. Called from `push` so the queue self-prunes.
154    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        // Collect keys to delete to avoid holding an iterator while
163        // mutating the table.
164        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                // Corrupt entry — drop it.
176                to_delete.push(key);
177            }
178        }
179        // Capacity: ensure post-insert size won't exceed `capacity`.
180        // We're called *before* the insert, so we need at most
181        // `capacity - 1` events after this eviction completes.
182        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            // Evict before insert to keep the queue bounded.
219            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            // Telemetry must never panic the host application — log
237            // and drop. Real adapters can poll a separate health
238            // signal if they care.
239            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                                // Synthesize a placeholder for corrupt
263                                // entries so the caller sees an event
264                                // count but the bad row gets removed.
265                                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            // Iterate in reverse order (newest first) by collecting
332            // and reversing — redb doesn't expose a reverse-iterator
333            // builder generically, but the table is small (<10k).
334            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            // Drop the queue without draining.
415        }
416        // Reopen and verify the events are still there.
417        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        // Insert, reopen, drain — order must match insertion order.
496        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}