2024-01-03 19:57:28 +00:00
|
|
|
use crate::tests::prelude::*;
|
2024-06-24 00:06:20 +00:00
|
|
|
use crate::topic_monitor::{GenerationsList, Topic, TopicMonitor};
|
2023-10-08 21:22:27 +00:00
|
|
|
use std::sync::{
|
|
|
|
atomic::{AtomicU32, AtomicU64, Ordering},
|
|
|
|
Arc,
|
|
|
|
};
|
|
|
|
|
2024-01-03 19:57:28 +00:00
|
|
|
#[test]
|
|
|
|
#[serial]
|
|
|
|
fn test_topic_monitor() {
|
2024-03-24 10:55:03 +00:00
|
|
|
let _cleanup = test_init();
|
2024-06-24 00:06:20 +00:00
|
|
|
let monitor = TopicMonitor::default();
|
2023-10-08 21:22:27 +00:00
|
|
|
let gens = GenerationsList::new();
|
2024-06-24 00:06:20 +00:00
|
|
|
let t = Topic::sigchld;
|
2023-10-08 21:22:27 +00:00
|
|
|
gens.sigchld.set(0);
|
|
|
|
assert_eq!(monitor.generation_for_topic(t), 0);
|
|
|
|
let changed = monitor.check(&gens, false /* wait */);
|
|
|
|
assert!(!changed);
|
|
|
|
assert_eq!(gens.sigchld.get(), 0);
|
|
|
|
|
|
|
|
monitor.post(t);
|
|
|
|
let changed = monitor.check(&gens, true /* wait */);
|
|
|
|
assert!(changed);
|
|
|
|
assert_eq!(gens.get(t), 1);
|
|
|
|
assert_eq!(monitor.generation_for_topic(t), 1);
|
|
|
|
|
|
|
|
monitor.post(t);
|
|
|
|
assert_eq!(monitor.generation_for_topic(t), 2);
|
|
|
|
let changed = monitor.check(&gens, true /* wait */);
|
|
|
|
assert!(changed);
|
|
|
|
assert_eq!(gens.sigchld.get(), 2);
|
2024-01-03 19:57:28 +00:00
|
|
|
}
|
2023-10-08 21:22:27 +00:00
|
|
|
|
2024-01-03 19:57:28 +00:00
|
|
|
#[test]
|
|
|
|
#[serial]
|
|
|
|
fn test_topic_monitor_torture() {
|
2024-03-24 10:55:03 +00:00
|
|
|
let _cleanup = test_init();
|
2024-06-24 00:06:20 +00:00
|
|
|
let monitor = Arc::new(TopicMonitor::default());
|
2023-10-08 21:22:27 +00:00
|
|
|
const THREAD_COUNT: usize = 64;
|
2024-06-24 00:06:20 +00:00
|
|
|
let t1 = Topic::sigchld;
|
|
|
|
let t2 = Topic::sighupint;
|
2023-10-08 21:22:27 +00:00
|
|
|
let mut gens_list = vec![GenerationsList::invalid(); THREAD_COUNT];
|
|
|
|
let post_count = Arc::new(AtomicU64::new(0));
|
|
|
|
for gen in &mut gens_list {
|
|
|
|
*gen = monitor.current_generations();
|
|
|
|
post_count.fetch_add(1, Ordering::Relaxed);
|
|
|
|
monitor.post(t1);
|
|
|
|
}
|
|
|
|
|
|
|
|
let completed = Arc::new(AtomicU32::new(0));
|
|
|
|
let mut threads = vec![];
|
|
|
|
|
|
|
|
for gens in gens_list {
|
|
|
|
let monitor = Arc::downgrade(&monitor);
|
|
|
|
let post_count = Arc::downgrade(&post_count);
|
|
|
|
let completed = Arc::downgrade(&completed);
|
|
|
|
threads.push(std::thread::spawn(move || {
|
|
|
|
for _ in 0..1 << 11 {
|
|
|
|
let before = gens.clone();
|
|
|
|
let _changed = monitor.upgrade().unwrap().check(&gens, true /* wait */);
|
|
|
|
assert!(before.get(t1) < gens.get(t1));
|
|
|
|
assert!(gens.get(t1) <= post_count.upgrade().unwrap().load(Ordering::Relaxed));
|
|
|
|
assert_eq!(gens.get(t2), 0);
|
|
|
|
}
|
|
|
|
let _amt = completed.upgrade().unwrap().fetch_add(1, Ordering::Relaxed);
|
|
|
|
}));
|
|
|
|
}
|
|
|
|
|
|
|
|
while completed.load(Ordering::Relaxed) < THREAD_COUNT.try_into().unwrap() {
|
|
|
|
post_count.fetch_add(1, Ordering::Relaxed);
|
|
|
|
monitor.post(t1);
|
|
|
|
std::thread::yield_now();
|
|
|
|
}
|
|
|
|
for t in threads {
|
|
|
|
t.join().unwrap();
|
|
|
|
}
|
2024-01-03 19:57:28 +00:00
|
|
|
}
|