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