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/// Convert a text filter value to a JSON literal string for use with
17/// SQLite's `json()` function. This ensures `json_each.value` (which returns
18/// native SQLite types) is compared against the correct type:
19///
20/// - `"true"` / `"false"` → JSON boolean (SQLite integer 1/0)
21/// - `"null"` → JSON null (SQLite NULL)
22/// - `"42"` / `"3.14"` → JSON number (SQLite integer/real)
23/// - `"hello"` → JSON string `"hello"`
24///
25/// Without this conversion, binding `8080` as text would never match
26/// `json_each.value` returning integer `8080`.
27fn text_to_json_literal(value: &str) -> String {
28    // Boolean
29    if value.eq_ignore_ascii_case("true") {
30        return "true".to_string();
31    }
32    if value.eq_ignore_ascii_case("false") {
33        return "false".to_string();
34    }
35    // Null
36    if value.eq_ignore_ascii_case("null") {
37        return "null".to_string();
38    }
39    // Number: validate as a JSON number (stricter than f64::parse, which
40    // accepts "+42", "1.", ".5", "inf", "nan" — none are valid JSON and
41    // would cause SQLite's json_extract to fail).
42    if let Ok(serde_json::Value::Number(_)) = serde_json::from_str(value) {
43        return value.to_string();
44    }
45    // String: JSON-escape and quote
46    serde_json::to_string(value).unwrap_or_else(|_| format!(r#""{value}""#))
47}
48
49/// Registers a `regexp` SQL function backed by the `regex` crate.
50///
51/// SQLite does not ship a REGEXP implementation by default. This function
52/// is registered on every new connection so that `message REGEXP ?` works
53/// consistently across queries.
54fn add_regexp_function(conn: &Connection) -> Result<()> {
55    use std::cell::RefCell;
56
57    // Cache compiled regexes per connection to avoid recompiling the same
58    // pattern on every row evaluation. The cache is small and thread-local
59    // because scalar functions are invoked on the connection's thread.
60    let cache: RefCell<lru::LruCache<String, regex::Regex>> =
61        RefCell::new(lru::LruCache::new(std::num::NonZeroUsize::new(32).unwrap()));
62
63    conn.create_scalar_function(
64        "regexp",
65        2,
66        rusqlite::functions::FunctionFlags::SQLITE_UTF8
67            | rusqlite::functions::FunctionFlags::SQLITE_DETERMINISTIC,
68        move |ctx| {
69            let pattern: String = ctx.get(0)?;
70            let text: String = ctx.get(1)?;
71
72            let mut cache = cache.borrow_mut();
73            let re = match cache.get(&pattern) {
74                Some(re) => re.clone(),
75                None => {
76                    let re = regex::Regex::new(&pattern)
77                        .map_err(|e| rusqlite::Error::UserFunctionError(e.to_string().into()))?;
78                    cache.put(pattern.clone(), re.clone());
79                    re
80                }
81            };
82            Ok(re.is_match(&text))
83        },
84    )
85    .into_diagnostic()
86}
87
88/// SQLite-backed log store with WAL mode for concurrent readers.
89pub struct SqliteLogStore {
90    conn: Mutex<Connection>,
91    path: PathBuf,
92}
93
94/// Minimum result-set size to trigger parallel query.
95///
96/// Each parallel shard opens a new read-only SQLite connection (~5-12ms
97/// overhead per connection for PRAGMA setup + regexp function registration).
98/// Below this threshold, single-threaded is faster because the connection
99/// overhead exceeds the query savings.
100const PARALLEL_QUERY_THRESHOLD: usize = 200_000;
101
102/// How long a statement waits for a contended lock before giving up. Generous
103/// enough to outlast any batch insert or retention prune, short enough that a
104/// genuinely stuck writer still surfaces as an error rather than hanging.
105const BUSY_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(5);
106
107/// Put the database into WAL mode, tolerating a lost race.
108///
109/// Switching journal modes needs an exclusive lock and, unlike ordinary
110/// statements, does not consult the busy handler — so concurrent opens collide
111/// here no matter how large `busy_timeout` is. WAL is recorded in the database
112/// header and persists across connections, so only the first open has to
113/// succeed: retry briefly in case a peer is mid-switch, and if it still has not
114/// taken, carry on rather than fail. A connection that never personally set WAL
115/// still works correctly; refusing to open the log store over a transient
116/// pragma race would be far worse than running one connection unoptimised.
117fn enable_wal(conn: &Connection) {
118    const ATTEMPTS: usize = 5;
119    for attempt in 0..ATTEMPTS {
120        match conn.query_row("PRAGMA journal_mode", [], |row| row.get::<_, String>(0)) {
121            Ok(mode) if mode.eq_ignore_ascii_case("wal") => return,
122            Ok(_) => {}
123            Err(e) => {
124                debug!("could not read journal_mode: {e}");
125                return;
126            }
127        }
128        if conn.execute_batch("PRAGMA journal_mode = WAL;").is_ok() {
129            return;
130        }
131        if attempt + 1 < ATTEMPTS {
132            std::thread::sleep(std::time::Duration::from_millis(20));
133        }
134    }
135    debug!("log store is not in WAL mode; another connection may be switching it");
136}
137
138impl SqliteLogStore {
139    /// Open or create the SQLite log store at the given path.
140    pub fn open(path: impl Into<PathBuf>) -> Result<Self> {
141        let path = path.into();
142        if let Some(parent) = path.parent() {
143            std::fs::create_dir_all(parent).into_diagnostic()?;
144        }
145        let conn = Connection::open(&path).into_diagnostic()?;
146
147        // A contended write must wait for the lock rather than fail: WAL
148        // allows concurrent readers but still serializes writers, and a writer
149        // with no timeout takes SQLITE_BUSY immediately and drops its work.
150        // rusqlite already defaults to this value, but state it explicitly so
151        // the requirement is visible here and does not rest on a dependency's
152        // default. Set before anything else, since the statements below take
153        // locks of their own.
154        conn.busy_timeout(BUSY_TIMEOUT).into_diagnostic()?;
155
156        add_regexp_function(&conn)?;
157
158        enable_wal(&conn);
159
160        conn.execute_batch(
161            "PRAGMA synchronous = NORMAL;
162             PRAGMA mmap_size = 268435456;",
163        )
164        .into_diagnostic()?;
165        conn.execute(
166            "CREATE TABLE IF NOT EXISTS log_entries (
167                id          INTEGER PRIMARY KEY AUTOINCREMENT,
168                daemon_id   TEXT    NOT NULL,
169                timestamp   INTEGER NOT NULL,
170                message     TEXT    NOT NULL,
171                level       TEXT,
172                msg         TEXT,
173                logger      TEXT,
174                fields_json TEXT
175            );",
176            [],
177        )
178        .into_diagnostic()?;
179        conn.execute(
180            "CREATE INDEX IF NOT EXISTS idx_daemon_ts ON log_entries(daemon_id, timestamp);",
181            [],
182        )
183        .into_diagnostic()?;
184        conn.execute(
185            "CREATE INDEX IF NOT EXISTS idx_daemon_id ON log_entries(daemon_id, id);",
186            [],
187        )
188        .into_diagnostic()?;
189        conn.execute(
190            "CREATE INDEX IF NOT EXISTS idx_timestamp ON log_entries(timestamp);",
191            [],
192        )
193        .into_diagnostic()?;
194
195        // Migrate existing tables: add columns introduced in this version.
196        // Must run BEFORE creating indexes that reference the new columns.
197        let existing_cols: Vec<String> = {
198            let mut stmt = conn
199                .prepare("PRAGMA table_info(log_entries)")
200                .into_diagnostic()?;
201            let rows = stmt
202                .query_map([], |row| row.get::<_, String>(1))
203                .into_diagnostic()?;
204            rows.filter_map(|r| r.ok()).collect()
205        };
206        for col in ["level", "msg", "logger", "fields_json"] {
207            if !existing_cols.iter().any(|c| c == col) {
208                conn.execute(
209                    &format!("ALTER TABLE log_entries ADD COLUMN {col} TEXT"),
210                    [],
211                )
212                .into_diagnostic()?;
213            }
214        }
215        conn.execute(
216            "CREATE INDEX IF NOT EXISTS idx_daemon_level_ts ON log_entries(daemon_id, level, timestamp);",
217            [],
218        )
219        .into_diagnostic()?;
220
221        conn.execute(
222            "CREATE TABLE IF NOT EXISTS log_clear_generations (
223                daemon_id TEXT PRIMARY KEY,
224                generation INTEGER NOT NULL DEFAULT 0
225            );",
226            [],
227        )
228        .into_diagnostic()?;
229        Ok(Self {
230            conn: Mutex::new(conn),
231            path,
232        })
233    }
234
235    fn row_to_entry(row: &rusqlite::Row) -> rusqlite::Result<LogEntry> {
236        let id: i64 = row.get(0)?;
237        let daemon_id: String = row.get(1)?;
238        let ts_millis: i64 = row.get(2)?;
239        let message: String = row.get(3)?;
240        let level: Option<String> = row.get(4)?;
241        let msg: Option<String> = row.get(5)?;
242        let logger: Option<String> = row.get(6)?;
243        let fields_json: Option<String> = row.get(7)?;
244        let timestamp = Local
245            .timestamp_millis_opt(ts_millis)
246            .single()
247            .unwrap_or_else(Local::now);
248        Ok(LogEntry {
249            id,
250            daemon_id,
251            timestamp,
252            message,
253            level,
254            msg,
255            logger,
256            fields_json,
257        })
258    }
259
260    fn archive_entries(
261        &self,
262        entries: &[LogEntry],
263        archive_hook: &ArchiveHook,
264        daemon_id: &DaemonId,
265        reason: &str,
266    ) -> Result<()> {
267        use crate::shell::HideConsoleWindow;
268        use std::process::{Command, Stdio};
269
270        if entries.is_empty() {
271            return Ok(());
272        }
273
274        // The archive hook is a shell command like any other daemon script, so
275        // it runs under the same resolved shell rather than a hardcoded `sh`,
276        // which does not exist on a stock Windows machine. Resolve here rather
277        // than when the hook is built: a bad setting must skip the batch, not
278        // drop the hook, or retention would prune entries that were never
279        // archived.
280        let shell =
281            crate::settings::resolve_shell().map_err(|e| miette::miette!("archive hook: {e}"))?;
282        let (shell_program, shell_args) = shell.split_first().unwrap();
283
284        for chunk in entries.chunks(archive_hook.batch_size.max(1)) {
285            let mut child = Command::new(shell_program)
286                .args(shell_args)
287                .arg(&archive_hook.command)
288                .stdin(Stdio::piped())
289                .stdout(Stdio::null())
290                .stderr(Stdio::piped())
291                .env("PITCHFORK_DAEMON_ID", daemon_id.qualified())
292                .env("PITCHFORK_ARCHIVE_REASON", reason)
293                .hide_console_window()
294                .spawn()
295                .into_diagnostic()
296                .map_err(|e| miette::miette!("failed to spawn archive hook: {e}"))?;
297
298            // Write entries to stdin. On error, kill and reap the child
299            // before propagating — Child::drop does not call wait(), so
300            // without this a crashed hook would leave a zombie process.
301            let write_result = {
302                let stdin = child.stdin.take().expect("piped stdin should be available");
303                let mut stdin = std::io::BufWriter::new(stdin);
304                let mut result = Ok(());
305                for entry in chunk {
306                    let line = serde_json::json!({
307                        "id": entry.id,
308                        "daemon_id": entry.daemon_id,
309                        "timestamp": entry.timestamp.to_rfc3339(),
310                        "message": entry.message,
311                    });
312                    if let Err(e) = writeln!(stdin, "{}", line) {
313                        result = Err(miette::miette!(
314                            "failed to write to archive hook stdin: {e}"
315                        ));
316                        break;
317                    }
318                }
319                // Explicitly flush so a buffer-drain failure (e.g. the hook
320                // exited early and closed stdin) is surfaced as an error
321                // rather than being silently swallowed by BufWriter::drop.
322                // On a small batch where all data fits in the 8 KB buffer no
323                // individual writeln! has touched the underlying pipe yet, so
324                // without this flush the failure would be lost and the entries
325                // deleted without ever being delivered to the hook.
326                if result.is_ok()
327                    && let Err(e) = stdin.flush()
328                {
329                    result = Err(miette::miette!("failed to flush archive hook stdin: {e}"));
330                }
331                result
332                // BufWriter + ChildStdin drop here, closing stdin (EOF signal).
333            };
334
335            if let Err(e) = write_result {
336                let _ = child.kill();
337                let _ = child.wait();
338                return Err(e);
339            }
340
341            let output = child.wait_with_output().into_diagnostic()?;
342            if !output.status.success() {
343                let stderr = String::from_utf8_lossy(&output.stderr);
344                return Err(miette::miette!(
345                    "archive hook failed with status {}: {stderr}",
346                    output.status
347                ));
348            }
349        }
350
351        Ok(())
352    }
353
354    /// Delete rows matching the given IDs, returning the number of rows deleted.
355    ///
356    /// Chunks at 999 to stay within SQLite's `SQLITE_MAX_VARIABLE_NUMBER`
357    /// limit on versions ≤ 3.31 (e.g. Ubuntu 20.04 LTS).
358    fn delete_by_ids(&self, ids: &[i64]) -> Result<u64> {
359        const SQLITE_MAX_VARS: usize = 999;
360
361        let mut total = 0u64;
362        let conn = self.conn.lock().unwrap();
363        for chunk in ids.chunks(SQLITE_MAX_VARS) {
364            if chunk.is_empty() {
365                continue;
366            }
367            let placeholders: Vec<String> = (1..=chunk.len()).map(|i| format!("?{i}")).collect();
368            let sql = format!(
369                "DELETE FROM log_entries WHERE id IN ({})",
370                placeholders.join(", ")
371            );
372            total += conn
373                .execute(&sql, rusqlite::params_from_iter(chunk.iter()))
374                .into_diagnostic()? as u64;
375        }
376        Ok(total)
377    }
378
379    /// Rotate (delete) old log entries for a specific daemon based on retention policy.
380    ///
381    /// When an archive hook is configured, entries are fetched in batches of
382    /// `batch_size`, the hook is invoked *without* holding the SQLite mutex
383    /// (so concurrent log appends are not blocked), and then deleted by ID.
384    /// Row-read errors are propagated rather than silently dropped.
385    pub fn rotate_by_age(
386        &self,
387        daemon_id: &DaemonId,
388        max_age: chrono::Duration,
389        archive_hook: Option<&ArchiveHook>,
390    ) -> Result<u64> {
391        let cutoff = (Local::now() - max_age).timestamp_millis();
392        let hook = archive_hook.filter(|h| h.is_enabled());
393
394        if let Some(hook) = hook {
395            let mut total_deleted = 0u64;
396            loop {
397                // Fetch one batch under lock, then release.
398                let entries: Vec<LogEntry> = {
399                    let conn = self.conn.lock().unwrap();
400                    let mut stmt = conn
401                        .prepare(
402                            "SELECT id, daemon_id, timestamp, message, level, msg, logger, fields_json FROM log_entries
403                             WHERE daemon_id = ?1 AND timestamp < ?2
404                             ORDER BY timestamp ASC, id ASC
405                             LIMIT ?3",
406                        )
407                        .into_diagnostic()?;
408                    stmt.query_map(
409                        params![daemon_id.qualified(), cutoff, hook.batch_size as i64],
410                        Self::row_to_entry,
411                    )
412                    .into_diagnostic()?
413                    .collect::<rusqlite::Result<Vec<_>>>()
414                    .into_diagnostic()?
415                };
416
417                if entries.is_empty() {
418                    break;
419                }
420
421                let batch_ids: Vec<i64> = entries.iter().map(|e| e.id).collect();
422
423                // Run the hook without holding the mutex.
424                self.archive_entries(&entries, hook, daemon_id, "age")?;
425
426                // Re-acquire lock and delete exactly the archived IDs.
427                let deleted = self.delete_by_ids(&batch_ids)?;
428                total_deleted += deleted;
429            }
430            Ok(total_deleted)
431        } else {
432            let conn = self.conn.lock().unwrap();
433            let rows = conn
434                .execute(
435                    "DELETE FROM log_entries WHERE daemon_id = ?1 AND timestamp < ?2",
436                    params![daemon_id.qualified(), cutoff],
437                )
438                .into_diagnostic()?;
439            Ok(rows as u64)
440        }
441    }
442
443    /// Rotate (delete) old log entries keeping only the most recent `max_count` rows
444    /// for a specific daemon.
445    ///
446    /// When an archive hook is configured, entries are fetched in batches of
447    /// `batch_size`, the hook is invoked *without* holding the SQLite mutex
448    /// (so concurrent log appends are not blocked), and then deleted by ID.
449    /// Row-read errors are propagated rather than silently dropped.
450    pub fn rotate_by_count(
451        &self,
452        daemon_id: &DaemonId,
453        max_count: u64,
454        archive_hook: Option<&ArchiveHook>,
455    ) -> Result<u64> {
456        let hook = archive_hook.filter(|h| h.is_enabled());
457
458        // Determine how many entries to delete.
459        let to_delete: i64 = {
460            let conn = self.conn.lock().unwrap();
461            let count: i64 = conn
462                .query_row(
463                    "SELECT COUNT(*) FROM log_entries WHERE daemon_id = ?1",
464                    [daemon_id.qualified()],
465                    |row| row.get(0),
466                )
467                .into_diagnostic()?;
468            count.saturating_sub(max_count as i64)
469        };
470
471        if to_delete <= 0 {
472            return Ok(0);
473        }
474
475        if let Some(hook) = hook {
476            let mut total_deleted = 0u64;
477            let mut remaining = to_delete;
478            loop {
479                let batch_len = remaining.min(hook.batch_size as i64);
480
481                // Fetch one batch under lock, then release.
482                let entries: Vec<LogEntry> = {
483                    let conn = self.conn.lock().unwrap();
484                    let mut stmt = conn
485                        .prepare(
486                            "SELECT id, daemon_id, timestamp, message, level, msg, logger, fields_json FROM log_entries
487                             WHERE daemon_id = ?1
488                             ORDER BY timestamp ASC, id ASC
489                             LIMIT ?2",
490                        )
491                        .into_diagnostic()?;
492                    stmt.query_map(
493                        params![daemon_id.qualified(), batch_len],
494                        Self::row_to_entry,
495                    )
496                    .into_diagnostic()?
497                    .collect::<rusqlite::Result<Vec<_>>>()
498                    .into_diagnostic()?
499                };
500
501                if entries.is_empty() {
502                    break;
503                }
504
505                let fetched = entries.len() as i64;
506                let batch_ids: Vec<i64> = entries.iter().map(|e| e.id).collect();
507
508                // Run the hook without holding the mutex.
509                self.archive_entries(&entries, hook, daemon_id, "count")?;
510
511                // Re-acquire lock and delete exactly the archived IDs.
512                let deleted = self.delete_by_ids(&batch_ids)?;
513                total_deleted += deleted;
514                remaining -= fetched;
515            }
516            Ok(total_deleted)
517        } else {
518            let conn = self.conn.lock().unwrap();
519            let rows = conn
520                .execute(
521                    "DELETE FROM log_entries WHERE id IN (
522                        SELECT id FROM log_entries WHERE daemon_id = ?1
523                        ORDER BY timestamp ASC, id ASC LIMIT ?2
524                    )",
525                    params![daemon_id.qualified(), to_delete],
526                )
527                .into_diagnostic()?;
528            Ok(rows as u64)
529        }
530    }
531
532    /// Migrate existing text logs for a daemon into SQLite.
533    ///
534    /// Reads the legacy text file line-by-line (streaming) and inserts in
535    /// batches of 1000 to avoid loading multi-GB files into memory at once.
536    pub fn migrate_daemon_text_logs(&self, daemon_id: &DaemonId) -> Result<u64> {
537        let text_path = daemon_id.log_path();
538        if !text_path.exists() {
539            return Ok(0);
540        }
541
542        let file = std::fs::File::open(&text_path).into_diagnostic()?;
543        let reader = BufReader::new(file);
544        let re = regex::Regex::new(r"^(\d{4}-\d{2}-\d{2} \d{2}:\d{2}:\d{2}) ([\w./-]+) (.*)$")
545            .expect("invalid regex");
546
547        let mut current_timestamp: Option<DateTime<Local>> = None;
548        let mut current_message = String::new();
549        let mut entries = Vec::with_capacity(1000);
550        let mut total_migrated: u64 = 0;
551
552        for line in reader.lines() {
553            let line = line.into_diagnostic()?;
554            if let Some(caps) = re.captures(&line) {
555                if let Some(ts) = current_timestamp.take() {
556                    entries.push((ts, std::mem::take(&mut current_message)));
557                }
558                let ts_str = caps.get(1).map(|m| m.as_str()).unwrap_or_default();
559                let msg = caps.get(3).map(|m| m.as_str()).unwrap_or_default();
560                if let Ok(naive) =
561                    chrono::NaiveDateTime::parse_from_str(ts_str, "%Y-%m-%d %H:%M:%S")
562                {
563                    current_timestamp = Local.from_local_datetime(&naive).single();
564                    current_message = msg.to_string();
565                }
566            } else if current_timestamp.is_some() {
567                current_message.push('\n');
568                current_message.push_str(&line);
569            }
570
571            if entries.len() >= 1000 {
572                total_migrated += self.insert_batch(daemon_id, &entries)?;
573                entries.clear();
574            }
575        }
576
577        if let Some(ts) = current_timestamp {
578            entries.push((ts, std::mem::take(&mut current_message)));
579        }
580
581        if !entries.is_empty() {
582            total_migrated += self.insert_batch(daemon_id, &entries)?;
583        }
584
585        if total_migrated > 0
586            && let Err(e) = std::fs::remove_file(&text_path)
587        {
588            log::warn!(
589                "failed to remove legacy log file after migration {}: {e}",
590                text_path.display()
591            );
592        }
593
594        Ok(total_migrated)
595    }
596
597    fn insert_batch(
598        &self,
599        daemon_id: &DaemonId,
600        entries: &[(DateTime<Local>, String)],
601    ) -> Result<u64> {
602        let mut conn = self.conn.lock().unwrap();
603        let tx = conn.transaction().into_diagnostic()?;
604        let mut count = 0u64;
605        {
606            let mut stmt = tx
607                .prepare(
608                    "INSERT INTO log_entries (daemon_id, timestamp, message) VALUES (?1, ?2, ?3)",
609                )
610                .into_diagnostic()?;
611            for (ts, msg) in entries {
612                stmt.execute(params![daemon_id.qualified(), ts.timestamp_millis(), msg])
613                    .into_diagnostic()?;
614                count += 1;
615            }
616        }
617        tx.commit().into_diagnostic()?;
618        Ok(count)
619    }
620
621    /// Build the SQL query string and parameters for the given options.
622    ///
623    /// `id_range` is used by `query_parallel` to shard the query by id range.
624    /// When `Some((start, end))`, adds `id > start AND id <= end` to the WHERE clause.
625    fn build_query_sql(
626        opts: &LogQuery,
627        id_range: Option<(i64, i64)>,
628    ) -> (String, Vec<Box<dyn rusqlite::ToSql>>) {
629        let mut conditions = Vec::new();
630        let mut query_params: Vec<Box<dyn rusqlite::ToSql>> = Vec::new();
631
632        if !opts.daemon_ids.is_empty() {
633            let placeholders: Vec<String> = (1..=opts.daemon_ids.len())
634                .map(|i| format!("?{}", i))
635                .collect();
636            conditions.push(format!("daemon_id IN ({})", placeholders.join(", ")));
637            for id in &opts.daemon_ids {
638                query_params.push(Box::new(id.clone()));
639            }
640        }
641
642        if let Some(from) = opts.from {
643            conditions.push(format!("timestamp >= ?{}", query_params.len() + 1));
644            query_params.push(Box::new(from.timestamp_millis()));
645        }
646
647        if let Some(to) = opts.to {
648            conditions.push(format!("timestamp <= ?{}", query_params.len() + 1));
649            query_params.push(Box::new(to.timestamp_millis()));
650        }
651
652        if let Some(after_id) = opts.after_id {
653            conditions.push(format!("id > ?{}", query_params.len() + 1));
654            query_params.push(Box::new(after_id));
655        }
656
657        if let Some(before_id) = opts.before_id {
658            conditions.push(format!("id < ?{}", query_params.len() + 1));
659            query_params.push(Box::new(before_id));
660        }
661
662        if let Some((start, end)) = id_range {
663            conditions.push(format!("id > ?{}", query_params.len() + 1));
664            query_params.push(Box::new(start));
665            conditions.push(format!("id <= ?{}", query_params.len() + 1));
666            query_params.push(Box::new(end));
667        }
668
669        let mut message_conditions = Vec::new();
670        for filter in &opts.message_filters {
671            match filter {
672                MessageFilter::Contains {
673                    pattern,
674                    case_sensitive,
675                } => {
676                    let param_index = query_params.len() + 1;
677                    if *case_sensitive {
678                        message_conditions
679                            .push(format!("INSTR(message, ?{idx}) > 0", idx = param_index));
680                        query_params.push(Box::new(pattern.clone()));
681                    } else {
682                        let escaped = escape_like_pattern(pattern);
683                        let param = format!("%{}%", escaped);
684                        message_conditions.push(format!(
685                            "LOWER(message) LIKE LOWER(?{idx}) ESCAPE '\\'",
686                            idx = param_index
687                        ));
688                        query_params.push(Box::new(param));
689                    }
690                }
691                MessageFilter::Regex { pattern } => {
692                    let param_index = query_params.len() + 1;
693                    message_conditions.push(format!("message REGEXP ?{param_index}"));
694                    query_params.push(Box::new(pattern.clone()));
695                }
696            }
697        }
698        if !message_conditions.is_empty() {
699            conditions.push(format!("({})", message_conditions.join(" OR ")));
700        }
701
702        for filter in &opts.field_filters {
703            match filter {
704                FieldFilter::LevelMin(level) => {
705                    let matching = crate::log_store::levels_at_or_above(level);
706                    if matching.is_empty() {
707                        // Unknown level: no results
708                        conditions.push("0".to_string());
709                    } else {
710                        let placeholders = matching
711                            .iter()
712                            .map(|l| {
713                                let idx = query_params.len() + 1;
714                                query_params.push(Box::new((*l).to_string()));
715                                format!("?{idx}")
716                            })
717                            .collect::<Vec<_>>()
718                            .join(", ");
719                        conditions.push(format!("level IN ({placeholders})"));
720                    }
721                }
722                FieldFilter::FieldEq { key, value } => {
723                    let key_idx = query_params.len() + 1;
724                    let val_idx = query_params.len() + 2;
725                    // Use json_each with parameterized key to avoid JSON path
726                    // interpolation. This correctly handles keys containing '.'
727                    // (e.g. "request.id") which json_extract would treat as
728                    // nested traversal, and prevents SQL injection from
729                    // untrusted key input (CLI --field doesn't validate keys).
730                    //
731                    // Match both the typed JSON value and the raw text.
732                    // json_each.value has no type affinity, so comparing it
733                    // IS ? (text) only matches string values, not integers
734                    // or booleans. This prevents a string field storing "42"
735                    // or "true" from being unreachable via field filters.
736                    //
737                    // Typed path:   json_each.value IS json_extract(json_literal, '$')
738                    //   matches numbers, booleans, null (native SQLite types)
739                    // Text path:    json_each.value IS ? (raw text)
740                    //   matches string fields whose value textually resembles
741                    //   a primitive (e.g. field storing "true", "42", "null")
742                    let json_literal = text_to_json_literal(value);
743                    let raw_idx = query_params.len() + 3;
744                    conditions.push(format!(
745                        "EXISTS (SELECT 1 FROM json_each(fields_json) \
746                         WHERE json_each.key = ?{key_idx} \
747                           AND (json_each.value IS json_extract(?{val_idx}, '$') \
748                                OR json_each.value IS ?{raw_idx}))"
749                    ));
750                    query_params.push(Box::new(key.clone()));
751                    query_params.push(Box::new(json_literal));
752                    query_params.push(Box::new(value.clone()));
753                }
754                FieldFilter::LoggerContains(pattern) => {
755                    let param_index = query_params.len() + 1;
756                    conditions.push(format!("logger LIKE ?{param_index} ESCAPE '\\'"));
757                    let escaped = crate::log_store::escape_like_pattern(pattern);
758                    query_params.push(Box::new(format!("%{escaped}%")));
759                }
760            }
761        }
762
763        let where_clause = if conditions.is_empty() {
764            String::new()
765        } else {
766            format!("WHERE {}", conditions.join(" AND "))
767        };
768
769        let order = if opts.order_desc { "DESC" } else { "ASC" };
770
771        let limit_clause = opts
772            .limit
773            .map(|n| format!("LIMIT {}", n))
774            .unwrap_or_default();
775
776        let columns = if opts.include_structured {
777            "id, daemon_id, timestamp, message, level, msg, logger, fields_json"
778        } else {
779            "id, daemon_id, timestamp, message, NULL, NULL, NULL, NULL"
780        };
781
782        let sql = format!(
783            "SELECT {columns} FROM log_entries {} ORDER BY timestamp {}, id {} {}",
784            where_clause, order, order, limit_clause
785        );
786
787        (sql, query_params)
788    }
789
790    /// Returns true if parallel query is beneficial for the given options.
791    fn should_parallelize(opts: &LogQuery) -> bool {
792        if opts.daemon_ids.len() != 1 {
793            return false;
794        }
795        // Skip parallel for incremental tail polls (after_id set) — they
796        // return small batches and the connection overhead dominates.
797        if opts.after_id.is_some() {
798            return false;
799        }
800        let limit = opts.limit.unwrap_or(usize::MAX);
801        if limit < PARALLEL_QUERY_THRESHOLD {
802            return false;
803        }
804        std::thread::available_parallelism()
805            .map(|n| n.get() >= 2)
806            .unwrap_or(false)
807    }
808
809    /// Query using multiple read-only connections, sharded by id range.
810    ///
811    /// For a single daemon, id order ≈ timestamp order (timestamps are
812    /// assigned in insertion order within a single writer). This lets us
813    /// shard by contiguous id ranges and merge by concatenating shards in
814    /// the right order, without a global sort.
815    fn query_parallel(&self, opts: &LogQuery) -> Result<Vec<LogEntry>> {
816        // Cap parallelism: each shard opens a new read-only connection with
817        // PRAGMA setup + regexp function registration (~5-12ms each). Beyond
818        // 2 threads the connection overhead dominates the query savings for
819        // typical log store sizes.
820        let max_threads = 2;
821        let num_threads = std::thread::available_parallelism()
822            .map(|n| n.get().min(max_threads))
823            .unwrap_or(1);
824
825        let max_id: Option<i64> = {
826            let conn = self.conn.lock().unwrap();
827            conn.query_row(
828                "SELECT MAX(id) FROM log_entries WHERE daemon_id = ?1",
829                params![&opts.daemon_ids[0]],
830                |row| row.get(0),
831            )
832            .ok()
833        };
834
835        let Some(max_id) = max_id else {
836            return Ok(Vec::new());
837        };
838        if max_id == 0 {
839            return Ok(Vec::new());
840        }
841
842        let shard_size = (max_id as usize).div_ceil(num_threads);
843        let path = self.path.clone();
844        let opts = opts.clone();
845        let needs_regexp = opts
846            .message_filters
847            .iter()
848            .any(|f| matches!(f, MessageFilter::Regex { .. }));
849
850        let shards: Vec<Result<Vec<LogEntry>>> = std::thread::scope(|s| {
851            (0..num_threads)
852                .map(|i| {
853                    let start = (i * shard_size) as i64;
854                    let end = if i == num_threads - 1 {
855                        max_id
856                    } else {
857                        ((i + 1) * shard_size) as i64
858                    };
859                    let opts = opts.clone();
860                    let path = path.clone();
861                    s.spawn(move || -> Result<Vec<LogEntry>> {
862                        let conn = Connection::open_with_flags(
863                            &path,
864                            OpenFlags::SQLITE_OPEN_READ_ONLY | OpenFlags::SQLITE_OPEN_NO_MUTEX,
865                        )
866                        .into_diagnostic()?;
867                        conn.execute_batch(
868                            "PRAGMA mmap_size = 268435456;
869                             PRAGMA query_only = 1;",
870                        )
871                        .into_diagnostic()?;
872                        if needs_regexp {
873                            add_regexp_function(&conn)?;
874                        }
875
876                        let (sql, query_params) = Self::build_query_sql(&opts, Some((start, end)));
877                        Self::execute_built_query(&conn, &sql, &query_params)
878                    })
879                })
880                .map(|h| h.join().unwrap())
881                .collect()
882        });
883
884        let mut merged = Vec::new();
885        if opts.order_desc {
886            for shard in shards.into_iter().rev() {
887                merged.extend(shard?);
888            }
889        } else {
890            for shard in shards {
891                merged.extend(shard?);
892            }
893        }
894
895        if let Some(limit) = opts.limit
896            && merged.len() > limit
897        {
898            merged.truncate(limit);
899        }
900
901        Ok(merged)
902    }
903
904    /// Execute a built SQL query and collect results into LogEntry.
905    fn execute_built_query(
906        conn: &Connection,
907        sql: &str,
908        query_params: &[Box<dyn rusqlite::ToSql>],
909    ) -> Result<Vec<LogEntry>> {
910        let mut stmt = conn.prepare(sql).into_diagnostic()?;
911        let params_ref: Vec<&dyn rusqlite::ToSql> =
912            query_params.iter().map(|p| p.as_ref()).collect();
913        let rows = stmt
914            .query_map(params_ref.as_slice(), Self::row_to_entry)
915            .into_diagnostic()?;
916        let mut entries = Vec::new();
917        for row in rows {
918            entries.push(row.into_diagnostic()?);
919        }
920        Ok(entries)
921    }
922
923    /// Return distinct non-null logger values for a daemon, sorted alphabetically.
924    ///
925    /// Scans only the most recent 5000 entries for performance — autocomplete
926    /// only needs a representative sample, not the full history.
927    pub fn distinct_loggers(&self, daemon_id: &str) -> Result<Vec<String>> {
928        let conn = self.conn.lock().unwrap();
929        let mut stmt = conn
930            .prepare(
931                "SELECT DISTINCT logger FROM ( \
932                  SELECT logger FROM log_entries \
933                  WHERE daemon_id = ?1 AND logger IS NOT NULL \
934                  ORDER BY id DESC LIMIT 5000 \
935                 ) ORDER BY logger",
936            )
937            .into_diagnostic()?;
938        let rows = stmt
939            .query_map(params![daemon_id], |row| row.get::<_, String>(0))
940            .into_diagnostic()?;
941        let mut loggers = Vec::new();
942        for row in rows {
943            loggers.push(row.into_diagnostic()?);
944        }
945        Ok(loggers)
946    }
947
948    /// Return distinct keys from the `fields_json` column for a daemon.
949    ///
950    /// Uses `json_each` to extract all object keys across recent structured log
951    /// entries, returning a sorted unique list. Useful for jq autocomplete.
952    /// Scans only the most recent 5000 entries for performance.
953    pub fn distinct_field_keys(&self, daemon_id: &str) -> Result<Vec<String>> {
954        let conn = self.conn.lock().unwrap();
955        let mut stmt = conn
956            .prepare(
957                "SELECT DISTINCT je.key \
958                 FROM ( \
959                   SELECT fields_json FROM log_entries \
960                   WHERE daemon_id = ?1 AND fields_json IS NOT NULL \
961                   ORDER BY id DESC LIMIT 5000 \
962                 ), json_each(fields_json) AS je \
963                 ORDER BY je.key",
964            )
965            .into_diagnostic()?;
966        let rows = stmt
967            .query_map(params![daemon_id], |row| row.get::<_, String>(0))
968            .into_diagnostic()?;
969        let mut keys = Vec::new();
970        for row in rows {
971            keys.push(row.into_diagnostic()?);
972        }
973        Ok(keys)
974    }
975}
976
977impl LogStore for SqliteLogStore {
978    fn append(&self, daemon_id: &DaemonId, message: &str) -> Result<()> {
979        let id = daemon_id.qualified();
980        let msg = message.to_string();
981
982        let conn = self.conn.lock().unwrap();
983        // Sample while holding the lock so timestamp order matches insertion
984        // order under concurrent writers.
985        let ts = Local::now().timestamp_millis();
986        let _ = conn
987            .execute(
988                "INSERT INTO log_entries (daemon_id, timestamp, message) VALUES (?1, ?2, ?3)",
989                params![id, ts, msg],
990            )
991            .into_diagnostic()?;
992        Ok(())
993    }
994
995    fn append_batch(&self, daemon_id: &DaemonId, messages: &[String]) -> Result<()> {
996        if messages.is_empty() {
997            return Ok(());
998        }
999        let id = daemon_id.qualified();
1000
1001        let mut conn = self.conn.lock().unwrap();
1002        // IMMEDIATE acquires the write lock at BEGIN, and we sample only after
1003        // that, so timestamp order matches commit order even across processes
1004        // (sink subprocesses and the supervisor each hold their own
1005        // SqliteLogStore instance on the same database).
1006        let tx = conn
1007            .transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)
1008            .into_diagnostic()?;
1009        let base_ts = Local::now().timestamp_millis();
1010        {
1011            let mut stmt = tx
1012                .prepare(
1013                    "INSERT INTO log_entries (daemon_id, timestamp, message) VALUES (?1, ?2, ?3)",
1014                )
1015                .into_diagnostic()?;
1016            for msg in messages.iter() {
1017                // Use the same timestamp for the whole batch: `id` is
1018                // AUTOINCREMENT (insertion order), and queries order by
1019                // (timestamp, id), so equal timestamps still sort by
1020                // insertion order. Staggering by `idx` instead overlaps with
1021                // the timestamps of a later batch flushed within this batch's
1022                // fabricated span, interleaving the two batches.
1023                let ts = base_ts;
1024                stmt.execute(params![id, ts, msg]).into_diagnostic()?;
1025            }
1026        }
1027        tx.commit().into_diagnostic()?;
1028        Ok(())
1029    }
1030
1031    fn append_structured(&self, daemon_id: &DaemonId, parsed: &ParsedLog) -> Result<()> {
1032        let id = daemon_id.qualified();
1033
1034        let conn = self.conn.lock().unwrap();
1035        // Sample while holding the lock so timestamp order matches insertion
1036        // order under concurrent writers.
1037        let ts = Local::now().timestamp_millis();
1038        let _ = conn
1039            .execute(
1040                "INSERT INTO log_entries (daemon_id, timestamp, message, level, msg, logger, fields_json) \
1041                 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)",
1042                params![id, ts, parsed.message, parsed.level, parsed.msg, parsed.logger, parsed.fields_json],
1043            )
1044            .into_diagnostic()?;
1045        Ok(())
1046    }
1047
1048    fn append_structured_batch(&self, daemon_id: &DaemonId, entries: &[ParsedLog]) -> Result<()> {
1049        if entries.is_empty() {
1050            return Ok(());
1051        }
1052        let id = daemon_id.qualified();
1053
1054        let mut conn = self.conn.lock().unwrap();
1055        // IMMEDIATE acquires the write lock at BEGIN, and we sample only after
1056        // that, so timestamp order matches commit order even across processes
1057        // (sink subprocesses and the supervisor each hold their own
1058        // SqliteLogStore instance on the same database).
1059        let tx = conn
1060            .transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)
1061            .into_diagnostic()?;
1062        let base_ts = Local::now().timestamp_millis();
1063        {
1064            let mut stmt = tx
1065                .prepare(
1066                    "INSERT INTO log_entries (daemon_id, timestamp, message, level, msg, logger, fields_json) \
1067                     VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)",
1068                )
1069                .into_diagnostic()?;
1070            for entry in entries.iter() {
1071                // Same timestamp for the whole batch; queries order by
1072                // (timestamp, id) and `id` is AUTOINCREMENT, so insertion
1073                // order is preserved without the cross-batch overlap that
1074                // staggered timestamps (`base_ts + idx`) caused.
1075                let ts = base_ts;
1076                stmt.execute(params![
1077                    id,
1078                    ts,
1079                    entry.message,
1080                    entry.level,
1081                    entry.msg,
1082                    entry.logger,
1083                    entry.fields_json
1084                ])
1085                .into_diagnostic()?;
1086            }
1087        }
1088        tx.commit().into_diagnostic()?;
1089        Ok(())
1090    }
1091
1092    fn query(&self, opts: &LogQuery) -> Result<Vec<LogEntry>> {
1093        // Delegate to parallel path for large single-daemon queries.
1094        if Self::should_parallelize(opts)
1095            && self.path.as_os_str() != ":memory:"
1096            && let Ok(entries) = self.query_parallel(opts)
1097        {
1098            return Ok(entries);
1099        }
1100        // Fall back to single-threaded on parallel failure.
1101
1102        // Single-threaded path.
1103        let conn = self.conn.lock().unwrap();
1104        let (sql, query_params) = Self::build_query_sql(opts, None);
1105        Self::execute_built_query(&conn, &sql, &query_params)
1106    }
1107
1108    fn tail(&self, daemon_id: &DaemonId, after_id: Option<i64>) -> Result<Vec<LogEntry>> {
1109        self.query(&LogQuery {
1110            daemon_ids: vec![daemon_id.qualified()],
1111            from: None,
1112            to: None,
1113            limit: None,
1114            order_desc: false,
1115            after_id,
1116            before_id: None,
1117            message_filters: Vec::new(),
1118            field_filters: Vec::new(),
1119            include_structured: false,
1120        })
1121    }
1122
1123    fn clear(&self, daemon_ids: &[DaemonId]) -> Result<()> {
1124        let mut conn = self.conn.lock().unwrap();
1125        let tx = conn.transaction().into_diagnostic()?;
1126        for id in daemon_ids {
1127            tx.execute(
1128                "DELETE FROM log_entries WHERE daemon_id = ?1",
1129                params![id.qualified()],
1130            )
1131            .into_diagnostic()?;
1132            tx.execute(
1133                "INSERT INTO log_clear_generations (daemon_id, generation)
1134                 VALUES (?1, 1)
1135                 ON CONFLICT(daemon_id) DO UPDATE SET generation = generation + 1",
1136                params![id.qualified()],
1137            )
1138            .into_diagnostic()?;
1139        }
1140        tx.commit().into_diagnostic()?;
1141        Ok(())
1142    }
1143
1144    fn last_id(&self, daemon_id: &DaemonId) -> Result<Option<i64>> {
1145        let conn = self.conn.lock().unwrap();
1146        // MAX(id) returns NULL when no rows exist for the daemon; the
1147        // Option<i64> decode maps NULL to None automatically.
1148        let id: Option<i64> = conn
1149            .query_row(
1150                "SELECT MAX(id) FROM log_entries WHERE daemon_id = ?1",
1151                params![daemon_id.qualified()],
1152                |row| row.get(0),
1153            )
1154            .into_diagnostic()?;
1155        Ok(id)
1156    }
1157
1158    fn list_daemon_ids(&self) -> Result<Vec<String>> {
1159        let conn = self.conn.lock().unwrap();
1160        let mut stmt = conn
1161            .prepare("SELECT DISTINCT daemon_id FROM log_entries")
1162            .into_diagnostic()?;
1163        let ids = stmt
1164            .query_map([], |row| {
1165                let id: String = row.get(0)?;
1166                Ok(id)
1167            })
1168            .into_diagnostic()?
1169            .filter_map(|r| r.ok())
1170            .collect();
1171        Ok(ids)
1172    }
1173
1174    fn apply_retention(
1175        &self,
1176        policy: &super::RetentionPolicy,
1177        excluded_daemon_ids: &[DaemonId],
1178        archive_hook: Option<&ArchiveHook>,
1179    ) -> Result<u64> {
1180        let daemon_ids = self.list_daemon_ids()?;
1181        let excluded: HashSet<String> = excluded_daemon_ids.iter().map(|d| d.qualified()).collect();
1182        let mut total = 0u64;
1183        for id_str in daemon_ids {
1184            if excluded.contains(&id_str) {
1185                continue;
1186            }
1187            let id = DaemonId::parse(&id_str).unwrap_or_else(|_| {
1188                DaemonId::try_new("global", &id_str).unwrap_or_else(|_| DaemonId::pitchfork())
1189            });
1190            if let Some(dur) = policy.age {
1191                total += self.rotate_by_age(&id, dur, archive_hook)?;
1192            }
1193            if let Some(n) = policy.count {
1194                total += self.rotate_by_count(&id, n, archive_hook)?;
1195            }
1196        }
1197        Ok(total)
1198    }
1199
1200    fn apply_retention_for_daemon(
1201        &self,
1202        daemon_id: &DaemonId,
1203        policy: &super::RetentionPolicy,
1204        archive_hook: Option<&ArchiveHook>,
1205    ) -> Result<u64> {
1206        let mut total = 0u64;
1207        if let Some(dur) = policy.age {
1208            total += self.rotate_by_age(daemon_id, dur, archive_hook)?;
1209        }
1210        if let Some(n) = policy.count {
1211            total += self.rotate_by_count(daemon_id, n, archive_hook)?;
1212        }
1213        Ok(total)
1214    }
1215
1216    fn last_clear_generation(&self, daemon_id: &DaemonId) -> Result<Option<u64>> {
1217        let conn = self.conn.lock().unwrap();
1218        let mut stmt = conn
1219            .prepare("SELECT generation FROM log_clear_generations WHERE daemon_id = ?1")
1220            .into_diagnostic()?;
1221        let generation: Option<i64> = stmt
1222            .query_row(params![daemon_id.qualified()], |row| row.get(0))
1223            .optional()
1224            .into_diagnostic()?;
1225        generation
1226            .map(|generation| {
1227                u64::try_from(generation)
1228                    .map_err(|_| miette::miette!("log clear generation cannot be negative"))
1229            })
1230            .transpose()
1231    }
1232
1233    fn query_with_generation(
1234        &self,
1235        opts: &LogQuery,
1236        daemon_id: &DaemonId,
1237    ) -> Result<(Vec<LogEntry>, Option<u64>)> {
1238        let conn = self.conn.lock().unwrap();
1239        // Wrap both reads in a single transaction so a concurrent clear
1240        // cannot interleave between them. In WAL mode, BEGIN acquires a
1241        // consistent snapshot that both statements share.
1242        //
1243        // Intentionally call build_query_sql + execute_built_query directly
1244        // instead of self.query(): the latter may delegate to query_parallel
1245        // which opens separate connections outside this transaction, breaking
1246        // the snapshot guarantee.
1247        conn.execute_batch("BEGIN").into_diagnostic()?;
1248        let result = (|| {
1249            let (sql, query_params) = Self::build_query_sql(opts, None);
1250            let entries = Self::execute_built_query(&conn, &sql, &query_params)?;
1251            let generation: Option<i64> = conn
1252                .query_row(
1253                    "SELECT generation FROM log_clear_generations WHERE daemon_id = ?1",
1254                    params![daemon_id.qualified()],
1255                    |row| row.get(0),
1256                )
1257                .optional()
1258                .into_diagnostic()?;
1259            let generation = generation
1260                .map(|g| {
1261                    u64::try_from(g)
1262                        .map_err(|_| miette::miette!("log clear generation cannot be negative"))
1263                })
1264                .transpose()?;
1265            Ok((entries, generation))
1266        })();
1267        // Always rollback (read transaction just needs to end).
1268        let _ = conn.execute_batch("ROLLBACK");
1269        result
1270    }
1271}
1272
1273/// Global singleton log store.
1274use once_cell::sync::Lazy;
1275use std::sync::Arc;
1276
1277pub static LOG_STORE: Lazy<Arc<SqliteLogStore>> = Lazy::new(|| {
1278    let path = crate::env::PITCHFORK_LOGS_DIR.join("logs.db");
1279    let mut is_fallback = false;
1280    let store = Arc::new(SqliteLogStore::open(&path).unwrap_or_else(|e| {
1281        error!(
1282            "failed to open log store at {}: {e}. Falling back to in-memory store; logs will not persist across restarts.",
1283            path.display()
1284        );
1285        is_fallback = true;
1286        SqliteLogStore::open(":memory:").expect("in-memory SQLite should always open")
1287    }));
1288
1289    // Auto-migrate any legacy text log files into SQLite on first access.
1290    // This runs once per process startup and is idempotent.
1291    // Skip migration when using the in-memory fallback to prevent data loss.
1292    if !is_fallback {
1293        if let Err(e) = auto_migrate_legacy_logs(&store) {
1294            warn!("legacy log auto-migration failed: {e}");
1295        }
1296    } else {
1297        warn!(
1298            "skipping legacy log auto-migration because log store is in-memory (no durable destination)"
1299        );
1300    }
1301
1302    store
1303});
1304
1305/// Auto-migrate legacy text log files into the SQLite log store.
1306///
1307/// Scans the logs directory for directories matching the new-format layout
1308/// (`namespace--name/namespace--name.log`), attempts to parse the directory
1309/// name as a valid safe-path daemon ID, and imports the content into SQLite.
1310/// This is idempotent: re-running it on already-migrated data is a no-op
1311/// because the legacy text files are deleted after successful import.
1312fn auto_migrate_legacy_logs(store: &SqliteLogStore) -> Result<()> {
1313    let logs_dir = &*crate::env::PITCHFORK_LOGS_DIR;
1314    if !logs_dir.exists() {
1315        return Ok(());
1316    }
1317
1318    let Ok(entries) = std::fs::read_dir(logs_dir) else {
1319        return Ok(());
1320    };
1321
1322    let mut total_migrated = 0u64;
1323    let mut migrated_ids = Vec::new();
1324
1325    for entry in entries.flatten() {
1326        let path = entry.path();
1327        if !path.is_dir() {
1328            continue;
1329        }
1330        // Skip the supervisor's own log directory.
1331        let file_name = path
1332            .file_name()
1333            .map_or(String::new(), |n| n.to_string_lossy().to_string());
1334        if file_name == "pitchfork" {
1335            continue;
1336        }
1337
1338        // Only consider directories that look like new-format safe-paths
1339        if !file_name.contains("--") {
1340            continue;
1341        }
1342        let log_file = path.join(format!("{file_name}.log"));
1343        if !log_file.exists() {
1344            continue;
1345        }
1346
1347        let daemon_id = match DaemonId::from_safe_path(&file_name) {
1348            Ok(id) => id,
1349            Err(_) => continue,
1350        };
1351
1352        // Skip the supervisor's own daemon; its log directory may be present
1353        // under logs/ but should never be imported into the user-facing log store.
1354        if daemon_id == DaemonId::pitchfork() {
1355            continue;
1356        }
1357
1358        match store.migrate_daemon_text_logs(&daemon_id) {
1359            Ok(0) => {}
1360            Ok(n) => {
1361                total_migrated += n;
1362                migrated_ids.push(daemon_id.qualified());
1363            }
1364            Err(e) => {
1365                warn!(
1366                    "failed to migrate text logs for {}: {e}",
1367                    daemon_id.qualified()
1368                );
1369            }
1370        }
1371    }
1372
1373    if total_migrated > 0 {
1374        warn!(
1375            "auto-migrated {total_migrated} legacy log entries from {count} daemon(s): {ids}",
1376            count = migrated_ids.len(),
1377            ids = migrated_ids.join(", ")
1378        );
1379    }
1380
1381    Ok(())
1382}
1383
1384#[cfg(test)]
1385mod tests {
1386    use super::text_to_json_literal;
1387    use super::*;
1388    use crate::log_store::LogStore;
1389
1390    #[test]
1391    fn test_json_literal_booleans() {
1392        assert_eq!(text_to_json_literal("true"), "true");
1393        assert_eq!(text_to_json_literal("TRUE"), "true");
1394        assert_eq!(text_to_json_literal("false"), "false");
1395        assert_eq!(text_to_json_literal("False"), "false");
1396    }
1397
1398    #[test]
1399    fn test_json_literal_null() {
1400        assert_eq!(text_to_json_literal("null"), "null");
1401        assert_eq!(text_to_json_literal("NULL"), "null");
1402    }
1403
1404    #[test]
1405    fn test_json_literal_valid_numbers() {
1406        assert_eq!(text_to_json_literal("42"), "42");
1407        assert_eq!(text_to_json_literal("-1"), "-1");
1408        assert_eq!(text_to_json_literal("0"), "0");
1409        assert_eq!(text_to_json_literal("3.14"), "3.14");
1410        assert_eq!(text_to_json_literal("1e10"), "1e10");
1411        assert_eq!(text_to_json_literal("1.5e-3"), "1.5e-3");
1412    }
1413
1414    #[test]
1415    fn test_json_literal_rejects_invalid_json_numbers() {
1416        // f64::parse accepts these but they are not valid JSON numbers,
1417        // so they must be treated as strings (quoted) instead.
1418        assert_eq!(text_to_json_literal("+42"), r#""+42""#);
1419        assert_eq!(text_to_json_literal("1."), r#""1.""#);
1420        assert_eq!(text_to_json_literal(".5"), r#"".5""#);
1421        assert_eq!(text_to_json_literal("inf"), r#""inf""#);
1422        assert_eq!(text_to_json_literal("nan"), r#""nan""#);
1423        assert_eq!(text_to_json_literal("infinity"), r#""infinity""#);
1424    }
1425
1426    #[test]
1427    fn test_json_literal_strings() {
1428        assert_eq!(text_to_json_literal("hello"), r#""hello""#);
1429        assert_eq!(text_to_json_literal("req_1"), r#""req_1""#);
1430        // String with special chars gets JSON-escaped
1431        assert_eq!(text_to_json_literal(r#"a"b"#), r#""a\"b""#);
1432    }
1433
1434    /// A writer that finds the lock held must wait for it, not give up. WAL
1435    /// serializes writers, so without `busy_timeout` the second writer takes
1436    /// SQLITE_BUSY immediately and silently discards its batch.
1437    #[test]
1438    fn write_waits_for_a_contended_lock() {
1439        let dir = tempfile::tempdir().unwrap();
1440        let path = dir.path().join("logs.db");
1441        let store = SqliteLogStore::open(&path).unwrap();
1442
1443        // Hold the write lock from a second connection, then release it after a
1444        // delay that a non-waiting writer could never survive.
1445        let blocker = Connection::open(&path).unwrap();
1446        blocker.busy_timeout(BUSY_TIMEOUT).unwrap();
1447        blocker.execute_batch("BEGIN IMMEDIATE").unwrap();
1448        let releaser = std::thread::spawn(move || {
1449            std::thread::sleep(std::time::Duration::from_millis(300));
1450            blocker.execute_batch("COMMIT").unwrap();
1451        });
1452
1453        let id = DaemonId::try_new("test", "blocked").unwrap();
1454        let entries = vec![crate::log_parse::parse("held-lock-line", "text")];
1455        store
1456            .append_structured_batch(&id, &entries)
1457            .expect("a contended write must wait for the lock, not fail");
1458
1459        releaser.join().unwrap();
1460        let found = store
1461            .query(&LogQuery {
1462                daemon_ids: vec![id.qualified()],
1463                ..Default::default()
1464            })
1465            .unwrap()
1466            .len();
1467        assert_eq!(found, 1, "the batch written under contention was lost");
1468    }
1469
1470    /// Two processes writing the same store is the normal case, not an edge
1471    /// case: daemon output and retention pruning already collide, and a
1472    /// per-daemon writer multiplies the opportunities. WAL serializes writers,
1473    /// so without `busy_timeout` the loser of that race takes SQLITE_BUSY
1474    /// immediately and silently discards its batch.
1475    #[test]
1476    fn concurrent_writers_do_not_lose_batches() {
1477        let dir = tempfile::tempdir().unwrap();
1478        let path = dir.path().join("logs.db");
1479
1480        const WRITERS: usize = 4;
1481        const BATCHES: usize = 15;
1482        const PER_BATCH: usize = 10;
1483
1484        // Release every worker into open() at the same moment. Without this
1485        // the scheduler is free to run them one after another, so the first
1486        // would enable WAL and the rest would never contend for the
1487        // journal-mode switch — leaving the race this guards unexercised.
1488        let barrier = std::sync::Barrier::new(WRITERS);
1489
1490        std::thread::scope(|scope| {
1491            for writer in 0..WRITERS {
1492                let path = path.clone();
1493                let barrier = &barrier;
1494                scope.spawn(move || {
1495                    // A separate connection per writer, as separate processes
1496                    // would have.
1497                    barrier.wait();
1498                    let store = SqliteLogStore::open(&path).unwrap();
1499                    let id = DaemonId::try_new("test", format!("w{writer}")).unwrap();
1500                    for batch in 0..BATCHES {
1501                        let entries: Vec<ParsedLog> = (0..PER_BATCH)
1502                            .map(|i| {
1503                                crate::log_parse::parse(&format!("w{writer}-{batch}-{i}"), "text")
1504                            })
1505                            .collect();
1506                        store
1507                            .append_structured_batch(&id, &entries)
1508                            .expect("concurrent batch write must not fail");
1509                    }
1510                });
1511            }
1512        });
1513
1514        let store = SqliteLogStore::open(&path).unwrap();
1515        for writer in 0..WRITERS {
1516            let id = DaemonId::try_new("test", format!("w{writer}")).unwrap();
1517            let found = store
1518                .query(&LogQuery {
1519                    daemon_ids: vec![id.qualified()],
1520                    ..Default::default()
1521                })
1522                .unwrap()
1523                .len();
1524            assert_eq!(
1525                found,
1526                BATCHES * PER_BATCH,
1527                "writer {writer} lost entries: got {found}"
1528            );
1529        }
1530    }
1531
1532    /// Regression for the "burst of log lines comes back out of order" bug:
1533    /// rows within one batch must share a single timestamp instead of being
1534    /// staggered by index. Staggered timestamps (`base_ts + idx`) overlap with
1535    /// the timestamps of the next batch flushed moments later, interleaving the
1536    /// two batches under `ORDER BY timestamp, id`.
1537    #[test]
1538    fn batch_rows_share_single_timestamp() {
1539        let dir = tempfile::tempdir().unwrap();
1540        let path = dir.path().join("logs.db");
1541        let store = SqliteLogStore::open(&path).unwrap();
1542        let id = DaemonId::try_new("test", "burst").unwrap();
1543
1544        let entries: Vec<ParsedLog> = (0..50)
1545            .map(|i| crate::log_parse::parse(&format!("line-{i}"), "text"))
1546            .collect();
1547        store.append_structured_batch(&id, &entries).unwrap();
1548
1549        let conn = store.conn.lock().unwrap();
1550        let timestamps: Vec<i64> = {
1551            let mut stmt = conn
1552                .prepare("SELECT timestamp FROM log_entries ORDER BY id")
1553                .unwrap();
1554            let rows = stmt.query_map([], |row| row.get::<_, i64>(0)).unwrap();
1555            rows.map(|r| r.unwrap()).collect()
1556        };
1557        assert_eq!(timestamps.len(), 50);
1558        let distinct: std::collections::HashSet<i64> = timestamps.into_iter().collect();
1559        assert_eq!(
1560            distinct.len(),
1561            1,
1562            "all rows in a batch must share one timestamp so (timestamp, id) ordering matches insertion order"
1563        );
1564    }
1565
1566    /// End-to-end regression for the reported repro (`run = "exec seq 500"`):
1567    /// several batches flushed back-to-back must still read back in insertion
1568    /// order. With staggered per-row timestamps the later batch's fabricated
1569    /// span overlaps the earlier one and interleaves the output.
1570    #[test]
1571    fn burst_batches_preserve_insertion_order() {
1572        let dir = tempfile::tempdir().unwrap();
1573        let path = dir.path().join("logs.db");
1574        let store = SqliteLogStore::open(&path).unwrap();
1575        let id = DaemonId::try_new("test", "burst").unwrap();
1576
1577        // 3 batches of 300 rows, appended as fast as possible — the staggered
1578        // spans (300ms each) are far wider than the inter-batch gap, so the
1579        // old code always overlapped and interleaved these.
1580        const BATCHES: usize = 3;
1581        const PER_BATCH: usize = 300;
1582        for batch in 0..BATCHES {
1583            let entries: Vec<ParsedLog> = (0..PER_BATCH)
1584                .map(|i| crate::log_parse::parse(&format!("{}", batch * PER_BATCH + i), "text"))
1585                .collect();
1586            store.append_structured_batch(&id, &entries).unwrap();
1587        }
1588
1589        let entries = store
1590            .query(&LogQuery {
1591                daemon_ids: vec![id.qualified()],
1592                ..Default::default()
1593            })
1594            .unwrap();
1595        assert_eq!(entries.len(), BATCHES * PER_BATCH);
1596        for (n, entry) in entries.iter().enumerate() {
1597            assert_eq!(
1598                entry.message,
1599                n.to_string(),
1600                "log line {n} read back out of order (got '{}')",
1601                entry.message
1602            );
1603        }
1604    }
1605}