Skip to main content

pitchfork_cli/cli/
logs.rs

1use crate::cli::json_output::{JsonLogEntry, print_json};
2use crate::daemon_id::DaemonId;
3use crate::log_store::sqlite::LOG_STORE;
4use crate::log_store::{LogQuery, LogStore, MessageFilter};
5use crate::pitchfork_toml::PitchforkToml;
6use crate::settings::settings;
7use crate::state_file::StateFile;
8use crate::ui::style::{edim, estyle, ndim};
9use crate::{Result, env};
10use chrono::{DateTime, Local, NaiveDateTime, NaiveTime, TimeZone};
11use console;
12use itertools::Itertools;
13use miette::IntoDiagnostic;
14use std::collections::BTreeSet;
15use std::io::{self, IsTerminal, Write};
16use std::process::{Child, Command, Stdio};
17use std::time::Duration;
18
19/// Pager configuration for displaying logs
20struct PagerConfig {
21    command: String,
22    args: Vec<String>,
23}
24
25impl PagerConfig {
26    /// Select and configure the appropriate pager.
27    /// Uses $PAGER environment variable if set, otherwise defaults to less.
28    fn new(start_at_end: bool) -> Self {
29        let command = std::env::var("PAGER").unwrap_or_else(|_| "less".to_string());
30        let args = Self::build_args(&command, start_at_end);
31        Self { command, args }
32    }
33
34    fn build_args(pager: &str, start_at_end: bool) -> Vec<String> {
35        let mut args = vec![];
36        if pager == "less" {
37            args.push("-R".to_string());
38            if start_at_end {
39                args.push("+G".to_string());
40            }
41        }
42        args
43    }
44
45    /// Spawn the pager with piped stdin
46    fn spawn_piped(&self) -> io::Result<Child> {
47        Command::new(&self.command)
48            .args(&self.args)
49            .stdin(Stdio::piped())
50            .spawn()
51    }
52}
53
54/// Format a single log line for output.
55/// When `single_daemon` is true, omits the daemon ID from the output.
56/// `id_width` is the display width used to pad the daemon name column
57/// so messages line up vertically across different daemon names.
58/// When `strip_ansi` is true, strips ANSI escape codes from the message.
59fn format_log_line(
60    date: &str,
61    id: &str,
62    msg: &str,
63    single_daemon: bool,
64    id_width: usize,
65    strip_ansi: bool,
66    show_timestamp: bool,
67) -> String {
68    let msg = if strip_ansi {
69        console::strip_ansi_codes(msg).to_string()
70    } else {
71        msg.to_string()
72    };
73    if single_daemon {
74        if show_timestamp {
75            format!("{} {}", ndim(date), msg)
76        } else {
77            msg
78        }
79    } else {
80        let colors_on = !strip_ansi && console::colors_enabled();
81        let colored = dimmed_id(id, colors_on);
82        let padded = console::pad_str(&colored, id_width, console::Alignment::Left, None);
83        if show_timestamp {
84            format!("{}  {} {}", padded, ndim(date), msg)
85        } else {
86            format!("{}  {}", padded, msg)
87        }
88    }
89}
90
91/// Return a dimmed, colorized daemon ID string for display.
92/// Each daemon gets a deterministic color via FNV-1a hash so that
93/// multiple daemons are visually distinguishable while remaining subtle.
94fn dimmed_id(id: &str, colors_enabled: bool) -> String {
95    if !colors_enabled {
96        return id.to_string();
97    }
98    let colors = [
99        (180, 120, 120), // dim red
100        (180, 160, 100), // dim yellow
101        (120, 180, 120), // dim green
102        (120, 180, 180), // dim cyan
103        (180, 120, 180), // dim magenta
104        (120, 160, 180), // dim blue
105    ];
106    let mut h: usize = 0x811C_9DC5; // FNV offset basis
107    for b in id.bytes() {
108        h = h.wrapping_mul(0x0100_0193).wrapping_add(b as usize);
109    }
110    let (r, g, b) = colors[h % colors.len()];
111    format!("\x1b[2;38;2;{};{};{}m{}\x1b[0m", r, g, b, id)
112}
113
114/// Return a colorized `[namespace/id]` label for display in progress jobs.
115/// Uses brighter colors than `dimmed_id` and includes the square brackets.
116pub fn colored_id_label(id: &str, colors_enabled: bool) -> String {
117    if !colors_enabled {
118        return format!("[{}]", id);
119    }
120    // Same palette as mise: Blue, Magenta, Cyan, Green
121    // Excludes Red/Yellow to avoid confusion with errors/warnings.
122    let colors: [u8; 4] = [34, 35, 36, 32]; // ANSI: Blue, Magenta, Cyan, Green
123    let mut h: usize = 0x811C_9DC5; // FNV offset basis
124    for b in id.bytes() {
125        h = h.wrapping_mul(0x0100_0193).wrapping_add(b as usize);
126    }
127    let color = colors[h % colors.len()];
128    format!("\x1b[{color}m[{id}]\x1b[0m")
129}
130
131/// Displays logs for daemon(s)
132#[derive(Debug, clap::Args)]
133#[clap(
134    visible_alias = "l",
135    verbatim_doc_comment,
136    long_about = "\
137Displays logs for daemon(s)
138
139Shows logs from managed daemons. Logs are stored in the pitchfork logs directory
140and include timestamps for filtering.
141
142Examples:
143  pitchfork logs api              Show all logs for 'api' (paged if needed)
144  pitchfork logs api worker       Show logs for multiple daemons
145  pitchfork logs                  Show logs for all daemons
146  pitchfork logs api -n 50        Show last 50 lines
147  pitchfork logs api --follow     Follow logs in real-time
148  pitchfork logs api --since '2024-01-15 10:00:00'
149                                  Show logs since a specific time (forward)
150  pitchfork logs api --since '10:30:00'
151                                  Show logs since 10:30:00 today
152  pitchfork logs api --since '10:30' --until '12:00'
153                                  Show logs since 10:30:00 until 12:00:00 today
154  pitchfork logs api --since 5min Show logs from last 5 minutes
155  pitchfork logs api --raw        Output raw log lines without formatting
156  pitchfork logs api --raw -n 100 Output last 100 raw log lines
157  pitchfork logs api --clear      Delete logs for 'api'
158  pitchfork logs --clear          Delete logs for all daemons"
159)]
160pub struct Logs {
161    /// Show only logs for the specified daemon(s)
162    id: Vec<String>,
163
164    /// Delete logs
165    #[clap(short, long)]
166    clear: bool,
167
168    /// Show last N lines of logs
169    ///
170    /// Only applies when --since/--until is not used.
171    /// Without this option, all logs are shown.
172    #[clap(short)]
173    n: Option<usize>,
174
175    /// Show logs in real-time
176    #[clap(short = 't', short_alias = 'f', long, visible_alias = "follow")]
177    tail: bool,
178
179    /// Show logs from this time
180    ///
181    /// Supports multiple formats:
182    /// - Full datetime: "YYYY-MM-DD HH:MM:SS" or "YYYY-MM-DD HH:MM"
183    /// - Time only: "HH:MM:SS" or "HH:MM" (uses today's date)
184    /// - Relative time: "5min", "2h", "1d" (e.g., last 5 minutes)
185    #[clap(short = 's', long)]
186    since: Option<String>,
187
188    /// Show logs until this time
189    ///
190    /// Supports multiple formats:
191    /// - Full datetime: "YYYY-MM-DD HH:MM:SS" or "YYYY-MM-DD HH:MM"
192    /// - Time only: "HH:MM:SS" or "HH:MM" (uses today's date)
193    #[clap(short = 'u', long)]
194    until: Option<String>,
195
196    /// Disable pager even in interactive terminal
197    #[clap(long)]
198    no_pager: bool,
199
200    /// Output raw log lines without color or formatting
201    #[clap(long)]
202    raw: bool,
203
204    /// Output in JSON format
205    #[clap(long, conflicts_with = "raw", conflicts_with = "tail")]
206    json: bool,
207
208    /// Filter logs by case-insensitive substring (can be repeated)
209    ///
210    /// Multiple --grep options are combined with OR.
211    #[clap(long)]
212    grep: Vec<String>,
213
214    /// Filter logs by regular expression
215    #[clap(long)]
216    regex: Option<String>,
217
218    /// Make --grep matching case-sensitive
219    #[clap(long)]
220    case_sensitive: bool,
221
222    /// Omit timestamps from log output
223    #[clap(long)]
224    no_timestamp: bool,
225}
226
227impl Logs {
228    pub async fn run(&self) -> Result<()> {
229        migrate_legacy_log_dirs();
230
231        let resolved_ids: Vec<DaemonId> = if self.id.is_empty() {
232            get_all_daemon_ids()?
233        } else {
234            PitchforkToml::resolve_ids(&self.id)?
235        };
236
237        if self.clear {
238            LOG_STORE.clear(&resolved_ids)?;
239            return Ok(());
240        }
241
242        let from = if let Some(since) = self.since.as_ref() {
243            Some(parse_time_input(since, true)?)
244        } else {
245            None
246        };
247        let to = if let Some(until) = self.until.as_ref() {
248            Some(parse_time_input(until, false)?)
249        } else {
250            None
251        };
252
253        let message_filters = self.build_message_filters()?;
254
255        if self.json {
256            return self.output_json(&resolved_ids, from, to, message_filters);
257        }
258
259        let single_daemon = resolved_ids.len() == 1;
260        let show_timestamp = settings().logs.timestamp && !self.no_timestamp;
261        let log_lines = self.fetch_log_lines(&resolved_ids, from, to, message_filters.clone())?;
262        let has_time_filter = from.is_some() || to.is_some();
263        self.output_logs(
264            log_lines,
265            single_daemon,
266            has_time_filter,
267            self.tail,
268            show_timestamp,
269        )?;
270        if self.tail {
271            tail_logs(
272                &resolved_ids,
273                single_daemon,
274                true,
275                message_filters,
276                show_timestamp,
277            )
278            .await?;
279        }
280
281        Ok(())
282    }
283
284    fn build_message_filters(&self) -> Result<Vec<MessageFilter>> {
285        if self.case_sensitive && self.grep.is_empty() {
286            warn!("--case-sensitive has no effect without --grep");
287        }
288        let mut filters = Vec::new();
289        for pattern in &self.grep {
290            filters.push(MessageFilter::Contains {
291                pattern: pattern.clone(),
292                case_sensitive: self.case_sensitive,
293            });
294        }
295        if let Some(pattern) = self.regex.as_ref() {
296            // Validate the regex early so the user gets a clear CLI error
297            // instead of a SQLite user-function failure at query time.
298            let _ = regex::Regex::new(pattern)
299                .into_diagnostic()
300                .map_err(|e| miette::miette!("invalid regex pattern: {e}"))?;
301            filters.push(MessageFilter::Regex {
302                pattern: pattern.clone(),
303            });
304        }
305        Ok(filters)
306    }
307
308    fn fetch_log_lines(
309        &self,
310        resolved_ids: &[DaemonId],
311        from: Option<DateTime<Local>>,
312        to: Option<DateTime<Local>>,
313        message_filters: Vec<MessageFilter>,
314    ) -> Result<Vec<(String, String, String)>> {
315        let daemon_ids: Vec<String> = resolved_ids.iter().map(|id| id.qualified()).collect();
316        let has_time_filter = from.is_some() || to.is_some();
317
318        let opts = LogQuery {
319            daemon_ids: daemon_ids.clone(),
320            from,
321            to,
322            limit: if !has_time_filter { self.n } else { None },
323            order_desc: !has_time_filter,
324            after_id: None,
325            message_filters,
326        };
327        let entries = LOG_STORE.query(&opts)?;
328        let log_lines: Vec<(String, String, String)> = entries
329            .into_iter()
330            .map(|e| {
331                let ts = e.timestamp.format("%Y-%m-%d %H:%M:%S").to_string();
332                (ts, e.daemon_id, e.message)
333            })
334            .collect();
335
336        let log_lines = if has_time_filter {
337            if let Some(n) = self.n {
338                let len = log_lines.len();
339                if len > n {
340                    log_lines.into_iter().skip(len - n).collect_vec()
341                } else {
342                    log_lines
343                }
344            } else {
345                log_lines
346            }
347        } else if let Some(n) = self.n {
348            let len = log_lines.len();
349            if len > n {
350                log_lines.into_iter().skip(len - n).rev().collect_vec()
351            } else {
352                log_lines.into_iter().rev().collect_vec()
353            }
354        } else {
355            log_lines.into_iter().rev().collect_vec()
356        };
357
358        Ok(log_lines)
359    }
360
361    fn output_json(
362        &self,
363        resolved_ids: &[DaemonId],
364        from: Option<DateTime<Local>>,
365        to: Option<DateTime<Local>>,
366        message_filters: Vec<MessageFilter>,
367    ) -> Result<()> {
368        let log_lines = self.fetch_log_lines(resolved_ids, from, to, message_filters)?;
369
370        let json_entries: Vec<JsonLogEntry> = log_lines
371            .into_iter()
372            .map(|(timestamp, daemon_id, message)| JsonLogEntry {
373                timestamp,
374                daemon_id,
375                message: console::strip_ansi_codes(&message).to_string(),
376            })
377            .collect();
378
379        print_json(&json_entries)
380    }
381
382    fn output_logs(
383        &self,
384        log_lines: Vec<(String, String, String)>,
385        single_daemon: bool,
386        has_time_filter: bool,
387        force_no_pager: bool,
388        show_timestamp: bool,
389    ) -> Result<()> {
390        if log_lines.is_empty() {
391            return Ok(());
392        }
393
394        let id_width = log_lines
395            .iter()
396            .map(|(_, id, _)| id.len())
397            .max()
398            .unwrap_or(0);
399        let strip_ansi = self.raw || !console::colors_enabled();
400
401        if self.raw {
402            for (date, id, msg) in log_lines {
403                let line = format_log_line(
404                    &date,
405                    &id,
406                    &msg,
407                    single_daemon,
408                    id_width,
409                    strip_ansi,
410                    show_timestamp,
411                );
412                println!("{line}");
413            }
414            return Ok(());
415        }
416
417        let use_pager = !force_no_pager && !self.no_pager && should_use_pager(log_lines.len());
418
419        if use_pager {
420            self.output_with_pager(
421                log_lines,
422                single_daemon,
423                id_width,
424                has_time_filter,
425                strip_ansi,
426                show_timestamp,
427            )?;
428        } else {
429            for (date, id, msg) in log_lines {
430                println!(
431                    "{}",
432                    format_log_line(
433                        &date,
434                        &id,
435                        &msg,
436                        single_daemon,
437                        id_width,
438                        strip_ansi,
439                        show_timestamp,
440                    )
441                );
442            }
443        }
444
445        Ok(())
446    }
447
448    fn output_with_pager(
449        &self,
450        log_lines: Vec<(String, String, String)>,
451        single_daemon: bool,
452        id_width: usize,
453        has_time_filter: bool,
454        strip_ansi: bool,
455        show_timestamp: bool,
456    ) -> Result<()> {
457        // When time filter is used, start at top; otherwise start at end
458        let pager_config = PagerConfig::new(!has_time_filter);
459
460        match pager_config.spawn_piped() {
461            Ok(mut child) => {
462                if let Some(stdin) = child.stdin.as_mut() {
463                    for (date, id, msg) in log_lines {
464                        let line = format!(
465                            "{}\n",
466                            format_log_line(
467                                &date,
468                                &id,
469                                &msg,
470                                single_daemon,
471                                id_width,
472                                strip_ansi,
473                                show_timestamp,
474                            )
475                        );
476                        if stdin.write_all(line.as_bytes()).is_err() {
477                            break;
478                        }
479                    }
480                    let _ = child.wait();
481                } else {
482                    debug!("Failed to get pager stdin, falling back to direct output");
483                    for (date, id, msg) in log_lines {
484                        println!(
485                            "{}",
486                            format_log_line(
487                                &date,
488                                &id,
489                                &msg,
490                                single_daemon,
491                                id_width,
492                                strip_ansi,
493                                show_timestamp,
494                            )
495                        );
496                    }
497                }
498            }
499            Err(e) => {
500                debug!("Failed to spawn pager: {e}, falling back to direct output");
501                for (date, id, msg) in log_lines {
502                    println!(
503                        "{}",
504                        format_log_line(
505                            &date,
506                            &id,
507                            &msg,
508                            single_daemon,
509                            id_width,
510                            strip_ansi,
511                            show_timestamp,
512                        )
513                    );
514                }
515            }
516        }
517
518        Ok(())
519    }
520}
521
522fn should_use_pager(line_count: usize) -> bool {
523    if !io::stdout().is_terminal() {
524        return false;
525    }
526
527    let terminal_height = get_terminal_height().unwrap_or(24);
528    line_count > terminal_height
529}
530
531fn get_terminal_height() -> Option<usize> {
532    if let Ok(rows) = std::env::var("LINES")
533        && let Ok(h) = rows.parse::<usize>()
534    {
535        return Some(h);
536    }
537
538    crossterm::terminal::size().ok().map(|(_, h)| h as usize)
539}
540
541/// Rename legacy log directories that predate namespace-qualified daemon IDs.
542///
543/// Old layout: `PITCHFORK_LOGS_DIR/<name>/<name>.log`
544/// New layout: `PITCHFORK_LOGS_DIR/legacy--<name>/legacy--<name>.log`
545///
546/// Only directories that clearly match the old layout are migrated:
547/// - directory name does not contain `"--"`
548/// - directory contains `<name>.log`
549/// - `<name>` is a valid daemon short name under current DaemonId rules
550fn migrate_legacy_log_dirs() {
551    let known_safe_paths = known_daemon_safe_paths();
552    let dirs = match xx::file::ls(&*env::PITCHFORK_LOGS_DIR) {
553        Ok(d) => d,
554        Err(_) => return,
555    };
556    for dir in dirs {
557        if dir.starts_with(".") || !dir.is_dir() {
558            continue;
559        }
560        let name = match dir.file_name().map(|f| f.to_string_lossy().to_string()) {
561            Some(n) => n,
562            None => continue,
563        };
564        // Skip the supervisor's own log directory.
565        if name == "pitchfork" {
566            continue;
567        }
568        // New-format directories usually contain "--". For safety, only treat
569        // them as new-format if they match a known daemon ID safe-path.
570        if name.contains("--") {
571            // If it parses as a valid safe-path, treat it as already migrated
572            // and keep idempotent behavior silent.
573            if DaemonId::from_safe_path(&name).is_ok() {
574                continue;
575            }
576            // Keep noisy warnings only for invalid/ambiguous names that cannot
577            // be interpreted as new-format IDs.
578            if known_safe_paths.contains(&name) {
579                continue;
580            }
581            warn!(
582                "Skipping invalid legacy log directory '{name}': contains '--' but is not a valid daemon safe-path"
583            );
584            continue;
585        }
586
587        // Migrate only explicit old-layout directories to avoid renaming
588        // unrelated folders under logs/.
589        let old_log = dir.join(format!("{name}.log"));
590        if !old_log.exists() {
591            continue;
592        }
593        if DaemonId::try_new("legacy", &name).is_err() {
594            warn!("Skipping invalid legacy log directory '{name}': not a valid daemon ID");
595            continue;
596        }
597
598        let new_name = format!("legacy--{name}");
599        let new_dir = env::PITCHFORK_LOGS_DIR.join(&new_name);
600        // Skip if a target directory already exists to avoid clobbering data.
601        if new_dir.exists() {
602            continue;
603        }
604        if std::fs::rename(&dir, &new_dir).is_err() {
605            continue;
606        }
607        // Also rename the log file inside the directory.
608        let old_log = new_dir.join(format!("{name}.log"));
609        let new_log = new_dir.join(format!("{new_name}.log"));
610        if old_log.exists() {
611            let _ = std::fs::rename(&old_log, &new_log);
612        }
613        debug!("Migrated legacy log dir '{name}' → '{new_name}'");
614    }
615}
616
617fn known_daemon_safe_paths() -> BTreeSet<String> {
618    let mut out = BTreeSet::new();
619
620    match StateFile::read(&*env::PITCHFORK_STATE_FILE) {
621        Ok(state) => {
622            for id in state.daemons.keys() {
623                out.insert(id.safe_path());
624            }
625        }
626        Err(e) => {
627            warn!("Failed to read state while checking known daemon IDs: {e}");
628        }
629    }
630
631    match PitchforkToml::all_merged() {
632        Ok(config) => {
633            for id in config.daemons.keys() {
634                out.insert(id.safe_path());
635            }
636        }
637        Err(e) => {
638            warn!("Failed to read config while checking known daemon IDs: {e}");
639        }
640    }
641
642    out
643}
644
645fn get_all_daemon_ids() -> Result<Vec<DaemonId>> {
646    let mut ids = BTreeSet::new();
647
648    match StateFile::read(&*env::PITCHFORK_STATE_FILE) {
649        Ok(state) => ids.extend(state.daemons.keys().cloned()),
650        Err(e) => warn!("Failed to read state for log daemon discovery: {e}"),
651    }
652
653    match PitchforkToml::all_merged() {
654        Ok(config) => ids.extend(config.daemons.keys().cloned()),
655        Err(e) => warn!("Failed to read config for log daemon discovery: {e}"),
656    }
657
658    let logged_ids: std::collections::HashSet<String> =
659        LOG_STORE.list_daemon_ids()?.into_iter().collect();
660    Ok(ids
661        .into_iter()
662        .filter(|id| logged_ids.contains(&id.qualified()))
663        .collect())
664}
665
666pub async fn tail_logs(
667    names: &[DaemonId],
668    single_daemon: bool,
669    start_from_end: bool,
670    message_filters: Vec<MessageFilter>,
671    show_timestamp: bool,
672) -> Result<()> {
673    // Poll SQLite log store for new entries since last known row id.
674    let id_width = names
675        .iter()
676        .map(|id| id.qualified().len())
677        .max()
678        .unwrap_or(0);
679
680    let strip_ansi = !console::colors_enabled();
681
682    let mut states: std::collections::HashMap<String, i64> = names
683        .iter()
684        .map(|id| {
685            let since = if start_from_end {
686                // Anchor to the last entry overall, not the last filtered entry,
687                // so --tail combined with a filter does not rescan from row 1
688                // on every poll when no message matches yet.
689                LOG_STORE.last_id(id).unwrap_or(None).unwrap_or(0)
690            } else {
691                0
692            };
693            (id.qualified(), since)
694        })
695        .collect();
696
697    let interval = tokio::time::interval(Duration::from_millis(200));
698    tokio::pin!(interval);
699
700    loop {
701        interval.tick().await;
702
703        let mut out = vec![];
704        for id in names {
705            let after_id = states.get(&id.qualified()).copied();
706            match LOG_STORE.query(&LogQuery {
707                daemon_ids: vec![id.qualified()],
708                from: None,
709                to: None,
710                limit: None,
711                order_desc: false,
712                after_id,
713                message_filters: message_filters.clone(),
714            }) {
715                Ok(entries) => {
716                    for entry in &entries {
717                        let ts = entry.timestamp.format("%Y-%m-%d %H:%M:%S").to_string();
718                        out.push((ts, entry.daemon_id.clone(), entry.message.clone()));
719                    }
720                    // Advance the cursor past rows already evaluated.
721                    //
722                    // When the query returned entries, entries.last().id is the
723                    // highest row id examined (within the same lock hold as the
724                    // query, so no race with concurrent append_batch). Using it
725                    // directly avoids a separate last_id call that could skip
726                    // matching rows written between the two lock acquisitions.
727                    //
728                    // When the filter is active and no entries matched, fall
729                    // back to last_id to advance past non-matching rows so they
730                    // are not re-evaluated on every poll. The narrow race here
731                    // (a matching row written between query and last_id) is
732                    // accepted as the cost of bounded re-scanning.
733                    let new_cursor = if !entries.is_empty() {
734                        entries.last().map(|e| e.id)
735                    } else if message_filters.is_empty() {
736                        // No filter and no entries: nothing to advance past.
737                        None
738                    } else {
739                        LOG_STORE.last_id(id).ok().flatten()
740                    };
741                    if let Some(last_id) = new_cursor {
742                        states.insert(id.qualified(), last_id);
743                    }
744                }
745                Err(e) => {
746                    error!("Failed to tail logs for {}: {e}", id.qualified());
747                }
748            }
749        }
750
751        if !out.is_empty() {
752            let out = out
753                .into_iter()
754                .sorted_by(|a, b| (&a.0, &a.1).cmp(&(&b.0, &b.1)))
755                .collect_vec();
756            for (date, name, msg) in out {
757                println!(
758                    "{}",
759                    format_log_line(
760                        &date,
761                        &name,
762                        &msg,
763                        single_daemon,
764                        id_width,
765                        strip_ansi,
766                        show_timestamp,
767                    )
768                );
769            }
770        }
771    }
772}
773
774fn parse_datetime(s: &str) -> Result<DateTime<Local>> {
775    let naive_dt = NaiveDateTime::parse_from_str(s, "%Y-%m-%d %H:%M:%S").into_diagnostic()?;
776    Local
777        .from_local_datetime(&naive_dt)
778        .single()
779        .ok_or_else(|| miette::miette!("Invalid or ambiguous datetime: '{}'. ", s))
780}
781
782/// Parse time input string into DateTime.
783///
784/// `is_since` indicates whether this is for --since (true) or --until (false).
785/// The "yesterday fallback" only applies to --since: if the time is in the future,
786/// assume the user meant yesterday. For --until, future times are kept as-is.
787fn parse_time_input(s: &str, is_since: bool) -> Result<DateTime<Local>> {
788    let s = s.trim();
789
790    // Try full datetime first (YYYY-MM-DD HH:MM:SS)
791    if let Ok(dt) = parse_datetime(s) {
792        return Ok(dt);
793    }
794
795    // Try datetime without seconds (YYYY-MM-DD HH:MM)
796    if let Ok(naive_dt) = NaiveDateTime::parse_from_str(s, "%Y-%m-%d %H:%M") {
797        return Local
798            .from_local_datetime(&naive_dt)
799            .single()
800            .ok_or_else(|| miette::miette!("Invalid or ambiguous datetime: '{}'", s));
801    }
802
803    // Try time-only format (HH:MM:SS or HH:MM)
804    // Note: This branch won't be reached for inputs like "10:30" that could match
805    // parse_datetime, because parse_datetime expects a full date prefix and will fail.
806    if let Ok(time) = parse_time_only(s) {
807        let now = Local::now();
808        let today = now.date_naive();
809        let mut naive_dt = NaiveDateTime::new(today, time);
810        let mut dt = Local
811            .from_local_datetime(&naive_dt)
812            .single()
813            .ok_or_else(|| miette::miette!("Invalid or ambiguous datetime: '{}'", s))?;
814
815        // If the interpreted time for today is in the future, assume the user meant yesterday
816        // BUT only for --since. For --until, a future time today is valid.
817        if is_since
818            && dt > now
819            && let Some(yesterday) = today.pred_opt()
820        {
821            naive_dt = NaiveDateTime::new(yesterday, time);
822            dt = Local
823                .from_local_datetime(&naive_dt)
824                .single()
825                .ok_or_else(|| miette::miette!("Invalid or ambiguous datetime: '{}'", s))?;
826        }
827        return Ok(dt);
828    }
829
830    if let Ok(duration) = humantime::parse_duration(s) {
831        let now = Local::now();
832        let target = now - chrono::Duration::from_std(duration).into_diagnostic()?;
833        return Ok(target);
834    }
835
836    Err(miette::miette!(
837        "Invalid time format: '{}'. Expected formats:\n\
838         - Full datetime: \"YYYY-MM-DD HH:MM:SS\" or \"YYYY-MM-DD HH:MM\"\n\
839         - Time only: \"HH:MM:SS\" or \"HH:MM\" (uses today's date)\n\
840         - Relative time: \"5min\", \"2h\", \"1d\" (e.g., last 5 minutes)",
841        s
842    ))
843}
844
845fn parse_time_only(s: &str) -> Result<NaiveTime> {
846    if let Ok(time) = NaiveTime::parse_from_str(s, "%H:%M:%S") {
847        return Ok(time);
848    }
849
850    if let Ok(time) = NaiveTime::parse_from_str(s, "%H:%M") {
851        return Ok(time);
852    }
853
854    Err(miette::miette!("Invalid time format: '{}'", s))
855}
856
857/// Prints error log lines in a styled block matching the startup logs format.
858///
859/// Format:
860/// ```text
861///  ERROR LOGS
862///  12:00:00 error message
863/// ```
864///
865/// Timestamps use dimmed red. The tag uses white text on red background.
866pub fn print_error_logs_block(log_lines: &[(String, String, String)]) {
867    if log_lines.is_empty() {
868        return;
869    }
870
871    let is_tty = std::io::stderr().is_terminal();
872    let format_msg = |msg: &str| -> String {
873        let stripped = strip_pty_controls(msg);
874        if is_tty {
875            stripped
876        } else {
877            console::strip_ansi_codes(&stripped).to_string()
878        }
879    };
880
881    let tag = estyle(" ERROR LOGS ").white().on_red();
882    eprintln!("\n{tag}");
883
884    // Determine if we need to show daemon IDs (same logic as startup logs)
885    let unique_ids: BTreeSet<&str> = log_lines.iter().map(|(_, id, _)| id.as_str()).collect();
886    let show_id = unique_ids.len() > 1;
887
888    if show_id {
889        let id_width = log_lines
890            .iter()
891            .map(|(_, id, _)| console::measure_text_width(id))
892            .max()
893            .unwrap_or(0);
894        for (date, id, msg) in log_lines {
895            let time = date.split(' ').nth(1).unwrap_or(date);
896            let colored = dimmed_id(id, is_tty && console::colors_enabled_stderr());
897            let padded = console::pad_str(&colored, id_width, console::Alignment::Left, None);
898            eprintln!(
899                "{}  {} {}",
900                padded,
901                estyle(time).red().dim(),
902                format_msg(msg)
903            );
904        }
905    } else {
906        for (date, _, msg) in log_lines {
907            let time = date.split(' ').nth(1).unwrap_or(date);
908            eprintln!("{} {}", estyle(time).red().dim(), format_msg(msg));
909        }
910    }
911}
912
913/// Describes the type of ready check being performed for display purposes.
914pub enum ReadyCheckType {
915    Output(String),
916    Http(String),
917    Port(u16),
918    Cmd(String),
919    Delay(u64),
920    Default,
921}
922
923impl std::fmt::Display for ReadyCheckType {
924    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
925        match self {
926            ReadyCheckType::Output(pattern) => write!(f, "output matching '{pattern}'"),
927            ReadyCheckType::Http(url) => write!(f, "HTTP {url}"),
928            ReadyCheckType::Port(port) => write!(f, "TCP port {port}"),
929            ReadyCheckType::Cmd(cmd) => write!(f, "command '{cmd}'"),
930            ReadyCheckType::Delay(secs) => write!(f, "delay ({secs}s)"),
931            ReadyCheckType::Default => write!(f, "default readiness check"),
932        }
933    }
934}
935
936/// Creates a progress job showing a spinner while waiting for a ready check.
937///
938/// Returns a `Arc<ProgressJob>` that the caller should update:
939/// - Set body to success message and status to `Done` when the daemon is ready
940/// - Set body to failure message and status to `Failed` when the daemon fails
941pub fn create_ready_check_job(
942    daemon_id: &DaemonId,
943    check_type: &ReadyCheckType,
944) -> std::sync::Arc<clx::progress::ProgressJob> {
945    use clx::progress::{ProgressJobBuilder, ProgressJobDoneBehavior, ProgressStatus};
946
947    let is_tty = std::io::stderr().is_terminal();
948    let colors_enabled = is_tty && console::colors_enabled_stderr();
949    let id_label = colored_id_label(&daemon_id.qualified(), colors_enabled);
950    let show_ts = crate::settings::settings().general.startup_log_timestamps;
951
952    // When timestamps are off, {{spinner()}} renders as an animated spinner
953    // (1 char wide) matching the "•" prefix used by println.  When on,
954    // we show a dim timestamp instead.
955    let prefix = if show_ts {
956        // The timestamp updates each refresh via the now() tera function.
957        // We use a fixed-width format (HH:MM:SS = 8 chars) for alignment.
958        edim(chrono::Local::now().format("%H:%M:%S").to_string()).to_string()
959    } else {
960        "{{spinner()}}".to_string()
961    };
962
963    ProgressJobBuilder::new()
964        .body(format!(
965            "{} {} waiting for {{{{ check_type }}}}...",
966            prefix, id_label
967        ))
968        .prop("check_type", &check_type.to_string())
969        .status(ProgressStatus::Running)
970        .on_done(ProgressJobDoneBehavior::Keep)
971        .start()
972}
973
974/// Collects startup log lines for a single daemon (does not print).
975///
976/// Returns a list of `(time, daemon_id_qualified, message)` tuples for log
977/// entries written after `from`.
978pub fn collect_startup_logs(
979    daemon_id: &DaemonId,
980    from: DateTime<Local>,
981) -> Result<Vec<(String, String, String)>> {
982    let entries = LOG_STORE.query(&LogQuery {
983        daemon_ids: vec![daemon_id.qualified()],
984        from: Some(from),
985        to: None,
986        limit: None,
987        order_desc: false,
988        after_id: None,
989        message_filters: Vec::new(),
990    })?;
991    let log_lines = entries
992        .into_iter()
993        .map(|e| {
994            let ts = e.timestamp.format("%Y-%m-%d %H:%M:%S").to_string();
995            (ts, e.daemon_id, e.message)
996        })
997        .collect();
998
999    Ok(log_lines)
1000}
1001
1002/// Stream startup logs for a daemon to a progress job in real-time.
1003///
1004/// Spawns a background tokio task that polls the daemon's log store
1005/// and calls `job.println()` for each new line. Returns a watch sender
1006/// that stops the streaming when sent `true`.
1007pub fn stream_startup_logs(
1008    daemon_id: &DaemonId,
1009    job: std::sync::Arc<clx::progress::ProgressJob>,
1010) -> (
1011    tokio::sync::watch::Sender<bool>,
1012    tokio::task::JoinHandle<()>,
1013) {
1014    let (tx, mut rx) = tokio::sync::watch::channel(false);
1015    let id = daemon_id.clone();
1016
1017    let show_ts = crate::settings::settings().general.startup_log_timestamps;
1018
1019    // Anchor to the daemon's current max log id *synchronously* before
1020    // spawning the streaming task. This must happen before ipc.run() starts
1021    // the daemon, otherwise early output could be written to the log store
1022    // before the anchor is established and get skipped as "already seen".
1023    let anchor_id: i64 = LOG_STORE
1024        .query(&LogQuery {
1025            daemon_ids: vec![id.qualified()],
1026            limit: Some(1),
1027            order_desc: true,
1028            ..Default::default()
1029        })
1030        .ok()
1031        .and_then(|entries| entries.last().map(|e| e.id))
1032        .unwrap_or(0);
1033
1034    let handle = tokio::spawn(async move {
1035        let is_tty = std::io::stderr().is_terminal();
1036        let colors_enabled = is_tty && console::colors_enabled_stderr();
1037        let id_label = colored_id_label(&id.qualified(), colors_enabled);
1038        let prefix = if show_ts {
1039            String::new()
1040        } else {
1041            edim("•").to_string()
1042        };
1043
1044        let mut last_id = anchor_id;
1045
1046        // Initial fetch: catch any logs already written since the anchor.
1047        if let Ok(entries) = LOG_STORE.tail(&id, Some(last_id)) {
1048            for entry in &entries {
1049                let time = entry.timestamp.format("%H:%M:%S").to_string();
1050                let msg = strip_pty_controls(&entry.message);
1051                let msg = if is_tty {
1052                    msg
1053                } else {
1054                    console::strip_ansi_codes(&msg).to_string()
1055                };
1056                let line_prefix = if show_ts {
1057                    edim(time).to_string()
1058                } else {
1059                    prefix.clone()
1060                };
1061                job.println(&format!("{} {} {}", line_prefix, id_label, msg));
1062            }
1063            if let Some(last) = entries.last() {
1064                last_id = last.id;
1065            }
1066        }
1067
1068        loop {
1069            tokio::select! {
1070                _ = tokio::time::sleep(Duration::from_millis(200)) => {
1071                    if let Ok(entries) = LOG_STORE.tail(&id, Some(last_id)) {
1072                        for entry in &entries {
1073                            let time = entry.timestamp.format("%H:%M:%S").to_string();
1074                            let msg = strip_pty_controls(&entry.message);
1075                            let msg = if is_tty {
1076                                msg
1077                            } else {
1078                                console::strip_ansi_codes(&msg).to_string()
1079                            };
1080                            let line_prefix = if show_ts {
1081                                edim(time).to_string()
1082                            } else {
1083                                prefix.clone()
1084                            };
1085                            job.println(&format!("{} {} {}", line_prefix, id_label, msg));
1086                        }
1087                        if let Some(last) = entries.last() {
1088                            last_id = last.id;
1089                        }
1090                    }
1091                }
1092                _ = rx.changed() => {
1093                    break;
1094                }
1095            }
1096        }
1097
1098        // Final drain
1099        if let Ok(entries) = LOG_STORE.tail(&id, Some(last_id)) {
1100            for entry in &entries {
1101                let time = entry.timestamp.format("%H:%M:%S").to_string();
1102                let msg = strip_pty_controls(&entry.message);
1103                let msg = if is_tty {
1104                    msg
1105                } else {
1106                    console::strip_ansi_codes(&msg).to_string()
1107                };
1108                let line_prefix = if show_ts {
1109                    edim(time).to_string()
1110                } else {
1111                    prefix.clone()
1112                };
1113                job.println(&format!("{} {} {}", line_prefix, id_label, msg));
1114            }
1115        }
1116    });
1117
1118    (tx, handle)
1119}
1120
1121/// Strips PTY control sequences from a string while preserving SGR (color/style) codes.
1122///
1123/// Removes CSI sequences that control cursor movement, screen clearing, erasing, etc.,
1124/// but keeps `\x1b[...m` (SGR) sequences so colors are retained.
1125fn strip_pty_controls(s: &str) -> String {
1126    struct Stripper {
1127        result: String,
1128    }
1129
1130    impl vte::Perform for Stripper {
1131        fn print(&mut self, c: char) {
1132            self.result.push(c);
1133        }
1134
1135        fn execute(&mut self, byte: u8) {
1136            // Keep \n and \t; drop other control characters (BEL, BS, CR, etc.)
1137            if byte == b'\n' || byte == b'\t' {
1138                self.result.push(byte as char);
1139            }
1140        }
1141
1142        fn csi_dispatch(
1143            &mut self,
1144            params: &vte::Params,
1145            _intermediates: &[u8],
1146            _ignore: bool,
1147            action: char,
1148        ) {
1149            // Keep SGR sequences (final byte 'm')
1150            if action == 'm' {
1151                self.result.push_str("\x1b[");
1152                let mut first = true;
1153                for sub in params.iter() {
1154                    if !first {
1155                        self.result.push(';');
1156                    }
1157                    first = false;
1158                    for (i, &p) in sub.iter().enumerate() {
1159                        if i > 0 {
1160                            self.result.push(':');
1161                        }
1162                        self.result.push_str(&p.to_string());
1163                    }
1164                }
1165                self.result.push('m');
1166            }
1167            // All other CSI sequences (cursor move, clear, erase, etc.) are dropped
1168        }
1169
1170        fn osc_dispatch(&mut self, _params: &[&[u8]], _bell_terminated: bool) {
1171            // Drop OSC sequences (e.g. window title)
1172        }
1173
1174        fn esc_dispatch(&mut self, _intermediates: &[u8], _ignore: bool, _byte: u8) {
1175            // Drop ESC sequences (e.g. ESC c = reset terminal)
1176        }
1177
1178        fn hook(
1179            &mut self,
1180            _params: &vte::Params,
1181            _intermediates: &[u8],
1182            _ignore: bool,
1183            _action: char,
1184        ) {
1185            // Drop DCS hooks
1186        }
1187
1188        fn put(&mut self, _byte: u8) {
1189            // Drop DCS data
1190        }
1191
1192        fn unhook(&mut self) {
1193            // Drop DCS unhook
1194        }
1195    }
1196
1197    let mut parser = vte::Parser::new();
1198    let mut stripper = Stripper {
1199        result: String::with_capacity(s.len()),
1200    };
1201    parser.advance(&mut stripper, s.as_bytes());
1202    stripper.result
1203}