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