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 add_regexp_function(conn: &Connection) -> Result<()> {
22 use std::cell::RefCell;
23
24 let cache: RefCell<lru::LruCache<String, regex::Regex>> =
28 RefCell::new(lru::LruCache::new(std::num::NonZeroUsize::new(32).unwrap()));
29
30 conn.create_scalar_function(
31 "regexp",
32 2,
33 rusqlite::functions::FunctionFlags::SQLITE_UTF8
34 | rusqlite::functions::FunctionFlags::SQLITE_DETERMINISTIC,
35 move |ctx| {
36 let pattern: String = ctx.get(0)?;
37 let text: String = ctx.get(1)?;
38
39 let mut cache = cache.borrow_mut();
40 let re = match cache.get(&pattern) {
41 Some(re) => re.clone(),
42 None => {
43 let re = regex::Regex::new(&pattern)
44 .map_err(|e| rusqlite::Error::UserFunctionError(e.to_string().into()))?;
45 cache.put(pattern.clone(), re.clone());
46 re
47 }
48 };
49 Ok(re.is_match(&text))
50 },
51 )
52 .into_diagnostic()
53}
54
55pub struct SqliteLogStore {
57 conn: Mutex<Connection>,
58 path: PathBuf,
59}
60
61const PARALLEL_QUERY_THRESHOLD: usize = 200_000;
68
69impl SqliteLogStore {
70 pub fn open(path: impl Into<PathBuf>) -> Result<Self> {
72 let path = path.into();
73 if let Some(parent) = path.parent() {
74 std::fs::create_dir_all(parent).into_diagnostic()?;
75 }
76 let conn = Connection::open(&path).into_diagnostic()?;
77 add_regexp_function(&conn)?;
78 conn.execute_batch(
79 "PRAGMA journal_mode = WAL;
80 PRAGMA synchronous = NORMAL;
81 PRAGMA mmap_size = 268435456;",
82 )
83 .into_diagnostic()?;
84 conn.execute(
85 "CREATE TABLE IF NOT EXISTS log_entries (
86 id INTEGER PRIMARY KEY AUTOINCREMENT,
87 daemon_id TEXT NOT NULL,
88 timestamp INTEGER NOT NULL,
89 message TEXT NOT NULL,
90 level TEXT,
91 msg TEXT,
92 logger TEXT,
93 fields_json TEXT
94 );",
95 [],
96 )
97 .into_diagnostic()?;
98 conn.execute(
99 "CREATE INDEX IF NOT EXISTS idx_daemon_ts ON log_entries(daemon_id, timestamp);",
100 [],
101 )
102 .into_diagnostic()?;
103 conn.execute(
104 "CREATE INDEX IF NOT EXISTS idx_daemon_id ON log_entries(daemon_id, id);",
105 [],
106 )
107 .into_diagnostic()?;
108 conn.execute(
109 "CREATE INDEX IF NOT EXISTS idx_timestamp ON log_entries(timestamp);",
110 [],
111 )
112 .into_diagnostic()?;
113
114 let existing_cols: Vec<String> = {
117 let mut stmt = conn
118 .prepare("PRAGMA table_info(log_entries)")
119 .into_diagnostic()?;
120 let rows = stmt
121 .query_map([], |row| row.get::<_, String>(1))
122 .into_diagnostic()?;
123 rows.filter_map(|r| r.ok()).collect()
124 };
125 for col in ["level", "msg", "logger", "fields_json"] {
126 if !existing_cols.iter().any(|c| c == col) {
127 conn.execute(
128 &format!("ALTER TABLE log_entries ADD COLUMN {col} TEXT"),
129 [],
130 )
131 .into_diagnostic()?;
132 }
133 }
134 conn.execute(
135 "CREATE INDEX IF NOT EXISTS idx_daemon_level_ts ON log_entries(daemon_id, level, timestamp);",
136 [],
137 )
138 .into_diagnostic()?;
139
140 conn.execute(
141 "CREATE TABLE IF NOT EXISTS log_clear_generations (
142 daemon_id TEXT PRIMARY KEY,
143 generation INTEGER NOT NULL DEFAULT 0
144 );",
145 [],
146 )
147 .into_diagnostic()?;
148 Ok(Self {
149 conn: Mutex::new(conn),
150 path,
151 })
152 }
153
154 fn row_to_entry(row: &rusqlite::Row) -> rusqlite::Result<LogEntry> {
155 let id: i64 = row.get(0)?;
156 let daemon_id: String = row.get(1)?;
157 let ts_millis: i64 = row.get(2)?;
158 let message: String = row.get(3)?;
159 let level: Option<String> = row.get(4)?;
160 let msg: Option<String> = row.get(5)?;
161 let logger: Option<String> = row.get(6)?;
162 let fields_json: Option<String> = row.get(7)?;
163 let timestamp = Local
164 .timestamp_millis_opt(ts_millis)
165 .single()
166 .unwrap_or_else(Local::now);
167 Ok(LogEntry {
168 id,
169 daemon_id,
170 timestamp,
171 message,
172 level,
173 msg,
174 logger,
175 fields_json,
176 })
177 }
178
179 fn archive_entries(
180 &self,
181 entries: &[LogEntry],
182 archive_hook: &ArchiveHook,
183 daemon_id: &DaemonId,
184 reason: &str,
185 ) -> Result<()> {
186 use std::process::{Command, Stdio};
187
188 if entries.is_empty() {
189 return Ok(());
190 }
191
192 for chunk in entries.chunks(archive_hook.batch_size.max(1)) {
193 let mut child = Command::new("sh")
194 .arg("-c")
195 .arg(&archive_hook.command)
196 .stdin(Stdio::piped())
197 .stdout(Stdio::null())
198 .stderr(Stdio::piped())
199 .env("PITCHFORK_DAEMON_ID", daemon_id.qualified())
200 .env("PITCHFORK_ARCHIVE_REASON", reason)
201 .spawn()
202 .into_diagnostic()
203 .map_err(|e| miette::miette!("failed to spawn archive hook: {e}"))?;
204
205 let write_result = {
209 let stdin = child.stdin.take().expect("piped stdin should be available");
210 let mut stdin = std::io::BufWriter::new(stdin);
211 let mut result = Ok(());
212 for entry in chunk {
213 let line = serde_json::json!({
214 "id": entry.id,
215 "daemon_id": entry.daemon_id,
216 "timestamp": entry.timestamp.to_rfc3339(),
217 "message": entry.message,
218 });
219 if let Err(e) = writeln!(stdin, "{}", line) {
220 result = Err(miette::miette!(
221 "failed to write to archive hook stdin: {e}"
222 ));
223 break;
224 }
225 }
226 if result.is_ok() {
234 if let Err(e) = stdin.flush() {
235 result = Err(miette::miette!("failed to flush archive hook stdin: {e}"));
236 }
237 }
238 result
239 };
241
242 if let Err(e) = write_result {
243 let _ = child.kill();
244 let _ = child.wait();
245 return Err(e);
246 }
247
248 let output = child.wait_with_output().into_diagnostic()?;
249 if !output.status.success() {
250 let stderr = String::from_utf8_lossy(&output.stderr);
251 return Err(miette::miette!(
252 "archive hook failed with status {}: {stderr}",
253 output.status
254 ));
255 }
256 }
257
258 Ok(())
259 }
260
261 fn delete_by_ids(&self, ids: &[i64]) -> Result<u64> {
266 const SQLITE_MAX_VARS: usize = 999;
267
268 let mut total = 0u64;
269 let conn = self.conn.lock().unwrap();
270 for chunk in ids.chunks(SQLITE_MAX_VARS) {
271 if chunk.is_empty() {
272 continue;
273 }
274 let placeholders: Vec<String> = (1..=chunk.len()).map(|i| format!("?{i}")).collect();
275 let sql = format!(
276 "DELETE FROM log_entries WHERE id IN ({})",
277 placeholders.join(", ")
278 );
279 total += conn
280 .execute(&sql, rusqlite::params_from_iter(chunk.iter()))
281 .into_diagnostic()? as u64;
282 }
283 Ok(total)
284 }
285
286 pub fn rotate_by_age(
293 &self,
294 daemon_id: &DaemonId,
295 max_age: chrono::Duration,
296 archive_hook: Option<&ArchiveHook>,
297 ) -> Result<u64> {
298 let cutoff = (Local::now() - max_age).timestamp_millis();
299 let hook = archive_hook.filter(|h| h.is_enabled());
300
301 if let Some(hook) = hook {
302 let mut total_deleted = 0u64;
303 loop {
304 let entries: Vec<LogEntry> = {
306 let conn = self.conn.lock().unwrap();
307 let mut stmt = conn
308 .prepare(
309 "SELECT id, daemon_id, timestamp, message, level, msg, logger, fields_json FROM log_entries
310 WHERE daemon_id = ?1 AND timestamp < ?2
311 ORDER BY timestamp ASC, id ASC
312 LIMIT ?3",
313 )
314 .into_diagnostic()?;
315 stmt.query_map(
316 params![daemon_id.qualified(), cutoff, hook.batch_size as i64],
317 Self::row_to_entry,
318 )
319 .into_diagnostic()?
320 .collect::<rusqlite::Result<Vec<_>>>()
321 .into_diagnostic()?
322 };
323
324 if entries.is_empty() {
325 break;
326 }
327
328 let batch_ids: Vec<i64> = entries.iter().map(|e| e.id).collect();
329
330 self.archive_entries(&entries, hook, daemon_id, "age")?;
332
333 let deleted = self.delete_by_ids(&batch_ids)?;
335 total_deleted += deleted;
336 }
337 Ok(total_deleted)
338 } else {
339 let conn = self.conn.lock().unwrap();
340 let rows = conn
341 .execute(
342 "DELETE FROM log_entries WHERE daemon_id = ?1 AND timestamp < ?2",
343 params![daemon_id.qualified(), cutoff],
344 )
345 .into_diagnostic()?;
346 Ok(rows as u64)
347 }
348 }
349
350 pub fn rotate_by_count(
358 &self,
359 daemon_id: &DaemonId,
360 max_count: u64,
361 archive_hook: Option<&ArchiveHook>,
362 ) -> Result<u64> {
363 let hook = archive_hook.filter(|h| h.is_enabled());
364
365 let to_delete: i64 = {
367 let conn = self.conn.lock().unwrap();
368 let count: i64 = conn
369 .query_row(
370 "SELECT COUNT(*) FROM log_entries WHERE daemon_id = ?1",
371 [daemon_id.qualified()],
372 |row| row.get(0),
373 )
374 .into_diagnostic()?;
375 count.saturating_sub(max_count as i64)
376 };
377
378 if to_delete <= 0 {
379 return Ok(0);
380 }
381
382 if let Some(hook) = hook {
383 let mut total_deleted = 0u64;
384 let mut remaining = to_delete;
385 loop {
386 let batch_len = remaining.min(hook.batch_size as i64);
387
388 let entries: Vec<LogEntry> = {
390 let conn = self.conn.lock().unwrap();
391 let mut stmt = conn
392 .prepare(
393 "SELECT id, daemon_id, timestamp, message, level, msg, logger, fields_json FROM log_entries
394 WHERE daemon_id = ?1
395 ORDER BY timestamp ASC, id ASC
396 LIMIT ?2",
397 )
398 .into_diagnostic()?;
399 stmt.query_map(
400 params![daemon_id.qualified(), batch_len],
401 Self::row_to_entry,
402 )
403 .into_diagnostic()?
404 .collect::<rusqlite::Result<Vec<_>>>()
405 .into_diagnostic()?
406 };
407
408 if entries.is_empty() {
409 break;
410 }
411
412 let fetched = entries.len() as i64;
413 let batch_ids: Vec<i64> = entries.iter().map(|e| e.id).collect();
414
415 self.archive_entries(&entries, hook, daemon_id, "count")?;
417
418 let deleted = self.delete_by_ids(&batch_ids)?;
420 total_deleted += deleted;
421 remaining -= fetched;
422 }
423 Ok(total_deleted)
424 } else {
425 let conn = self.conn.lock().unwrap();
426 let rows = conn
427 .execute(
428 "DELETE FROM log_entries WHERE id IN (
429 SELECT id FROM log_entries WHERE daemon_id = ?1
430 ORDER BY timestamp ASC, id ASC LIMIT ?2
431 )",
432 params![daemon_id.qualified(), to_delete],
433 )
434 .into_diagnostic()?;
435 Ok(rows as u64)
436 }
437 }
438
439 pub fn migrate_daemon_text_logs(&self, daemon_id: &DaemonId) -> Result<u64> {
444 let text_path = daemon_id.log_path();
445 if !text_path.exists() {
446 return Ok(0);
447 }
448
449 let file = std::fs::File::open(&text_path).into_diagnostic()?;
450 let reader = BufReader::new(file);
451 let re = regex::Regex::new(r"^(\d{4}-\d{2}-\d{2} \d{2}:\d{2}:\d{2}) ([\w./-]+) (.*)$")
452 .expect("invalid regex");
453
454 let mut current_timestamp: Option<DateTime<Local>> = None;
455 let mut current_message = String::new();
456 let mut entries = Vec::with_capacity(1000);
457 let mut total_migrated: u64 = 0;
458
459 for line in reader.lines() {
460 let line = line.into_diagnostic()?;
461 if let Some(caps) = re.captures(&line) {
462 if let Some(ts) = current_timestamp.take() {
463 entries.push((ts, std::mem::take(&mut current_message)));
464 }
465 let ts_str = caps.get(1).map(|m| m.as_str()).unwrap_or_default();
466 let msg = caps.get(3).map(|m| m.as_str()).unwrap_or_default();
467 if let Ok(naive) =
468 chrono::NaiveDateTime::parse_from_str(ts_str, "%Y-%m-%d %H:%M:%S")
469 {
470 current_timestamp = Local.from_local_datetime(&naive).single();
471 current_message = msg.to_string();
472 }
473 } else if current_timestamp.is_some() {
474 current_message.push('\n');
475 current_message.push_str(&line);
476 }
477
478 if entries.len() >= 1000 {
479 total_migrated += self.insert_batch(daemon_id, &entries)?;
480 entries.clear();
481 }
482 }
483
484 if let Some(ts) = current_timestamp {
485 entries.push((ts, std::mem::take(&mut current_message)));
486 }
487
488 if !entries.is_empty() {
489 total_migrated += self.insert_batch(daemon_id, &entries)?;
490 }
491
492 if total_migrated > 0 {
493 if let Err(e) = std::fs::remove_file(&text_path) {
494 log::warn!(
495 "failed to remove legacy log file after migration {}: {e}",
496 text_path.display()
497 );
498 }
499 }
500
501 Ok(total_migrated)
502 }
503
504 fn insert_batch(
505 &self,
506 daemon_id: &DaemonId,
507 entries: &[(DateTime<Local>, String)],
508 ) -> Result<u64> {
509 let mut conn = self.conn.lock().unwrap();
510 let tx = conn.transaction().into_diagnostic()?;
511 let mut count = 0u64;
512 {
513 let mut stmt = tx
514 .prepare(
515 "INSERT INTO log_entries (daemon_id, timestamp, message) VALUES (?1, ?2, ?3)",
516 )
517 .into_diagnostic()?;
518 for (ts, msg) in entries {
519 stmt.execute(params![daemon_id.qualified(), ts.timestamp_millis(), msg])
520 .into_diagnostic()?;
521 count += 1;
522 }
523 }
524 tx.commit().into_diagnostic()?;
525 Ok(count)
526 }
527
528 fn build_query_sql(
533 opts: &LogQuery,
534 id_range: Option<(i64, i64)>,
535 ) -> (String, Vec<Box<dyn rusqlite::ToSql>>) {
536 let mut conditions = Vec::new();
537 let mut query_params: Vec<Box<dyn rusqlite::ToSql>> = Vec::new();
538
539 if !opts.daemon_ids.is_empty() {
540 let placeholders: Vec<String> = (1..=opts.daemon_ids.len())
541 .map(|i| format!("?{}", i))
542 .collect();
543 conditions.push(format!("daemon_id IN ({})", placeholders.join(", ")));
544 for id in &opts.daemon_ids {
545 query_params.push(Box::new(id.clone()));
546 }
547 }
548
549 if let Some(from) = opts.from {
550 conditions.push(format!("timestamp >= ?{}", query_params.len() + 1));
551 query_params.push(Box::new(from.timestamp_millis()));
552 }
553
554 if let Some(to) = opts.to {
555 conditions.push(format!("timestamp <= ?{}", query_params.len() + 1));
556 query_params.push(Box::new(to.timestamp_millis()));
557 }
558
559 if let Some(after_id) = opts.after_id {
560 conditions.push(format!("id > ?{}", query_params.len() + 1));
561 query_params.push(Box::new(after_id));
562 }
563
564 if let Some((start, end)) = id_range {
565 conditions.push(format!("id > ?{}", query_params.len() + 1));
566 query_params.push(Box::new(start));
567 conditions.push(format!("id <= ?{}", query_params.len() + 1));
568 query_params.push(Box::new(end));
569 }
570
571 let mut message_conditions = Vec::new();
572 for filter in &opts.message_filters {
573 match filter {
574 MessageFilter::Contains {
575 pattern,
576 case_sensitive,
577 } => {
578 let param_index = query_params.len() + 1;
579 if *case_sensitive {
580 message_conditions
581 .push(format!("INSTR(message, ?{idx}) > 0", idx = param_index));
582 query_params.push(Box::new(pattern.clone()));
583 } else {
584 let escaped = escape_like_pattern(pattern);
585 let param = format!("%{}%", escaped);
586 message_conditions.push(format!(
587 "LOWER(message) LIKE LOWER(?{idx}) ESCAPE '\\'",
588 idx = param_index
589 ));
590 query_params.push(Box::new(param));
591 }
592 }
593 MessageFilter::Regex { pattern } => {
594 let param_index = query_params.len() + 1;
595 message_conditions.push(format!("message REGEXP ?{param_index}"));
596 query_params.push(Box::new(pattern.clone()));
597 }
598 }
599 }
600 if !message_conditions.is_empty() {
601 conditions.push(format!("({})", message_conditions.join(" OR ")));
602 }
603
604 for filter in &opts.field_filters {
605 match filter {
606 FieldFilter::LevelMin(level) => {
607 let matching = crate::log_store::levels_at_or_above(level);
608 if matching.is_empty() {
609 conditions.push("0".to_string());
611 } else {
612 let placeholders = matching
613 .iter()
614 .map(|l| {
615 let idx = query_params.len() + 1;
616 query_params.push(Box::new((*l).to_string()));
617 format!("?{idx}")
618 })
619 .collect::<Vec<_>>()
620 .join(", ");
621 conditions.push(format!("level IN ({placeholders})"));
622 }
623 }
624 FieldFilter::FieldEq { key, value } => {
625 let param_index = query_params.len() + 1;
626 conditions.push(format!(
627 "json_extract(fields_json, '$.{key}') = ?{param_index}"
628 ));
629 query_params.push(Box::new(value.clone()));
630 }
631 }
632 }
633
634 let where_clause = if conditions.is_empty() {
635 String::new()
636 } else {
637 format!("WHERE {}", conditions.join(" AND "))
638 };
639
640 let order = if opts.order_desc { "DESC" } else { "ASC" };
641
642 let limit_clause = opts
643 .limit
644 .map(|n| format!("LIMIT {}", n))
645 .unwrap_or_default();
646
647 let columns = if opts.include_structured {
648 "id, daemon_id, timestamp, message, level, msg, logger, fields_json"
649 } else {
650 "id, daemon_id, timestamp, message, NULL, NULL, NULL, NULL"
651 };
652
653 let sql = format!(
654 "SELECT {columns} FROM log_entries {} ORDER BY timestamp {}, id {} {}",
655 where_clause, order, order, limit_clause
656 );
657
658 (sql, query_params)
659 }
660
661 fn should_parallelize(opts: &LogQuery) -> bool {
663 if opts.daemon_ids.len() != 1 {
664 return false;
665 }
666 if opts.after_id.is_some() {
669 return false;
670 }
671 let limit = opts.limit.unwrap_or(usize::MAX);
672 if limit < PARALLEL_QUERY_THRESHOLD {
673 return false;
674 }
675 std::thread::available_parallelism()
676 .map(|n| n.get() >= 2)
677 .unwrap_or(false)
678 }
679
680 fn query_parallel(&self, opts: &LogQuery) -> Result<Vec<LogEntry>> {
687 let max_threads = 2;
692 let num_threads = std::thread::available_parallelism()
693 .map(|n| n.get().min(max_threads))
694 .unwrap_or(1);
695
696 let max_id: Option<i64> = {
697 let conn = self.conn.lock().unwrap();
698 conn.query_row(
699 "SELECT MAX(id) FROM log_entries WHERE daemon_id = ?1",
700 params![&opts.daemon_ids[0]],
701 |row| row.get(0),
702 )
703 .ok()
704 };
705
706 let Some(max_id) = max_id else {
707 return Ok(Vec::new());
708 };
709 if max_id == 0 {
710 return Ok(Vec::new());
711 }
712
713 let shard_size = (max_id as usize).div_ceil(num_threads);
714 let path = self.path.clone();
715 let opts = opts.clone();
716 let needs_regexp = opts
717 .message_filters
718 .iter()
719 .any(|f| matches!(f, MessageFilter::Regex { .. }));
720
721 let shards: Vec<Result<Vec<LogEntry>>> = std::thread::scope(|s| {
722 (0..num_threads)
723 .map(|i| {
724 let start = (i * shard_size) as i64;
725 let end = if i == num_threads - 1 {
726 max_id
727 } else {
728 ((i + 1) * shard_size) as i64
729 };
730 let opts = opts.clone();
731 let path = path.clone();
732 s.spawn(move || -> Result<Vec<LogEntry>> {
733 let conn = Connection::open_with_flags(
734 &path,
735 OpenFlags::SQLITE_OPEN_READ_ONLY | OpenFlags::SQLITE_OPEN_NO_MUTEX,
736 )
737 .into_diagnostic()?;
738 conn.execute_batch(
739 "PRAGMA mmap_size = 268435456;
740 PRAGMA query_only = 1;",
741 )
742 .into_diagnostic()?;
743 if needs_regexp {
744 add_regexp_function(&conn)?;
745 }
746
747 let (sql, query_params) = Self::build_query_sql(&opts, Some((start, end)));
748 Self::execute_built_query(&conn, &sql, &query_params)
749 })
750 })
751 .map(|h| h.join().unwrap())
752 .collect()
753 });
754
755 let mut merged = Vec::new();
756 if opts.order_desc {
757 for shard in shards.into_iter().rev() {
758 merged.extend(shard?);
759 }
760 } else {
761 for shard in shards {
762 merged.extend(shard?);
763 }
764 }
765
766 if let Some(limit) = opts.limit {
767 if merged.len() > limit {
768 merged.truncate(limit);
769 }
770 }
771
772 Ok(merged)
773 }
774
775 fn execute_built_query(
777 conn: &Connection,
778 sql: &str,
779 query_params: &[Box<dyn rusqlite::ToSql>],
780 ) -> Result<Vec<LogEntry>> {
781 let mut stmt = conn.prepare(sql).into_diagnostic()?;
782 let params_ref: Vec<&dyn rusqlite::ToSql> =
783 query_params.iter().map(|p| p.as_ref()).collect();
784 let rows = stmt
785 .query_map(params_ref.as_slice(), Self::row_to_entry)
786 .into_diagnostic()?;
787 let mut entries = Vec::new();
788 for row in rows {
789 entries.push(row.into_diagnostic()?);
790 }
791 Ok(entries)
792 }
793}
794
795impl LogStore for SqliteLogStore {
796 fn append(&self, daemon_id: &DaemonId, message: &str) -> Result<()> {
797 let ts = Local::now().timestamp_millis();
798 let id = daemon_id.qualified();
799 let msg = message.to_string();
800
801 let conn = self.conn.lock().unwrap();
802 let _ = conn
803 .execute(
804 "INSERT INTO log_entries (daemon_id, timestamp, message) VALUES (?1, ?2, ?3)",
805 params![id, ts, msg],
806 )
807 .into_diagnostic()?;
808 Ok(())
809 }
810
811 fn append_batch(&self, daemon_id: &DaemonId, messages: &[String]) -> Result<()> {
812 if messages.is_empty() {
813 return Ok(());
814 }
815 let base_ts = Local::now().timestamp_millis();
816 let id = daemon_id.qualified();
817
818 let mut conn = self.conn.lock().unwrap();
819 let tx = conn.transaction().into_diagnostic()?;
820 {
821 let mut stmt = tx
822 .prepare(
823 "INSERT INTO log_entries (daemon_id, timestamp, message) VALUES (?1, ?2, ?3)",
824 )
825 .into_diagnostic()?;
826 for (idx, msg) in messages.iter().enumerate() {
827 let ts = base_ts + idx as i64;
831 stmt.execute(params![id, ts, msg]).into_diagnostic()?;
832 }
833 }
834 tx.commit().into_diagnostic()?;
835 Ok(())
836 }
837
838 fn append_structured(&self, daemon_id: &DaemonId, parsed: &ParsedLog) -> Result<()> {
839 let ts = Local::now().timestamp_millis();
840 let id = daemon_id.qualified();
841
842 let conn = self.conn.lock().unwrap();
843 let _ = conn
844 .execute(
845 "INSERT INTO log_entries (daemon_id, timestamp, message, level, msg, logger, fields_json) \
846 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)",
847 params![id, ts, parsed.message, parsed.level, parsed.msg, parsed.logger, parsed.fields_json],
848 )
849 .into_diagnostic()?;
850 Ok(())
851 }
852
853 fn append_structured_batch(&self, daemon_id: &DaemonId, entries: &[ParsedLog]) -> Result<()> {
854 if entries.is_empty() {
855 return Ok(());
856 }
857 let base_ts = Local::now().timestamp_millis();
858 let id = daemon_id.qualified();
859
860 let mut conn = self.conn.lock().unwrap();
861 let tx = conn.transaction().into_diagnostic()?;
862 {
863 let mut stmt = tx
864 .prepare(
865 "INSERT INTO log_entries (daemon_id, timestamp, message, level, msg, logger, fields_json) \
866 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)",
867 )
868 .into_diagnostic()?;
869 for (idx, entry) in entries.iter().enumerate() {
870 let ts = base_ts + idx as i64;
871 stmt.execute(params![
872 id,
873 ts,
874 entry.message,
875 entry.level,
876 entry.msg,
877 entry.logger,
878 entry.fields_json
879 ])
880 .into_diagnostic()?;
881 }
882 }
883 tx.commit().into_diagnostic()?;
884 Ok(())
885 }
886
887 fn query(&self, opts: &LogQuery) -> Result<Vec<LogEntry>> {
888 if Self::should_parallelize(opts) && self.path.as_os_str() != ":memory:" {
890 if let Ok(entries) = self.query_parallel(opts) {
891 return Ok(entries);
892 }
893 }
895
896 let conn = self.conn.lock().unwrap();
898 let (sql, query_params) = Self::build_query_sql(opts, None);
899 Self::execute_built_query(&conn, &sql, &query_params)
900 }
901
902 fn tail(&self, daemon_id: &DaemonId, after_id: Option<i64>) -> Result<Vec<LogEntry>> {
903 self.query(&LogQuery {
904 daemon_ids: vec![daemon_id.qualified()],
905 from: None,
906 to: None,
907 limit: None,
908 order_desc: false,
909 after_id,
910 message_filters: Vec::new(),
911 field_filters: Vec::new(),
912 include_structured: false,
913 })
914 }
915
916 fn clear(&self, daemon_ids: &[DaemonId]) -> Result<()> {
917 let mut conn = self.conn.lock().unwrap();
918 let tx = conn.transaction().into_diagnostic()?;
919 for id in daemon_ids {
920 tx.execute(
921 "DELETE FROM log_entries WHERE daemon_id = ?1",
922 params![id.qualified()],
923 )
924 .into_diagnostic()?;
925 tx.execute(
926 "INSERT INTO log_clear_generations (daemon_id, generation)
927 VALUES (?1, 1)
928 ON CONFLICT(daemon_id) DO UPDATE SET generation = generation + 1",
929 params![id.qualified()],
930 )
931 .into_diagnostic()?;
932 }
933 tx.commit().into_diagnostic()?;
934 Ok(())
935 }
936
937 fn last_id(&self, daemon_id: &DaemonId) -> Result<Option<i64>> {
938 let conn = self.conn.lock().unwrap();
939 let id: Option<i64> = conn
942 .query_row(
943 "SELECT MAX(id) FROM log_entries WHERE daemon_id = ?1",
944 params![daemon_id.qualified()],
945 |row| row.get(0),
946 )
947 .into_diagnostic()?;
948 Ok(id)
949 }
950
951 fn list_daemon_ids(&self) -> Result<Vec<String>> {
952 let conn = self.conn.lock().unwrap();
953 let mut stmt = conn
954 .prepare("SELECT DISTINCT daemon_id FROM log_entries")
955 .into_diagnostic()?;
956 let ids = stmt
957 .query_map([], |row| {
958 let id: String = row.get(0)?;
959 Ok(id)
960 })
961 .into_diagnostic()?
962 .filter_map(|r| r.ok())
963 .collect();
964 Ok(ids)
965 }
966
967 fn apply_retention(
968 &self,
969 policy: &super::RetentionPolicy,
970 excluded_daemon_ids: &[DaemonId],
971 archive_hook: Option<&ArchiveHook>,
972 ) -> Result<u64> {
973 let daemon_ids = self.list_daemon_ids()?;
974 let excluded: HashSet<String> = excluded_daemon_ids.iter().map(|d| d.qualified()).collect();
975 let mut total = 0u64;
976 for id_str in daemon_ids {
977 if excluded.contains(&id_str) {
978 continue;
979 }
980 let id = DaemonId::parse(&id_str).unwrap_or_else(|_| {
981 DaemonId::try_new("global", &id_str).unwrap_or_else(|_| DaemonId::pitchfork())
982 });
983 if let Some(dur) = policy.age {
984 total += self.rotate_by_age(&id, dur, archive_hook)?;
985 }
986 if let Some(n) = policy.count {
987 total += self.rotate_by_count(&id, n, archive_hook)?;
988 }
989 }
990 Ok(total)
991 }
992
993 fn apply_retention_for_daemon(
994 &self,
995 daemon_id: &DaemonId,
996 policy: &super::RetentionPolicy,
997 archive_hook: Option<&ArchiveHook>,
998 ) -> Result<u64> {
999 let mut total = 0u64;
1000 if let Some(dur) = policy.age {
1001 total += self.rotate_by_age(daemon_id, dur, archive_hook)?;
1002 }
1003 if let Some(n) = policy.count {
1004 total += self.rotate_by_count(daemon_id, n, archive_hook)?;
1005 }
1006 Ok(total)
1007 }
1008
1009 fn last_clear_generation(&self, daemon_id: &DaemonId) -> Result<Option<u64>> {
1010 let conn = self.conn.lock().unwrap();
1011 let mut stmt = conn
1012 .prepare("SELECT generation FROM log_clear_generations WHERE daemon_id = ?1")
1013 .into_diagnostic()?;
1014 let generation: Option<i64> = stmt
1015 .query_row(params![daemon_id.qualified()], |row| row.get(0))
1016 .optional()
1017 .into_diagnostic()?;
1018 generation
1019 .map(|generation| {
1020 u64::try_from(generation)
1021 .map_err(|_| miette::miette!("log clear generation cannot be negative"))
1022 })
1023 .transpose()
1024 }
1025}
1026
1027use once_cell::sync::Lazy;
1029use std::sync::Arc;
1030
1031pub static LOG_STORE: Lazy<Arc<SqliteLogStore>> = Lazy::new(|| {
1032 let path = crate::env::PITCHFORK_LOGS_DIR.join("logs.db");
1033 let mut is_fallback = false;
1034 let store = Arc::new(SqliteLogStore::open(&path).unwrap_or_else(|e| {
1035 error!(
1036 "failed to open log store at {}: {e}. Falling back to in-memory store; logs will not persist across restarts.",
1037 path.display()
1038 );
1039 is_fallback = true;
1040 SqliteLogStore::open(":memory:").expect("in-memory SQLite should always open")
1041 }));
1042
1043 if !is_fallback {
1047 if let Err(e) = auto_migrate_legacy_logs(&store) {
1048 warn!("legacy log auto-migration failed: {e}");
1049 }
1050 } else {
1051 warn!(
1052 "skipping legacy log auto-migration because log store is in-memory (no durable destination)"
1053 );
1054 }
1055
1056 store
1057});
1058
1059fn auto_migrate_legacy_logs(store: &SqliteLogStore) -> Result<()> {
1067 let logs_dir = &*crate::env::PITCHFORK_LOGS_DIR;
1068 if !logs_dir.exists() {
1069 return Ok(());
1070 }
1071
1072 let Ok(entries) = std::fs::read_dir(logs_dir) else {
1073 return Ok(());
1074 };
1075
1076 let mut total_migrated = 0u64;
1077 let mut migrated_ids = Vec::new();
1078
1079 for entry in entries.flatten() {
1080 let path = entry.path();
1081 if !path.is_dir() {
1082 continue;
1083 }
1084 let file_name = path
1086 .file_name()
1087 .map_or(String::new(), |n| n.to_string_lossy().to_string());
1088 if file_name == "pitchfork" {
1089 continue;
1090 }
1091
1092 if !file_name.contains("--") {
1094 continue;
1095 }
1096 let log_file = path.join(format!("{file_name}.log"));
1097 if !log_file.exists() {
1098 continue;
1099 }
1100
1101 let daemon_id = match DaemonId::from_safe_path(&file_name) {
1102 Ok(id) => id,
1103 Err(_) => continue,
1104 };
1105
1106 if daemon_id == DaemonId::pitchfork() {
1109 continue;
1110 }
1111
1112 match store.migrate_daemon_text_logs(&daemon_id) {
1113 Ok(0) => {}
1114 Ok(n) => {
1115 total_migrated += n;
1116 migrated_ids.push(daemon_id.qualified());
1117 }
1118 Err(e) => {
1119 warn!(
1120 "failed to migrate text logs for {}: {e}",
1121 daemon_id.qualified()
1122 );
1123 }
1124 }
1125 }
1126
1127 if total_migrated > 0 {
1128 warn!(
1129 "auto-migrated {total_migrated} legacy log entries from {count} daemon(s): {ids}",
1130 count = migrated_ids.len(),
1131 ids = migrated_ids.join(", ")
1132 );
1133 }
1134
1135 Ok(())
1136}