use std::collections::VecDeque;
use std::sync::{Arc, Mutex, MutexGuard, PoisonError};
use std::time::Duration;
#[derive(Clone, Copy)]
pub enum Channel {
Shared,
Keyed,
}
impl Channel {
const fn as_str(self) -> &'static str {
match self {
Self::Shared => "shared",
Self::Keyed => "keyed",
}
}
}
pub fn batch(pulled: usize, updated: usize, shared_pending: usize) {
tracing::trace!(
target: "tears::runtime::load",
pulled,
updated,
shared_pending,
"processed message batch",
);
}
pub fn capacity_wait(channel: Channel, waited: Duration) {
tracing::debug!(
target: "tears::runtime::load",
channel = channel.as_str(),
wait_us = u64::try_from(waited.as_micros()).unwrap_or(u64::MAX),
"capacity wait",
);
}
#[derive(Clone, Default)]
pub struct LoadObserver {
gauges: Arc<Mutex<Gauges>>,
}
#[derive(Default)]
struct Gauges {
seq: u64,
subscriptions: usize,
unkeyed_commands: usize,
keyed_commands: usize,
blocked: usize,
pending: VecDeque<GaugeSnapshot>,
draining: bool,
}
impl Gauges {
const fn capture(&mut self) -> GaugeSnapshot {
self.seq = self.seq.wrapping_add(1);
GaugeSnapshot {
seq: self.seq,
subscriptions: self.subscriptions,
unkeyed_commands: self.unkeyed_commands,
keyed_commands: self.keyed_commands,
blocked: self.blocked,
}
}
}
#[derive(Clone, Copy)]
struct GaugeSnapshot {
seq: u64,
subscriptions: usize,
unkeyed_commands: usize,
keyed_commands: usize,
blocked: usize,
}
impl GaugeSnapshot {
fn dispatch(self) {
tracing::debug!(
target: "tears::runtime::load",
seq = self.seq,
subscriptions = self.subscriptions,
unkeyed_commands = self.unkeyed_commands,
keyed_commands = self.keyed_commands,
blocked = self.blocked,
"producer gauges",
);
}
}
#[derive(Clone, Copy)]
enum Field {
Subscriptions,
UnkeyedCommands,
Blocked,
}
impl Field {
const fn counter_mut(self, gauges: &mut Gauges) -> &mut usize {
match self {
Self::Subscriptions => &mut gauges.subscriptions,
Self::UnkeyedCommands => &mut gauges.unkeyed_commands,
Self::Blocked => &mut gauges.blocked,
}
}
}
impl LoadObserver {
#[must_use]
pub fn track_subscription(&self) -> GaugeGuard {
self.enter(Field::Subscriptions)
}
#[must_use]
pub fn track_unkeyed_command(&self) -> GaugeGuard {
self.enter(Field::UnkeyedCommands)
}
#[must_use]
pub fn track_blocked(&self) -> GaugeGuard {
self.enter(Field::Blocked)
}
pub fn set_keyed_entries(&self, count: usize) {
self.emit(|gauges| {
if gauges.keyed_commands == count {
return false;
}
gauges.keyed_commands = count;
true
});
}
fn enter(&self, field: Field) -> GaugeGuard {
self.step(field, 1);
GaugeGuard {
observer: self.clone(),
field,
}
}
fn step(&self, field: Field, delta: isize) {
self.emit(|gauges| {
let counter = field.counter_mut(gauges);
*counter = counter.wrapping_add_signed(delta);
true
});
}
fn emit(&self, mutate: impl FnOnce(&mut Gauges) -> bool) {
let enabled =
tracing::event_enabled!(target: "tears::runtime::load", tracing::Level::DEBUG);
let first = {
let mut gauges = self.lock();
if !mutate(&mut gauges) {
return;
}
if gauges.draining {
let snapshot = gauges.capture();
gauges.pending.push_back(snapshot);
return;
}
if !enabled {
return;
}
gauges.draining = true;
gauges.capture()
};
let mut release = DrainGuard {
observer: self,
armed: true,
};
let mut next = first;
loop {
next.dispatch();
let popped = {
let mut gauges = self.lock();
let snapshot = gauges.pending.pop_front();
if snapshot.is_none() {
gauges.draining = false;
}
snapshot
};
let Some(snapshot) = popped else {
release.armed = false;
return;
};
next = snapshot;
}
}
fn lock(&self) -> MutexGuard<'_, Gauges> {
self.gauges.lock().unwrap_or_else(PoisonError::into_inner)
}
}
pub struct GaugeGuard {
observer: LoadObserver,
field: Field,
}
impl Drop for GaugeGuard {
fn drop(&mut self) {
self.observer.step(self.field, -1);
}
}
struct DrainGuard<'a> {
observer: &'a LoadObserver,
armed: bool,
}
impl Drop for DrainGuard<'_> {
fn drop(&mut self) {
if self.armed {
let mut gauges = self.observer.lock();
gauges.draining = false;
gauges.pending.clear();
}
}
}
#[cfg(test)]
mod tests {
use std::fmt::Debug;
use std::panic::{self, AssertUnwindSafe};
use std::sync::atomic::{AtomicBool, Ordering};
use tracing::field::{Field, Visit};
use tracing::span::{Attributes, Id, Record};
use tracing::{Event, Level, Metadata, Subscriber};
use super::*;
use crate::test_support::{TraceRecorder, set_default_subscriber, with_silent_panic_hook};
#[test]
fn gauge_event_carries_the_full_field_set() {
let recorder = TraceRecorder::new().with_target("tears::runtime::load");
let _guard = recorder.set_default();
let observer = LoadObserver::default();
let subscription = observer.track_subscription();
observer.set_keyed_entries(2);
drop(subscription);
let gauge_events: Vec<_> = recorder
.field_name_sets()
.into_iter()
.filter(|fields| fields.iter().any(|name| name == "subscriptions"))
.collect();
assert!(!gauge_events.is_empty(), "gauge events should have fired");
for fields in gauge_events {
for required in [
"seq",
"subscriptions",
"unkeyed_commands",
"keyed_commands",
"blocked",
] {
assert!(
fields.iter().any(|name| name == required),
"a gauge event is missing `{required}`: {fields:?}"
);
}
}
}
#[test]
fn each_gauge_change_carries_a_monotone_seq() {
let recorder = TraceRecorder::new().with_target("tears::runtime::load");
let _guard = recorder.set_default();
let observer = LoadObserver::default();
let first = observer.track_subscription();
let second = observer.track_subscription();
drop(second);
drop(first);
assert_eq!(
recorder.u64_values("seq"),
vec![1, 2, 3, 4],
"each of the four gauge changes emits one event with the next `seq`"
);
}
#[test]
fn gauge_changes_made_while_unsubscribed_are_not_lost() {
let observer = LoadObserver::default();
let first = observer.track_subscription();
let second = observer.track_subscription();
drop(first);
let third = observer.track_subscription();
let recorder = TraceRecorder::new().with_target("tears::runtime::load");
let _guard = recorder.set_default();
let fourth = observer.track_subscription();
assert_eq!(
recorder.u64_values("subscriptions"),
vec![3],
"the first event after a subscriber attaches must report the true \
current count (second, third, fourth still held), not a count \
that missed the unobserved changes"
);
drop(second);
drop(third);
drop(fourth);
assert_eq!(
recorder.u64_values("subscriptions"),
vec![3, 2, 1, 0],
"every value reached while subscribed is still emitted, and the \
count never wraps from an unmatched decrement"
);
}
#[test]
fn gauge_events_reach_a_subscriber_that_filters_on_is_event() {
let recorder = TraceRecorder::new().with_target("tears::runtime::load");
let subscriber = EventOnlySubscriber(recorder.clone());
let _guard = set_default_subscriber(subscriber);
let observer = LoadObserver::default();
drop(observer.track_subscription());
assert_eq!(
recorder.u64_values("subscriptions"),
vec![1, 0],
"a subscriber that only answers enabled() for genuine events must \
still see both gauge changes"
);
}
struct EventOnlySubscriber(TraceRecorder);
impl Subscriber for EventOnlySubscriber {
fn enabled(&self, metadata: &Metadata<'_>) -> bool {
metadata.is_event() && self.0.enabled(metadata)
}
fn new_span(&self, span: &Attributes<'_>) -> Id {
self.0.new_span(span)
}
fn record(&self, span: &Id, values: &Record<'_>) {
self.0.record(span, values);
}
fn record_follows_from(&self, span: &Id, follows: &Id) {
self.0.record_follows_from(span, follows);
}
fn event(&self, event: &Event<'_>) {
self.0.event(event);
}
fn enter(&self, span: &Id) {
self.0.enter(span);
}
fn exit(&self, span: &Id) {
self.0.exit(span);
}
}
#[test]
fn each_gauge_change_emits_the_value_it_reached() {
let recorder = TraceRecorder::new().with_target("tears::runtime::load");
let _guard = recorder.set_default();
let observer = LoadObserver::default();
let first = observer.track_subscription();
let second = observer.track_subscription();
drop(second);
drop(first);
assert_eq!(
recorder.u64_values("subscriptions"),
vec![1, 2, 1, 0],
"every reached value, including the peak of 2, is emitted in order"
);
}
#[test]
fn schema_events_fire_at_their_declared_levels() {
let at_trace = TraceRecorder::new()
.with_target("tears::runtime::load")
.with_level(Level::TRACE);
{
let _guard = at_trace.set_default();
batch(1, 1, 0);
}
assert_eq!(
at_trace.u64_values("pulled"),
vec![1],
"batch event is TRACE"
);
let at_debug = TraceRecorder::new()
.with_target("tears::runtime::load")
.with_level(Level::DEBUG);
{
let _guard = at_debug.set_default();
batch(1, 1, 0);
}
assert!(
at_debug.u64_values("pulled").is_empty(),
"batch event is not DEBUG"
);
{
let _guard = at_debug.set_default();
capacity_wait(Channel::Shared, Duration::from_micros(1));
}
assert_eq!(
at_debug.str_values("channel"),
vec!["shared".to_owned()],
"capacity-wait event is DEBUG"
);
{
let _guard = at_debug.set_default();
LoadObserver::default().set_keyed_entries(1);
}
assert_eq!(
at_debug.u64_values("keyed_commands"),
vec![1],
"gauge event is DEBUG"
);
{
let _guard = at_trace.set_default();
LoadObserver::default().set_keyed_entries(1);
}
assert!(
at_trace.u64_values("keyed_commands").is_empty(),
"gauge event is not TRACE"
);
}
#[test]
fn reentrant_gauge_change_from_a_subscriber_is_delivered_not_dropped() {
let observer = LoadObserver::default();
let seen = Arc::new(Mutex::new(Vec::new()));
let subscriber = ReentrantGaugeSubscriber {
observer: observer.clone(),
reentered: Arc::new(AtomicBool::new(false)),
seen: Arc::clone(&seen),
};
let _guard = set_default_subscriber(subscriber);
let _subscription = observer.track_subscription();
let seen = seen
.lock()
.expect("reentrancy seen log mutex should not be poisoned")
.clone();
assert_eq!(
seen,
vec![(1, 0), (2, 1)],
"the re-entrant gauge change must be delivered as its own event with \
a distinct seq and its reached value, neither dropped nor deadlocked"
);
}
struct ReentrantGaugeSubscriber {
observer: LoadObserver,
reentered: Arc<AtomicBool>,
seen: Arc<Mutex<Vec<(u64, u64)>>>,
}
impl Subscriber for ReentrantGaugeSubscriber {
fn enabled(&self, _metadata: &Metadata<'_>) -> bool {
true
}
fn new_span(&self, _span: &Attributes<'_>) -> Id {
Id::from_u64(1)
}
fn record(&self, _span: &Id, _values: &Record<'_>) {}
fn record_follows_from(&self, _span: &Id, _follows: &Id) {}
fn event(&self, event: &Event<'_>) {
if event.metadata().target() != "tears::runtime::load" {
return;
}
let mut visitor = GaugeVisitor::default();
event.record(&mut visitor);
let Some(seq) = visitor.seq else { return };
self.seen
.lock()
.expect("reentrancy seen log mutex should not be poisoned")
.push((seq, visitor.keyed_commands.unwrap_or_default()));
if !self.reentered.swap(true, Ordering::SeqCst) {
self.observer.set_keyed_entries(1);
}
}
fn enter(&self, _span: &Id) {}
fn exit(&self, _span: &Id) {}
}
#[derive(Default)]
struct GaugeVisitor {
seq: Option<u64>,
keyed_commands: Option<u64>,
}
impl Visit for GaugeVisitor {
fn record_u64(&mut self, field: &Field, value: u64) {
match field.name() {
"seq" => self.seq = Some(value),
"keyed_commands" => self.keyed_commands = Some(value),
_ => {}
}
}
fn record_debug(&mut self, _field: &Field, _value: &dyn Debug) {}
}
#[tokio::test(flavor = "current_thread")]
async fn a_subscriber_panic_mid_dispatch_does_not_wedge_the_funnel() {
let observer = LoadObserver::default();
let seen_after = Arc::new(Mutex::new(Vec::new()));
let subscriber = PanicOnceGaugeSubscriber {
panicked: Arc::new(AtomicBool::new(false)),
seen_after: Arc::clone(&seen_after),
};
let _guard = set_default_subscriber(subscriber);
let seen_after = with_silent_panic_hook(async {
let outcome = panic::catch_unwind(AssertUnwindSafe(|| {
observer.set_keyed_entries(1);
}));
assert!(outcome.is_err(), "the subscriber panic must propagate");
observer.set_keyed_entries(2);
seen_after
.lock()
.expect("panic-recovery seen log mutex should not be poisoned")
.clone()
})
.await;
assert_eq!(
seen_after,
vec![2],
"after a subscriber panics mid-dispatch, later gauge events must \
dispatch again rather than pile up behind a wedged drainer"
);
}
struct PanicOnceGaugeSubscriber {
panicked: Arc<AtomicBool>,
seen_after: Arc<Mutex<Vec<u64>>>,
}
impl Subscriber for PanicOnceGaugeSubscriber {
fn enabled(&self, _metadata: &Metadata<'_>) -> bool {
true
}
fn new_span(&self, _span: &Attributes<'_>) -> Id {
Id::from_u64(1)
}
fn record(&self, _span: &Id, _values: &Record<'_>) {}
fn record_follows_from(&self, _span: &Id, _follows: &Id) {}
#[expect(
clippy::panic,
clippy::manual_assert,
reason = "the subscriber intentionally panics on its first event"
)]
fn event(&self, event: &Event<'_>) {
if event.metadata().target() != "tears::runtime::load" {
return;
}
let mut visitor = GaugeVisitor::default();
event.record(&mut visitor);
let Some(seq) = visitor.seq else { return };
if !self.panicked.swap(true, Ordering::SeqCst) {
panic!("subscriber panic mid-dispatch");
}
self.seen_after
.lock()
.expect("panic-recovery seen log mutex should not be poisoned")
.push(seq);
}
fn enter(&self, _span: &Id) {}
fn exit(&self, _span: &Id) {}
}
}