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
327 pitchfork logs api Show all logs for 'api' (paged if needed)
328 pitchfork logs api worker Show logs for multiple daemons
329 pitchfork logs Show logs for all daemons
330 pitchfork logs api -n 50 Show last 50 lines
331 pitchfork logs api --follow Follow logs in real-time
332 pitchfork logs api --since '2024-01-15 10:00:00'
333 Show logs since a specific time (forward)
334 pitchfork logs api --since '10:30:00'
335 Show logs since 10:30:00 today
336 pitchfork logs api --since '10:30' --until '12:00'
337 Show logs since 10:30:00 until 12:00:00 today
338 pitchfork logs api --since 5min Show logs from last 5 minutes
339 pitchfork logs api --raw Output raw log lines without formatting
340 pitchfork logs api --raw -n 100 Output last 100 raw log lines
341 pitchfork logs api --clear Delete logs for 'api'
342 pitchfork logs --clear Delete logs for all daemons"
343)]
344pub struct Logs {
345 id: Vec<String>,
347
348 #[clap(short, long)]
350 clear: bool,
351
352 #[clap(short)]
357 n: Option<usize>,
358
359 #[clap(short = 't', short_alias = 'f', long, visible_alias = "follow")]
361 tail: bool,
362
363 #[clap(short = 's', long)]
370 since: Option<String>,
371
372 #[clap(short = 'u', long)]
378 until: Option<String>,
379
380 #[clap(long)]
382 no_pager: bool,
383
384 #[clap(long)]
386 raw: bool,
387
388 #[clap(long, conflicts_with = "raw", conflicts_with = "tail")]
390 json: bool,
391
392 #[clap(long)]
396 grep: Vec<String>,
397
398 #[clap(long)]
400 regex: Option<String>,
401
402 #[clap(long)]
404 case_sensitive: bool,
405
406 #[clap(long)]
412 level: Option<String>,
413
414 #[clap(long, value_name = "KEY=VALUE")]
419 field: Vec<String>,
420
421 #[clap(long, value_name = "EXPR")]
427 jq: Option<String>,
428
429 #[clap(long)]
431 no_timestamp: bool,
432}
433
434impl Logs {
435 pub async fn run(&self) -> Result<()> {
436 migrate_legacy_log_dirs();
437
438 let resolved_ids: Vec<DaemonId> = if self.id.is_empty() {
439 get_all_daemon_ids()?
440 } else {
441 PitchforkToml::resolve_ids(&self.id)?
442 };
443
444 if self.clear {
445 LOG_STORE.clear(&resolved_ids)?;
446 return Ok(());
447 }
448
449 let from = if let Some(since) = self.since.as_ref() {
450 Some(parse_time_input(since, true)?)
451 } else {
452 None
453 };
454 let to = if let Some(until) = self.until.as_ref() {
455 Some(parse_time_input(until, false)?)
456 } else {
457 None
458 };
459
460 let message_filters = self.build_message_filters()?;
461 let field_filters = self.build_field_filters()?;
462
463 let jq_filter = match self.jq.as_deref() {
465 Some(expr) => Some(crate::log_jq::JqFilter::new(expr)?),
466 None => None,
467 };
468
469 if self.json {
470 return self.output_json(
471 &resolved_ids,
472 from,
473 to,
474 message_filters,
475 field_filters,
476 jq_filter.as_ref(),
477 );
478 }
479
480 let single_daemon = resolved_ids.len() == 1 && !self.id.is_empty();
485 let show_timestamp = settings().logs.timestamp && !self.no_timestamp && !self.raw;
486 let has_time_filter = from.is_some() || to.is_some();
487
488 self.query_and_output(
489 &resolved_ids,
490 from,
491 to,
492 message_filters.clone(),
493 field_filters.clone(),
494 jq_filter.as_ref(),
495 single_daemon,
496 has_time_filter,
497 show_timestamp,
498 )?;
499
500 if self.tail {
501 tail_logs(
502 &resolved_ids,
503 single_daemon,
504 true,
505 message_filters,
506 field_filters,
507 jq_filter.as_ref(),
508 show_timestamp,
509 self.raw,
510 )
511 .await?;
512 }
513
514 Ok(())
515 }
516
517 fn build_message_filters(&self) -> Result<Vec<MessageFilter>> {
518 if self.case_sensitive && self.grep.is_empty() {
519 warn!("--case-sensitive has no effect without --grep");
520 }
521 let mut filters = Vec::new();
522 for pattern in &self.grep {
523 filters.push(MessageFilter::Contains {
524 pattern: pattern.clone(),
525 case_sensitive: self.case_sensitive,
526 });
527 }
528 if let Some(pattern) = self.regex.as_ref() {
529 let _ = regex::Regex::new(pattern)
532 .into_diagnostic()
533 .map_err(|e| miette::miette!("invalid regex pattern: {e}"))?;
534 filters.push(MessageFilter::Regex {
535 pattern: pattern.clone(),
536 });
537 }
538 Ok(filters)
539 }
540
541 fn build_field_filters(&self) -> Result<Vec<FieldFilter>> {
542 let mut filters = Vec::new();
543 if let Some(level) = self.level.as_ref() {
544 let normalized = crate::log_parse::normalize_level_str(level).ok_or_else(|| {
547 miette::miette!(
548 "invalid level '{level}'; expected one of: \
549 error, err, fatal, critical, panic, alert, emerg, \
550 warn, warning, info, inf, information, notice, \
551 debug, dbg, trace, trc"
552 )
553 })?;
554 filters.push(FieldFilter::LevelMin(normalized));
555 }
556 for pair in &self.field {
557 let (key, value) = pair
558 .split_once('=')
559 .ok_or_else(|| miette::miette!("--field expects KEY=VALUE, got: {pair}"))?;
560 if key.is_empty() {
561 miette::bail!("--field key cannot be empty: {pair}");
562 }
563 if !key
566 .chars()
567 .all(|c| c.is_alphanumeric() || c == '_' || c == '.')
568 {
569 miette::bail!(
570 "--field key may only contain alphanumeric, underscore, and dot; got: {key}"
571 );
572 }
573 filters.push(FieldFilter::FieldEq {
574 key: key.to_string(),
575 value: value.to_string(),
576 });
577 }
578 Ok(filters)
579 }
580
581 #[allow(clippy::too_many_arguments)]
587 fn query_and_output(
588 &self,
589 resolved_ids: &[DaemonId],
590 from: Option<DateTime<Local>>,
591 to: Option<DateTime<Local>>,
592 message_filters: Vec<MessageFilter>,
593 field_filters: Vec<FieldFilter>,
594 jq_filter: Option<&crate::log_jq::JqFilter>,
595 single_daemon: bool,
596 has_time_filter: bool,
597 show_timestamp: bool,
598 ) -> Result<()> {
599 let daemon_ids: Vec<String> = resolved_ids.iter().map(|id| id.qualified()).collect();
600
601 let opts = LogQuery {
602 daemon_ids,
603 from,
604 to,
605 limit: if !has_time_filter { self.n } else { None },
606 order_desc: !has_time_filter,
607 after_id: None,
608 before_id: None,
609 message_filters,
610 field_filters,
611 include_structured: jq_filter.is_some() || !self.raw,
612 };
613 let mut entries = LOG_STORE.query(&opts)?;
614
615 if let Some(jq) = jq_filter {
617 entries = jq.filter(entries);
618 }
619
620 if has_time_filter {
622 if let Some(n) = self.n {
623 let len = entries.len();
624 if len > n {
625 entries = entries.split_off(len - n);
626 }
627 }
628 } else {
629 entries.reverse();
630 }
631
632 if entries.is_empty() {
633 return Ok(());
634 }
635
636 let strip_ansi = self.raw || !console::colors_enabled();
637 let use_pager = !self.tail && !self.no_pager && should_use_pager(entries.len());
638 let ts_format = &settings().logs.timestamp_format;
639
640 let mut date_buf = String::with_capacity(ts_format.len() + 6);
642
643 let mut write_entries = |w: &mut dyn Write| -> io::Result<()> {
644 for entry in &entries {
645 date_buf.clear();
646 write!(date_buf, "{}", entry.timestamp.format(ts_format))
647 .map_err(io::Error::other)?;
648 write_formatted_log(
649 w,
650 entry,
651 &date_buf,
652 single_daemon,
653 strip_ansi,
654 show_timestamp,
655 self.raw,
656 )?;
657 }
658 Ok(())
659 };
660
661 if use_pager {
662 let pager_config = PagerConfig::new(!has_time_filter);
663 match pager_config.spawn_piped() {
664 Ok(mut child) => {
665 if let Some(stdin) = child.stdin.as_mut() {
666 if let Err(e) = write_entries(stdin) {
667 debug!("pager write error: {e}");
668 }
669 let _ = child.wait();
670 } else {
671 debug!("Failed to get pager stdin, falling back to direct output");
672 let stdout = io::stdout();
673 let mut buf = io::BufWriter::new(stdout.lock());
674 write_entries(&mut buf).into_diagnostic()?;
675 }
676 }
677 Err(e) => {
678 debug!("Failed to spawn pager: {e}, falling back to direct output");
679 let stdout = io::stdout();
680 let mut buf = io::BufWriter::new(stdout.lock());
681 write_entries(&mut buf).into_diagnostic()?;
682 }
683 }
684 } else {
685 let stdout = io::stdout();
686 let mut buf = io::BufWriter::new(stdout.lock());
687 write_entries(&mut buf).into_diagnostic()?;
688 }
689
690 Ok(())
691 }
692
693 fn output_json(
694 &self,
695 resolved_ids: &[DaemonId],
696 from: Option<DateTime<Local>>,
697 to: Option<DateTime<Local>>,
698 message_filters: Vec<MessageFilter>,
699 field_filters: Vec<FieldFilter>,
700 jq_filter: Option<&crate::log_jq::JqFilter>,
701 ) -> Result<()> {
702 let daemon_ids: Vec<String> = resolved_ids.iter().map(|id| id.qualified()).collect();
703 let has_time_filter = from.is_some() || to.is_some();
704
705 let opts = LogQuery {
706 daemon_ids,
707 from,
708 to,
709 limit: if !has_time_filter { self.n } else { None },
710 order_desc: !has_time_filter,
711 after_id: None,
712 before_id: None,
713 message_filters,
714 field_filters,
715 include_structured: true,
716 };
717 let entries = LOG_STORE.query(&opts)?;
718
719 let entries = match jq_filter {
721 Some(jq) => jq.filter(entries),
722 None => entries,
723 };
724
725 let mut entries: Vec<_> = if has_time_filter {
729 entries
730 } else {
731 entries.into_iter().rev().collect()
732 };
733
734 if has_time_filter
737 && let Some(n) = self.n
738 && entries.len() > n
739 {
740 entries = entries.split_off(entries.len() - n);
741 }
742
743 let json_entries: Vec<JsonLogEntry> = entries.into_iter().map(Into::into).collect();
744
745 print_json(&json_entries)
746 }
747}
748
749fn should_use_pager(line_count: usize) -> bool {
750 if !io::stdout().is_terminal() {
751 return false;
752 }
753
754 let terminal_height = get_terminal_height().unwrap_or(24);
755 line_count > terminal_height
756}
757
758fn get_terminal_height() -> Option<usize> {
759 if let Ok(rows) = std::env::var("LINES")
760 && let Ok(h) = rows.parse::<usize>()
761 {
762 return Some(h);
763 }
764
765 crossterm::terminal::size().ok().map(|(_, h)| h as usize)
766}
767
768fn migrate_legacy_log_dirs() {
778 let known_safe_paths = known_daemon_safe_paths();
779 let dirs = match xx::file::ls(&*env::PITCHFORK_LOGS_DIR) {
780 Ok(d) => d,
781 Err(_) => return,
782 };
783 for dir in dirs {
784 if dir.starts_with(".") || !dir.is_dir() {
785 continue;
786 }
787 let name = match dir.file_name().map(|f| f.to_string_lossy().to_string()) {
788 Some(n) => n,
789 None => continue,
790 };
791 if name == "pitchfork" {
793 continue;
794 }
795 if name.contains("--") {
798 if DaemonId::from_safe_path(&name).is_ok() {
801 continue;
802 }
803 if known_safe_paths.contains(&name) {
806 continue;
807 }
808 warn!(
809 "Skipping invalid legacy log directory '{name}': contains '--' but is not a valid daemon safe-path"
810 );
811 continue;
812 }
813
814 let old_log = dir.join(format!("{name}.log"));
817 if !old_log.exists() {
818 continue;
819 }
820 if DaemonId::try_new("legacy", &name).is_err() {
821 warn!("Skipping invalid legacy log directory '{name}': not a valid daemon ID");
822 continue;
823 }
824
825 let new_name = format!("legacy--{name}");
826 let new_dir = env::PITCHFORK_LOGS_DIR.join(&new_name);
827 if new_dir.exists() {
829 continue;
830 }
831 if std::fs::rename(&dir, &new_dir).is_err() {
832 continue;
833 }
834 let old_log = new_dir.join(format!("{name}.log"));
836 let new_log = new_dir.join(format!("{new_name}.log"));
837 if old_log.exists() {
838 let _ = std::fs::rename(&old_log, &new_log);
839 }
840 debug!("Migrated legacy log dir '{name}' → '{new_name}'");
841 }
842}
843
844fn known_daemon_safe_paths() -> BTreeSet<String> {
845 let mut out = BTreeSet::new();
846
847 match StateFile::read(&*env::PITCHFORK_STATE_FILE) {
848 Ok(state) => {
849 for id in state.daemons.keys() {
850 out.insert(id.safe_path());
851 }
852 }
853 Err(e) => {
854 warn!("Failed to read state while checking known daemon IDs: {e}");
855 }
856 }
857
858 match PitchforkToml::all_merged() {
859 Ok(config) => {
860 for id in config.daemons.keys() {
861 out.insert(id.safe_path());
862 }
863 }
864 Err(e) => {
865 warn!("Failed to read config while checking known daemon IDs: {e}");
866 }
867 }
868
869 out
870}
871
872fn get_all_daemon_ids() -> Result<Vec<DaemonId>> {
873 let mut ids = BTreeSet::new();
874
875 match StateFile::read(&*env::PITCHFORK_STATE_FILE) {
876 Ok(state) => ids.extend(state.daemons.keys().cloned()),
877 Err(e) => warn!("Failed to read state for log daemon discovery: {e}"),
878 }
879
880 match PitchforkToml::all_merged() {
881 Ok(config) => ids.extend(config.daemons.keys().cloned()),
882 Err(e) => warn!("Failed to read config for log daemon discovery: {e}"),
883 }
884
885 let logged_ids: std::collections::HashSet<String> =
886 LOG_STORE.list_daemon_ids()?.into_iter().collect();
887 Ok(ids
888 .into_iter()
889 .filter(|id| logged_ids.contains(&id.qualified()))
890 .collect())
891}
892
893#[allow(clippy::too_many_arguments)]
894pub async fn tail_logs(
895 names: &[DaemonId],
896 single_daemon: bool,
897 start_from_end: bool,
898 message_filters: Vec<MessageFilter>,
899 field_filters: Vec<FieldFilter>,
900 jq_filter: Option<&crate::log_jq::JqFilter>,
901 show_timestamp: bool,
902 raw: bool,
903) -> Result<()> {
904 let strip_ansi = raw || !console::colors_enabled();
906
907 let mut states: std::collections::HashMap<String, i64> = names
908 .iter()
909 .map(|id| {
910 let since = if start_from_end {
911 LOG_STORE.last_id(id).unwrap_or(None).unwrap_or(0)
915 } else {
916 0
917 };
918 (id.qualified(), since)
919 })
920 .collect();
921
922 let interval = tokio::time::interval(Duration::from_millis(200));
923 tokio::pin!(interval);
924
925 loop {
926 interval.tick().await;
927
928 let mut out = vec![];
929 for id in names {
930 let after_id = states.get(&id.qualified()).copied();
931 match LOG_STORE.query(&LogQuery {
932 daemon_ids: vec![id.qualified()],
933 from: None,
934 to: None,
935 limit: None,
936 order_desc: false,
937 after_id,
938 before_id: None,
939 message_filters: message_filters.clone(),
940 field_filters: field_filters.clone(),
941 include_structured: jq_filter.is_some() || !raw,
942 }) {
943 Ok(raw_entries) => {
944 let last_raw_id = raw_entries.last().map(|e| e.id);
948
949 let entries = match jq_filter {
951 Some(jq) => jq.filter(raw_entries),
952 None => raw_entries,
953 };
954 out.extend(entries);
955 let has_sql_filter = !message_filters.is_empty() || !field_filters.is_empty();
967 let new_cursor = if let Some(id) = last_raw_id {
968 Some(id)
969 } else if has_sql_filter || jq_filter.is_some() {
970 LOG_STORE.last_id(id).ok().flatten()
971 } else {
972 None
973 };
974 if let Some(last_id) = new_cursor {
975 states.insert(id.qualified(), last_id);
976 }
977 }
978 Err(e) => {
979 error!("Failed to tail logs for {}: {e}", id.qualified());
980 }
981 }
982 }
983
984 if !out.is_empty() {
985 let out: Vec<LogEntry> = if single_daemon {
988 out
989 } else {
990 out.into_iter()
991 .sorted_by(|a, b| a.timestamp.cmp(&b.timestamp).then(a.id.cmp(&b.id)))
992 .collect()
993 };
994 let stdout = io::stdout();
995 let mut buf = io::BufWriter::new(stdout.lock());
996 let ts_format = &settings().logs.timestamp_format;
997 let mut date_buf = String::with_capacity(ts_format.len() + 6);
998 for entry in &out {
999 date_buf.clear();
1000 write!(date_buf, "{}", entry.timestamp.format(ts_format))
1001 .map_err(io::Error::other)
1002 .into_diagnostic()?;
1003 write_formatted_log(
1004 &mut buf,
1005 entry,
1006 &date_buf,
1007 single_daemon,
1008 strip_ansi,
1009 show_timestamp,
1010 raw,
1011 )
1012 .into_diagnostic()?;
1013 }
1014 }
1015 }
1016}
1017
1018fn parse_datetime(s: &str) -> Result<DateTime<Local>> {
1019 let naive_dt = NaiveDateTime::parse_from_str(s, "%Y-%m-%d %H:%M:%S").into_diagnostic()?;
1020 Local
1021 .from_local_datetime(&naive_dt)
1022 .single()
1023 .ok_or_else(|| miette::miette!("Invalid or ambiguous datetime: '{}'. ", s))
1024}
1025
1026fn parse_time_input(s: &str, is_since: bool) -> Result<DateTime<Local>> {
1032 let s = s.trim();
1033
1034 if let Ok(dt) = parse_datetime(s) {
1036 return Ok(dt);
1037 }
1038
1039 if let Ok(naive_dt) = NaiveDateTime::parse_from_str(s, "%Y-%m-%d %H:%M") {
1041 return Local
1042 .from_local_datetime(&naive_dt)
1043 .single()
1044 .ok_or_else(|| miette::miette!("Invalid or ambiguous datetime: '{}'", s));
1045 }
1046
1047 if let Ok(time) = parse_time_only(s) {
1051 let now = Local::now();
1052 let today = now.date_naive();
1053 let mut naive_dt = NaiveDateTime::new(today, time);
1054 let mut dt = Local
1055 .from_local_datetime(&naive_dt)
1056 .single()
1057 .ok_or_else(|| miette::miette!("Invalid or ambiguous datetime: '{}'", s))?;
1058
1059 if is_since
1062 && dt > now
1063 && let Some(yesterday) = today.pred_opt()
1064 {
1065 naive_dt = NaiveDateTime::new(yesterday, time);
1066 dt = Local
1067 .from_local_datetime(&naive_dt)
1068 .single()
1069 .ok_or_else(|| miette::miette!("Invalid or ambiguous datetime: '{}'", s))?;
1070 }
1071 return Ok(dt);
1072 }
1073
1074 if let Ok(duration) = humantime::parse_duration(s) {
1075 let now = Local::now();
1076 let target = now - chrono::Duration::from_std(duration).into_diagnostic()?;
1077 return Ok(target);
1078 }
1079
1080 Err(miette::miette!(
1081 "Invalid time format: '{}'. Expected formats:\n\
1082 - Full datetime: \"YYYY-MM-DD HH:MM:SS\" or \"YYYY-MM-DD HH:MM\"\n\
1083 - Time only: \"HH:MM:SS\" or \"HH:MM\" (uses today's date)\n\
1084 - Relative time: \"5min\", \"2h\", \"1d\" (e.g., last 5 minutes)",
1085 s
1086 ))
1087}
1088
1089fn parse_time_only(s: &str) -> Result<NaiveTime> {
1090 if let Ok(time) = NaiveTime::parse_from_str(s, "%H:%M:%S") {
1091 return Ok(time);
1092 }
1093
1094 if let Ok(time) = NaiveTime::parse_from_str(s, "%H:%M") {
1095 return Ok(time);
1096 }
1097
1098 Err(miette::miette!("Invalid time format: '{}'", s))
1099}
1100
1101pub fn print_error_logs_block(log_lines: &[(String, String, String)]) {
1111 if log_lines.is_empty() {
1112 return;
1113 }
1114
1115 let is_tty = std::io::stderr().is_terminal();
1116 let format_msg = |msg: &str| -> String {
1117 let stripped = strip_pty_controls(msg);
1118 if is_tty {
1119 stripped
1120 } else {
1121 console::strip_ansi_codes(&stripped).to_string()
1122 }
1123 };
1124
1125 let tag = estyle(" ERROR LOGS ").white().on_red();
1126 eprintln!("\n{tag}");
1127
1128 let unique_ids: BTreeSet<&str> = log_lines.iter().map(|(_, id, _)| id.as_str()).collect();
1130 let show_id = unique_ids.len() > 1;
1131
1132 if show_id {
1133 for (date, id, msg) in log_lines {
1134 let time = date.split(' ').nth(1).unwrap_or(date);
1135 let colored = colored_id_label(id, is_tty && console::colors_enabled_stderr());
1136 eprintln!(
1137 "{} {} {}",
1138 estyle(time).red().dim(),
1139 colored,
1140 format_msg(msg)
1141 );
1142 }
1143 } else {
1144 for (date, _, msg) in log_lines {
1145 let time = date.split(' ').nth(1).unwrap_or(date);
1146 eprintln!("{} {}", estyle(time).red().dim(), format_msg(msg));
1147 }
1148 }
1149}
1150
1151pub enum ReadyCheckType {
1153 Output(String),
1154 Http(String),
1155 Port(u16),
1156 Cmd(String),
1157 Delay(u64),
1158 Default,
1159}
1160
1161impl std::fmt::Display for ReadyCheckType {
1162 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1163 match self {
1164 ReadyCheckType::Output(pattern) => write!(f, "output matching '{pattern}'"),
1165 ReadyCheckType::Http(url) => write!(f, "HTTP {url}"),
1166 ReadyCheckType::Port(port) => write!(f, "TCP port {port}"),
1167 ReadyCheckType::Cmd(cmd) => write!(f, "command '{cmd}'"),
1168 ReadyCheckType::Delay(secs) => write!(f, "delay ({secs}s)"),
1169 ReadyCheckType::Default => write!(f, "default readiness check"),
1170 }
1171 }
1172}
1173
1174pub fn create_ready_check_job(
1180 daemon_id: &DaemonId,
1181 check_type: &ReadyCheckType,
1182) -> std::sync::Arc<clx::progress::ProgressJob> {
1183 use clx::progress::{ProgressJobBuilder, ProgressJobDoneBehavior, ProgressStatus};
1184
1185 let is_tty = std::io::stderr().is_terminal();
1186 let colors_enabled = is_tty && console::colors_enabled_stderr();
1187 let id_label = colored_id_label(&daemon_id.qualified(), colors_enabled);
1188 let show_ts = crate::settings::settings().general.startup_log_timestamps;
1189
1190 let prefix = if show_ts {
1194 edim(chrono::Local::now().format("%H:%M:%S").to_string()).to_string()
1197 } else {
1198 "{{spinner()}}".to_string()
1199 };
1200
1201 ProgressJobBuilder::new()
1202 .body(format!(
1203 "{} {} waiting for {{{{ check_type }}}}...",
1204 prefix, id_label
1205 ))
1206 .prop("check_type", &check_type.to_string())
1207 .status(ProgressStatus::Running)
1208 .on_done(ProgressJobDoneBehavior::Keep)
1209 .start()
1210}
1211
1212pub fn collect_startup_logs(
1217 daemon_id: &DaemonId,
1218 from: DateTime<Local>,
1219) -> Result<Vec<(String, String, String)>> {
1220 let entries = LOG_STORE.query(&LogQuery {
1221 daemon_ids: vec![daemon_id.qualified()],
1222 from: Some(from),
1223 to: None,
1224 limit: None,
1225 order_desc: false,
1226 after_id: None,
1227 before_id: None,
1228 message_filters: Vec::new(),
1229 field_filters: Vec::new(),
1230 include_structured: false,
1231 })?;
1232 let log_lines = entries
1233 .into_iter()
1234 .map(|e| {
1235 let ts = e.timestamp.format("%Y-%m-%d %H:%M:%S").to_string();
1236 (ts, e.daemon_id, e.message)
1237 })
1238 .collect();
1239
1240 Ok(log_lines)
1241}
1242
1243pub fn stream_startup_logs(
1249 daemon_id: &DaemonId,
1250 job: std::sync::Arc<clx::progress::ProgressJob>,
1251) -> (
1252 tokio::sync::watch::Sender<bool>,
1253 tokio::task::JoinHandle<()>,
1254) {
1255 let (tx, mut rx) = tokio::sync::watch::channel(false);
1256 let id = daemon_id.clone();
1257
1258 let show_ts = crate::settings::settings().general.startup_log_timestamps;
1259
1260 let anchor_id: i64 = LOG_STORE
1265 .query(&LogQuery {
1266 daemon_ids: vec![id.qualified()],
1267 limit: Some(1),
1268 order_desc: true,
1269 ..Default::default()
1270 })
1271 .ok()
1272 .and_then(|entries| entries.last().map(|e| e.id))
1273 .unwrap_or(0);
1274
1275 let handle = tokio::spawn(async move {
1276 let is_tty = std::io::stderr().is_terminal();
1277 let colors_enabled = is_tty && console::colors_enabled_stderr();
1278 let id_label = colored_id_label(&id.qualified(), colors_enabled);
1279 let prefix = if show_ts {
1280 String::new()
1281 } else {
1282 edim("•").to_string()
1283 };
1284
1285 let mut last_id = anchor_id;
1286
1287 if let Ok(entries) = LOG_STORE.tail(&id, Some(last_id)) {
1289 for entry in &entries {
1290 let time = entry.timestamp.format("%H:%M:%S").to_string();
1291 let msg = strip_pty_controls(&entry.message);
1292 let msg = if is_tty {
1293 msg
1294 } else {
1295 console::strip_ansi_codes(&msg).to_string()
1296 };
1297 let line_prefix = if show_ts {
1298 edim(time).to_string()
1299 } else {
1300 prefix.clone()
1301 };
1302 job.println(&format!("{} {} {}", line_prefix, id_label, msg));
1303 }
1304 if let Some(last) = entries.last() {
1305 last_id = last.id;
1306 }
1307 }
1308
1309 loop {
1310 tokio::select! {
1311 _ = tokio::time::sleep(Duration::from_millis(200)) => {
1312 if let Ok(entries) = LOG_STORE.tail(&id, Some(last_id)) {
1313 for entry in &entries {
1314 let time = entry.timestamp.format("%H:%M:%S").to_string();
1315 let msg = strip_pty_controls(&entry.message);
1316 let msg = if is_tty {
1317 msg
1318 } else {
1319 console::strip_ansi_codes(&msg).to_string()
1320 };
1321 let line_prefix = if show_ts {
1322 edim(time).to_string()
1323 } else {
1324 prefix.clone()
1325 };
1326 job.println(&format!("{} {} {}", line_prefix, id_label, msg));
1327 }
1328 if let Some(last) = entries.last() {
1329 last_id = last.id;
1330 }
1331 }
1332 }
1333 _ = rx.changed() => {
1334 break;
1335 }
1336 }
1337 }
1338
1339 if let Ok(entries) = LOG_STORE.tail(&id, Some(last_id)) {
1341 for entry in &entries {
1342 let time = entry.timestamp.format("%H:%M:%S").to_string();
1343 let msg = strip_pty_controls(&entry.message);
1344 let msg = if is_tty {
1345 msg
1346 } else {
1347 console::strip_ansi_codes(&msg).to_string()
1348 };
1349 let line_prefix = if show_ts {
1350 edim(time).to_string()
1351 } else {
1352 prefix.clone()
1353 };
1354 job.println(&format!("{} {} {}", line_prefix, id_label, msg));
1355 }
1356 }
1357 });
1358
1359 (tx, handle)
1360}
1361
1362fn strip_pty_controls(s: &str) -> String {
1367 struct Stripper {
1368 result: String,
1369 }
1370
1371 impl vte::Perform for Stripper {
1372 fn print(&mut self, c: char) {
1373 self.result.push(c);
1374 }
1375
1376 fn execute(&mut self, byte: u8) {
1377 if byte == b'\n' || byte == b'\t' {
1379 self.result.push(byte as char);
1380 }
1381 }
1382
1383 fn csi_dispatch(
1384 &mut self,
1385 params: &vte::Params,
1386 _intermediates: &[u8],
1387 _ignore: bool,
1388 action: char,
1389 ) {
1390 if action == 'm' {
1392 self.result.push_str("\x1b[");
1393 let mut first = true;
1394 for sub in params.iter() {
1395 if !first {
1396 self.result.push(';');
1397 }
1398 first = false;
1399 for (i, &p) in sub.iter().enumerate() {
1400 if i > 0 {
1401 self.result.push(':');
1402 }
1403 self.result.push_str(&p.to_string());
1404 }
1405 }
1406 self.result.push('m');
1407 }
1408 }
1410
1411 fn osc_dispatch(&mut self, _params: &[&[u8]], _bell_terminated: bool) {
1412 }
1414
1415 fn esc_dispatch(&mut self, _intermediates: &[u8], _ignore: bool, _byte: u8) {
1416 }
1418
1419 fn hook(
1420 &mut self,
1421 _params: &vte::Params,
1422 _intermediates: &[u8],
1423 _ignore: bool,
1424 _action: char,
1425 ) {
1426 }
1428
1429 fn put(&mut self, _byte: u8) {
1430 }
1432
1433 fn unhook(&mut self) {
1434 }
1436 }
1437
1438 let mut parser = vte::Parser::new();
1439 let mut stripper = Stripper {
1440 result: String::with_capacity(s.len()),
1441 };
1442 parser.advance(&mut stripper, s.as_bytes());
1443 stripper.result
1444}