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 ts = Local::now().timestamp_millis();
968 let id = daemon_id.qualified();
969 let msg = message.to_string();
970
971 let conn = self.conn.lock().unwrap();
972 let _ = conn
973 .execute(
974 "INSERT INTO log_entries (daemon_id, timestamp, message) VALUES (?1, ?2, ?3)",
975 params![id, ts, msg],
976 )
977 .into_diagnostic()?;
978 Ok(())
979 }
980
981 fn append_batch(&self, daemon_id: &DaemonId, messages: &[String]) -> Result<()> {
982 if messages.is_empty() {
983 return Ok(());
984 }
985 let base_ts = Local::now().timestamp_millis();
986 let id = daemon_id.qualified();
987
988 let mut conn = self.conn.lock().unwrap();
989 let tx = conn.transaction().into_diagnostic()?;
990 {
991 let mut stmt = tx
992 .prepare(
993 "INSERT INTO log_entries (daemon_id, timestamp, message) VALUES (?1, ?2, ?3)",
994 )
995 .into_diagnostic()?;
996 for (idx, msg) in messages.iter().enumerate() {
997 let ts = base_ts + idx as i64;
1001 stmt.execute(params![id, ts, msg]).into_diagnostic()?;
1002 }
1003 }
1004 tx.commit().into_diagnostic()?;
1005 Ok(())
1006 }
1007
1008 fn append_structured(&self, daemon_id: &DaemonId, parsed: &ParsedLog) -> Result<()> {
1009 let ts = Local::now().timestamp_millis();
1010 let id = daemon_id.qualified();
1011
1012 let conn = self.conn.lock().unwrap();
1013 let _ = conn
1014 .execute(
1015 "INSERT INTO log_entries (daemon_id, timestamp, message, level, msg, logger, fields_json) \
1016 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)",
1017 params![id, ts, parsed.message, parsed.level, parsed.msg, parsed.logger, parsed.fields_json],
1018 )
1019 .into_diagnostic()?;
1020 Ok(())
1021 }
1022
1023 fn append_structured_batch(&self, daemon_id: &DaemonId, entries: &[ParsedLog]) -> Result<()> {
1024 if entries.is_empty() {
1025 return Ok(());
1026 }
1027 let base_ts = Local::now().timestamp_millis();
1028 let id = daemon_id.qualified();
1029
1030 let mut conn = self.conn.lock().unwrap();
1031 let tx = conn.transaction().into_diagnostic()?;
1032 {
1033 let mut stmt = tx
1034 .prepare(
1035 "INSERT INTO log_entries (daemon_id, timestamp, message, level, msg, logger, fields_json) \
1036 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)",
1037 )
1038 .into_diagnostic()?;
1039 for (idx, entry) in entries.iter().enumerate() {
1040 let ts = base_ts + idx as i64;
1041 stmt.execute(params![
1042 id,
1043 ts,
1044 entry.message,
1045 entry.level,
1046 entry.msg,
1047 entry.logger,
1048 entry.fields_json
1049 ])
1050 .into_diagnostic()?;
1051 }
1052 }
1053 tx.commit().into_diagnostic()?;
1054 Ok(())
1055 }
1056
1057 fn query(&self, opts: &LogQuery) -> Result<Vec<LogEntry>> {
1058 if Self::should_parallelize(opts)
1060 && self.path.as_os_str() != ":memory:"
1061 && let Ok(entries) = self.query_parallel(opts)
1062 {
1063 return Ok(entries);
1064 }
1065 let conn = self.conn.lock().unwrap();
1069 let (sql, query_params) = Self::build_query_sql(opts, None);
1070 Self::execute_built_query(&conn, &sql, &query_params)
1071 }
1072
1073 fn tail(&self, daemon_id: &DaemonId, after_id: Option<i64>) -> Result<Vec<LogEntry>> {
1074 self.query(&LogQuery {
1075 daemon_ids: vec![daemon_id.qualified()],
1076 from: None,
1077 to: None,
1078 limit: None,
1079 order_desc: false,
1080 after_id,
1081 before_id: None,
1082 message_filters: Vec::new(),
1083 field_filters: Vec::new(),
1084 include_structured: false,
1085 })
1086 }
1087
1088 fn clear(&self, daemon_ids: &[DaemonId]) -> Result<()> {
1089 let mut conn = self.conn.lock().unwrap();
1090 let tx = conn.transaction().into_diagnostic()?;
1091 for id in daemon_ids {
1092 tx.execute(
1093 "DELETE FROM log_entries WHERE daemon_id = ?1",
1094 params![id.qualified()],
1095 )
1096 .into_diagnostic()?;
1097 tx.execute(
1098 "INSERT INTO log_clear_generations (daemon_id, generation)
1099 VALUES (?1, 1)
1100 ON CONFLICT(daemon_id) DO UPDATE SET generation = generation + 1",
1101 params![id.qualified()],
1102 )
1103 .into_diagnostic()?;
1104 }
1105 tx.commit().into_diagnostic()?;
1106 Ok(())
1107 }
1108
1109 fn last_id(&self, daemon_id: &DaemonId) -> Result<Option<i64>> {
1110 let conn = self.conn.lock().unwrap();
1111 let id: Option<i64> = conn
1114 .query_row(
1115 "SELECT MAX(id) FROM log_entries WHERE daemon_id = ?1",
1116 params![daemon_id.qualified()],
1117 |row| row.get(0),
1118 )
1119 .into_diagnostic()?;
1120 Ok(id)
1121 }
1122
1123 fn list_daemon_ids(&self) -> Result<Vec<String>> {
1124 let conn = self.conn.lock().unwrap();
1125 let mut stmt = conn
1126 .prepare("SELECT DISTINCT daemon_id FROM log_entries")
1127 .into_diagnostic()?;
1128 let ids = stmt
1129 .query_map([], |row| {
1130 let id: String = row.get(0)?;
1131 Ok(id)
1132 })
1133 .into_diagnostic()?
1134 .filter_map(|r| r.ok())
1135 .collect();
1136 Ok(ids)
1137 }
1138
1139 fn apply_retention(
1140 &self,
1141 policy: &super::RetentionPolicy,
1142 excluded_daemon_ids: &[DaemonId],
1143 archive_hook: Option<&ArchiveHook>,
1144 ) -> Result<u64> {
1145 let daemon_ids = self.list_daemon_ids()?;
1146 let excluded: HashSet<String> = excluded_daemon_ids.iter().map(|d| d.qualified()).collect();
1147 let mut total = 0u64;
1148 for id_str in daemon_ids {
1149 if excluded.contains(&id_str) {
1150 continue;
1151 }
1152 let id = DaemonId::parse(&id_str).unwrap_or_else(|_| {
1153 DaemonId::try_new("global", &id_str).unwrap_or_else(|_| DaemonId::pitchfork())
1154 });
1155 if let Some(dur) = policy.age {
1156 total += self.rotate_by_age(&id, dur, archive_hook)?;
1157 }
1158 if let Some(n) = policy.count {
1159 total += self.rotate_by_count(&id, n, archive_hook)?;
1160 }
1161 }
1162 Ok(total)
1163 }
1164
1165 fn apply_retention_for_daemon(
1166 &self,
1167 daemon_id: &DaemonId,
1168 policy: &super::RetentionPolicy,
1169 archive_hook: Option<&ArchiveHook>,
1170 ) -> Result<u64> {
1171 let mut total = 0u64;
1172 if let Some(dur) = policy.age {
1173 total += self.rotate_by_age(daemon_id, dur, archive_hook)?;
1174 }
1175 if let Some(n) = policy.count {
1176 total += self.rotate_by_count(daemon_id, n, archive_hook)?;
1177 }
1178 Ok(total)
1179 }
1180
1181 fn last_clear_generation(&self, daemon_id: &DaemonId) -> Result<Option<u64>> {
1182 let conn = self.conn.lock().unwrap();
1183 let mut stmt = conn
1184 .prepare("SELECT generation FROM log_clear_generations WHERE daemon_id = ?1")
1185 .into_diagnostic()?;
1186 let generation: Option<i64> = stmt
1187 .query_row(params![daemon_id.qualified()], |row| row.get(0))
1188 .optional()
1189 .into_diagnostic()?;
1190 generation
1191 .map(|generation| {
1192 u64::try_from(generation)
1193 .map_err(|_| miette::miette!("log clear generation cannot be negative"))
1194 })
1195 .transpose()
1196 }
1197
1198 fn query_with_generation(
1199 &self,
1200 opts: &LogQuery,
1201 daemon_id: &DaemonId,
1202 ) -> Result<(Vec<LogEntry>, Option<u64>)> {
1203 let conn = self.conn.lock().unwrap();
1204 conn.execute_batch("BEGIN").into_diagnostic()?;
1213 let result = (|| {
1214 let (sql, query_params) = Self::build_query_sql(opts, None);
1215 let entries = Self::execute_built_query(&conn, &sql, &query_params)?;
1216 let generation: Option<i64> = conn
1217 .query_row(
1218 "SELECT generation FROM log_clear_generations WHERE daemon_id = ?1",
1219 params![daemon_id.qualified()],
1220 |row| row.get(0),
1221 )
1222 .optional()
1223 .into_diagnostic()?;
1224 let generation = generation
1225 .map(|g| {
1226 u64::try_from(g)
1227 .map_err(|_| miette::miette!("log clear generation cannot be negative"))
1228 })
1229 .transpose()?;
1230 Ok((entries, generation))
1231 })();
1232 let _ = conn.execute_batch("ROLLBACK");
1234 result
1235 }
1236}
1237
1238use once_cell::sync::Lazy;
1240use std::sync::Arc;
1241
1242pub static LOG_STORE: Lazy<Arc<SqliteLogStore>> = Lazy::new(|| {
1243 let path = crate::env::PITCHFORK_LOGS_DIR.join("logs.db");
1244 let mut is_fallback = false;
1245 let store = Arc::new(SqliteLogStore::open(&path).unwrap_or_else(|e| {
1246 error!(
1247 "failed to open log store at {}: {e}. Falling back to in-memory store; logs will not persist across restarts.",
1248 path.display()
1249 );
1250 is_fallback = true;
1251 SqliteLogStore::open(":memory:").expect("in-memory SQLite should always open")
1252 }));
1253
1254 if !is_fallback {
1258 if let Err(e) = auto_migrate_legacy_logs(&store) {
1259 warn!("legacy log auto-migration failed: {e}");
1260 }
1261 } else {
1262 warn!(
1263 "skipping legacy log auto-migration because log store is in-memory (no durable destination)"
1264 );
1265 }
1266
1267 store
1268});
1269
1270fn auto_migrate_legacy_logs(store: &SqliteLogStore) -> Result<()> {
1278 let logs_dir = &*crate::env::PITCHFORK_LOGS_DIR;
1279 if !logs_dir.exists() {
1280 return Ok(());
1281 }
1282
1283 let Ok(entries) = std::fs::read_dir(logs_dir) else {
1284 return Ok(());
1285 };
1286
1287 let mut total_migrated = 0u64;
1288 let mut migrated_ids = Vec::new();
1289
1290 for entry in entries.flatten() {
1291 let path = entry.path();
1292 if !path.is_dir() {
1293 continue;
1294 }
1295 let file_name = path
1297 .file_name()
1298 .map_or(String::new(), |n| n.to_string_lossy().to_string());
1299 if file_name == "pitchfork" {
1300 continue;
1301 }
1302
1303 if !file_name.contains("--") {
1305 continue;
1306 }
1307 let log_file = path.join(format!("{file_name}.log"));
1308 if !log_file.exists() {
1309 continue;
1310 }
1311
1312 let daemon_id = match DaemonId::from_safe_path(&file_name) {
1313 Ok(id) => id,
1314 Err(_) => continue,
1315 };
1316
1317 if daemon_id == DaemonId::pitchfork() {
1320 continue;
1321 }
1322
1323 match store.migrate_daemon_text_logs(&daemon_id) {
1324 Ok(0) => {}
1325 Ok(n) => {
1326 total_migrated += n;
1327 migrated_ids.push(daemon_id.qualified());
1328 }
1329 Err(e) => {
1330 warn!(
1331 "failed to migrate text logs for {}: {e}",
1332 daemon_id.qualified()
1333 );
1334 }
1335 }
1336 }
1337
1338 if total_migrated > 0 {
1339 warn!(
1340 "auto-migrated {total_migrated} legacy log entries from {count} daemon(s): {ids}",
1341 count = migrated_ids.len(),
1342 ids = migrated_ids.join(", ")
1343 );
1344 }
1345
1346 Ok(())
1347}
1348
1349#[cfg(test)]
1350mod tests {
1351 use super::text_to_json_literal;
1352 use super::*;
1353 use crate::log_store::LogStore;
1354
1355 #[test]
1356 fn test_json_literal_booleans() {
1357 assert_eq!(text_to_json_literal("true"), "true");
1358 assert_eq!(text_to_json_literal("TRUE"), "true");
1359 assert_eq!(text_to_json_literal("false"), "false");
1360 assert_eq!(text_to_json_literal("False"), "false");
1361 }
1362
1363 #[test]
1364 fn test_json_literal_null() {
1365 assert_eq!(text_to_json_literal("null"), "null");
1366 assert_eq!(text_to_json_literal("NULL"), "null");
1367 }
1368
1369 #[test]
1370 fn test_json_literal_valid_numbers() {
1371 assert_eq!(text_to_json_literal("42"), "42");
1372 assert_eq!(text_to_json_literal("-1"), "-1");
1373 assert_eq!(text_to_json_literal("0"), "0");
1374 assert_eq!(text_to_json_literal("3.14"), "3.14");
1375 assert_eq!(text_to_json_literal("1e10"), "1e10");
1376 assert_eq!(text_to_json_literal("1.5e-3"), "1.5e-3");
1377 }
1378
1379 #[test]
1380 fn test_json_literal_rejects_invalid_json_numbers() {
1381 assert_eq!(text_to_json_literal("+42"), r#""+42""#);
1384 assert_eq!(text_to_json_literal("1."), r#""1.""#);
1385 assert_eq!(text_to_json_literal(".5"), r#"".5""#);
1386 assert_eq!(text_to_json_literal("inf"), r#""inf""#);
1387 assert_eq!(text_to_json_literal("nan"), r#""nan""#);
1388 assert_eq!(text_to_json_literal("infinity"), r#""infinity""#);
1389 }
1390
1391 #[test]
1392 fn test_json_literal_strings() {
1393 assert_eq!(text_to_json_literal("hello"), r#""hello""#);
1394 assert_eq!(text_to_json_literal("req_1"), r#""req_1""#);
1395 assert_eq!(text_to_json_literal(r#"a"b"#), r#""a\"b""#);
1397 }
1398
1399 #[test]
1403 fn write_waits_for_a_contended_lock() {
1404 let dir = tempfile::tempdir().unwrap();
1405 let path = dir.path().join("logs.db");
1406 let store = SqliteLogStore::open(&path).unwrap();
1407
1408 let blocker = Connection::open(&path).unwrap();
1411 blocker.busy_timeout(BUSY_TIMEOUT).unwrap();
1412 blocker.execute_batch("BEGIN IMMEDIATE").unwrap();
1413 let releaser = std::thread::spawn(move || {
1414 std::thread::sleep(std::time::Duration::from_millis(300));
1415 blocker.execute_batch("COMMIT").unwrap();
1416 });
1417
1418 let id = DaemonId::try_new("test", "blocked").unwrap();
1419 let entries = vec![crate::log_parse::parse("held-lock-line", "text")];
1420 store
1421 .append_structured_batch(&id, &entries)
1422 .expect("a contended write must wait for the lock, not fail");
1423
1424 releaser.join().unwrap();
1425 let found = store
1426 .query(&LogQuery {
1427 daemon_ids: vec![id.qualified()],
1428 ..Default::default()
1429 })
1430 .unwrap()
1431 .len();
1432 assert_eq!(found, 1, "the batch written under contention was lost");
1433 }
1434
1435 #[test]
1441 fn concurrent_writers_do_not_lose_batches() {
1442 let dir = tempfile::tempdir().unwrap();
1443 let path = dir.path().join("logs.db");
1444
1445 const WRITERS: usize = 4;
1446 const BATCHES: usize = 15;
1447 const PER_BATCH: usize = 10;
1448
1449 let barrier = std::sync::Barrier::new(WRITERS);
1454
1455 std::thread::scope(|scope| {
1456 for writer in 0..WRITERS {
1457 let path = path.clone();
1458 let barrier = &barrier;
1459 scope.spawn(move || {
1460 barrier.wait();
1463 let store = SqliteLogStore::open(&path).unwrap();
1464 let id = DaemonId::try_new("test", format!("w{writer}")).unwrap();
1465 for batch in 0..BATCHES {
1466 let entries: Vec<ParsedLog> = (0..PER_BATCH)
1467 .map(|i| {
1468 crate::log_parse::parse(&format!("w{writer}-{batch}-{i}"), "text")
1469 })
1470 .collect();
1471 store
1472 .append_structured_batch(&id, &entries)
1473 .expect("concurrent batch write must not fail");
1474 }
1475 });
1476 }
1477 });
1478
1479 let store = SqliteLogStore::open(&path).unwrap();
1480 for writer in 0..WRITERS {
1481 let id = DaemonId::try_new("test", format!("w{writer}")).unwrap();
1482 let found = store
1483 .query(&LogQuery {
1484 daemon_ids: vec![id.qualified()],
1485 ..Default::default()
1486 })
1487 .unwrap()
1488 .len();
1489 assert_eq!(
1490 found,
1491 BATCHES * PER_BATCH,
1492 "writer {writer} lost entries: got {found}"
1493 );
1494 }
1495 }
1496}