#[cfg(test)]
mod test_stats {
use anyhow::Result;
use kerf::{
ArcEvent, Config, Event, EventType, Kerf, Level, Location, Match, MatcherSet, StatsConfig,
StatsTracker,
};
use std::{
collections::HashMap,
sync::Arc,
time::{Duration, UNIX_EPOCH},
};
use tokio::sync::Mutex;
#[allow(clippy::too_many_arguments)]
fn create_test_event_for_stats(
id: u64,
level: Level,
message: &str,
module_path: Option<&str>,
file: Option<&str>,
line: Option<u32>,
span_name: Option<&str>,
) -> ArcEvent {
let mut event = Event {
id,
timestamp: chrono::Local::now(),
level,
target: "test_target".to_string(),
name: "test_event".to_string(),
module_path: module_path.map(|s| s.to_string()),
file: file.map(|s| s.to_string()),
line,
message: message.to_string(),
fields: HashMap::new(),
span_name: span_name.map(|s| s.to_string()),
span_hierarchy: span_name.map(|s| s.to_string()),
};
event
.fields
.insert("test_field".to_string(), "test_value".to_string());
Arc::new(event)
}
async fn send_and_process_event(tracer: &Kerf, event: ArcEvent) {
let _ = tracer._get_sender_for_testing().send(event);
tokio::time::sleep(Duration::from_millis(50)).await;
}
#[test]
fn test_stats_config_defaults() {
let config = StatsConfig::default();
assert!(config.track_by_location);
assert!(config.track_by_module);
assert!(config.track_by_level);
assert_eq!(config.max_locations, 10_000);
assert_eq!(config.max_modules, 1_000);
}
#[test]
fn test_stats_config_serialization() -> Result<()> {
let config = StatsConfig {
track_by_location: true,
track_by_module: false,
track_by_level: true,
max_locations: 500,
max_modules: 100,
};
let json = serde_json::to_string(&config)?;
let deserialized: StatsConfig = serde_json::from_str(&json)?;
assert_eq!(config.track_by_location, deserialized.track_by_location);
assert_eq!(config.track_by_module, deserialized.track_by_module);
assert_eq!(config.track_by_level, deserialized.track_by_level);
assert_eq!(config.max_locations, deserialized.max_locations);
assert_eq!(config.max_modules, deserialized.max_modules);
Ok(())
}
#[test]
fn test_location_creation_and_display() {
let location = Location {
file: "test.rs".to_string(),
line: 42,
};
assert_eq!(location.to_string(), "test.rs:42");
assert_eq!(location.file, "test.rs");
assert_eq!(location.line, 42);
let event = create_test_event_for_stats(
1,
Level::INFO,
"Test",
Some("test_module"),
Some("test.rs"),
Some(42),
None,
);
let location_from_event = Location::from_trace_data(&event);
assert!(location_from_event.is_some());
let loc = location_from_event.unwrap();
assert_eq!(loc.file, "test.rs");
assert_eq!(loc.line, 42);
let event_no_location = create_test_event_for_stats(
2,
Level::INFO,
"Test",
Some("test_module"),
None,
None,
None,
);
let no_location = Location::from_trace_data(&event_no_location);
assert!(no_location.is_none());
}
#[test]
fn test_event_type_display() {
assert_eq!(EventType::Captured.to_string(), "Captured");
assert_eq!(EventType::Silenced.to_string(), "Silenced");
assert_eq!(EventType::Dropped.to_string(), "Dropped");
}
#[test]
fn test_atomic_counter() {
let counter = kerf::AtomicCounter::new();
assert_eq!(counter.get(), 0);
let first_increment = counter.increment();
assert_eq!(first_increment, 0); assert_eq!(counter.get(), 1);
let second_increment = counter.increment();
assert_eq!(second_increment, 1);
assert_eq!(counter.get(), 2);
counter.reset();
assert_eq!(counter.get(), 0);
}
#[test]
fn test_stat_entry() {
let entry = kerf::StatEntry::new();
assert_eq!(entry.get_count(), 0);
let initial_time = entry.get_last_seen();
assert!(initial_time >= UNIX_EPOCH);
entry.record_event();
assert_eq!(entry.get_count(), 1);
let updated_time = entry.get_last_seen();
assert!(updated_time >= initial_time);
entry.record_event();
assert_eq!(entry.get_count(), 2);
}
#[test]
fn test_stats_tracker_basic() {
let config = StatsConfig::default();
let tracker = StatsTracker::new(config);
let event = create_test_event_for_stats(
1,
Level::INFO,
"Test event",
Some("test_module"),
Some("test.rs"),
Some(42),
None,
);
tracker.record_event(EventType::Captured, &event);
assert_eq!(tracker.get_total_count(EventType::Captured, Level::INFO), 1);
assert_eq!(tracker.get_total_count(EventType::Silenced, Level::INFO), 0);
assert_eq!(tracker.get_total_count(EventType::Dropped, Level::INFO), 0);
assert_eq!(tracker.get_total_count_by_type(EventType::Captured), 1);
assert_eq!(tracker.get_total_count_by_type(EventType::Silenced), 0);
assert_eq!(tracker.get_total_count_by_type(EventType::Dropped), 0);
tracker.record_event(EventType::Silenced, &event);
assert_eq!(tracker.get_total_count(EventType::Captured, Level::INFO), 1);
assert_eq!(tracker.get_total_count(EventType::Silenced, Level::INFO), 1);
assert_eq!(tracker.get_total_count_by_type(EventType::Captured), 1);
assert_eq!(tracker.get_total_count_by_type(EventType::Silenced), 1);
}
#[test]
fn test_stats_tracker_snapshot() {
let config = StatsConfig::default();
let tracker = StatsTracker::new(config);
let event1 = create_test_event_for_stats(
1,
Level::INFO,
"Event 1",
Some("module_a"),
Some("file_a.rs"),
Some(10),
None,
);
let event2 = create_test_event_for_stats(
2,
Level::ERROR,
"Event 2",
Some("module_b"),
Some("file_b.rs"),
Some(20),
None,
);
let event3 = create_test_event_for_stats(
3,
Level::INFO,
"Event 3",
Some("module_a"),
Some("file_a.rs"),
Some(10), None,
);
tracker.record_event(EventType::Captured, &event1);
tracker.record_event(EventType::Silenced, &event2);
tracker.record_event(EventType::Captured, &event3);
let snapshot = tracker.get_snapshot();
assert!(snapshot.total_entries > 0);
assert_eq!(snapshot.location_stats.len(), 2);
let location_a = Location {
file: "file_a.rs".to_string(),
line: 10,
};
let location_b = Location {
file: "file_b.rs".to_string(),
line: 20,
};
assert!(snapshot.location_stats.contains_key(&location_a));
assert!(snapshot.location_stats.contains_key(&location_b));
let loc_a_stats = &snapshot.location_stats[&location_a];
assert_eq!(loc_a_stats.get_total_for_type(EventType::Captured), 2); assert_eq!(loc_a_stats.get_total_for_type(EventType::Silenced), 0);
let loc_b_stats = &snapshot.location_stats[&location_b];
assert_eq!(loc_b_stats.get_total_for_type(EventType::Captured), 0);
assert_eq!(loc_b_stats.get_total_for_type(EventType::Silenced), 1);
assert_eq!(snapshot.module_stats.len(), 2);
assert!(snapshot.module_stats.contains_key("module_a"));
assert!(snapshot.module_stats.contains_key("module_b"));
let module_a_stats = &snapshot.module_stats["module_a"];
assert_eq!(module_a_stats.get_total_for_type(EventType::Captured), 2); assert_eq!(module_a_stats.get_total_for_type(EventType::Silenced), 0);
let module_b_stats = &snapshot.module_stats["module_b"];
assert_eq!(module_b_stats.get_total_for_type(EventType::Captured), 0);
assert_eq!(module_b_stats.get_total_for_type(EventType::Silenced), 1);
assert_eq!(
snapshot.level_event_counts[&(Level::INFO, EventType::Captured)],
2
); assert_eq!(
snapshot.level_event_counts[&(Level::ERROR, EventType::Silenced)],
1
); }
#[test]
fn test_stats_tracker_limits() {
let config = StatsConfig {
track_by_location: true,
track_by_module: true,
track_by_level: true,
max_locations: 2, max_modules: 2, };
let tracker = StatsTracker::new(config);
for i in 0..5 {
let event = create_test_event_for_stats(
i,
Level::INFO,
&format!("Event {i}"),
Some(&format!("module_{i}")),
Some(&format!("file_{i}.rs")),
Some(i as u32),
None,
);
tracker.record_event(EventType::Captured, &event);
}
let snapshot = tracker.get_snapshot();
assert!(snapshot.location_stats.len() <= 2);
assert!(snapshot.module_stats.len() <= 2);
assert_eq!(tracker.get_total_count_by_type(EventType::Captured), 5);
}
#[test]
fn test_stats_tracker_clear() {
let config = StatsConfig::default();
let tracker = StatsTracker::new(config);
let event = create_test_event_for_stats(
1,
Level::INFO,
"Test event",
Some("test_module"),
Some("test.rs"),
Some(42),
None,
);
tracker.record_event(EventType::Captured, &event);
tracker.record_event(EventType::Silenced, &event);
assert_eq!(tracker.get_total_count_by_type(EventType::Captured), 1);
assert_eq!(tracker.get_total_count_by_type(EventType::Silenced), 1);
assert!(tracker.get_snapshot().total_entries > 0);
tracker.clear();
assert_eq!(tracker.get_total_count_by_type(EventType::Captured), 0);
assert_eq!(tracker.get_total_count_by_type(EventType::Silenced), 0);
assert_eq!(tracker.get_snapshot().total_entries, 0);
assert!(tracker.get_snapshot().location_stats.is_empty());
assert!(tracker.get_snapshot().module_stats.is_empty());
assert!(tracker.get_snapshot().level_event_counts.is_empty());
}
#[test]
fn test_stats_with_disabled_tracking() {
let config = StatsConfig {
track_by_location: false,
track_by_module: false,
track_by_level: true,
max_locations: 100,
max_modules: 100,
};
let tracker = StatsTracker::new(config);
let event = create_test_event_for_stats(
1,
Level::INFO,
"Test event",
Some("test_module"),
Some("test.rs"),
Some(42),
None,
);
tracker.record_event(EventType::Captured, &event);
assert_eq!(tracker.get_total_count_by_type(EventType::Captured), 1);
let snapshot = tracker.get_snapshot();
assert_eq!(snapshot.total_entries, 0);
assert!(snapshot.location_stats.is_empty());
assert!(snapshot.module_stats.is_empty());
}
#[tokio::test]
async fn test_stats_integration_with_tracer() -> Result<()> {
let stats_config = StatsConfig {
track_by_location: true,
track_by_module: true,
track_by_level: true,
max_locations: 1000,
max_modules: 100,
};
let matcher_set = MatcherSet::from_matcher(Match::info().module_pattern("test_*"));
let config = Config::from_tab(("test_tab", matcher_set)).with_stats(stats_config);
let tracer = Kerf::new_with_config(config);
let captured_events = Arc::new(Mutex::new(Vec::<String>::new()));
let captured_clone = captured_events.clone();
tracer
.set_callback(move |event, _tabs| {
let captured = captured_clone.clone();
let msg = event.message.clone();
tokio::spawn(async move {
let mut lock = captured.lock().await;
lock.push(msg);
});
})?
.await??;
let captured_event = create_test_event_for_stats(
1,
Level::INFO,
"Captured event",
Some("test_module"),
Some("test.rs"),
Some(42),
None,
);
let dropped_event = create_test_event_for_stats(
2,
Level::INFO,
"Dropped event",
Some("other_module"),
Some("other.rs"),
Some(43),
None,
);
send_and_process_event(&tracer, captured_event).await;
send_and_process_event(&tracer, dropped_event).await;
{
let captured = captured_events.lock().await;
assert_eq!(captured.len(), 1);
assert!(captured.contains(&"Captured event".to_string()));
}
if let Ok(stats_rx) = tracer.get_stats() {
if let Ok(Some(snapshot)) = stats_rx.await {
assert!(snapshot.total_entries > 0);
let mut total_captured = 0;
let mut total_dropped = 0;
for ((_level, event_type), count) in &snapshot.level_event_counts {
match event_type {
EventType::Captured => total_captured += count,
EventType::Dropped => total_dropped += count,
_ => {}
}
}
assert_eq!(total_captured, 1);
assert_eq!(total_dropped, 1);
let test_location = Location {
file: "test.rs".to_string(),
line: 42,
};
assert!(snapshot.location_stats.contains_key(&test_location));
let loc_stats = &snapshot.location_stats[&test_location];
assert_eq!(loc_stats.get_total_for_type(EventType::Captured), 1);
assert!(snapshot.module_stats.contains_key("test_module"));
let mod_stats = &snapshot.module_stats["test_module"];
assert_eq!(mod_stats.get_total_for_type(EventType::Captured), 1);
} else {
panic!("Failed to get stats snapshot");
}
} else {
panic!("Failed to request stats");
}
Ok(())
}
#[tokio::test]
async fn test_stats_with_silenced_events() -> Result<()> {
let stats_config = StatsConfig::default();
let mut silencing_matcher = MatcherSet::empty();
silencing_matcher.add_matcher(Match::info().module_pattern("test_*"));
silencing_matcher.add_matcher(Match::info().exclude().module_pattern("test_silenced*"));
let config =
Config::from_tab(("silencing_tab", silencing_matcher)).with_stats(stats_config);
let tracer = Kerf::new_with_config(config);
let silenced_events = Arc::new(Mutex::new(Vec::<String>::new()));
let silenced_clone = silenced_events.clone();
tracer
.set_silenced_callback(move |event, _silencers| {
let silenced = silenced_clone.clone();
let msg = event.message.clone();
tokio::spawn(async move {
let mut lock = silenced.lock().await;
lock.push(msg);
});
})?
.await??;
let silenced_event = create_test_event_for_stats(
1,
Level::INFO,
"Silenced event",
Some("test_silenced"),
Some("silenced.rs"),
Some(42),
None,
);
send_and_process_event(&tracer, silenced_event).await;
{
let silenced = silenced_events.lock().await;
assert_eq!(silenced.len(), 1);
assert!(silenced.contains(&"Silenced event".to_string()));
}
if let Ok(stats_rx) = tracer.get_stats() {
if let Ok(Some(snapshot)) = stats_rx.await {
let mut total_silenced = 0;
for ((_level, event_type), count) in &snapshot.level_event_counts {
if *event_type == EventType::Silenced {
total_silenced += count;
}
}
assert_eq!(total_silenced, 1);
let silenced_location = Location {
file: "silenced.rs".to_string(),
line: 42,
};
assert!(snapshot.location_stats.contains_key(&silenced_location));
let loc_stats = &snapshot.location_stats[&silenced_location];
assert_eq!(loc_stats.get_total_for_type(EventType::Silenced), 1);
}
}
Ok(())
}
#[tokio::test]
async fn test_stats_clear_functionality() -> Result<()> {
let stats_config = StatsConfig::default();
let matcher_set = MatcherSet::from_matcher(Match::debug().all_modules());
let config = Config::from_tab(("test_tab", matcher_set)).with_stats(stats_config);
let tracer = Kerf::new_with_config(config);
for i in 0..5 {
let event = create_test_event_for_stats(
i,
Level::INFO,
&format!("Event {i}"),
Some(&format!("module_{i}")),
Some(&format!("file_{i}.rs")),
Some(i as u32),
None,
);
send_and_process_event(&tracer, event).await;
}
if let Ok(stats_rx) = tracer.get_stats() {
if let Ok(Some(snapshot)) = stats_rx.await {
assert!(snapshot.total_entries > 0);
assert!(!snapshot.location_stats.is_empty());
assert!(!snapshot.module_stats.is_empty());
let mut total_events = 0;
for ((_, _), count) in &snapshot.level_event_counts {
total_events += count;
}
assert_eq!(total_events, 5);
}
}
tracer.clear_stats()?.await??;
if let Ok(stats_rx) = tracer.get_stats() {
if let Ok(Some(snapshot)) = stats_rx.await {
assert_eq!(snapshot.total_entries, 0);
assert!(snapshot.location_stats.is_empty());
assert!(snapshot.module_stats.is_empty());
let mut total_events = 0;
for ((_, _), count) in &snapshot.level_event_counts {
total_events += count;
}
assert_eq!(total_events, 0);
}
}
Ok(())
}
#[tokio::test]
async fn test_stats_multiple_levels_and_types() -> Result<()> {
let stats_config = StatsConfig::default();
let matcher_set = MatcherSet::from_matcher(Match::trace().all_modules());
let config = Config::from_tab(("test_tab", matcher_set)).with_stats(stats_config);
let tracer = Kerf::new_with_config(config);
let levels = [
Level::TRACE,
Level::DEBUG,
Level::INFO,
Level::WARN,
Level::ERROR,
];
for (i, level) in levels.iter().enumerate() {
let event = create_test_event_for_stats(
i as u64,
*level,
&format!("{level:?} level event"),
Some("test_module"),
Some("test.rs"),
Some((i + 1) as u32),
None,
);
send_and_process_event(&tracer, event).await;
}
if let Ok(stats_rx) = tracer.get_stats() {
if let Ok(Some(snapshot)) = stats_rx.await {
for level in &levels {
let count = snapshot
.level_event_counts
.get(&(*level, EventType::Captured))
.unwrap_or(&0);
assert_eq!(*count, 1, "Should have 1 event for level {level:?}");
}
assert_eq!(snapshot.location_stats.len(), 5);
assert_eq!(snapshot.module_stats.len(), 1); let mod_stats = &snapshot.module_stats["test_module"];
assert_eq!(mod_stats.get_total_for_type(EventType::Captured), 5);
}
}
Ok(())
}
#[tokio::test]
async fn test_stats_high_volume() -> Result<()> {
let stats_config = StatsConfig {
track_by_location: true,
track_by_module: true,
track_by_level: true,
max_locations: 200,
max_modules: 50,
};
let matcher_set = MatcherSet::from_matcher(Match::debug().all_modules());
let config = Config::from_tab(("test_tab", matcher_set)).with_stats(stats_config);
let tracer = Kerf::new_with_config(config);
const EVENT_COUNT: u64 = 100;
const MODULE_COUNT: u64 = 10;
const LINES_PER_MODULE: u64 = 20;
for i in 0..EVENT_COUNT {
let module_idx = i % MODULE_COUNT;
let line_idx = i % LINES_PER_MODULE;
let event = create_test_event_for_stats(
i,
Level::INFO,
&format!("High volume event {i}"),
Some(&format!("module_{module_idx}")),
Some(&format!("file_{module_idx}.rs")),
Some(line_idx as u32),
None,
);
send_and_process_event(&tracer, event).await;
}
if let Ok(stats_rx) = tracer.get_stats() {
if let Ok(Some(snapshot)) = stats_rx.await {
let mut total_captured = 0;
for ((_level, event_type), count) in &snapshot.level_event_counts {
if *event_type == EventType::Captured {
total_captured += count;
}
}
assert_eq!(total_captured, EVENT_COUNT);
assert_eq!(snapshot.module_stats.len(), MODULE_COUNT as usize);
for i in 0..MODULE_COUNT {
let module_name = format!("module_{i}");
assert!(snapshot.module_stats.contains_key(&module_name));
let mod_stats = &snapshot.module_stats[&module_name];
assert_eq!(
mod_stats.get_total_for_type(EventType::Captured),
EVENT_COUNT / MODULE_COUNT
);
}
assert_eq!(snapshot.location_stats.len(), LINES_PER_MODULE as usize);
}
}
Ok(())
}
#[tokio::test]
async fn test_stats_with_different_event_types_comprehensive() -> Result<()> {
let stats_config = StatsConfig::default();
let mut capture_matcher = MatcherSet::empty();
capture_matcher.add_matcher(Match::info().module_pattern("capture_*"));
let mut silence_matcher = MatcherSet::empty();
silence_matcher.add_matcher(Match::info().module_pattern("silence_*"));
silence_matcher.add_matcher(Match::info().exclude().module_pattern("silence_this*"));
let config = Config::empty()
.with_tab("capture_tab", capture_matcher)
.with_tab("silence_tab", silence_matcher)
.with_stats(stats_config);
let tracer = Kerf::new_with_config(config);
let captured_events = Arc::new(Mutex::new(Vec::<String>::new()));
let silenced_events = Arc::new(Mutex::new(Vec::<String>::new()));
let dropped_events = Arc::new(Mutex::new(Vec::<String>::new()));
let captured_clone = captured_events.clone();
let silenced_clone = silenced_events.clone();
let dropped_clone = dropped_events.clone();
tracer
.set_callback(move |event, _tabs| {
let captured = captured_clone.clone();
let msg = event.message.clone();
tokio::spawn(async move {
let mut lock = captured.lock().await;
lock.push(msg);
});
})?
.await??;
tracer
.set_silenced_callback(move |event, _silencers| {
let silenced = silenced_clone.clone();
let msg = event.message.clone();
tokio::spawn(async move {
let mut lock = silenced.lock().await;
lock.push(msg);
});
})?
.await??;
tracer
.set_dropped_callback(move |event| {
let dropped = dropped_clone.clone();
let msg = event.message.clone();
tokio::spawn(async move {
let mut lock = dropped.lock().await;
lock.push(msg);
});
})?
.await??;
let captured_event1 = create_test_event_for_stats(
1,
Level::INFO,
"Captured event 1",
Some("capture_module"),
Some("capture.rs"),
Some(10),
None,
);
let captured_event2 = create_test_event_for_stats(
2,
Level::WARN,
"Captured event 2",
Some("capture_other"),
Some("capture.rs"),
Some(20),
None,
);
let silenced_event = create_test_event_for_stats(
3,
Level::INFO,
"Silenced event",
Some("silence_this"),
Some("silence.rs"),
Some(30),
None,
);
let dropped_event = create_test_event_for_stats(
4,
Level::INFO,
"Dropped event",
Some("other_module"),
Some("other.rs"),
Some(40),
None,
);
send_and_process_event(&tracer, captured_event1).await;
send_and_process_event(&tracer, captured_event2).await;
send_and_process_event(&tracer, silenced_event).await;
send_and_process_event(&tracer, dropped_event).await;
tokio::time::sleep(Duration::from_millis(200)).await;
{
let captured = captured_events.lock().await;
let silenced = silenced_events.lock().await;
let dropped = dropped_events.lock().await;
assert_eq!(captured.len(), 2);
assert!(captured.contains(&"Captured event 1".to_string()));
assert!(captured.contains(&"Captured event 2".to_string()));
assert_eq!(silenced.len(), 1);
assert!(silenced.contains(&"Silenced event".to_string()));
assert_eq!(dropped.len(), 1);
assert!(dropped.contains(&"Dropped event".to_string()));
}
if let Ok(stats_rx) = tracer.get_stats() {
if let Ok(Some(snapshot)) = stats_rx.await {
let mut total_captured = 0;
let mut total_silenced = 0;
let mut total_dropped = 0;
for ((_level, event_type), count) in &snapshot.level_event_counts {
match event_type {
EventType::Captured => total_captured += count,
EventType::Silenced => total_silenced += count,
EventType::Dropped => total_dropped += count,
}
}
assert_eq!(total_captured, 2);
assert_eq!(total_silenced, 1);
assert_eq!(total_dropped, 1);
assert_eq!(
snapshot.level_event_counts[&(Level::INFO, EventType::Captured)],
1
); assert_eq!(
snapshot.level_event_counts[&(Level::WARN, EventType::Captured)],
1
); assert_eq!(
snapshot.level_event_counts[&(Level::INFO, EventType::Silenced)],
1
); assert_eq!(
snapshot.level_event_counts[&(Level::INFO, EventType::Dropped)],
1
);
assert_eq!(snapshot.location_stats.len(), 4);
let capture_location = Location {
file: "capture.rs".to_string(),
line: 10,
};
let silence_location = Location {
file: "silence.rs".to_string(),
line: 30,
};
let dropped_location = Location {
file: "other.rs".to_string(),
line: 40,
};
assert!(snapshot.location_stats.contains_key(&capture_location));
assert!(snapshot.location_stats.contains_key(&silence_location));
assert!(snapshot.location_stats.contains_key(&dropped_location));
assert!(snapshot.module_stats.contains_key("capture_module"));
assert!(snapshot.module_stats.contains_key("silence_this"));
assert!(snapshot.module_stats.contains_key("other_module"));
let capture_mod_stats = &snapshot.module_stats["capture_module"];
assert_eq!(capture_mod_stats.get_total_for_type(EventType::Captured), 1);
let silence_mod_stats = &snapshot.module_stats["silence_this"];
assert_eq!(silence_mod_stats.get_total_for_type(EventType::Silenced), 1);
let dropped_mod_stats = &snapshot.module_stats["other_module"];
assert_eq!(dropped_mod_stats.get_total_for_type(EventType::Dropped), 1);
}
}
Ok(())
}
#[tokio::test]
async fn test_stats_concurrency_stress() -> Result<()> {
let stats_config = StatsConfig::default();
let matcher_set = MatcherSet::from_matcher(Match::debug().all_modules());
let config = Config::from_tab(("stress_tab", matcher_set)).with_stats(stats_config);
let tracer = Kerf::new_with_config(config);
const CONCURRENT_TASKS: usize = 10;
const EVENTS_PER_TASK: usize = 50;
let mut handles = Vec::new();
for task_id in 0..CONCURRENT_TASKS {
let tracer_clone = tracer.clone();
let handle = tokio::spawn(async move {
for event_id in 0..EVENTS_PER_TASK {
let event = create_test_event_for_stats(
(task_id * EVENTS_PER_TASK + event_id) as u64,
Level::INFO,
&format!("Concurrent event {event_id} from task {task_id}"),
Some(&format!("task_module_{task_id}")),
Some(&format!("task_{task_id}.rs")),
Some(event_id as u32),
None,
);
let _ = tracer_clone._get_sender_for_testing().send(event);
tokio::time::sleep(Duration::from_millis(1)).await;
}
});
handles.push(handle);
}
for handle in handles {
handle.await?;
}
tokio::time::sleep(Duration::from_millis(500)).await;
if let Ok(stats_rx) = tracer.get_stats() {
if let Ok(Some(snapshot)) = stats_rx.await {
let mut total_captured = 0;
for ((_level, event_type), count) in &snapshot.level_event_counts {
if *event_type == EventType::Captured {
total_captured += count;
}
}
assert_eq!(total_captured, (CONCURRENT_TASKS * EVENTS_PER_TASK) as u64);
assert_eq!(snapshot.module_stats.len(), CONCURRENT_TASKS);
for task_id in 0..CONCURRENT_TASKS {
let module_name = format!("task_module_{task_id}");
assert!(snapshot.module_stats.contains_key(&module_name));
let mod_stats = &snapshot.module_stats[&module_name];
assert_eq!(
mod_stats.get_total_for_type(EventType::Captured),
EVENTS_PER_TASK as u64
);
}
assert_eq!(
snapshot.location_stats.len(),
CONCURRENT_TASKS * EVENTS_PER_TASK
);
}
}
Ok(())
}
}