use chrono::Utc;
use std::sync::Arc;
use sz_orm_scheduler::{
CounterJobHandler, CronScheduler, JobHandler, RecordingJobHandler, ScheduledTask, Scheduler,
};
#[test]
fn stress_scheduler_10k_tasks() {
let scheduler = CronScheduler::new();
let n: usize = 10_000;
for i in 0..n {
let task = ScheduledTask::new(format!("task-{}", i), format!("name-{}", i), "* * * * *");
scheduler.schedule(task).unwrap();
}
assert_eq!(scheduler.list_tasks().len(), n);
for i in 0..n {
scheduler.cancel(&format!("task-{}", i)).unwrap();
}
assert_eq!(scheduler.list_tasks().len(), 0);
}
#[test]
fn stress_scheduler_1000_tasks_fire_simultaneously() {
let scheduler = CronScheduler::new();
let n: usize = 1000;
let counter = Arc::new(CounterJobHandler::new());
let counter_clone = counter.clone();
for i in 0..n {
let task = ScheduledTask::new(format!("task-{}", i), format!("name-{}", i), "* * * * *");
scheduler.schedule(task).unwrap();
let c = counter_clone.clone();
let wrapper = Arc::new(SharedCounterHandler::new(c));
scheduler.register_handler(format!("task-{}", i), wrapper);
}
let now = Utc::now();
let fired = scheduler.try_fire_due(now);
assert_eq!(fired, n, "all {} tasks should fire", n);
assert_eq!(counter.count(), n as u64);
}
struct SharedCounterHandler {
counter: Arc<CounterJobHandler>,
}
impl SharedCounterHandler {
fn new(counter: Arc<CounterJobHandler>) -> Self {
Self { counter }
}
}
impl JobHandler for SharedCounterHandler {
fn handle(&self, _task: &ScheduledTask) -> Result<(), String> {
self.counter.handle(_task)
}
}
#[test]
fn stress_scheduler_parse_cron_various() {
let scheduler = CronScheduler::new();
let valid_exprs = vec![
"* * * * *",
"0 * * * *",
"0 0 * * *",
"0 0 1 * *",
"0 0 1 1 *",
"*/5 * * * *",
"0 9-17 * * *",
"0,15,30,45 * * * *",
"0 0 1,15 * *",
];
let invalid_field_count_exprs = vec!["", "* * * *", "* * * * * *", "abc def * * *"];
let from = Utc::now();
for _ in 0..3 {
for expr in &valid_exprs {
let result = scheduler.next_run_time(expr, from);
assert!(result.is_ok(), "expr should be valid: {}", expr);
}
}
for _ in 0..10 {
for expr in &invalid_field_count_exprs {
let result = scheduler.next_run_time(expr, from);
assert!(
result.is_err(),
"expr should be invalid (field count): {}",
expr
);
}
}
}
#[test]
fn stress_scheduler_next_run_time_repeated() {
let scheduler = CronScheduler::new();
let exprs = vec![
"* * * * *",
"0 * * * *",
"0 0 * * *",
"*/5 * * * *",
"0 9-17 * * *",
];
let from = Utc::now();
for _ in 0..5 {
for expr in &exprs {
let result = scheduler.next_run_time(expr, from);
assert!(result.is_ok(), "next_run_time failed for expr: {}", expr);
}
}
}
#[test]
fn stress_scheduler_pause_resume() {
let scheduler = CronScheduler::new();
let n: usize = 1000;
for i in 0..n {
let task = ScheduledTask::new(format!("task-{}", i), format!("name-{}", i), "* * * * *");
scheduler.schedule(task).unwrap();
}
for i in 0..n / 2 {
scheduler.pause(&format!("task-{}", i)).unwrap();
}
let now = Utc::now();
let fired = scheduler.try_fire_due(now);
assert_eq!(fired, n / 2, "only non-paused tasks should fire");
for i in 0..n / 2 {
scheduler.resume(&format!("task-{}", i)).unwrap();
}
let fired = scheduler.try_fire_due(now);
assert_eq!(fired, n, "all tasks should fire after resume");
}
#[test]
fn stress_scheduler_cancel_nonexistent() {
let scheduler = CronScheduler::new();
for i in 0..100 {
let result = scheduler.cancel(&format!("nonexistent-{}", i));
assert!(result.is_err());
}
}
#[test]
fn stress_scheduler_schedule_duplicate_overwrites() {
let scheduler = CronScheduler::new();
for _ in 0..1000 {
let task = ScheduledTask::new("dup-id", "name", "* * * * *");
scheduler.schedule(task).unwrap();
}
assert_eq!(scheduler.list_tasks().len(), 1);
}
#[test]
fn stress_scheduler_recording_handler_order() {
let scheduler = CronScheduler::new();
let recording = Arc::new(RecordingJobHandler::new());
let task = ScheduledTask::new("task-1", "name", "* * * * *");
scheduler.schedule(task).unwrap();
scheduler.register_handler("task-1", recording.clone());
for _ in 0..1000 {
scheduler.try_fire_due(Utc::now());
}
let handled = recording.handled_ids();
assert_eq!(handled.len(), 1000);
for id in &handled {
assert_eq!(id, "task-1");
}
assert_eq!(recording.handled_ids().len(), 1000);
}
#[test]
fn stress_scheduler_disabled_task_does_not_fire() {
let scheduler = CronScheduler::new();
let counter = Arc::new(CounterJobHandler::new());
let counter_clone = counter.clone();
let task = ScheduledTask::new("task-1", "name", "* * * * *").disable();
scheduler.schedule(task).unwrap();
scheduler.register_handler("task-1", Arc::new(SharedCounterHandler::new(counter_clone)));
for _ in 0..100 {
scheduler.try_fire_due(Utc::now());
}
assert_eq!(counter.count(), 0, "disabled task should not fire");
}
#[test]
fn stress_scheduler_start_stop_cycle() {
let scheduler = Arc::new(CronScheduler::new());
for cycle in 0..10 {
let s = scheduler.clone();
let handle = std::thread::spawn(move || {
s.start(10).unwrap();
std::thread::sleep(std::time::Duration::from_millis(50));
s.stop().unwrap();
});
handle.join().unwrap();
let _ = cycle;
}
assert!(!scheduler.is_running());
}