use std::sync::Arc;
use std::thread;
use tiny_counter::EventStore;
#[test]
fn concurrent_record_from_multiple_threads() {
let store = Arc::new(EventStore::new());
let handles: Vec<_> = (0..10)
.map(|i| {
let store = store.clone();
thread::spawn(move || {
for _ in 0..100 {
store.record(format!("event_{}", i));
}
})
})
.collect();
for h in handles {
h.join().unwrap();
}
for i in 0..10 {
let sum = store
.query(format!("event_{}", i))
.last_hours(1)
.sum()
.unwrap();
assert_eq!(sum, 100);
}
}
#[test]
fn concurrent_query_while_recording() {
let store = Arc::new(EventStore::new());
let writer = {
let store = store.clone();
thread::spawn(move || {
for _ in 0..1000 {
store.record("event");
}
})
};
let readers: Vec<_> = (0..5)
.map(|_| {
let store = store.clone();
thread::spawn(move || {
for _ in 0..100 {
let _ = store.query("event").last_hours(1).sum();
}
})
})
.collect();
writer.join().unwrap();
for r in readers {
r.join().unwrap();
}
assert_eq!(store.query("event").last_hours(1).sum(), Some(1000));
}
#[test]
fn stress_high_contention_single_event() {
let store = Arc::new(EventStore::new());
let num_threads = 100;
let events_per_thread = 1000;
let handles: Vec<_> = (0..num_threads)
.map(|_| {
let store = store.clone();
thread::spawn(move || {
for _ in 0..events_per_thread {
store.record("hotspot");
}
})
})
.collect();
for h in handles {
h.join().unwrap();
}
let total = store.query("hotspot").last_hours(1).sum().unwrap();
assert_eq!(total, num_threads * events_per_thread);
}
#[test]
fn stress_concurrent_record_and_query() {
let store = Arc::new(EventStore::new());
let num_writers = 8;
let num_readers = 8;
let writes_per_thread = 500;
let reads_per_thread = 500;
let mut handles = vec![];
for i in 0..num_writers {
let store = store.clone();
handles.push(thread::spawn(move || {
for _ in 0..writes_per_thread {
store.record(format!("event_{}", i % 4));
}
}));
}
for i in 0..num_readers {
let store = store.clone();
handles.push(thread::spawn(move || {
for _ in 0..reads_per_thread {
let _ = store.query(format!("event_{}", i % 4)).last_hours(1).sum();
}
}));
}
for h in handles {
h.join().unwrap();
}
let total_writes = num_writers * writes_per_thread;
let expected_per_event = total_writes / 4;
for i in 0..4 {
let sum = store
.query(format!("event_{}", i))
.last_hours(1)
.sum()
.unwrap();
assert_eq!(sum, expected_per_event);
}
}
#[test]
fn stress_concurrent_query_many() {
let store = Arc::new(EventStore::new());
for i in 0..10 {
store.record_count(format!("event_{}", i), 100);
}
let readers: Vec<_> = (0..10)
.map(|_| {
let store = store.clone();
thread::spawn(move || {
for _ in 0..1000 {
let event_ids: Vec<String> = (0..5).map(|i| format!("event_{}", i)).collect();
let event_refs: Vec<&str> = event_ids.iter().map(|s| s.as_str()).collect();
let _ = store.query_many(&event_refs).last_hours(1).sum();
}
})
})
.collect();
for r in readers {
r.join().unwrap();
}
let event_ids: Vec<String> = (0..5).map(|i| format!("event_{}", i)).collect();
let event_refs: Vec<&str> = event_ids.iter().map(|s| s.as_str()).collect();
let sum = store.query_many(&event_refs).last_hours(1).sum();
assert_eq!(sum, Some(500));
}
#[test]
fn stress_concurrent_query_ratio() {
let store = Arc::new(EventStore::new());
store.record_count("numerator", 1000);
store.record_count("denominator", 500);
let readers: Vec<_> = (0..10)
.map(|_| {
let store = store.clone();
thread::spawn(move || {
for _ in 0..1000 {
let ratio = store.query_ratio("numerator", "denominator").last_hours(1);
assert_eq!(ratio, Some(2.0));
}
})
})
.collect();
for r in readers {
r.join().unwrap();
}
}
#[test]
fn stress_concurrent_query_delta() {
let store = Arc::new(EventStore::new());
store.record_count("positive", 1000);
store.record_count("negative", 300);
let readers: Vec<_> = (0..10)
.map(|_| {
let store = store.clone();
thread::spawn(move || {
for _ in 0..1000 {
let delta = store
.query_delta("positive", "negative")
.last_hours(1)
.sum();
assert_eq!(delta, 700);
}
})
})
.collect();
for r in readers {
r.join().unwrap();
}
}
#[test]
fn stress_mixed_workload() {
let store = Arc::new(EventStore::new());
for i in 0..5 {
store.record_count(format!("event_{}", i), 50);
}
let mut handles = vec![];
for i in 0..4 {
let store = store.clone();
handles.push(thread::spawn(move || {
for _ in 0..500 {
store.record(format!("event_{}", i));
}
}));
}
for i in 0..4 {
let store = store.clone();
handles.push(thread::spawn(move || {
for _ in 0..250 {
let _ = store.query(format!("event_{}", i)).last_hours(1).sum();
}
}));
}
for _ in 0..2 {
let store = store.clone();
handles.push(thread::spawn(move || {
for _ in 0..250 {
let event_ids: Vec<String> = (0..3).map(|i| format!("event_{}", i)).collect();
let event_refs: Vec<&str> = event_ids.iter().map(|s| s.as_str()).collect();
let _ = store.query_many(&event_refs).last_hours(1).sum();
}
}));
}
for _ in 0..2 {
let store = store.clone();
handles.push(thread::spawn(move || {
for _ in 0..250 {
let _ = store.query_ratio("event_0", "event_1").last_hours(1);
}
}));
}
for h in handles {
h.join().unwrap();
}
for i in 0..4 {
let sum = store
.query(format!("event_{}", i))
.last_hours(1)
.sum()
.unwrap();
assert_eq!(sum, 550); }
}
#[test]
#[cfg(feature = "serde")]
fn stress_reservation_commit_cancel() {
use tiny_counter::TimeUnit;
let store = Arc::new(EventStore::new());
let num_threads = 50;
let ops_per_thread = 100;
let handles: Vec<_> = (0..num_threads)
.map(|i| {
let store = store.clone();
thread::spawn(move || {
for j in 0..ops_per_thread {
let result = store
.limit()
.at_most("limited_event", 10000, TimeUnit::Hours)
.reserve("limited_event");
if let Ok(reservation) = result {
if (i + j) % 10 < 7 {
reservation.commit();
} else {
reservation.cancel();
}
}
}
})
})
.collect();
for h in handles {
h.join().unwrap();
}
let total = store.query("limited_event").last_hours(1).sum().unwrap();
assert!(total > 0);
let result = store
.limit()
.at_most("limited_event", 10000, TimeUnit::Hours)
.reserve("limited_event");
assert!(result.is_ok());
}
#[test]
#[cfg(feature = "serde")]
fn stress_reservation_respects_limits() {
use tiny_counter::TimeUnit;
let store = Arc::new(EventStore::new());
let limit = 1000u32;
let num_threads = 100;
let attempts_per_thread = 50;
let handles: Vec<_> = (0..num_threads)
.map(|_| {
let store = store.clone();
thread::spawn(move || {
for _ in 0..attempts_per_thread {
let result = store
.limit()
.at_most("limited", limit, TimeUnit::Hours)
.reserve("limited");
if let Ok(reservation) = result {
reservation.commit();
}
}
})
})
.collect();
for h in handles {
h.join().unwrap();
}
let total = store.query("limited").last_hours(1).sum().unwrap();
assert!(
total <= limit,
"total {} should not exceed limit {}",
total,
limit
);
}
#[test]
#[cfg(feature = "serde")]
fn stress_persistence_under_load() {
use tiny_counter::storage::MemoryStorage;
let store = EventStore::builder()
.with_storage(MemoryStorage::new())
.build()
.unwrap();
let store = Arc::new(store);
let num_writers = 10;
let num_persisters = 5;
let writes_per_thread = 1000;
let persists_per_thread = 20;
let writes_completed = Arc::new(std::sync::atomic::AtomicU64::new(0));
let mut handles = vec![];
for i in 0..num_writers {
let store = store.clone();
let writes_completed = writes_completed.clone();
handles.push(thread::spawn(move || {
for _ in 0..writes_per_thread {
store.record(format!("event_{}", i % 3));
writes_completed.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
}
}));
}
for _ in 0..num_persisters {
let store = store.clone();
handles.push(thread::spawn(move || {
for _ in 0..persists_per_thread {
let _ = store.persist();
thread::sleep(std::time::Duration::from_millis(5));
}
}));
}
for h in handles {
h.join().unwrap();
}
let actual_writes = writes_completed.load(std::sync::atomic::Ordering::Relaxed);
assert_eq!(
actual_writes,
num_writers * writes_per_thread,
"Not all write operations completed - test infrastructure bug"
);
store.persist().unwrap();
let mut event_counts = Vec::new();
for i in 0..3 {
let sum = store
.query(format!("event_{}", i))
.last_hours(1)
.sum()
.unwrap_or(0);
event_counts.push((i, sum));
}
let total: u32 = event_counts.iter().map(|(_, count)| count).sum();
let expected = (num_writers * writes_per_thread) as u32;
assert_eq!(
total,
expected,
"Event count mismatch!\n\
Expected: {} total events\n\
Got: {} total events\n\
Lost: {} events ({:.1}%)\n\
Breakdown: {:?}\n\
This suggests a race condition in persist/record logic.",
expected,
total,
expected - total,
(expected - total) as f64 / expected as f64 * 100.0,
event_counts
);
}
#[test]
#[cfg(feature = "serde")]
fn test_get_counter_race_with_concurrent_record() {
use std::sync::Barrier;
use tiny_counter::storage::MemoryStorage;
let store = EventStore::builder()
.with_storage(MemoryStorage::new())
.build()
.unwrap();
let store = Arc::new(store);
let barrier = Arc::new(Barrier::new(20));
let writes_per_thread = 500;
let handles: Vec<_> = (0..20)
.map(|_| {
let store = store.clone();
let barrier = barrier.clone();
thread::spawn(move || {
barrier.wait(); for _ in 0..writes_per_thread {
store.record("target_event");
}
})
})
.collect();
for h in handles {
h.join().unwrap();
}
let expected = 20 * writes_per_thread;
let count = store.query("target_event").last_hours(1).sum().unwrap();
assert_eq!(
count,
expected,
"Lost {} events ({:.1}%) due to TOCTOU race in get_counter_for_record",
expected - count,
(expected - count) as f64 / expected as f64 * 100.0
);
}