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 id = daemon_id.qualified();
968        let msg = message.to_string();
969
970        let conn = self.conn.lock().unwrap();
971        // Sample while holding the lock so timestamp order matches insertion
972        // order under concurrent writers.
973        let ts = Local::now().timestamp_millis();
974        let _ = conn
975            .execute(
976                "INSERT INTO log_entries (daemon_id, timestamp, message) VALUES (?1, ?2, ?3)",
977                params![id, ts, msg],
978            )
979            .into_diagnostic()?;
980        Ok(())
981    }
982
983    fn append_batch(&self, daemon_id: &DaemonId, messages: &[String]) -> Result<()> {
984        if messages.is_empty() {
985            return Ok(());
986        }
987        let id = daemon_id.qualified();
988
989        let mut conn = self.conn.lock().unwrap();
990        // IMMEDIATE acquires the write lock at BEGIN, and we sample only after
991        // that, so timestamp order matches commit order even across processes
992        // (sink subprocesses and the supervisor each hold their own
993        // SqliteLogStore instance on the same database).
994        let tx = conn
995            .transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)
996            .into_diagnostic()?;
997        let base_ts = Local::now().timestamp_millis();
998        {
999            let mut stmt = tx
1000                .prepare(
1001                    "INSERT INTO log_entries (daemon_id, timestamp, message) VALUES (?1, ?2, ?3)",
1002                )
1003                .into_diagnostic()?;
1004            for msg in messages.iter() {
1005                // Use the same timestamp for the whole batch: `id` is
1006                // AUTOINCREMENT (insertion order), and queries order by
1007                // (timestamp, id), so equal timestamps still sort by
1008                // insertion order. Staggering by `idx` instead overlaps with
1009                // the timestamps of a later batch flushed within this batch's
1010                // fabricated span, interleaving the two batches.
1011                let ts = base_ts;
1012                stmt.execute(params![id, ts, msg]).into_diagnostic()?;
1013            }
1014        }
1015        tx.commit().into_diagnostic()?;
1016        Ok(())
1017    }
1018
1019    fn append_structured(&self, daemon_id: &DaemonId, parsed: &ParsedLog) -> Result<()> {
1020        let id = daemon_id.qualified();
1021
1022        let conn = self.conn.lock().unwrap();
1023        // Sample while holding the lock so timestamp order matches insertion
1024        // order under concurrent writers.
1025        let ts = Local::now().timestamp_millis();
1026        let _ = conn
1027            .execute(
1028                "INSERT INTO log_entries (daemon_id, timestamp, message, level, msg, logger, fields_json) \
1029                 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)",
1030                params![id, ts, parsed.message, parsed.level, parsed.msg, parsed.logger, parsed.fields_json],
1031            )
1032            .into_diagnostic()?;
1033        Ok(())
1034    }
1035
1036    fn append_structured_batch(&self, daemon_id: &DaemonId, entries: &[ParsedLog]) -> Result<()> {
1037        if entries.is_empty() {
1038            return Ok(());
1039        }
1040        let id = daemon_id.qualified();
1041
1042        let mut conn = self.conn.lock().unwrap();
1043        // IMMEDIATE acquires the write lock at BEGIN, and we sample only after
1044        // that, so timestamp order matches commit order even across processes
1045        // (sink subprocesses and the supervisor each hold their own
1046        // SqliteLogStore instance on the same database).
1047        let tx = conn
1048            .transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)
1049            .into_diagnostic()?;
1050        let base_ts = Local::now().timestamp_millis();
1051        {
1052            let mut stmt = tx
1053                .prepare(
1054                    "INSERT INTO log_entries (daemon_id, timestamp, message, level, msg, logger, fields_json) \
1055                     VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)",
1056                )
1057                .into_diagnostic()?;
1058            for entry in entries.iter() {
1059                // Same timestamp for the whole batch; queries order by
1060                // (timestamp, id) and `id` is AUTOINCREMENT, so insertion
1061                // order is preserved without the cross-batch overlap that
1062                // staggered timestamps (`base_ts + idx`) caused.
1063                let ts = base_ts;
1064                stmt.execute(params![
1065                    id,
1066                    ts,
1067                    entry.message,
1068                    entry.level,
1069                    entry.msg,
1070                    entry.logger,
1071                    entry.fields_json
1072                ])
1073                .into_diagnostic()?;
1074            }
1075        }
1076        tx.commit().into_diagnostic()?;
1077        Ok(())
1078    }
1079
1080    fn query(&self, opts: &LogQuery) -> Result<Vec<LogEntry>> {
1081        // Delegate to parallel path for large single-daemon queries.
1082        if Self::should_parallelize(opts)
1083            && self.path.as_os_str() != ":memory:"
1084            && let Ok(entries) = self.query_parallel(opts)
1085        {
1086            return Ok(entries);
1087        }
1088        // Fall back to single-threaded on parallel failure.
1089
1090        // Single-threaded path.
1091        let conn = self.conn.lock().unwrap();
1092        let (sql, query_params) = Self::build_query_sql(opts, None);
1093        Self::execute_built_query(&conn, &sql, &query_params)
1094    }
1095
1096    fn tail(&self, daemon_id: &DaemonId, after_id: Option<i64>) -> Result<Vec<LogEntry>> {
1097        self.query(&LogQuery {
1098            daemon_ids: vec![daemon_id.qualified()],
1099            from: None,
1100            to: None,
1101            limit: None,
1102            order_desc: false,
1103            after_id,
1104            before_id: None,
1105            message_filters: Vec::new(),
1106            field_filters: Vec::new(),
1107            include_structured: false,
1108        })
1109    }
1110
1111    fn clear(&self, daemon_ids: &[DaemonId]) -> Result<()> {
1112        let mut conn = self.conn.lock().unwrap();
1113        let tx = conn.transaction().into_diagnostic()?;
1114        for id in daemon_ids {
1115            tx.execute(
1116                "DELETE FROM log_entries WHERE daemon_id = ?1",
1117                params![id.qualified()],
1118            )
1119            .into_diagnostic()?;
1120            tx.execute(
1121                "INSERT INTO log_clear_generations (daemon_id, generation)
1122                 VALUES (?1, 1)
1123                 ON CONFLICT(daemon_id) DO UPDATE SET generation = generation + 1",
1124                params![id.qualified()],
1125            )
1126            .into_diagnostic()?;
1127        }
1128        tx.commit().into_diagnostic()?;
1129        Ok(())
1130    }
1131
1132    fn last_id(&self, daemon_id: &DaemonId) -> Result<Option<i64>> {
1133        let conn = self.conn.lock().unwrap();
1134        // MAX(id) returns NULL when no rows exist for the daemon; the
1135        // Option<i64> decode maps NULL to None automatically.
1136        let id: Option<i64> = conn
1137            .query_row(
1138                "SELECT MAX(id) FROM log_entries WHERE daemon_id = ?1",
1139                params![daemon_id.qualified()],
1140                |row| row.get(0),
1141            )
1142            .into_diagnostic()?;
1143        Ok(id)
1144    }
1145
1146    fn list_daemon_ids(&self) -> Result<Vec<String>> {
1147        let conn = self.conn.lock().unwrap();
1148        let mut stmt = conn
1149            .prepare("SELECT DISTINCT daemon_id FROM log_entries")
1150            .into_diagnostic()?;
1151        let ids = stmt
1152            .query_map([], |row| {
1153                let id: String = row.get(0)?;
1154                Ok(id)
1155            })
1156            .into_diagnostic()?
1157            .filter_map(|r| r.ok())
1158            .collect();
1159        Ok(ids)
1160    }
1161
1162    fn apply_retention(
1163        &self,
1164        policy: &super::RetentionPolicy,
1165        excluded_daemon_ids: &[DaemonId],
1166        archive_hook: Option<&ArchiveHook>,
1167    ) -> Result<u64> {
1168        let daemon_ids = self.list_daemon_ids()?;
1169        let excluded: HashSet<String> = excluded_daemon_ids.iter().map(|d| d.qualified()).collect();
1170        let mut total = 0u64;
1171        for id_str in daemon_ids {
1172            if excluded.contains(&id_str) {
1173                continue;
1174            }
1175            let id = DaemonId::parse(&id_str).unwrap_or_else(|_| {
1176                DaemonId::try_new("global", &id_str).unwrap_or_else(|_| DaemonId::pitchfork())
1177            });
1178            if let Some(dur) = policy.age {
1179                total += self.rotate_by_age(&id, dur, archive_hook)?;
1180            }
1181            if let Some(n) = policy.count {
1182                total += self.rotate_by_count(&id, n, archive_hook)?;
1183            }
1184        }
1185        Ok(total)
1186    }
1187
1188    fn apply_retention_for_daemon(
1189        &self,
1190        daemon_id: &DaemonId,
1191        policy: &super::RetentionPolicy,
1192        archive_hook: Option<&ArchiveHook>,
1193    ) -> Result<u64> {
1194        let mut total = 0u64;
1195        if let Some(dur) = policy.age {
1196            total += self.rotate_by_age(daemon_id, dur, archive_hook)?;
1197        }
1198        if let Some(n) = policy.count {
1199            total += self.rotate_by_count(daemon_id, n, archive_hook)?;
1200        }
1201        Ok(total)
1202    }
1203
1204    fn last_clear_generation(&self, daemon_id: &DaemonId) -> Result<Option<u64>> {
1205        let conn = self.conn.lock().unwrap();
1206        let mut stmt = conn
1207            .prepare("SELECT generation FROM log_clear_generations WHERE daemon_id = ?1")
1208            .into_diagnostic()?;
1209        let generation: Option<i64> = stmt
1210            .query_row(params![daemon_id.qualified()], |row| row.get(0))
1211            .optional()
1212            .into_diagnostic()?;
1213        generation
1214            .map(|generation| {
1215                u64::try_from(generation)
1216                    .map_err(|_| miette::miette!("log clear generation cannot be negative"))
1217            })
1218            .transpose()
1219    }
1220
1221    fn query_with_generation(
1222        &self,
1223        opts: &LogQuery,
1224        daemon_id: &DaemonId,
1225    ) -> Result<(Vec<LogEntry>, Option<u64>)> {
1226        let conn = self.conn.lock().unwrap();
1227        // Wrap both reads in a single transaction so a concurrent clear
1228        // cannot interleave between them. In WAL mode, BEGIN acquires a
1229        // consistent snapshot that both statements share.
1230        //
1231        // Intentionally call build_query_sql + execute_built_query directly
1232        // instead of self.query(): the latter may delegate to query_parallel
1233        // which opens separate connections outside this transaction, breaking
1234        // the snapshot guarantee.
1235        conn.execute_batch("BEGIN").into_diagnostic()?;
1236        let result = (|| {
1237            let (sql, query_params) = Self::build_query_sql(opts, None);
1238            let entries = Self::execute_built_query(&conn, &sql, &query_params)?;
1239            let generation: Option<i64> = conn
1240                .query_row(
1241                    "SELECT generation FROM log_clear_generations WHERE daemon_id = ?1",
1242                    params![daemon_id.qualified()],
1243                    |row| row.get(0),
1244                )
1245                .optional()
1246                .into_diagnostic()?;
1247            let generation = generation
1248                .map(|g| {
1249                    u64::try_from(g)
1250                        .map_err(|_| miette::miette!("log clear generation cannot be negative"))
1251                })
1252                .transpose()?;
1253            Ok((entries, generation))
1254        })();
1255        // Always rollback (read transaction just needs to end).
1256        let _ = conn.execute_batch("ROLLBACK");
1257        result
1258    }
1259}
1260
1261/// Global singleton log store.
1262use once_cell::sync::Lazy;
1263use std::sync::Arc;
1264
1265pub static LOG_STORE: Lazy<Arc<SqliteLogStore>> = Lazy::new(|| {
1266    let path = crate::env::PITCHFORK_LOGS_DIR.join("logs.db");
1267    let mut is_fallback = false;
1268    let store = Arc::new(SqliteLogStore::open(&path).unwrap_or_else(|e| {
1269        error!(
1270            "failed to open log store at {}: {e}. Falling back to in-memory store; logs will not persist across restarts.",
1271            path.display()
1272        );
1273        is_fallback = true;
1274        SqliteLogStore::open(":memory:").expect("in-memory SQLite should always open")
1275    }));
1276
1277    // Auto-migrate any legacy text log files into SQLite on first access.
1278    // This runs once per process startup and is idempotent.
1279    // Skip migration when using the in-memory fallback to prevent data loss.
1280    if !is_fallback {
1281        if let Err(e) = auto_migrate_legacy_logs(&store) {
1282            warn!("legacy log auto-migration failed: {e}");
1283        }
1284    } else {
1285        warn!(
1286            "skipping legacy log auto-migration because log store is in-memory (no durable destination)"
1287        );
1288    }
1289
1290    store
1291});
1292
1293/// Auto-migrate legacy text log files into the SQLite log store.
1294///
1295/// Scans the logs directory for directories matching the new-format layout
1296/// (`namespace--name/namespace--name.log`), attempts to parse the directory
1297/// name as a valid safe-path daemon ID, and imports the content into SQLite.
1298/// This is idempotent: re-running it on already-migrated data is a no-op
1299/// because the legacy text files are deleted after successful import.
1300fn auto_migrate_legacy_logs(store: &SqliteLogStore) -> Result<()> {
1301    let logs_dir = &*crate::env::PITCHFORK_LOGS_DIR;
1302    if !logs_dir.exists() {
1303        return Ok(());
1304    }
1305
1306    let Ok(entries) = std::fs::read_dir(logs_dir) else {
1307        return Ok(());
1308    };
1309
1310    let mut total_migrated = 0u64;
1311    let mut migrated_ids = Vec::new();
1312
1313    for entry in entries.flatten() {
1314        let path = entry.path();
1315        if !path.is_dir() {
1316            continue;
1317        }
1318        // Skip the supervisor's own log directory.
1319        let file_name = path
1320            .file_name()
1321            .map_or(String::new(), |n| n.to_string_lossy().to_string());
1322        if file_name == "pitchfork" {
1323            continue;
1324        }
1325
1326        // Only consider directories that look like new-format safe-paths
1327        if !file_name.contains("--") {
1328            continue;
1329        }
1330        let log_file = path.join(format!("{file_name}.log"));
1331        if !log_file.exists() {
1332            continue;
1333        }
1334
1335        let daemon_id = match DaemonId::from_safe_path(&file_name) {
1336            Ok(id) => id,
1337            Err(_) => continue,
1338        };
1339
1340        // Skip the supervisor's own daemon; its log directory may be present
1341        // under logs/ but should never be imported into the user-facing log store.
1342        if daemon_id == DaemonId::pitchfork() {
1343            continue;
1344        }
1345
1346        match store.migrate_daemon_text_logs(&daemon_id) {
1347            Ok(0) => {}
1348            Ok(n) => {
1349                total_migrated += n;
1350                migrated_ids.push(daemon_id.qualified());
1351            }
1352            Err(e) => {
1353                warn!(
1354                    "failed to migrate text logs for {}: {e}",
1355                    daemon_id.qualified()
1356                );
1357            }
1358        }
1359    }
1360
1361    if total_migrated > 0 {
1362        warn!(
1363            "auto-migrated {total_migrated} legacy log entries from {count} daemon(s): {ids}",
1364            count = migrated_ids.len(),
1365            ids = migrated_ids.join(", ")
1366        );
1367    }
1368
1369    Ok(())
1370}
1371
1372#[cfg(test)]
1373mod tests {
1374    use super::text_to_json_literal;
1375    use super::*;
1376    use crate::log_store::LogStore;
1377
1378    #[test]
1379    fn test_json_literal_booleans() {
1380        assert_eq!(text_to_json_literal("true"), "true");
1381        assert_eq!(text_to_json_literal("TRUE"), "true");
1382        assert_eq!(text_to_json_literal("false"), "false");
1383        assert_eq!(text_to_json_literal("False"), "false");
1384    }
1385
1386    #[test]
1387    fn test_json_literal_null() {
1388        assert_eq!(text_to_json_literal("null"), "null");
1389        assert_eq!(text_to_json_literal("NULL"), "null");
1390    }
1391
1392    #[test]
1393    fn test_json_literal_valid_numbers() {
1394        assert_eq!(text_to_json_literal("42"), "42");
1395        assert_eq!(text_to_json_literal("-1"), "-1");
1396        assert_eq!(text_to_json_literal("0"), "0");
1397        assert_eq!(text_to_json_literal("3.14"), "3.14");
1398        assert_eq!(text_to_json_literal("1e10"), "1e10");
1399        assert_eq!(text_to_json_literal("1.5e-3"), "1.5e-3");
1400    }
1401
1402    #[test]
1403    fn test_json_literal_rejects_invalid_json_numbers() {
1404        // f64::parse accepts these but they are not valid JSON numbers,
1405        // so they must be treated as strings (quoted) instead.
1406        assert_eq!(text_to_json_literal("+42"), r#""+42""#);
1407        assert_eq!(text_to_json_literal("1."), r#""1.""#);
1408        assert_eq!(text_to_json_literal(".5"), r#"".5""#);
1409        assert_eq!(text_to_json_literal("inf"), r#""inf""#);
1410        assert_eq!(text_to_json_literal("nan"), r#""nan""#);
1411        assert_eq!(text_to_json_literal("infinity"), r#""infinity""#);
1412    }
1413
1414    #[test]
1415    fn test_json_literal_strings() {
1416        assert_eq!(text_to_json_literal("hello"), r#""hello""#);
1417        assert_eq!(text_to_json_literal("req_1"), r#""req_1""#);
1418        // String with special chars gets JSON-escaped
1419        assert_eq!(text_to_json_literal(r#"a"b"#), r#""a\"b""#);
1420    }
1421
1422    /// A writer that finds the lock held must wait for it, not give up. WAL
1423    /// serializes writers, so without `busy_timeout` the second writer takes
1424    /// SQLITE_BUSY immediately and silently discards its batch.
1425    #[test]
1426    fn write_waits_for_a_contended_lock() {
1427        let dir = tempfile::tempdir().unwrap();
1428        let path = dir.path().join("logs.db");
1429        let store = SqliteLogStore::open(&path).unwrap();
1430
1431        // Hold the write lock from a second connection, then release it after a
1432        // delay that a non-waiting writer could never survive.
1433        let blocker = Connection::open(&path).unwrap();
1434        blocker.busy_timeout(BUSY_TIMEOUT).unwrap();
1435        blocker.execute_batch("BEGIN IMMEDIATE").unwrap();
1436        let releaser = std::thread::spawn(move || {
1437            std::thread::sleep(std::time::Duration::from_millis(300));
1438            blocker.execute_batch("COMMIT").unwrap();
1439        });
1440
1441        let id = DaemonId::try_new("test", "blocked").unwrap();
1442        let entries = vec![crate::log_parse::parse("held-lock-line", "text")];
1443        store
1444            .append_structured_batch(&id, &entries)
1445            .expect("a contended write must wait for the lock, not fail");
1446
1447        releaser.join().unwrap();
1448        let found = store
1449            .query(&LogQuery {
1450                daemon_ids: vec![id.qualified()],
1451                ..Default::default()
1452            })
1453            .unwrap()
1454            .len();
1455        assert_eq!(found, 1, "the batch written under contention was lost");
1456    }
1457
1458    /// Two processes writing the same store is the normal case, not an edge
1459    /// case: daemon output and retention pruning already collide, and a
1460    /// per-daemon writer multiplies the opportunities. WAL serializes writers,
1461    /// so without `busy_timeout` the loser of that race takes SQLITE_BUSY
1462    /// immediately and silently discards its batch.
1463    #[test]
1464    fn concurrent_writers_do_not_lose_batches() {
1465        let dir = tempfile::tempdir().unwrap();
1466        let path = dir.path().join("logs.db");
1467
1468        const WRITERS: usize = 4;
1469        const BATCHES: usize = 15;
1470        const PER_BATCH: usize = 10;
1471
1472        // Release every worker into open() at the same moment. Without this
1473        // the scheduler is free to run them one after another, so the first
1474        // would enable WAL and the rest would never contend for the
1475        // journal-mode switch — leaving the race this guards unexercised.
1476        let barrier = std::sync::Barrier::new(WRITERS);
1477
1478        std::thread::scope(|scope| {
1479            for writer in 0..WRITERS {
1480                let path = path.clone();
1481                let barrier = &barrier;
1482                scope.spawn(move || {
1483                    // A separate connection per writer, as separate processes
1484                    // would have.
1485                    barrier.wait();
1486                    let store = SqliteLogStore::open(&path).unwrap();
1487                    let id = DaemonId::try_new("test", format!("w{writer}")).unwrap();
1488                    for batch in 0..BATCHES {
1489                        let entries: Vec<ParsedLog> = (0..PER_BATCH)
1490                            .map(|i| {
1491                                crate::log_parse::parse(&format!("w{writer}-{batch}-{i}"), "text")
1492                            })
1493                            .collect();
1494                        store
1495                            .append_structured_batch(&id, &entries)
1496                            .expect("concurrent batch write must not fail");
1497                    }
1498                });
1499            }
1500        });
1501
1502        let store = SqliteLogStore::open(&path).unwrap();
1503        for writer in 0..WRITERS {
1504            let id = DaemonId::try_new("test", format!("w{writer}")).unwrap();
1505            let found = store
1506                .query(&LogQuery {
1507                    daemon_ids: vec![id.qualified()],
1508                    ..Default::default()
1509                })
1510                .unwrap()
1511                .len();
1512            assert_eq!(
1513                found,
1514                BATCHES * PER_BATCH,
1515                "writer {writer} lost entries: got {found}"
1516            );
1517        }
1518    }
1519
1520    /// Regression for the "burst of log lines comes back out of order" bug:
1521    /// rows within one batch must share a single timestamp instead of being
1522    /// staggered by index. Staggered timestamps (`base_ts + idx`) overlap with
1523    /// the timestamps of the next batch flushed moments later, interleaving the
1524    /// two batches under `ORDER BY timestamp, id`.
1525    #[test]
1526    fn batch_rows_share_single_timestamp() {
1527        let dir = tempfile::tempdir().unwrap();
1528        let path = dir.path().join("logs.db");
1529        let store = SqliteLogStore::open(&path).unwrap();
1530        let id = DaemonId::try_new("test", "burst").unwrap();
1531
1532        let entries: Vec<ParsedLog> = (0..50)
1533            .map(|i| crate::log_parse::parse(&format!("line-{i}"), "text"))
1534            .collect();
1535        store.append_structured_batch(&id, &entries).unwrap();
1536
1537        let conn = store.conn.lock().unwrap();
1538        let timestamps: Vec<i64> = {
1539            let mut stmt = conn
1540                .prepare("SELECT timestamp FROM log_entries ORDER BY id")
1541                .unwrap();
1542            let rows = stmt.query_map([], |row| row.get::<_, i64>(0)).unwrap();
1543            rows.map(|r| r.unwrap()).collect()
1544        };
1545        assert_eq!(timestamps.len(), 50);
1546        let distinct: std::collections::HashSet<i64> = timestamps.into_iter().collect();
1547        assert_eq!(
1548            distinct.len(),
1549            1,
1550            "all rows in a batch must share one timestamp so (timestamp, id) ordering matches insertion order"
1551        );
1552    }
1553
1554    /// End-to-end regression for the reported repro (`run = "exec seq 500"`):
1555    /// several batches flushed back-to-back must still read back in insertion
1556    /// order. With staggered per-row timestamps the later batch's fabricated
1557    /// span overlaps the earlier one and interleaves the output.
1558    #[test]
1559    fn burst_batches_preserve_insertion_order() {
1560        let dir = tempfile::tempdir().unwrap();
1561        let path = dir.path().join("logs.db");
1562        let store = SqliteLogStore::open(&path).unwrap();
1563        let id = DaemonId::try_new("test", "burst").unwrap();
1564
1565        // 3 batches of 300 rows, appended as fast as possible — the staggered
1566        // spans (300ms each) are far wider than the inter-batch gap, so the
1567        // old code always overlapped and interleaved these.
1568        const BATCHES: usize = 3;
1569        const PER_BATCH: usize = 300;
1570        for batch in 0..BATCHES {
1571            let entries: Vec<ParsedLog> = (0..PER_BATCH)
1572                .map(|i| crate::log_parse::parse(&format!("{}", batch * PER_BATCH + i), "text"))
1573                .collect();
1574            store.append_structured_batch(&id, &entries).unwrap();
1575        }
1576
1577        let entries = store
1578            .query(&LogQuery {
1579                daemon_ids: vec![id.qualified()],
1580                ..Default::default()
1581            })
1582            .unwrap();
1583        assert_eq!(entries.len(), BATCHES * PER_BATCH);
1584        for (n, entry) in entries.iter().enumerate() {
1585            assert_eq!(
1586                entry.message,
1587                n.to_string(),
1588                "log line {n} read back out of order (got '{}')",
1589                entry.message
1590            );
1591        }
1592    }
1593}