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 before_id: None,
608 message_filters,
609 field_filters,
610 include_structured: jq_filter.is_some() || !self.raw,
611 };
612 let mut entries = LOG_STORE.query(&opts)?;
613
614 if let Some(jq) = jq_filter {
616 entries = jq.filter(entries);
617 }
618
619 if has_time_filter {
621 if let Some(n) = self.n {
622 let len = entries.len();
623 if len > n {
624 entries = entries.split_off(len - n);
625 }
626 }
627 } else {
628 entries.reverse();
629 }
630
631 if entries.is_empty() {
632 return Ok(());
633 }
634
635 let strip_ansi = self.raw || !console::colors_enabled();
636 let use_pager = !self.tail && !self.no_pager && should_use_pager(entries.len());
637 let ts_format = &settings().logs.timestamp_format;
638
639 let mut date_buf = String::with_capacity(ts_format.len() + 6);
641
642 let mut write_entries = |w: &mut dyn Write| -> io::Result<()> {
643 for entry in &entries {
644 date_buf.clear();
645 write!(date_buf, "{}", entry.timestamp.format(ts_format))
646 .map_err(io::Error::other)?;
647 write_formatted_log(
648 w,
649 entry,
650 &date_buf,
651 single_daemon,
652 strip_ansi,
653 show_timestamp,
654 self.raw,
655 )?;
656 }
657 Ok(())
658 };
659
660 if use_pager {
661 let pager_config = PagerConfig::new(!has_time_filter);
662 match pager_config.spawn_piped() {
663 Ok(mut child) => {
664 if let Some(stdin) = child.stdin.as_mut() {
665 if let Err(e) = write_entries(stdin) {
666 debug!("pager write error: {e}");
667 }
668 let _ = child.wait();
669 } else {
670 debug!("Failed to get pager stdin, falling back to direct output");
671 let stdout = io::stdout();
672 let mut buf = io::BufWriter::new(stdout.lock());
673 write_entries(&mut buf).into_diagnostic()?;
674 }
675 }
676 Err(e) => {
677 debug!("Failed to spawn pager: {e}, falling back to direct output");
678 let stdout = io::stdout();
679 let mut buf = io::BufWriter::new(stdout.lock());
680 write_entries(&mut buf).into_diagnostic()?;
681 }
682 }
683 } else {
684 let stdout = io::stdout();
685 let mut buf = io::BufWriter::new(stdout.lock());
686 write_entries(&mut buf).into_diagnostic()?;
687 }
688
689 Ok(())
690 }
691
692 fn output_json(
693 &self,
694 resolved_ids: &[DaemonId],
695 from: Option<DateTime<Local>>,
696 to: Option<DateTime<Local>>,
697 message_filters: Vec<MessageFilter>,
698 field_filters: Vec<FieldFilter>,
699 jq_filter: Option<&crate::log_jq::JqFilter>,
700 ) -> Result<()> {
701 let daemon_ids: Vec<String> = resolved_ids.iter().map(|id| id.qualified()).collect();
702 let has_time_filter = from.is_some() || to.is_some();
703
704 let opts = LogQuery {
705 daemon_ids,
706 from,
707 to,
708 limit: if !has_time_filter { self.n } else { None },
709 order_desc: !has_time_filter,
710 after_id: None,
711 before_id: None,
712 message_filters,
713 field_filters,
714 include_structured: true,
715 };
716 let entries = LOG_STORE.query(&opts)?;
717
718 let entries = match jq_filter {
720 Some(jq) => jq.filter(entries),
721 None => entries,
722 };
723
724 let mut entries: Vec<_> = if has_time_filter {
728 entries
729 } else {
730 entries.into_iter().rev().collect()
731 };
732
733 if has_time_filter
736 && let Some(n) = self.n
737 && entries.len() > n
738 {
739 entries = entries.split_off(entries.len() - n);
740 }
741
742 let json_entries: Vec<JsonLogEntry> = entries.into_iter().map(Into::into).collect();
743
744 print_json(&json_entries)
745 }
746}
747
748fn should_use_pager(line_count: usize) -> bool {
749 if !io::stdout().is_terminal() {
750 return false;
751 }
752
753 let terminal_height = get_terminal_height().unwrap_or(24);
754 line_count > terminal_height
755}
756
757fn get_terminal_height() -> Option<usize> {
758 if let Ok(rows) = std::env::var("LINES")
759 && let Ok(h) = rows.parse::<usize>()
760 {
761 return Some(h);
762 }
763
764 crossterm::terminal::size().ok().map(|(_, h)| h as usize)
765}
766
767fn migrate_legacy_log_dirs() {
777 let known_safe_paths = known_daemon_safe_paths();
778 let dirs = match xx::file::ls(&*env::PITCHFORK_LOGS_DIR) {
779 Ok(d) => d,
780 Err(_) => return,
781 };
782 for dir in dirs {
783 if dir.starts_with(".") || !dir.is_dir() {
784 continue;
785 }
786 let name = match dir.file_name().map(|f| f.to_string_lossy().to_string()) {
787 Some(n) => n,
788 None => continue,
789 };
790 if name == "pitchfork" {
792 continue;
793 }
794 if name.contains("--") {
797 if DaemonId::from_safe_path(&name).is_ok() {
800 continue;
801 }
802 if known_safe_paths.contains(&name) {
805 continue;
806 }
807 warn!(
808 "Skipping invalid legacy log directory '{name}': contains '--' but is not a valid daemon safe-path"
809 );
810 continue;
811 }
812
813 let old_log = dir.join(format!("{name}.log"));
816 if !old_log.exists() {
817 continue;
818 }
819 if DaemonId::try_new("legacy", &name).is_err() {
820 warn!("Skipping invalid legacy log directory '{name}': not a valid daemon ID");
821 continue;
822 }
823
824 let new_name = format!("legacy--{name}");
825 let new_dir = env::PITCHFORK_LOGS_DIR.join(&new_name);
826 if new_dir.exists() {
828 continue;
829 }
830 if std::fs::rename(&dir, &new_dir).is_err() {
831 continue;
832 }
833 let old_log = new_dir.join(format!("{name}.log"));
835 let new_log = new_dir.join(format!("{new_name}.log"));
836 if old_log.exists() {
837 let _ = std::fs::rename(&old_log, &new_log);
838 }
839 debug!("Migrated legacy log dir '{name}' → '{new_name}'");
840 }
841}
842
843fn known_daemon_safe_paths() -> BTreeSet<String> {
844 let mut out = BTreeSet::new();
845
846 match StateFile::read(&*env::PITCHFORK_STATE_FILE) {
847 Ok(state) => {
848 for id in state.daemons.keys() {
849 out.insert(id.safe_path());
850 }
851 }
852 Err(e) => {
853 warn!("Failed to read state while checking known daemon IDs: {e}");
854 }
855 }
856
857 match PitchforkToml::all_merged() {
858 Ok(config) => {
859 for id in config.daemons.keys() {
860 out.insert(id.safe_path());
861 }
862 }
863 Err(e) => {
864 warn!("Failed to read config while checking known daemon IDs: {e}");
865 }
866 }
867
868 out
869}
870
871fn get_all_daemon_ids() -> Result<Vec<DaemonId>> {
872 let mut ids = BTreeSet::new();
873
874 match StateFile::read(&*env::PITCHFORK_STATE_FILE) {
875 Ok(state) => ids.extend(state.daemons.keys().cloned()),
876 Err(e) => warn!("Failed to read state for log daemon discovery: {e}"),
877 }
878
879 match PitchforkToml::all_merged() {
880 Ok(config) => ids.extend(config.daemons.keys().cloned()),
881 Err(e) => warn!("Failed to read config for log daemon discovery: {e}"),
882 }
883
884 let logged_ids: std::collections::HashSet<String> =
885 LOG_STORE.list_daemon_ids()?.into_iter().collect();
886 Ok(ids
887 .into_iter()
888 .filter(|id| logged_ids.contains(&id.qualified()))
889 .collect())
890}
891
892#[allow(clippy::too_many_arguments)]
893pub async fn tail_logs(
894 names: &[DaemonId],
895 single_daemon: bool,
896 start_from_end: bool,
897 message_filters: Vec<MessageFilter>,
898 field_filters: Vec<FieldFilter>,
899 jq_filter: Option<&crate::log_jq::JqFilter>,
900 show_timestamp: bool,
901 raw: bool,
902) -> Result<()> {
903 let strip_ansi = raw || !console::colors_enabled();
905
906 let mut states: std::collections::HashMap<String, i64> = names
907 .iter()
908 .map(|id| {
909 let since = if start_from_end {
910 LOG_STORE.last_id(id).unwrap_or(None).unwrap_or(0)
914 } else {
915 0
916 };
917 (id.qualified(), since)
918 })
919 .collect();
920
921 let interval = tokio::time::interval(Duration::from_millis(200));
922 tokio::pin!(interval);
923
924 loop {
925 interval.tick().await;
926
927 let mut out = vec![];
928 for id in names {
929 let after_id = states.get(&id.qualified()).copied();
930 match LOG_STORE.query(&LogQuery {
931 daemon_ids: vec![id.qualified()],
932 from: None,
933 to: None,
934 limit: None,
935 order_desc: false,
936 after_id,
937 before_id: None,
938 message_filters: message_filters.clone(),
939 field_filters: field_filters.clone(),
940 include_structured: jq_filter.is_some() || !raw,
941 }) {
942 Ok(raw_entries) => {
943 let last_raw_id = raw_entries.last().map(|e| e.id);
947
948 let entries = match jq_filter {
950 Some(jq) => jq.filter(raw_entries),
951 None => raw_entries,
952 };
953 out.extend(entries);
954 let has_sql_filter = !message_filters.is_empty() || !field_filters.is_empty();
966 let new_cursor = if let Some(id) = last_raw_id {
967 Some(id)
968 } else if has_sql_filter || jq_filter.is_some() {
969 LOG_STORE.last_id(id).ok().flatten()
970 } else {
971 None
972 };
973 if let Some(last_id) = new_cursor {
974 states.insert(id.qualified(), last_id);
975 }
976 }
977 Err(e) => {
978 error!("Failed to tail logs for {}: {e}", id.qualified());
979 }
980 }
981 }
982
983 if !out.is_empty() {
984 let out: Vec<LogEntry> = if single_daemon {
987 out
988 } else {
989 out.into_iter()
990 .sorted_by(|a, b| a.timestamp.cmp(&b.timestamp).then(a.id.cmp(&b.id)))
991 .collect()
992 };
993 let stdout = io::stdout();
994 let mut buf = io::BufWriter::new(stdout.lock());
995 let ts_format = &settings().logs.timestamp_format;
996 let mut date_buf = String::with_capacity(ts_format.len() + 6);
997 for entry in &out {
998 date_buf.clear();
999 write!(date_buf, "{}", entry.timestamp.format(ts_format))
1000 .map_err(io::Error::other)
1001 .into_diagnostic()?;
1002 write_formatted_log(
1003 &mut buf,
1004 entry,
1005 &date_buf,
1006 single_daemon,
1007 strip_ansi,
1008 show_timestamp,
1009 raw,
1010 )
1011 .into_diagnostic()?;
1012 }
1013 }
1014 }
1015}
1016
1017fn parse_datetime(s: &str) -> Result<DateTime<Local>> {
1018 let naive_dt = NaiveDateTime::parse_from_str(s, "%Y-%m-%d %H:%M:%S").into_diagnostic()?;
1019 Local
1020 .from_local_datetime(&naive_dt)
1021 .single()
1022 .ok_or_else(|| miette::miette!("Invalid or ambiguous datetime: '{}'. ", s))
1023}
1024
1025fn parse_time_input(s: &str, is_since: bool) -> Result<DateTime<Local>> {
1031 let s = s.trim();
1032
1033 if let Ok(dt) = parse_datetime(s) {
1035 return Ok(dt);
1036 }
1037
1038 if let Ok(naive_dt) = NaiveDateTime::parse_from_str(s, "%Y-%m-%d %H:%M") {
1040 return Local
1041 .from_local_datetime(&naive_dt)
1042 .single()
1043 .ok_or_else(|| miette::miette!("Invalid or ambiguous datetime: '{}'", s));
1044 }
1045
1046 if let Ok(time) = parse_time_only(s) {
1050 let now = Local::now();
1051 let today = now.date_naive();
1052 let mut naive_dt = NaiveDateTime::new(today, time);
1053 let mut dt = Local
1054 .from_local_datetime(&naive_dt)
1055 .single()
1056 .ok_or_else(|| miette::miette!("Invalid or ambiguous datetime: '{}'", s))?;
1057
1058 if is_since
1061 && dt > now
1062 && let Some(yesterday) = today.pred_opt()
1063 {
1064 naive_dt = NaiveDateTime::new(yesterday, time);
1065 dt = Local
1066 .from_local_datetime(&naive_dt)
1067 .single()
1068 .ok_or_else(|| miette::miette!("Invalid or ambiguous datetime: '{}'", s))?;
1069 }
1070 return Ok(dt);
1071 }
1072
1073 if let Ok(duration) = humantime::parse_duration(s) {
1074 let now = Local::now();
1075 let target = now - chrono::Duration::from_std(duration).into_diagnostic()?;
1076 return Ok(target);
1077 }
1078
1079 Err(miette::miette!(
1080 "Invalid time format: '{}'. Expected formats:\n\
1081 - Full datetime: \"YYYY-MM-DD HH:MM:SS\" or \"YYYY-MM-DD HH:MM\"\n\
1082 - Time only: \"HH:MM:SS\" or \"HH:MM\" (uses today's date)\n\
1083 - Relative time: \"5min\", \"2h\", \"1d\" (e.g., last 5 minutes)",
1084 s
1085 ))
1086}
1087
1088fn parse_time_only(s: &str) -> Result<NaiveTime> {
1089 if let Ok(time) = NaiveTime::parse_from_str(s, "%H:%M:%S") {
1090 return Ok(time);
1091 }
1092
1093 if let Ok(time) = NaiveTime::parse_from_str(s, "%H:%M") {
1094 return Ok(time);
1095 }
1096
1097 Err(miette::miette!("Invalid time format: '{}'", s))
1098}
1099
1100pub fn print_error_logs_block(log_lines: &[(String, String, String)]) {
1110 if log_lines.is_empty() {
1111 return;
1112 }
1113
1114 let is_tty = std::io::stderr().is_terminal();
1115 let format_msg = |msg: &str| -> String {
1116 let stripped = strip_pty_controls(msg);
1117 if is_tty {
1118 stripped
1119 } else {
1120 console::strip_ansi_codes(&stripped).to_string()
1121 }
1122 };
1123
1124 let tag = estyle(" ERROR LOGS ").white().on_red();
1125 eprintln!("\n{tag}");
1126
1127 let unique_ids: BTreeSet<&str> = log_lines.iter().map(|(_, id, _)| id.as_str()).collect();
1129 let show_id = unique_ids.len() > 1;
1130
1131 if show_id {
1132 for (date, id, msg) in log_lines {
1133 let time = date.split(' ').nth(1).unwrap_or(date);
1134 let colored = colored_id_label(id, is_tty && console::colors_enabled_stderr());
1135 eprintln!(
1136 "{} {} {}",
1137 estyle(time).red().dim(),
1138 colored,
1139 format_msg(msg)
1140 );
1141 }
1142 } else {
1143 for (date, _, msg) in log_lines {
1144 let time = date.split(' ').nth(1).unwrap_or(date);
1145 eprintln!("{} {}", estyle(time).red().dim(), format_msg(msg));
1146 }
1147 }
1148}
1149
1150pub enum ReadyCheckType {
1152 Output(String),
1153 Http(String),
1154 Port(u16),
1155 Cmd(String),
1156 Delay(u64),
1157 Default,
1158}
1159
1160impl std::fmt::Display for ReadyCheckType {
1161 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1162 match self {
1163 ReadyCheckType::Output(pattern) => write!(f, "output matching '{pattern}'"),
1164 ReadyCheckType::Http(url) => write!(f, "HTTP {url}"),
1165 ReadyCheckType::Port(port) => write!(f, "TCP port {port}"),
1166 ReadyCheckType::Cmd(cmd) => write!(f, "command '{cmd}'"),
1167 ReadyCheckType::Delay(secs) => write!(f, "delay ({secs}s)"),
1168 ReadyCheckType::Default => write!(f, "default readiness check"),
1169 }
1170 }
1171}
1172
1173pub fn create_ready_check_job(
1179 daemon_id: &DaemonId,
1180 check_type: &ReadyCheckType,
1181) -> std::sync::Arc<clx::progress::ProgressJob> {
1182 use clx::progress::{ProgressJobBuilder, ProgressJobDoneBehavior, ProgressStatus};
1183
1184 let is_tty = std::io::stderr().is_terminal();
1185 let colors_enabled = is_tty && console::colors_enabled_stderr();
1186 let id_label = colored_id_label(&daemon_id.qualified(), colors_enabled);
1187 let show_ts = crate::settings::settings().general.startup_log_timestamps;
1188
1189 let prefix = if show_ts {
1193 edim(chrono::Local::now().format("%H:%M:%S").to_string()).to_string()
1196 } else {
1197 "{{spinner()}}".to_string()
1198 };
1199
1200 ProgressJobBuilder::new()
1201 .body(format!(
1202 "{} {} waiting for {{{{ check_type }}}}...",
1203 prefix, id_label
1204 ))
1205 .prop("check_type", &check_type.to_string())
1206 .status(ProgressStatus::Running)
1207 .on_done(ProgressJobDoneBehavior::Keep)
1208 .start()
1209}
1210
1211pub fn collect_startup_logs(
1216 daemon_id: &DaemonId,
1217 from: DateTime<Local>,
1218) -> Result<Vec<(String, String, String)>> {
1219 let entries = LOG_STORE.query(&LogQuery {
1220 daemon_ids: vec![daemon_id.qualified()],
1221 from: Some(from),
1222 to: None,
1223 limit: None,
1224 order_desc: false,
1225 after_id: None,
1226 before_id: None,
1227 message_filters: Vec::new(),
1228 field_filters: Vec::new(),
1229 include_structured: false,
1230 })?;
1231 let log_lines = entries
1232 .into_iter()
1233 .map(|e| {
1234 let ts = e.timestamp.format("%Y-%m-%d %H:%M:%S").to_string();
1235 (ts, e.daemon_id, e.message)
1236 })
1237 .collect();
1238
1239 Ok(log_lines)
1240}
1241
1242pub fn stream_startup_logs(
1248 daemon_id: &DaemonId,
1249 job: std::sync::Arc<clx::progress::ProgressJob>,
1250) -> (
1251 tokio::sync::watch::Sender<bool>,
1252 tokio::task::JoinHandle<()>,
1253) {
1254 let (tx, mut rx) = tokio::sync::watch::channel(false);
1255 let id = daemon_id.clone();
1256
1257 let show_ts = crate::settings::settings().general.startup_log_timestamps;
1258
1259 let anchor_id: i64 = LOG_STORE
1264 .query(&LogQuery {
1265 daemon_ids: vec![id.qualified()],
1266 limit: Some(1),
1267 order_desc: true,
1268 ..Default::default()
1269 })
1270 .ok()
1271 .and_then(|entries| entries.last().map(|e| e.id))
1272 .unwrap_or(0);
1273
1274 let handle = tokio::spawn(async move {
1275 let is_tty = std::io::stderr().is_terminal();
1276 let colors_enabled = is_tty && console::colors_enabled_stderr();
1277 let id_label = colored_id_label(&id.qualified(), colors_enabled);
1278 let prefix = if show_ts {
1279 String::new()
1280 } else {
1281 edim("•").to_string()
1282 };
1283
1284 let mut last_id = anchor_id;
1285
1286 if let Ok(entries) = LOG_STORE.tail(&id, Some(last_id)) {
1288 for entry in &entries {
1289 let time = entry.timestamp.format("%H:%M:%S").to_string();
1290 let msg = strip_pty_controls(&entry.message);
1291 let msg = if is_tty {
1292 msg
1293 } else {
1294 console::strip_ansi_codes(&msg).to_string()
1295 };
1296 let line_prefix = if show_ts {
1297 edim(time).to_string()
1298 } else {
1299 prefix.clone()
1300 };
1301 job.println(&format!("{} {} {}", line_prefix, id_label, msg));
1302 }
1303 if let Some(last) = entries.last() {
1304 last_id = last.id;
1305 }
1306 }
1307
1308 loop {
1309 tokio::select! {
1310 _ = tokio::time::sleep(Duration::from_millis(200)) => {
1311 if let Ok(entries) = LOG_STORE.tail(&id, Some(last_id)) {
1312 for entry in &entries {
1313 let time = entry.timestamp.format("%H:%M:%S").to_string();
1314 let msg = strip_pty_controls(&entry.message);
1315 let msg = if is_tty {
1316 msg
1317 } else {
1318 console::strip_ansi_codes(&msg).to_string()
1319 };
1320 let line_prefix = if show_ts {
1321 edim(time).to_string()
1322 } else {
1323 prefix.clone()
1324 };
1325 job.println(&format!("{} {} {}", line_prefix, id_label, msg));
1326 }
1327 if let Some(last) = entries.last() {
1328 last_id = last.id;
1329 }
1330 }
1331 }
1332 _ = rx.changed() => {
1333 break;
1334 }
1335 }
1336 }
1337
1338 if let Ok(entries) = LOG_STORE.tail(&id, Some(last_id)) {
1340 for entry in &entries {
1341 let time = entry.timestamp.format("%H:%M:%S").to_string();
1342 let msg = strip_pty_controls(&entry.message);
1343 let msg = if is_tty {
1344 msg
1345 } else {
1346 console::strip_ansi_codes(&msg).to_string()
1347 };
1348 let line_prefix = if show_ts {
1349 edim(time).to_string()
1350 } else {
1351 prefix.clone()
1352 };
1353 job.println(&format!("{} {} {}", line_prefix, id_label, msg));
1354 }
1355 }
1356 });
1357
1358 (tx, handle)
1359}
1360
1361fn strip_pty_controls(s: &str) -> String {
1366 struct Stripper {
1367 result: String,
1368 }
1369
1370 impl vte::Perform for Stripper {
1371 fn print(&mut self, c: char) {
1372 self.result.push(c);
1373 }
1374
1375 fn execute(&mut self, byte: u8) {
1376 if byte == b'\n' || byte == b'\t' {
1378 self.result.push(byte as char);
1379 }
1380 }
1381
1382 fn csi_dispatch(
1383 &mut self,
1384 params: &vte::Params,
1385 _intermediates: &[u8],
1386 _ignore: bool,
1387 action: char,
1388 ) {
1389 if action == 'm' {
1391 self.result.push_str("\x1b[");
1392 let mut first = true;
1393 for sub in params.iter() {
1394 if !first {
1395 self.result.push(';');
1396 }
1397 first = false;
1398 for (i, &p) in sub.iter().enumerate() {
1399 if i > 0 {
1400 self.result.push(':');
1401 }
1402 self.result.push_str(&p.to_string());
1403 }
1404 }
1405 self.result.push('m');
1406 }
1407 }
1409
1410 fn osc_dispatch(&mut self, _params: &[&[u8]], _bell_terminated: bool) {
1411 }
1413
1414 fn esc_dispatch(&mut self, _intermediates: &[u8], _ignore: bool, _byte: u8) {
1415 }
1417
1418 fn hook(
1419 &mut self,
1420 _params: &vte::Params,
1421 _intermediates: &[u8],
1422 _ignore: bool,
1423 _action: char,
1424 ) {
1425 }
1427
1428 fn put(&mut self, _byte: u8) {
1429 }
1431
1432 fn unhook(&mut self) {
1433 }
1435 }
1436
1437 let mut parser = vte::Parser::new();
1438 let mut stripper = Stripper {
1439 result: String::with_capacity(s.len()),
1440 };
1441 parser.advance(&mut stripper, s.as_bytes());
1442 stripper.result
1443}