Skip to content

Commit a22a37a

Browse files
committed
feat: gate tracing behind feature flag
1 parent c1d9f12 commit a22a37a

10 files changed

Lines changed: 541 additions & 23 deletions

File tree

Cargo.lock

Lines changed: 319 additions & 0 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

Cargo.toml

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -46,20 +46,22 @@ inventory = { version = "0.3", optional = true }
4646
jaeb-macros = { version = "0.2", path = "jaeb-macros", optional = true }
4747
metrics = { version = "0.24", optional = true }
4848
tokio = { version = "1.51", features = ["rt", "rt-multi-thread", "sync", "macros", "time"] }
49-
tracing = { version = "0.1", features = ["std", "log"] }
49+
tracing = { version = "0.1", features = ["std", "log"], optional = true }
5050

5151
[features]
5252
default = []
5353
macros = ["dep:jaeb-macros", "dep:inventory"]
5454
metrics = ["dep:metrics"]
5555
test-utils = []
56+
trace = ["dep:tracing"]
5657

5758
[dev-dependencies]
5859
criterion = { version = "0.8", features = ["async_tokio"] }
5960
tokio = { version = "1.51", features = ["full", "test-util"] }
6061
async-trait = "0.1"
6162
eventbuzz = { version = "0.2", features = ["asynchronous"] }
6263
evno = "1.0.2"
64+
eventador = { version = "0.0.18", features = ["async"] }
6365

6466
[[bench]]
6567
name = "throughput"

benches/comparison.rs

Lines changed: 99 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,8 +1,10 @@
11
use std::sync::Arc;
2+
use std::sync::atomic::{AtomicBool, Ordering};
23
use std::time::Duration;
34

45
use async_trait::async_trait;
56
use criterion::{Criterion, criterion_group, criterion_main};
7+
use eventador::Eventador;
68
use eventbuzz::asynchronous::prelude::{ApplicationEvent, AsyncApplicationEventListener, AsyncEventbus};
79
use evno::{Bus as EvnoBus, Close as EvnoClose, Emit as EvnoEmit, Guard as EvnoGuard, from_fn as evno_from_fn};
810
use jaeb::{EventBus, EventHandler, HandlerResult};
@@ -99,6 +101,32 @@ fn bench_async_single_listener(c: &mut Criterion) {
99101
})
100102
});
101103

104+
group.bench_function("eventador_publish", |b| {
105+
b.to_async(&rt).iter_custom(|iters| async move {
106+
let bus = Eventador::new(1024).expect("create eventador");
107+
let sub = bus.subscribe::<u64>();
108+
let stop = Arc::new(AtomicBool::new(false));
109+
let stop2 = Arc::clone(&stop);
110+
111+
let drainer = std::thread::spawn(move || {
112+
while !stop2.load(Ordering::Relaxed) {
113+
let _ = sub.recv();
114+
}
115+
});
116+
117+
let start = std::time::Instant::now();
118+
for i in 0..iters {
119+
bus.publish(i);
120+
}
121+
let elapsed = start.elapsed();
122+
123+
stop.store(true, Ordering::Relaxed);
124+
bus.publish(0u64); // sentinel to wake blocked recv
125+
drainer.join().expect("eventador drainer join");
126+
elapsed
127+
})
128+
});
129+
102130
group.finish();
103131
}
104132

@@ -168,6 +196,36 @@ fn bench_async_fanout_10(c: &mut Criterion) {
168196
})
169197
});
170198

