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