use crate::cli::json_output::{JsonLogEntry, print_json};
use crate::daemon_id::DaemonId;
use crate::log_store::sqlite::LOG_STORE;
use crate::log_store::{FieldFilter, LogEntry, LogQuery, LogStore, MessageFilter};
use crate::pitchfork_toml::PitchforkToml;
use crate::settings::settings;
use crate::state_file::StateFile;
use crate::ui::style::{edim, estyle, ndim};
use crate::{Result, env};
use chrono::{DateTime, Local, NaiveDateTime, NaiveTime, TimeZone};
use console;
use itertools::Itertools;
use miette::IntoDiagnostic;
use std::collections::BTreeSet;
use std::fmt::Write as _;
use std::io::{self, IsTerminal, Write};
use std::process::{Child, Command, Stdio};
use std::time::Duration;
struct PagerConfig {
command: String,
args: Vec<String>,
}
impl PagerConfig {
fn new(start_at_end: bool) -> Self {
let command = std::env::var("PAGER").unwrap_or_else(|_| "less".to_string());
let args = Self::build_args(&command, start_at_end);
Self { command, args }
}
fn build_args(pager: &str, start_at_end: bool) -> Vec<String> {
let mut args = vec![];
if pager == "less" {
args.push("-R".to_string());
if start_at_end {
args.push("+G".to_string());
}
}
args
}
fn spawn_piped(&self) -> io::Result<Child> {
Command::new(&self.command)
.args(&self.args)
.stdin(Stdio::piped())
.spawn()
}
}
const KNOWN_FIELD_KEYS: &[&str] = &[
"level",
"severity",
"lvl",
"PRIORITY",
"@level",
"msg",
"message",
"event",
"@message",
"logger",
"name",
"component",
"module",
"timestamp",
"ts",
"time",
"@timestamp",
];
fn level_badge(level: &str) -> String {
let label = match level {
"error" => console::style("ERR").red().bold(),
"warn" => console::style("WRN").yellow().bold(),
"info" => console::style("INF").cyan().bold(),
"debug" => console::style("DBG").magenta().bold(),
"trace" => console::style("TRC").dim().bold(),
_ => return String::new(),
};
format!("{}{}{}", ndim("["), label, ndim("]"))
}
fn format_field_value(value: &serde_json::Value) -> String {
match value {
serde_json::Value::String(s) => {
let needs_quotes = s.is_empty()
|| s.contains(' ')
|| s.contains('=')
|| s.contains('\t')
|| s == "true"
|| s == "false"
|| s == "null"
|| s.parse::<f64>().is_ok();
if needs_quotes {
format!("\"{s}\"")
} else {
s.clone()
}
}
serde_json::Value::Number(n) => console::style(n).green().to_string(),
serde_json::Value::Bool(true) => console::style("true").yellow().to_string(),
serde_json::Value::Bool(false) => console::style("false").red().to_string(),
serde_json::Value::Null => console::style("null").dim().to_string(),
serde_json::Value::Object(_) | serde_json::Value::Array(_) => {
serde_json::to_string(value).unwrap_or_else(|_| value.to_string())
}
}
}
fn format_fields(fields_json: &str) -> String {
let Ok(serde_json::Value::Object(obj)) = serde_json::from_str(fields_json) else {
return String::new();
};
let mut parts = Vec::new();
for (key, value) in &obj {
if KNOWN_FIELD_KEYS.contains(&key.as_str()) {
continue;
}
let key_styled = console::style(key.as_str()).blue().to_string();
let value_styled = format_field_value(value);
parts.push(format!("{key_styled}={value_styled}"));
}
parts.join(" ")
}
#[allow(clippy::too_many_arguments)]
fn write_formatted_log(
w: &mut dyn Write,
entry: &LogEntry,
date: &str,
single_daemon: bool,
strip_ansi: bool,
show_timestamp: bool,
raw: bool,
) -> io::Result<()> {
let clean_msg = strip_pty_controls(&entry.message);
let raw_msg: std::borrow::Cow<'_, str> = if strip_ansi {
console::strip_ansi_codes(&clean_msg)
} else {
std::borrow::Cow::Owned(clean_msg)
};
let mut out = String::with_capacity(256);
if raw {
if show_timestamp {
out.push_str(&ndim(date).to_string());
out.push(' ');
}
if !single_daemon {
let colors_on = !strip_ansi && console::colors_enabled();
out.push_str(&colored_id_label(&entry.daemon_id, colors_on));
out.push(' ');
}
out.push_str(&raw_msg);
} else if entry.fields_json.is_some() {
let level = entry.level.as_deref();
let accent = match level {
Some("error") => "error",
Some("warn") => "warn",
_ => "",
};
if show_timestamp {
let ts = match accent {
"error" => console::style(date).red().to_string(),
"warn" => console::style(date).yellow().to_string(),
_ => ndim(date).to_string(),
};
out.push_str(&ts);
out.push(' ');
}
if !single_daemon {
let colors_on = !strip_ansi && console::colors_enabled();
out.push_str(&colored_id_label(&entry.daemon_id, colors_on));
out.push(' ');
}
let mut need_sep = false;
if let Some(lvl) = level {
let badge = level_badge(lvl);
if !badge.is_empty() {
out.push_str(&badge);
need_sep = true;
}
}
let sep = ndim(":").to_string();
let arrow = ndim(" > ").to_string();
if let Some(logger) = &entry.logger {
if need_sep {
out.push(' ');
}
out.push_str(&console::style(logger).italic().dim().to_string());
out.push_str(&sep);
need_sep = true;
}
let msg_cow = entry.msg.as_deref().filter(|s| !s.is_empty()).map(|s| {
let cleaned = strip_pty_controls(s);
if strip_ansi {
std::borrow::Cow::<str>::Owned(console::strip_ansi_codes(&cleaned).to_string())
} else {
std::borrow::Cow::<str>::Owned(cleaned)
}
});
let fields_str = entry
.fields_json
.as_deref()
.map(format_fields)
.filter(|s| !s.is_empty());
if msg_cow.is_some() || fields_str.is_some() {
if need_sep {
out.push(' ');
}
if let Some(msg) = &msg_cow {
let styled = match accent {
"error" => console::style(msg.as_ref()).red().bold().to_string(),
"warn" => console::style(msg.as_ref()).yellow().bold().to_string(),
_ => console::style(msg.as_ref()).bold().to_string(),
};
out.push_str(&styled);
if fields_str.is_some() {
out.push_str(&arrow);
}
}
if let Some(fields) = &fields_str {
out.push_str(fields);
}
} else if !need_sep {
out.push_str(&raw_msg);
}
} else {
if show_timestamp {
out.push_str(&ndim(date).to_string());
out.push(' ');
}
if !single_daemon {
let colors_on = !strip_ansi && console::colors_enabled();
out.push_str(&colored_id_label(&entry.daemon_id, colors_on));
out.push(' ');
}
out.push_str(&raw_msg);
}
out.push('\n');
w.write_all(out.as_bytes())
}
pub fn colored_id_label(id: &str, colors_enabled: bool) -> String {
if !colors_enabled {
return format!("[{}]", id);
}
let colors: [u8; 4] = [34, 35, 36, 32]; let mut h: usize = 0x811C_9DC5; for b in id.bytes() {
h = h.wrapping_mul(0x0100_0193).wrapping_add(b as usize);
}
let color = colors[h % colors.len()];
format!("\x1b[{color}m[{id}]\x1b[0m")
}
#[derive(Debug, clap::Args)]
#[clap(
visible_alias = "l",
verbatim_doc_comment,
long_about = "\
Displays logs for daemon(s)
Shows logs from managed daemons. Logs are stored in the pitchfork logs directory
and include timestamps for filtering.
Examples:
pitchfork logs api Show all logs for 'api' (paged if needed)
pitchfork logs api worker Show logs for multiple daemons
pitchfork logs Show logs for all daemons
pitchfork logs api -n 50 Show last 50 lines
pitchfork logs api --follow Follow logs in real-time
pitchfork logs api --since '2024-01-15 10:00:00'
Show logs since a specific time (forward)
pitchfork logs api --since '10:30:00'
Show logs since 10:30:00 today
pitchfork logs api --since '10:30' --until '12:00'
Show logs since 10:30:00 until 12:00:00 today
pitchfork logs api --since 5min Show logs from last 5 minutes
pitchfork logs api --raw Output raw log lines without formatting
pitchfork logs api --raw -n 100 Output last 100 raw log lines
pitchfork logs api --clear Delete logs for 'api'
pitchfork logs --clear Delete logs for all daemons"
)]
pub struct Logs {
id: Vec<String>,
#[clap(short, long)]
clear: bool,
#[clap(short)]
n: Option<usize>,
#[clap(short = 't', short_alias = 'f', long, visible_alias = "follow")]
tail: bool,
#[clap(short = 's', long)]
since: Option<String>,
#[clap(short = 'u', long)]
until: Option<String>,
#[clap(long)]
no_pager: bool,
#[clap(long)]
raw: bool,
#[clap(long, conflicts_with = "raw", conflicts_with = "tail")]
json: bool,
#[clap(long)]
grep: Vec<String>,
#[clap(long)]
regex: Option<String>,
#[clap(long)]
case_sensitive: bool,
#[clap(long)]
level: Option<String>,
#[clap(long, value_name = "KEY=VALUE")]
field: Vec<String>,
#[clap(long, value_name = "EXPR")]
jq: Option<String>,
#[clap(long)]
no_timestamp: bool,
}
impl Logs {
pub async fn run(&self) -> Result<()> {
migrate_legacy_log_dirs();
let resolved_ids: Vec<DaemonId> = if self.id.is_empty() {
get_all_daemon_ids()?
} else {
PitchforkToml::resolve_ids(&self.id)?
};
if self.clear {
LOG_STORE.clear(&resolved_ids)?;
return Ok(());
}
let from = if let Some(since) = self.since.as_ref() {
Some(parse_time_input(since, true)?)
} else {
None
};
let to = if let Some(until) = self.until.as_ref() {
Some(parse_time_input(until, false)?)
} else {
None
};
let message_filters = self.build_message_filters()?;
let field_filters = self.build_field_filters()?;
let jq_filter = match self.jq.as_deref() {
Some(expr) => Some(crate::log_jq::JqFilter::new(expr)?),
None => None,
};
if self.json {
return self.output_json(
&resolved_ids,
from,
to,
message_filters,
field_filters,
jq_filter.as_ref(),
);
}
let single_daemon = resolved_ids.len() == 1 && !self.id.is_empty();
let show_timestamp = settings().logs.timestamp && !self.no_timestamp && !self.raw;
let has_time_filter = from.is_some() || to.is_some();
self.query_and_output(
&resolved_ids,
from,
to,
message_filters.clone(),
field_filters.clone(),
jq_filter.as_ref(),
single_daemon,
has_time_filter,
show_timestamp,
)?;
if self.tail {
tail_logs(
&resolved_ids,
single_daemon,
true,
message_filters,
field_filters,
jq_filter.as_ref(),
show_timestamp,
self.raw,
)
.await?;
}
Ok(())
}
fn build_message_filters(&self) -> Result<Vec<MessageFilter>> {
if self.case_sensitive && self.grep.is_empty() {
warn!("--case-sensitive has no effect without --grep");
}
let mut filters = Vec::new();
for pattern in &self.grep {
filters.push(MessageFilter::Contains {
pattern: pattern.clone(),
case_sensitive: self.case_sensitive,
});
}
if let Some(pattern) = self.regex.as_ref() {
let _ = regex::Regex::new(pattern)
.into_diagnostic()
.map_err(|e| miette::miette!("invalid regex pattern: {e}"))?;
filters.push(MessageFilter::Regex {
pattern: pattern.clone(),
});
}
Ok(filters)
}
fn build_field_filters(&self) -> Result<Vec<FieldFilter>> {
let mut filters = Vec::new();
if let Some(level) = self.level.as_ref() {
let normalized = crate::log_parse::normalize_level_str(level).ok_or_else(|| {
miette::miette!(
"invalid level '{level}'; expected one of: \
error, err, fatal, critical, panic, alert, emerg, \
warn, warning, info, inf, information, notice, \
debug, dbg, trace, trc"
)
})?;
filters.push(FieldFilter::LevelMin(normalized));
}
for pair in &self.field {
let (key, value) = pair
.split_once('=')
.ok_or_else(|| miette::miette!("--field expects KEY=VALUE, got: {pair}"))?;
if key.is_empty() {
miette::bail!("--field key cannot be empty: {pair}");
}
if !key
.chars()
.all(|c| c.is_alphanumeric() || c == '_' || c == '.')
{
miette::bail!(
"--field key may only contain alphanumeric, underscore, and dot; got: {key}"
);
}
filters.push(FieldFilter::FieldEq {
key: key.to_string(),
value: value.to_string(),
});
}
Ok(filters)
}
#[allow(clippy::too_many_arguments)]
fn query_and_output(
&self,
resolved_ids: &[DaemonId],
from: Option<DateTime<Local>>,
to: Option<DateTime<Local>>,
message_filters: Vec<MessageFilter>,
field_filters: Vec<FieldFilter>,
jq_filter: Option<&crate::log_jq::JqFilter>,
single_daemon: bool,
has_time_filter: bool,
show_timestamp: bool,
) -> Result<()> {
let daemon_ids: Vec<String> = resolved_ids.iter().map(|id| id.qualified()).collect();
let opts = LogQuery {
daemon_ids,
from,
to,
limit: if !has_time_filter { self.n } else { None },
order_desc: !has_time_filter,
after_id: None,
before_id: None,
message_filters,
field_filters,
include_structured: jq_filter.is_some() || !self.raw,
};
let mut entries = LOG_STORE.query(&opts)?;
if let Some(jq) = jq_filter {
entries = jq.filter(entries);
}
if has_time_filter {
if let Some(n) = self.n {
let len = entries.len();
if len > n {
entries = entries.split_off(len - n);
}
}
} else {
entries.reverse();
}
if entries.is_empty() {
return Ok(());
}
let strip_ansi = self.raw || !console::colors_enabled();
let use_pager = !self.tail && !self.no_pager && should_use_pager(entries.len());
let ts_format = &settings().logs.timestamp_format;
let mut date_buf = String::with_capacity(ts_format.len() + 6);
let mut write_entries = |w: &mut dyn Write| -> io::Result<()> {
for entry in &entries {
date_buf.clear();
write!(date_buf, "{}", entry.timestamp.format(ts_format))
.map_err(io::Error::other)?;
write_formatted_log(
w,
entry,
&date_buf,
single_daemon,
strip_ansi,
show_timestamp,
self.raw,
)?;
}
Ok(())
};
if use_pager {
let pager_config = PagerConfig::new(!has_time_filter);
match pager_config.spawn_piped() {
Ok(mut child) => {
if let Some(stdin) = child.stdin.as_mut() {
if let Err(e) = write_entries(stdin) {
debug!("pager write error: {e}");
}
let _ = child.wait();
} else {
debug!("Failed to get pager stdin, falling back to direct output");
let stdout = io::stdout();
let mut buf = io::BufWriter::new(stdout.lock());
write_entries(&mut buf).into_diagnostic()?;
}
}
Err(e) => {
debug!("Failed to spawn pager: {e}, falling back to direct output");
let stdout = io::stdout();
let mut buf = io::BufWriter::new(stdout.lock());
write_entries(&mut buf).into_diagnostic()?;
}
}
} else {
let stdout = io::stdout();
let mut buf = io::BufWriter::new(stdout.lock());
write_entries(&mut buf).into_diagnostic()?;
}
Ok(())
}
fn output_json(
&self,
resolved_ids: &[DaemonId],
from: Option<DateTime<Local>>,
to: Option<DateTime<Local>>,
message_filters: Vec<MessageFilter>,
field_filters: Vec<FieldFilter>,
jq_filter: Option<&crate::log_jq::JqFilter>,
) -> Result<()> {
let daemon_ids: Vec<String> = resolved_ids.iter().map(|id| id.qualified()).collect();
let has_time_filter = from.is_some() || to.is_some();
let opts = LogQuery {
daemon_ids,
from,
to,
limit: if !has_time_filter { self.n } else { None },
order_desc: !has_time_filter,
after_id: None,
before_id: None,
message_filters,
field_filters,
include_structured: true,
};
let entries = LOG_STORE.query(&opts)?;
let entries = match jq_filter {
Some(jq) => jq.filter(entries),
None => entries,
};
let mut entries: Vec<_> = if has_time_filter {
entries
} else {
entries.into_iter().rev().collect()
};
if has_time_filter
&& let Some(n) = self.n
&& entries.len() > n
{
entries = entries.split_off(entries.len() - n);
}
let json_entries: Vec<JsonLogEntry> = entries.into_iter().map(Into::into).collect();
print_json(&json_entries)
}
}
fn should_use_pager(line_count: usize) -> bool {
if !io::stdout().is_terminal() {
return false;
}
let terminal_height = get_terminal_height().unwrap_or(24);
line_count > terminal_height
}
fn get_terminal_height() -> Option<usize> {
if let Ok(rows) = std::env::var("LINES")
&& let Ok(h) = rows.parse::<usize>()
{
return Some(h);
}
crossterm::terminal::size().ok().map(|(_, h)| h as usize)
}
fn migrate_legacy_log_dirs() {
let known_safe_paths = known_daemon_safe_paths();
let dirs = match xx::file::ls(&*env::PITCHFORK_LOGS_DIR) {
Ok(d) => d,
Err(_) => return,
};
for dir in dirs {
if dir.starts_with(".") || !dir.is_dir() {
continue;
}
let name = match dir.file_name().map(|f| f.to_string_lossy().to_string()) {
Some(n) => n,
None => continue,
};
if name == "pitchfork" {
continue;
}
if name.contains("--") {
if DaemonId::from_safe_path(&name).is_ok() {
continue;
}
if known_safe_paths.contains(&name) {
continue;
}
warn!(
"Skipping invalid legacy log directory '{name}': contains '--' but is not a valid daemon safe-path"
);
continue;
}
let old_log = dir.join(format!("{name}.log"));
if !old_log.exists() {
continue;
}
if DaemonId::try_new("legacy", &name).is_err() {
warn!("Skipping invalid legacy log directory '{name}': not a valid daemon ID");
continue;
}
let new_name = format!("legacy--{name}");
let new_dir = env::PITCHFORK_LOGS_DIR.join(&new_name);
if new_dir.exists() {
continue;
}
if std::fs::rename(&dir, &new_dir).is_err() {
continue;
}
let old_log = new_dir.join(format!("{name}.log"));
let new_log = new_dir.join(format!("{new_name}.log"));
if old_log.exists() {
let _ = std::fs::rename(&old_log, &new_log);
}
debug!("Migrated legacy log dir '{name}' → '{new_name}'");
}
}
fn known_daemon_safe_paths() -> BTreeSet<String> {
let mut out = BTreeSet::new();
match StateFile::read(&*env::PITCHFORK_STATE_FILE) {
Ok(state) => {
for id in state.daemons.keys() {
out.insert(id.safe_path());
}
}
Err(e) => {
warn!("Failed to read state while checking known daemon IDs: {e}");
}
}
match PitchforkToml::all_merged() {
Ok(config) => {
for id in config.daemons.keys() {
out.insert(id.safe_path());
}
}
Err(e) => {
warn!("Failed to read config while checking known daemon IDs: {e}");
}
}
out
}
fn get_all_daemon_ids() -> Result<Vec<DaemonId>> {
let mut ids = BTreeSet::new();
match StateFile::read(&*env::PITCHFORK_STATE_FILE) {
Ok(state) => ids.extend(state.daemons.keys().cloned()),
Err(e) => warn!("Failed to read state for log daemon discovery: {e}"),
}
match PitchforkToml::all_merged() {
Ok(config) => ids.extend(config.daemons.keys().cloned()),
Err(e) => warn!("Failed to read config for log daemon discovery: {e}"),
}
let logged_ids: std::collections::HashSet<String> =
LOG_STORE.list_daemon_ids()?.into_iter().collect();
Ok(ids
.into_iter()
.filter(|id| logged_ids.contains(&id.qualified()))
.collect())
}
#[allow(clippy::too_many_arguments)]
pub async fn tail_logs(
names: &[DaemonId],
single_daemon: bool,
start_from_end: bool,
message_filters: Vec<MessageFilter>,
field_filters: Vec<FieldFilter>,
jq_filter: Option<&crate::log_jq::JqFilter>,
show_timestamp: bool,
raw: bool,
) -> Result<()> {
let strip_ansi = raw || !console::colors_enabled();
let mut states: std::collections::HashMap<String, i64> = names
.iter()
.map(|id| {
let since = if start_from_end {
LOG_STORE.last_id(id).unwrap_or(None).unwrap_or(0)
} else {
0
};
(id.qualified(), since)
})
.collect();
let interval = tokio::time::interval(Duration::from_millis(200));
tokio::pin!(interval);
loop {
interval.tick().await;
let mut out = vec![];
for id in names {
let after_id = states.get(&id.qualified()).copied();
match LOG_STORE.query(&LogQuery {
daemon_ids: vec![id.qualified()],
from: None,
to: None,
limit: None,
order_desc: false,
after_id,
before_id: None,
message_filters: message_filters.clone(),
field_filters: field_filters.clone(),
include_structured: jq_filter.is_some() || !raw,
}) {
Ok(raw_entries) => {
let last_raw_id = raw_entries.last().map(|e| e.id);
let entries = match jq_filter {
Some(jq) => jq.filter(raw_entries),
None => raw_entries,
};
out.extend(entries);
let has_sql_filter = !message_filters.is_empty() || !field_filters.is_empty();
let new_cursor = if let Some(id) = last_raw_id {
Some(id)
} else if has_sql_filter || jq_filter.is_some() {
LOG_STORE.last_id(id).ok().flatten()
} else {
None
};
if let Some(last_id) = new_cursor {
states.insert(id.qualified(), last_id);
}
}
Err(e) => {
error!("Failed to tail logs for {}: {e}", id.qualified());
}
}
}
if !out.is_empty() {
let out: Vec<LogEntry> = if single_daemon {
out
} else {
out.into_iter()
.sorted_by(|a, b| a.timestamp.cmp(&b.timestamp).then(a.id.cmp(&b.id)))
.collect()
};
let stdout = io::stdout();
let mut buf = io::BufWriter::new(stdout.lock());
let ts_format = &settings().logs.timestamp_format;
let mut date_buf = String::with_capacity(ts_format.len() + 6);
for entry in &out {
date_buf.clear();
write!(date_buf, "{}", entry.timestamp.format(ts_format))
.map_err(io::Error::other)
.into_diagnostic()?;
write_formatted_log(
&mut buf,
entry,
&date_buf,
single_daemon,
strip_ansi,
show_timestamp,
raw,
)
.into_diagnostic()?;
}
}
}
}
fn parse_datetime(s: &str) -> Result<DateTime<Local>> {
let naive_dt = NaiveDateTime::parse_from_str(s, "%Y-%m-%d %H:%M:%S").into_diagnostic()?;
Local
.from_local_datetime(&naive_dt)
.single()
.ok_or_else(|| miette::miette!("Invalid or ambiguous datetime: '{}'. ", s))
}
fn parse_time_input(s: &str, is_since: bool) -> Result<DateTime<Local>> {
let s = s.trim();
if let Ok(dt) = parse_datetime(s) {
return Ok(dt);
}
if let Ok(naive_dt) = NaiveDateTime::parse_from_str(s, "%Y-%m-%d %H:%M") {
return Local
.from_local_datetime(&naive_dt)
.single()
.ok_or_else(|| miette::miette!("Invalid or ambiguous datetime: '{}'", s));
}
if let Ok(time) = parse_time_only(s) {
let now = Local::now();
let today = now.date_naive();
let mut naive_dt = NaiveDateTime::new(today, time);
let mut dt = Local
.from_local_datetime(&naive_dt)
.single()
.ok_or_else(|| miette::miette!("Invalid or ambiguous datetime: '{}'", s))?;
if is_since
&& dt > now
&& let Some(yesterday) = today.pred_opt()
{
naive_dt = NaiveDateTime::new(yesterday, time);
dt = Local
.from_local_datetime(&naive_dt)
.single()
.ok_or_else(|| miette::miette!("Invalid or ambiguous datetime: '{}'", s))?;
}
return Ok(dt);
}
if let Ok(duration) = humantime::parse_duration(s) {
let now = Local::now();
let target = now - chrono::Duration::from_std(duration).into_diagnostic()?;
return Ok(target);
}
Err(miette::miette!(
"Invalid time format: '{}'. Expected formats:\n\
- Full datetime: \"YYYY-MM-DD HH:MM:SS\" or \"YYYY-MM-DD HH:MM\"\n\
- Time only: \"HH:MM:SS\" or \"HH:MM\" (uses today's date)\n\
- Relative time: \"5min\", \"2h\", \"1d\" (e.g., last 5 minutes)",
s
))
}
fn parse_time_only(s: &str) -> Result<NaiveTime> {
if let Ok(time) = NaiveTime::parse_from_str(s, "%H:%M:%S") {
return Ok(time);
}
if let Ok(time) = NaiveTime::parse_from_str(s, "%H:%M") {
return Ok(time);
}
Err(miette::miette!("Invalid time format: '{}'", s))
}
pub fn print_error_logs_block(log_lines: &[(String, String, String)]) {
if log_lines.is_empty() {
return;
}
let is_tty = std::io::stderr().is_terminal();
let format_msg = |msg: &str| -> String {
let stripped = strip_pty_controls(msg);
if is_tty {
stripped
} else {
console::strip_ansi_codes(&stripped).to_string()
}
};
let tag = estyle(" ERROR LOGS ").white().on_red();
eprintln!("\n{tag}");
let unique_ids: BTreeSet<&str> = log_lines.iter().map(|(_, id, _)| id.as_str()).collect();
let show_id = unique_ids.len() > 1;
if show_id {
for (date, id, msg) in log_lines {
let time = date.split(' ').nth(1).unwrap_or(date);
let colored = colored_id_label(id, is_tty && console::colors_enabled_stderr());
eprintln!(
"{} {} {}",
estyle(time).red().dim(),
colored,
format_msg(msg)
);
}
} else {
for (date, _, msg) in log_lines {
let time = date.split(' ').nth(1).unwrap_or(date);
eprintln!("{} {}", estyle(time).red().dim(), format_msg(msg));
}
}
}
pub enum ReadyCheckType {
Output(String),
Http(String),
Port(u16),
Cmd(String),
Delay(u64),
Default,
}
impl std::fmt::Display for ReadyCheckType {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
ReadyCheckType::Output(pattern) => write!(f, "output matching '{pattern}'"),
ReadyCheckType::Http(url) => write!(f, "HTTP {url}"),
ReadyCheckType::Port(port) => write!(f, "TCP port {port}"),
ReadyCheckType::Cmd(cmd) => write!(f, "command '{cmd}'"),
ReadyCheckType::Delay(secs) => write!(f, "delay ({secs}s)"),
ReadyCheckType::Default => write!(f, "default readiness check"),
}
}
}
pub fn create_ready_check_job(
daemon_id: &DaemonId,
check_type: &ReadyCheckType,
) -> std::sync::Arc<clx::progress::ProgressJob> {
use clx::progress::{ProgressJobBuilder, ProgressJobDoneBehavior, ProgressStatus};
let is_tty = std::io::stderr().is_terminal();
let colors_enabled = is_tty && console::colors_enabled_stderr();
let id_label = colored_id_label(&daemon_id.qualified(), colors_enabled);
let show_ts = crate::settings::settings().general.startup_log_timestamps;
let prefix = if show_ts {
edim(chrono::Local::now().format("%H:%M:%S").to_string()).to_string()
} else {
"{{spinner()}}".to_string()
};
ProgressJobBuilder::new()
.body(format!(
"{} {} waiting for {{{{ check_type }}}}...",
prefix, id_label
))
.prop("check_type", &check_type.to_string())
.status(ProgressStatus::Running)
.on_done(ProgressJobDoneBehavior::Keep)
.start()
}
pub fn collect_startup_logs(
daemon_id: &DaemonId,
from: DateTime<Local>,
) -> Result<Vec<(String, String, String)>> {
let entries = LOG_STORE.query(&LogQuery {
daemon_ids: vec![daemon_id.qualified()],
from: Some(from),
to: None,
limit: None,
order_desc: false,
after_id: None,
before_id: None,
message_filters: Vec::new(),
field_filters: Vec::new(),
include_structured: false,
})?;
let log_lines = entries
.into_iter()
.map(|e| {
let ts = e.timestamp.format("%Y-%m-%d %H:%M:%S").to_string();
(ts, e.daemon_id, e.message)
})
.collect();
Ok(log_lines)
}
pub fn stream_startup_logs(
daemon_id: &DaemonId,
job: std::sync::Arc<clx::progress::ProgressJob>,
) -> (
tokio::sync::watch::Sender<bool>,
tokio::task::JoinHandle<()>,
) {
let (tx, mut rx) = tokio::sync::watch::channel(false);
let id = daemon_id.clone();
let show_ts = crate::settings::settings().general.startup_log_timestamps;
let anchor_id: i64 = LOG_STORE
.query(&LogQuery {
daemon_ids: vec![id.qualified()],
limit: Some(1),
order_desc: true,
..Default::default()
})
.ok()
.and_then(|entries| entries.last().map(|e| e.id))
.unwrap_or(0);
let handle = tokio::spawn(async move {
let is_tty = std::io::stderr().is_terminal();
let colors_enabled = is_tty && console::colors_enabled_stderr();
let id_label = colored_id_label(&id.qualified(), colors_enabled);
let prefix = if show_ts {
String::new()
} else {
edim("•").to_string()
};
let mut last_id = anchor_id;
if let Ok(entries) = LOG_STORE.tail(&id, Some(last_id)) {
for entry in &entries {
let time = entry.timestamp.format("%H:%M:%S").to_string();
let msg = strip_pty_controls(&entry.message);
let msg = if is_tty {
msg
} else {
console::strip_ansi_codes(&msg).to_string()
};
let line_prefix = if show_ts {
edim(time).to_string()
} else {
prefix.clone()
};
job.println(&format!("{} {} {}", line_prefix, id_label, msg));
}
if let Some(last) = entries.last() {
last_id = last.id;
}
}
loop {
tokio::select! {
_ = tokio::time::sleep(Duration::from_millis(200)) => {
if let Ok(entries) = LOG_STORE.tail(&id, Some(last_id)) {
for entry in &entries {
let time = entry.timestamp.format("%H:%M:%S").to_string();
let msg = strip_pty_controls(&entry.message);
let msg = if is_tty {
msg
} else {
console::strip_ansi_codes(&msg).to_string()
};
let line_prefix = if show_ts {
edim(time).to_string()
} else {
prefix.clone()
};
job.println(&format!("{} {} {}", line_prefix, id_label, msg));
}
if let Some(last) = entries.last() {
last_id = last.id;
}
}
}
_ = rx.changed() => {
break;
}
}
}
if let Ok(entries) = LOG_STORE.tail(&id, Some(last_id)) {
for entry in &entries {
let time = entry.timestamp.format("%H:%M:%S").to_string();
let msg = strip_pty_controls(&entry.message);
let msg = if is_tty {
msg
} else {
console::strip_ansi_codes(&msg).to_string()
};
let line_prefix = if show_ts {
edim(time).to_string()
} else {
prefix.clone()
};
job.println(&format!("{} {} {}", line_prefix, id_label, msg));
}
}
});
(tx, handle)
}
fn strip_pty_controls(s: &str) -> String {
struct Stripper {
result: String,
}
impl vte::Perform for Stripper {
fn print(&mut self, c: char) {
self.result.push(c);
}
fn execute(&mut self, byte: u8) {
if byte == b'\n' || byte == b'\t' {
self.result.push(byte as char);
}
}
fn csi_dispatch(
&mut self,
params: &vte::Params,
_intermediates: &[u8],
_ignore: bool,
action: char,
) {
if action == 'm' {
self.result.push_str("\x1b[");
let mut first = true;
for sub in params.iter() {
if !first {
self.result.push(';');
}
first = false;
for (i, &p) in sub.iter().enumerate() {
if i > 0 {
self.result.push(':');
}
self.result.push_str(&p.to_string());
}
}
self.result.push('m');
}
}
fn osc_dispatch(&mut self, _params: &[&[u8]], _bell_terminated: bool) {
}
fn esc_dispatch(&mut self, _intermediates: &[u8], _ignore: bool, _byte: u8) {
}
fn hook(
&mut self,
_params: &vte::Params,
_intermediates: &[u8],
_ignore: bool,
_action: char,
) {
}
fn put(&mut self, _byte: u8) {
}
fn unhook(&mut self) {
}
}
let mut parser = vte::Parser::new();
let mut stripper = Stripper {
result: String::with_capacity(s.len()),
};
parser.advance(&mut stripper, s.as_bytes());
stripper.result
}