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 for chunk in entries.chunks(archive_hook.batch_size.max(1)) {
275 let mut child = Command::new("sh")
276 .arg("-c")
277 .arg(&archive_hook.command)
278 .stdin(Stdio::piped())
279 .stdout(Stdio::null())
280 .stderr(Stdio::piped())
281 .env("PITCHFORK_DAEMON_ID", daemon_id.qualified())
282 .env("PITCHFORK_ARCHIVE_REASON", reason)
283 .hide_console_window()
284 .spawn()
285 .into_diagnostic()
286 .map_err(|e| miette::miette!("failed to spawn archive hook: {e}"))?;
287
288 let write_result = {
292 let stdin = child.stdin.take().expect("piped stdin should be available");
293 let mut stdin = std::io::BufWriter::new(stdin);
294 let mut result = Ok(());
295 for entry in chunk {
296 let line = serde_json::json!({
297 "id": entry.id,
298 "daemon_id": entry.daemon_id,
299 "timestamp": entry.timestamp.to_rfc3339(),
300 "message": entry.message,
301 });
302 if let Err(e) = writeln!(stdin, "{}", line) {
303 result = Err(miette::miette!(
304 "failed to write to archive hook stdin: {e}"
305 ));
306 break;
307 }
308 }
309 if result.is_ok()
317 && let Err(e) = stdin.flush()
318 {
319 result = Err(miette::miette!("failed to flush archive hook stdin: {e}"));
320 }
321 result
322 };
324
325 if let Err(e) = write_result {
326 let _ = child.kill();
327 let _ = child.wait();
328 return Err(e);
329 }
330
331 let output = child.wait_with_output().into_diagnostic()?;
332 if !output.status.success() {
333 let stderr = String::from_utf8_lossy(&output.stderr);
334 return Err(miette::miette!(
335 "archive hook failed with status {}: {stderr}",
336 output.status
337 ));
338 }
339 }
340
341 Ok(())
342 }
343
344 fn delete_by_ids(&self, ids: &[i64]) -> Result<u64> {
349 const SQLITE_MAX_VARS: usize = 999;
350
351 let mut total = 0u64;
352 let conn = self.conn.lock().unwrap();
353 for chunk in ids.chunks(SQLITE_MAX_VARS) {
354 if chunk.is_empty() {
355 continue;
356 }
357 let placeholders: Vec<String> = (1..=chunk.len()).map(|i| format!("?{i}")).collect();
358 let sql = format!(
359 "DELETE FROM log_entries WHERE id IN ({})",
360 placeholders.join(", ")
361 );
362 total += conn
363 .execute(&sql, rusqlite::params_from_iter(chunk.iter()))
364 .into_diagnostic()? as u64;
365 }
366 Ok(total)
367 }
368
369 pub fn rotate_by_age(
376 &self,
377 daemon_id: &DaemonId,
378 max_age: chrono::Duration,
379 archive_hook: Option<&ArchiveHook>,
380 ) -> Result<u64> {
381 let cutoff = (Local::now() - max_age).timestamp_millis();
382 let hook = archive_hook.filter(|h| h.is_enabled());
383
384 if let Some(hook) = hook {
385 let mut total_deleted = 0u64;
386 loop {
387 let entries: Vec<LogEntry> = {
389 let conn = self.conn.lock().unwrap();
390 let mut stmt = conn
391 .prepare(
392 "SELECT id, daemon_id, timestamp, message, level, msg, logger, fields_json FROM log_entries
393 WHERE daemon_id = ?1 AND timestamp < ?2
394 ORDER BY timestamp ASC, id ASC
395 LIMIT ?3",
396 )
397 .into_diagnostic()?;
398 stmt.query_map(
399 params![daemon_id.qualified(), cutoff, hook.batch_size as i64],
400 Self::row_to_entry,
401 )
402 .into_diagnostic()?
403 .collect::<rusqlite::Result<Vec<_>>>()
404 .into_diagnostic()?
405 };
406
407 if entries.is_empty() {
408 break;
409 }
410
411 let batch_ids: Vec<i64> = entries.iter().map(|e| e.id).collect();
412
413 self.archive_entries(&entries, hook, daemon_id, "age")?;
415
416 let deleted = self.delete_by_ids(&batch_ids)?;
418 total_deleted += deleted;
419 }
420 Ok(total_deleted)
421 } else {
422 let conn = self.conn.lock().unwrap();
423 let rows = conn
424 .execute(
425 "DELETE FROM log_entries WHERE daemon_id = ?1 AND timestamp < ?2",
426 params![daemon_id.qualified(), cutoff],
427 )
428 .into_diagnostic()?;
429 Ok(rows as u64)
430 }
431 }
432
433 pub fn rotate_by_count(
441 &self,
442 daemon_id: &DaemonId,
443 max_count: u64,
444 archive_hook: Option<&ArchiveHook>,
445 ) -> Result<u64> {
446 let hook = archive_hook.filter(|h| h.is_enabled());
447
448 let to_delete: i64 = {
450 let conn = self.conn.lock().unwrap();
451 let count: i64 = conn
452 .query_row(
453 "SELECT COUNT(*) FROM log_entries WHERE daemon_id = ?1",
454 [daemon_id.qualified()],
455 |row| row.get(0),
456 )
457 .into_diagnostic()?;
458 count.saturating_sub(max_count as i64)
459 };
460
461 if to_delete <= 0 {
462 return Ok(0);
463 }
464
465 if let Some(hook) = hook {
466 let mut total_deleted = 0u64;
467 let mut remaining = to_delete;
468 loop {
469 let batch_len = remaining.min(hook.batch_size as i64);
470
471 let entries: Vec<LogEntry> = {
473 let conn = self.conn.lock().unwrap();
474 let mut stmt = conn
475 .prepare(
476 "SELECT id, daemon_id, timestamp, message, level, msg, logger, fields_json FROM log_entries
477 WHERE daemon_id = ?1
478 ORDER BY timestamp ASC, id ASC
479 LIMIT ?2",
480 )
481 .into_diagnostic()?;
482 stmt.query_map(
483 params![daemon_id.qualified(), batch_len],
484 Self::row_to_entry,
485 )
486 .into_diagnostic()?
487 .collect::<rusqlite::Result<Vec<_>>>()
488 .into_diagnostic()?
489 };
490
491 if entries.is_empty() {
492 break;
493 }
494
495 let fetched = entries.len() as i64;
496 let batch_ids: Vec<i64> = entries.iter().map(|e| e.id).collect();
497
498 self.archive_entries(&entries, hook, daemon_id, "count")?;
500
501 let deleted = self.delete_by_ids(&batch_ids)?;
503 total_deleted += deleted;
504 remaining -= fetched;
505 }
506 Ok(total_deleted)
507 } else {
508 let conn = self.conn.lock().unwrap();
509 let rows = conn
510 .execute(
511 "DELETE FROM log_entries WHERE id IN (
512 SELECT id FROM log_entries WHERE daemon_id = ?1
513 ORDER BY timestamp ASC, id ASC LIMIT ?2
514 )",
515 params![daemon_id.qualified(), to_delete],
516 )
517 .into_diagnostic()?;
518 Ok(rows as u64)
519 }
520 }
521
522 pub fn migrate_daemon_text_logs(&self, daemon_id: &DaemonId) -> Result<u64> {
527 let text_path = daemon_id.log_path();
528 if !text_path.exists() {
529 return Ok(0);
530 }
531
532 let file = std::fs::File::open(&text_path).into_diagnostic()?;
533 let reader = BufReader::new(file);
534 let re = regex::Regex::new(r"^(\d{4}-\d{2}-\d{2} \d{2}:\d{2}:\d{2}) ([\w./-]+) (.*)$")
535 .expect("invalid regex");
536
537 let mut current_timestamp: Option<DateTime<Local>> = None;
538 let mut current_message = String::new();
539 let mut entries = Vec::with_capacity(1000);
540 let mut total_migrated: u64 = 0;
541
542 for line in reader.lines() {
543 let line = line.into_diagnostic()?;
544 if let Some(caps) = re.captures(&line) {
545 if let Some(ts) = current_timestamp.take() {
546 entries.push((ts, std::mem::take(&mut current_message)));
547 }
548 let ts_str = caps.get(1).map(|m| m.as_str()).unwrap_or_default();
549 let msg = caps.get(3).map(|m| m.as_str()).unwrap_or_default();
550 if let Ok(naive) =
551 chrono::NaiveDateTime::parse_from_str(ts_str, "%Y-%m-%d %H:%M:%S")
552 {
553 current_timestamp = Local.from_local_datetime(&naive).single();
554 current_message = msg.to_string();
555 }
556 } else if current_timestamp.is_some() {
557 current_message.push('\n');
558 current_message.push_str(&line);
559 }
560
561 if entries.len() >= 1000 {
562 total_migrated += self.insert_batch(daemon_id, &entries)?;
563 entries.clear();
564 }
565 }
566
567 if let Some(ts) = current_timestamp {
568 entries.push((ts, std::mem::take(&mut current_message)));
569 }
570
571 if !entries.is_empty() {
572 total_migrated += self.insert_batch(daemon_id, &entries)?;
573 }
574
575 if total_migrated > 0
576 && let Err(e) = std::fs::remove_file(&text_path)
577 {
578 log::warn!(
579 "failed to remove legacy log file after migration {}: {e}",
580 text_path.display()
581 );
582 }
583
584 Ok(total_migrated)
585 }
586
587 fn insert_batch(
588 &self,
589 daemon_id: &DaemonId,
590 entries: &[(DateTime<Local>, String)],
591 ) -> Result<u64> {
592 let mut conn = self.conn.lock().unwrap();
593 let tx = conn.transaction().into_diagnostic()?;
594 let mut count = 0u64;
595 {
596 let mut stmt = tx
597 .prepare(
598 "INSERT INTO log_entries (daemon_id, timestamp, message) VALUES (?1, ?2, ?3)",
599 )
600 .into_diagnostic()?;
601 for (ts, msg) in entries {
602 stmt.execute(params![daemon_id.qualified(), ts.timestamp_millis(), msg])
603 .into_diagnostic()?;
604 count += 1;
605 }
606 }
607 tx.commit().into_diagnostic()?;
608 Ok(count)
609 }
610
611 fn build_query_sql(
616 opts: &LogQuery,
617 id_range: Option<(i64, i64)>,
618 ) -> (String, Vec<Box<dyn rusqlite::ToSql>>) {
619 let mut conditions = Vec::new();
620 let mut query_params: Vec<Box<dyn rusqlite::ToSql>> = Vec::new();
621
622 if !opts.daemon_ids.is_empty() {
623 let placeholders: Vec<String> = (1..=opts.daemon_ids.len())
624 .map(|i| format!("?{}", i))
625 .collect();
626 conditions.push(format!("daemon_id IN ({})", placeholders.join(", ")));
627 for id in &opts.daemon_ids {
628 query_params.push(Box::new(id.clone()));
629 }
630 }
631
632 if let Some(from) = opts.from {
633 conditions.push(format!("timestamp >= ?{}", query_params.len() + 1));
634 query_params.push(Box::new(from.timestamp_millis()));
635 }
636
637 if let Some(to) = opts.to {
638 conditions.push(format!("timestamp <= ?{}", query_params.len() + 1));
639 query_params.push(Box::new(to.timestamp_millis()));
640 }
641
642 if let Some(after_id) = opts.after_id {
643 conditions.push(format!("id > ?{}", query_params.len() + 1));
644 query_params.push(Box::new(after_id));
645 }
646
647 if let Some(before_id) = opts.before_id {
648 conditions.push(format!("id < ?{}", query_params.len() + 1));
649 query_params.push(Box::new(before_id));
650 }
651
652 if let Some((start, end)) = id_range {
653 conditions.push(format!("id > ?{}", query_params.len() + 1));
654 query_params.push(Box::new(start));
655 conditions.push(format!("id <= ?{}", query_params.len() + 1));
656 query_params.push(Box::new(end));
657 }
658
659 let mut message_conditions = Vec::new();
660 for filter in &opts.message_filters {
661 match filter {
662 MessageFilter::Contains {
663 pattern,
664 case_sensitive,
665 } => {
666 let param_index = query_params.len() + 1;
667 if *case_sensitive {
668 message_conditions
669 .push(format!("INSTR(message, ?{idx}) > 0", idx = param_index));
670 query_params.push(Box::new(pattern.clone()));
671 } else {
672 let escaped = escape_like_pattern(pattern);
673 let param = format!("%{}%", escaped);
674 message_conditions.push(format!(
675 "LOWER(message) LIKE LOWER(?{idx}) ESCAPE '\\'",
676 idx = param_index
677 ));
678 query_params.push(Box::new(param));
679 }
680 }
681 MessageFilter::Regex { pattern } => {
682 let param_index = query_params.len() + 1;
683 message_conditions.push(format!("message REGEXP ?{param_index}"));
684 query_params.push(Box::new(pattern.clone()));
685 }
686 }
687 }
688 if !message_conditions.is_empty() {
689 conditions.push(format!("({})", message_conditions.join(" OR ")));
690 }
691
692 for filter in &opts.field_filters {
693 match filter {
694 FieldFilter::LevelMin(level) => {
695 let matching = crate::log_store::levels_at_or_above(level);
696 if matching.is_empty() {
697 conditions.push("0".to_string());
699 } else {
700 let placeholders = matching
701 .iter()
702 .map(|l| {
703 let idx = query_params.len() + 1;
704 query_params.push(Box::new((*l).to_string()));
705 format!("?{idx}")
706 })
707 .collect::<Vec<_>>()
708 .join(", ");
709 conditions.push(format!("level IN ({placeholders})"));
710 }
711 }
712 FieldFilter::FieldEq { key, value } => {
713 let key_idx = query_params.len() + 1;
714 let val_idx = query_params.len() + 2;
715 let json_literal = text_to_json_literal(value);
733 let raw_idx = query_params.len() + 3;
734 conditions.push(format!(
735 "EXISTS (SELECT 1 FROM json_each(fields_json) \
736 WHERE json_each.key = ?{key_idx} \
737 AND (json_each.value IS json_extract(?{val_idx}, '$') \
738 OR json_each.value IS ?{raw_idx}))"
739 ));
740 query_params.push(Box::new(key.clone()));
741 query_params.push(Box::new(json_literal));
742 query_params.push(Box::new(value.clone()));
743 }
744 FieldFilter::LoggerContains(pattern) => {
745 let param_index = query_params.len() + 1;
746 conditions.push(format!("logger LIKE ?{param_index} ESCAPE '\\'"));
747 let escaped = crate::log_store::escape_like_pattern(pattern);
748 query_params.push(Box::new(format!("%{escaped}%")));
749 }
750 }
751 }
752
753 let where_clause = if conditions.is_empty() {
754 String::new()
755 } else {
756 format!("WHERE {}", conditions.join(" AND "))
757 };
758
759 let order = if opts.order_desc { "DESC" } else { "ASC" };
760
761 let limit_clause = opts
762 .limit
763 .map(|n| format!("LIMIT {}", n))
764 .unwrap_or_default();
765
766 let columns = if opts.include_structured {
767 "id, daemon_id, timestamp, message, level, msg, logger, fields_json"
768 } else {
769 "id, daemon_id, timestamp, message, NULL, NULL, NULL, NULL"
770 };
771
772 let sql = format!(
773 "SELECT {columns} FROM log_entries {} ORDER BY timestamp {}, id {} {}",
774 where_clause, order, order, limit_clause
775 );
776
777 (sql, query_params)
778 }
779
780 fn should_parallelize(opts: &LogQuery) -> bool {
782 if opts.daemon_ids.len() != 1 {
783 return false;
784 }
785 if opts.after_id.is_some() {
788 return false;
789 }
790 let limit = opts.limit.unwrap_or(usize::MAX);
791 if limit < PARALLEL_QUERY_THRESHOLD {
792 return false;
793 }
794 std::thread::available_parallelism()
795 .map(|n| n.get() >= 2)
796 .unwrap_or(false)
797 }
798
799 fn query_parallel(&self, opts: &LogQuery) -> Result<Vec<LogEntry>> {
806 let max_threads = 2;
811 let num_threads = std::thread::available_parallelism()
812 .map(|n| n.get().min(max_threads))
813 .unwrap_or(1);
814
815 let max_id: Option<i64> = {
816 let conn = self.conn.lock().unwrap();
817 conn.query_row(
818 "SELECT MAX(id) FROM log_entries WHERE daemon_id = ?1",
819 params![&opts.daemon_ids[0]],
820 |row| row.get(0),
821 )
822 .ok()
823 };
824
825 let Some(max_id) = max_id else {
826 return Ok(Vec::new());
827 };
828 if max_id == 0 {
829 return Ok(Vec::new());
830 }
831
832 let shard_size = (max_id as usize).div_ceil(num_threads);
833 let path = self.path.clone();
834 let opts = opts.clone();
835 let needs_regexp = opts
836 .message_filters
837 .iter()
838 .any(|f| matches!(f, MessageFilter::Regex { .. }));
839
840 let shards: Vec<Result<Vec<LogEntry>>> = std::thread::scope(|s| {
841 (0..num_threads)
842 .map(|i| {
843 let start = (i * shard_size) as i64;
844 let end = if i == num_threads - 1 {
845 max_id
846 } else {
847 ((i + 1) * shard_size) as i64
848 };
849 let opts = opts.clone();
850 let path = path.clone();
851 s.spawn(move || -> Result<Vec<LogEntry>> {
852 let conn = Connection::open_with_flags(
853 &path,
854 OpenFlags::SQLITE_OPEN_READ_ONLY | OpenFlags::SQLITE_OPEN_NO_MUTEX,
855 )
856 .into_diagnostic()?;
857 conn.execute_batch(
858 "PRAGMA mmap_size = 268435456;
859 PRAGMA query_only = 1;",
860 )
861 .into_diagnostic()?;
862 if needs_regexp {
863 add_regexp_function(&conn)?;
864 }
865
866 let (sql, query_params) = Self::build_query_sql(&opts, Some((start, end)));
867 Self::execute_built_query(&conn, &sql, &query_params)
868 })
869 })
870 .map(|h| h.join().unwrap())
871 .collect()
872 });
873
874 let mut merged = Vec::new();
875 if opts.order_desc {
876 for shard in shards.into_iter().rev() {
877 merged.extend(shard?);
878 }
879 } else {
880 for shard in shards {
881 merged.extend(shard?);
882 }
883 }
884
885 if let Some(limit) = opts.limit
886 && merged.len() > limit
887 {
888 merged.truncate(limit);
889 }
890
891 Ok(merged)
892 }
893
894 fn execute_built_query(
896 conn: &Connection,
897 sql: &str,
898 query_params: &[Box<dyn rusqlite::ToSql>],
899 ) -> Result<Vec<LogEntry>> {
900 let mut stmt = conn.prepare(sql).into_diagnostic()?;
901 let params_ref: Vec<&dyn rusqlite::ToSql> =
902 query_params.iter().map(|p| p.as_ref()).collect();
903 let rows = stmt
904 .query_map(params_ref.as_slice(), Self::row_to_entry)
905 .into_diagnostic()?;
906 let mut entries = Vec::new();
907 for row in rows {
908 entries.push(row.into_diagnostic()?);
909 }
910 Ok(entries)
911 }
912
913 pub fn distinct_loggers(&self, daemon_id: &str) -> Result<Vec<String>> {
918 let conn = self.conn.lock().unwrap();
919 let mut stmt = conn
920 .prepare(
921 "SELECT DISTINCT logger FROM ( \
922 SELECT logger FROM log_entries \
923 WHERE daemon_id = ?1 AND logger IS NOT NULL \
924 ORDER BY id DESC LIMIT 5000 \
925 ) ORDER BY logger",
926 )
927 .into_diagnostic()?;
928 let rows = stmt
929 .query_map(params![daemon_id], |row| row.get::<_, String>(0))
930 .into_diagnostic()?;
931 let mut loggers = Vec::new();
932 for row in rows {
933 loggers.push(row.into_diagnostic()?);
934 }
935 Ok(loggers)
936 }
937
938 pub fn distinct_field_keys(&self, daemon_id: &str) -> Result<Vec<String>> {
944 let conn = self.conn.lock().unwrap();
945 let mut stmt = conn
946 .prepare(
947 "SELECT DISTINCT je.key \
948 FROM ( \
949 SELECT fields_json FROM log_entries \
950 WHERE daemon_id = ?1 AND fields_json IS NOT NULL \
951 ORDER BY id DESC LIMIT 5000 \
952 ), json_each(fields_json) AS je \
953 ORDER BY je.key",
954 )
955 .into_diagnostic()?;
956 let rows = stmt
957 .query_map(params![daemon_id], |row| row.get::<_, String>(0))
958 .into_diagnostic()?;
959 let mut keys = Vec::new();
960 for row in rows {
961 keys.push(row.into_diagnostic()?);
962 }
963 Ok(keys)
964 }
965}
966
967impl LogStore for SqliteLogStore {
968 fn append(&self, daemon_id: &DaemonId, message: &str) -> Result<()> {
969 let id = daemon_id.qualified();
970 let msg = message.to_string();
971
972 let conn = self.conn.lock().unwrap();
973 let ts = Local::now().timestamp_millis();
976 let _ = conn
977 .execute(
978 "INSERT INTO log_entries (daemon_id, timestamp, message) VALUES (?1, ?2, ?3)",
979 params![id, ts, msg],
980 )
981 .into_diagnostic()?;
982 Ok(())
983 }
984
985 fn append_batch(&self, daemon_id: &DaemonId, messages: &[String]) -> Result<()> {
986 if messages.is_empty() {
987 return Ok(());
988 }
989 let id = daemon_id.qualified();
990
991 let mut conn = self.conn.lock().unwrap();
992 let tx = conn
997 .transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)
998 .into_diagnostic()?;
999 let base_ts = Local::now().timestamp_millis();
1000 {
1001 let mut stmt = tx
1002 .prepare(
1003 "INSERT INTO log_entries (daemon_id, timestamp, message) VALUES (?1, ?2, ?3)",
1004 )
1005 .into_diagnostic()?;
1006 for msg in messages.iter() {
1007 let ts = base_ts;
1014 stmt.execute(params![id, ts, msg]).into_diagnostic()?;
1015 }
1016 }
1017 tx.commit().into_diagnostic()?;
1018 Ok(())
1019 }
1020
1021 fn append_structured(&self, daemon_id: &DaemonId, parsed: &ParsedLog) -> Result<()> {
1022 let id = daemon_id.qualified();
1023
1024 let conn = self.conn.lock().unwrap();
1025 let ts = Local::now().timestamp_millis();
1028 let _ = conn
1029 .execute(
1030 "INSERT INTO log_entries (daemon_id, timestamp, message, level, msg, logger, fields_json) \
1031 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)",
1032 params![id, ts, parsed.message, parsed.level, parsed.msg, parsed.logger, parsed.fields_json],
1033 )
1034 .into_diagnostic()?;
1035 Ok(())
1036 }
1037
1038 fn append_structured_batch(&self, daemon_id: &DaemonId, entries: &[ParsedLog]) -> Result<()> {
1039 if entries.is_empty() {
1040 return Ok(());
1041 }
1042 let id = daemon_id.qualified();
1043
1044 let mut conn = self.conn.lock().unwrap();
1045 let tx = conn
1050 .transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)
1051 .into_diagnostic()?;
1052 let base_ts = Local::now().timestamp_millis();
1053 {
1054 let mut stmt = tx
1055 .prepare(
1056 "INSERT INTO log_entries (daemon_id, timestamp, message, level, msg, logger, fields_json) \
1057 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)",
1058 )
1059 .into_diagnostic()?;
1060 for entry in entries.iter() {
1061 let ts = base_ts;
1066 stmt.execute(params![
1067 id,
1068 ts,
1069 entry.message,
1070 entry.level,
1071 entry.msg,
1072 entry.logger,
1073 entry.fields_json
1074 ])
1075 .into_diagnostic()?;
1076 }
1077 }
1078 tx.commit().into_diagnostic()?;
1079 Ok(())
1080 }
1081
1082 fn query(&self, opts: &LogQuery) -> Result<Vec<LogEntry>> {
1083 if Self::should_parallelize(opts)
1085 && self.path.as_os_str() != ":memory:"
1086 && let Ok(entries) = self.query_parallel(opts)
1087 {
1088 return Ok(entries);
1089 }
1090 let conn = self.conn.lock().unwrap();
1094 let (sql, query_params) = Self::build_query_sql(opts, None);
1095 Self::execute_built_query(&conn, &sql, &query_params)
1096 }
1097
1098 fn tail(&self, daemon_id: &DaemonId, after_id: Option<i64>) -> Result<Vec<LogEntry>> {
1099 self.query(&LogQuery {
1100 daemon_ids: vec![daemon_id.qualified()],
1101 from: None,
1102 to: None,
1103 limit: None,
1104 order_desc: false,
1105 after_id,
1106 before_id: None,
1107 message_filters: Vec::new(),
1108 field_filters: Vec::new(),
1109 include_structured: false,
1110 })
1111 }
1112
1113 fn clear(&self, daemon_ids: &[DaemonId]) -> Result<()> {
1114 let mut conn = self.conn.lock().unwrap();
1115 let tx = conn.transaction().into_diagnostic()?;
1116 for id in daemon_ids {
1117 tx.execute(
1118 "DELETE FROM log_entries WHERE daemon_id = ?1",
1119 params![id.qualified()],
1120 )
1121 .into_diagnostic()?;
1122 tx.execute(
1123 "INSERT INTO log_clear_generations (daemon_id, generation)
1124 VALUES (?1, 1)
1125 ON CONFLICT(daemon_id) DO UPDATE SET generation = generation + 1",
1126 params![id.qualified()],
1127 )
1128 .into_diagnostic()?;
1129 }
1130 tx.commit().into_diagnostic()?;
1131 Ok(())
1132 }
1133
1134 fn last_id(&self, daemon_id: &DaemonId) -> Result<Option<i64>> {
1135 let conn = self.conn.lock().unwrap();
1136 let id: Option<i64> = conn
1139 .query_row(
1140 "SELECT MAX(id) FROM log_entries WHERE daemon_id = ?1",
1141 params![daemon_id.qualified()],
1142 |row| row.get(0),
1143 )
1144 .into_diagnostic()?;
1145 Ok(id)
1146 }
1147
1148 fn list_daemon_ids(&self) -> Result<Vec<String>> {
1149 let conn = self.conn.lock().unwrap();
1150 let mut stmt = conn
1151 .prepare("SELECT DISTINCT daemon_id FROM log_entries")
1152 .into_diagnostic()?;
1153 let ids = stmt
1154 .query_map([], |row| {
1155 let id: String = row.get(0)?;
1156 Ok(id)
1157 })
1158 .into_diagnostic()?
1159 .filter_map(|r| r.ok())
1160 .collect();
1161 Ok(ids)
1162 }
1163
1164 fn apply_retention(
1165 &self,
1166 policy: &super::RetentionPolicy,
1167 excluded_daemon_ids: &[DaemonId],
1168 archive_hook: Option<&ArchiveHook>,
1169 ) -> Result<u64> {
1170 let daemon_ids = self.list_daemon_ids()?;
1171 let excluded: HashSet<String> = excluded_daemon_ids.iter().map(|d| d.qualified()).collect();
1172 let mut total = 0u64;
1173 for id_str in daemon_ids {
1174 if excluded.contains(&id_str) {
1175 continue;
1176 }
1177 let id = DaemonId::parse(&id_str).unwrap_or_else(|_| {
1178 DaemonId::try_new("global", &id_str).unwrap_or_else(|_| DaemonId::pitchfork())
1179 });
1180 if let Some(dur) = policy.age {
1181 total += self.rotate_by_age(&id, dur, archive_hook)?;
1182 }
1183 if let Some(n) = policy.count {
1184 total += self.rotate_by_count(&id, n, archive_hook)?;
1185 }
1186 }
1187 Ok(total)
1188 }
1189
1190 fn apply_retention_for_daemon(
1191 &self,
1192 daemon_id: &DaemonId,
1193 policy: &super::RetentionPolicy,
1194 archive_hook: Option<&ArchiveHook>,
1195 ) -> Result<u64> {
1196 let mut total = 0u64;
1197 if let Some(dur) = policy.age {
1198 total += self.rotate_by_age(daemon_id, dur, archive_hook)?;
1199 }
1200 if let Some(n) = policy.count {
1201 total += self.rotate_by_count(daemon_id, n, archive_hook)?;
1202 }
1203 Ok(total)
1204 }
1205
1206 fn last_clear_generation(&self, daemon_id: &DaemonId) -> Result<Option<u64>> {
1207 let conn = self.conn.lock().unwrap();
1208 let mut stmt = conn
1209 .prepare("SELECT generation FROM log_clear_generations WHERE daemon_id = ?1")
1210 .into_diagnostic()?;
1211 let generation: Option<i64> = stmt
1212 .query_row(params![daemon_id.qualified()], |row| row.get(0))
1213 .optional()
1214 .into_diagnostic()?;
1215 generation
1216 .map(|generation| {
1217 u64::try_from(generation)
1218 .map_err(|_| miette::miette!("log clear generation cannot be negative"))
1219 })
1220 .transpose()
1221 }
1222
1223 fn query_with_generation(
1224 &self,
1225 opts: &LogQuery,
1226 daemon_id: &DaemonId,
1227 ) -> Result<(Vec<LogEntry>, Option<u64>)> {
1228 let conn = self.conn.lock().unwrap();
1229 conn.execute_batch("BEGIN").into_diagnostic()?;
1238 let result = (|| {
1239 let (sql, query_params) = Self::build_query_sql(opts, None);
1240 let entries = Self::execute_built_query(&conn, &sql, &query_params)?;
1241 let generation: Option<i64> = conn
1242 .query_row(
1243 "SELECT generation FROM log_clear_generations WHERE daemon_id = ?1",
1244 params![daemon_id.qualified()],
1245 |row| row.get(0),
1246 )
1247 .optional()
1248 .into_diagnostic()?;
1249 let generation = generation
1250 .map(|g| {
1251 u64::try_from(g)
1252 .map_err(|_| miette::miette!("log clear generation cannot be negative"))
1253 })
1254 .transpose()?;
1255 Ok((entries, generation))
1256 })();
1257 let _ = conn.execute_batch("ROLLBACK");
1259 result
1260 }
1261}
1262
1263use once_cell::sync::Lazy;
1265use std::sync::Arc;
1266
1267pub static LOG_STORE: Lazy<Arc<SqliteLogStore>> = Lazy::new(|| {
1268 let path = crate::env::PITCHFORK_LOGS_DIR.join("logs.db");
1269 let mut is_fallback = false;
1270 let store = Arc::new(SqliteLogStore::open(&path).unwrap_or_else(|e| {
1271 error!(
1272 "failed to open log store at {}: {e}. Falling back to in-memory store; logs will not persist across restarts.",
1273 path.display()
1274 );
1275 is_fallback = true;
1276 SqliteLogStore::open(":memory:").expect("in-memory SQLite should always open")
1277 }));
1278
1279 if !is_fallback {
1283 if let Err(e) = auto_migrate_legacy_logs(&store) {
1284 warn!("legacy log auto-migration failed: {e}");
1285 }
1286 } else {
1287 warn!(
1288 "skipping legacy log auto-migration because log store is in-memory (no durable destination)"
1289 );
1290 }
1291
1292 store
1293});
1294
1295fn auto_migrate_legacy_logs(store: &SqliteLogStore) -> Result<()> {
1303 let logs_dir = &*crate::env::PITCHFORK_LOGS_DIR;
1304 if !logs_dir.exists() {
1305 return Ok(());
1306 }
1307
1308 let Ok(entries) = std::fs::read_dir(logs_dir) else {
1309 return Ok(());
1310 };
1311
1312 let mut total_migrated = 0u64;
1313 let mut migrated_ids = Vec::new();
1314
1315 for entry in entries.flatten() {
1316 let path = entry.path();
1317 if !path.is_dir() {
1318 continue;
1319 }
1320 let file_name = path
1322 .file_name()
1323 .map_or(String::new(), |n| n.to_string_lossy().to_string());
1324 if file_name == "pitchfork" {
1325 continue;
1326 }
1327
1328 if !file_name.contains("--") {
1330 continue;
1331 }
1332 let log_file = path.join(format!("{file_name}.log"));
1333 if !log_file.exists() {
1334 continue;
1335 }
1336
1337 let daemon_id = match DaemonId::from_safe_path(&file_name) {
1338 Ok(id) => id,
1339 Err(_) => continue,
1340 };
1341
1342 if daemon_id == DaemonId::pitchfork() {
1345 continue;
1346 }
1347
1348 match store.migrate_daemon_text_logs(&daemon_id) {
1349 Ok(0) => {}
1350 Ok(n) => {
1351 total_migrated += n;
1352 migrated_ids.push(daemon_id.qualified());
1353 }
1354 Err(e) => {
1355 warn!(
1356 "failed to migrate text logs for {}: {e}",
1357 daemon_id.qualified()
1358 );
1359 }
1360 }
1361 }
1362
1363 if total_migrated > 0 {
1364 warn!(
1365 "auto-migrated {total_migrated} legacy log entries from {count} daemon(s): {ids}",
1366 count = migrated_ids.len(),
1367 ids = migrated_ids.join(", ")
1368 );
1369 }
1370
1371 Ok(())
1372}
1373
1374#[cfg(test)]
1375mod tests {
1376 use super::text_to_json_literal;
1377 use super::*;
1378 use crate::log_store::LogStore;
1379
1380 #[test]
1381 fn test_json_literal_booleans() {
1382 assert_eq!(text_to_json_literal("true"), "true");
1383 assert_eq!(text_to_json_literal("TRUE"), "true");
1384 assert_eq!(text_to_json_literal("false"), "false");
1385 assert_eq!(text_to_json_literal("False"), "false");
1386 }
1387
1388 #[test]
1389 fn test_json_literal_null() {
1390 assert_eq!(text_to_json_literal("null"), "null");
1391 assert_eq!(text_to_json_literal("NULL"), "null");
1392 }
1393
1394 #[test]
1395 fn test_json_literal_valid_numbers() {
1396 assert_eq!(text_to_json_literal("42"), "42");
1397 assert_eq!(text_to_json_literal("-1"), "-1");
1398 assert_eq!(text_to_json_literal("0"), "0");
1399 assert_eq!(text_to_json_literal("3.14"), "3.14");
1400 assert_eq!(text_to_json_literal("1e10"), "1e10");
1401 assert_eq!(text_to_json_literal("1.5e-3"), "1.5e-3");
1402 }
1403
1404 #[test]
1405 fn test_json_literal_rejects_invalid_json_numbers() {
1406 assert_eq!(text_to_json_literal("+42"), r#""+42""#);
1409 assert_eq!(text_to_json_literal("1."), r#""1.""#);
1410 assert_eq!(text_to_json_literal(".5"), r#"".5""#);
1411 assert_eq!(text_to_json_literal("inf"), r#""inf""#);
1412 assert_eq!(text_to_json_literal("nan"), r#""nan""#);
1413 assert_eq!(text_to_json_literal("infinity"), r#""infinity""#);
1414 }
1415
1416 #[test]
1417 fn test_json_literal_strings() {
1418 assert_eq!(text_to_json_literal("hello"), r#""hello""#);
1419 assert_eq!(text_to_json_literal("req_1"), r#""req_1""#);
1420 assert_eq!(text_to_json_literal(r#"a"b"#), r#""a\"b""#);
1422 }
1423
1424 #[test]
1428 fn write_waits_for_a_contended_lock() {
1429 let dir = tempfile::tempdir().unwrap();
1430 let path = dir.path().join("logs.db");
1431 let store = SqliteLogStore::open(&path).unwrap();
1432
1433 let blocker = Connection::open(&path).unwrap();
1436 blocker.busy_timeout(BUSY_TIMEOUT).unwrap();
1437 blocker.execute_batch("BEGIN IMMEDIATE").unwrap();
1438 let releaser = std::thread::spawn(move || {
1439 std::thread::sleep(std::time::Duration::from_millis(300));
1440 blocker.execute_batch("COMMIT").unwrap();
1441 });
1442
1443 let id = DaemonId::try_new("test", "blocked").unwrap();
1444 let entries = vec![crate::log_parse::parse("held-lock-line", "text")];
1445 store
1446 .append_structured_batch(&id, &entries)
1447 .expect("a contended write must wait for the lock, not fail");
1448
1449 releaser.join().unwrap();
1450 let found = store
1451 .query(&LogQuery {
1452 daemon_ids: vec![id.qualified()],
1453 ..Default::default()
1454 })
1455 .unwrap()
1456 .len();
1457 assert_eq!(found, 1, "the batch written under contention was lost");
1458 }
1459
1460 #[test]
1466 fn concurrent_writers_do_not_lose_batches() {
1467 let dir = tempfile::tempdir().unwrap();
1468 let path = dir.path().join("logs.db");
1469
1470 const WRITERS: usize = 4;
1471 const BATCHES: usize = 15;
1472 const PER_BATCH: usize = 10;
1473
1474 let barrier = std::sync::Barrier::new(WRITERS);
1479
1480 std::thread::scope(|scope| {
1481 for writer in 0..WRITERS {
1482 let path = path.clone();
1483 let barrier = &barrier;
1484 scope.spawn(move || {
1485 barrier.wait();
1488 let store = SqliteLogStore::open(&path).unwrap();
1489 let id = DaemonId::try_new("test", format!("w{writer}")).unwrap();
1490 for batch in 0..BATCHES {
1491 let entries: Vec<ParsedLog> = (0..PER_BATCH)
1492 .map(|i| {
1493 crate::log_parse::parse(&format!("w{writer}-{batch}-{i}"), "text")
1494 })
1495 .collect();
1496 store
1497 .append_structured_batch(&id, &entries)
1498 .expect("concurrent batch write must not fail");
1499 }
1500 });
1501 }
1502 });
1503
1504 let store = SqliteLogStore::open(&path).unwrap();
1505 for writer in 0..WRITERS {
1506 let id = DaemonId::try_new("test", format!("w{writer}")).unwrap();
1507 let found = store
1508 .query(&LogQuery {
1509 daemon_ids: vec![id.qualified()],
1510 ..Default::default()
1511 })
1512 .unwrap()
1513 .len();
1514 assert_eq!(
1515 found,
1516 BATCHES * PER_BATCH,
1517 "writer {writer} lost entries: got {found}"
1518 );
1519 }
1520 }
1521
1522 #[test]
1528 fn batch_rows_share_single_timestamp() {
1529 let dir = tempfile::tempdir().unwrap();
1530 let path = dir.path().join("logs.db");
1531 let store = SqliteLogStore::open(&path).unwrap();
1532 let id = DaemonId::try_new("test", "burst").unwrap();
1533
1534 let entries: Vec<ParsedLog> = (0..50)
1535 .map(|i| crate::log_parse::parse(&format!("line-{i}"), "text"))
1536 .collect();
1537 store.append_structured_batch(&id, &entries).unwrap();
1538
1539 let conn = store.conn.lock().unwrap();
1540 let timestamps: Vec<i64> = {
1541 let mut stmt = conn
1542 .prepare("SELECT timestamp FROM log_entries ORDER BY id")
1543 .unwrap();
1544 let rows = stmt.query_map([], |row| row.get::<_, i64>(0)).unwrap();
1545 rows.map(|r| r.unwrap()).collect()
1546 };
1547 assert_eq!(timestamps.len(), 50);
1548 let distinct: std::collections::HashSet<i64> = timestamps.into_iter().collect();
1549 assert_eq!(
1550 distinct.len(),
1551 1,
1552 "all rows in a batch must share one timestamp so (timestamp, id) ordering matches insertion order"
1553 );
1554 }
1555
1556 #[test]
1561 fn burst_batches_preserve_insertion_order() {
1562 let dir = tempfile::tempdir().unwrap();
1563 let path = dir.path().join("logs.db");
1564 let store = SqliteLogStore::open(&path).unwrap();
1565 let id = DaemonId::try_new("test", "burst").unwrap();
1566
1567 const BATCHES: usize = 3;
1571 const PER_BATCH: usize = 300;
1572 for batch in 0..BATCHES {
1573 let entries: Vec<ParsedLog> = (0..PER_BATCH)
1574 .map(|i| crate::log_parse::parse(&format!("{}", batch * PER_BATCH + i), "text"))
1575 .collect();
1576 store.append_structured_batch(&id, &entries).unwrap();
1577 }
1578
1579 let entries = store
1580 .query(&LogQuery {
1581 daemon_ids: vec![id.qualified()],
1582 ..Default::default()
1583 })
1584 .unwrap();
1585 assert_eq!(entries.len(), BATCHES * PER_BATCH);
1586 for (n, entry) in entries.iter().enumerate() {
1587 assert_eq!(
1588 entry.message,
1589 n.to_string(),
1590 "log line {n} read back out of order (got '{}')",
1591 entry.message
1592 );
1593 }
1594 }
1595}