1use std::io::Write as _;
16use std::path::{Path, PathBuf};
17use std::time::Duration;
18
19use tracing_appender::rolling::{Builder as AppenderBuilder, Rotation};
20use tracing_subscriber::fmt::MakeWriter;
21use tracing_subscriber::layer::SubscriberExt;
22use tracing_subscriber::util::SubscriberInitExt;
23use tracing_subscriber::{reload, EnvFilter, Layer};
24
25use crate::error::{Error, Result};
26use crate::layer::{DigJsonLayer, OwnedStatics};
27use crate::writer::{LossyWriter, WriterGuard};
28use crate::{correlation, dirs, filter, janitor, Service};
29
30const MAINTENANCE_INTERVAL: Duration = Duration::from_secs(3600);
32
33type FilterSetter = Box<dyn Fn(&str) -> Result<()> + Send + Sync>;
35
36pub struct LogGuard {
39 _writer: Option<WriterGuard>,
40 dir: PathBuf,
41 file_error: Option<String>,
42 set_filter: FilterSetter,
43}
44
45impl LogGuard {
46 pub fn log_dir(&self) -> &std::path::Path {
49 &self.dir
50 }
51
52 pub fn file_error(&self) -> Option<&str> {
56 self.file_error.as_deref()
57 }
58
59 pub fn set_filter(&self, directive: &str) -> Result<()> {
61 (self.set_filter)(directive)
62 }
63}
64
65struct FileSink {
68 layer: DigJsonLayer<LossyWriter>,
69 guard: WriterGuard,
70 writer: LossyWriter,
71}
72
73pub fn init(service: Service) -> Result<LogGuard> {
80 init_with_console(service, dirs::log_dir(service.name), std::io::stderr)
81}
82
83fn init_with_console<W>(service: Service, dir: PathBuf, console: W) -> Result<LogGuard>
86where
87 W: for<'w> MakeWriter<'w> + Send + Sync + 'static,
88{
89 let max_bytes = janitor::max_bytes(|key: &str| std::env::var(key).ok());
90
91 let (file_sink, file_error) = match open_file_sink(&dir, service, max_bytes) {
92 Ok(sink) => (Some(sink), None),
93 Err(error) => {
94 warn_file_logging_disabled(&console, &dir, &error);
95 (None, Some(error.to_string()))
96 }
97 };
98
99 let directive = filter::resolve_filter_from_env(filter::read_persisted_level(&dir).as_deref());
100 let env_filter = EnvFilter::try_new(&directive).map_err(|e| Error::Filter {
101 directive: directive.clone(),
102 message: e.to_string(),
103 })?;
104 let (filter_layer, reload_handle) = reload::Layer::new(env_filter);
105
106 let (json_layer, writer_guard, file_writer) = match file_sink {
107 Some(sink) => (Some(sink.layer), Some(sink.guard), Some(sink.writer)),
108 None => (None, None, None),
109 };
110 let console_layer = tracing_subscriber::fmt::layer()
111 .with_writer(console)
112 .compact();
113
114 tracing_subscriber::registry()
115 .with(filter_layer)
116 .with(json_layer)
117 .with(console_layer.boxed())
118 .try_init()
119 .map_err(|_| Error::AlreadyInitialized)?;
120
121 if let Some(writer) = file_writer {
122 spawn_maintenance(dir.clone(), service.name, max_bytes, writer);
123 }
124
125 let set_filter = Box::new(move |directive: &str| -> Result<()> {
126 let new = EnvFilter::try_new(directive).map_err(|e| Error::Filter {
127 directive: directive.to_string(),
128 message: e.to_string(),
129 })?;
130 reload_handle.reload(new).map_err(|e| Error::Filter {
131 directive: directive.to_string(),
132 message: e.to_string(),
133 })
134 });
135
136 Ok(LogGuard {
137 _writer: writer_guard,
138 dir,
139 file_error,
140 set_filter,
141 })
142}
143
144fn open_file_sink(
148 dir: &Path,
149 service: Service,
150 max_bytes: u64,
151) -> std::result::Result<FileSink, Error> {
152 std::fs::create_dir_all(dir).map_err(|source| Error::LogDir {
153 path: dir.to_path_buf(),
154 source,
155 })?;
156
157 let retention = janitor::retention_days(|key: &str| std::env::var(key).ok());
158 janitor::enforce_byte_cap(dir, service.name, max_bytes);
159
160 let appender = AppenderBuilder::new()
161 .rotation(Rotation::DAILY)
162 .filename_prefix(format!("{}.jsonl", service.name))
163 .max_log_files(retention)
164 .build(dir)
165 .map_err(|source| Error::Appender {
166 path: dir.to_path_buf(),
167 source,
168 })?;
169 let (writer, guard) = crate::writer::spawn(appender);
170
171 let statics = OwnedStatics {
172 service: service.name.to_string(),
173 service_version: service.version.to_string(),
174 run_context: service.run_context.as_str().to_string(),
175 run_id: correlation::new_run_id(),
176 parent_op_id: correlation::parent_op_id_from_env(),
177 };
178
179 Ok(FileSink {
180 layer: DigJsonLayer::new(statics, writer.clone()),
181 guard,
182 writer,
183 })
184}
185
186fn warn_file_logging_disabled<W>(console: &W, dir: &Path, error: &Error)
190where
191 W: for<'w> MakeWriter<'w>,
192{
193 let _ = writeln!(
195 console.make_writer(),
196 "WARN dig-logging: file logging is DISABLED for {} ({}). Console logging continues; set \
197 DIG_LOG_DIR to a writable directory to restore JSONL log files.",
198 dir.display(),
199 error
200 );
201}
202
203fn spawn_maintenance(dir: PathBuf, service: &'static str, max_bytes: u64, writer: LossyWriter) {
206 std::thread::Builder::new()
207 .name("dig-logging-maintenance".into())
208 .spawn(move || {
209 let mut last_dropped = 0u64;
210 loop {
211 std::thread::sleep(MAINTENANCE_INTERVAL);
212 janitor::enforce_byte_cap(&dir, service, max_bytes);
213 let dropped = writer.dropped();
214 if dropped > last_dropped {
215 tracing::warn!(
216 target: "dig_logging",
217 dropped,
218 "log lines dropped under backpressure since start"
219 );
220 last_dropped = dropped;
221 }
222 }
223 })
224 .expect("spawn dig-logging maintenance thread");
225}
226
227#[cfg(test)]
228mod tests {
229 use super::*;
230 use std::sync::{Arc, Mutex};
231
232 #[derive(Clone, Default)]
234 struct Captured(Arc<Mutex<Vec<u8>>>);
235
236 impl Captured {
237 fn text(&self) -> String {
238 String::from_utf8_lossy(&self.0.lock().unwrap()).into_owned()
239 }
240 }
241
242 impl std::io::Write for Captured {
243 fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
244 self.0.lock().unwrap().extend_from_slice(buf);
245 Ok(buf.len())
246 }
247 fn flush(&mut self) -> std::io::Result<()> {
248 Ok(())
249 }
250 }
251
252 impl<'a> MakeWriter<'a> for Captured {
253 type Writer = Captured;
254 fn make_writer(&'a self) -> Self::Writer {
255 self.clone()
256 }
257 }
258
259 #[test]
267 fn unwritable_log_dir_degrades_to_console_instead_of_silencing_logging() {
268 let tmp = tempfile::tempdir().unwrap();
269 let blocked = tmp.path().join("blocked");
272 std::fs::write(&blocked, b"not a directory").unwrap();
273
274 let console = Captured::default();
275 let guard = init_with_console(
276 Service {
277 name: "dig-node",
278 version: "9.9.9",
279 run_context: crate::RunContext::Cli,
280 },
281 blocked.clone(),
282 console.clone(),
283 )
284 .expect("an unwritable log dir must not fail init");
285
286 let warning = console.text();
287 assert!(
288 warning.contains(&blocked.display().to_string()),
289 "the warning names the path that failed; got: {warning:?}"
290 );
291 assert!(
292 guard.file_error().is_some(),
293 "the degrade is reportable to the caller"
294 );
295
296 tracing::info!(target: "dig_logging_test", "console still receives this record");
297
298 let logged = console.text();
299 assert!(
300 logged.contains("console still receives this record"),
301 "records must still reach the console sink; got: {logged:?}"
302 );
303 }
304}