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