Skip to main content

pitchfork_cli/log_store/
sqlite.rs

1use crate::Result;
2use crate::daemon_id::DaemonId;
3use crate::log_store::{
4    ArchiveHook, LogEntry, LogQuery, LogStore, MessageFilter, escape_like_pattern,
5};
6use chrono::{DateTime, Local, TimeZone};
7use log::error;
8use miette::IntoDiagnostic;
9use rusqlite::{Connection, OptionalExtension, params};
10use std::collections::HashSet;
11use std::io::{BufRead, BufReader, Write};
12use std::path::PathBuf;
13use std::sync::Mutex;
14
15/// Registers a `regexp` SQL function backed by the `regex` crate.
16///
17/// SQLite does not ship a REGEXP implementation by default. This function
18/// is registered on every new connection so that `message REGEXP ?` works
19/// consistently across queries.
20fn add_regexp_function(conn: &Connection) -> Result<()> {
21    use std::cell::RefCell;
22
23    // Cache compiled regexes per connection to avoid recompiling the same
24    // pattern on every row evaluation. The cache is small and thread-local
25    // because scalar functions are invoked on the connection's thread.
26    let cache: RefCell<lru::LruCache<String, regex::Regex>> =
27        RefCell::new(lru::LruCache::new(std::num::NonZeroUsize::new(32).unwrap()));
28
29    conn.create_scalar_function(
30        "regexp",
31        2,
32        rusqlite::functions::FunctionFlags::SQLITE_UTF8
33            | rusqlite::functions::FunctionFlags::SQLITE_DETERMINISTIC,
34        move |ctx| {
35            let pattern: String = ctx.get(0)?;
36            let text: String = ctx.get(1)?;
37
38            let mut cache = cache.borrow_mut();
39            let re = match cache.get(&pattern) {
40                Some(re) => re.clone(),
41                None => {
42                    let re = regex::Regex::new(&pattern)
43                        .map_err(|e| rusqlite::Error::UserFunctionError(e.to_string().into()))?;
44                    cache.put(pattern.clone(), re.clone());
45                    re
46                }
47            };
48            Ok(re.is_match(&text))
49        },
50    )
51    .into_diagnostic()
52}
53
54/// SQLite-backed log store with WAL mode for concurrent readers.
55pub struct SqliteLogStore {
56    conn: Mutex<Connection>,
57}
58
59impl SqliteLogStore {
60    /// Open or create the SQLite log store at the given path.
61    pub fn open(path: impl Into<PathBuf>) -> Result<Self> {
62        let path = path.into();
63        if let Some(parent) = path.parent() {
64            std::fs::create_dir_all(parent).into_diagnostic()?;
65        }
66        let conn = Connection::open(&path).into_diagnostic()?;
67        add_regexp_function(&conn)?;
68        conn.execute_batch(
69            "PRAGMA journal_mode = WAL;
70             PRAGMA synchronous = NORMAL;",
71        )
72        .into_diagnostic()?;
73        conn.execute(
74            "CREATE TABLE IF NOT EXISTS log_entries (
75                id          INTEGER PRIMARY KEY AUTOINCREMENT,
76                daemon_id   TEXT    NOT NULL,
77                timestamp   INTEGER NOT NULL,
78                message     TEXT    NOT NULL
79            );",
80            [],
81        )
82        .into_diagnostic()?;
83        conn.execute(
84            "CREATE INDEX IF NOT EXISTS idx_daemon_ts ON log_entries(daemon_id, timestamp);",
85            [],
86        )
87        .into_diagnostic()?;
88        conn.execute(
89            "CREATE INDEX IF NOT EXISTS idx_daemon_id ON log_entries(daemon_id, id);",
90            [],
91        )
92        .into_diagnostic()?;
93        conn.execute(
94            "CREATE INDEX IF NOT EXISTS idx_timestamp ON log_entries(timestamp);",
95            [],
96        )
97        .into_diagnostic()?;
98        conn.execute(
99            "CREATE TABLE IF NOT EXISTS log_clear_generations (
100                daemon_id TEXT PRIMARY KEY,
101                generation INTEGER NOT NULL DEFAULT 0
102            );",
103            [],
104        )
105        .into_diagnostic()?;
106        Ok(Self {
107            conn: Mutex::new(conn),
108        })
109    }
110
111    fn row_to_entry(row: &rusqlite::Row) -> rusqlite::Result<LogEntry> {
112        let id: i64 = row.get(0)?;
113        let daemon_id: String = row.get(1)?;
114        let ts_millis: i64 = row.get(2)?;
115        let message: String = row.get(3)?;
116        let timestamp = Local
117            .timestamp_millis_opt(ts_millis)
118            .single()
119            .unwrap_or_else(Local::now);
120        Ok(LogEntry {
121            id,
122            daemon_id,
123            timestamp,
124            message,
125        })
126    }
127
128    fn archive_entries(
129        &self,
130        entries: &[LogEntry],
131        archive_hook: &ArchiveHook,
132        daemon_id: &DaemonId,
133        reason: &str,
134    ) -> Result<()> {
135        use std::process::{Command, Stdio};
136
137        if entries.is_empty() {
138            return Ok(());
139        }
140
141        for chunk in entries.chunks(archive_hook.batch_size.max(1)) {
142            let mut child = Command::new("sh")
143                .arg("-c")
144                .arg(&archive_hook.command)
145                .stdin(Stdio::piped())
146                .stdout(Stdio::null())
147                .stderr(Stdio::piped())
148                .env("PITCHFORK_DAEMON_ID", daemon_id.qualified())
149                .env("PITCHFORK_ARCHIVE_REASON", reason)
150                .spawn()
151                .into_diagnostic()
152                .map_err(|e| miette::miette!("failed to spawn archive hook: {e}"))?;
153
154            // Write entries to stdin. On error, kill and reap the child
155            // before propagating — Child::drop does not call wait(), so
156            // without this a crashed hook would leave a zombie process.
157            let write_result = {
158                let stdin = child.stdin.take().expect("piped stdin should be available");
159                let mut stdin = std::io::BufWriter::new(stdin);
160                let mut result = Ok(());
161                for entry in chunk {
162                    let line = serde_json::json!({
163                        "id": entry.id,
164                        "daemon_id": entry.daemon_id,
165                        "timestamp": entry.timestamp.to_rfc3339(),
166                        "message": entry.message,
167                    });
168                    if let Err(e) = writeln!(stdin, "{}", line) {
169                        result = Err(miette::miette!(
170                            "failed to write to archive hook stdin: {e}"
171                        ));
172                        break;
173                    }
174                }
175                // Explicitly flush so a buffer-drain failure (e.g. the hook
176                // exited early and closed stdin) is surfaced as an error
177                // rather than being silently swallowed by BufWriter::drop.
178                // On a small batch where all data fits in the 8 KB buffer no
179                // individual writeln! has touched the underlying pipe yet, so
180                // without this flush the failure would be lost and the entries
181                // deleted without ever being delivered to the hook.
182                if result.is_ok() {
183                    if let Err(e) = stdin.flush() {
184                        result = Err(miette::miette!("failed to flush archive hook stdin: {e}"));
185                    }
186                }
187                result
188                // BufWriter + ChildStdin drop here, closing stdin (EOF signal).
189            };
190
191            if let Err(e) = write_result {
192                let _ = child.kill();
193                let _ = child.wait();
194                return Err(e);
195            }
196
197            let output = child.wait_with_output().into_diagnostic()?;
198            if !output.status.success() {
199                let stderr = String::from_utf8_lossy(&output.stderr);
200                return Err(miette::miette!(
201                    "archive hook failed with status {}: {stderr}",
202                    output.status
203                ));
204            }
205        }
206
207        Ok(())
208    }
209
210    /// Delete rows matching the given IDs, returning the number of rows deleted.
211    ///
212    /// Chunks at 999 to stay within SQLite's `SQLITE_MAX_VARIABLE_NUMBER`
213    /// limit on versions ≤ 3.31 (e.g. Ubuntu 20.04 LTS).
214    fn delete_by_ids(&self, ids: &[i64]) -> Result<u64> {
215        const SQLITE_MAX_VARS: usize = 999;
216
217        let mut total = 0u64;
218        let conn = self.conn.lock().unwrap();
219        for chunk in ids.chunks(SQLITE_MAX_VARS) {
220            if chunk.is_empty() {
221                continue;
222            }
223            let placeholders: Vec<String> = (1..=chunk.len()).map(|i| format!("?{i}")).collect();
224            let sql = format!(
225                "DELETE FROM log_entries WHERE id IN ({})",
226                placeholders.join(", ")
227            );
228            total += conn
229                .execute(&sql, rusqlite::params_from_iter(chunk.iter()))
230                .into_diagnostic()? as u64;
231        }
232        Ok(total)
233    }
234
235    /// Rotate (delete) old log entries for a specific daemon based on retention policy.
236    ///
237    /// When an archive hook is configured, entries are fetched in batches of
238    /// `batch_size`, the hook is invoked *without* holding the SQLite mutex
239    /// (so concurrent log appends are not blocked), and then deleted by ID.
240    /// Row-read errors are propagated rather than silently dropped.
241    pub fn rotate_by_age(
242        &self,
243        daemon_id: &DaemonId,
244        max_age: chrono::Duration,
245        archive_hook: Option<&ArchiveHook>,
246    ) -> Result<u64> {
247        let cutoff = (Local::now() - max_age).timestamp_millis();
248        let hook = archive_hook.filter(|h| h.is_enabled());
249
250        if let Some(hook) = hook {
251            let mut total_deleted = 0u64;
252            loop {
253                // Fetch one batch under lock, then release.
254                let entries: Vec<LogEntry> = {
255                    let conn = self.conn.lock().unwrap();
256                    let mut stmt = conn
257                        .prepare(
258                            "SELECT id, daemon_id, timestamp, message FROM log_entries
259                             WHERE daemon_id = ?1 AND timestamp < ?2
260                             ORDER BY timestamp ASC, id ASC
261                             LIMIT ?3",
262                        )
263                        .into_diagnostic()?;
264                    stmt.query_map(
265                        params![daemon_id.qualified(), cutoff, hook.batch_size as i64],
266                        Self::row_to_entry,
267                    )
268                    .into_diagnostic()?
269                    .collect::<rusqlite::Result<Vec<_>>>()
270                    .into_diagnostic()?
271                };
272
273                if entries.is_empty() {
274                    break;
275                }
276
277                let batch_ids: Vec<i64> = entries.iter().map(|e| e.id).collect();
278
279                // Run the hook without holding the mutex.
280                self.archive_entries(&entries, hook, daemon_id, "age")?;
281
282                // Re-acquire lock and delete exactly the archived IDs.
283                let deleted = self.delete_by_ids(&batch_ids)?;
284                total_deleted += deleted;
285            }
286            Ok(total_deleted)
287        } else {
288            let conn = self.conn.lock().unwrap();
289            let rows = conn
290                .execute(
291                    "DELETE FROM log_entries WHERE daemon_id = ?1 AND timestamp < ?2",
292                    params![daemon_id.qualified(), cutoff],
293                )
294                .into_diagnostic()?;
295            Ok(rows as u64)
296        }
297    }
298
299    /// Rotate (delete) old log entries keeping only the most recent `max_count` rows
300    /// for a specific daemon.
301    ///
302    /// When an archive hook is configured, entries are fetched in batches of
303    /// `batch_size`, the hook is invoked *without* holding the SQLite mutex
304    /// (so concurrent log appends are not blocked), and then deleted by ID.
305    /// Row-read errors are propagated rather than silently dropped.
306    pub fn rotate_by_count(
307        &self,
308        daemon_id: &DaemonId,
309        max_count: u64,
310        archive_hook: Option<&ArchiveHook>,
311    ) -> Result<u64> {
312        let hook = archive_hook.filter(|h| h.is_enabled());
313
314        // Determine how many entries to delete.
315        let to_delete: i64 = {
316            let conn = self.conn.lock().unwrap();
317            let count: i64 = conn
318                .query_row(
319                    "SELECT COUNT(*) FROM log_entries WHERE daemon_id = ?1",
320                    [daemon_id.qualified()],
321                    |row| row.get(0),
322                )
323                .into_diagnostic()?;
324            count.saturating_sub(max_count as i64)
325        };
326
327        if to_delete <= 0 {
328            return Ok(0);
329        }
330
331        if let Some(hook) = hook {
332            let mut total_deleted = 0u64;
333            let mut remaining = to_delete;
334            loop {
335                let batch_len = remaining.min(hook.batch_size as i64);
336
337                // Fetch one batch under lock, then release.
338                let entries: Vec<LogEntry> = {
339                    let conn = self.conn.lock().unwrap();
340                    let mut stmt = conn
341                        .prepare(
342                            "SELECT id, daemon_id, timestamp, message FROM log_entries
343                             WHERE daemon_id = ?1
344                             ORDER BY timestamp ASC, id ASC
345                             LIMIT ?2",
346                        )
347                        .into_diagnostic()?;
348                    stmt.query_map(
349                        params![daemon_id.qualified(), batch_len],
350                        Self::row_to_entry,
351                    )
352                    .into_diagnostic()?
353                    .collect::<rusqlite::Result<Vec<_>>>()
354                    .into_diagnostic()?
355                };
356
357                if entries.is_empty() {
358                    break;
359                }
360
361                let fetched = entries.len() as i64;
362                let batch_ids: Vec<i64> = entries.iter().map(|e| e.id).collect();
363
364                // Run the hook without holding the mutex.
365                self.archive_entries(&entries, hook, daemon_id, "count")?;
366
367                // Re-acquire lock and delete exactly the archived IDs.
368                let deleted = self.delete_by_ids(&batch_ids)?;
369                total_deleted += deleted;
370                remaining -= fetched;
371            }
372            Ok(total_deleted)
373        } else {
374            let conn = self.conn.lock().unwrap();
375            let rows = conn
376                .execute(
377                    "DELETE FROM log_entries WHERE id IN (
378                        SELECT id FROM log_entries WHERE daemon_id = ?1
379                        ORDER BY timestamp ASC, id ASC LIMIT ?2
380                    )",
381                    params![daemon_id.qualified(), to_delete],
382                )
383                .into_diagnostic()?;
384            Ok(rows as u64)
385        }
386    }
387
388    /// Migrate existing text logs for a daemon into SQLite.
389    ///
390    /// Reads the legacy text file line-by-line (streaming) and inserts in
391    /// batches of 1000 to avoid loading multi-GB files into memory at once.
392    pub fn migrate_daemon_text_logs(&self, daemon_id: &DaemonId) -> Result<u64> {
393        let text_path = daemon_id.log_path();
394        if !text_path.exists() {
395            return Ok(0);
396        }
397
398        let file = std::fs::File::open(&text_path).into_diagnostic()?;
399        let reader = BufReader::new(file);
400        let re = regex::Regex::new(r"^(\d{4}-\d{2}-\d{2} \d{2}:\d{2}:\d{2}) ([\w./-]+) (.*)$")
401            .expect("invalid regex");
402
403        let mut current_timestamp: Option<DateTime<Local>> = None;
404        let mut current_message = String::new();
405        let mut entries = Vec::with_capacity(1000);
406        let mut total_migrated: u64 = 0;
407
408        for line in reader.lines() {
409            let line = line.into_diagnostic()?;
410            if let Some(caps) = re.captures(&line) {
411                if let Some(ts) = current_timestamp.take() {
412                    entries.push((ts, std::mem::take(&mut current_message)));
413                }
414                let ts_str = caps.get(1).map(|m| m.as_str()).unwrap_or_default();
415                let msg = caps.get(3).map(|m| m.as_str()).unwrap_or_default();
416                if let Ok(naive) =
417                    chrono::NaiveDateTime::parse_from_str(ts_str, "%Y-%m-%d %H:%M:%S")
418                {
419                    current_timestamp = Local.from_local_datetime(&naive).single();
420                    current_message = msg.to_string();
421                }
422            } else if current_timestamp.is_some() {
423                current_message.push('\n');
424                current_message.push_str(&line);
425            }
426
427            if entries.len() >= 1000 {
428                total_migrated += self.insert_batch(daemon_id, &entries)?;
429                entries.clear();
430            }
431        }
432
433        if let Some(ts) = current_timestamp {
434            entries.push((ts, std::mem::take(&mut current_message)));
435        }
436
437        if !entries.is_empty() {
438            total_migrated += self.insert_batch(daemon_id, &entries)?;
439        }
440
441        if total_migrated > 0 {
442            if let Err(e) = std::fs::remove_file(&text_path) {
443                log::warn!(
444                    "failed to remove legacy log file after migration {}: {e}",
445                    text_path.display()
446                );
447            }
448        }
449
450        Ok(total_migrated)
451    }
452
453    fn insert_batch(
454        &self,
455        daemon_id: &DaemonId,
456        entries: &[(DateTime<Local>, String)],
457    ) -> Result<u64> {
458        let mut conn = self.conn.lock().unwrap();
459        let tx = conn.transaction().into_diagnostic()?;
460        let mut count = 0u64;
461        {
462            let mut stmt = tx
463                .prepare(
464                    "INSERT INTO log_entries (daemon_id, timestamp, message) VALUES (?1, ?2, ?3)",
465                )
466                .into_diagnostic()?;
467            for (ts, msg) in entries {
468                stmt.execute(params![daemon_id.qualified(), ts.timestamp_millis(), msg])
469                    .into_diagnostic()?;
470                count += 1;
471            }
472        }
473        tx.commit().into_diagnostic()?;
474        Ok(count)
475    }
476}
477
478impl LogStore for SqliteLogStore {
479    fn append(&self, daemon_id: &DaemonId, message: &str) -> Result<()> {
480        let ts = Local::now().timestamp_millis();
481        let id = daemon_id.qualified();
482        let msg = message.to_string();
483
484        let conn = self.conn.lock().unwrap();
485        let _ = conn
486            .execute(
487                "INSERT INTO log_entries (daemon_id, timestamp, message) VALUES (?1, ?2, ?3)",
488                params![id, ts, msg],
489            )
490            .into_diagnostic()?;
491        Ok(())
492    }
493
494    fn append_batch(&self, daemon_id: &DaemonId, messages: &[String]) -> Result<()> {
495        if messages.is_empty() {
496            return Ok(());
497        }
498        let base_ts = Local::now().timestamp_millis();
499        let id = daemon_id.qualified();
500
501        let mut conn = self.conn.lock().unwrap();
502        let tx = conn.transaction().into_diagnostic()?;
503        {
504            let mut stmt = tx
505                .prepare(
506                    "INSERT INTO log_entries (daemon_id, timestamp, message) VALUES (?1, ?2, ?3)",
507                )
508                .into_diagnostic()?;
509            for (idx, msg) in messages.iter().enumerate() {
510                // Slightly stagger timestamps within a batch so ordering by
511                // (timestamp, id) preserves insertion order without paying for
512                // a separate per-row clock read.
513                let ts = base_ts + idx as i64;
514                stmt.execute(params![id, ts, msg]).into_diagnostic()?;
515            }
516        }
517        tx.commit().into_diagnostic()?;
518        Ok(())
519    }
520
521    fn query(&self, opts: &LogQuery) -> Result<Vec<LogEntry>> {
522        let conn = self.conn.lock().unwrap();
523        let mut conditions = Vec::new();
524        let mut query_params: Vec<Box<dyn rusqlite::ToSql>> = Vec::new();
525
526        if !opts.daemon_ids.is_empty() {
527            let placeholders: Vec<String> = (1..=opts.daemon_ids.len())
528                .map(|i| format!("?{}", i))
529                .collect();
530            conditions.push(format!("daemon_id IN ({})", placeholders.join(", ")));
531            for id in &opts.daemon_ids {
532                query_params.push(Box::new(id.clone()));
533            }
534        }
535
536        if let Some(from) = opts.from {
537            conditions.push(format!("timestamp >= ?{}", query_params.len() + 1));
538            query_params.push(Box::new(from.timestamp_millis()));
539        }
540
541        if let Some(to) = opts.to {
542            conditions.push(format!("timestamp <= ?{}", query_params.len() + 1));
543            query_params.push(Box::new(to.timestamp_millis()));
544        }
545
546        if let Some(after_id) = opts.after_id {
547            conditions.push(format!("id > ?{}", query_params.len() + 1));
548            query_params.push(Box::new(after_id));
549        }
550
551        let mut message_conditions = Vec::new();
552        for filter in &opts.message_filters {
553            match filter {
554                MessageFilter::Contains {
555                    pattern,
556                    case_sensitive,
557                } => {
558                    let param_index = query_params.len() + 1;
559                    if *case_sensitive {
560                        // SQLite LIKE is ASCII-case-insensitive by default, so use
561                        // INSTR to enforce literal case-sensitive substring matches.
562                        // INSTR performs a plain byte-level substring search, so no
563                        // LIKE escaping or % delimiters should be added.
564                        message_conditions
565                            .push(format!("INSTR(message, ?{idx}) > 0", idx = param_index));
566                        query_params.push(Box::new(pattern.clone()));
567                    } else {
568                        let escaped = escape_like_pattern(pattern);
569                        let param = format!("%{}%", escaped);
570                        message_conditions.push(format!(
571                            "LOWER(message) LIKE LOWER(?{idx}) ESCAPE '\\'",
572                            idx = param_index
573                        ));
574                        query_params.push(Box::new(param));
575                    }
576                }
577                MessageFilter::Regex { pattern } => {
578                    let param_index = query_params.len() + 1;
579                    message_conditions.push(format!("message REGEXP ?{param_index}"));
580                    query_params.push(Box::new(pattern.clone()));
581                }
582            }
583        }
584        if !message_conditions.is_empty() {
585            conditions.push(format!("({})", message_conditions.join(" OR ")));
586        }
587
588        let where_clause = if conditions.is_empty() {
589            String::new()
590        } else {
591            format!("WHERE {}", conditions.join(" AND "))
592        };
593
594        let order = if opts.order_desc { "DESC" } else { "ASC" };
595
596        let limit_clause = opts
597            .limit
598            .map(|n| format!("LIMIT {}", n))
599            .unwrap_or_default();
600
601        let sql = format!(
602            "SELECT id, daemon_id, timestamp, message FROM log_entries {} ORDER BY timestamp {}, id {} {}",
603            where_clause, order, order, limit_clause
604        );
605
606        let mut stmt = conn.prepare(&sql).into_diagnostic()?;
607        let params_ref: Vec<&dyn rusqlite::ToSql> =
608            query_params.iter().map(|p| p.as_ref()).collect();
609        let rows = stmt
610            .query_map(params_ref.as_slice(), |row| {
611                let id: i64 = row.get(0)?;
612                let daemon_id: String = row.get(1)?;
613                let ts_millis: i64 = row.get(2)?;
614                let message: String = row.get(3)?;
615                let timestamp = Local
616                    .timestamp_millis_opt(ts_millis)
617                    .single()
618                    .unwrap_or_else(Local::now);
619                Ok(LogEntry {
620                    id,
621                    daemon_id,
622                    timestamp,
623                    message,
624                })
625            })
626            .into_diagnostic()?;
627
628        let mut entries = Vec::new();
629        for row in rows {
630            entries.push(row.into_diagnostic()?);
631        }
632        Ok(entries)
633    }
634
635    fn tail(&self, daemon_id: &DaemonId, after_id: Option<i64>) -> Result<Vec<LogEntry>> {
636        self.query(&LogQuery {
637            daemon_ids: vec![daemon_id.qualified()],
638            from: None,
639            to: None,
640            limit: None,
641            order_desc: false,
642            after_id,
643            message_filters: Vec::new(),
644        })
645    }
646
647    fn clear(&self, daemon_ids: &[DaemonId]) -> Result<()> {
648        let mut conn = self.conn.lock().unwrap();
649        let tx = conn.transaction().into_diagnostic()?;
650        for id in daemon_ids {
651            tx.execute(
652                "DELETE FROM log_entries WHERE daemon_id = ?1",
653                params![id.qualified()],
654            )
655            .into_diagnostic()?;
656            tx.execute(
657                "INSERT INTO log_clear_generations (daemon_id, generation)
658                 VALUES (?1, 1)
659                 ON CONFLICT(daemon_id) DO UPDATE SET generation = generation + 1",
660                params![id.qualified()],
661            )
662            .into_diagnostic()?;
663        }
664        tx.commit().into_diagnostic()?;
665        Ok(())
666    }
667
668    fn last_id(&self, daemon_id: &DaemonId) -> Result<Option<i64>> {
669        let conn = self.conn.lock().unwrap();
670        // MAX(id) returns NULL when no rows exist for the daemon; the
671        // Option<i64> decode maps NULL to None automatically.
672        let id: Option<i64> = conn
673            .query_row(
674                "SELECT MAX(id) FROM log_entries WHERE daemon_id = ?1",
675                params![daemon_id.qualified()],
676                |row| row.get(0),
677            )
678            .into_diagnostic()?;
679        Ok(id)
680    }
681
682    fn list_daemon_ids(&self) -> Result<Vec<String>> {
683        let conn = self.conn.lock().unwrap();
684        let mut stmt = conn
685            .prepare("SELECT DISTINCT daemon_id FROM log_entries")
686            .into_diagnostic()?;
687        let ids = stmt
688            .query_map([], |row| {
689                let id: String = row.get(0)?;
690                Ok(id)
691            })
692            .into_diagnostic()?
693            .filter_map(|r| r.ok())
694            .collect();
695        Ok(ids)
696    }
697
698    fn apply_retention(
699        &self,
700        policy: &super::RetentionPolicy,
701        excluded_daemon_ids: &[DaemonId],
702        archive_hook: Option<&ArchiveHook>,
703    ) -> Result<u64> {
704        let daemon_ids = self.list_daemon_ids()?;
705        let excluded: HashSet<String> = excluded_daemon_ids.iter().map(|d| d.qualified()).collect();
706        let mut total = 0u64;
707        for id_str in daemon_ids {
708            if excluded.contains(&id_str) {
709                continue;
710            }
711            let id = DaemonId::parse(&id_str).unwrap_or_else(|_| {
712                DaemonId::try_new("global", &id_str).unwrap_or_else(|_| DaemonId::pitchfork())
713            });
714            if let Some(dur) = policy.age {
715                total += self.rotate_by_age(&id, dur, archive_hook)?;
716            }
717            if let Some(n) = policy.count {
718                total += self.rotate_by_count(&id, n, archive_hook)?;
719            }
720        }
721        Ok(total)
722    }
723
724    fn apply_retention_for_daemon(
725        &self,
726        daemon_id: &DaemonId,
727        policy: &super::RetentionPolicy,
728        archive_hook: Option<&ArchiveHook>,
729    ) -> Result<u64> {
730        let mut total = 0u64;
731        if let Some(dur) = policy.age {
732            total += self.rotate_by_age(daemon_id, dur, archive_hook)?;
733        }
734        if let Some(n) = policy.count {
735            total += self.rotate_by_count(daemon_id, n, archive_hook)?;
736        }
737        Ok(total)
738    }
739
740    fn last_clear_generation(&self, daemon_id: &DaemonId) -> Result<Option<u64>> {
741        let conn = self.conn.lock().unwrap();
742        let mut stmt = conn
743            .prepare("SELECT generation FROM log_clear_generations WHERE daemon_id = ?1")
744            .into_diagnostic()?;
745        let generation: Option<i64> = stmt
746            .query_row(params![daemon_id.qualified()], |row| row.get(0))
747            .optional()
748            .into_diagnostic()?;
749        generation
750            .map(|generation| {
751                u64::try_from(generation)
752                    .map_err(|_| miette::miette!("log clear generation cannot be negative"))
753            })
754            .transpose()
755    }
756}
757
758/// Global singleton log store.
759use once_cell::sync::Lazy;
760use std::sync::Arc;
761
762pub static LOG_STORE: Lazy<Arc<SqliteLogStore>> = Lazy::new(|| {
763    let path = crate::env::PITCHFORK_LOGS_DIR.join("logs.db");
764    let mut is_fallback = false;
765    let store = Arc::new(SqliteLogStore::open(&path).unwrap_or_else(|e| {
766        error!(
767            "failed to open log store at {}: {e}. Falling back to in-memory store; logs will not persist across restarts.",
768            path.display()
769        );
770        is_fallback = true;
771        SqliteLogStore::open(":memory:").expect("in-memory SQLite should always open")
772    }));
773
774    // Auto-migrate any legacy text log files into SQLite on first access.
775    // This runs once per process startup and is idempotent.
776    // Skip migration when using the in-memory fallback to prevent data loss.
777    if !is_fallback {
778        if let Err(e) = auto_migrate_legacy_logs(&store) {
779            warn!("legacy log auto-migration failed: {e}");
780        }
781    } else {
782        warn!(
783            "skipping legacy log auto-migration because log store is in-memory (no durable destination)"
784        );
785    }
786
787    store
788});
789
790/// Auto-migrate legacy text log files into the SQLite log store.
791///
792/// Scans the logs directory for directories matching the new-format layout
793/// (`namespace--name/namespace--name.log`), attempts to parse the directory
794/// name as a valid safe-path daemon ID, and imports the content into SQLite.
795/// This is idempotent: re-running it on already-migrated data is a no-op
796/// because the legacy text files are deleted after successful import.
797fn auto_migrate_legacy_logs(store: &SqliteLogStore) -> Result<()> {
798    let logs_dir = &*crate::env::PITCHFORK_LOGS_DIR;
799    if !logs_dir.exists() {
800        return Ok(());
801    }
802
803    let Ok(entries) = std::fs::read_dir(logs_dir) else {
804        return Ok(());
805    };
806
807    let mut total_migrated = 0u64;
808    let mut migrated_ids = Vec::new();
809
810    for entry in entries.flatten() {
811        let path = entry.path();
812        if !path.is_dir() {
813            continue;
814        }
815        // Skip the supervisor's own log directory.
816        let file_name = path
817            .file_name()
818            .map_or(String::new(), |n| n.to_string_lossy().to_string());
819        if file_name == "pitchfork" {
820            continue;
821        }
822
823        // Only consider directories that look like new-format safe-paths
824        if !file_name.contains("--") {
825            continue;
826        }
827        let log_file = path.join(format!("{file_name}.log"));
828        if !log_file.exists() {
829            continue;
830        }
831
832        let daemon_id = match DaemonId::from_safe_path(&file_name) {
833            Ok(id) => id,
834            Err(_) => continue,
835        };
836
837        // Skip the supervisor's own daemon; its log directory may be present
838        // under logs/ but should never be imported into the user-facing log store.
839        if daemon_id == DaemonId::pitchfork() {
840            continue;
841        }
842
843        match store.migrate_daemon_text_logs(&daemon_id) {
844            Ok(0) => {}
845            Ok(n) => {
846                total_migrated += n;
847                migrated_ids.push(daemon_id.qualified());
848            }
849            Err(e) => {
850                warn!(
851                    "failed to migrate text logs for {}: {e}",
852                    daemon_id.qualified()
853                );
854            }
855        }
856    }
857
858    if total_migrated > 0 {
859        warn!(
860            "auto-migrated {total_migrated} legacy log entries from {count} daemon(s): {ids}",
861            count = migrated_ids.len(),
862            ids = migrated_ids.join(", ")
863        );
864    }
865
866    Ok(())
867}