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