Skip to main content

pitchfork_cli/log_store/
sqlite.rs

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
13/// SQLite-backed log store with WAL mode for concurrent readers.
14pub struct SqliteLogStore {
15    conn: Mutex<Connection>,
16}
17
18impl SqliteLogStore {
19    /// Open or create the SQLite log store at the given path.
20    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            // Write entries to stdin. On error, kill and reap the child
113            // before propagating — Child::drop does not call wait(), so
114            // without this a crashed hook would leave a zombie process.
115            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                // Explicitly flush so a buffer-drain failure (e.g. the hook
134                // exited early and closed stdin) is surfaced as an error
135                // rather than being silently swallowed by BufWriter::drop.
136                // On a small batch where all data fits in the 8 KB buffer no
137                // individual writeln! has touched the underlying pipe yet, so
138                // without this flush the failure would be lost and the entries
139                // deleted without ever being delivered to the hook.
140                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                // BufWriter + ChildStdin drop here, closing stdin (EOF signal).
147            };
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    /// Delete rows matching the given IDs, returning the number of rows deleted.
169    ///
170    /// Chunks at 999 to stay within SQLite's `SQLITE_MAX_VARIABLE_NUMBER`
171    /// limit on versions ≤ 3.31 (e.g. Ubuntu 20.04 LTS).
172    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    /// Rotate (delete) old log entries for a specific daemon based on retention policy.
194    ///
195    /// When an archive hook is configured, entries are fetched in batches of
196    /// `batch_size`, the hook is invoked *without* holding the SQLite mutex
197    /// (so concurrent log appends are not blocked), and then deleted by ID.
198    /// Row-read errors are propagated rather than silently dropped.
199    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                // Fetch one batch under lock, then release.
212                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                // Run the hook without holding the mutex.
238                self.archive_entries(&entries, hook, daemon_id, "age")?;
239
240                // Re-acquire lock and delete exactly the archived IDs.
241                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    /// Rotate (delete) old log entries keeping only the most recent `max_count` rows
258    /// for a specific daemon.
259    ///
260    /// When an archive hook is configured, entries are fetched in batches of
261    /// `batch_size`, the hook is invoked *without* holding the SQLite mutex
262    /// (so concurrent log appends are not blocked), and then deleted by ID.
263    /// Row-read errors are propagated rather than silently dropped.
264    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        // Determine how many entries to delete.
273        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                // Fetch one batch under lock, then release.
296                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                // Run the hook without holding the mutex.
323                self.archive_entries(&entries, hook, daemon_id, "count")?;
324
325                // Re-acquire lock and delete exactly the archived IDs.
326                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    /// Migrate existing text logs for a daemon into SQLite.
347    ///
348    /// Reads the legacy text file line-by-line (streaming) and inserts in
349    /// batches of 1000 to avoid loading multi-GB files into memory at once.
350    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
637/// Global singleton log store.
638use 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    // Auto-migrate any legacy text log files into SQLite on first access.
654    // This runs once per process startup and is idempotent.
655    // Skip migration when using the in-memory fallback to prevent data loss.
656    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
669/// Auto-migrate legacy text log files into the SQLite log store.
670///
671/// Scans the logs directory for directories matching the new-format layout
672/// (`namespace--name/namespace--name.log`), attempts to parse the directory
673/// name as a valid safe-path daemon ID, and imports the content into SQLite.
674/// This is idempotent: re-running it on already-migrated data is a no-op
675/// because the legacy text files are deleted after successful import.
676fn 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        // Skip the supervisor's own log directory.
695        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        // Only consider directories that look like new-format safe-paths
703        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        // Skip the supervisor's own daemon; its log directory may be present
717        // under logs/ but should never be imported into the user-facing log store.
718        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}