199+
group.bench_function("eventador_publish", |b| {
200+
b.to_async(&rt).iter_custom(|iters| async move {
201+
let bus = Eventador::new(1024).expect("create eventador");
202+
let stop = Arc::new(AtomicBool::new(false));
203+
let mut drainers = Vec::with_capacity(10);
204+
for _ in 0..10 {
205+
let sub = bus.subscribe::<u64>();
206+
let stop2 = Arc::clone(&stop);
207+
drainers.push(std::thread::spawn(move || {
208+
while !stop2.load(Ordering::Relaxed) {
209+
let _ = sub.recv();
210+
}
211+
}));
212+
}
213+
214+
let start = std::time::Instant::now();
215+
for i in 0..iters {
216+
bus.publish(i);
217+
}
218+
let elapsed = start.elapsed();
219+
220+
stop.store(true, Ordering::Relaxed);
221+
bus.publish(0u64); // sentinel to wake blocked recv
222+
for d in drainers {
223+
d.join().expect("eventador drainer join");
224+
}
225+
elapsed
226+
})
227+
});
228+
171229
group.finish();
172230
}
173231

@@ -247,6 +305,47 @@ fn bench_contention_4_publishers(c: &mut Criterion) {
247305
// timeouts) reduce the probability but do not eliminate it — the
248306
// standalone benchmark in `benches/evno_contention.rs` still triggers
249307
// timeouts ~2% of samples. See BENCHMARK.md for the full analysis.
308+
309+
group.bench_function("eventador_publish", |b| {
310+
b.to_async(&rt).iter_custom(|iters| async move {
311+
let workers = 4usize;
312+
let per = (iters as usize) / workers;
313+
let extra = (iters as usize) % workers;
314+
315+
let bus = Eventador::new(4096).expect("create eventador");
316+
let sub = bus.subscribe::<u64>();
317+
let stop = Arc::new(AtomicBool::new(false));
318+
let stop2 = Arc::clone(&stop);
319+
320+
let drainer = std::thread::spawn(move || {
321+
while !stop2.load(Ordering::Relaxed) {
322+
let _ = sub.recv();
323+
}
324+
});
325+
326+
let start = std::time::Instant::now();
327+
let mut joins = Vec::with_capacity(workers);
328+
for worker_idx in 0..workers {
329+
let bus = bus.clone();
330+
let n = per + usize::from(worker_idx < extra);
331+
joins.push(tokio::spawn(async move {
332+
for i in 0..n {
333+
bus.publish(i as u64);
334+
}
335+
}));
336+
}
337+
for join in joins {
338+
join.await.expect("join");
339+
}
340+
let elapsed = start.elapsed();
341+
342+
stop.store(true, Ordering::Relaxed);
343+
bus.publish(0u64); // sentinel to wake blocked recv
344+
drainer.join().expect("eventador drainer join");
345+
elapsed
346+
})
347+
});
348+
250349
group.finish();
251350
}
252351

examples/jaeb-demo/Cargo.toml

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -5,7 +5,7 @@ edition = "2024"
55
publish = false
66

77
[dependencies]
8-
jaeb = { path = "../..", features = ["metrics"] }
8+
jaeb = { path = "../..", features = ["metrics", "trace"] }
99
tokio = { version = "1", features = ["full"] }
1010
tracing-subscriber = { version = "0.3", features = ["fmt", "env-filter", "json"] }
1111
metrics-exporter-prometheus = "0.18"

src/bus.rs

Lines changed: 11 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -9,6 +9,7 @@ use std::time::Duration;
99
use arc_swap::ArcSwap;
1010
use tokio::sync::{Mutex, Semaphore, mpsc, oneshot};
1111
use tokio::task::JoinHandle;
12+
#[cfg(feature = "trace")]
1213
use tracing::{error, trace, warn};
1314

1415
use crate::error::{ConfigError, EventBusError};
@@ -349,6 +350,7 @@ impl EventBus {
349350
return Err(EventBusError::Stopped);
350351
}
351352

353+
#[cfg(feature = "trace")]
352354
trace!("event_bus.subscribe {:?}", &registered.name);
353355
let id = self.next_subscription_id();
354356
let listener = (registered.register)(id, subscription_policy, once);
@@ -357,7 +359,7 @@ impl EventBus {
357359
registry.add_listener(TypeId::of::<E>(), std::any::type_name::<E>(), listener);
358360
self.refresh_snapshot_locked(&registry).await;
359361

360-
Ok(Subscription::new(id, self.clone()))
362+
Ok(Subscription::new(id, registered.name, self.clone()))
361363
}
362364

