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
69const BUSY_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(5);
73
74fn enable_wal(conn: &Connection) {
85 const ATTEMPTS: usize = 5;
86 for attempt in 0..ATTEMPTS {
87 match conn.query_row("PRAGMA journal_mode", [], |row| row.get::<_, String>(0)) {
88 Ok(mode) if mode.eq_ignore_ascii_case("wal") => return,
89 Ok(_) => {}
90 Err(e) => {
91 debug!("could not read journal_mode: {e}");
92 return;
93 }
94 }
95 if conn.execute_batch("PRAGMA journal_mode = WAL;").is_ok() {
96 return;
97 }
98 if attempt + 1 < ATTEMPTS {
99 std::thread::sleep(std::time::Duration::from_millis(20));
100 }
101 }
102 debug!("log store is not in WAL mode; another connection may be switching it");
103}
104
105impl SqliteLogStore {
106 pub fn open(path: impl Into<PathBuf>) -> Result<Self> {
108 let path = path.into();
109 if let Some(parent) = path.parent() {
110 std::fs::create_dir_all(parent).into_diagnostic()?;
111 }
112 let conn = Connection::open(&path).into_diagnostic()?;
113
114 conn.busy_timeout(BUSY_TIMEOUT).into_diagnostic()?;
122
123 add_regexp_function(&conn)?;
124
125 enable_wal(&conn);
126
127 conn.execute_batch(
128 "PRAGMA synchronous = NORMAL;
129 PRAGMA mmap_size = 268435456;",
130 )
131 .into_diagnostic()?;
132 conn.execute(
133 "CREATE TABLE IF NOT EXISTS log_entries (
134 id INTEGER PRIMARY KEY AUTOINCREMENT,
135 daemon_id TEXT NOT NULL,
136 timestamp INTEGER NOT NULL,
137 message TEXT NOT NULL,
138 level TEXT,
139 msg TEXT,
140 logger TEXT,
141 fields_json TEXT
142 );",
143 [],
144 )
145 .into_diagnostic()?;
146 conn.execute(
147 "CREATE INDEX IF NOT EXISTS idx_daemon_ts ON log_entries(daemon_id, timestamp);",
148 [],
149 )
150 .into_diagnostic()?;
151 conn.execute(
152 "CREATE INDEX IF NOT EXISTS idx_daemon_id ON log_entries(daemon_id, id);",
153 [],
154 )
155 .into_diagnostic()?;
156 conn.execute(
157 "CREATE INDEX IF NOT EXISTS idx_timestamp ON log_entries(timestamp);",
158 [],
159 )
160 .into_diagnostic()?;
161
162 let existing_cols: Vec<String> = {
165 let mut stmt = conn
166 .prepare("PRAGMA table_info(log_entries)")
167 .into_diagnostic()?;
168 let rows = stmt
169 .query_map([], |row| row.get::<_, String>(1))
170 .into_diagnostic()?;
171 rows.filter_map(|r| r.ok()).collect()
172 };
173 for col in ["level", "msg", "logger", "fields_json"] {
174 if !existing_cols.iter().any(|c| c == col) {
175 conn.execute(
176 &format!("ALTER TABLE log_entries ADD COLUMN {col} TEXT"),
177 [],
178 )
179 .into_diagnostic()?;
180 }
181 }
182 conn.execute(
183 "CREATE INDEX IF NOT EXISTS idx_daemon_level_ts ON log_entries(daemon_id, level, timestamp);",
184 [],
185 )
186 .into_diagnostic()?;
187
188 conn.execute(
189 "CREATE TABLE IF NOT EXISTS log_clear_generations (
190 daemon_id TEXT PRIMARY KEY,
191 generation INTEGER NOT NULL DEFAULT 0
192 );",
193 [],
194 )
195 .into_diagnostic()?;
196 Ok(Self {
197 conn: Mutex::new(conn),
198 path,
199 })
200 }
201
202 fn row_to_entry(row: &rusqlite::Row) -> rusqlite::Result<LogEntry> {
203 let id: i64 = row.get(0)?;
204 let daemon_id: String = row.get(1)?;
205 let ts_millis: i64 = row.get(2)?;
206 let message: String = row.get(3)?;
207 let level: Option<String> = row.get(4)?;
208 let msg: Option<String> = row.get(5)?;
209 let logger: Option<String> = row.get(6)?;
210 let fields_json: Option<String> = row.get(7)?;
211 let timestamp = Local
212 .timestamp_millis_opt(ts_millis)
213 .single()
214 .unwrap_or_else(Local::now);
215 Ok(LogEntry {
216 id,
217 daemon_id,
218 timestamp,
219 message,
220 level,
221 msg,
222 logger,
223 fields_json,
224 })
225 }
226
227 fn archive_entries(
228 &self,
229 entries: &[LogEntry],
230 archive_hook: &ArchiveHook,
231 daemon_id: &DaemonId,
232 reason: &str,
233 ) -> Result<()> {
234 use std::process::{Command, Stdio};
235
236 if entries.is_empty() {
237 return Ok(());
238 }
239
240 for chunk in entries.chunks(archive_hook.batch_size.max(1)) {
241 let mut child = Command::new("sh")
242 .arg("-c")
243 .arg(&archive_hook.command)
244 .stdin(Stdio::piped())
245 .stdout(Stdio::null())
246 .stderr(Stdio::piped())
247 .env("PITCHFORK_DAEMON_ID", daemon_id.qualified())
248 .env("PITCHFORK_ARCHIVE_REASON", reason)
249 .spawn()
250 .into_diagnostic()
251 .map_err(|e| miette::miette!("failed to spawn archive hook: {e}"))?;
252
253 let write_result = {
257 let stdin = child.stdin.take().expect("piped stdin should be available");
258 let mut stdin = std::io::BufWriter::new(stdin);
259 let mut result = Ok(());
260 for entry in chunk {
261 let line = serde_json::json!({
262 "id": entry.id,
263 "daemon_id": entry.daemon_id,
264 "timestamp": entry.timestamp.to_rfc3339(),
265 "message": entry.message,
266 });
267 if let Err(e) = writeln!(stdin, "{}", line) {
268 result = Err(miette::miette!(
269 "failed to write to archive hook stdin: {e}"
270 ));
271 break;
272 }
273 }
274 if result.is_ok()
282 && let Err(e) = stdin.flush()
283 {
284 result = Err(miette::miette!("failed to flush archive hook stdin: {e}"));
285 }
286 result
287 };
289
290 if let Err(e) = write_result {
291 let _ = child.kill();
292 let _ = child.wait();
293 return Err(e);
294 }
295
296 let output = child.wait_with_output().into_diagnostic()?;
297 if !output.status.success() {
298 let stderr = String::from_utf8_lossy(&output.stderr);
299 return Err(miette::miette!(
300 "archive hook failed with status {}: {stderr}",
301 output.status
302 ));
303 }
304 }
305
306 Ok(())
307 }
308
309 fn delete_by_ids(&self, ids: &[i64]) -> Result<u64> {
314 const SQLITE_MAX_VARS: usize = 999;
315
316 let mut total = 0u64;
317 let conn = self.conn.lock().unwrap();
318 for chunk in ids.chunks(SQLITE_MAX_VARS) {
319 if chunk.is_empty() {
320 continue;
321 }
322 let placeholders: Vec<String> = (1..=chunk.len()).map(|i| format!("?{i}")).collect();
323 let sql = format!(
324 "DELETE FROM log_entries WHERE id IN ({})",
325 placeholders.join(", ")
326 );
327 total += conn
328 .execute(&sql, rusqlite::params_from_iter(chunk.iter()))
329 .into_diagnostic()? as u64;
330 }
331 Ok(total)
332 }
333
334 pub fn rotate_by_age(
341 &self,
342 daemon_id: &DaemonId,
343 max_age: chrono::Duration,
344 archive_hook: Option<&ArchiveHook>,
345 ) -> Result<u64> {
346 let cutoff = (Local::now() - max_age).timestamp_millis();
347 let hook = archive_hook.filter(|h| h.is_enabled());
348
349 if let Some(hook) = hook {
350 let mut total_deleted = 0u64;
351 loop {
352 let entries: Vec<LogEntry> = {
354 let conn = self.conn.lock().unwrap();
355 let mut stmt = conn
356 .prepare(
357 "SELECT id, daemon_id, timestamp, message, level, msg, logger, fields_json FROM log_entries
358 WHERE daemon_id = ?1 AND timestamp < ?2
359 ORDER BY timestamp ASC, id ASC
360 LIMIT ?3",
361 )
362 .into_diagnostic()?;
363 stmt.query_map(
364 params![daemon_id.qualified(), cutoff, hook.batch_size as i64],
365 Self::row_to_entry,
366 )
367 .into_diagnostic()?
368 .collect::<rusqlite::Result<Vec<_>>>()
369 .into_diagnostic()?
370 };
371
372 if entries.is_empty() {
373 break;
374 }
375
376 let batch_ids: Vec<i64> = entries.iter().map(|e| e.id).collect();
377
378 self.archive_entries(&entries, hook, daemon_id, "age")?;
380
381 let deleted = self.delete_by_ids(&batch_ids)?;
383 total_deleted += deleted;
384 }
385 Ok(total_deleted)
386 } else {
387 let conn = self.conn.lock().unwrap();
388 let rows = conn
389 .execute(
390 "DELETE FROM log_entries WHERE daemon_id = ?1 AND timestamp < ?2",
391 params![daemon_id.qualified(), cutoff],
392 )
393 .into_diagnostic()?;
394 Ok(rows as u64)
395 }
396 }
397
398 pub fn rotate_by_count(
406 &self,
407 daemon_id: &DaemonId,
408 max_count: u64,
409 archive_hook: Option<&ArchiveHook>,
410 ) -> Result<u64> {
411 let hook = archive_hook.filter(|h| h.is_enabled());
412
413 let to_delete: i64 = {
415 let conn = self.conn.lock().unwrap();
416 let count: i64 = conn
417 .query_row(
418 "SELECT COUNT(*) FROM log_entries WHERE daemon_id = ?1",
419 [daemon_id.qualified()],
420 |row| row.get(0),
421 )
422 .into_diagnostic()?;
423 count.saturating_sub(max_count as i64)
424 };
425
426 if to_delete <= 0 {
427 return Ok(0);
428 }
429
430 if let Some(hook) = hook {
431 let mut total_deleted = 0u64;
432 let mut remaining = to_delete;
433 loop {
434 let batch_len = remaining.min(hook.batch_size as i64);
435
436 let entries: Vec<LogEntry> = {
438 let conn = self.conn.lock().unwrap();
439 let mut stmt = conn
440 .prepare(
441 "SELECT id, daemon_id, timestamp, message, level, msg, logger, fields_json FROM log_entries
442 WHERE daemon_id = ?1
443 ORDER BY timestamp ASC, id ASC
444 LIMIT ?2",
445 )
446 .into_diagnostic()?;
447 stmt.query_map(
448 params![daemon_id.qualified(), batch_len],
449 Self::row_to_entry,
450 )
451 .into_diagnostic()?
452 .collect::<rusqlite::Result<Vec<_>>>()
453 .into_diagnostic()?
454 };
455
456 if entries.is_empty() {
457 break;
458 }
459
460 let fetched = entries.len() as i64;
461 let batch_ids: Vec<i64> = entries.iter().map(|e| e.id).collect();
462
463 self.archive_entries(&entries, hook, daemon_id, "count")?;
465
466 let deleted = self.delete_by_ids(&batch_ids)?;
468 total_deleted += deleted;
469 remaining -= fetched;
470 }
471 Ok(total_deleted)
472 } else {
473 let conn = self.conn.lock().unwrap();
474 let rows = conn
475 .execute(
476 "DELETE FROM log_entries WHERE id IN (
477 SELECT id FROM log_entries WHERE daemon_id = ?1
478 ORDER BY timestamp ASC, id ASC LIMIT ?2
479 )",
480 params![daemon_id.qualified(), to_delete],
481 )
482 .into_diagnostic()?;
483 Ok(rows as u64)
484 }
485 }
486
487 pub fn migrate_daemon_text_logs(&self, daemon_id: &DaemonId) -> Result<u64> {
492 let text_path = daemon_id.log_path();
493 if !text_path.exists() {
494 return Ok(0);
495 }
496
497 let file = std::fs::File::open(&text_path).into_diagnostic()?;
498 let reader = BufReader::new(file);
499 let re = regex::Regex::new(r"^(\d{4}-\d{2}-\d{2} \d{2}:\d{2}:\d{2}) ([\w./-]+) (.*)$")
500 .expect("invalid regex");
501
502 let mut current_timestamp: Option<DateTime<Local>> = None;
503 let mut current_message = String::new();
504 let mut entries = Vec::with_capacity(1000);
505 let mut total_migrated: u64 = 0;
506
507 for line in reader.lines() {
508 let line = line.into_diagnostic()?;
509 if let Some(caps) = re.captures(&line) {
510 if let Some(ts) = current_timestamp.take() {
511 entries.push((ts, std::mem::take(&mut current_message)));
512 }
513 let ts_str = caps.get(1).map(|m| m.as_str()).unwrap_or_default();
514 let msg = caps.get(3).map(|m| m.as_str()).unwrap_or_default();
515 if let Ok(naive) =
516 chrono::NaiveDateTime::parse_from_str(ts_str, "%Y-%m-%d %H:%M:%S")
517 {
518 current_timestamp = Local.from_local_datetime(&naive).single();
519 current_message = msg.to_string();
520 }
521 } else if current_timestamp.is_some() {
522 current_message.push('\n');
523 current_message.push_str(&line);
524 }
525
526 if entries.len() >= 1000 {
527 total_migrated += self.insert_batch(daemon_id, &entries)?;
528 entries.clear();
529 }
530 }
531
532 if let Some(ts) = current_timestamp {
533 entries.push((ts, std::mem::take(&mut current_message)));
534 }
535
536 if !entries.is_empty() {
537 total_migrated += self.insert_batch(daemon_id, &entries)?;
538 }
539
540 if total_migrated > 0
541 && let Err(e) = std::fs::remove_file(&text_path)
542 {
543 log::warn!(
544 "failed to remove legacy log file after migration {}: {e}",
545 text_path.display()
546 );
547 }
548
549 Ok(total_migrated)
550 }
551
552 fn insert_batch(
553 &self,
554 daemon_id: &DaemonId,
555 entries: &[(DateTime<Local>, String)],
556 ) -> Result<u64> {
557 let mut conn = self.conn.lock().unwrap();
558 let tx = conn.transaction().into_diagnostic()?;
559 let mut count = 0u64;
560 {
561 let mut stmt = tx
562 .prepare(
563 "INSERT INTO log_entries (daemon_id, timestamp, message) VALUES (?1, ?2, ?3)",
564 )
565 .into_diagnostic()?;
566 for (ts, msg) in entries {
567 stmt.execute(params![daemon_id.qualified(), ts.timestamp_millis(), msg])
568 .into_diagnostic()?;
569 count += 1;
570 }
571 }
572 tx.commit().into_diagnostic()?;
573 Ok(count)
574 }
575
576 fn build_query_sql(
581 opts: &LogQuery,
582 id_range: Option<(i64, i64)>,
583 ) -> (String, Vec<Box<dyn rusqlite::ToSql>>) {
584 let mut conditions = Vec::new();
585 let mut query_params: Vec<Box<dyn rusqlite::ToSql>> = Vec::new();
586
587 if !opts.daemon_ids.is_empty() {
588 let placeholders: Vec<String> = (1..=opts.daemon_ids.len())
589 .map(|i| format!("?{}", i))
590 .collect();
591 conditions.push(format!("daemon_id IN ({})", placeholders.join(", ")));
592 for id in &opts.daemon_ids {
593 query_params.push(Box::new(id.clone()));
594 }
595 }
596
597 if let Some(from) = opts.from {
598 conditions.push(format!("timestamp >= ?{}", query_params.len() + 1));
599 query_params.push(Box::new(from.timestamp_millis()));
600 }
601
602 if let Some(to) = opts.to {
603 conditions.push(format!("timestamp <= ?{}", query_params.len() + 1));
604 query_params.push(Box::new(to.timestamp_millis()));
605 }
606
607 if let Some(after_id) = opts.after_id {
608 conditions.push(format!("id > ?{}", query_params.len() + 1));
609 query_params.push(Box::new(after_id));
610 }
611
612 if let Some((start, end)) = id_range {
613 conditions.push(format!("id > ?{}", query_params.len() + 1));
614 query_params.push(Box::new(start));
615 conditions.push(format!("id <= ?{}", query_params.len() + 1));
616 query_params.push(Box::new(end));
617 }
618
619 let mut message_conditions = Vec::new();
620 for filter in &opts.message_filters {
621 match filter {
622 MessageFilter::Contains {
623 pattern,
624 case_sensitive,
625 } => {
626 let param_index = query_params.len() + 1;
627 if *case_sensitive {
628 message_conditions
629 .push(format!("INSTR(message, ?{idx}) > 0", idx = param_index));
630 query_params.push(Box::new(pattern.clone()));
631 } else {
632 let escaped = escape_like_pattern(pattern);
633 let param = format!("%{}%", escaped);
634 message_conditions.push(format!(
635 "LOWER(message) LIKE LOWER(?{idx}) ESCAPE '\\'",
636 idx = param_index
637 ));
638 query_params.push(Box::new(param));
639 }
640 }
641 MessageFilter::Regex { pattern } => {
642 let param_index = query_params.len() + 1;
643 message_conditions.push(format!("message REGEXP ?{param_index}"));
644 query_params.push(Box::new(pattern.clone()));
645 }
646 }
647 }
648 if !message_conditions.is_empty() {
649 conditions.push(format!("({})", message_conditions.join(" OR ")));
650 }
651
652 for filter in &opts.field_filters {
653 match filter {
654 FieldFilter::LevelMin(level) => {
655 let matching = crate::log_store::levels_at_or_above(level);
656 if matching.is_empty() {
657 conditions.push("0".to_string());
659 } else {
660 let placeholders = matching
661 .iter()
662 .map(|l| {
663 let idx = query_params.len() + 1;
664 query_params.push(Box::new((*l).to_string()));
665 format!("?{idx}")
666 })
667 .collect::<Vec<_>>()
668 .join(", ");
669 conditions.push(format!("level IN ({placeholders})"));
670 }
671 }
672 FieldFilter::FieldEq { key, value } => {
673 let param_index = query_params.len() + 1;
674 conditions.push(format!(
675 "json_extract(fields_json, '$.{key}') = ?{param_index}"
676 ));
677 query_params.push(Box::new(value.clone()));
678 }
679 }
680 }
681
682 let where_clause = if conditions.is_empty() {
683 String::new()
684 } else {
685 format!("WHERE {}", conditions.join(" AND "))
686 };
687
688 let order = if opts.order_desc { "DESC" } else { "ASC" };
689
690 let limit_clause = opts
691 .limit
692 .map(|n| format!("LIMIT {}", n))
693 .unwrap_or_default();
694
695 let columns = if opts.include_structured {
696 "id, daemon_id, timestamp, message, level, msg, logger, fields_json"
697 } else {
698 "id, daemon_id, timestamp, message, NULL, NULL, NULL, NULL"
699 };
700
701 let sql = format!(
702 "SELECT {columns} FROM log_entries {} ORDER BY timestamp {}, id {} {}",
703 where_clause, order, order, limit_clause
704 );
705
706 (sql, query_params)
707 }
708
709 fn should_parallelize(opts: &LogQuery) -> bool {
711 if opts.daemon_ids.len() != 1 {
712 return false;
713 }
714 if opts.after_id.is_some() {
717 return false;
718 }
719 let limit = opts.limit.unwrap_or(usize::MAX);
720 if limit < PARALLEL_QUERY_THRESHOLD {
721 return false;
722 }
723 std::thread::available_parallelism()
724 .map(|n| n.get() >= 2)
725 .unwrap_or(false)
726 }
727
728 fn query_parallel(&self, opts: &LogQuery) -> Result<Vec<LogEntry>> {
735 let max_threads = 2;
740 let num_threads = std::thread::available_parallelism()
741 .map(|n| n.get().min(max_threads))
742 .unwrap_or(1);
743
744 let max_id: Option<i64> = {
745 let conn = self.conn.lock().unwrap();
746 conn.query_row(
747 "SELECT MAX(id) FROM log_entries WHERE daemon_id = ?1",
748 params![&opts.daemon_ids[0]],
749 |row| row.get(0),
750 )
751 .ok()
752 };
753
754 let Some(max_id) = max_id else {
755 return Ok(Vec::new());
756 };
757 if max_id == 0 {
758 return Ok(Vec::new());
759 }
760
761 let shard_size = (max_id as usize).div_ceil(num_threads);
762 let path = self.path.clone();
763 let opts = opts.clone();
764 let needs_regexp = opts
765 .message_filters
766 .iter()
767 .any(|f| matches!(f, MessageFilter::Regex { .. }));
768
769 let shards: Vec<Result<Vec<LogEntry>>> = std::thread::scope(|s| {
770 (0..num_threads)
771 .map(|i| {
772 let start = (i * shard_size) as i64;
773 let end = if i == num_threads - 1 {
774 max_id
775 } else {
776 ((i + 1) * shard_size) as i64
777 };
778 let opts = opts.clone();
779 let path = path.clone();
780 s.spawn(move || -> Result<Vec<LogEntry>> {
781 let conn = Connection::open_with_flags(
782 &path,
783 OpenFlags::SQLITE_OPEN_READ_ONLY | OpenFlags::SQLITE_OPEN_NO_MUTEX,
784 )
785 .into_diagnostic()?;
786 conn.execute_batch(
787 "PRAGMA mmap_size = 268435456;
788 PRAGMA query_only = 1;",
789 )
790 .into_diagnostic()?;
791 if needs_regexp {
792 add_regexp_function(&conn)?;
793 }
794
795 let (sql, query_params) = Self::build_query_sql(&opts, Some((start, end)));
796 Self::execute_built_query(&conn, &sql, &query_params)
797 })
798 })
799 .map(|h| h.join().unwrap())
800 .collect()
801 });
802
803 let mut merged = Vec::new();
804 if opts.order_desc {
805 for shard in shards.into_iter().rev() {
806 merged.extend(shard?);
807 }
808 } else {
809 for shard in shards {
810 merged.extend(shard?);
811 }
812 }
813
814 if let Some(limit) = opts.limit
815 && merged.len() > limit
816 {
817 merged.truncate(limit);
818 }
819
820 Ok(merged)
821 }
822
823 fn execute_built_query(
825 conn: &Connection,
826 sql: &str,
827 query_params: &[Box<dyn rusqlite::ToSql>],
828 ) -> Result<Vec<LogEntry>> {
829 let mut stmt = conn.prepare(sql).into_diagnostic()?;
830 let params_ref: Vec<&dyn rusqlite::ToSql> =
831 query_params.iter().map(|p| p.as_ref()).collect();
832 let rows = stmt
833 .query_map(params_ref.as_slice(), Self::row_to_entry)
834 .into_diagnostic()?;
835 let mut entries = Vec::new();
836 for row in rows {
837 entries.push(row.into_diagnostic()?);
838 }
839 Ok(entries)
840 }
841}
842
843impl LogStore for SqliteLogStore {
844 fn append(&self, daemon_id: &DaemonId, message: &str) -> Result<()> {
845 let ts = Local::now().timestamp_millis();
846 let id = daemon_id.qualified();
847 let msg = message.to_string();
848
849 let conn = self.conn.lock().unwrap();
850 let _ = conn
851 .execute(
852 "INSERT INTO log_entries (daemon_id, timestamp, message) VALUES (?1, ?2, ?3)",
853 params![id, ts, msg],
854 )
855 .into_diagnostic()?;
856 Ok(())
857 }
858
859 fn append_batch(&self, daemon_id: &DaemonId, messages: &[String]) -> Result<()> {
860 if messages.is_empty() {
861 return Ok(());
862 }
863 let base_ts = Local::now().timestamp_millis();
864 let id = daemon_id.qualified();
865
866 let mut conn = self.conn.lock().unwrap();
867 let tx = conn.transaction().into_diagnostic()?;
868 {
869 let mut stmt = tx
870 .prepare(
871 "INSERT INTO log_entries (daemon_id, timestamp, message) VALUES (?1, ?2, ?3)",
872 )
873 .into_diagnostic()?;
874 for (idx, msg) in messages.iter().enumerate() {
875 let ts = base_ts + idx as i64;
879 stmt.execute(params![id, ts, msg]).into_diagnostic()?;
880 }
881 }
882 tx.commit().into_diagnostic()?;
883 Ok(())
884 }
885
886 fn append_structured(&self, daemon_id: &DaemonId, parsed: &ParsedLog) -> Result<()> {
887 let ts = Local::now().timestamp_millis();
888 let id = daemon_id.qualified();
889
890 let conn = self.conn.lock().unwrap();
891 let _ = conn
892 .execute(
893 "INSERT INTO log_entries (daemon_id, timestamp, message, level, msg, logger, fields_json) \
894 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)",
895 params![id, ts, parsed.message, parsed.level, parsed.msg, parsed.logger, parsed.fields_json],
896 )
897 .into_diagnostic()?;
898 Ok(())
899 }
900
901 fn append_structured_batch(&self, daemon_id: &DaemonId, entries: &[ParsedLog]) -> Result<()> {
902 if entries.is_empty() {
903 return Ok(());
904 }
905 let base_ts = Local::now().timestamp_millis();
906 let id = daemon_id.qualified();
907
908 let mut conn = self.conn.lock().unwrap();
909 let tx = conn.transaction().into_diagnostic()?;
910 {
911 let mut stmt = tx
912 .prepare(
913 "INSERT INTO log_entries (daemon_id, timestamp, message, level, msg, logger, fields_json) \
914 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)",
915 )
916 .into_diagnostic()?;
917 for (idx, entry) in entries.iter().enumerate() {
918 let ts = base_ts + idx as i64;
919 stmt.execute(params![
920 id,
921 ts,
922 entry.message,
923 entry.level,
924 entry.msg,
925 entry.logger,
926 entry.fields_json
927 ])
928 .into_diagnostic()?;
929 }
930 }
931 tx.commit().into_diagnostic()?;
932 Ok(())
933 }
934
935 fn query(&self, opts: &LogQuery) -> Result<Vec<LogEntry>> {
936 if Self::should_parallelize(opts)
938 && self.path.as_os_str() != ":memory:"
939 && let Ok(entries) = self.query_parallel(opts)
940 {
941 return Ok(entries);
942 }
943 let conn = self.conn.lock().unwrap();
947 let (sql, query_params) = Self::build_query_sql(opts, None);
948 Self::execute_built_query(&conn, &sql, &query_params)
949 }
950
951 fn tail(&self, daemon_id: &DaemonId, after_id: Option<i64>) -> Result<Vec<LogEntry>> {
952 self.query(&LogQuery {
953 daemon_ids: vec![daemon_id.qualified()],
954 from: None,
955 to: None,
956 limit: None,
957 order_desc: false,
958 after_id,
959 message_filters: Vec::new(),
960 field_filters: Vec::new(),
961 include_structured: false,
962 })
963 }
964
965 fn clear(&self, daemon_ids: &[DaemonId]) -> Result<()> {
966 let mut conn = self.conn.lock().unwrap();
967 let tx = conn.transaction().into_diagnostic()?;
968 for id in daemon_ids {
969 tx.execute(
970 "DELETE FROM log_entries WHERE daemon_id = ?1",
971 params![id.qualified()],
972 )
973 .into_diagnostic()?;
974 tx.execute(
975 "INSERT INTO log_clear_generations (daemon_id, generation)
976 VALUES (?1, 1)
977 ON CONFLICT(daemon_id) DO UPDATE SET generation = generation + 1",
978 params![id.qualified()],
979 )
980 .into_diagnostic()?;
981 }
982 tx.commit().into_diagnostic()?;
983 Ok(())
984 }
985
986 fn last_id(&self, daemon_id: &DaemonId) -> Result<Option<i64>> {
987 let conn = self.conn.lock().unwrap();
988 let id: Option<i64> = conn
991 .query_row(
992 "SELECT MAX(id) FROM log_entries WHERE daemon_id = ?1",
993 params![daemon_id.qualified()],
994 |row| row.get(0),
995 )
996 .into_diagnostic()?;
997 Ok(id)
998 }
999
1000 fn list_daemon_ids(&self) -> Result<Vec<String>> {
1001 let conn = self.conn.lock().unwrap();
1002 let mut stmt = conn
1003 .prepare("SELECT DISTINCT daemon_id FROM log_entries")
1004 .into_diagnostic()?;
1005 let ids = stmt
1006 .query_map([], |row| {
1007 let id: String = row.get(0)?;
1008 Ok(id)
1009 })
1010 .into_diagnostic()?
1011 .filter_map(|r| r.ok())
1012 .collect();
1013 Ok(ids)
1014 }
1015
1016 fn apply_retention(
1017 &self,
1018 policy: &super::RetentionPolicy,
1019 excluded_daemon_ids: &[DaemonId],
1020 archive_hook: Option<&ArchiveHook>,
1021 ) -> Result<u64> {
1022 let daemon_ids = self.list_daemon_ids()?;
1023 let excluded: HashSet<String> = excluded_daemon_ids.iter().map(|d| d.qualified()).collect();
1024 let mut total = 0u64;
1025 for id_str in daemon_ids {
1026 if excluded.contains(&id_str) {
1027 continue;
1028 }
1029 let id = DaemonId::parse(&id_str).unwrap_or_else(|_| {
1030 DaemonId::try_new("global", &id_str).unwrap_or_else(|_| DaemonId::pitchfork())
1031 });
1032 if let Some(dur) = policy.age {
1033 total += self.rotate_by_age(&id, dur, archive_hook)?;
1034 }
1035 if let Some(n) = policy.count {
1036 total += self.rotate_by_count(&id, n, archive_hook)?;
1037 }
1038 }
1039 Ok(total)
1040 }
1041
1042 fn apply_retention_for_daemon(
1043 &self,
1044 daemon_id: &DaemonId,
1045 policy: &super::RetentionPolicy,
1046 archive_hook: Option<&ArchiveHook>,
1047 ) -> Result<u64> {
1048 let mut total = 0u64;
1049 if let Some(dur) = policy.age {
1050 total += self.rotate_by_age(daemon_id, dur, archive_hook)?;
1051 }
1052 if let Some(n) = policy.count {
1053 total += self.rotate_by_count(daemon_id, n, archive_hook)?;
1054 }
1055 Ok(total)
1056 }
1057
1058 fn last_clear_generation(&self, daemon_id: &DaemonId) -> Result<Option<u64>> {
1059 let conn = self.conn.lock().unwrap();
1060 let mut stmt = conn
1061 .prepare("SELECT generation FROM log_clear_generations WHERE daemon_id = ?1")
1062 .into_diagnostic()?;
1063 let generation: Option<i64> = stmt
1064 .query_row(params![daemon_id.qualified()], |row| row.get(0))
1065 .optional()
1066 .into_diagnostic()?;
1067 generation
1068 .map(|generation| {
1069 u64::try_from(generation)
1070 .map_err(|_| miette::miette!("log clear generation cannot be negative"))
1071 })
1072 .transpose()
1073 }
1074}
1075
1076use once_cell::sync::Lazy;
1078use std::sync::Arc;
1079
1080pub static LOG_STORE: Lazy<Arc<SqliteLogStore>> = Lazy::new(|| {
1081 let path = crate::env::PITCHFORK_LOGS_DIR.join("logs.db");
1082 let mut is_fallback = false;
1083 let store = Arc::new(SqliteLogStore::open(&path).unwrap_or_else(|e| {
1084 error!(
1085 "failed to open log store at {}: {e}. Falling back to in-memory store; logs will not persist across restarts.",
1086 path.display()
1087 );
1088 is_fallback = true;
1089 SqliteLogStore::open(":memory:").expect("in-memory SQLite should always open")
1090 }));
1091
1092 if !is_fallback {
1096 if let Err(e) = auto_migrate_legacy_logs(&store) {
1097 warn!("legacy log auto-migration failed: {e}");
1098 }
1099 } else {
1100 warn!(
1101 "skipping legacy log auto-migration because log store is in-memory (no durable destination)"
1102 );
1103 }
1104
1105 store
1106});
1107
1108fn auto_migrate_legacy_logs(store: &SqliteLogStore) -> Result<()> {
1116 let logs_dir = &*crate::env::PITCHFORK_LOGS_DIR;
1117 if !logs_dir.exists() {
1118 return Ok(());
1119 }
1120
1121 let Ok(entries) = std::fs::read_dir(logs_dir) else {
1122 return Ok(());
1123 };
1124
1125 let mut total_migrated = 0u64;
1126 let mut migrated_ids = Vec::new();
1127
1128 for entry in entries.flatten() {
1129 let path = entry.path();
1130 if !path.is_dir() {
1131 continue;
1132 }
1133 let file_name = path
1135 .file_name()
1136 .map_or(String::new(), |n| n.to_string_lossy().to_string());
1137 if file_name == "pitchfork" {
1138 continue;
1139 }
1140
1141 if !file_name.contains("--") {
1143 continue;
1144 }
1145 let log_file = path.join(format!("{file_name}.log"));
1146 if !log_file.exists() {
1147 continue;
1148 }
1149
1150 let daemon_id = match DaemonId::from_safe_path(&file_name) {
1151 Ok(id) => id,
1152 Err(_) => continue,
1153 };
1154
1155 if daemon_id == DaemonId::pitchfork() {
1158 continue;
1159 }
1160
1161 match store.migrate_daemon_text_logs(&daemon_id) {
1162 Ok(0) => {}
1163 Ok(n) => {
1164 total_migrated += n;
1165 migrated_ids.push(daemon_id.qualified());
1166 }
1167 Err(e) => {
1168 warn!(
1169 "failed to migrate text logs for {}: {e}",
1170 daemon_id.qualified()
1171 );
1172 }
1173 }
1174 }
1175
1176 if total_migrated > 0 {
1177 warn!(
1178 "auto-migrated {total_migrated} legacy log entries from {count} daemon(s): {ids}",
1179 count = migrated_ids.len(),
1180 ids = migrated_ids.join(", ")
1181 );
1182 }
1183
1184 Ok(())
1185}
1186
1187#[cfg(test)]
1188mod tests {
1189 use super::*;
1190 use crate::log_store::LogStore;
1191
1192 #[test]
1196 fn write_waits_for_a_contended_lock() {
1197 let dir = tempfile::tempdir().unwrap();
1198 let path = dir.path().join("logs.db");
1199 let store = SqliteLogStore::open(&path).unwrap();
1200
1201 let blocker = Connection::open(&path).unwrap();
1204 blocker.busy_timeout(BUSY_TIMEOUT).unwrap();
1205 blocker.execute_batch("BEGIN IMMEDIATE").unwrap();
1206 let releaser = std::thread::spawn(move || {
1207 std::thread::sleep(std::time::Duration::from_millis(300));
1208 blocker.execute_batch("COMMIT").unwrap();
1209 });
1210
1211 let id = DaemonId::try_new("test", "blocked").unwrap();
1212 let entries = vec![crate::log_parse::parse("held-lock-line", "text")];
1213 store
1214 .append_structured_batch(&id, &entries)
1215 .expect("a contended write must wait for the lock, not fail");
1216
1217 releaser.join().unwrap();
1218 let found = store
1219 .query(&LogQuery {
1220 daemon_ids: vec![id.qualified()],
1221 ..Default::default()
1222 })
1223 .unwrap()
1224 .len();
1225 assert_eq!(found, 1, "the batch written under contention was lost");
1226 }
1227
1228 #[test]
1234 fn concurrent_writers_do_not_lose_batches() {
1235 let dir = tempfile::tempdir().unwrap();
1236 let path = dir.path().join("logs.db");
1237
1238 const WRITERS: usize = 4;
1239 const BATCHES: usize = 15;
1240 const PER_BATCH: usize = 10;
1241
1242 let barrier = std::sync::Barrier::new(WRITERS);
1247
1248 std::thread::scope(|scope| {
1249 for writer in 0..WRITERS {
1250 let path = path.clone();
1251 let barrier = &barrier;
1252 scope.spawn(move || {
1253 barrier.wait();
1256 let store = SqliteLogStore::open(&path).unwrap();
1257 let id = DaemonId::try_new("test", format!("w{writer}")).unwrap();
1258 for batch in 0..BATCHES {
1259 let entries: Vec<ParsedLog> = (0..PER_BATCH)
1260 .map(|i| {
1261 crate::log_parse::parse(&format!("w{writer}-{batch}-{i}"), "text")
1262 })
1263 .collect();
1264 store
1265 .append_structured_batch(&id, &entries)
1266 .expect("concurrent batch write must not fail");
1267 }
1268 });
1269 }
1270 });
1271
1272 let store = SqliteLogStore::open(&path).unwrap();
1273 for writer in 0..WRITERS {
1274 let id = DaemonId::try_new("test", format!("w{writer}")).unwrap();
1275 let found = store
1276 .query(&LogQuery {
1277 daemon_ids: vec![id.qualified()],
1278 ..Default::default()
1279 })
1280 .unwrap()
1281 .len();
1282 assert_eq!(
1283 found,
1284 BATCHES * PER_BATCH,
1285 "writer {writer} lost entries: got {found}"
1286 );
1287 }
1288 }
1289}