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