363365
/// Register a sync handler that receives [`DeadLetter`] events.
@@ -450,7 +452,7 @@ impl EventBus {
450452
let mut registry = self.inner.registry.lock().await;
451453
registry.add_global_middleware(id, erased);
452454
self.refresh_snapshot_locked(&registry).await;
453-
Ok(Subscription::new(id, self.clone()))
455+
Ok(Subscription::new(id, None, self.clone()))
454456
}
455457

456458
/// Add a global **sync** middleware that intercepts all event types.
@@ -474,7 +476,7 @@ impl EventBus {
474476
let mut registry = self.inner.registry.lock().await;
475477
registry.add_global_middleware(id, erased);
476478
self.refresh_snapshot_locked(&registry).await;
477-
Ok(Subscription::new(id, self.clone()))
479+
Ok(Subscription::new(id, None, self.clone()))
478480
}
479481

480482
/// Add an async middleware scoped to a single event type `E`.
@@ -514,7 +516,7 @@ impl EventBus {
514516
let mut registry = self.inner.registry.lock().await;
515517
registry.add_typed_middleware(TypeId::of::<E>(), std::any::type_name::<E>(), slot.clone());
516518
self.refresh_snapshot_locked(&registry).await;
517-
Ok(Subscription::new(slot.id, self.clone()))
519+
Ok(Subscription::new(slot.id, None, self.clone()))
518520
}
519521

520522
/// Add a **sync** middleware scoped to a single event type `E`.
@@ -549,7 +551,7 @@ impl EventBus {
549551
let mut registry = self.inner.registry.lock().await;
550552
registry.add_typed_middleware(TypeId::of::<E>(), std::any::type_name::<E>(), slot.clone());
551553
self.refresh_snapshot_locked(&registry).await;
552-
Ok(Subscription::new(slot.id, self.clone()))
554+
Ok(Subscription::new(slot.id, None, self.clone()))
553555
}
554556

555557
async fn publish_erased(
@@ -645,11 +647,12 @@ impl EventBus {
645647
tokio::spawn(async move {
646648
let _keep = permit;
647649
let dispatch_ctx = bus.inner.full_dispatch_context();
648-
if let Err(err) = bus
650+
if let Err(_err) = bus
649651
.publish_erased(TypeId::of::<E>(), Arc::new(event), std::any::type_name::<E>(), &dispatch_ctx)
650652
.await
651653
{
652-
error!(error = %err, "event_bus.try_publish.dispatch_failed");
654+
#[cfg(feature = "trace")]
655+
error!(error = %_err, "event_bus.try_publish.dispatch_failed");
653656
}
654657
});
655658
Ok(())
@@ -786,5 +789,6 @@ async fn control_loop(inner: std::sync::Weak<Inner>, mut notify_rx: mpsc::Unboun
786789
}
787790
}
788791
}
792+
#[cfg(feature = "trace")]
789793
warn!("event_bus.control_loop.stopped");
790794
}

