use std::io::Write as _;
use std::path::{Path, PathBuf};
use std::time::Duration;
use tracing_appender::rolling::{Builder as AppenderBuilder, Rotation};
use tracing_subscriber::fmt::MakeWriter;
use tracing_subscriber::layer::SubscriberExt;
use tracing_subscriber::util::SubscriberInitExt;
use tracing_subscriber::{reload, EnvFilter, Layer};
use crate::error::{Error, Result};
use crate::layer::{DigJsonLayer, OwnedStatics};
use crate::writer::{LossyWriter, WriterGuard};
use crate::{correlation, dirs, filter, janitor, Service};
const MAINTENANCE_INTERVAL: Duration = Duration::from_secs(3600);
type FilterSetter = Box<dyn Fn(&str) -> Result<()> + Send + Sync>;
pub struct LogGuard {
_writer: Option<WriterGuard>,
dir: PathBuf,
file_error: Option<String>,
set_filter: FilterSetter,
}
impl LogGuard {
pub fn log_dir(&self) -> &std::path::Path {
&self.dir
}
pub fn file_error(&self) -> Option<&str> {
self.file_error.as_deref()
}
pub fn set_filter(&self, directive: &str) -> Result<()> {
(self.set_filter)(directive)
}
}
struct FileSink {
layer: DigJsonLayer<LossyWriter>,
guard: WriterGuard,
writer: LossyWriter,
}
pub fn init(service: Service) -> Result<LogGuard> {
init_with_console(service, dirs::log_dir(service.name), std::io::stderr)
}
fn init_with_console<W>(service: Service, dir: PathBuf, console: W) -> Result<LogGuard>
where
W: for<'w> MakeWriter<'w> + Send + Sync + 'static,
{
let max_bytes = janitor::max_bytes(|key: &str| std::env::var(key).ok());
let (file_sink, file_error) = match open_file_sink(&dir, service, max_bytes) {
Ok(sink) => (Some(sink), None),
Err(error) => {
warn_file_logging_disabled(&console, &dir, &error);
(None, Some(error.to_string()))
}
};
let directive = filter::resolve_filter_from_env(filter::read_persisted_level(&dir).as_deref());
let env_filter = EnvFilter::try_new(&directive).map_err(|e| Error::Filter {
directive: directive.clone(),
message: e.to_string(),
})?;
let (filter_layer, reload_handle) = reload::Layer::new(env_filter);
let (json_layer, writer_guard, file_writer) = match file_sink {
Some(sink) => (Some(sink.layer), Some(sink.guard), Some(sink.writer)),
None => (None, None, None),
};
let console_layer = tracing_subscriber::fmt::layer()
.with_writer(console)
.compact();
tracing_subscriber::registry()
.with(filter_layer)
.with(json_layer)
.with(console_layer.boxed())
.try_init()
.map_err(|_| Error::AlreadyInitialized)?;
if let Some(writer) = file_writer {
spawn_maintenance(dir.clone(), service.name, max_bytes, writer);
}
let set_filter = Box::new(move |directive: &str| -> Result<()> {
let new = EnvFilter::try_new(directive).map_err(|e| Error::Filter {
directive: directive.to_string(),
message: e.to_string(),
})?;
reload_handle.reload(new).map_err(|e| Error::Filter {
directive: directive.to_string(),
message: e.to_string(),
})
});
Ok(LogGuard {
_writer: writer_guard,
dir,
file_error,
set_filter,
})
}
fn open_file_sink(
dir: &Path,
service: Service,
max_bytes: u64,
) -> std::result::Result<FileSink, Error> {
std::fs::create_dir_all(dir).map_err(|source| Error::LogDir {
path: dir.to_path_buf(),
source,
})?;
let retention = janitor::retention_days(|key: &str| std::env::var(key).ok());
janitor::enforce_byte_cap(dir, service.name, max_bytes);
let appender = AppenderBuilder::new()
.rotation(Rotation::DAILY)
.filename_prefix(format!("{}.jsonl", service.name))
.max_log_files(retention)
.build(dir)
.map_err(|source| Error::Appender {
path: dir.to_path_buf(),
source,
})?;
let (writer, guard) = crate::writer::spawn(appender);
let statics = OwnedStatics {
service: service.name.to_string(),
service_version: service.version.to_string(),
run_context: service.run_context.as_str().to_string(),
run_id: correlation::new_run_id(),
parent_op_id: correlation::parent_op_id_from_env(),
};
Ok(FileSink {
layer: DigJsonLayer::new(statics, writer.clone()),
guard,
writer,
})
}
fn warn_file_logging_disabled<W>(console: &W, dir: &Path, error: &Error)
where
W: for<'w> MakeWriter<'w>,
{
let _ = writeln!(
console.make_writer(),
"WARN dig-logging: file logging is DISABLED for {} ({}). Console logging continues; set \
DIG_LOG_DIR to a writable directory to restore JSONL log files.",
dir.display(),
error
);
}
fn spawn_maintenance(dir: PathBuf, service: &'static str, max_bytes: u64, writer: LossyWriter) {
std::thread::Builder::new()
.name("dig-logging-maintenance".into())
.spawn(move || {
let mut last_dropped = 0u64;
loop {
std::thread::sleep(MAINTENANCE_INTERVAL);
janitor::enforce_byte_cap(&dir, service, max_bytes);
let dropped = writer.dropped();
if dropped > last_dropped {
tracing::warn!(
target: "dig_logging",
dropped,
"log lines dropped under backpressure since start"
);
last_dropped = dropped;
}
}
})
.expect("spawn dig-logging maintenance thread");
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::{Arc, Mutex};
#[derive(Clone, Default)]
struct Captured(Arc<Mutex<Vec<u8>>>);
impl Captured {
fn text(&self) -> String {
String::from_utf8_lossy(&self.0.lock().unwrap()).into_owned()
}
}
impl std::io::Write for Captured {
fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
self.0.lock().unwrap().extend_from_slice(buf);
Ok(buf.len())
}
fn flush(&mut self) -> std::io::Result<()> {
Ok(())
}
}
impl<'a> MakeWriter<'a> for Captured {
type Writer = Captured;
fn make_writer(&'a self) -> Self::Writer {
self.clone()
}
}
#[test]
fn unwritable_log_dir_degrades_to_console_instead_of_silencing_logging() {
let tmp = tempfile::tempdir().unwrap();
let blocked = tmp.path().join("blocked");
std::fs::write(&blocked, b"not a directory").unwrap();
let console = Captured::default();
let guard = init_with_console(
Service {
name: "dig-node",
version: "9.9.9",
run_context: crate::RunContext::Cli,
},
blocked.clone(),
console.clone(),
)
.expect("an unwritable log dir must not fail init");
let warning = console.text();
assert!(
warning.contains(&blocked.display().to_string()),
"the warning names the path that failed; got: {warning:?}"
);
assert!(
guard.file_error().is_some(),
"the degrade is reportable to the caller"
);
tracing::info!(target: "dig_logging_test", "console still receives this record");
let logged = console.text();
assert!(
logged.contains("console still receives this record"),
"records must still reach the console sink; got: {logged:?}"
);
}
}