pub mod metrics;
pub(crate) mod panic;
pub(crate) mod resources;
pub(crate) mod transcript;
use std::path::{Path, PathBuf};
use tracing_error::ErrorLayer;
use tracing_subscriber::layer::SubscriberExt;
use tracing_subscriber::util::SubscriberInitExt;
use tracing_subscriber::{EnvFilter, fmt};
pub const LOG_DIR_ENV_VAR: &str = "PIGEON_LOG_DIR";
const LOG_FILE_NAME: &str = "pigeon.jsonl";
fn default_log_dir() -> Result<PathBuf, String> {
if let Ok(dir) = std::env::var(LOG_DIR_ENV_VAR) {
return Ok(PathBuf::from(dir));
}
let project_dirs = directories::ProjectDirs::from("", "", "pigeon")
.ok_or("could not determine the data directory for this platform")?;
Ok(project_dirs.data_local_dir().join("logs"))
}
fn resolve_log_path(log_file: Option<&Path>) -> Result<(PathBuf, String), String> {
match log_file {
Some(path) => {
let dir = path
.parent()
.filter(|parent| !parent.as_os_str().is_empty())
.map(Path::to_path_buf)
.unwrap_or_else(|| PathBuf::from("."));
let file_name = path
.file_name()
.map(|name| name.to_string_lossy().into_owned())
.unwrap_or_else(|| LOG_FILE_NAME.to_string());
Ok((dir, file_name))
}
None => Ok((default_log_dir()?, LOG_FILE_NAME.to_string())),
}
}
pub fn init(
log_level: Option<&str>,
log_file: Option<&Path>,
) -> Result<tracing_appender::non_blocking::WorkerGuard, String> {
let (dir, file_name) = resolve_log_path(log_file)?;
std::fs::create_dir_all(&dir)
.map_err(|err| format!("failed to create {}: {err}", dir.display()))?;
let _ = LOG_FILE_PATH.set(dir.join(&file_name));
let appender = tracing_appender::rolling::never(&dir, &file_name);
let (writer, guard) = tracing_appender::non_blocking(appender);
let directive = log_level
.map(str::to_string)
.or_else(|| std::env::var("RUST_LOG").ok())
.unwrap_or_else(|| "warn,pigeon=info".to_string());
let filter = EnvFilter::try_new(&directive)
.map_err(|err| format!("invalid log filter '{directive}': {err}"))?;
let json_layer = fmt::layer()
.json()
.with_writer(writer)
.with_span_events(fmt::format::FmtSpan::CLOSE);
tracing_subscriber::registry()
.with(filter)
.with(json_layer)
.with(ErrorLayer::default())
.init();
Ok(guard)
}
pub fn install_panic_hook() {
panic::install_panic_hook();
}
pub(crate) fn instance() -> &'static str {
static INSTANCE: std::sync::OnceLock<String> = std::sync::OnceLock::new();
INSTANCE.get_or_init(|| sysinfo::System::host_name().unwrap_or_else(|| "unknown".to_string()))
}
static LOG_FILE_PATH: std::sync::OnceLock<PathBuf> = std::sync::OnceLock::new();
pub(crate) fn log_file_path() -> Result<PathBuf, String> {
LOG_FILE_PATH
.get()
.cloned()
.ok_or_else(|| "log file path not resolved -- observability::init must run first".into())
}
pub(crate) fn run_instrumented(command_name: &'static str, f: impl FnOnce() -> i32) -> i32 {
let span = tracing::info_span!("command", command = command_name, instance = instance());
let _guard = span.enter();
tracing::info!("command started");
let start = std::time::Instant::now();
let exit_code = f();
let elapsed = start.elapsed();
tracing::info!(
exit_code,
elapsed_ms = elapsed.as_millis() as u64,
"command finished"
);
let status = if exit_code == 0 { "success" } else { "failure" };
::metrics::histogram!(
"pigeon_command_duration_seconds",
"command" => command_name,
"instance" => instance(),
)
.record(elapsed.as_secs_f64());
::metrics::counter!(
"pigeon_command_runs_total",
"command" => command_name,
"status" => status,
"instance" => instance(),
)
.increment(1);
exit_code
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn resolve_log_path_splits_a_given_file_into_dir_and_name() {
let (dir, file_name) = resolve_log_path(Some(Path::new("/tmp/foo/out.jsonl"))).unwrap();
assert_eq!(dir, PathBuf::from("/tmp/foo"));
assert_eq!(file_name, "out.jsonl");
}
#[test]
fn resolve_log_path_falls_back_to_the_default_dir_and_name_when_omitted() {
let dir = tempfile::tempdir().unwrap();
unsafe {
std::env::set_var(LOG_DIR_ENV_VAR, dir.path());
}
let (resolved_dir, file_name) = resolve_log_path(None).unwrap();
unsafe {
std::env::remove_var(LOG_DIR_ENV_VAR);
}
assert_eq!(resolved_dir, dir.path());
assert_eq!(file_name, LOG_FILE_NAME);
}
#[test]
fn run_instrumented_returns_the_wrapped_closure_s_exit_code() {
assert_eq!(run_instrumented("test.command", || 0), 0);
assert_eq!(run_instrumented("test.command", || 1), 1);
}
}