Skip to main content

pitchfork_cli/log_store/
sqlite.rs

1use crate::Result;
2use crate::daemon_id::DaemonId;
3use crate::log_parse::ParsedLog;
4use crate::log_store::{
5    ArchiveHook, FieldFilter, LogEntry, LogQuery, LogStore, MessageFilter, escape_like_pattern,
6};
7use chrono::{DateTime, Local, TimeZone};
8use log::error;
9use miette::IntoDiagnostic;
10use rusqlite::{Connection, OpenFlags, OptionalExtension, params};
11use std::collections::HashSet;
12use std::io::{BufRead, BufReader, Write};
13use std::path::PathBuf;
14use std::sync::Mutex;
15
16/// Registers a `regexp` SQL function backed by the `regex` crate.
17///
18/// SQLite does not ship a REGEXP implementation by default. This function
19/// is registered on every new connection so that `message REGEXP ?` works
20/// consistently across queries.
21fn add_regexp_function(conn: &Connection) -> Result<()> {
22    use std::cell::RefCell;
23
24    // Cache compiled regexes per connection to avoid recompiling the same
25    // pattern on every row evaluation. The cache is small and thread-local
26    // because scalar functions are invoked on the connection's thread.
27    let cache: RefCell<lru::LruCache<String, regex::Regex>> =
28        RefCell::new(lru::LruCache::new(std::num::NonZeroUsize::new(32).unwrap()));
29
30    conn.create_scalar_function(
31        "regexp",
32        2,
33        rusqlite::functions::FunctionFlags::SQLITE_UTF8
34            | rusqlite::functions::FunctionFlags::SQLITE_DETERMINISTIC,
35        move |ctx| {
36            let pattern: String = ctx.get(0)?;
37            let text: String = ctx.get(1)?;
38
39            let mut cache = cache.borrow_mut();
40            let re = match cache.get(&pattern) {
41                Some(re) => re.clone(),
42                None => {
43                    let re = regex::Regex::new(&pattern)
44                        .map_err(|e| rusqlite::Error::UserFunctionError(e.to_string().into()))?;
45                    cache.put(pattern.clone(), re.clone());
46                    re
47                }
48            };
49            Ok(re.is_match(&text))
50        },
51    )
52    .into_diagnostic()
53}
54
55/// SQLite-backed log store with WAL mode for concurrent readers.
56pub struct SqliteLogStore {
57    conn: Mutex<Connection>,
58    path: PathBuf,
59}
60
61/// Minimum result-set size to trigger parallel query.
62///
63/// Each parallel shard opens a new read-only SQLite connection (~5-12ms
64/// overhead per connection for PRAGMA setup + regexp function registration).
65/// Below this threshold, single-threaded is faster because the connection
66/// overhead exceeds the query savings.
67const PARALLEL_QUERY_THRESHOLD: usize = 200_000;
68
69impl SqliteLogStore {
70    /// Open or create the SQLite log store at the given path.
71    pub fn open(path: impl Into<PathBuf>) -> Result<Self> {
72        let path = path.into();
73        if let Some(parent) = path.parent() {
74            std::fs::create_dir_all(parent).into_diagnostic()?;
75        }
76        let conn = Connection::open(&path).into_diagnostic()?;
77        add_regexp_function(&conn)?;
78        conn.execute_batch(
79            "PRAGMA journal_mode = WAL;
80             PRAGMA synchronous = NORMAL;
81             PRAGMA mmap_size = 268435456;",
82        )
83        .into_diagnostic()?;
84        conn.execute(
85            "CREATE TABLE IF NOT EXISTS log_entries (
86                id          INTEGER PRIMARY KEY AUTOINCREMENT,
87                daemon_id   TEXT    NOT NULL,
88                timestamp   INTEGER NOT NULL,
89                message     TEXT    NOT NULL,
90                level       TEXT,
91                msg         TEXT,
92                logger      TEXT,
93                fields_json TEXT
94            );",
95            [],
96        )
97        .into_diagnostic()?;
98        conn.execute(
99            "CREATE INDEX IF NOT EXISTS idx_daemon_ts ON log_entries(daemon_id, timestamp);",
100            [],
101        )
102        .into_diagnostic()?;
103        conn.execute(
104            "CREATE INDEX IF NOT EXISTS idx_daemon_id ON log_entries(daemon_id, id);",
105            [],
106        )
107        .into_diagnostic()?;
108        conn.execute(
109            "CREATE INDEX IF NOT EXISTS idx_timestamp ON log_entries(timestamp);",
110            [],
111        )
112        .into_diagnostic()?;
113
114        // Migrate existing tables: add columns introduced in this version.
115        // Must run BEFORE creating indexes that reference the new columns.
116        let existing_cols: Vec<String> = {
117            let mut stmt = conn
118                .prepare("PRAGMA table_info(log_entries)")
119                .into_diagnostic()?;
120            let rows = stmt
121                .query_map([], |row| row.get::<_, String>(1))
122                .into_diagnostic()?;
123            rows.filter_map(|r| r.ok()).collect()
124        };
125        for col in ["level", "msg", "logger", "fields_json"] {
126            if !existing_cols.iter().any(|c| c == col) {
127                conn.execute(
128                    &format!("ALTER TABLE log_entries ADD COLUMN {col} TEXT"),
129                    [],
130                )
131                .into_diagnostic()?;
132            }
133        }
134        conn.execute(
135            "CREATE INDEX IF NOT EXISTS idx_daemon_level_ts ON log_entries(daemon_id, level, timestamp);",
136            [],
137        )
138        .into_diagnostic()?;
139
140        conn.execute(
141            "CREATE TABLE IF NOT EXISTS log_clear_generations (
142                daemon_id TEXT PRIMARY KEY,
143                generation INTEGER NOT NULL DEFAULT 0
144            );",
145            [],
146        )
147        .into_diagnostic()?;
148        Ok(Self {
149            conn: Mutex::new(conn),
150            path,
151        })
152    }
153
154    fn row_to_entry(row: &rusqlite::Row) -> rusqlite::Result<LogEntry> {
155        let id: i64 = row.get(0)?;
156        let daemon_id: String = row.get(1)?;
157        let ts_millis: i64 = row.get(2)?;
158        let message: String = row.get(3)?;
159        let level: Option<String> = row.get(4)?;
160        let msg: Option<String> = row.get(5)?;
161        let logger: Option<String> = row.get(6)?;
162        let fields_json: Option<String> = row.get(7)?;
163        let timestamp = Local
164            .timestamp_millis_opt(ts_millis)
165            .single()
166            .unwrap_or_else(Local::now);
167        Ok(LogEntry {
168            id,
169            daemon_id,
170            timestamp,
171            message,
172            level,
173            msg,
174            logger,
175            fields_json,
176        })
177    }
178
179    fn archive_entries(
180        &self,
181        entries: &[LogEntry],
182        archive_hook: &ArchiveHook,
183        daemon_id: &DaemonId,
184        reason: &str,
185    ) -> Result<()> {
186        use std::process::{Command, Stdio};
187
188        if entries.is_empty() {
189            return Ok(());
190        }
191
192        for chunk in entries.chunks(archive_hook.batch_size.max(1)) {
193            let mut child = Command::new("sh")
194                .arg("-c")
195                .arg(&archive_hook.command)
196                .stdin(Stdio::piped())
197                .stdout(Stdio::null())
198                .stderr(Stdio::piped())
199                .env("PITCHFORK_DAEMON_ID", daemon_id.qualified())
200                .env("PITCHFORK_ARCHIVE_REASON", reason)
201                .spawn()
202                .into_diagnostic()
203                .map_err(|e| miette::miette!("failed to spawn archive hook: {e}"))?;
204
205            // Write entries to stdin. On error, kill and reap the child
206            // before propagating — Child::drop does not call wait(), so
207            // without this a crashed hook would leave a zombie process.
208            let write_result = {
209                let stdin = child.stdin.take().expect("piped stdin should be available");
210                let mut stdin = std::io::BufWriter::new(stdin);
211                let mut result = Ok(());
212                for entry in chunk {
213                    let line = serde_json::json!({
214                        "id": entry.id,
215                        "daemon_id": entry.daemon_id,
216                        "timestamp": entry.timestamp.to_rfc3339(),
217                        "message": entry.message,
218                    });
219                    if let Err(e) = writeln!(stdin, "{}", line) {
220                        result = Err(miette::miette!(
221                            "failed to write to archive hook stdin: {e}"
222                        ));
223                        break;
224                    }
225                }
226                // Explicitly flush so a buffer-drain failure (e.g. the hook
227                // exited early and closed stdin) is surfaced as an error
228                // rather than being silently swallowed by BufWriter::drop.
229                // On a small batch where all data fits in the 8 KB buffer no
230                // individual writeln! has touched the underlying pipe yet, so
231                // without this flush the failure would be lost and the entries
232                // deleted without ever being delivered to the hook.
233                if result.is_ok() {
234                    if let Err(e) = stdin.flush() {
235                        result = Err(miette::miette!("failed to flush archive hook stdin: {e}"));
236                    }
237                }
238                result
239                // BufWriter + ChildStdin drop here, closing stdin (EOF signal).
240            };
241
242            if let Err(e) = write_result {
243                let _ = child.kill();
244                let _ = child.wait();
245                return Err(e);
246            }
247
248            let output = child.wait_with_output().into_diagnostic()?;
249            if !output.status.success() {
250                let stderr = String::from_utf8_lossy(&output.stderr);
251                return Err(miette::miette!(
252                    "archive hook failed with status {}: {stderr}",
253                    output.status
254                ));
255            }
256        }
257
258        Ok(())
259    }
260
261    /// Delete rows matching the given IDs, returning the number of rows deleted.
262    ///
263    /// Chunks at 999 to stay within SQLite's `SQLITE_MAX_VARIABLE_NUMBER`
264    /// limit on versions ≤ 3.31 (e.g. Ubuntu 20.04 LTS).
265    fn delete_by_ids(&self, ids: &[i64]) -> Result<u64> {
266        const SQLITE_MAX_VARS: usize = 999;
267
268        let mut total = 0u64;
269        let conn = self.conn.lock().unwrap();
270        for chunk in ids.chunks(SQLITE_MAX_VARS) {
271            if chunk.is_empty() {
272                continue;
273            }
274            let placeholders: Vec<String> = (1..=chunk.len()).map(|i| format!("?{i}")).collect();
275            let sql = format!(
276                "DELETE FROM log_entries WHERE id IN ({})",
277                placeholders.join(", ")
278            );
279            total += conn
280                .execute(&sql, rusqlite::params_from_iter(chunk.iter()))
281                .into_diagnostic()? as u64;
282        }
283        Ok(total)
284    }
285
286    /// Rotate (delete) old log entries for a specific daemon based on retention policy.
287    ///
288    /// When an archive hook is configured, entries are fetched in batches of
289    /// `batch_size`, the hook is invoked *without* holding the SQLite mutex
290    /// (so concurrent log appends are not blocked), and then deleted by ID.
291    /// Row-read errors are propagated rather than silently dropped.
292    pub fn rotate_by_age(
293        &self,
294        daemon_id: &DaemonId,
295        max_age: chrono::Duration,
296        archive_hook: Option<&ArchiveHook>,
297    ) -> Result<u64> {
298        let cutoff = (Local::now() - max_age).timestamp_millis();
299        let hook = archive_hook.filter(|h| h.is_enabled());
300
301        if let Some(hook) = hook {
302            let mut total_deleted = 0u64;
303            loop {
304                // Fetch one batch under lock, then release.
305                let entries: Vec<LogEntry> = {
306                    let conn = self.conn.lock().unwrap();
307                    let mut stmt = conn
308                        .prepare(
309                            "SELECT id, daemon_id, timestamp, message, level, msg, logger, fields_json FROM log_entries
310                             WHERE daemon_id = ?1 AND timestamp < ?2
311                             ORDER BY timestamp ASC, id ASC
312                             LIMIT ?3",
313                        )
314                        .into_diagnostic()?;
315                    stmt.query_map(
316                        params![daemon_id.qualified(), cutoff, hook.batch_size as i64],
317                        Self::row_to_entry,
318                    )
319                    .into_diagnostic()?
320                    .collect::<rusqlite::Result<Vec<_>>>()
321                    .into_diagnostic()?
322                };
323
324                if entries.is_empty() {
325                    break;
326                }
327
328                let batch_ids: Vec<i64> = entries.iter().map(|e| e.id).collect();
329
330                // Run the hook without holding the mutex.
331                self.archive_entries(&entries, hook, daemon_id, "age")?;
332
333                // Re-acquire lock and delete exactly the archived IDs.
334                let deleted = self.delete_by_ids(&batch_ids)?;
335                total_deleted += deleted;
336            }
337            Ok(total_deleted)
338        } else {
339            let conn = self.conn.lock().unwrap();
340            let rows = conn
341                .execute(
342                    "DELETE FROM log_entries WHERE daemon_id = ?1 AND timestamp < ?2",
343                    params![daemon_id.qualified(), cutoff],
344                )
345                .into_diagnostic()?;
346            Ok(rows as u64)
347        }
348    }
349
350    /// Rotate (delete) old log entries keeping only the most recent `max_count` rows
351    /// for a specific daemon.
352    ///
353    /// When an archive hook is configured, entries are fetched in batches of
354    /// `batch_size`, the hook is invoked *without* holding the SQLite mutex
355    /// (so concurrent log appends are not blocked), and then deleted by ID.
356    /// Row-read errors are propagated rather than silently dropped.
357    pub fn rotate_by_count(
358        &self,
359        daemon_id: &DaemonId,
360        max_count: u64,
361        archive_hook: Option<&ArchiveHook>,
362    ) -> Result<u64> {
363        let hook = archive_hook.filter(|h| h.is_enabled());
364
365        // Determine how many entries to delete.
366        let to_delete: i64 = {
367            let conn = self.conn.lock().unwrap();
368            let count: i64 = conn
369                .query_row(
370                    "SELECT COUNT(*) FROM log_entries WHERE daemon_id = ?1",
371                    [daemon_id.qualified()],
372                    |row| row.get(0),
373                )
374                .into_diagnostic()?;
375            count.saturating_sub(max_count as i64)
376        };
377
378        if to_delete <= 0 {
379            return Ok(0);
380        }
381
382        if let Some(hook) = hook {
383            let mut total_deleted = 0u64;
384            let mut remaining = to_delete;
385            loop {
386                let batch_len = remaining.min(hook.batch_size as i64);
387
388                // Fetch one batch under lock, then release.
389                let entries: Vec<LogEntry> = {
390                    let conn = self.conn.lock().unwrap();
391                    let mut stmt = conn
392                        .prepare(
393                            "SELECT id, daemon_id, timestamp, message, level, msg, logger, fields_json FROM log_entries
394                             WHERE daemon_id = ?1
395                             ORDER BY timestamp ASC, id ASC
396                             LIMIT ?2",
397                        )
398                        .into_diagnostic()?;
399                    stmt.query_map(
400                        params![daemon_id.qualified(), batch_len],
401                        Self::row_to_entry,
402                    )
403                    .into_diagnostic()?
404                    .collect::<rusqlite::Result<Vec<_>>>()
405                    .into_diagnostic()?
406                };
407
408                if entries.is_empty() {
409                    break;
410                }
411
412                let fetched = entries.len() as i64;
413                let batch_ids: Vec<i64> = entries.iter().map(|e| e.id).collect();
414
415                // Run the hook without holding the mutex.
416                self.archive_entries(&entries, hook, daemon_id, "count")?;
417
418                // Re-acquire lock and delete exactly the archived IDs.
419                let deleted = self.delete_by_ids(&batch_ids)?;
420                total_deleted += deleted;
421                remaining -= fetched;
422            }
423            Ok(total_deleted)
424        } else {
425            let conn = self.conn.lock().unwrap();
426            let rows = conn
427                .execute(
428                    "DELETE FROM log_entries WHERE id IN (
429                        SELECT id FROM log_entries WHERE daemon_id = ?1
430                        ORDER BY timestamp ASC, id ASC LIMIT ?2
431                    )",
432                    params![daemon_id.qualified(), to_delete],
433                )
434                .into_diagnostic()?;
435            Ok(rows as u64)
436        }
437    }
438
439    /// Migrate existing text logs for a daemon into SQLite.
440    ///
441    /// Reads the legacy text file line-by-line (streaming) and inserts in
442    /// batches of 1000 to avoid loading multi-GB files into memory at once.
443    pub fn migrate_daemon_text_logs(&self, daemon_id: &DaemonId) -> Result<u64> {
444        let text_path = daemon_id.log_path();
445        if !text_path.exists() {
446            return Ok(0);
447        }
448
449        let file = std::fs::File::open(&text_path).into_diagnostic()?;
450        let reader = BufReader::new(file);
451        let re = regex::Regex::new(r"^(\d{4}-\d{2}-\d{2} \d{2}:\d{2}:\d{2}) ([\w./-]+) (.*)$")
452            .expect("invalid regex");
453
454        let mut current_timestamp: Option<DateTime<Local>> = None;
455        let mut current_message = String::new();
456        let mut entries = Vec::with_capacity(1000);
457        let mut total_migrated: u64 = 0;
458
459        for line in reader.lines() {
460            let line = line.into_diagnostic()?;
461            if let Some(caps) = re.captures(&line) {
462                if let Some(ts) = current_timestamp.take() {
463                    entries.push((ts, std::mem::take(&mut current_message)));
464                }
465                let ts_str = caps.get(1).map(|m| m.as_str()).unwrap_or_default();
466                let msg = caps.get(3).map(|m| m.as_str()).unwrap_or_default();
467                if let Ok(naive) =
468                    chrono::NaiveDateTime::parse_from_str(ts_str, "%Y-%m-%d %H:%M:%S")
469                {
470                    current_timestamp = Local.from_local_datetime(&naive).single();
471                    current_message = msg.to_string();
472                }
473            } else if current_timestamp.is_some() {
474                current_message.push('\n');
475                current_message.push_str(&line);
476            }
477
478            if entries.len() >= 1000 {
479                total_migrated += self.insert_batch(daemon_id, &entries)?;
480                entries.clear();
481            }
482        }
483
484        if let Some(ts) = current_timestamp {
485            entries.push((ts, std::mem::take(&mut current_message)));
486        }
487
488        if !entries.is_empty() {
489            total_migrated += self.insert_batch(daemon_id, &entries)?;
490        }
491
492        if total_migrated > 0 {
493            if let Err(e) = std::fs::remove_file(&text_path) {
494                log::warn!(
495                    "failed to remove legacy log file after migration {}: {e}",
496                    text_path.display()
497                );
498            }
499        }
500
501        Ok(total_migrated)
502    }
503
504    fn insert_batch(
505        &self,
506        daemon_id: &DaemonId,
507        entries: &[(DateTime<Local>, String)],
508    ) -> Result<u64> {
509        let mut conn = self.conn.lock().unwrap();
510        let tx = conn.transaction().into_diagnostic()?;
511        let mut count = 0u64;
512        {
513            let mut stmt = tx
514                .prepare(
515                    "INSERT INTO log_entries (daemon_id, timestamp, message) VALUES (?1, ?2, ?3)",
516                )
517                .into_diagnostic()?;
518            for (ts, msg) in entries {
519                stmt.execute(params![daemon_id.qualified(), ts.timestamp_millis(), msg])
520                    .into_diagnostic()?;
521                count += 1;
522            }
523        }
524        tx.commit().into_diagnostic()?;
525        Ok(count)
526    }
527
528    /// Build the SQL query string and parameters for the given options.
529    ///
530    /// `id_range` is used by `query_parallel` to shard the query by id range.
531    /// When `Some((start, end))`, adds `id > start AND id <= end` to the WHERE clause.
532    fn build_query_sql(
533        opts: &LogQuery,
534        id_range: Option<(i64, i64)>,
535    ) -> (String, Vec<Box<dyn rusqlite::ToSql>>) {
536        let mut conditions = Vec::new();
537        let mut query_params: Vec<Box<dyn rusqlite::ToSql>> = Vec::new();
538
539        if !opts.daemon_ids.is_empty() {
540            let placeholders: Vec<String> = (1..=opts.daemon_ids.len())
541                .map(|i| format!("?{}", i))
542                .collect();
543            conditions.push(format!("daemon_id IN ({})", placeholders.join(", ")));
544            for id in &opts.daemon_ids {
545                query_params.push(Box::new(id.clone()));
546            }
547        }
548
549        if let Some(from) = opts.from {
550            conditions.push(format!("timestamp >= ?{}", query_params.len() + 1));
551            query_params.push(Box::new(from.timestamp_millis()));
552        }
553
554        if let Some(to) = opts.to {
555            conditions.push(format!("timestamp <= ?{}", query_params.len() + 1));
556            query_params.push(Box::new(to.timestamp_millis()));
557        }
558
559        if let Some(after_id) = opts.after_id {
560            conditions.push(format!("id > ?{}", query_params.len() + 1));
561            query_params.push(Box::new(after_id));
562        }
563
564        if let Some((start, end)) = id_range {
565            conditions.push(format!("id > ?{}", query_params.len() + 1));
566            query_params.push(Box::new(start));
567            conditions.push(format!("id <= ?{}", query_params.len() + 1));
568            query_params.push(Box::new(end));
569        }
570
571        let mut message_conditions = Vec::new();
572        for filter in &opts.message_filters {
573            match filter {
574                MessageFilter::Contains {
575                    pattern,
576                    case_sensitive,
577                } => {
578                    let param_index = query_params.len() + 1;
579                    if *case_sensitive {
580                        message_conditions
581                            .push(format!("INSTR(message, ?{idx}) > 0", idx = param_index));
582                        query_params.push(Box::new(pattern.clone()));
583                    } else {
584                        let escaped = escape_like_pattern(pattern);
585                        let param = format!("%{}%", escaped);
586                        message_conditions.push(format!(
587                            "LOWER(message) LIKE LOWER(?{idx}) ESCAPE '\\'",
588                            idx = param_index
589                        ));
590                        query_params.push(Box::new(param));
591                    }
592                }
593                MessageFilter::Regex { pattern } => {
594                    let param_index = query_params.len() + 1;
595                    message_conditions.push(format!("message REGEXP ?{param_index}"));
596                    query_params.push(Box::new(pattern.clone()));
597                }
598            }
599        }
600        if !message_conditions.is_empty() {
601            conditions.push(format!("({})", message_conditions.join(" OR ")));
602        }
603
604        for filter in &opts.field_filters {
605            match filter {
606                FieldFilter::LevelMin(level) => {
607                    let matching = crate::log_store::levels_at_or_above(level);
608                    if matching.is_empty() {
609                        // Unknown level: no results
610                        conditions.push("0".to_string());
611                    } else {
612                        let placeholders = matching
613                            .iter()
614                            .map(|l| {
615                                let idx = query_params.len() + 1;
616                                query_params.push(Box::new((*l).to_string()));
617                                format!("?{idx}")
618                            })
619                            .collect::<Vec<_>>()
620                            .join(", ");
621                        conditions.push(format!("level IN ({placeholders})"));
622                    }
623                }
624                FieldFilter::FieldEq { key, value } => {
625                    let param_index = query_params.len() + 1;
626                    conditions.push(format!(
627                        "json_extract(fields_json, '$.{key}') = ?{param_index}"
628                    ));
629                    query_params.push(Box::new(value.clone()));
630                }
631            }
632        }
633
634        let where_clause = if conditions.is_empty() {
635            String::new()
636        } else {
637            format!("WHERE {}", conditions.join(" AND "))
638        };
639
640        let order = if opts.order_desc { "DESC" } else { "ASC" };
641
642        let limit_clause = opts
643            .limit
644            .map(|n| format!("LIMIT {}", n))
645            .unwrap_or_default();
646
647        let columns = if opts.include_structured {
648            "id, daemon_id, timestamp, message, level, msg, logger, fields_json"
649        } else {
650            "id, daemon_id, timestamp, message, NULL, NULL, NULL, NULL"
651        };
652
653        let sql = format!(
654            "SELECT {columns} FROM log_entries {} ORDER BY timestamp {}, id {} {}",
655            where_clause, order, order, limit_clause
656        );
657
658        (sql, query_params)
659    }
660
661    /// Returns true if parallel query is beneficial for the given options.
662    fn should_parallelize(opts: &LogQuery) -> bool {
663        if opts.daemon_ids.len() != 1 {
664            return false;
665        }
666        // Skip parallel for incremental tail polls (after_id set) — they
667        // return small batches and the connection overhead dominates.
668        if opts.after_id.is_some() {
669            return false;
670        }
671        let limit = opts.limit.unwrap_or(usize::MAX);
672        if limit < PARALLEL_QUERY_THRESHOLD {
673            return false;
674        }
675        std::thread::available_parallelism()
676            .map(|n| n.get() >= 2)
677            .unwrap_or(false)
678    }
679
680    /// Query using multiple read-only connections, sharded by id range.
681    ///
682    /// For a single daemon, id order ≈ timestamp order (timestamps are
683    /// assigned in insertion order within a single writer). This lets us
684    /// shard by contiguous id ranges and merge by concatenating shards in
685    /// the right order, without a global sort.
686    fn query_parallel(&self, opts: &LogQuery) -> Result<Vec<LogEntry>> {
687        // Cap parallelism: each shard opens a new read-only connection with
688        // PRAGMA setup + regexp function registration (~5-12ms each). Beyond
689        // 2 threads the connection overhead dominates the query savings for
690        // typical log store sizes.
691        let max_threads = 2;
692        let num_threads = std::thread::available_parallelism()
693            .map(|n| n.get().min(max_threads))
694            .unwrap_or(1);
695
696        let max_id: Option<i64> = {
697            let conn = self.conn.lock().unwrap();
698            conn.query_row(
699                "SELECT MAX(id) FROM log_entries WHERE daemon_id = ?1",
700                params![&opts.daemon_ids[0]],
701                |row| row.get(0),
702            )
703            .ok()
704        };
705
706        let Some(max_id) = max_id else {
707            return Ok(Vec::new());
708        };
709        if max_id == 0 {
710            return Ok(Vec::new());
711        }
712
713        let shard_size = (max_id as usize).div_ceil(num_threads);
714        let path = self.path.clone();
715        let opts = opts.clone();
716        let needs_regexp = opts
717            .message_filters
718            .iter()
719            .any(|f| matches!(f, MessageFilter::Regex { .. }));
720
721        let shards: Vec<Result<Vec<LogEntry>>> = std::thread::scope(|s| {
722            (0..num_threads)
723                .map(|i| {
724                    let start = (i * shard_size) as i64;
725                    let end = if i == num_threads - 1 {
726                        max_id
727                    } else {
728                        ((i + 1) * shard_size) as i64
729                    };
730                    let opts = opts.clone();
731                    let path = path.clone();
732                    s.spawn(move || -> Result<Vec<LogEntry>> {
733                        let conn = Connection::open_with_flags(
734                            &path,
735                            OpenFlags::SQLITE_OPEN_READ_ONLY | OpenFlags::SQLITE_OPEN_NO_MUTEX,
736                        )
737                        .into_diagnostic()?;
738                        conn.execute_batch(
739                            "PRAGMA mmap_size = 268435456;
740                             PRAGMA query_only = 1;",
741                        )
742                        .into_diagnostic()?;
743                        if needs_regexp {
744                            add_regexp_function(&conn)?;
745                        }
746
747                        let (sql, query_params) = Self::build_query_sql(&opts, Some((start, end)));
748                        Self::execute_built_query(&conn, &sql, &query_params)
749                    })
750                })
751                .map(|h| h.join().unwrap())
752                .collect()
753        });
754
755        let mut merged = Vec::new();
756        if opts.order_desc {
757            for shard in shards.into_iter().rev() {
758                merged.extend(shard?);
759            }
760        } else {
761            for shard in shards {
762                merged.extend(shard?);
763            }
764        }
765
766        if let Some(limit) = opts.limit {
767            if merged.len() > limit {
768                merged.truncate(limit);
769            }
770        }
771
772        Ok(merged)
773    }
774
775    /// Execute a built SQL query and collect results into LogEntry.
776    fn execute_built_query(
777        conn: &Connection,
778        sql: &str,
779        query_params: &[Box<dyn rusqlite::ToSql>],
780    ) -> Result<Vec<LogEntry>> {
781        let mut stmt = conn.prepare(sql).into_diagnostic()?;
782        let params_ref: Vec<&dyn rusqlite::ToSql> =
783            query_params.iter().map(|p| p.as_ref()).collect();
784        let rows = stmt
785            .query_map(params_ref.as_slice(), Self::row_to_entry)
786            .into_diagnostic()?;
787        let mut entries = Vec::new();
788        for row in rows {
789            entries.push(row.into_diagnostic()?);
790        }
791        Ok(entries)
792    }
793}
794
795impl LogStore for SqliteLogStore {
796    fn append(&self, daemon_id: &DaemonId, message: &str) -> Result<()> {
797        let ts = Local::now().timestamp_millis();
798        let id = daemon_id.qualified();
799        let msg = message.to_string();
800
801        let conn = self.conn.lock().unwrap();
802        let _ = conn
803            .execute(
804                "INSERT INTO log_entries (daemon_id, timestamp, message) VALUES (?1, ?2, ?3)",
805                params![id, ts, msg],
806            )
807            .into_diagnostic()?;
808        Ok(())
809    }
810
811    fn append_batch(&self, daemon_id: &DaemonId, messages: &[String]) -> Result<()> {
812        if messages.is_empty() {
813            return Ok(());
814        }
815        let base_ts = Local::now().timestamp_millis();
816        let id = daemon_id.qualified();
817
818        let mut conn = self.conn.lock().unwrap();
819        let tx = conn.transaction().into_diagnostic()?;
820        {
821            let mut stmt = tx
822                .prepare(
823                    "INSERT INTO log_entries (daemon_id, timestamp, message) VALUES (?1, ?2, ?3)",
824                )
825                .into_diagnostic()?;
826            for (idx, msg) in messages.iter().enumerate() {
827                // Slightly stagger timestamps within a batch so ordering by
828                // (timestamp, id) preserves insertion order without paying for
829                // a separate per-row clock read.
830                let ts = base_ts + idx as i64;
831                stmt.execute(params![id, ts, msg]).into_diagnostic()?;
832            }
833        }
834        tx.commit().into_diagnostic()?;
835        Ok(())
836    }
837
838    fn append_structured(&self, daemon_id: &DaemonId, parsed: &ParsedLog) -> Result<()> {
839        let ts = Local::now().timestamp_millis();
840        let id = daemon_id.qualified();
841
842        let conn = self.conn.lock().unwrap();
843        let _ = conn
844            .execute(
845                "INSERT INTO log_entries (daemon_id, timestamp, message, level, msg, logger, fields_json) \
846                 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)",
847                params![id, ts, parsed.message, parsed.level, parsed.msg, parsed.logger, parsed.fields_json],
848            )
849            .into_diagnostic()?;
850        Ok(())
851    }
852
853    fn append_structured_batch(&self, daemon_id: &DaemonId, entries: &[ParsedLog]) -> Result<()> {
854        if entries.is_empty() {
855            return Ok(());
856        }
857        let base_ts = Local::now().timestamp_millis();
858        let id = daemon_id.qualified();
859
860        let mut conn = self.conn.lock().unwrap();
861        let tx = conn.transaction().into_diagnostic()?;
862        {
863            let mut stmt = tx
864                .prepare(
865                    "INSERT INTO log_entries (daemon_id, timestamp, message, level, msg, logger, fields_json) \
866                     VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)",
867                )
868                .into_diagnostic()?;
869            for (idx, entry) in entries.iter().enumerate() {
870                let ts = base_ts + idx as i64;
871                stmt.execute(params![
872                    id,
873                    ts,
874                    entry.message,
875                    entry.level,
876                    entry.msg,
877                    entry.logger,
878                    entry.fields_json
879                ])
880                .into_diagnostic()?;
881            }
882        }
883        tx.commit().into_diagnostic()?;
884        Ok(())
885    }
886
887    fn query(&self, opts: &LogQuery) -> Result<Vec<LogEntry>> {
888        // Delegate to parallel path for large single-daemon queries.
889        if Self::should_parallelize(opts) && self.path.as_os_str() != ":memory:" {
890            if let Ok(entries) = self.query_parallel(opts) {
891                return Ok(entries);
892            }
893            // Fall back to single-threaded on parallel failure.
894        }
895
896        // Single-threaded path.
897        let conn = self.conn.lock().unwrap();
898        let (sql, query_params) = Self::build_query_sql(opts, None);
899        Self::execute_built_query(&conn, &sql, &query_params)
900    }
901
902    fn tail(&self, daemon_id: &DaemonId, after_id: Option<i64>) -> Result<Vec<LogEntry>> {
903        self.query(&LogQuery {
904            daemon_ids: vec![daemon_id.qualified()],
905            from: None,
906            to: None,
907            limit: None,
908            order_desc: false,
909            after_id,
910            message_filters: Vec::new(),
911            field_filters: Vec::new(),
912            include_structured: false,
913        })
914    }
915
916    fn clear(&self, daemon_ids: &[DaemonId]) -> Result<()> {
917        let mut conn = self.conn.lock().unwrap();
918        let tx = conn.transaction().into_diagnostic()?;
919        for id in daemon_ids {
920            tx.execute(
921                "DELETE FROM log_entries WHERE daemon_id = ?1",
922                params![id.qualified()],
923            )
924            .into_diagnostic()?;
925            tx.execute(
926                "INSERT INTO log_clear_generations (daemon_id, generation)
927                 VALUES (?1, 1)
928                 ON CONFLICT(daemon_id) DO UPDATE SET generation = generation + 1",
929                params![id.qualified()],
930            )
931            .into_diagnostic()?;
932        }
933        tx.commit().into_diagnostic()?;
934        Ok(())
935    }
936
937    fn last_id(&self, daemon_id: &DaemonId) -> Result<Option<i64>> {
938        let conn = self.conn.lock().unwrap();
939        // MAX(id) returns NULL when no rows exist for the daemon; the
940        // Option<i64> decode maps NULL to None automatically.
941        let id: Option<i64> = conn
942            .query_row(
943                "SELECT MAX(id) FROM log_entries WHERE daemon_id = ?1",
944                params![daemon_id.qualified()],
945                |row| row.get(0),
946            )
947            .into_diagnostic()?;
948        Ok(id)
949    }
950
951    fn list_daemon_ids(&self) -> Result<Vec<String>> {
952        let conn = self.conn.lock().unwrap();
953        let mut stmt = conn
954            .prepare("SELECT DISTINCT daemon_id FROM log_entries")
955            .into_diagnostic()?;
956        let ids = stmt
957            .query_map([], |row| {
958                let id: String = row.get(0)?;
959                Ok(id)
960            })
961            .into_diagnostic()?
962            .filter_map(|r| r.ok())
963            .collect();
964        Ok(ids)
965    }
966
967    fn apply_retention(
968        &self,
969        policy: &super::RetentionPolicy,
970        excluded_daemon_ids: &[DaemonId],
971        archive_hook: Option<&ArchiveHook>,
972    ) -> Result<u64> {
973        let daemon_ids = self.list_daemon_ids()?;
974        let excluded: HashSet<String> = excluded_daemon_ids.iter().map(|d| d.qualified()).collect();
975        let mut total = 0u64;
976        for id_str in daemon_ids {
977            if excluded.contains(&id_str) {
978                continue;
979            }
980            let id = DaemonId::parse(&id_str).unwrap_or_else(|_| {
981                DaemonId::try_new("global", &id_str).unwrap_or_else(|_| DaemonId::pitchfork())
982            });
983            if let Some(dur) = policy.age {
984                total += self.rotate_by_age(&id, dur, archive_hook)?;
985            }
986            if let Some(n) = policy.count {
987                total += self.rotate_by_count(&id, n, archive_hook)?;
988            }
989        }
990        Ok(total)
991    }
992
993    fn apply_retention_for_daemon(
994        &self,
995        daemon_id: &DaemonId,
996        policy: &super::RetentionPolicy,
997        archive_hook: Option<&ArchiveHook>,
998    ) -> Result<u64> {
999        let mut total = 0u64;
1000        if let Some(dur) = policy.age {
1001            total += self.rotate_by_age(daemon_id, dur, archive_hook)?;
1002        }
1003        if let Some(n) = policy.count {
1004            total += self.rotate_by_count(daemon_id, n, archive_hook)?;
1005        }
1006        Ok(total)
1007    }
1008
1009    fn last_clear_generation(&self, daemon_id: &DaemonId) -> Result<Option<u64>> {
1010        let conn = self.conn.lock().unwrap();
1011        let mut stmt = conn
1012            .prepare("SELECT generation FROM log_clear_generations WHERE daemon_id = ?1")
1013            .into_diagnostic()?;
1014        let generation: Option<i64> = stmt
1015            .query_row(params![daemon_id.qualified()], |row| row.get(0))
1016            .optional()
1017            .into_diagnostic()?;
1018        generation
1019            .map(|generation| {
1020                u64::try_from(generation)
1021                    .map_err(|_| miette::miette!("log clear generation cannot be negative"))
1022            })
1023            .transpose()
1024    }
1025}
1026
1027/// Global singleton log store.
1028use once_cell::sync::Lazy;
1029use std::sync::Arc;
1030
1031pub static LOG_STORE: Lazy<Arc<SqliteLogStore>> = Lazy::new(|| {
1032    let path = crate::env::PITCHFORK_LOGS_DIR.join("logs.db");
1033    let mut is_fallback = false;
1034    let store = Arc::new(SqliteLogStore::open(&path).unwrap_or_else(|e| {
1035        error!(
1036            "failed to open log store at {}: {e}. Falling back to in-memory store; logs will not persist across restarts.",
1037            path.display()
1038        );
1039        is_fallback = true;
1040        SqliteLogStore::open(":memory:").expect("in-memory SQLite should always open")
1041    }));
1042
1043    // Auto-migrate any legacy text log files into SQLite on first access.
1044    // This runs once per process startup and is idempotent.
1045    // Skip migration when using the in-memory fallback to prevent data loss.
1046    if !is_fallback {
1047        if let Err(e) = auto_migrate_legacy_logs(&store) {
1048            warn!("legacy log auto-migration failed: {e}");
1049        }
1050    } else {
1051        warn!(
1052            "skipping legacy log auto-migration because log store is in-memory (no durable destination)"
1053        );
1054    }
1055
1056    store
1057});
1058
1059/// Auto-migrate legacy text log files into the SQLite log store.
1060///
1061/// Scans the logs directory for directories matching the new-format layout
1062/// (`namespace--name/namespace--name.log`), attempts to parse the directory
1063/// name as a valid safe-path daemon ID, and imports the content into SQLite.
1064/// This is idempotent: re-running it on already-migrated data is a no-op
1065/// because the legacy text files are deleted after successful import.
1066fn auto_migrate_legacy_logs(store: &SqliteLogStore) -> Result<()> {
1067    let logs_dir = &*crate::env::PITCHFORK_LOGS_DIR;
1068    if !logs_dir.exists() {
1069        return Ok(());
1070    }
1071
1072    let Ok(entries) = std::fs::read_dir(logs_dir) else {
1073        return Ok(());
1074    };
1075
1076    let mut total_migrated = 0u64;
1077    let mut migrated_ids = Vec::new();
1078
1079    for entry in entries.flatten() {
1080        let path = entry.path();
1081        if !path.is_dir() {
1082            continue;
1083        }
1084        // Skip the supervisor's own log directory.
1085        let file_name = path
1086            .file_name()
1087            .map_or(String::new(), |n| n.to_string_lossy().to_string());
1088        if file_name == "pitchfork" {
1089            continue;
1090        }
1091
1092        // Only consider directories that look like new-format safe-paths
1093        if !file_name.contains("--") {
1094            continue;
1095        }
1096        let log_file = path.join(format!("{file_name}.log"));
1097        if !log_file.exists() {
1098            continue;
1099        }
1100
1101        let daemon_id = match DaemonId::from_safe_path(&file_name) {
1102            Ok(id) => id,
1103            Err(_) => continue,
1104        };
1105
1106        // Skip the supervisor's own daemon; its log directory may be present
1107        // under logs/ but should never be imported into the user-facing log store.
1108        if daemon_id == DaemonId::pitchfork() {
1109            continue;
1110        }
1111
1112        match store.migrate_daemon_text_logs(&daemon_id) {
1113            Ok(0) => {}
1114            Ok(n) => {
1115                total_migrated += n;
1116                migrated_ids.push(daemon_id.qualified());
1117            }
1118            Err(e) => {
1119                warn!(
1120                    "failed to migrate text logs for {}: {e}",
1121                    daemon_id.qualified()
1122                );
1123            }
1124        }
1125    }
1126
1127    if total_migrated > 0 {
1128        warn!(
1129            "auto-migrated {total_migrated} legacy log entries from {count} daemon(s): {ids}",
1130            count = migrated_ids.len(),
1131            ids = migrated_ids.join(", ")
1132        );
1133    }
1134
1135    Ok(())
1136}