1use crate::Result;
2use crate::daemon_id::DaemonId;
3use crate::log_store::{ArchiveHook, LogEntry, LogQuery, LogStore};
4use chrono::{DateTime, Local, TimeZone};
5use log::error;
6use miette::IntoDiagnostic;
7use rusqlite::{Connection, OptionalExtension, params};
8use std::collections::HashSet;
9use std::io::{BufRead, BufReader, Write};
10use std::path::PathBuf;
11use std::sync::Mutex;
12
13pub struct SqliteLogStore {
15 conn: Mutex<Connection>,
16}
17
18impl SqliteLogStore {
19 pub fn open(path: impl Into<PathBuf>) -> Result<Self> {
21 let path = path.into();
22 if let Some(parent) = path.parent() {
23 std::fs::create_dir_all(parent).into_diagnostic()?;
24 }
25 let conn = Connection::open(&path).into_diagnostic()?;
26 conn.execute_batch(
27 "PRAGMA journal_mode = WAL;
28 PRAGMA synchronous = NORMAL;",
29 )
30 .into_diagnostic()?;
31 conn.execute(
32 "CREATE TABLE IF NOT EXISTS log_entries (
33 id INTEGER PRIMARY KEY AUTOINCREMENT,
34 daemon_id TEXT NOT NULL,
35 timestamp INTEGER NOT NULL,
36 message TEXT NOT NULL
37 );",
38 [],
39 )
40 .into_diagnostic()?;
41 conn.execute(
42 "CREATE INDEX IF NOT EXISTS idx_daemon_ts ON log_entries(daemon_id, timestamp);",
43 [],
44 )
45 .into_diagnostic()?;
46 conn.execute(
47 "CREATE INDEX IF NOT EXISTS idx_daemon_id ON log_entries(daemon_id, id);",
48 [],
49 )
50 .into_diagnostic()?;
51 conn.execute(
52 "CREATE INDEX IF NOT EXISTS idx_timestamp ON log_entries(timestamp);",
53 [],
54 )
55 .into_diagnostic()?;
56 conn.execute(
57 "CREATE TABLE IF NOT EXISTS log_clear_generations (
58 daemon_id TEXT PRIMARY KEY,
59 generation INTEGER NOT NULL DEFAULT 0
60 );",
61 [],
62 )
63 .into_diagnostic()?;
64 Ok(Self {
65 conn: Mutex::new(conn),
66 })
67 }
68
69 fn row_to_entry(row: &rusqlite::Row) -> rusqlite::Result<LogEntry> {
70 let id: i64 = row.get(0)?;
71 let daemon_id: String = row.get(1)?;
72 let ts_millis: i64 = row.get(2)?;
73 let message: String = row.get(3)?;
74 let timestamp = Local
75 .timestamp_millis_opt(ts_millis)
76 .single()
77 .unwrap_or_else(Local::now);
78 Ok(LogEntry {
79 id,
80 daemon_id,
81 timestamp,
82 message,
83 })
84 }
85
86 fn archive_entries(
87 &self,
88 entries: &[LogEntry],
89 archive_hook: &ArchiveHook,
90 daemon_id: &DaemonId,
91 reason: &str,
92 ) -> Result<()> {
93 use std::process::{Command, Stdio};
94
95 if entries.is_empty() {
96 return Ok(());
97 }
98
99 for chunk in entries.chunks(archive_hook.batch_size.max(1)) {
100 let mut child = Command::new("sh")
101 .arg("-c")
102 .arg(&archive_hook.command)
103 .stdin(Stdio::piped())
104 .stdout(Stdio::null())
105 .stderr(Stdio::piped())
106 .env("PITCHFORK_DAEMON_ID", daemon_id.qualified())
107 .env("PITCHFORK_ARCHIVE_REASON", reason)
108 .spawn()
109 .into_diagnostic()
110 .map_err(|e| miette::miette!("failed to spawn archive hook: {e}"))?;
111
112 let write_result = {
116 let stdin = child.stdin.take().expect("piped stdin should be available");
117 let mut stdin = std::io::BufWriter::new(stdin);
118 let mut result = Ok(());
119 for entry in chunk {
120 let line = serde_json::json!({
121 "id": entry.id,
122 "daemon_id": entry.daemon_id,
123 "timestamp": entry.timestamp.to_rfc3339(),
124 "message": entry.message,
125 });
126 if let Err(e) = writeln!(stdin, "{}", line) {
127 result = Err(miette::miette!(
128 "failed to write to archive hook stdin: {e}"
129 ));
130 break;
131 }
132 }
133 if result.is_ok() {
141 if let Err(e) = stdin.flush() {
142 result = Err(miette::miette!("failed to flush archive hook stdin: {e}"));
143 }
144 }
145 result
146 };
148
149 if let Err(e) = write_result {
150 let _ = child.kill();
151 let _ = child.wait();
152 return Err(e);
153 }
154
155 let output = child.wait_with_output().into_diagnostic()?;
156 if !output.status.success() {
157 let stderr = String::from_utf8_lossy(&output.stderr);
158 return Err(miette::miette!(
159 "archive hook failed with status {}: {stderr}",
160 output.status
161 ));
162 }
163 }
164
165 Ok(())
166 }
167
168 fn delete_by_ids(&self, ids: &[i64]) -> Result<u64> {
173 const SQLITE_MAX_VARS: usize = 999;
174
175 let mut total = 0u64;
176 let conn = self.conn.lock().unwrap();
177 for chunk in ids.chunks(SQLITE_MAX_VARS) {
178 if chunk.is_empty() {
179 continue;
180 }
181 let placeholders: Vec<String> = (1..=chunk.len()).map(|i| format!("?{i}")).collect();
182 let sql = format!(
183 "DELETE FROM log_entries WHERE id IN ({})",
184 placeholders.join(", ")
185 );
186 total += conn
187 .execute(&sql, rusqlite::params_from_iter(chunk.iter()))
188 .into_diagnostic()? as u64;
189 }
190 Ok(total)
191 }
192
193 pub fn rotate_by_age(
200 &self,
201 daemon_id: &DaemonId,
202 max_age: chrono::Duration,
203 archive_hook: Option<&ArchiveHook>,
204 ) -> Result<u64> {
205 let cutoff = (Local::now() - max_age).timestamp_millis();
206 let hook = archive_hook.filter(|h| h.is_enabled());
207
208 if let Some(hook) = hook {
209 let mut total_deleted = 0u64;
210 loop {
211 let entries: Vec<LogEntry> = {
213 let conn = self.conn.lock().unwrap();
214 let mut stmt = conn
215 .prepare(
216 "SELECT id, daemon_id, timestamp, message FROM log_entries
217 WHERE daemon_id = ?1 AND timestamp < ?2
218 ORDER BY timestamp ASC, id ASC
219 LIMIT ?3",
220 )
221 .into_diagnostic()?;
222 stmt.query_map(
223 params![daemon_id.qualified(), cutoff, hook.batch_size as i64],
224 Self::row_to_entry,
225 )
226 .into_diagnostic()?
227 .collect::<rusqlite::Result<Vec<_>>>()
228 .into_diagnostic()?
229 };
230
231 if entries.is_empty() {
232 break;
233 }
234
235 let batch_ids: Vec<i64> = entries.iter().map(|e| e.id).collect();
236
237 self.archive_entries(&entries, hook, daemon_id, "age")?;
239
240 let deleted = self.delete_by_ids(&batch_ids)?;
242 total_deleted += deleted;
243 }
244 Ok(total_deleted)
245 } else {
246 let conn = self.conn.lock().unwrap();
247 let rows = conn
248 .execute(
249 "DELETE FROM log_entries WHERE daemon_id = ?1 AND timestamp < ?2",
250 params![daemon_id.qualified(), cutoff],
251 )
252 .into_diagnostic()?;
253 Ok(rows as u64)
254 }
255 }
256
257 pub fn rotate_by_count(
265 &self,
266 daemon_id: &DaemonId,
267 max_count: u64,
268 archive_hook: Option<&ArchiveHook>,
269 ) -> Result<u64> {
270 let hook = archive_hook.filter(|h| h.is_enabled());
271
272 let to_delete: i64 = {
274 let conn = self.conn.lock().unwrap();
275 let count: i64 = conn
276 .query_row(
277 "SELECT COUNT(*) FROM log_entries WHERE daemon_id = ?1",
278 [daemon_id.qualified()],
279 |row| row.get(0),
280 )
281 .into_diagnostic()?;
282 count.saturating_sub(max_count as i64)
283 };
284
285 if to_delete <= 0 {
286 return Ok(0);
287 }
288
289 if let Some(hook) = hook {
290 let mut total_deleted = 0u64;
291 let mut remaining = to_delete;
292 loop {
293 let batch_len = remaining.min(hook.batch_size as i64);
294
295 let entries: Vec<LogEntry> = {
297 let conn = self.conn.lock().unwrap();
298 let mut stmt = conn
299 .prepare(
300 "SELECT id, daemon_id, timestamp, message FROM log_entries
301 WHERE daemon_id = ?1
302 ORDER BY timestamp ASC, id ASC
303 LIMIT ?2",
304 )
305 .into_diagnostic()?;
306 stmt.query_map(
307 params![daemon_id.qualified(), batch_len],
308 Self::row_to_entry,
309 )
310 .into_diagnostic()?
311 .collect::<rusqlite::Result<Vec<_>>>()
312 .into_diagnostic()?
313 };
314
315 if entries.is_empty() {
316 break;
317 }
318
319 let fetched = entries.len() as i64;
320 let batch_ids: Vec<i64> = entries.iter().map(|e| e.id).collect();
321
322 self.archive_entries(&entries, hook, daemon_id, "count")?;
324
325 let deleted = self.delete_by_ids(&batch_ids)?;
327 total_deleted += deleted;
328 remaining -= fetched;
329 }
330 Ok(total_deleted)
331 } else {
332 let conn = self.conn.lock().unwrap();
333 let rows = conn
334 .execute(
335 "DELETE FROM log_entries WHERE id IN (
336 SELECT id FROM log_entries WHERE daemon_id = ?1
337 ORDER BY timestamp ASC, id ASC LIMIT ?2
338 )",
339 params![daemon_id.qualified(), to_delete],
340 )
341 .into_diagnostic()?;
342 Ok(rows as u64)
343 }
344 }
345
346 pub fn migrate_daemon_text_logs(&self, daemon_id: &DaemonId) -> Result<u64> {
351 let text_path = daemon_id.log_path();
352 if !text_path.exists() {
353 return Ok(0);
354 }
355
356 let file = std::fs::File::open(&text_path).into_diagnostic()?;
357 let reader = BufReader::new(file);
358 let re = regex::Regex::new(r"^(\d{4}-\d{2}-\d{2} \d{2}:\d{2}:\d{2}) ([\w./-]+) (.*)$")
359 .expect("invalid regex");
360
361 let mut current_timestamp: Option<DateTime<Local>> = None;
362 let mut current_message = String::new();
363 let mut entries = Vec::with_capacity(1000);
364 let mut total_migrated: u64 = 0;
365
366 for line in reader.lines() {
367 let line = line.into_diagnostic()?;
368 if let Some(caps) = re.captures(&line) {
369 if let Some(ts) = current_timestamp.take() {
370 entries.push((ts, std::mem::take(&mut current_message)));
371 }
372 let ts_str = caps.get(1).map(|m| m.as_str()).unwrap_or_default();
373 let msg = caps.get(3).map(|m| m.as_str()).unwrap_or_default();
374 if let Ok(naive) =
375 chrono::NaiveDateTime::parse_from_str(ts_str, "%Y-%m-%d %H:%M:%S")
376 {
377 current_timestamp = Local.from_local_datetime(&naive).single();
378 current_message = msg.to_string();
379 }
380 } else if current_timestamp.is_some() {
381 current_message.push('\n');
382 current_message.push_str(&line);
383 }
384
385 if entries.len() >= 1000 {
386 total_migrated += self.insert_batch(daemon_id, &entries)?;
387 entries.clear();
388 }
389 }
390
391 if let Some(ts) = current_timestamp {
392 entries.push((ts, std::mem::take(&mut current_message)));
393 }
394
395 if !entries.is_empty() {
396 total_migrated += self.insert_batch(daemon_id, &entries)?;
397 }
398
399 if total_migrated > 0 {
400 if let Err(e) = std::fs::remove_file(&text_path) {
401 log::warn!(
402 "failed to remove legacy log file after migration {}: {e}",
403 text_path.display()
404 );
405 }
406 }
407
408 Ok(total_migrated)
409 }
410
411 fn insert_batch(
412 &self,
413 daemon_id: &DaemonId,
414 entries: &[(DateTime<Local>, String)],
415 ) -> Result<u64> {
416 let mut conn = self.conn.lock().unwrap();
417 let tx = conn.transaction().into_diagnostic()?;
418 let mut count = 0u64;
419 {
420 let mut stmt = tx
421 .prepare(
422 "INSERT INTO log_entries (daemon_id, timestamp, message) VALUES (?1, ?2, ?3)",
423 )
424 .into_diagnostic()?;
425 for (ts, msg) in entries {
426 stmt.execute(params![daemon_id.qualified(), ts.timestamp_millis(), msg])
427 .into_diagnostic()?;
428 count += 1;
429 }
430 }
431 tx.commit().into_diagnostic()?;
432 Ok(count)
433 }
434}
435
436impl LogStore for SqliteLogStore {
437 fn append(&self, daemon_id: &DaemonId, message: &str) -> Result<()> {
438 let ts = Local::now().timestamp_millis();
439 let id = daemon_id.qualified();
440 let msg = message.to_string();
441
442 let conn = self.conn.lock().unwrap();
443 let _ = conn
444 .execute(
445 "INSERT INTO log_entries (daemon_id, timestamp, message) VALUES (?1, ?2, ?3)",
446 params![id, ts, msg],
447 )
448 .into_diagnostic()?;
449 Ok(())
450 }
451
452 fn query(&self, opts: &LogQuery) -> Result<Vec<LogEntry>> {
453 let conn = self.conn.lock().unwrap();
454 let mut conditions = Vec::new();
455 let mut query_params: Vec<Box<dyn rusqlite::ToSql>> = Vec::new();
456
457 if !opts.daemon_ids.is_empty() {
458 let placeholders: Vec<String> = (1..=opts.daemon_ids.len())
459 .map(|i| format!("?{}", i))
460 .collect();
461 conditions.push(format!("daemon_id IN ({})", placeholders.join(", ")));
462 for id in &opts.daemon_ids {
463 query_params.push(Box::new(id.clone()));
464 }
465 }
466
467 if let Some(from) = opts.from {
468 conditions.push(format!("timestamp >= ?{}", query_params.len() + 1));
469 query_params.push(Box::new(from.timestamp_millis()));
470 }
471
472 if let Some(to) = opts.to {
473 conditions.push(format!("timestamp <= ?{}", query_params.len() + 1));
474 query_params.push(Box::new(to.timestamp_millis()));
475 }
476
477 if let Some(after_id) = opts.after_id {
478 conditions.push(format!("id > ?{}", query_params.len() + 1));
479 query_params.push(Box::new(after_id));
480 }
481
482 let where_clause = if conditions.is_empty() {
483 String::new()
484 } else {
485 format!("WHERE {}", conditions.join(" AND "))
486 };
487
488 let order = if opts.order_desc { "DESC" } else { "ASC" };
489
490 let limit_clause = opts
491 .limit
492 .map(|n| format!("LIMIT {}", n))
493 .unwrap_or_default();
494
495 let sql = format!(
496 "SELECT id, daemon_id, timestamp, message FROM log_entries {} ORDER BY timestamp {}, id {} {}",
497 where_clause, order, order, limit_clause
498 );
499
500 let mut stmt = conn.prepare(&sql).into_diagnostic()?;
501 let params_ref: Vec<&dyn rusqlite::ToSql> =
502 query_params.iter().map(|p| p.as_ref()).collect();
503 let rows = stmt
504 .query_map(params_ref.as_slice(), |row| {
505 let id: i64 = row.get(0)?;
506 let daemon_id: String = row.get(1)?;
507 let ts_millis: i64 = row.get(2)?;
508 let message: String = row.get(3)?;
509 let timestamp = Local
510 .timestamp_millis_opt(ts_millis)
511 .single()
512 .unwrap_or_else(Local::now);
513 Ok(LogEntry {
514 id,
515 daemon_id,
516 timestamp,
517 message,
518 })
519 })
520 .into_diagnostic()?;
521
522 let mut entries = Vec::new();
523 for row in rows {
524 entries.push(row.into_diagnostic()?);
525 }
526 Ok(entries)
527 }
528
529 fn tail(&self, daemon_id: &DaemonId, after_id: Option<i64>) -> Result<Vec<LogEntry>> {
530 self.query(&LogQuery {
531 daemon_ids: vec![daemon_id.qualified()],
532 from: None,
533 to: None,
534 limit: None,
535 order_desc: false,
536 after_id,
537 })
538 }
539
540 fn clear(&self, daemon_ids: &[DaemonId]) -> Result<()> {
541 let mut conn = self.conn.lock().unwrap();
542 let tx = conn.transaction().into_diagnostic()?;
543 for id in daemon_ids {
544 tx.execute(
545 "DELETE FROM log_entries WHERE daemon_id = ?1",
546 params![id.qualified()],
547 )
548 .into_diagnostic()?;
549 tx.execute(
550 "INSERT INTO log_clear_generations (daemon_id, generation)
551 VALUES (?1, 1)
552 ON CONFLICT(daemon_id) DO UPDATE SET generation = generation + 1",
553 params![id.qualified()],
554 )
555 .into_diagnostic()?;
556 }
557 tx.commit().into_diagnostic()?;
558 Ok(())
559 }
560
561 fn list_daemon_ids(&self) -> Result<Vec<String>> {
562 let conn = self.conn.lock().unwrap();
563 let mut stmt = conn
564 .prepare("SELECT DISTINCT daemon_id FROM log_entries")
565 .into_diagnostic()?;
566 let ids = stmt
567 .query_map([], |row| {
568 let id: String = row.get(0)?;
569 Ok(id)
570 })
571 .into_diagnostic()?
572 .filter_map(|r| r.ok())
573 .collect();
574 Ok(ids)
575 }
576
577 fn apply_retention(
578 &self,
579 policy: &super::RetentionPolicy,
580 excluded_daemon_ids: &[DaemonId],
581 archive_hook: Option<&ArchiveHook>,
582 ) -> Result<u64> {
583 let daemon_ids = self.list_daemon_ids()?;
584 let excluded: HashSet<String> = excluded_daemon_ids.iter().map(|d| d.qualified()).collect();
585 let mut total = 0u64;
586 for id_str in daemon_ids {
587 if excluded.contains(&id_str) {
588 continue;
589 }
590 let id = DaemonId::parse(&id_str).unwrap_or_else(|_| {
591 DaemonId::try_new("global", &id_str).unwrap_or_else(|_| DaemonId::pitchfork())
592 });
593 if let Some(dur) = policy.age {
594 total += self.rotate_by_age(&id, dur, archive_hook)?;
595 }
596 if let Some(n) = policy.count {
597 total += self.rotate_by_count(&id, n, archive_hook)?;
598 }
599 }
600 Ok(total)
601 }
602
603 fn apply_retention_for_daemon(
604 &self,
605 daemon_id: &DaemonId,
606 policy: &super::RetentionPolicy,
607 archive_hook: Option<&ArchiveHook>,
608 ) -> Result<u64> {
609 let mut total = 0u64;
610 if let Some(dur) = policy.age {
611 total += self.rotate_by_age(daemon_id, dur, archive_hook)?;
612 }
613 if let Some(n) = policy.count {
614 total += self.rotate_by_count(daemon_id, n, archive_hook)?;
615 }
616 Ok(total)
617 }
618
619 fn last_clear_generation(&self, daemon_id: &DaemonId) -> Result<Option<u64>> {
620 let conn = self.conn.lock().unwrap();
621 let mut stmt = conn
622 .prepare("SELECT generation FROM log_clear_generations WHERE daemon_id = ?1")
623 .into_diagnostic()?;
624 let generation: Option<i64> = stmt
625 .query_row(params![daemon_id.qualified()], |row| row.get(0))
626 .optional()
627 .into_diagnostic()?;
628 generation
629 .map(|generation| {
630 u64::try_from(generation)
631 .map_err(|_| miette::miette!("log clear generation cannot be negative"))
632 })
633 .transpose()
634 }
635}
636
637use once_cell::sync::Lazy;
639use std::sync::Arc;
640
641pub static LOG_STORE: Lazy<Arc<SqliteLogStore>> = Lazy::new(|| {
642 let path = crate::env::PITCHFORK_LOGS_DIR.join("logs.db");
643 let mut is_fallback = false;
644 let store = Arc::new(SqliteLogStore::open(&path).unwrap_or_else(|e| {
645 error!(
646 "failed to open log store at {}: {e}. Falling back to in-memory store; logs will not persist across restarts.",
647 path.display()
648 );
649 is_fallback = true;
650 SqliteLogStore::open(":memory:").expect("in-memory SQLite should always open")
651 }));
652
653 if !is_fallback {
657 if let Err(e) = auto_migrate_legacy_logs(&store) {
658 warn!("legacy log auto-migration failed: {e}");
659 }
660 } else {
661 warn!(
662 "skipping legacy log auto-migration because log store is in-memory (no durable destination)"
663 );
664 }
665
666 store
667});
668
669fn auto_migrate_legacy_logs(store: &SqliteLogStore) -> Result<()> {
677 let logs_dir = &*crate::env::PITCHFORK_LOGS_DIR;
678 if !logs_dir.exists() {
679 return Ok(());
680 }
681
682 let Ok(entries) = std::fs::read_dir(logs_dir) else {
683 return Ok(());
684 };
685
686 let mut total_migrated = 0u64;
687 let mut migrated_ids = Vec::new();
688
689 for entry in entries.flatten() {
690 let path = entry.path();
691 if !path.is_dir() {
692 continue;
693 }
694 let file_name = path
696 .file_name()
697 .map_or(String::new(), |n| n.to_string_lossy().to_string());
698 if file_name == "pitchfork" {
699 continue;
700 }
701
702 if !file_name.contains("--") {
704 continue;
705 }
706 let log_file = path.join(format!("{file_name}.log"));
707 if !log_file.exists() {
708 continue;
709 }
710
711 let daemon_id = match DaemonId::from_safe_path(&file_name) {
712 Ok(id) => id,
713 Err(_) => continue,
714 };
715
716 if daemon_id == DaemonId::pitchfork() {
719 continue;
720 }
721
722 match store.migrate_daemon_text_logs(&daemon_id) {
723 Ok(0) => {}
724 Ok(n) => {
725 total_migrated += n;
726 migrated_ids.push(daemon_id.qualified());
727 }
728 Err(e) => {
729 warn!(
730 "failed to migrate text logs for {}: {e}",
731 daemon_id.qualified()
732 );
733 }
734 }
735 }
736
737 if total_migrated > 0 {
738 warn!(
739 "auto-migrated {total_migrated} legacy log entries from {count} daemon(s): {ids}",
740 count = migrated_ids.len(),
741 ids = migrated_ids.join(", ")
742 );
743 }
744
745 Ok(())
746}