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
19struct PagerConfig {
21 command: String,
22 args: Vec<String>,
23}
24
25impl PagerConfig {
26 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 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
54fn 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
91fn 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), (180, 160, 100), (120, 180, 120), (120, 180, 180), (180, 120, 180), (120, 160, 180), ];
106 let mut h: usize = 0x811C_9DC5; 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
114pub fn colored_id_label(id: &str, colors_enabled: bool) -> String {
117 if !colors_enabled {
118 return format!("[{}]", id);
119 }
120 let colors: [u8; 4] = [34, 35, 36, 32]; let mut h: usize = 0x811C_9DC5; 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#[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 id: Vec<String>,
163
164 #[clap(short, long)]
166 clear: bool,
167
168 #[clap(short)]
173 n: Option<usize>,
174
175 #[clap(short = 't', short_alias = 'f', long, visible_alias = "follow")]
177 tail: bool,
178
179 #[clap(short = 's', long)]
186 since: Option<String>,
187
188 #[clap(short = 'u', long)]
194 until: Option<String>,
195
196 #[clap(long)]
198 no_pager: bool,
199
200 #[clap(long)]
202 raw: bool,
203
204 #[clap(long, conflicts_with = "raw", conflicts_with = "tail")]
206 json: bool,
207
208 #[clap(long)]
212 grep: Vec<String>,
213
214 #[clap(long)]
216 regex: Option<String>,
217
218 #[clap(long)]
220 case_sensitive: bool,
221
222 #[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 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 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
541fn 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 if name == "pitchfork" {
566 continue;
567 }
568 if name.contains("--") {
571 if DaemonId::from_safe_path(&name).is_ok() {
574 continue;
575 }
576 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 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 if new_dir.exists() {
602 continue;
603 }
604 if std::fs::rename(&dir, &new_dir).is_err() {
605 continue;
606 }
607 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 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 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 let new_cursor = if !entries.is_empty() {
734 entries.last().map(|e| e.id)
735 } else if message_filters.is_empty() {
736 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
782fn parse_time_input(s: &str, is_since: bool) -> Result<DateTime<Local>> {
788 let s = s.trim();
789
790 if let Ok(dt) = parse_datetime(s) {
792 return Ok(dt);
793 }
794
795 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 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 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
857pub 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 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
913pub 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
936pub 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 let prefix = if show_ts {
956 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
974pub 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
1002pub 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 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 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 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
1121fn 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 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 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 }
1169
1170 fn osc_dispatch(&mut self, _params: &[&[u8]], _bell_terminated: bool) {
1171 }
1173
1174 fn esc_dispatch(&mut self, _intermediates: &[u8], _ignore: bool, _byte: u8) {
1175 }
1177
1178 fn hook(
1179 &mut self,
1180 _params: &vte::Params,
1181 _intermediates: &[u8],
1182 _ignore: bool,
1183 _action: char,
1184 ) {
1185 }
1187
1188 fn put(&mut self, _byte: u8) {
1189 }
1191
1192 fn unhook(&mut self) {
1193 }
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}