macro_rules! vlog {
($($arg:tt)*) => {
if $crate::verbose() {
eprintln!($($arg)*);
}
};
}
pub(crate) use vlog;
pub(crate) fn verbose() -> bool {
use std::sync::OnceLock;
static VERBOSE: OnceLock<bool> = OnceLock::new();
*VERBOSE.get_or_init(|| std::env::var_os("FIBRE_LOGGING_VERBOSE").is_some_and(|v| v != "0"))
}
pub mod config;
pub mod encoders;
pub mod error;
pub mod error_handling;
pub mod init;
pub mod model;
mod roller;
pub mod subscriber;
#[cfg(debug_assertions)]
pub mod debug_report;
#[cfg(not(debug_assertions))]
pub mod debug_report {
#[inline(always)]
pub fn print_debug_report() {}
#[inline(always)]
pub fn clear_debug_report() {}
#[inline(always)]
pub fn debug_report_totals() -> (usize, usize) {
(0, 0)
}
}
pub use error::{Error, Result};
pub use error_handling::{InternalErrorReport, InternalErrorSource};
pub use model::{LogValue, LogEvent};
use std::{
collections::HashMap,
sync::{atomic::AtomicBool, Arc},
thread::JoinHandle,
time::{Duration, Instant},
};
pub type CustomEventReceiver = fibre::mpsc::BoundedSyncReceiver<LogEvent>;
pub type AppenderTaskHandle = JoinHandle<()>;
pub use init::{find_config_file, init_from_file};
const DEFAULT_SHUTDOWN_TIMEOUT: Duration = Duration::from_secs(5);
#[must_use = "The InitResult and its guards must be kept alive for logging to work correctly and flush on exit"]
pub struct InitResult {
pub appender_task_handles: Vec<AppenderTaskHandle>,
pub(crate) appender_task_names: Vec<String>,
pub(crate) shutdown_signal: Arc<AtomicBool>,
pub(crate) processor: Option<Arc<subscriber::EventProcessor>>,
pub internal_error_rx: Option<fibre::mpsc::BoundedSyncReceiver<InternalErrorReport>>,
pub custom_streams: HashMap<String, CustomEventReceiver>,
}
impl InitResult {
pub fn shutdown(mut self, timeout: Duration) {
self.shutdown_impl(timeout);
}
fn shutdown_impl(&mut self, timeout: Duration) {
if self.processor.is_none() && self.appender_task_handles.is_empty() {
return;
}
self
.shutdown_signal
.store(true, std::sync::atomic::Ordering::SeqCst);
if let Some(processor) = self.processor.take() {
processor.close_channels();
}
if self.appender_task_handles.is_empty() {
vlog!("[fibre_logging] Shutdown complete.");
return;
}
vlog!("[fibre_logging] Shutting down. Waiting for appender tasks to flush...");
let mut names = std::mem::take(&mut self.appender_task_names).into_iter();
let mut pending: Vec<(String, AppenderTaskHandle)> = self
.appender_task_handles
.drain(..)
.map(|handle| {
(
names.next().unwrap_or_else(|| "appender".to_string()),
handle,
)
})
.collect();
let deadline = Instant::now() + timeout;
loop {
let mut i = 0;
while i < pending.len() {
if pending[i].1.is_finished() {
let (name, handle) = pending.swap_remove(i);
if let Err(e) = handle.join() {
eprintln!(
"[fibre_logging:ERROR] Appender task '{}' panicked during shutdown: {:?}",
name, e
);
}
} else {
i += 1;
}
}
if pending.is_empty() || Instant::now() >= deadline {
break;
}
std::thread::sleep(Duration::from_millis(10));
}
for (name, _handle) in pending {
eprintln!(
"[fibre_logging:ERROR] Appender task '{}' did not shut down within {:?}; abandoning it.",
name, timeout
);
}
vlog!("[fibre_logging] Shutdown complete.");
}
}
impl Drop for InitResult {
fn drop(&mut self) {
self.shutdown_impl(DEFAULT_SHUTDOWN_TIMEOUT);
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::atomic::Ordering;
use std::thread;
fn make_init_result(
handles: Vec<AppenderTaskHandle>,
names: Vec<String>,
signal: Arc<AtomicBool>,
) -> InitResult {
InitResult {
appender_task_handles: handles,
appender_task_names: names,
shutdown_signal: signal,
processor: None,
internal_error_rx: None,
custom_streams: HashMap::new(),
}
}
#[test]
fn shutdown_joins_cooperative_tasks_promptly() {
let signal = Arc::new(AtomicBool::new(false));
let task_signal = Arc::clone(&signal);
let handle = thread::spawn(move || {
while !task_signal.load(Ordering::Relaxed) {
thread::sleep(Duration::from_millis(5));
}
});
let result = make_init_result(vec![handle], vec!["coop".to_string()], signal);
let start = Instant::now();
result.shutdown(Duration::from_secs(5));
assert!(
start.elapsed() < Duration::from_secs(1),
"cooperative task should be joined well before the timeout"
);
}
#[test]
fn shutdown_abandons_stuck_tasks_at_timeout() {
let signal = Arc::new(AtomicBool::new(false));
let handle = thread::spawn(|| loop {
thread::sleep(Duration::from_secs(60));
});
let result = make_init_result(vec![handle], vec!["stuck".to_string()], signal);
let start = Instant::now();
result.shutdown(Duration::from_millis(200));
let elapsed = start.elapsed();
assert!(elapsed >= Duration::from_millis(200));
assert!(
elapsed < Duration::from_secs(2),
"shutdown must return at the timeout instead of hanging: {:?}",
elapsed
);
}
}