1use crate::cli::json_output::{JsonLogEntry, print_json};
2use crate::daemon_id::DaemonId;
3use crate::log_store::sqlite::LOG_STORE;
4use crate::log_store::{FieldFilter, LogEntry, 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::fmt::Write as _;
16use std::io::{self, IsTerminal, Write};
17use std::process::{Child, Command, Stdio};
18use std::time::Duration;
19
20struct PagerConfig {
22 command: String,
23 args: Vec<String>,
24}
25
26impl PagerConfig {
27 fn new(start_at_end: bool) -> Self {
30 let command = std::env::var("PAGER").unwrap_or_else(|_| "less".to_string());
31 let args = Self::build_args(&command, start_at_end);
32 Self { command, args }
33 }
34
35 fn build_args(pager: &str, start_at_end: bool) -> Vec<String> {
36 let mut args = vec![];
37 if pager == "less" {
38 args.push("-R".to_string());
39 if start_at_end {
40 args.push("+G".to_string());
41 }
42 }
43 args
44 }
45
46 fn spawn_piped(&self) -> io::Result<Child> {
48 Command::new(&self.command)
49 .args(&self.args)
50 .stdin(Stdio::piped())
51 .spawn()
52 }
53}
54
55const KNOWN_FIELD_KEYS: &[&str] = &[
58 "level",
59 "severity",
60 "lvl",
61 "PRIORITY",
62 "@level",
63 "msg",
64 "message",
65 "event",
66 "@message",
67 "logger",
68 "name",
69 "component",
70 "module",
71 "timestamp",
72 "ts",
73 "time",
74 "@timestamp",
75];
76
77fn level_badge(level: &str) -> String {
82 let label = match level {
83 "error" => console::style("ERR").red().bold(),
84 "warn" => console::style("WRN").yellow().bold(),
85 "info" => console::style("INF").cyan().bold(),
86 "debug" => console::style("DBG").magenta().bold(),
87 "trace" => console::style("TRC").dim().bold(),
88 _ => return String::new(),
89 };
90 format!("{}{}{}", ndim("["), label, ndim("]"))
91}
92
93fn format_field_value(value: &serde_json::Value) -> String {
99 match value {
100 serde_json::Value::String(s) => {
101 let needs_quotes = s.is_empty()
102 || s.contains(' ')
103 || s.contains('=')
104 || s.contains('\t')
105 || s == "true"
106 || s == "false"
107 || s == "null"
108 || s.parse::<f64>().is_ok();
109 if needs_quotes {
110 format!("\"{s}\"")
111 } else {
112 s.clone()
113 }
114 }
115 serde_json::Value::Number(n) => console::style(n).green().to_string(),
116 serde_json::Value::Bool(true) => console::style("true").yellow().to_string(),
117 serde_json::Value::Bool(false) => console::style("false").red().to_string(),
118 serde_json::Value::Null => console::style("null").dim().to_string(),
119 serde_json::Value::Object(_) | serde_json::Value::Array(_) => {
120 serde_json::to_string(value).unwrap_or_else(|_| value.to_string())
121 }
122 }
123}
124
125fn format_fields(fields_json: &str) -> String {
127 let Ok(serde_json::Value::Object(obj)) = serde_json::from_str(fields_json) else {
128 return String::new();
129 };
130 let mut parts = Vec::new();
131 for (key, value) in &obj {
132 if KNOWN_FIELD_KEYS.contains(&key.as_str()) {
133 continue;
134 }
135 let key_styled = console::style(key.as_str()).blue().to_string();
136 let value_styled = format_field_value(value);
137 parts.push(format!("{key_styled}={value_styled}"));
138 }
139 parts.join(" ")
140}
141
142#[allow(clippy::too_many_arguments)]
154fn write_formatted_log(
155 w: &mut dyn Write,
156 entry: &LogEntry,
157 date: &str,
158 single_daemon: bool,
159 strip_ansi: bool,
160 show_timestamp: bool,
161 raw: bool,
162) -> io::Result<()> {
163 let clean_msg = strip_pty_controls(&entry.message);
167 let raw_msg: std::borrow::Cow<'_, str> = if strip_ansi {
168 console::strip_ansi_codes(&clean_msg)
169 } else {
170 std::borrow::Cow::Owned(clean_msg)
171 };
172
173 let mut out = String::with_capacity(256);
174
175 if raw {
176 if show_timestamp {
180 out.push_str(&ndim(date).to_string());
181 out.push(' ');
182 }
183 if !single_daemon {
184 let colors_on = !strip_ansi && console::colors_enabled();
185 out.push_str(&colored_id_label(&entry.daemon_id, colors_on));
186 out.push(' ');
187 }
188 out.push_str(&raw_msg);
189 } else if entry.fields_json.is_some() {
190 let level = entry.level.as_deref();
192 let accent = match level {
193 Some("error") => "error",
194 Some("warn") => "warn",
195 _ => "",
196 };
197
198 if show_timestamp {
200 let ts = match accent {
201 "error" => console::style(date).red().to_string(),
202 "warn" => console::style(date).yellow().to_string(),
203 _ => ndim(date).to_string(),
204 };
205 out.push_str(&ts);
206 out.push(' ');
207 }
208
209 if !single_daemon {
211 let colors_on = !strip_ansi && console::colors_enabled();
212 out.push_str(&colored_id_label(&entry.daemon_id, colors_on));
213 out.push(' ');
214 }
215
216 let mut need_sep = false;
217
218 if let Some(lvl) = level {
220 let badge = level_badge(lvl);
221 if !badge.is_empty() {
222 out.push_str(&badge);
223 need_sep = true;
224 }
225 }
226
227 let sep = ndim(":").to_string();
229 let arrow = ndim(" > ").to_string();
230
231 if let Some(logger) = &entry.logger {
233 if need_sep {
234 out.push(' ');
235 }
236 out.push_str(&console::style(logger).italic().dim().to_string());
237 out.push_str(&sep);
238 need_sep = true;
239 }
240
241 let msg_cow = entry.msg.as_deref().filter(|s| !s.is_empty()).map(|s| {
243 let cleaned = strip_pty_controls(s);
244 if strip_ansi {
245 std::borrow::Cow::<str>::Owned(console::strip_ansi_codes(&cleaned).to_string())
246 } else {
247 std::borrow::Cow::<str>::Owned(cleaned)
248 }
249 });
250 let fields_str = entry
251 .fields_json
252 .as_deref()
253 .map(format_fields)
254 .filter(|s| !s.is_empty());
255
256 if msg_cow.is_some() || fields_str.is_some() {
257 if need_sep {
258 out.push(' ');
259 }
260 if let Some(msg) = &msg_cow {
261 let styled = match accent {
262 "error" => console::style(msg.as_ref()).red().bold().to_string(),
263 "warn" => console::style(msg.as_ref()).yellow().bold().to_string(),
264 _ => console::style(msg.as_ref()).bold().to_string(),
265 };
266 out.push_str(&styled);
267 if fields_str.is_some() {
268 out.push_str(&arrow);
269 }
270 }
271 if let Some(fields) = &fields_str {
272 out.push_str(fields);
273 }
274 } else if !need_sep {
275 out.push_str(&raw_msg);
278 }
279 } else {
280 if show_timestamp {
282 out.push_str(&ndim(date).to_string());
283 out.push(' ');
284 }
285 if !single_daemon {
286 let colors_on = !strip_ansi && console::colors_enabled();
287 out.push_str(&colored_id_label(&entry.daemon_id, colors_on));
288 out.push(' ');
289 }
290 out.push_str(&raw_msg);
291 }
292
293 out.push('\n');
294 w.write_all(out.as_bytes())
295}
296
297pub fn colored_id_label(id: &str, colors_enabled: bool) -> String {
300 if !colors_enabled {
301 return format!("[{}]", id);
302 }
303 let colors: [u8; 4] = [34, 35, 36, 32]; let mut h: usize = 0x811C_9DC5; for b in id.bytes() {
308 h = h.wrapping_mul(0x0100_0193).wrapping_add(b as usize);
309 }
310 let color = colors[h % colors.len()];
311 format!("\x1b[{color}m[{id}]\x1b[0m")
312}
313
314#[derive(Debug, clap::Args)]
316#[clap(
317 visible_alias = "l",
318 verbatim_doc_comment,
319 long_about = "\
320Displays logs for daemon(s)
321
322Shows logs from managed daemons. Logs are stored in the pitchfork logs directory
323and include timestamps for filtering.
324
325Examples:
326 pitchfork logs api Show all logs for 'api' (paged if needed)
327 pitchfork logs api worker Show logs for multiple daemons
328 pitchfork logs Show logs for all daemons
329 pitchfork logs api -n 50 Show last 50 lines
330 pitchfork logs api --follow Follow logs in real-time
331 pitchfork logs api --since '2024-01-15 10:00:00'
332 Show logs since a specific time (forward)
333 pitchfork logs api --since '10:30:00'
334 Show logs since 10:30:00 today
335 pitchfork logs api --since '10:30' --until '12:00'
336 Show logs since 10:30:00 until 12:00:00 today
337 pitchfork logs api --since 5min Show logs from last 5 minutes
338 pitchfork logs api --raw Output raw log lines without formatting
339 pitchfork logs api --raw -n 100 Output last 100 raw log lines
340 pitchfork logs api --clear Delete logs for 'api'
341 pitchfork logs --clear Delete logs for all daemons"
342)]
343pub struct Logs {
344 id: Vec<String>,
346
347 #[clap(short, long)]
349 clear: bool,
350
351 #[clap(short)]
356 n: Option<usize>,
357
358 #[clap(short = 't', short_alias = 'f', long, visible_alias = "follow")]
360 tail: bool,
361
362 #[clap(short = 's', long)]
369 since: Option<String>,
370
371 #[clap(short = 'u', long)]
377 until: Option<String>,
378
379 #[clap(long)]
381 no_pager: bool,
382
383 #[clap(long)]
385 raw: bool,
386
387 #[clap(long, conflicts_with = "raw", conflicts_with = "tail")]
389 json: bool,
390
391 #[clap(long)]
395 grep: Vec<String>,
396
397 #[clap(long)]
399 regex: Option<String>,
400
401 #[clap(long)]
403 case_sensitive: bool,
404
405 #[clap(long)]
411 level: Option<String>,
412
413 #[clap(long, value_name = "KEY=VALUE")]
418 field: Vec<String>,
419
420 #[clap(long, value_name = "EXPR")]
426 jq: Option<String>,
427
428 #[clap(long)]
430 no_timestamp: bool,
431}
432
433impl Logs {
434 pub async fn run(&self) -> Result<()> {
435 migrate_legacy_log_dirs();
436
437 let resolved_ids: Vec<DaemonId> = if self.id.is_empty() {
438 get_all_daemon_ids()?
439 } else {
440 PitchforkToml::resolve_ids(&self.id)?
441 };
442
443 if self.clear {
444 LOG_STORE.clear(&resolved_ids)?;
445 return Ok(());
446 }
447
448 let from = if let Some(since) = self.since.as_ref() {
449 Some(parse_time_input(since, true)?)
450 } else {
451 None
452 };
453 let to = if let Some(until) = self.until.as_ref() {
454 Some(parse_time_input(until, false)?)
455 } else {
456 None
457 };
458
459 let message_filters = self.build_message_filters()?;
460 let field_filters = self.build_field_filters()?;
461
462 let jq_filter = match self.jq.as_deref() {
464 Some(expr) => Some(crate::log_jq::JqFilter::new(expr)?),
465 None => None,
466 };
467
468 if self.json {
469 return self.output_json(
470 &resolved_ids,
471 from,
472 to,
473 message_filters,
474 field_filters,
475 jq_filter.as_ref(),
476 );
477 }
478
479 let single_daemon = resolved_ids.len() == 1 && !self.id.is_empty();
484 let show_timestamp = settings().logs.timestamp && !self.no_timestamp && !self.raw;
485 let has_time_filter = from.is_some() || to.is_some();
486
487 self.query_and_output(
488 &resolved_ids,
489 from,
490 to,
491 message_filters.clone(),
492 field_filters.clone(),
493 jq_filter.as_ref(),
494 single_daemon,
495 has_time_filter,
496 show_timestamp,
497 )?;
498
499 if self.tail {
500 tail_logs(
501 &resolved_ids,
502 single_daemon,
503 true,
504 message_filters,
505 field_filters,
506 jq_filter.as_ref(),
507 show_timestamp,
508 self.raw,
509 )
510 .await?;
511 }
512
513 Ok(())
514 }
515
516 fn build_message_filters(&self) -> Result<Vec<MessageFilter>> {
517 if self.case_sensitive && self.grep.is_empty() {
518 warn!("--case-sensitive has no effect without --grep");
519 }
520 let mut filters = Vec::new();
521 for pattern in &self.grep {
522 filters.push(MessageFilter::Contains {
523 pattern: pattern.clone(),
524 case_sensitive: self.case_sensitive,
525 });
526 }
527 if let Some(pattern) = self.regex.as_ref() {
528 let _ = regex::Regex::new(pattern)
531 .into_diagnostic()
532 .map_err(|e| miette::miette!("invalid regex pattern: {e}"))?;
533 filters.push(MessageFilter::Regex {
534 pattern: pattern.clone(),
535 });
536 }
537 Ok(filters)
538 }
539
540 fn build_field_filters(&self) -> Result<Vec<FieldFilter>> {
541 let mut filters = Vec::new();
542 if let Some(level) = self.level.as_ref() {
543 let normalized = crate::log_parse::normalize_level_str(level).ok_or_else(|| {
546 miette::miette!(
547 "invalid level '{level}'; expected one of: \
548 error, err, fatal, critical, panic, alert, emerg, \
549 warn, warning, info, inf, information, notice, \
550 debug, dbg, trace, trc"
551 )
552 })?;
553 filters.push(FieldFilter::LevelMin(normalized));
554 }
555 for pair in &self.field {
556 let (key, value) = pair
557 .split_once('=')
558 .ok_or_else(|| miette::miette!("--field expects KEY=VALUE, got: {pair}"))?;
559 if key.is_empty() {
560 miette::bail!("--field key cannot be empty: {pair}");
561 }
562 if !key
565 .chars()
566 .all(|c| c.is_alphanumeric() || c == '_' || c == '.')
567 {
568 miette::bail!(
569 "--field key may only contain alphanumeric, underscore, and dot; got: {key}"
570 );
571 }
572 filters.push(FieldFilter::FieldEq {
573 key: key.to_string(),
574 value: value.to_string(),
575 });
576 }
577 Ok(filters)
578 }
579
580 #[allow(clippy::too_many_arguments)]
586 fn query_and_output(
587 &self,
588 resolved_ids: &[DaemonId],
589 from: Option<DateTime<Local>>,
590 to: Option<DateTime<Local>>,
591 message_filters: Vec<MessageFilter>,
592 field_filters: Vec<FieldFilter>,
593 jq_filter: Option<&crate::log_jq::JqFilter>,
594 single_daemon: bool,
595 has_time_filter: bool,
596 show_timestamp: bool,
597 ) -> Result<()> {
598 let daemon_ids: Vec<String> = resolved_ids.iter().map(|id| id.qualified()).collect();
599
600 let opts = LogQuery {
601 daemon_ids,
602 from,
603 to,
604 limit: if !has_time_filter { self.n } else { None },
605 order_desc: !has_time_filter,
606 after_id: None,
607 message_filters,
608 field_filters,
609 include_structured: jq_filter.is_some() || !self.raw,
610 };
611 let mut entries = LOG_STORE.query(&opts)?;
612
613 if let Some(jq) = jq_filter {
615 entries = jq.filter(entries);
616 }
617
618 if has_time_filter {
620 if let Some(n) = self.n {
621 let len = entries.len();
622 if len > n {
623 entries = entries.split_off(len - n);
624 }
625 }
626 } else {
627 entries.reverse();
628 }
629
630 if entries.is_empty() {
631 return Ok(());
632 }
633
634 let strip_ansi = self.raw || !console::colors_enabled();
635 let use_pager = !self.tail && !self.no_pager && should_use_pager(entries.len());
636 let ts_format = &settings().logs.timestamp_format;
637
638 let mut date_buf = String::with_capacity(ts_format.len() + 6);
640
641 let mut write_entries = |w: &mut dyn Write| -> io::Result<()> {
642 for entry in &entries {
643 date_buf.clear();
644 write!(date_buf, "{}", entry.timestamp.format(ts_format))
645 .map_err(io::Error::other)?;
646 write_formatted_log(
647 w,
648 entry,
649 &date_buf,
650 single_daemon,
651 strip_ansi,
652 show_timestamp,
653 self.raw,
654 )?;
655 }
656 Ok(())
657 };
658
659 if use_pager {
660 let pager_config = PagerConfig::new(!has_time_filter);
661 match pager_config.spawn_piped() {
662 Ok(mut child) => {
663 if let Some(stdin) = child.stdin.as_mut() {
664 if let Err(e) = write_entries(stdin) {
665 debug!("pager write error: {e}");
666 }
667 let _ = child.wait();
668 } else {
669 debug!("Failed to get pager stdin, falling back to direct output");
670 let stdout = io::stdout();
671 let mut buf = io::BufWriter::new(stdout.lock());
672 write_entries(&mut buf).into_diagnostic()?;
673 }
674 }
675 Err(e) => {
676 debug!("Failed to spawn pager: {e}, falling back to direct output");
677 let stdout = io::stdout();
678 let mut buf = io::BufWriter::new(stdout.lock());
679 write_entries(&mut buf).into_diagnostic()?;
680 }
681 }
682 } else {
683 let stdout = io::stdout();
684 let mut buf = io::BufWriter::new(stdout.lock());
685 write_entries(&mut buf).into_diagnostic()?;
686 }
687
688 Ok(())
689 }
690
691 fn output_json(
692 &self,
693 resolved_ids: &[DaemonId],
694 from: Option<DateTime<Local>>,
695 to: Option<DateTime<Local>>,
696 message_filters: Vec<MessageFilter>,
697 field_filters: Vec<FieldFilter>,
698 jq_filter: Option<&crate::log_jq::JqFilter>,
699 ) -> Result<()> {
700 let daemon_ids: Vec<String> = resolved_ids.iter().map(|id| id.qualified()).collect();
701 let has_time_filter = from.is_some() || to.is_some();
702
703 let opts = LogQuery {
704 daemon_ids,
705 from,
706 to,
707 limit: if !has_time_filter { self.n } else { None },
708 order_desc: !has_time_filter,
709 after_id: None,
710 message_filters,
711 field_filters,
712 include_structured: true,
713 };
714 let entries = LOG_STORE.query(&opts)?;
715
716 let entries = match jq_filter {
718 Some(jq) => jq.filter(entries),
719 None => entries,
720 };
721
722 let mut entries: Vec<_> = if has_time_filter {
726 entries
727 } else {
728 entries.into_iter().rev().collect()
729 };
730
731 if has_time_filter
734 && let Some(n) = self.n
735 && entries.len() > n
736 {
737 entries = entries.split_off(entries.len() - n);
738 }
739
740 let json_entries: Vec<JsonLogEntry> = entries
741 .into_iter()
742 .map(|e| {
743 let fields = e
744 .fields_json
745 .as_deref()
746 .and_then(|s| serde_json::from_str(s).ok());
747 JsonLogEntry {
748 timestamp: e.timestamp.format("%Y-%m-%d %H:%M:%S").to_string(),
749 daemon_id: e.daemon_id,
750 message: console::strip_ansi_codes(&e.message).to_string(),
751 level: e.level,
752 msg: e.msg,
753 logger: e.logger,
754 fields,
755 }
756 })
757 .collect();
758
759 print_json(&json_entries)
760 }
761}
762
763fn should_use_pager(line_count: usize) -> bool {
764 if !io::stdout().is_terminal() {
765 return false;
766 }
767
768 let terminal_height = get_terminal_height().unwrap_or(24);
769 line_count > terminal_height
770}
771
772fn get_terminal_height() -> Option<usize> {
773 if let Ok(rows) = std::env::var("LINES")
774 && let Ok(h) = rows.parse::<usize>()
775 {
776 return Some(h);
777 }
778
779 crossterm::terminal::size().ok().map(|(_, h)| h as usize)
780}
781
782fn migrate_legacy_log_dirs() {
792 let known_safe_paths = known_daemon_safe_paths();
793 let dirs = match xx::file::ls(&*env::PITCHFORK_LOGS_DIR) {
794 Ok(d) => d,
795 Err(_) => return,
796 };
797 for dir in dirs {
798 if dir.starts_with(".") || !dir.is_dir() {
799 continue;
800 }
801 let name = match dir.file_name().map(|f| f.to_string_lossy().to_string()) {
802 Some(n) => n,
803 None => continue,
804 };
805 if name == "pitchfork" {
807 continue;
808 }
809 if name.contains("--") {
812 if DaemonId::from_safe_path(&name).is_ok() {
815 continue;
816 }
817 if known_safe_paths.contains(&name) {
820 continue;
821 }
822 warn!(
823 "Skipping invalid legacy log directory '{name}': contains '--' but is not a valid daemon safe-path"
824 );
825 continue;
826 }
827
828 let old_log = dir.join(format!("{name}.log"));
831 if !old_log.exists() {
832 continue;
833 }
834 if DaemonId::try_new("legacy", &name).is_err() {
835 warn!("Skipping invalid legacy log directory '{name}': not a valid daemon ID");
836 continue;
837 }
838
839 let new_name = format!("legacy--{name}");
840 let new_dir = env::PITCHFORK_LOGS_DIR.join(&new_name);
841 if new_dir.exists() {
843 continue;
844 }
845 if std::fs::rename(&dir, &new_dir).is_err() {
846 continue;
847 }
848 let old_log = new_dir.join(format!("{name}.log"));
850 let new_log = new_dir.join(format!("{new_name}.log"));
851 if old_log.exists() {
852 let _ = std::fs::rename(&old_log, &new_log);
853 }
854 debug!("Migrated legacy log dir '{name}' → '{new_name}'");
855 }
856}
857
858fn known_daemon_safe_paths() -> BTreeSet<String> {
859 let mut out = BTreeSet::new();
860
861 match StateFile::read(&*env::PITCHFORK_STATE_FILE) {
862 Ok(state) => {
863 for id in state.daemons.keys() {
864 out.insert(id.safe_path());
865 }
866 }
867 Err(e) => {
868 warn!("Failed to read state while checking known daemon IDs: {e}");
869 }
870 }
871
872 match PitchforkToml::all_merged() {
873 Ok(config) => {
874 for id in config.daemons.keys() {
875 out.insert(id.safe_path());
876 }
877 }
878 Err(e) => {
879 warn!("Failed to read config while checking known daemon IDs: {e}");
880 }
881 }
882
883 out
884}
885
886fn get_all_daemon_ids() -> Result<Vec<DaemonId>> {
887 let mut ids = BTreeSet::new();
888
889 match StateFile::read(&*env::PITCHFORK_STATE_FILE) {
890 Ok(state) => ids.extend(state.daemons.keys().cloned()),
891 Err(e) => warn!("Failed to read state for log daemon discovery: {e}"),
892 }
893
894 match PitchforkToml::all_merged() {
895 Ok(config) => ids.extend(config.daemons.keys().cloned()),
896 Err(e) => warn!("Failed to read config for log daemon discovery: {e}"),
897 }
898
899 let logged_ids: std::collections::HashSet<String> =
900 LOG_STORE.list_daemon_ids()?.into_iter().collect();
901 Ok(ids
902 .into_iter()
903 .filter(|id| logged_ids.contains(&id.qualified()))
904 .collect())
905}
906
907#[allow(clippy::too_many_arguments)]
908pub async fn tail_logs(
909 names: &[DaemonId],
910 single_daemon: bool,
911 start_from_end: bool,
912 message_filters: Vec<MessageFilter>,
913 field_filters: Vec<FieldFilter>,
914 jq_filter: Option<&crate::log_jq::JqFilter>,
915 show_timestamp: bool,
916 raw: bool,
917) -> Result<()> {
918 let strip_ansi = raw || !console::colors_enabled();
920
921 let mut states: std::collections::HashMap<String, i64> = names
922 .iter()
923 .map(|id| {
924 let since = if start_from_end {
925 LOG_STORE.last_id(id).unwrap_or(None).unwrap_or(0)
929 } else {
930 0
931 };
932 (id.qualified(), since)
933 })
934 .collect();
935
936 let interval = tokio::time::interval(Duration::from_millis(200));
937 tokio::pin!(interval);
938
939 loop {
940 interval.tick().await;
941
942 let mut out = vec![];
943 for id in names {
944 let after_id = states.get(&id.qualified()).copied();
945 match LOG_STORE.query(&LogQuery {
946 daemon_ids: vec![id.qualified()],
947 from: None,
948 to: None,
949 limit: None,
950 order_desc: false,
951 after_id,
952 message_filters: message_filters.clone(),
953 field_filters: field_filters.clone(),
954 include_structured: jq_filter.is_some() || !raw,
955 }) {
956 Ok(raw_entries) => {
957 let last_raw_id = raw_entries.last().map(|e| e.id);
961
962 let entries = match jq_filter {
964 Some(jq) => jq.filter(raw_entries),
965 None => raw_entries,
966 };
967 out.extend(entries);
968 let has_sql_filter = !message_filters.is_empty() || !field_filters.is_empty();
980 let new_cursor = if let Some(id) = last_raw_id {
981 Some(id)
982 } else if has_sql_filter || jq_filter.is_some() {
983 LOG_STORE.last_id(id).ok().flatten()
984 } else {
985 None
986 };
987 if let Some(last_id) = new_cursor {
988 states.insert(id.qualified(), last_id);
989 }
990 }
991 Err(e) => {
992 error!("Failed to tail logs for {}: {e}", id.qualified());
993 }
994 }
995 }
996
997 if !out.is_empty() {
998 let out: Vec<LogEntry> = if single_daemon {
1001 out
1002 } else {
1003 out.into_iter()
1004 .sorted_by(|a, b| a.timestamp.cmp(&b.timestamp).then(a.id.cmp(&b.id)))
1005 .collect()
1006 };
1007 let stdout = io::stdout();
1008 let mut buf = io::BufWriter::new(stdout.lock());
1009 let ts_format = &settings().logs.timestamp_format;
1010 let mut date_buf = String::with_capacity(ts_format.len() + 6);
1011 for entry in &out {
1012 date_buf.clear();
1013 write!(date_buf, "{}", entry.timestamp.format(ts_format))
1014 .map_err(io::Error::other)
1015 .into_diagnostic()?;
1016 write_formatted_log(
1017 &mut buf,
1018 entry,
1019 &date_buf,
1020 single_daemon,
1021 strip_ansi,
1022 show_timestamp,
1023 raw,
1024 )
1025 .into_diagnostic()?;
1026 }
1027 }
1028 }
1029}
1030
1031fn parse_datetime(s: &str) -> Result<DateTime<Local>> {
1032 let naive_dt = NaiveDateTime::parse_from_str(s, "%Y-%m-%d %H:%M:%S").into_diagnostic()?;
1033 Local
1034 .from_local_datetime(&naive_dt)
1035 .single()
1036 .ok_or_else(|| miette::miette!("Invalid or ambiguous datetime: '{}'. ", s))
1037}
1038
1039fn parse_time_input(s: &str, is_since: bool) -> Result<DateTime<Local>> {
1045 let s = s.trim();
1046
1047 if let Ok(dt) = parse_datetime(s) {
1049 return Ok(dt);
1050 }
1051
1052 if let Ok(naive_dt) = NaiveDateTime::parse_from_str(s, "%Y-%m-%d %H:%M") {
1054 return Local
1055 .from_local_datetime(&naive_dt)
1056 .single()
1057 .ok_or_else(|| miette::miette!("Invalid or ambiguous datetime: '{}'", s));
1058 }
1059
1060 if let Ok(time) = parse_time_only(s) {
1064 let now = Local::now();
1065 let today = now.date_naive();
1066 let mut naive_dt = NaiveDateTime::new(today, time);
1067 let mut dt = Local
1068 .from_local_datetime(&naive_dt)
1069 .single()
1070 .ok_or_else(|| miette::miette!("Invalid or ambiguous datetime: '{}'", s))?;
1071
1072 if is_since
1075 && dt > now
1076 && let Some(yesterday) = today.pred_opt()
1077 {
1078 naive_dt = NaiveDateTime::new(yesterday, time);
1079 dt = Local
1080 .from_local_datetime(&naive_dt)
1081 .single()
1082 .ok_or_else(|| miette::miette!("Invalid or ambiguous datetime: '{}'", s))?;
1083 }
1084 return Ok(dt);
1085 }
1086
1087 if let Ok(duration) = humantime::parse_duration(s) {
1088 let now = Local::now();
1089 let target = now - chrono::Duration::from_std(duration).into_diagnostic()?;
1090 return Ok(target);
1091 }
1092
1093 Err(miette::miette!(
1094 "Invalid time format: '{}'. Expected formats:\n\
1095 - Full datetime: \"YYYY-MM-DD HH:MM:SS\" or \"YYYY-MM-DD HH:MM\"\n\
1096 - Time only: \"HH:MM:SS\" or \"HH:MM\" (uses today's date)\n\
1097 - Relative time: \"5min\", \"2h\", \"1d\" (e.g., last 5 minutes)",
1098 s
1099 ))
1100}
1101
1102fn parse_time_only(s: &str) -> Result<NaiveTime> {
1103 if let Ok(time) = NaiveTime::parse_from_str(s, "%H:%M:%S") {
1104 return Ok(time);
1105 }
1106
1107 if let Ok(time) = NaiveTime::parse_from_str(s, "%H:%M") {
1108 return Ok(time);
1109 }
1110
1111 Err(miette::miette!("Invalid time format: '{}'", s))
1112}
1113
1114pub fn print_error_logs_block(log_lines: &[(String, String, String)]) {
1124 if log_lines.is_empty() {
1125 return;
1126 }
1127
1128 let is_tty = std::io::stderr().is_terminal();
1129 let format_msg = |msg: &str| -> String {
1130 let stripped = strip_pty_controls(msg);
1131 if is_tty {
1132 stripped
1133 } else {
1134 console::strip_ansi_codes(&stripped).to_string()
1135 }
1136 };
1137
1138 let tag = estyle(" ERROR LOGS ").white().on_red();
1139 eprintln!("\n{tag}");
1140
1141 let unique_ids: BTreeSet<&str> = log_lines.iter().map(|(_, id, _)| id.as_str()).collect();
1143 let show_id = unique_ids.len() > 1;
1144
1145 if show_id {
1146 for (date, id, msg) in log_lines {
1147 let time = date.split(' ').nth(1).unwrap_or(date);
1148 let colored = colored_id_label(id, is_tty && console::colors_enabled_stderr());
1149 eprintln!(
1150 "{} {} {}",
1151 estyle(time).red().dim(),
1152 colored,
1153 format_msg(msg)
1154 );
1155 }
1156 } else {
1157 for (date, _, msg) in log_lines {
1158 let time = date.split(' ').nth(1).unwrap_or(date);
1159 eprintln!("{} {}", estyle(time).red().dim(), format_msg(msg));
1160 }
1161 }
1162}
1163
1164pub enum ReadyCheckType {
1166 Output(String),
1167 Http(String),
1168 Port(u16),
1169 Cmd(String),
1170 Delay(u64),
1171 Default,
1172}
1173
1174impl std::fmt::Display for ReadyCheckType {
1175 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1176 match self {
1177 ReadyCheckType::Output(pattern) => write!(f, "output matching '{pattern}'"),
1178 ReadyCheckType::Http(url) => write!(f, "HTTP {url}"),
1179 ReadyCheckType::Port(port) => write!(f, "TCP port {port}"),
1180 ReadyCheckType::Cmd(cmd) => write!(f, "command '{cmd}'"),
1181 ReadyCheckType::Delay(secs) => write!(f, "delay ({secs}s)"),
1182 ReadyCheckType::Default => write!(f, "default readiness check"),
1183 }
1184 }
1185}
1186
1187pub fn create_ready_check_job(
1193 daemon_id: &DaemonId,
1194 check_type: &ReadyCheckType,
1195) -> std::sync::Arc<clx::progress::ProgressJob> {
1196 use clx::progress::{ProgressJobBuilder, ProgressJobDoneBehavior, ProgressStatus};
1197
1198 let is_tty = std::io::stderr().is_terminal();
1199 let colors_enabled = is_tty && console::colors_enabled_stderr();
1200 let id_label = colored_id_label(&daemon_id.qualified(), colors_enabled);
1201 let show_ts = crate::settings::settings().general.startup_log_timestamps;
1202
1203 let prefix = if show_ts {
1207 edim(chrono::Local::now().format("%H:%M:%S").to_string()).to_string()
1210 } else {
1211 "{{spinner()}}".to_string()
1212 };
1213
1214 ProgressJobBuilder::new()
1215 .body(format!(
1216 "{} {} waiting for {{{{ check_type }}}}...",
1217 prefix, id_label
1218 ))
1219 .prop("check_type", &check_type.to_string())
1220 .status(ProgressStatus::Running)
1221 .on_done(ProgressJobDoneBehavior::Keep)
1222 .start()
1223}
1224
1225pub fn collect_startup_logs(
1230 daemon_id: &DaemonId,
1231 from: DateTime<Local>,
1232) -> Result<Vec<(String, String, String)>> {
1233 let entries = LOG_STORE.query(&LogQuery {
1234 daemon_ids: vec![daemon_id.qualified()],
1235 from: Some(from),
1236 to: None,
1237 limit: None,
1238 order_desc: false,
1239 after_id: None,
1240 message_filters: Vec::new(),
1241 field_filters: Vec::new(),
1242 include_structured: false,
1243 })?;
1244 let log_lines = entries
1245 .into_iter()
1246 .map(|e| {
1247 let ts = e.timestamp.format("%Y-%m-%d %H:%M:%S").to_string();
1248 (ts, e.daemon_id, e.message)
1249 })
1250 .collect();
1251
1252 Ok(log_lines)
1253}
1254
1255pub fn stream_startup_logs(
1261 daemon_id: &DaemonId,
1262 job: std::sync::Arc<clx::progress::ProgressJob>,
1263) -> (
1264 tokio::sync::watch::Sender<bool>,
1265 tokio::task::JoinHandle<()>,
1266) {
1267 let (tx, mut rx) = tokio::sync::watch::channel(false);
1268 let id = daemon_id.clone();
1269
1270 let show_ts = crate::settings::settings().general.startup_log_timestamps;
1271
1272 let anchor_id: i64 = LOG_STORE
1277 .query(&LogQuery {
1278 daemon_ids: vec![id.qualified()],
1279 limit: Some(1),
1280 order_desc: true,
1281 ..Default::default()
1282 })
1283 .ok()
1284 .and_then(|entries| entries.last().map(|e| e.id))
1285 .unwrap_or(0);
1286
1287 let handle = tokio::spawn(async move {
1288 let is_tty = std::io::stderr().is_terminal();
1289 let colors_enabled = is_tty && console::colors_enabled_stderr();
1290 let id_label = colored_id_label(&id.qualified(), colors_enabled);
1291 let prefix = if show_ts {
1292 String::new()
1293 } else {
1294 edim("•").to_string()
1295 };
1296
1297 let mut last_id = anchor_id;
1298
1299 if let Ok(entries) = LOG_STORE.tail(&id, Some(last_id)) {
1301 for entry in &entries {
1302 let time = entry.timestamp.format("%H:%M:%S").to_string();
1303 let msg = strip_pty_controls(&entry.message);
1304 let msg = if is_tty {
1305 msg
1306 } else {
1307 console::strip_ansi_codes(&msg).to_string()
1308 };
1309 let line_prefix = if show_ts {
1310 edim(time).to_string()
1311 } else {
1312 prefix.clone()
1313 };
1314 job.println(&format!("{} {} {}", line_prefix, id_label, msg));
1315 }
1316 if let Some(last) = entries.last() {
1317 last_id = last.id;
1318 }
1319 }
1320
1321 loop {
1322 tokio::select! {
1323 _ = tokio::time::sleep(Duration::from_millis(200)) => {
1324 if let Ok(entries) = LOG_STORE.tail(&id, Some(last_id)) {
1325 for entry in &entries {
1326 let time = entry.timestamp.format("%H:%M:%S").to_string();
1327 let msg = strip_pty_controls(&entry.message);
1328 let msg = if is_tty {
1329 msg
1330 } else {
1331 console::strip_ansi_codes(&msg).to_string()
1332 };
1333 let line_prefix = if show_ts {
1334 edim(time).to_string()
1335 } else {
1336 prefix.clone()
1337 };
1338 job.println(&format!("{} {} {}", line_prefix, id_label, msg));
1339 }
1340 if let Some(last) = entries.last() {
1341 last_id = last.id;
1342 }
1343 }
1344 }
1345 _ = rx.changed() => {
1346 break;
1347 }
1348 }
1349 }
1350
1351 if let Ok(entries) = LOG_STORE.tail(&id, Some(last_id)) {
1353 for entry in &entries {
1354 let time = entry.timestamp.format("%H:%M:%S").to_string();
1355 let msg = strip_pty_controls(&entry.message);
1356 let msg = if is_tty {
1357 msg
1358 } else {
1359 console::strip_ansi_codes(&msg).to_string()
1360 };
1361 let line_prefix = if show_ts {
1362 edim(time).to_string()
1363 } else {
1364 prefix.clone()
1365 };
1366 job.println(&format!("{} {} {}", line_prefix, id_label, msg));
1367 }
1368 }
1369 });
1370
1371 (tx, handle)
1372}
1373
1374fn strip_pty_controls(s: &str) -> String {
1379 struct Stripper {
1380 result: String,
1381 }
1382
1383 impl vte::Perform for Stripper {
1384 fn print(&mut self, c: char) {
1385 self.result.push(c);
1386 }
1387
1388 fn execute(&mut self, byte: u8) {
1389 if byte == b'\n' || byte == b'\t' {
1391 self.result.push(byte as char);
1392 }
1393 }
1394
1395 fn csi_dispatch(
1396 &mut self,
1397 params: &vte::Params,
1398 _intermediates: &[u8],
1399 _ignore: bool,
1400 action: char,
1401 ) {
1402 if action == 'm' {
1404 self.result.push_str("\x1b[");
1405 let mut first = true;
1406 for sub in params.iter() {
1407 if !first {
1408 self.result.push(';');
1409 }
1410 first = false;
1411 for (i, &p) in sub.iter().enumerate() {
1412 if i > 0 {
1413 self.result.push(':');
1414 }
1415 self.result.push_str(&p.to_string());
1416 }
1417 }
1418 self.result.push('m');
1419 }
1420 }
1422
1423 fn osc_dispatch(&mut self, _params: &[&[u8]], _bell_terminated: bool) {
1424 }
1426
1427 fn esc_dispatch(&mut self, _intermediates: &[u8], _ignore: bool, _byte: u8) {
1428 }
1430
1431 fn hook(
1432 &mut self,
1433 _params: &vte::Params,
1434 _intermediates: &[u8],
1435 _ignore: bool,
1436 _action: char,
1437 ) {
1438 }
1440
1441 fn put(&mut self, _byte: u8) {
1442 }
1444
1445 fn unhook(&mut self) {
1446 }
1448 }
1449
1450 let mut parser = vte::Parser::new();
1451 let mut stripper = Stripper {
1452 result: String::with_capacity(s.len()),
1453 };
1454 parser.advance(&mut stripper, s.as_bytes());
1455 stripper.result
1456}