use std::fmt::Write;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, LazyLock, Mutex, Once};
const CAPTURE_LIMIT: usize = 64 * 1024;
#[derive(Default)]
struct Transcript {
events: Vec<Box<str>>,
bytes: usize,
truncated: bool,
}
impl Transcript {
fn push(&mut self, text: &str) {
match self.bytes.checked_add(text.len()) {
Some(total) if total <= CAPTURE_LIMIT => {
self.bytes = total;
self.events.push(text.into());
}
_ => self.truncated = true,
}
}
}
struct Subscription {
needle: Box<str>,
events: Arc<Mutex<Transcript>>,
}
static SUBSCRIPTIONS: LazyLock<Mutex<Vec<Subscription>>> = LazyLock::new(|| Mutex::new(Vec::new()));
static INSTALL: Once = Once::new();
static INSTALL_REFUSED: AtomicBool = AtomicBool::new(false);
fn subscriptions() -> &'static Mutex<Vec<Subscription>> {
&SUBSCRIPTIONS
}
fn install_bus() {
INSTALL.call_once(|| {
use tracing_subscriber::layer::SubscriberExt;
let subscriber = tracing_subscriber::registry().with(CaptureBus);
let installed = camber::tracing::subscriber::set_global_default(subscriber);
INSTALL_REFUSED.store(installed.is_err(), Ordering::Release);
});
assert!(
!INSTALL_REFUSED.load(Ordering::Acquire),
"another subscriber already owns this test binary, so no capture can record anything"
);
}
struct CaptureBus;
impl<S> tracing_subscriber::Layer<S> for CaptureBus
where
S: camber::tracing::Subscriber,
{
fn on_event(
&self,
event: &camber::tracing::Event<'_>,
_context: tracing_subscriber::layer::Context<'_, S>,
) {
let subscriptions = subscriptions()
.lock()
.unwrap_or_else(|error| error.into_inner());
match subscriptions.is_empty() {
true => {}
false => fan_out(event, &subscriptions),
}
}
}
fn fan_out(event: &camber::tracing::Event<'_>, subscriptions: &[Subscription]) {
let mut fields = FieldText(String::new());
event.record(&mut fields);
let text = fields.0;
for subscription in subscriptions
.iter()
.filter(|subscription| text.contains(subscription.needle.as_ref()))
{
subscription
.events
.lock()
.unwrap_or_else(|error| error.into_inner())
.push(text.as_str());
}
}
struct FieldText(String);
impl FieldText {
fn open(&mut self, field: &camber::tracing::field::Field) -> &mut String {
match self.0.is_empty() {
true => {}
false => self.0.push(' '),
}
self.0.push_str(field.name());
self.0.push('=');
&mut self.0
}
}
impl camber::tracing::field::Visit for FieldText {
fn record_debug(&mut self, field: &camber::tracing::field::Field, value: &dyn std::fmt::Debug) {
write!(self.open(field), "{value:?}").expect("a write into a String cannot fail");
}
fn record_str(&mut self, field: &camber::tracing::field::Field, value: &str) {
self.open(field).push_str(value);
}
}
pub struct TraceCapture {
events: Arc<Mutex<Transcript>>,
}
impl TraceCapture {
pub fn events(&self) -> Box<[Box<str>]> {
let transcript = self.lock();
assert!(
!transcript.truncated,
"the captured transcript exceeded its {CAPTURE_LIMIT}-byte budget, so it is not the whole record"
);
transcript.events.iter().cloned().collect()
}
pub fn len(&self) -> usize {
self.lock().events.len()
}
pub fn is_empty(&self) -> bool {
self.len() == 0
}
pub fn truncated(&self) -> bool {
self.lock().truncated
}
pub fn recorded(&self, fields: &[&str]) -> bool {
let transcript = self.lock();
let found = transcript
.events
.iter()
.any(|event| fields.iter().all(|field| event.contains(field)));
assert!(
found || !transcript.truncated,
"the captured transcript exceeded its {CAPTURE_LIMIT}-byte budget, so it cannot report {fields:?} as absent"
);
found
}
fn lock(&self) -> std::sync::MutexGuard<'_, Transcript> {
self.events
.lock()
.unwrap_or_else(|error| error.into_inner())
}
}
impl Drop for TraceCapture {
fn drop(&mut self) {
let mut subscriptions = subscriptions()
.lock()
.unwrap_or_else(|error| error.into_inner());
subscriptions.retain(|subscription| !Arc::ptr_eq(&subscription.events, &self.events));
}
}
fn events_saying<'a>(events: &'a [Box<str>], sentence: &str) -> Box<[&'a str]> {
events
.iter()
.map(Box::as_ref)
.filter(|event| event.contains(sentence))
.collect()
}
pub fn only_event<'a>(events: &'a [Box<str>], sentence: &str, label: &str) -> &'a str {
let matching = events_saying(events, sentence);
assert_eq!(
matching.len(),
1,
"{label}: exactly one {sentence:?} event was expected, but {} of the {} captured carry it: {matching:?}",
matching.len(),
events.len()
);
matching[0]
}
pub fn assert_fields(event: &str, fields: &[&str], label: &str) {
assert!(
!fields.is_empty(),
"{label}: no fields were declared, so this event proves nothing: {event}"
);
for field in fields {
assert!(
event.contains(field),
"{label}: the event does not carry {field:?}: {event}"
);
}
}
pub fn assert_field_value(event: &str, name: &str, value: &str, label: &str) {
let recorded = field_value(event, name)
.unwrap_or_else(|| panic!("{label}: the event records no {name} field: {event}"));
assert_eq!(
recorded, value,
"{label}: the event records {name}={recorded}, not {name}={value}: {event}"
);
}
pub fn field_value<'a>(event: &'a str, name: &str) -> Option<&'a str> {
let needle = format!("{name}=");
let (at, _) = event
.match_indices(needle.as_str())
.find(|(at, _)| *at == 0 || event.as_bytes()[at - 1] == b' ')?;
let rest = &event[at + needle.len()..];
Some(rest.split_once(' ').map_or(rest, |(value, _)| value))
}
pub fn assert_message_is_fixed(event: &str, sentence: &str, label: &str) {
let at = event
.find(sentence)
.unwrap_or_else(|| panic!("{label}: the event does not carry {sentence:?}: {event}"));
let tail = &event[at + sentence.len()..];
assert!(
tail.is_empty() || opens_a_field(tail),
"{label}: the fixed message carries interpolated text: {event}"
);
}
fn opens_a_field(tail: &str) -> bool {
tail.starts_with(' ')
&& tail
.split_whitespace()
.next()
.and_then(|token| token.split_once('='))
.is_some_and(|(name, _)| !name.is_empty())
}
pub fn captures_for(needle: &str) -> usize {
subscriptions()
.lock()
.unwrap_or_else(|error| error.into_inner())
.iter()
.filter(|subscription| subscription.needle.as_ref() == needle)
.count()
}
pub fn capture_events(needle: &str) -> TraceCapture {
install_bus();
let events = Arc::new(Mutex::new(Transcript::default()));
subscriptions()
.lock()
.unwrap_or_else(|error| error.into_inner())
.push(Subscription {
needle: needle.into(),
events: Arc::clone(&events),
});
TraceCapture { events }
}