use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
use std::sync::{Arc, Mutex, MutexGuard};
use std::time::Instant;
use tracing::Level;
use tracing_subscriber::Layer;
use tracing_subscriber::layer::Context;
use tracing_subscriber::prelude::*;
#[derive(Debug, Clone)]
pub struct LogEvent {
pub at: Instant,
pub level: Level,
pub target: String,
pub message: String,
}
pub const LOG_WINDOW_CAP: usize = 10_000;
const MARKER_TARGET: &str = "camel_integration_test::log_capture";
struct WindowEntry {
id: u64,
opened_at: Instant,
buffer: Arc<Mutex<Vec<LogEvent>>>,
}
static WINDOWS: Mutex<Vec<WindowEntry>> = Mutex::new(Vec::new());
static NEXT_ID: AtomicU64 = AtomicU64::new(0);
static OWN_INSTALL: AtomicBool = AtomicBool::new(false);
pub struct WindowHandle {
id: u64,
buffer: Arc<Mutex<Vec<LogEvent>>>,
}
impl WindowHandle {
pub fn close(self) -> Vec<LogEvent> {
unregister(self.id);
let mut buffer = lock(&self.buffer);
std::mem::take(&mut *buffer)
}
}
impl Drop for WindowHandle {
fn drop(&mut self) {
unregister(self.id);
}
}
pub fn open_window() -> WindowHandle {
let id = NEXT_ID.fetch_add(1, Ordering::Relaxed);
let buffer = Arc::new(Mutex::new(Vec::new()));
let mut windows = lock(&WINDOWS);
windows.push(WindowEntry {
id,
opened_at: Instant::now(),
buffer: Arc::clone(&buffer),
});
WindowHandle { id, buffer }
}
fn unregister(id: u64) {
let mut windows = lock(&WINDOWS);
windows.retain(|entry| entry.id != id);
}
fn lock<T>(mutex: &Mutex<T>) -> MutexGuard<'_, T> {
mutex
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
}
pub fn capture_installed() -> bool {
OWN_INSTALL.load(Ordering::Acquire)
}
pub fn ensure_capture_subscriber() {
if OWN_INSTALL.load(Ordering::Acquire) {
return;
}
let capture = tracing_subscriber::registry()
.with(CaptureLayer)
.with(tracing_subscriber::fmt::layer());
if capture.try_init().is_ok() {
OWN_INSTALL.store(true, Ordering::Release);
}
}
#[cfg(test)]
pub(crate) fn scoped_capture_dispatch() -> tracing::Dispatch {
tracing::Dispatch::new(tracing_subscriber::registry().with(CaptureLayer))
}
struct CaptureLayer;
impl<S> Layer<S> for CaptureLayer
where
S: tracing::Subscriber,
{
fn on_event(&self, event: &tracing::Event<'_>, _ctx: Context<'_, S>) {
let at = Instant::now();
let mut visitor = MessageVisitor::default();
event.record(&mut visitor);
let log_event = LogEvent {
at,
level: *event.metadata().level(),
target: event.metadata().target().to_string(),
message: visitor.message.unwrap_or_default(),
};
let windows = lock(&WINDOWS);
for window in windows.iter() {
if window.opened_at > at {
continue;
}
let mut buffer = lock(&window.buffer);
if buffer.len() + 2 > LOG_WINDOW_CAP {
let drop_count = buffer.len() + 2 - LOG_WINDOW_CAP;
buffer.drain(..drop_count);
buffer.insert(
0,
LogEvent {
at,
level: Level::TRACE,
target: MARKER_TARGET.to_string(),
message: format!(
"window cap {LOG_WINDOW_CAP} reached: earlier events dropped"
),
},
);
}
buffer.push(log_event.clone());
}
}
}
#[derive(Default)]
struct MessageVisitor {
message: Option<String>,
}
impl tracing::field::Visit for MessageVisitor {
fn record_debug(&mut self, field: &tracing::field::Field, value: &dyn std::fmt::Debug) {
if field.name() == "message" {
self.message = Some(format!("{value:?}"));
}
}
}
#[cfg(test)]
mod tests {
use super::*;
fn scoped_capture() -> tracing::Dispatch {
tracing::Dispatch::new(tracing_subscriber::registry().with(CaptureLayer))
}
#[test]
fn concurrent_windows_attribute_conservatively() {
let dispatch = scoped_capture();
let first = open_window();
let second = open_window();
tracing::dispatcher::with_default(&dispatch, || tracing::info!("shared marker"));
let late = open_window();
tracing::dispatcher::with_default(&dispatch, || tracing::info!("later marker"));
let first_events = first.close();
let late_events = late.close();
let second_events = second.close();
assert!(
first_events
.iter()
.any(|e| e.message.contains("shared marker")),
"first window (open at event time) must capture: {first_events:?}"
);
assert!(
second_events
.iter()
.any(|e| e.message.contains("shared marker")),
"second window (open at event time) must capture: {second_events:?}"
);
assert!(
first_events
.iter()
.any(|e| e.message.contains("later marker"))
);
assert!(
second_events
.iter()
.any(|e| e.message.contains("later marker"))
);
assert!(
!late_events
.iter()
.any(|e| e.message.contains("shared marker")),
"conservative attribution: a window opened after the event never sees it: {late_events:?}"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn spawned_task_events_counted() {
let dispatch = scoped_capture();
let window = open_window();
let task = tokio::spawn(async move {
tracing::dispatcher::with_default(&dispatch, || tracing::warn!("spawned task marker"));
});
task.await.expect("spawned task completes");
let events = window.close();
assert!(
events
.iter()
.any(|e| e.level == Level::WARN && e.message.contains("spawned task marker")),
"the spawned task's warn must land in the open window: {events:?}"
);
}
}