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
16fn text_to_json_literal(value: &str) -> String {
28 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 if value.eq_ignore_ascii_case("null") {
37 return "null".to_string();
38 }
39 if let Ok(serde_json::Value::Number(_)) = serde_json::from_str(value) {
43 return value.to_string();
44 }
45 serde_json::to_string(value).unwrap_or_else(|_| format!(r#""{value}""#))
47}
48
49fn add_regexp_function(conn: &Connection) -> Result<()> {
55 use std::cell::RefCell;
56
57 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
88pub struct SqliteLogStore {
90 conn: Mutex<Connection>,
91 path: PathBuf,
92}
93
94const PARALLEL_QUERY_THRESHOLD: usize = 200_000;
101
102const BUSY_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(5);
106
107fn 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 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 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 let existing_cols: Vec<String> = {
198 let mut stmt = conn
199 .prepare("PRAGMA table_info(log_entries)")
200 .into_diagnostic()?;
201 let rows = stmt
202 .query_map([], |row| row.get::<_, String>(1))
203 .into_diagnostic()?;
204 rows.filter_map(|r| r.ok()).collect()
205 };
206 for col in ["level", "msg", "logger", "fields_json"] {
207 if !existing_cols.iter().any(|c| c == col) {
208 conn.execute(
209 &format!("ALTER TABLE log_entries ADD COLUMN {col} TEXT"),
210 [],
211 )
212 .into_diagnostic()?;
213 }
214 }
215 conn.execute(
216 "CREATE INDEX IF NOT EXISTS idx_daemon_level_ts ON log_entries(daemon_id, level, timestamp);",
217 [],
218 )
219 .into_diagnostic()?;
220
221 conn.execute(
222 "CREATE TABLE IF NOT EXISTS log_clear_generations (
223 daemon_id TEXT PRIMARY KEY,
224 generation INTEGER NOT NULL DEFAULT 0
225 );",
226 [],
227 )
228 .into_diagnostic()?;
229 Ok(Self {
230 conn: Mutex::new(conn),
231 path,
232 })
233 }
234
235 fn row_to_entry(row: &rusqlite::Row) -> rusqlite::Result<LogEntry> {
236 let id: i64 = row.get(0)?;
237 let daemon_id: String = row.get(1)?;
238 let ts_millis: i64 = row.get(2)?;
239 let message: String = row.get(3)?;
240 let level: Option<String> = row.get(4)?;
241 let msg: Option<String> = row.get(5)?;
242 let logger: Option<String> = row.get(6)?;
243 let fields_json: Option<String> = row.get(7)?;
244 let timestamp = Local
245 .timestamp_millis_opt(ts_millis)
246 .single()
247 .unwrap_or_else(Local::now);
248 Ok(LogEntry {
249 id,
250 daemon_id,
251 timestamp,
252 message,
253 level,
254 msg,
255 logger,
256 fields_json,
257 })
258 }
259
260 fn archive_entries(
261 &self,
262 entries: &[LogEntry],
263 archive_hook: &ArchiveHook,
264 daemon_id: &DaemonId,
265 reason: &str,
266 ) -> Result<()> {
267 use crate::shell::HideConsoleWindow;
268 use std::process::{Command, Stdio};
269
270 if entries.is_empty() {
271 return Ok(());
272 }
273
274 let shell =
281 crate::settings::resolve_shell().map_err(|e| miette::miette!("archive hook: {e}"))?;
282 let (shell_program, shell_args) = shell.split_first().unwrap();
283
284 for chunk in entries.chunks(archive_hook.batch_size.max(1)) {
285 let mut child = Command::new(shell_program)
286 .args(shell_args)
287 .arg(&archive_hook.command)
288 .stdin(Stdio::piped())
289 .stdout(Stdio::null())
290 .stderr(Stdio::piped())
291 .env("PITCHFORK_DAEMON_ID", daemon_id.qualified())
292 .env("PITCHFORK_ARCHIVE_REASON", reason)
293 .hide_console_window()
294 .spawn()
295 .into_diagnostic()
296 .map_err(|e| miette::miette!("failed to spawn archive hook: {e}"))?;
297
298 let write_result = {
302 let stdin = child.stdin.take().expect("piped stdin should be available");
303 let mut stdin = std::io::BufWriter::new(stdin);
304 let mut result = Ok(());
305 for entry in chunk {
306 let line = serde_json::json!({
307 "id": entry.id,
308 "daemon_id": entry.daemon_id,
309 "timestamp": entry.timestamp.to_rfc3339(),
310 "message": entry.message,
311 });
312 if let Err(e) = writeln!(stdin, "{}", line) {
313 result = Err(miette::miette!(
314 "failed to write to archive hook stdin: {e}"
315 ));
316 break;
317 }
318 }
319 if result.is_ok()
327 && let Err(e) = stdin.flush()
328 {
329 result = Err(miette::miette!("failed to flush archive hook stdin: {e}"));
330 }
331 result
332 };
334
335 if let Err(e) = write_result {
336 let _ = child.kill();
337 let _ = child.wait();
338 return Err(e);
339 }
340
341 let output = child.wait_with_output().into_diagnostic()?;
342 if !output.status.success() {
343 let stderr = String::from_utf8_lossy(&output.stderr);
344 return Err(miette::miette!(
345 "archive hook failed with status {}: {stderr}",
346 output.status
347 ));
348 }
349 }
350
351 Ok(())
352 }
353
354 fn delete_by_ids(&self, ids: &[i64]) -> Result<u64> {
359 const SQLITE_MAX_VARS: usize = 999;
360
361 let mut total = 0u64;
362 let conn = self.conn.lock().unwrap();
363 for chunk in ids.chunks(SQLITE_MAX_VARS) {
364 if chunk.is_empty() {
365 continue;
366 }
367 let placeholders: Vec<String> = (1..=chunk.len()).map(|i| format!("?{i}")).collect();
368 let sql = format!(
369 "DELETE FROM log_entries WHERE id IN ({})",
370 placeholders.join(", ")
371 );
372 total += conn
373 .execute(&sql, rusqlite::params_from_iter(chunk.iter()))
374 .into_diagnostic()? as u64;
375 }
376 Ok(total)
377 }
378
379 pub fn rotate_by_age(
386 &self,
387 daemon_id: &DaemonId,
388 max_age: chrono::Duration,
389 archive_hook: Option<&ArchiveHook>,
390 ) -> Result<u64> {
391 let cutoff = (Local::now() - max_age).timestamp_millis();
392 let hook = archive_hook.filter(|h| h.is_enabled());
393
394 if let Some(hook) = hook {
395 let mut total_deleted = 0u64;
396 loop {
397 let entries: Vec<LogEntry> = {
399 let conn = self.conn.lock().unwrap();
400 let mut stmt = conn
401 .prepare(
402 "SELECT id, daemon_id, timestamp, message, level, msg, logger, fields_json FROM log_entries
403 WHERE daemon_id = ?1 AND timestamp < ?2
404 ORDER BY timestamp ASC, id ASC
405 LIMIT ?3",
406 )
407 .into_diagnostic()?;
408 stmt.query_map(
409 params![daemon_id.qualified(), cutoff, hook.batch_size as i64],
410 Self::row_to_entry,
411 )
412 .into_diagnostic()?
413 .collect::<rusqlite::Result<Vec<_>>>()
414 .into_diagnostic()?
415 };
416
417 if entries.is_empty() {
418 break;
419 }
420
421 let batch_ids: Vec<i64> = entries.iter().map(|e| e.id).collect();
422
423 self.archive_entries(&entries, hook, daemon_id, "age")?;
425
426 let deleted = self.delete_by_ids(&batch_ids)?;
428 total_deleted += deleted;
429 }
430 Ok(total_deleted)
431 } else {
432 let conn = self.conn.lock().unwrap();
433 let rows = conn
434 .execute(
435 "DELETE FROM log_entries WHERE daemon_id = ?1 AND timestamp < ?2",
436 params![daemon_id.qualified(), cutoff],
437 )
438 .into_diagnostic()?;
439 Ok(rows as u64)
440 }
441 }
442
443 pub fn rotate_by_count(
451 &self,
452 daemon_id: &DaemonId,
453 max_count: u64,
454 archive_hook: Option<&ArchiveHook>,
455 ) -> Result<u64> {
456 let hook = archive_hook.filter(|h| h.is_enabled());
457
458 let to_delete: i64 = {
460 let conn = self.conn.lock().unwrap();
461 let count: i64 = conn
462 .query_row(
463 "SELECT COUNT(*) FROM log_entries WHERE daemon_id = ?1",
464 [daemon_id.qualified()],
465 |row| row.get(0),
466 )
467 .into_diagnostic()?;
468 count.saturating_sub(max_count as i64)
469 };
470
471 if to_delete <= 0 {
472 return Ok(0);
473 }
474
475 if let Some(hook) = hook {
476 let mut total_deleted = 0u64;
477 let mut remaining = to_delete;
478 loop {
479 let batch_len = remaining.min(hook.batch_size as i64);
480
481 let entries: Vec<LogEntry> = {
483 let conn = self.conn.lock().unwrap();
484 let mut stmt = conn
485 .prepare(
486 "SELECT id, daemon_id, timestamp, message, level, msg, logger, fields_json FROM log_entries
487 WHERE daemon_id = ?1
488 ORDER BY timestamp ASC, id ASC
489 LIMIT ?2",
490 )
491 .into_diagnostic()?;
492 stmt.query_map(
493 params![daemon_id.qualified(), batch_len],
494 Self::row_to_entry,
495 )
496 .into_diagnostic()?
497 .collect::<rusqlite::Result<Vec<_>>>()
498 .into_diagnostic()?
499 };
500
501 if entries.is_empty() {
502 break;
503 }
504
505 let fetched = entries.len() as i64;
506 let batch_ids: Vec<i64> = entries.iter().map(|e| e.id).collect();
507
508 self.archive_entries(&entries, hook, daemon_id, "count")?;
510
511 let deleted = self.delete_by_ids(&batch_ids)?;
513 total_deleted += deleted;
514 remaining -= fetched;
515 }
516 Ok(total_deleted)
517 } else {
518 let conn = self.conn.lock().unwrap();
519 let rows = conn
520 .execute(
521 "DELETE FROM log_entries WHERE id IN (
522 SELECT id FROM log_entries WHERE daemon_id = ?1
523 ORDER BY timestamp ASC, id ASC LIMIT ?2
524 )",
525 params![daemon_id.qualified(), to_delete],
526 )
527 .into_diagnostic()?;
528 Ok(rows as u64)
529 }
530 }
531
532 pub fn migrate_daemon_text_logs(&self, daemon_id: &DaemonId) -> Result<u64> {
537 let text_path = daemon_id.log_path();
538 if !text_path.exists() {
539 return Ok(0);
540 }
541
542 let file = std::fs::File::open(&text_path).into_diagnostic()?;
543 let reader = BufReader::new(file);
544 let re = regex::Regex::new(r"^(\d{4}-\d{2}-\d{2} \d{2}:\d{2}:\d{2}) ([\w./-]+) (.*)$")
545 .expect("invalid regex");
546
547 let mut current_timestamp: Option<DateTime<Local>> = None;
548 let mut current_message = String::new();
549 let mut entries = Vec::with_capacity(1000);
550 let mut total_migrated: u64 = 0;
551
552 for line in reader.lines() {
553 let line = line.into_diagnostic()?;
554 if let Some(caps) = re.captures(&line) {
555 if let Some(ts) = current_timestamp.take() {
556 entries.push((ts, std::mem::take(&mut current_message)));
557 }
558 let ts_str = caps.get(1).map(|m| m.as_str()).unwrap_or_default();
559 let msg = caps.get(3).map(|m| m.as_str()).unwrap_or_default();
560 if let Ok(naive) =
561 chrono::NaiveDateTime::parse_from_str(ts_str, "%Y-%m-%d %H:%M:%S")
562 {
563 current_timestamp = Local.from_local_datetime(&naive).single();
564 current_message = msg.to_string();
565 }
566 } else if current_timestamp.is_some() {
567 current_message.push('\n');
568 current_message.push_str(&line);
569 }
570
571 if entries.len() >= 1000 {
572 total_migrated += self.insert_batch(daemon_id, &entries)?;
573 entries.clear();
574 }
575 }
576
577 if let Some(ts) = current_timestamp {
578 entries.push((ts, std::mem::take(&mut current_message)));
579 }
580
581 if !entries.is_empty() {
582 total_migrated += self.insert_batch(daemon_id, &entries)?;
583 }
584
585 if total_migrated > 0
586 && let Err(e) = std::fs::remove_file(&text_path)
587 {
588 log::warn!(
589 "failed to remove legacy log file after migration {}: {e}",
590 text_path.display()
591 );
592 }
593
594 Ok(total_migrated)
595 }
596
597 fn insert_batch(
598 &self,
599 daemon_id: &DaemonId,
600 entries: &[(DateTime<Local>, String)],
601 ) -> Result<u64> {
602 let mut conn = self.conn.lock().unwrap();
603 let tx = conn.transaction().into_diagnostic()?;
604 let mut count = 0u64;
605 {
606 let mut stmt = tx
607 .prepare(
608 "INSERT INTO log_entries (daemon_id, timestamp, message) VALUES (?1, ?2, ?3)",
609 )
610 .into_diagnostic()?;
611 for (ts, msg) in entries {
612 stmt.execute(params![daemon_id.qualified(), ts.timestamp_millis(), msg])
613 .into_diagnostic()?;
614 count += 1;
615 }
616 }
617 tx.commit().into_diagnostic()?;
618 Ok(count)
619 }
620
621 fn build_query_sql(
626 opts: &LogQuery,
627 id_range: Option<(i64, i64)>,
628 ) -> (String, Vec<Box<dyn rusqlite::ToSql>>) {
629 let mut conditions = Vec::new();
630 let mut query_params: Vec<Box<dyn rusqlite::ToSql>> = Vec::new();
631
632 if !opts.daemon_ids.is_empty() {
633 let placeholders: Vec<String> = (1..=opts.daemon_ids.len())
634 .map(|i| format!("?{}", i))
635 .collect();
636 conditions.push(format!("daemon_id IN ({})", placeholders.join(", ")));
637 for id in &opts.daemon_ids {
638 query_params.push(Box::new(id.clone()));
639 }
640 }
641
642 if let Some(from) = opts.from {
643 conditions.push(format!("timestamp >= ?{}", query_params.len() + 1));
644 query_params.push(Box::new(from.timestamp_millis()));
645 }
646
647 if let Some(to) = opts.to {
648 conditions.push(format!("timestamp <= ?{}", query_params.len() + 1));
649 query_params.push(Box::new(to.timestamp_millis()));
650 }
651
652 if let Some(after_id) = opts.after_id {
653 conditions.push(format!("id > ?{}", query_params.len() + 1));
654 query_params.push(Box::new(after_id));
655 }
656
657 if let Some(before_id) = opts.before_id {
658 conditions.push(format!("id < ?{}", query_params.len() + 1));
659 query_params.push(Box::new(before_id));
660 }
661
662 if let Some((start, end)) = id_range {
663 conditions.push(format!("id > ?{}", query_params.len() + 1));
664 query_params.push(Box::new(start));
665 conditions.push(format!("id <= ?{}", query_params.len() + 1));
666 query_params.push(Box::new(end));
667 }
668
669 let mut message_conditions = Vec::new();
670 for filter in &opts.message_filters {
671 match filter {
672 MessageFilter::Contains {
673 pattern,
674 case_sensitive,
675 } => {
676 let param_index = query_params.len() + 1;
677 if *case_sensitive {
678 message_conditions
679 .push(format!("INSTR(message, ?{idx}) > 0", idx = param_index));
680 query_params.push(Box::new(pattern.clone()));
681 } else {
682 let escaped = escape_like_pattern(pattern);
683 let param = format!("%{}%", escaped);
684 message_conditions.push(format!(
685 "LOWER(message) LIKE LOWER(?{idx}) ESCAPE '\\'",
686 idx = param_index
687 ));
688 query_params.push(Box::new(param));
689 }
690 }
691 MessageFilter::Regex { pattern } => {
692 let param_index = query_params.len() + 1;
693 message_conditions.push(format!("message REGEXP ?{param_index}"));
694 query_params.push(Box::new(pattern.clone()));
695 }
696 }
697 }
698 if !message_conditions.is_empty() {
699 conditions.push(format!("({})", message_conditions.join(" OR ")));
700 }
701
702 for filter in &opts.field_filters {
703 match filter {
704 FieldFilter::LevelMin(level) => {
705 let matching = crate::log_store::levels_at_or_above(level);
706 if matching.is_empty() {
707 conditions.push("0".to_string());
709 } else {
710 let placeholders = matching
711 .iter()
712 .map(|l| {
713 let idx = query_params.len() + 1;
714 query_params.push(Box::new((*l).to_string()));
715 format!("?{idx}")
716 })
717 .collect::<Vec<_>>()
718 .join(", ");
719 conditions.push(format!("level IN ({placeholders})"));
720 }
721 }
722 FieldFilter::FieldEq { key, value } => {
723 let key_idx = query_params.len() + 1;
724 let val_idx = query_params.len() + 2;
725 let json_literal = text_to_json_literal(value);
743 let raw_idx = query_params.len() + 3;
744 conditions.push(format!(
745 "EXISTS (SELECT 1 FROM json_each(fields_json) \
746 WHERE json_each.key = ?{key_idx} \
747 AND (json_each.value IS json_extract(?{val_idx}, '$') \
748 OR json_each.value IS ?{raw_idx}))"
749 ));
750 query_params.push(Box::new(key.clone()));
751 query_params.push(Box::new(json_literal));
752 query_params.push(Box::new(value.clone()));
753 }
754 FieldFilter::LoggerContains(pattern) => {
755 let param_index = query_params.len() + 1;
756 conditions.push(format!("logger LIKE ?{param_index} ESCAPE '\\'"));
757 let escaped = crate::log_store::escape_like_pattern(pattern);
758 query_params.push(Box::new(format!("%{escaped}%")));
759 }
760 }
761 }
762
763 let where_clause = if conditions.is_empty() {
764 String::new()
765 } else {
766 format!("WHERE {}", conditions.join(" AND "))
767 };
768
769 let order = if opts.order_desc { "DESC" } else { "ASC" };
770
771 let limit_clause = opts
772 .limit
773 .map(|n| format!("LIMIT {}", n))
774 .unwrap_or_default();
775
776 let columns = if opts.include_structured {
777 "id, daemon_id, timestamp, message, level, msg, logger, fields_json"
778 } else {
779 "id, daemon_id, timestamp, message, NULL, NULL, NULL, NULL"
780 };
781
782 let sql = format!(
783 "SELECT {columns} FROM log_entries {} ORDER BY timestamp {}, id {} {}",
784 where_clause, order, order, limit_clause
785 );
786
787 (sql, query_params)
788 }
789
790 fn should_parallelize(opts: &LogQuery) -> bool {
792 if opts.daemon_ids.len() != 1 {
793 return false;
794 }
795 if opts.after_id.is_some() {
798 return false;
799 }
800 let limit = opts.limit.unwrap_or(usize::MAX);
801 if limit < PARALLEL_QUERY_THRESHOLD {
802 return false;
803 }
804 std::thread::available_parallelism()
805 .map(|n| n.get() >= 2)
806 .unwrap_or(false)
807 }
808
809 fn query_parallel(&self, opts: &LogQuery) -> Result<Vec<LogEntry>> {
816 let max_threads = 2;
821 let num_threads = std::thread::available_parallelism()
822 .map(|n| n.get().min(max_threads))
823 .unwrap_or(1);
824
825 let max_id: Option<i64> = {
826 let conn = self.conn.lock().unwrap();
827 conn.query_row(
828 "SELECT MAX(id) FROM log_entries WHERE daemon_id = ?1",
829 params![&opts.daemon_ids[0]],
830 |row| row.get(0),
831 )
832 .ok()
833 };
834
835 let Some(max_id) = max_id else {
836 return Ok(Vec::new());
837 };
838 if max_id == 0 {
839 return Ok(Vec::new());
840 }
841
842 let shard_size = (max_id as usize).div_ceil(num_threads);
843 let path = self.path.clone();
844 let opts = opts.clone();
845 let needs_regexp = opts
846 .message_filters
847 .iter()
848 .any(|f| matches!(f, MessageFilter::Regex { .. }));
849
850 let shards: Vec<Result<Vec<LogEntry>>> = std::thread::scope(|s| {
851 (0..num_threads)
852 .map(|i| {
853 let start = (i * shard_size) as i64;
854 let end = if i == num_threads - 1 {
855 max_id
856 } else {
857 ((i + 1) * shard_size) as i64
858 };
859 let opts = opts.clone();
860 let path = path.clone();
861 s.spawn(move || -> Result<Vec<LogEntry>> {
862 let conn = Connection::open_with_flags(
863 &path,
864 OpenFlags::SQLITE_OPEN_READ_ONLY | OpenFlags::SQLITE_OPEN_NO_MUTEX,
865 )
866 .into_diagnostic()?;
867 conn.execute_batch(
868 "PRAGMA mmap_size = 268435456;
869 PRAGMA query_only = 1;",
870 )
871 .into_diagnostic()?;
872 if needs_regexp {
873 add_regexp_function(&conn)?;
874 }
875
876 let (sql, query_params) = Self::build_query_sql(&opts, Some((start, end)));
877 Self::execute_built_query(&conn, &sql, &query_params)
878 })
879 })
880 .map(|h| h.join().unwrap())
881 .collect()
882 });
883
884 let mut merged = Vec::new();
885 if opts.order_desc {
886 for shard in shards.into_iter().rev() {
887 merged.extend(shard?);
888 }
889 } else {
890 for shard in shards {
891 merged.extend(shard?);
892 }
893 }
894
895 if let Some(limit) = opts.limit
896 && merged.len() > limit
897 {
898 merged.truncate(limit);
899 }
900
901 Ok(merged)
902 }
903
904 fn execute_built_query(
906 conn: &Connection,
907 sql: &str,
908 query_params: &[Box<dyn rusqlite::ToSql>],
909 ) -> Result<Vec<LogEntry>> {
910 let mut stmt = conn.prepare(sql).into_diagnostic()?;
911 let params_ref: Vec<&dyn rusqlite::ToSql> =
912 query_params.iter().map(|p| p.as_ref()).collect();
913 let rows = stmt
914 .query_map(params_ref.as_slice(), Self::row_to_entry)
915 .into_diagnostic()?;
916 let mut entries = Vec::new();
917 for row in rows {
918 entries.push(row.into_diagnostic()?);
919 }
920 Ok(entries)
921 }
922
923 pub fn distinct_loggers(&self, daemon_id: &str) -> Result<Vec<String>> {
928 let conn = self.conn.lock().unwrap();
929 let mut stmt = conn
930 .prepare(
931 "SELECT DISTINCT logger FROM ( \
932 SELECT logger FROM log_entries \
933 WHERE daemon_id = ?1 AND logger IS NOT NULL \
934 ORDER BY id DESC LIMIT 5000 \
935 ) ORDER BY logger",
936 )
937 .into_diagnostic()?;
938 let rows = stmt
939 .query_map(params![daemon_id], |row| row.get::<_, String>(0))
940 .into_diagnostic()?;
941 let mut loggers = Vec::new();
942 for row in rows {
943 loggers.push(row.into_diagnostic()?);
944 }
945 Ok(loggers)
946 }
947
948 pub fn distinct_field_keys(&self, daemon_id: &str) -> Result<Vec<String>> {
954 let conn = self.conn.lock().unwrap();
955 let mut stmt = conn
956 .prepare(
957 "SELECT DISTINCT je.key \
958 FROM ( \
959 SELECT fields_json FROM log_entries \
960 WHERE daemon_id = ?1 AND fields_json IS NOT NULL \
961 ORDER BY id DESC LIMIT 5000 \
962 ), json_each(fields_json) AS je \
963 ORDER BY je.key",
964 )
965 .into_diagnostic()?;
966 let rows = stmt
967 .query_map(params![daemon_id], |row| row.get::<_, String>(0))
968 .into_diagnostic()?;
969 let mut keys = Vec::new();
970 for row in rows {
971 keys.push(row.into_diagnostic()?);
972 }
973 Ok(keys)
974 }
975}
976
977impl LogStore for SqliteLogStore {
978 fn append(&self, daemon_id: &DaemonId, message: &str) -> Result<()> {
979 let id = daemon_id.qualified();
980 let msg = message.to_string();
981
982 let conn = self.conn.lock().unwrap();
983 let ts = Local::now().timestamp_millis();
986 let _ = conn
987 .execute(
988 "INSERT INTO log_entries (daemon_id, timestamp, message) VALUES (?1, ?2, ?3)",
989 params![id, ts, msg],
990 )
991 .into_diagnostic()?;
992 Ok(())
993 }
994
995 fn append_batch(&self, daemon_id: &DaemonId, messages: &[String]) -> Result<()> {
996 if messages.is_empty() {
997 return Ok(());
998 }
999 let id = daemon_id.qualified();
1000
1001 let mut conn = self.conn.lock().unwrap();
1002 let tx = conn
1007 .transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)
1008 .into_diagnostic()?;
1009 let base_ts = Local::now().timestamp_millis();
1010 {
1011 let mut stmt = tx
1012 .prepare(
1013 "INSERT INTO log_entries (daemon_id, timestamp, message) VALUES (?1, ?2, ?3)",
1014 )
1015 .into_diagnostic()?;
1016 for msg in messages.iter() {
1017 let ts = base_ts;
1024 stmt.execute(params![id, ts, msg]).into_diagnostic()?;
1025 }
1026 }
1027 tx.commit().into_diagnostic()?;
1028 Ok(())
1029 }
1030
1031 fn append_structured(&self, daemon_id: &DaemonId, parsed: &ParsedLog) -> Result<()> {
1032 let id = daemon_id.qualified();
1033
1034 let conn = self.conn.lock().unwrap();
1035 let ts = Local::now().timestamp_millis();
1038 let _ = conn
1039 .execute(
1040 "INSERT INTO log_entries (daemon_id, timestamp, message, level, msg, logger, fields_json) \
1041 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)",
1042 params![id, ts, parsed.message, parsed.level, parsed.msg, parsed.logger, parsed.fields_json],
1043 )
1044 .into_diagnostic()?;
1045 Ok(())
1046 }
1047
1048 fn append_structured_batch(&self, daemon_id: &DaemonId, entries: &[ParsedLog]) -> Result<()> {
1049 if entries.is_empty() {
1050 return Ok(());
1051 }
1052 let id = daemon_id.qualified();
1053
1054 let mut conn = self.conn.lock().unwrap();
1055 let tx = conn
1060 .transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)
1061 .into_diagnostic()?;
1062 let base_ts = Local::now().timestamp_millis();
1063 {
1064 let mut stmt = tx
1065 .prepare(
1066 "INSERT INTO log_entries (daemon_id, timestamp, message, level, msg, logger, fields_json) \
1067 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)",
1068 )
1069 .into_diagnostic()?;
1070 for entry in entries.iter() {
1071 let ts = base_ts;
1076 stmt.execute(params![
1077 id,
1078 ts,
1079 entry.message,
1080 entry.level,
1081 entry.msg,
1082 entry.logger,
1083 entry.fields_json
1084 ])
1085 .into_diagnostic()?;
1086 }
1087 }
1088 tx.commit().into_diagnostic()?;
1089 Ok(())
1090 }
1091
1092 fn query(&self, opts: &LogQuery) -> Result<Vec<LogEntry>> {
1093 if Self::should_parallelize(opts)
1095 && self.path.as_os_str() != ":memory:"
1096 && let Ok(entries) = self.query_parallel(opts)
1097 {
1098 return Ok(entries);
1099 }
1100 let conn = self.conn.lock().unwrap();
1104 let (sql, query_params) = Self::build_query_sql(opts, None);
1105 Self::execute_built_query(&conn, &sql, &query_params)
1106 }
1107
1108 fn tail(&self, daemon_id: &DaemonId, after_id: Option<i64>) -> Result<Vec<LogEntry>> {
1109 self.query(&LogQuery {
1110 daemon_ids: vec![daemon_id.qualified()],
1111 from: None,
1112 to: None,
1113 limit: None,
1114 order_desc: false,
1115 after_id,
1116 before_id: None,
1117 message_filters: Vec::new(),
1118 field_filters: Vec::new(),
1119 include_structured: false,
1120 })
1121 }
1122
1123 fn clear(&self, daemon_ids: &[DaemonId]) -> Result<()> {
1124 let mut conn = self.conn.lock().unwrap();
1125 let tx = conn.transaction().into_diagnostic()?;
1126 for id in daemon_ids {
1127 tx.execute(
1128 "DELETE FROM log_entries WHERE daemon_id = ?1",
1129 params![id.qualified()],
1130 )
1131 .into_diagnostic()?;
1132 tx.execute(
1133 "INSERT INTO log_clear_generations (daemon_id, generation)
1134 VALUES (?1, 1)
1135 ON CONFLICT(daemon_id) DO UPDATE SET generation = generation + 1",
1136 params![id.qualified()],
1137 )
1138 .into_diagnostic()?;
1139 }
1140 tx.commit().into_diagnostic()?;
1141 Ok(())
1142 }
1143
1144 fn last_id(&self, daemon_id: &DaemonId) -> Result<Option<i64>> {
1145 let conn = self.conn.lock().unwrap();
1146 let id: Option<i64> = conn
1149 .query_row(
1150 "SELECT MAX(id) FROM log_entries WHERE daemon_id = ?1",
1151 params![daemon_id.qualified()],
1152 |row| row.get(0),
1153 )
1154 .into_diagnostic()?;
1155 Ok(id)
1156 }
1157
1158 fn list_daemon_ids(&self) -> Result<Vec<String>> {
1159 let conn = self.conn.lock().unwrap();
1160 let mut stmt = conn
1161 .prepare("SELECT DISTINCT daemon_id FROM log_entries")
1162 .into_diagnostic()?;
1163 let ids = stmt
1164 .query_map([], |row| {
1165 let id: String = row.get(0)?;
1166 Ok(id)
1167 })
1168 .into_diagnostic()?
1169 .filter_map(|r| r.ok())
1170 .collect();
1171 Ok(ids)
1172 }
1173
1174 fn apply_retention(
1175 &self,
1176 policy: &super::RetentionPolicy,
1177 excluded_daemon_ids: &[DaemonId],
1178 archive_hook: Option<&ArchiveHook>,
1179 ) -> Result<u64> {
1180 let daemon_ids = self.list_daemon_ids()?;
1181 let excluded: HashSet<String> = excluded_daemon_ids.iter().map(|d| d.qualified()).collect();
1182 let mut total = 0u64;
1183 for id_str in daemon_ids {
1184 if excluded.contains(&id_str) {
1185 continue;
1186 }
1187 let id = DaemonId::parse(&id_str).unwrap_or_else(|_| {
1188 DaemonId::try_new("global", &id_str).unwrap_or_else(|_| DaemonId::pitchfork())
1189 });
1190 if let Some(dur) = policy.age {
1191 total += self.rotate_by_age(&id, dur, archive_hook)?;
1192 }
1193 if let Some(n) = policy.count {
1194 total += self.rotate_by_count(&id, n, archive_hook)?;
1195 }
1196 }
1197 Ok(total)
1198 }
1199
1200 fn apply_retention_for_daemon(
1201 &self,
1202 daemon_id: &DaemonId,
1203 policy: &super::RetentionPolicy,
1204 archive_hook: Option<&ArchiveHook>,
1205 ) -> Result<u64> {
1206 let mut total = 0u64;
1207 if let Some(dur) = policy.age {
1208 total += self.rotate_by_age(daemon_id, dur, archive_hook)?;
1209 }
1210 if let Some(n) = policy.count {
1211 total += self.rotate_by_count(daemon_id, n, archive_hook)?;
1212 }
1213 Ok(total)
1214 }
1215
1216 fn last_clear_generation(&self, daemon_id: &DaemonId) -> Result<Option<u64>> {
1217 let conn = self.conn.lock().unwrap();
1218 let mut stmt = conn
1219 .prepare("SELECT generation FROM log_clear_generations WHERE daemon_id = ?1")
1220 .into_diagnostic()?;
1221 let generation: Option<i64> = stmt
1222 .query_row(params![daemon_id.qualified()], |row| row.get(0))
1223 .optional()
1224 .into_diagnostic()?;
1225 generation
1226 .map(|generation| {
1227 u64::try_from(generation)
1228 .map_err(|_| miette::miette!("log clear generation cannot be negative"))
1229 })
1230 .transpose()
1231 }
1232
1233 fn query_with_generation(
1234 &self,
1235 opts: &LogQuery,
1236 daemon_id: &DaemonId,
1237 ) -> Result<(Vec<LogEntry>, Option<u64>)> {
1238 let conn = self.conn.lock().unwrap();
1239 conn.execute_batch("BEGIN").into_diagnostic()?;
1248 let result = (|| {
1249 let (sql, query_params) = Self::build_query_sql(opts, None);
1250 let entries = Self::execute_built_query(&conn, &sql, &query_params)?;
1251 let generation: Option<i64> = conn
1252 .query_row(
1253 "SELECT generation FROM log_clear_generations WHERE daemon_id = ?1",
1254 params![daemon_id.qualified()],
1255 |row| row.get(0),
1256 )
1257 .optional()
1258 .into_diagnostic()?;
1259 let generation = generation
1260 .map(|g| {
1261 u64::try_from(g)
1262 .map_err(|_| miette::miette!("log clear generation cannot be negative"))
1263 })
1264 .transpose()?;
1265 Ok((entries, generation))
1266 })();
1267 let _ = conn.execute_batch("ROLLBACK");
1269 result
1270 }
1271}
1272
1273use once_cell::sync::Lazy;
1275use std::sync::Arc;
1276
1277pub static LOG_STORE: Lazy<Arc<SqliteLogStore>> = Lazy::new(|| {
1278 let path = crate::env::PITCHFORK_LOGS_DIR.join("logs.db");
1279 let mut is_fallback = false;
1280 let store = Arc::new(SqliteLogStore::open(&path).unwrap_or_else(|e| {
1281 error!(
1282 "failed to open log store at {}: {e}. Falling back to in-memory store; logs will not persist across restarts.",
1283 path.display()
1284 );
1285 is_fallback = true;
1286 SqliteLogStore::open(":memory:").expect("in-memory SQLite should always open")
1287 }));
1288
1289 if !is_fallback {
1293 if let Err(e) = auto_migrate_legacy_logs(&store) {
1294 warn!("legacy log auto-migration failed: {e}");
1295 }
1296 } else {
1297 warn!(
1298 "skipping legacy log auto-migration because log store is in-memory (no durable destination)"
1299 );
1300 }
1301
1302 store
1303});
1304
1305fn auto_migrate_legacy_logs(store: &SqliteLogStore) -> Result<()> {
1313 let logs_dir = &*crate::env::PITCHFORK_LOGS_DIR;
1314 if !logs_dir.exists() {
1315 return Ok(());
1316 }
1317
1318 let Ok(entries) = std::fs::read_dir(logs_dir) else {
1319 return Ok(());
1320 };
1321
1322 let mut total_migrated = 0u64;
1323 let mut migrated_ids = Vec::new();
1324
1325 for entry in entries.flatten() {
1326 let path = entry.path();
1327 if !path.is_dir() {
1328 continue;
1329 }
1330 let file_name = path
1332 .file_name()
1333 .map_or(String::new(), |n| n.to_string_lossy().to_string());
1334 if file_name == "pitchfork" {
1335 continue;
1336 }
1337
1338 if !file_name.contains("--") {
1340 continue;
1341 }
1342 let log_file = path.join(format!("{file_name}.log"));
1343 if !log_file.exists() {
1344 continue;
1345 }
1346
1347 let daemon_id = match DaemonId::from_safe_path(&file_name) {
1348 Ok(id) => id,
1349 Err(_) => continue,
1350 };
1351
1352 if daemon_id == DaemonId::pitchfork() {
1355 continue;
1356 }
1357
1358 match store.migrate_daemon_text_logs(&daemon_id) {
1359 Ok(0) => {}
1360 Ok(n) => {
1361 total_migrated += n;
1362 migrated_ids.push(daemon_id.qualified());
1363 }
1364 Err(e) => {
1365 warn!(
1366 "failed to migrate text logs for {}: {e}",
1367 daemon_id.qualified()
1368 );
1369 }
1370 }
1371 }
1372
1373 if total_migrated > 0 {
1374 warn!(
1375 "auto-migrated {total_migrated} legacy log entries from {count} daemon(s): {ids}",
1376 count = migrated_ids.len(),
1377 ids = migrated_ids.join(", ")
1378 );
1379 }
1380
1381 Ok(())
1382}
1383
1384#[cfg(test)]
1385mod tests {
1386 use super::text_to_json_literal;
1387 use super::*;
1388 use crate::log_store::LogStore;
1389
1390 #[test]
1391 fn test_json_literal_booleans() {
1392 assert_eq!(text_to_json_literal("true"), "true");
1393 assert_eq!(text_to_json_literal("TRUE"), "true");
1394 assert_eq!(text_to_json_literal("false"), "false");
1395 assert_eq!(text_to_json_literal("False"), "false");
1396 }
1397
1398 #[test]
1399 fn test_json_literal_null() {
1400 assert_eq!(text_to_json_literal("null"), "null");
1401 assert_eq!(text_to_json_literal("NULL"), "null");
1402 }
1403
1404 #[test]
1405 fn test_json_literal_valid_numbers() {
1406 assert_eq!(text_to_json_literal("42"), "42");
1407 assert_eq!(text_to_json_literal("-1"), "-1");
1408 assert_eq!(text_to_json_literal("0"), "0");
1409 assert_eq!(text_to_json_literal("3.14"), "3.14");
1410 assert_eq!(text_to_json_literal("1e10"), "1e10");
1411 assert_eq!(text_to_json_literal("1.5e-3"), "1.5e-3");
1412 }
1413
1414 #[test]
1415 fn test_json_literal_rejects_invalid_json_numbers() {
1416 assert_eq!(text_to_json_literal("+42"), r#""+42""#);
1419 assert_eq!(text_to_json_literal("1."), r#""1.""#);
1420 assert_eq!(text_to_json_literal(".5"), r#"".5""#);
1421 assert_eq!(text_to_json_literal("inf"), r#""inf""#);
1422 assert_eq!(text_to_json_literal("nan"), r#""nan""#);
1423 assert_eq!(text_to_json_literal("infinity"), r#""infinity""#);
1424 }
1425
1426 #[test]
1427 fn test_json_literal_strings() {
1428 assert_eq!(text_to_json_literal("hello"), r#""hello""#);
1429 assert_eq!(text_to_json_literal("req_1"), r#""req_1""#);
1430 assert_eq!(text_to_json_literal(r#"a"b"#), r#""a\"b""#);
1432 }
1433
1434 #[test]
1438 fn write_waits_for_a_contended_lock() {
1439 let dir = tempfile::tempdir().unwrap();
1440 let path = dir.path().join("logs.db");
1441 let store = SqliteLogStore::open(&path).unwrap();
1442
1443 let blocker = Connection::open(&path).unwrap();
1446 blocker.busy_timeout(BUSY_TIMEOUT).unwrap();
1447 blocker.execute_batch("BEGIN IMMEDIATE").unwrap();
1448 let releaser = std::thread::spawn(move || {
1449 std::thread::sleep(std::time::Duration::from_millis(300));
1450 blocker.execute_batch("COMMIT").unwrap();
1451 });
1452
1453 let id = DaemonId::try_new("test", "blocked").unwrap();
1454 let entries = vec![crate::log_parse::parse("held-lock-line", "text")];
1455 store
1456 .append_structured_batch(&id, &entries)
1457 .expect("a contended write must wait for the lock, not fail");
1458
1459 releaser.join().unwrap();
1460 let found = store
1461 .query(&LogQuery {
1462 daemon_ids: vec![id.qualified()],
1463 ..Default::default()
1464 })
1465 .unwrap()
1466 .len();
1467 assert_eq!(found, 1, "the batch written under contention was lost");
1468 }
1469
1470 #[test]
1476 fn concurrent_writers_do_not_lose_batches() {
1477 let dir = tempfile::tempdir().unwrap();
1478 let path = dir.path().join("logs.db");
1479
1480 const WRITERS: usize = 4;
1481 const BATCHES: usize = 15;
1482 const PER_BATCH: usize = 10;
1483
1484 let barrier = std::sync::Barrier::new(WRITERS);
1489
1490 std::thread::scope(|scope| {
1491 for writer in 0..WRITERS {
1492 let path = path.clone();
1493 let barrier = &barrier;
1494 scope.spawn(move || {
1495 barrier.wait();
1498 let store = SqliteLogStore::open(&path).unwrap();
1499 let id = DaemonId::try_new("test", format!("w{writer}")).unwrap();
1500 for batch in 0..BATCHES {
1501 let entries: Vec<ParsedLog> = (0..PER_BATCH)
1502 .map(|i| {
1503 crate::log_parse::parse(&format!("w{writer}-{batch}-{i}"), "text")
1504 })
1505 .collect();
1506 store
1507 .append_structured_batch(&id, &entries)
1508 .expect("concurrent batch write must not fail");
1509 }
1510 });
1511 }
1512 });
1513
1514 let store = SqliteLogStore::open(&path).unwrap();
1515 for writer in 0..WRITERS {
1516 let id = DaemonId::try_new("test", format!("w{writer}")).unwrap();
1517 let found = store
1518 .query(&LogQuery {
1519 daemon_ids: vec![id.qualified()],
1520 ..Default::default()
1521 })
1522 .unwrap()
1523 .len();
1524 assert_eq!(
1525 found,
1526 BATCHES * PER_BATCH,
1527 "writer {writer} lost entries: got {found}"
1528 );
1529 }
1530 }
1531
1532 #[test]
1538 fn batch_rows_share_single_timestamp() {
1539 let dir = tempfile::tempdir().unwrap();
1540 let path = dir.path().join("logs.db");
1541 let store = SqliteLogStore::open(&path).unwrap();
1542 let id = DaemonId::try_new("test", "burst").unwrap();
1543
1544 let entries: Vec<ParsedLog> = (0..50)
1545 .map(|i| crate::log_parse::parse(&format!("line-{i}"), "text"))
1546 .collect();
1547 store.append_structured_batch(&id, &entries).unwrap();
1548
1549 let conn = store.conn.lock().unwrap();
1550 let timestamps: Vec<i64> = {
1551 let mut stmt = conn
1552 .prepare("SELECT timestamp FROM log_entries ORDER BY id")
1553 .unwrap();
1554 let rows = stmt.query_map([], |row| row.get::<_, i64>(0)).unwrap();
1555 rows.map(|r| r.unwrap()).collect()
1556 };
1557 assert_eq!(timestamps.len(), 50);
1558 let distinct: std::collections::HashSet<i64> = timestamps.into_iter().collect();
1559 assert_eq!(
1560 distinct.len(),
1561 1,
1562 "all rows in a batch must share one timestamp so (timestamp, id) ordering matches insertion order"
1563 );
1564 }
1565
1566 #[test]
1571 fn burst_batches_preserve_insertion_order() {
1572 let dir = tempfile::tempdir().unwrap();
1573 let path = dir.path().join("logs.db");
1574 let store = SqliteLogStore::open(&path).unwrap();
1575 let id = DaemonId::try_new("test", "burst").unwrap();
1576
1577 const BATCHES: usize = 3;
1581 const PER_BATCH: usize = 300;
1582 for batch in 0..BATCHES {
1583 let entries: Vec<ParsedLog> = (0..PER_BATCH)
1584 .map(|i| crate::log_parse::parse(&format!("{}", batch * PER_BATCH + i), "text"))
1585 .collect();
1586 store.append_structured_batch(&id, &entries).unwrap();
1587 }
1588
1589 let entries = store
1590 .query(&LogQuery {
1591 daemon_ids: vec![id.qualified()],
1592 ..Default::default()
1593 })
1594 .unwrap();
1595 assert_eq!(entries.len(), BATCHES * PER_BATCH);
1596 for (n, entry) in entries.iter().enumerate() {
1597 assert_eq!(
1598 entry.message,
1599 n.to_string(),
1600 "log line {n} read back out of order (got '{}')",
1601 entry.message
1602 );
1603 }
1604 }
1605}