use std::path::PathBuf;
use std::time::{SystemTime, UNIX_EPOCH};
use anyhow::Context as _;
use tracing::{Event, Level, Subscriber};
use tracing_appender::non_blocking::WorkerGuard;
use tracing_subscriber::fmt::format::{FormatEvent, FormatFields, Writer};
use tracing_subscriber::fmt::{self, FmtContext, FormattedFields};
use tracing_subscriber::layer::SubscriberExt as _;
use tracing_subscriber::registry::LookupSpan;
use tracing_subscriber::util::SubscriberInitExt as _;
use tracing_subscriber::{EnvFilter, Layer, Registry};
use crate::config::{self, TelemetryConfig};
const ENV_FILTER: &str = "ROTEIRO_LOG";
type BoxLayer = Box<dyn Layer<Registry> + Send + Sync + 'static>;
#[derive(Debug, Default, Clone)]
pub struct Overrides {
pub file: Option<String>,
pub enable_default: bool,
pub rotation: Option<String>,
pub format: Option<String>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum Rotation {
Daily,
Hourly,
Minutely,
Never,
}
impl Rotation {
fn parse(s: &str) -> anyhow::Result<Self> {
match s.trim().to_ascii_lowercase().as_str() {
"daily" => Ok(Self::Daily),
"hourly" => Ok(Self::Hourly),
"minutely" => Ok(Self::Minutely),
"never" => Ok(Self::Never),
other => anyhow::bail!(
"invalid telemetry rotation {other:?}: expected daily|hourly|minutely|never"
),
}
}
fn appender(self) -> tracing_appender::rolling::Rotation {
use tracing_appender::rolling::Rotation as R;
match self {
Self::Daily => R::DAILY,
Self::Hourly => R::HOURLY,
Self::Minutely => R::MINUTELY,
Self::Never => R::NEVER,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum Format {
Otel,
Text,
}
impl Format {
fn parse(s: &str) -> anyhow::Result<Self> {
match s.trim().to_ascii_lowercase().as_str() {
"otel" | "json" => Ok(Self::Otel),
"text" => Ok(Self::Text),
other => anyhow::bail!("invalid telemetry format {other:?}: expected otel|json|text"),
}
}
}
#[derive(Debug, Clone)]
struct Settings {
path: Option<PathBuf>,
rotation: Rotation,
format: Format,
}
#[derive(Debug)]
#[must_use = "hold the telemetry guard for the process lifetime; dropping it stops file logging"]
pub struct Guard(#[allow(dead_code)] Option<WorkerGuard>);
fn resolve(overrides: &Overrides, cfg: &TelemetryConfig) -> anyhow::Result<Settings> {
let raw_path = overrides
.file
.clone()
.or_else(|| cfg.file.clone())
.map(|p| resolve_path(&p));
let path = match raw_path {
Some(p) => Some(p),
None if overrides.enable_default => Some(
config::default_log_path()
.context("cannot resolve the default log path (no ROTEIRO_HOME or home dir)")?,
),
None => None,
};
let rotation = overrides
.rotation
.clone()
.or_else(|| cfg.rotation.clone())
.map_or(Ok(Rotation::Daily), |s| Rotation::parse(&s))?;
let format = overrides
.format
.clone()
.or_else(|| cfg.format.clone())
.map_or(Ok(Format::Otel), |s| Format::parse(&s))?;
Ok(Settings {
path,
rotation,
format,
})
}
fn resolve_path(raw: &str) -> PathBuf {
resolve_path_with(
raw,
config::home_dir().as_deref(),
config::roteiro_home().as_deref(),
)
}
fn resolve_path_with(
raw: &str,
home: Option<&std::path::Path>,
roteiro_home: Option<&std::path::Path>,
) -> PathBuf {
if raw == "~" || raw.starts_with("~/") {
let rest = raw.strip_prefix("~/").unwrap_or("");
if let Some(h) = home {
return if rest.is_empty() {
h.to_path_buf()
} else {
h.join(rest)
};
}
return roteiro_home.map_or_else(|| PathBuf::from(rest), |rh| rh.join(rest));
}
let p = PathBuf::from(raw);
if p.is_absolute() {
return p;
}
roteiro_home.map_or(p.clone(), |rh| rh.join(&p))
}
pub fn init(overrides: &Overrides, cfg: &TelemetryConfig) -> anyhow::Result<Guard> {
let settings = resolve(overrides, cfg)?;
let stdout = fmt::layer()
.with_writer(std::io::stdout)
.with_filter(env_filter(Level::WARN))
.boxed();
let mut layers: Vec<BoxLayer> = vec![stdout];
let mut guard = None;
if let Some(path) = &settings.path {
let (layer, worker) = file_layer(path, settings.rotation, settings.format)?;
layers.push(layer);
guard = Some(worker);
}
Registry::default()
.with(layers)
.try_init()
.context("installing the global tracing subscriber")?;
if settings.path.is_some() {
tracing::info!(
version = env!("CARGO_PKG_VERSION"),
"roteiro telemetry initialised"
);
}
Ok(Guard(guard))
}
fn env_filter(default: Level) -> EnvFilter {
EnvFilter::builder()
.with_default_directive(default.into())
.with_env_var(ENV_FILTER)
.from_env_lossy()
}
fn file_layer(
path: &std::path::Path,
rotation: Rotation,
format: Format,
) -> anyhow::Result<(BoxLayer, WorkerGuard)> {
let dir = path.parent().filter(|p| !p.as_os_str().is_empty());
let file_name = path
.file_name()
.context("telemetry log path has no file name")?;
let dir = dir.map_or_else(|| PathBuf::from("."), PathBuf::from);
std::fs::create_dir_all(&dir)
.with_context(|| format!("creating log directory {}", dir.display()))?;
let appender = tracing_appender::rolling::RollingFileAppender::new(
rotation.appender(),
&dir,
std::path::Path::new(file_name),
);
let (writer, worker) = tracing_appender::non_blocking(appender);
let layer = match format {
Format::Otel => fmt::layer()
.event_format(OtelJson)
.with_ansi(false)
.with_writer(writer)
.with_filter(env_filter(Level::INFO))
.boxed(),
Format::Text => fmt::layer()
.with_ansi(false)
.with_writer(writer)
.with_filter(env_filter(Level::INFO))
.boxed(),
};
Ok((layer, worker))
}
struct OtelJson;
impl<S, N> FormatEvent<S, N> for OtelJson
where
S: Subscriber + for<'a> LookupSpan<'a>,
N: for<'a> FormatFields<'a> + 'static,
{
fn format_event(
&self,
ctx: &FmtContext<'_, S, N>,
mut writer: Writer<'_>,
event: &Event<'_>,
) -> std::fmt::Result {
let meta = event.metadata();
let mut visitor = JsonVisitor::default();
event.record(&mut visitor);
let mut attributes = visitor.attributes;
attributes.insert("code.namespace".into(), meta.target().into());
if let Some(file) = meta.file() {
attributes.insert("code.filepath".into(), file.into());
}
if let Some(line) = meta.line() {
attributes.insert("code.lineno".into(), line.into());
}
if let Some(scope) = ctx.event_scope() {
let mut names = Vec::new();
for span in scope.from_root() {
names.push(span.name().to_owned());
let ext = span.extensions();
if let Some(fields) = ext.get::<FormattedFields<N>>()
&& !fields.fields.is_empty()
{
attributes.insert(
format!("span.{}.fields", span.name()),
fields.fields.as_str().into(),
);
}
}
if let Some(leaf) = names.last() {
attributes.insert("span.name".into(), leaf.clone().into());
}
attributes.insert("span.path".into(), names.join(" > ").into());
}
let now = SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap_or_default();
let unix_nano = u64::try_from(now.as_nanos()).unwrap_or(u64::MAX);
let (severity_number, severity_text) = severity(*meta.level());
let record = serde_json::json!({
"time_unix_nano": unix_nano,
"observed_time_unix_nano": unix_nano,
"severity_number": severity_number,
"severity_text": severity_text,
"body": visitor.body.unwrap_or_default(),
"attributes": attributes,
"resource": {
"service.name": "roteiro",
"service.version": env!("CARGO_PKG_VERSION"),
},
});
let line = serde_json::to_string(&record).map_err(|_| std::fmt::Error)?;
writeln!(writer, "{line}")
}
}
fn severity(level: Level) -> (u8, &'static str) {
match level {
Level::TRACE => (1, "TRACE"),
Level::DEBUG => (5, "DEBUG"),
Level::INFO => (9, "INFO"),
Level::WARN => (13, "WARN"),
Level::ERROR => (17, "ERROR"),
}
}
#[derive(Default)]
struct JsonVisitor {
body: Option<String>,
attributes: serde_json::Map<String, serde_json::Value>,
}
impl JsonVisitor {
fn put(&mut self, field: &tracing::field::Field, value: serde_json::Value) {
if field.name() == "message" {
self.body = value
.as_str()
.map(ToOwned::to_owned)
.or_else(|| Some(value.to_string()));
} else {
self.attributes.insert(field.name().to_owned(), value);
}
}
}
impl tracing::field::Visit for JsonVisitor {
fn record_str(&mut self, field: &tracing::field::Field, value: &str) {
self.put(field, value.into());
}
fn record_bool(&mut self, field: &tracing::field::Field, value: bool) {
self.put(field, value.into());
}
fn record_i64(&mut self, field: &tracing::field::Field, value: i64) {
self.put(field, value.into());
}
fn record_u64(&mut self, field: &tracing::field::Field, value: u64) {
self.put(field, value.into());
}
fn record_f64(&mut self, field: &tracing::field::Field, value: f64) {
self.put(field, value.into());
}
fn record_debug(&mut self, field: &tracing::field::Field, value: &dyn std::fmt::Debug) {
self.put(field, format!("{value:?}").into());
}
}
#[cfg(test)]
mod tests {
use super::{Format, Overrides, Rotation, file_layer, resolve, resolve_path_with};
use crate::config::TelemetryConfig;
use std::path::{Path, PathBuf};
use tracing_subscriber::layer::SubscriberExt as _;
#[test]
fn resolve_path_expands_tilde_and_falls_back_to_roteiro_home() {
let home = Path::new("/home/u");
let rh = Path::new("/iso/.roteiro");
assert_eq!(
resolve_path_with("~/.roteiro/logs/roteiro.log", Some(home), Some(rh)),
PathBuf::from("/home/u/.roteiro/logs/roteiro.log")
);
assert_eq!(resolve_path_with("~", Some(home), Some(rh)), home);
let got = resolve_path_with("~/logs/roteiro.log", None, Some(rh));
assert_eq!(got, PathBuf::from("/iso/.roteiro/logs/roteiro.log"));
assert!(
!got.components().any(|c| c.as_os_str() == "~"),
"no literal ~ segment survives: {got:?}"
);
let got = resolve_path_with("~/logs/x.log", None, None);
assert_eq!(got, PathBuf::from("logs/x.log"));
assert_eq!(
resolve_path_with("roteiro.log", Some(home), Some(rh)),
PathBuf::from("/iso/.roteiro/roteiro.log")
);
assert_eq!(
resolve_path_with("/var/log/roteiro.log", Some(home), Some(rh)),
PathBuf::from("/var/log/roteiro.log")
);
}
#[test]
fn resolve_precedence_and_defaults() {
let s = resolve(&Overrides::default(), &TelemetryConfig::default()).expect("resolve");
assert!(s.path.is_none(), "no file/flag ⇒ file logging off");
assert_eq!(s.rotation, Rotation::Daily);
assert_eq!(s.format, Format::Otel);
let cfg = TelemetryConfig {
file: Some("/var/log/roteiro.log".to_owned()),
rotation: Some("hourly".to_owned()),
format: Some("json".to_owned()),
};
let over = Overrides {
rotation: Some("never".to_owned()),
..Default::default()
};
let s = resolve(&over, &cfg).expect("resolve");
assert_eq!(
s.path.as_deref(),
Some(std::path::Path::new("/var/log/roteiro.log"))
);
assert_eq!(s.rotation, Rotation::Never, "override beats config");
assert_eq!(s.format, Format::Otel, "json is an alias of otel");
let over = Overrides {
file: Some("/tmp/explicit.log".to_owned()),
..Default::default()
};
let s = resolve(&over, &cfg).expect("resolve");
assert_eq!(
s.path.as_deref(),
Some(std::path::Path::new("/tmp/explicit.log"))
);
let bad = TelemetryConfig {
rotation: Some("weekly".to_owned()),
..Default::default()
};
assert!(
resolve(&Overrides::default(), &bad).is_err(),
"bad rotation errors"
);
}
#[test]
fn otel_file_layer_writes_parsable_line_with_otel_fields() {
let dir = std::env::temp_dir().join(format!("roteiro-telemetry-{}", std::process::id()));
std::fs::remove_dir_all(&dir).ok();
let path = dir.join("roteiro.log");
let (layer, guard) =
file_layer(&path, Rotation::Never, Format::Otel).expect("build file layer");
let subscriber = tracing_subscriber::registry().with(layer);
tracing::subscriber::with_default(subscriber, || {
let span = tracing::info_span!("ingest", repo = "demo");
let _e = span.enter();
tracing::info!(nodes = 42, "sync complete");
});
drop(guard);
let contents = std::fs::read_to_string(&path).expect("log file exists");
let line = contents
.lines()
.find(|l| l.contains("sync complete"))
.expect("our line");
let v: serde_json::Value = serde_json::from_str(line).expect("valid JSON line");
assert_eq!(v["body"], "sync complete", "message ⇒ OTEL body");
assert_eq!(v["severity_text"], "INFO");
assert_eq!(v["severity_number"], 9);
assert!(
v["time_unix_nano"].as_u64().is_some_and(|t| t > 0),
"timestamp present as ns since epoch"
);
assert_eq!(
v["attributes"]["nodes"], 42,
"non-message field ⇒ attribute"
);
assert_eq!(
v["attributes"]["span.name"], "ingest",
"span context surfaced"
);
assert_eq!(v["resource"]["service.name"], "roteiro");
std::fs::remove_dir_all(&dir).ok();
}
#[derive(Clone, Default)]
struct CountingLayer(std::sync::Arc<std::sync::atomic::AtomicUsize>);
impl<S: tracing::Subscriber> tracing_subscriber::Layer<S> for CountingLayer {
fn on_event(
&self,
_event: &tracing::Event<'_>,
_ctx: tracing_subscriber::layer::Context<'_, S>,
) {
self.0.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
}
}
fn native_events_passing(default: tracing::Level, event_level: tracing::Level) -> usize {
use tracing_subscriber::Layer as _;
let counter = CountingLayer::default();
let filter = tracing_subscriber::EnvFilter::builder()
.with_default_directive(default.into())
.parse_lossy("");
let subscriber = tracing_subscriber::registry().with(counter.clone().with_filter(filter));
tracing::subscriber::with_default(subscriber, || match event_level {
tracing::Level::INFO => {
tracing::info!(target: "llama.cpp", "llama_model_loader: loaded meta data");
}
tracing::Level::WARN => tracing::warn!(target: "ggml", "ggml warning"),
_ => tracing::debug!(target: "llama.cpp", "create_tensor: loading"),
});
counter.0.load(std::sync::atomic::Ordering::SeqCst)
}
#[test]
fn native_logs_are_level_gated_like_any_event() {
use tracing::Level;
assert_eq!(
native_events_passing(Level::WARN, Level::INFO),
0,
"native INFO wall is quiet on stdout by default (warn filter)"
);
assert_eq!(
native_events_passing(Level::INFO, Level::INFO),
1,
"native INFO wall is captured when the level is info (file default / --log)"
);
assert_eq!(
native_events_passing(Level::DEBUG, Level::DEBUG),
1,
"native DEBUG lines surface at ROTEIRO_LOG=debug"
);
assert_eq!(
native_events_passing(Level::WARN, Level::DEBUG),
0,
"native DEBUG lines stay hidden by default"
);
assert_eq!(
native_events_passing(Level::WARN, Level::WARN),
1,
"native WARN/ERROR always surface"
);
}
}