1use super::super::config::{ConsoleFormat, LoggingConfig, Section};
2use anyhow::Context;
3use std::io::Write;
4use std::path::Path;
5use std::sync::{Arc, OnceLock};
6
7use parking_lot::Mutex;
8use tracing_subscriber::{Layer, fmt};
9
10#[cfg(feature = "otel")]
11use super::log_correlation;
12
13#[cfg(feature = "otel")]
15pub type OtelLayer = tracing_opentelemetry::OpenTelemetryLayer<
16 tracing_subscriber::Registry,
17 opentelemetry_sdk::trace::Tracer,
18>;
19#[cfg(not(feature = "otel"))]
20pub type OtelLayer = ();
21
22#[cfg(feature = "otel")]
31type JsonEventFormat<T> = log_correlation::TraceIdJson<T>;
32#[cfg(not(feature = "otel"))]
33type JsonEventFormat<T> = fmt::format::Format<fmt::format::Json, T>;
34
35#[cfg_attr(not(feature = "otel"), allow(unused_variables))]
37fn json_event_format<T>(
38 base: fmt::format::Format<fmt::format::Json, T>,
39 inject_trace_ids: bool,
40) -> JsonEventFormat<T> {
41 #[cfg(feature = "otel")]
42 {
43 log_correlation::TraceIdJson::new(base, inject_trace_ids)
44 }
45 #[cfg(not(feature = "otel"))]
46 {
47 base
48 }
49}
50
51fn matches_crate_prefix(target: &str, crate_name: &str) -> bool {
55 target == crate_name
56 || (target.starts_with(crate_name) && target[crate_name.len()..].starts_with("::"))
57}
58
59use file_rotate::{
62 ContentLimit, FileRotate,
63 compression::Compression,
64 suffix::{AppendTimestamp, FileLimit},
65};
66
67#[derive(Clone)]
68struct RotWriter(Arc<Mutex<FileRotate<AppendTimestamp>>>);
69
70impl<'a> fmt::MakeWriter<'a> for RotWriter {
71 type Writer = RotWriterHandle;
72 fn make_writer(&'a self) -> Self::Writer {
73 RotWriterHandle(self.0.clone())
74 }
75}
76
77#[derive(Clone)]
78struct RotWriterHandle(Arc<Mutex<FileRotate<AppendTimestamp>>>);
79
80impl Write for RotWriterHandle {
81 fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
84 self.0.lock().write(buf)
85 }
86 fn flush(&mut self) -> std::io::Result<()> {
87 self.0.lock().flush()
88 }
89 fn write_all(&mut self, buf: &[u8]) -> std::io::Result<()> {
90 self.0.lock().write_all(buf)
91 }
92 fn write_fmt(&mut self, args: std::fmt::Arguments<'_>) -> std::io::Result<()> {
93 self.0.lock().write_fmt(args)
94 }
95}
96
97#[derive(Clone)]
99struct RoutedWriterHandle(Option<RotWriterHandle>);
100
101impl Write for RoutedWriterHandle {
102 fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
103 if let Some(w) = &mut self.0 {
104 w.write(buf)
105 } else {
106 Ok(buf.len())
107 }
108 }
109 fn flush(&mut self) -> std::io::Result<()> {
110 if let Some(w) = &mut self.0 {
111 w.flush()
112 } else {
113 Ok(())
114 }
115 }
116 fn write_all(&mut self, buf: &[u8]) -> std::io::Result<()> {
117 if let Some(w) = &mut self.0 {
118 w.write_all(buf)
119 } else {
120 Ok(())
121 }
122 }
123 fn write_fmt(&mut self, args: std::fmt::Arguments<'_>) -> std::io::Result<()> {
124 if let Some(w) = &mut self.0 {
125 w.write_fmt(args)
126 } else {
127 Ok(())
128 }
129 }
130}
131
132#[derive(Clone)]
135struct MultiFileRouter {
136 default: Option<RotWriter>, by_prefix: Vec<(String, RotWriter)>, }
139
140impl MultiFileRouter {
141 fn resolve_for(&self, target: &str) -> Option<RotWriterHandle> {
142 for (crate_name, wr) in &self.by_prefix {
143 if matches_crate_prefix(target, crate_name) {
144 return Some(RotWriterHandle(wr.0.clone()));
145 }
146 }
147 self.default.as_ref().map(|w| RotWriterHandle(w.0.clone()))
148 }
149
150 fn is_empty(&self) -> bool {
151 self.default.is_none() && self.by_prefix.is_empty()
152 }
153}
154
155impl<'a> fmt::MakeWriter<'a> for MultiFileRouter {
156 type Writer = RoutedWriterHandle;
157
158 fn make_writer(&'a self) -> Self::Writer {
159 RoutedWriterHandle(self.default.as_ref().map(|w| RotWriterHandle(w.0.clone())))
160 }
161
162 fn make_writer_for(&'a self, meta: &tracing::Metadata<'_>) -> Self::Writer {
163 let target = meta.target();
164 RoutedWriterHandle(self.resolve_for(target))
165 }
166}
167
168struct ConfigData<'a> {
171 default_section: Option<&'a Section>,
172 crate_sections: Vec<(String, &'a Section)>,
173}
174
175fn extract_config_data(cfg: &LoggingConfig) -> ConfigData<'_> {
176 let crate_sections = cfg
177 .iter()
178 .filter(|(k, _)| k.as_str() != "default")
179 .map(|(k, v)| (k.clone(), v))
180 .collect::<Vec<_>>();
181
182 ConfigData {
183 default_section: cfg.get("default"),
184 crate_sections,
185 }
186}
187
188fn create_rotating_writer_at_path(
191 log_path: &Path,
192 max_bytes: usize,
193 max_age_days: Option<u32>,
194 max_backups: Option<usize>,
195) -> Result<RotWriter, Box<dyn std::error::Error + Send + Sync>> {
196 if let Some(parent) = log_path.parent() {
197 std::fs::create_dir_all(parent)?;
198 }
199
200 let max_age_days = max_age_days.unwrap_or(1);
202 let age = chrono::Duration::try_days(i64::from(max_age_days))
203 .with_context(|| format!("Invalid max_age_days: {max_age_days}"))?;
204 let limit = if let Some(n) = max_backups {
205 FileLimit::MaxFiles(n)
206 } else {
207 FileLimit::Age(age)
208 };
209
210 let rot = FileRotate::new(
211 log_path,
212 AppendTimestamp::default(limit),
213 ContentLimit::BytesSurpassed(max_bytes),
214 Compression::None,
215 None,
216 );
217
218 Ok(RotWriter(Arc::new(Mutex::new(rot))))
219}
220
221static CONSOLE_GUARD: OnceLock<WorkerGuard> = OnceLock::new();
227
228#[allow(unknown_lints, de1301_no_print_macros)] pub fn init_logging_unified(
231 cfg: &LoggingConfig,
232 base_dir: &Path,
233 otel_layer: Option<OtelLayer>,
234 inject_trace_ids: bool,
235) {
236 CONSOLE_GUARD.get_or_init(|| {
237 if let Err(e) = tracing_log::LogTracer::init() {
239 eprintln!("LogTracer init skipped: {e}");
240 }
241
242 let data = extract_config_data(cfg);
243
244 if data.crate_sections.is_empty() && data.default_section.is_none() {
245 return init_minimal(otel_layer);
247 }
248
249 let file_router = build_file_router(&data, base_dir);
251
252 let console_targets = build_target_console(&data);
253 let file_targets = build_target_file(&data, file_router.default.is_some());
254
255 let console_format = data
256 .default_section
257 .map(|s| s.console_format)
258 .unwrap_or_default();
259
260 install_subscriber(
261 &console_targets,
262 &file_targets,
263 file_router,
264 console_format,
265 otel_layer,
266 inject_trace_ids,
267 )
268 });
269}
270
271use tracing::level_filters::LevelFilter;
274use tracing_appender::non_blocking::WorkerGuard;
275use tracing_subscriber::filter::Targets;
276
277const NOISY_CRATES: &[&str] = &["h2"];
279
280fn build_target_console(config: &ConfigData) -> Targets {
281 let default_level = config
283 .default_section
284 .and_then(|s| s.console_level)
285 .map_or(LevelFilter::INFO, LevelFilter::from_level);
286
287 let mut targets = Targets::new().with_default(default_level);
289
290 for crate_name in NOISY_CRATES {
292 targets = targets.with_target(*crate_name, LevelFilter::WARN);
293 }
294
295 for (crate_name, section) in &config.crate_sections {
297 if let Some(level) = section.console_level.map(LevelFilter::from_level) {
298 targets = targets.with_target(crate_name.clone(), level);
299 }
300 }
301
302 targets
303}
304
305fn build_target_file(config: &ConfigData, has_default_file: bool) -> Targets {
306 let default_level = if has_default_file {
308 config
309 .default_section
310 .and_then(Section::file_level)
311 .map_or(LevelFilter::INFO, LevelFilter::from_level)
312 } else {
313 LevelFilter::OFF
314 };
315
316 let mut targets = Targets::new().with_default(default_level);
317
318 for (crate_name, section) in &config.crate_sections {
320 if let Some(level) = section.file_level().map(LevelFilter::from_level) {
321 targets = targets.with_target(crate_name.clone(), level);
322 }
323 }
324
325 targets
326}
327
328fn build_file_router(config: &ConfigData, base_dir: &Path) -> MultiFileRouter {
331 let mut router = MultiFileRouter {
332 default: None,
333 by_prefix: Vec::with_capacity(config.crate_sections.len()),
334 };
335
336 if let Some(section) = config.default_section {
337 router.default = create_file_writer(None, section, base_dir);
338 }
339
340 for (crate_name, section) in &config.crate_sections {
341 if let Some(writer) = create_file_writer(Some(crate_name), section, base_dir) {
342 router.by_prefix.push((crate_name.clone(), writer));
343 }
344 }
345
346 router
349 .by_prefix
350 .sort_by(|a, b| b.0.len().cmp(&a.0.len()).then_with(|| a.0.cmp(&b.0)));
351
352 router
353}
354
355trait HasMaxSizeBytes {
356 fn max_size_bytes(&self) -> usize;
357}
358
359const DEFAULT_SECTION_MAX_SIZE_MB: usize = 100;
360
361impl HasMaxSizeBytes for Section {
362 fn max_size_bytes(&self) -> usize {
363 self.max_size_mb
364 .map(|mb| mb * 1024 * 1024)
365 .and_then(|b| usize::try_from(b).ok())
366 .unwrap_or(DEFAULT_SECTION_MAX_SIZE_MB * 1024 * 1024)
367 }
368}
369
370#[allow(unknown_lints, de1301_no_print_macros)] fn create_file_writer(
372 crate_name: Option<&str>,
373 section: &Section,
374 base_dir: &Path,
375) -> Option<RotWriter> {
376 let file = section.file()?;
377
378 let max_bytes = section.max_size_bytes();
379
380 let p = Path::new(file);
381 let log_path = if p.is_absolute() {
382 p.to_path_buf()
383 } else {
384 base_dir.join(p)
385 };
386
387 match create_rotating_writer_at_path(
388 &log_path,
389 max_bytes,
390 section.max_age_days,
391 section.max_backups,
392 ) {
393 Ok(writer) => Some(writer),
394 Err(e) => {
395 match crate_name {
396 Some(crate_name) => eprintln!(
397 "Failed to init log file for subsystem '{}': {} ({})",
398 crate_name,
399 log_path.to_string_lossy(),
400 e,
401 ),
402 None => eprintln!(
403 "Failed to initialize default log file '{}'",
404 log_path.to_string_lossy()
405 ),
406 }
407 None
408 }
409 }
410}
411
412fn stderr_supports_ansi() -> bool {
417 _ = enable_ansi_support::enable_ansi_support();
418 supports_color::on(supports_color::Stream::Stderr).is_some_and(|level| level.has_basic)
419}
420
421#[allow(unknown_lints, de1301_no_print_macros)]
428fn install_subscriber(
429 console_targets: &tracing_subscriber::filter::Targets,
430 file_targets: &tracing_subscriber::filter::Targets,
431 file_router: MultiFileRouter,
432 console_format: ConsoleFormat,
433 #[cfg_attr(not(feature = "otel"), allow(unused_variables))] otel_layer: Option<OtelLayer>,
434 #[cfg_attr(not(feature = "otel"), allow(unused_variables))] inject_trace_ids: bool,
435) -> WorkerGuard {
436 use tracing_subscriber::{EnvFilter, Registry, fmt, layer::SubscriberExt};
437
438 let env: Option<EnvFilter> = EnvFilter::try_from_default_env().ok();
441
442 let (nb_stderr, guard) = tracing_appender::non_blocking(std::io::stderr());
444
445 let base_json_format = || {
449 fmt::format()
450 .json()
451 .with_ansi(false)
452 .with_target(true)
453 .with_level(true)
454 .with_timer(fmt::time::UtcTime::rfc_3339())
455 };
456
457 let (console_text, console_json) = match console_format {
460 ConsoleFormat::Text => (
461 Some(
462 fmt::layer()
463 .with_writer(nb_stderr)
464 .with_ansi(stderr_supports_ansi())
465 .with_target(true)
466 .with_level(true)
467 .with_timer(fmt::time::UtcTime::rfc_3339())
468 .with_filter(console_targets.clone()),
469 ),
470 None,
471 ),
472 ConsoleFormat::Json => (
473 None,
474 Some(
475 fmt::layer()
476 .json()
477 .event_format(json_event_format(base_json_format(), inject_trace_ids))
478 .with_writer(nb_stderr)
479 .with_filter(console_targets.clone()),
480 ),
481 ),
482 };
483
484 let file_layer_opt = if file_router.is_empty() {
486 None
487 } else {
488 Some(
489 fmt::layer()
490 .json()
491 .event_format(json_event_format(base_json_format(), inject_trace_ids))
492 .with_writer(file_router)
493 .with_filter(file_targets.clone()),
494 )
495 };
496
497 let subscriber = {
503 let base = Registry::default();
504
505 #[cfg(feature = "otel")]
506 let base = {
507 let otel_opt = otel_layer.map(|otel| otel.with_filter(console_targets.clone()));
508 base.with(otel_opt)
509 };
510 #[cfg(not(feature = "otel"))]
511 let base = base;
512
513 let base = base.with(env);
514 base.with(console_text)
515 .with(console_json)
516 .with(file_layer_opt)
517 };
518
519 if let Err(e) = tracing::subscriber::set_global_default(subscriber) {
520 eprintln!("tracing subscriber init failed: {e}");
521 }
522
523 guard
524}
525#[allow(unknown_lints, de1301_no_print_macros)]
527fn init_minimal(
532 #[cfg_attr(not(feature = "otel"), allow(unused_variables))] otel: Option<OtelLayer>,
533) -> WorkerGuard {
534 use tracing_subscriber::{EnvFilter, Registry, fmt, layer::SubscriberExt};
535
536 let env = EnvFilter::try_from_default_env().ok();
538
539 let (nb_stderr, guard) = tracing_appender::non_blocking(std::io::stderr());
541
542 let fmt_layer = fmt::layer()
543 .with_writer(nb_stderr)
544 .with_ansi(stderr_supports_ansi())
545 .with_target(true)
546 .with_timer(fmt::time::UtcTime::rfc_3339());
547
548 let subscriber = {
550 let base = Registry::default();
551
552 #[cfg(feature = "otel")]
553 let base = base.with(otel);
554 #[cfg(not(feature = "otel"))]
555 let base = base;
556
557 base.with(env).with(fmt_layer)
558 };
559
560 if let Err(e) = tracing::subscriber::set_global_default(subscriber) {
563 eprintln!("tracing subscriber init failed (minimal): {e}");
564 }
565
566 guard
567}
568
569#[cfg(test)]
570mod tests {
571 use super::*;
572 use std::io::Write;
573
574 fn assert_concurrent_writes<'a, W>(writer: &'a W, log_path: &Path)
576 where
577 W: fmt::MakeWriter<'a> + Sync,
578 W::Writer: Write,
579 {
580 const NUM_THREADS: usize = 8;
581 const LINES_PER_THREAD: usize = 500;
582 const TOTAL_LINES: usize = NUM_THREADS * LINES_PER_THREAD;
583
584 std::thread::scope(|s| {
585 for thread_id in 0..NUM_THREADS {
586 s.spawn(move || {
587 for line_no in 0..LINES_PER_THREAD {
588 let mut handle = writer.make_writer();
589 writeln!(handle, "thread={thread_id} line={line_no}")
590 .expect("write must not fail");
591 }
592 });
593 }
594 });
595
596 let content = std::fs::read_to_string(log_path).expect("failed to read log file");
597 let lines: Vec<&str> = content.lines().collect();
598
599 assert_eq!(
600 lines.len(),
601 TOTAL_LINES,
602 "expected {TOTAL_LINES} lines but found {} ({} records {})",
603 lines.len(),
604 lines.len().abs_diff(TOTAL_LINES),
605 if lines.len() < TOTAL_LINES {
606 "lost"
607 } else {
608 "extra"
609 },
610 );
611
612 for (i, line) in lines.iter().enumerate() {
614 assert!(
615 line.starts_with("thread=") && line.contains(" line="),
616 "corrupted line {i}: {line:?}",
617 );
618 }
619 }
620
621 #[test]
622 fn concurrent_writes_are_not_dropped() {
623 let dir = tempfile::tempdir().expect("failed to create temp dir");
624 let log_path = dir.path().join("test.log");
625
626 let writer = create_rotating_writer_at_path(&log_path, 50 * 1024 * 1024, None, Some(1))
627 .expect("failed to create rotating writer");
628
629 assert_concurrent_writes(&writer, &log_path);
630 }
631
632 #[test]
633 fn concurrent_writes_through_routed_writer() {
634 let dir = tempfile::tempdir().expect("failed to create temp dir");
635 let log_path = dir.path().join("routed.log");
636
637 let rot = create_rotating_writer_at_path(&log_path, 50 * 1024 * 1024, None, Some(1))
638 .expect("failed to create rotating writer");
639
640 let router = MultiFileRouter {
641 default: Some(rot),
642 by_prefix: Vec::new(),
643 };
644
645 assert_concurrent_writes(&router, &log_path);
646 }
647
648 fn tmp_writer(dir: &Path, name: &str) -> (RotWriter, std::path::PathBuf) {
650 let p = dir.join(name);
651 let w = create_rotating_writer_at_path(&p, 50 * 1024 * 1024, None, Some(1))
652 .expect("failed to create rotating writer");
653 (w, p)
654 }
655
656 #[test]
657 fn resolve_for_picks_longest_matching_prefix() {
658 let dir = tempfile::tempdir().expect("failed to create temp dir");
659
660 let (broad, broad_path) = tmp_writer(dir.path(), "broad.log");
661 let (specific, specific_path) = tmp_writer(dir.path(), "specific.log");
662
663 broad.0.lock().write_all(b"BROAD\n").unwrap();
665 specific.0.lock().write_all(b"SPECIFIC\n").unwrap();
666
667 let router = MultiFileRouter {
668 default: None,
669 by_prefix: vec![
670 ("cf-gears::api_gateway".into(), specific),
671 ("cf-gears".into(), broad),
672 ],
673 };
674
675 let mut handle = router
677 .resolve_for("cf-gears::api_gateway::handler")
678 .expect("should resolve");
679 handle.write_all(b"routed\n").unwrap();
680 handle.flush().unwrap();
681
682 let specific_content = std::fs::read_to_string(&specific_path).unwrap();
683 assert!(
684 specific_content.contains("routed"),
685 "expected write to land in specific log, got: {specific_content:?}"
686 );
687
688 let broad_content = std::fs::read_to_string(&broad_path).unwrap();
689 assert!(
690 !broad_content.contains("routed"),
691 "write should NOT appear in broad log, got: {broad_content:?}"
692 );
693 }
694
695 #[test]
699 fn build_file_router_sorts_prefixes_longest_match_wins() {
700 use crate::bootstrap::config::SectionFile;
701
702 let dir = tempfile::tempdir().expect("failed to create temp dir");
703
704 let broad_section = Section {
705 console_format: ConsoleFormat::default(),
706 console_level: None,
707 section_file: Some(SectionFile {
708 file: "broad.log".to_owned(),
709 file_level: None,
710 }),
711 max_age_days: None,
712 max_backups: Some(1),
713 max_size_mb: None,
714 };
715 let specific_section = Section {
716 console_format: ConsoleFormat::default(),
717 console_level: None,
718 section_file: Some(SectionFile {
719 file: "specific.log".to_owned(),
720 file_level: None,
721 }),
722 max_age_days: None,
723 max_backups: Some(1),
724 max_size_mb: None,
725 };
726
727 let config = ConfigData {
730 default_section: None,
731 crate_sections: vec![
732 ("cf-gears".to_owned(), &broad_section),
733 ("cf-gears::api_gateway".to_owned(), &specific_section),
734 ],
735 };
736
737 let router = build_file_router(&config, dir.path());
738
739 let mut handle = router
740 .resolve_for("cf-gears::api_gateway::handler")
741 .expect("should resolve");
742 handle.write_all(b"routed\n").unwrap();
743 handle.flush().unwrap();
744
745 let specific_content = std::fs::read_to_string(dir.path().join("specific.log")).unwrap();
746 assert!(
747 specific_content.contains("routed"),
748 "expected write to land in specific log, got: {specific_content:?}"
749 );
750
751 let broad_content = std::fs::read_to_string(dir.path().join("broad.log")).unwrap();
752 assert!(
753 !broad_content.contains("routed"),
754 "write should NOT appear in broad log, got: {broad_content:?}"
755 );
756 }
757
758 #[test]
759 fn resolve_for_exact_match() {
760 let dir = tempfile::tempdir().expect("failed to create temp dir");
761 let (writer, _) = tmp_writer(dir.path(), "exact.log");
762
763 let router = MultiFileRouter {
764 default: None,
765 by_prefix: vec![("cf-gears".into(), writer)],
766 };
767
768 assert!(
770 router.resolve_for("cf-gears").is_some(),
771 "exact target should match"
772 );
773 assert!(
775 router.resolve_for("cf-gears::sub").is_some(),
776 "submodule target should match"
777 );
778 assert!(
780 router.resolve_for("cf_gears_extra").is_none(),
781 "non-prefix target should not match"
782 );
783 assert!(
784 router.resolve_for("other").is_none(),
785 "unrelated target should not match"
786 );
787 }
788
789 #[test]
790 fn resolve_for_falls_back_to_default() {
791 let dir = tempfile::tempdir().expect("failed to create temp dir");
792 let (default_writer, default_path) = tmp_writer(dir.path(), "default.log");
793
794 default_writer.0.lock().write_all(b"DEFAULT\n").unwrap();
795
796 let router = MultiFileRouter {
797 default: Some(default_writer),
798 by_prefix: vec![],
799 };
800
801 let mut handle = router
803 .resolve_for("unknown_crate::gear")
804 .expect("should fall back to default");
805 handle.write_all(b"fallback\n").unwrap();
806 handle.flush().unwrap();
807
808 let content = std::fs::read_to_string(&default_path).unwrap();
809 assert!(
810 content.contains("fallback"),
811 "expected write to land in default log, got: {content:?}"
812 );
813 }
814}