src/handler.rs

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -239,7 +239,7 @@ where
239239
let Some(event) = event.downcast_ref::<E>() else {
240240
return Err("event type mismatch".into());
241241
};
242-
(handler)(event)
242+
handler(event)
243243
});
244244
ListenerEntry {
245245
id,
@@ -275,7 +275,7 @@ where
275275
Box::pin(async move {
276276
let event = event.map_err(|_| "event type mismatch")?;
277277
let event = (*event).clone();
278-
(handler)(event).await
278+
handler(event).await
279279
})
280280
});
281281
ListenerEntry {

src/metrics.rs

Lines changed: 8 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -6,15 +6,17 @@ pub(crate) struct TimerGuard {
66
start: std::time::Instant,
77
name: &'static str,
88
event: &'static str,
9+
listener: Option<&'static str>,
910
}
1011

1112
#[cfg(feature = "metrics")]
1213
impl TimerGuard {
13-
pub fn start(name: &'static str, event: &'static str) -> Self {
14+
pub fn start(name: &'static str, event: &'static str, listener: Option<&'static str>) -> Self {
1415
Self {
1516
start: std::time::Instant::now(),
1617
name,
1718
event,
19+
listener,
1820
}
1921
}
2022
}
@@ -23,7 +25,11 @@ impl TimerGuard {
2325
impl Drop for TimerGuard {
2426
fn drop(&mut self) {
2527
let dur = self.start.elapsed();
26-
let histogram = histogram!(self.name, "event" => self.event);
28+
let histogram = if let Some(listener) = self.listener {
29+
histogram!(self.name, "event" => self.event, "listener" => listener)
30+
} else {
31+
histogram!(self.name, "event" => self.event)
32+
};
2733
histogram.record(dur.as_secs_f64());
2834
}
2935
}

src/registry.rs

Lines changed: 12 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -11,6 +11,7 @@ use std::time::Duration;
1111

1212
use tokio::sync::{Mutex, Notify, Semaphore, mpsc, oneshot};
1313
use tokio::task::AbortHandle;
14+
#[cfg(feature = "trace")]
1415
use tracing::{error, warn};
1516

1617
#[cfg(feature = "metrics")]
@@ -506,7 +507,7 @@ async fn execute_async_listener(
506507
attempts += 1;
507508

508509
#[cfg(feature = "metrics")]
509-
let _timer = TimerGuard::start("eventbus.handler.duration", event_name);
510+
let _timer = TimerGuard::start("eventbus.handler.duration", event_name, listener.listener_name);
510511

511512
let handler_future = CatchUnwindFuture::new(handler(Arc::clone(&event)));
512513

@@ -539,9 +540,11 @@ async fn execute_async_listener(
539540
}
540541

541542
retries_left -= 1;
543+
#[cfg(feature = "trace")]
542544
warn!(
543545
event = event_name,
544546
listener_id = listener.subscription_id.as_u64(),
547+
listener_name = listener.listener_name,
545548
attempts,
546549
retries_left,
547550
error = %error_message,
@@ -625,7 +628,7 @@ async fn dispatch_slot(slot: &TypeSlot, event: &EventType, event_name: &'static
625628
}
626629

627630
#[cfg(feature = "metrics")]
628-
let _timer = TimerGuard::start("eventbus.handler.duration", event_name);
631+
let _timer = TimerGuard::start("eventbus.handler.duration", event_name, listener.name);
629632

630633
let result = catch_unwind(AssertUnwindSafe(|| {
631634
// Safe because we only pass shared references to listener code.
@@ -696,16 +699,22 @@ pub(crate) async fn dispatch_with_snapshot(
696699
}
697700

698701
pub(crate) fn dead_letter_from_failure(failure: &ListenerFailure) -> Option<DeadLetter> {
702+
#[cfg(feature = "trace")]
699703
error!(
700704
event = failure.event_name,
701705
listener_id = failure.subscription_id.as_u64(),
706+
listener_name = failure.listener_name,
702707
attempts = failure.attempts,
703708
error = %failure.error,
704709
"handler.failed"
705710
);
706711

707712
#[cfg(feature = "metrics")]
708-
counter!("eventbus.handler.error", "event" => failure.event_name).increment(1);
713+
if let Some(name) = failure.listener_name {
714+
counter!("eventbus.handler.error", "event" => failure.event_name, "listener" => name).increment(1);
715+
} else {
716+
counter!("eventbus.handler.error", "event" => failure.event_name).increment(1);
717+
}
709718

710719
let dead_letter_type = std::any::type_name::<DeadLetter>();
711720
if failure.dead_letter && failure.event_name != dead_letter_type {

0 commit comments

Comments
 (0)