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