use busybeaver::{
listener, listener_with_error, work, Beaver, BeaverResult, RuntimeError, TimeIntervalBuilder,
WorkResult,
};
use std::sync::atomic::{AtomicBool, AtomicU32, Ordering};
use std::sync::Arc;
use std::time::{Duration, Instant};
#[tokio::test]
async fn test_basic_time_interval_task() -> BeaverResult<()> {
let beaver = Beaver::new("test", 256);
let counter = Arc::new(AtomicU32::new(0));
let counter_clone = Arc::clone(&counter);
let task = TimeIntervalBuilder::new(work(move || {
let c = Arc::clone(&counter_clone);
async move {
c.fetch_add(1, Ordering::SeqCst);
WorkResult::NeedRetry
}
}))
.intervals_millis([0, 0, 0]) .build()?;
beaver.enqueue(task).await?;
tokio::time::sleep(Duration::from_millis(200)).await;
assert_eq!(
counter.load(Ordering::SeqCst),
3,
"Should execute exactly 3 times (intervals length)"
);
Ok(())
}
#[tokio::test]
async fn test_time_interval_respects_delays() -> BeaverResult<()> {
let beaver = Beaver::new("test", 256);
let timestamps = Arc::new(std::sync::Mutex::new(Vec::new()));
let ts_clone = Arc::clone(×tamps);
let task = TimeIntervalBuilder::new(work(move || {
let ts = Arc::clone(&ts_clone);
async move {
ts.lock().unwrap().push(Instant::now());
WorkResult::NeedRetry
}
}))
.intervals_millis([0, 1000, 1000]) .build()?;
let start = Instant::now();
beaver.enqueue(task).await?;
tokio::time::sleep(Duration::from_secs(3)).await;
let ts = timestamps.lock().unwrap();
assert_eq!(ts.len(), 3, "Should have 3 timestamps");
let first_delay = ts[0] - start;
assert!(
first_delay < Duration::from_millis(100),
"First execution should be immediate"
);
let second_delay = ts[1] - ts[0];
assert!(
second_delay >= Duration::from_millis(900) && second_delay < Duration::from_millis(1200),
"Second delay should be ~1 second"
);
let third_delay = ts[2] - ts[1];
assert!(
third_delay >= Duration::from_millis(900) && third_delay < Duration::from_millis(1200),
"Third delay should be ~1 second"
);
Ok(())
}
#[tokio::test]
async fn test_time_interval_stops_on_done() -> BeaverResult<()> {
let beaver = Beaver::new("test", 256);
let counter = Arc::new(AtomicU32::new(0));
let counter_clone = Arc::clone(&counter);
let task = TimeIntervalBuilder::new(work(move || {
let c = Arc::clone(&counter_clone);
async move {
let count = c.fetch_add(1, Ordering::SeqCst) + 1;
if count == 2 {
WorkResult::Done(())
} else {
WorkResult::NeedRetry
}
}
}))
.intervals_millis([0, 0, 0, 0, 0]) .build()?;
beaver.enqueue(task).await?;
tokio::time::sleep(Duration::from_millis(200)).await;
assert_eq!(
counter.load(Ordering::SeqCst),
2,
"Should stop after returning Done"
);
Ok(())
}
#[tokio::test]
async fn test_time_interval_default_intervals() -> BeaverResult<()> {
let beaver = Beaver::new("test", 256);
let counter = Arc::new(AtomicU32::new(0));
let counter_clone = Arc::clone(&counter);
let task = TimeIntervalBuilder::new(work(move || {
let c = Arc::clone(&counter_clone);
async move {
c.fetch_add(1, Ordering::SeqCst);
WorkResult::NeedRetry
}
}))
.build()?;
let start = Instant::now();
beaver.enqueue(task).await?;
tokio::time::sleep(Duration::from_millis(1500)).await;
assert_eq!(
counter.load(Ordering::SeqCst),
1,
"Default interval should execute once"
);
assert!(start.elapsed() >= Duration::from_secs(1));
Ok(())
}
#[tokio::test]
async fn test_time_interval_empty_intervals() -> BeaverResult<()> {
let beaver = Beaver::new("test", 256);
let counter = Arc::new(AtomicU32::new(0));
let counter_clone = Arc::clone(&counter);
let task = TimeIntervalBuilder::new(work(move || {
let c = Arc::clone(&counter_clone);
async move {
c.fetch_add(1, Ordering::SeqCst);
WorkResult::NeedRetry
}
}))
.intervals_millis(Vec::<u64>::new()) .build()?;
beaver.enqueue(task).await?;
tokio::time::sleep(Duration::from_millis(100)).await;
assert_eq!(
counter.load(Ordering::SeqCst),
1,
"Empty intervals should execute once immediately"
);
Ok(())
}
#[tokio::test]
async fn test_exponential_backoff_pattern() -> BeaverResult<()> {
let beaver = Beaver::new("test", 256);
let timestamps = Arc::new(std::sync::Mutex::new(Vec::new()));
let ts_clone = Arc::clone(×tamps);
let task = TimeIntervalBuilder::new(work(move || {
let ts = Arc::clone(&ts_clone);
async move {
ts.lock().unwrap().push(Instant::now());
WorkResult::NeedRetry
}
}))
.intervals_millis([0, 1000, 2000])
.build()?;
let start = Instant::now();
beaver.enqueue(task).await?;
tokio::time::sleep(Duration::from_secs(4)).await;
let ts = timestamps.lock().unwrap();
assert_eq!(ts.len(), 3);
let first_delay = ts[0] - start;
let second_delay = ts[1] - ts[0];
let third_delay = ts[2] - ts[1];
assert!(
first_delay < Duration::from_millis(200),
"First should be immediate"
);
assert!(
second_delay >= Duration::from_millis(900),
"Second should be ~1s"
);
assert!(
third_delay >= Duration::from_millis(1900),
"Third should be ~2s"
);
Ok(())
}
#[tokio::test]
async fn test_linear_backoff_pattern() -> BeaverResult<()> {
let beaver = Beaver::new("test", 256);
let counter = Arc::new(AtomicU32::new(0));
let counter_clone = Arc::clone(&counter);
let task = TimeIntervalBuilder::new(work(move || {
let c = Arc::clone(&counter_clone);
async move {
c.fetch_add(1, Ordering::SeqCst);
WorkResult::NeedRetry
}
}))
.intervals_millis([0, 1000, 1000, 1000]) .build()?;
beaver.enqueue(task).await?;
tokio::time::sleep(Duration::from_millis(3500)).await;
assert_eq!(counter.load(Ordering::SeqCst), 4);
Ok(())
}
#[tokio::test]
async fn test_time_interval_on_error_retries_exhausted() -> BeaverResult<()> {
let beaver = Beaver::new("test", 256);
let completed = Arc::new(AtomicBool::new(false));
let completed_clone = Arc::clone(&completed);
let exhausted = Arc::new(AtomicBool::new(false));
let exhausted_clone = Arc::clone(&exhausted);
let task = TimeIntervalBuilder::new(work(|| async { WorkResult::NeedRetry }))
.intervals_millis([0, 0])
.listener(listener_with_error(
move || {
completed_clone.store(true, Ordering::SeqCst);
},
|| {},
move |e: RuntimeError| {
if matches!(e, RuntimeError::RetriesExhausted) {
exhausted_clone.store(true, Ordering::SeqCst);
}
},
))
.build()?;
beaver.enqueue(task).await?;
tokio::time::sleep(Duration::from_millis(200)).await;
assert!(
exhausted.load(Ordering::SeqCst),
"exhaustion should report on_error(RetriesExhausted)"
);
assert!(
!completed.load(Ordering::SeqCst),
"exhaustion should NOT fire on_complete"
);
Ok(())
}
#[tokio::test]
async fn test_time_interval_on_complete_on_done() -> BeaverResult<()> {
let beaver = Beaver::new("test", 256);
let completed = Arc::new(AtomicBool::new(false));
let completed_clone = Arc::clone(&completed);
let task = TimeIntervalBuilder::new(work(|| async {
WorkResult::Done(()) }))
.intervals_millis([0, 0, 0])
.listener(listener(
move || {
completed_clone.store(true, Ordering::SeqCst);
},
|| {},
))
.build()?;
beaver.enqueue(task).await?;
tokio::time::sleep(Duration::from_millis(200)).await;
assert!(
completed.load(Ordering::SeqCst),
"on_complete should be called when Done is returned"
);
Ok(())
}
#[tokio::test]
async fn test_time_interval_on_interrupt() -> BeaverResult<()> {
let beaver = Beaver::new("test", 256);
let interrupted = Arc::new(AtomicBool::new(false));
let execution_count = Arc::new(AtomicU32::new(0));
let interrupted_clone = Arc::clone(&interrupted);
let ec_clone = Arc::clone(&execution_count);
let task = TimeIntervalBuilder::new(work(move || {
let ec = Arc::clone(&ec_clone);
async move {
ec.fetch_add(1, Ordering::SeqCst);
WorkResult::NeedRetry
}
}))
.intervals_millis([0, 10_000]) .listener(listener(
|| {},
move || {
interrupted_clone.store(true, Ordering::SeqCst);
},
))
.build()?;
beaver.enqueue(task).await?;
tokio::time::sleep(Duration::from_millis(100)).await;
assert_eq!(
execution_count.load(Ordering::SeqCst),
1,
"First execution should have completed"
);
beaver.cancel_all().await?;
Ok(())
}
#[tokio::test]
async fn test_time_interval_with_tag() -> BeaverResult<()> {
let task = TimeIntervalBuilder::new(work(|| async { WorkResult::Done(()) }))
.intervals_millis([0])
.tag("my-backoff-task")
.build()?;
assert_eq!(task.tag(), "my-backoff-task");
Ok(())
}
#[tokio::test]
async fn test_time_interval_without_tag() -> BeaverResult<()> {
let task = TimeIntervalBuilder::new(work(|| async { WorkResult::Done(()) }))
.intervals_millis([0])
.build()?;
assert_eq!(task.tag(), "");
Ok(())
}
#[tokio::test]
async fn test_time_interval_interrupt_during_sleep() -> BeaverResult<()> {
let beaver = Beaver::new("test", 256);
let counter = Arc::new(AtomicU32::new(0));
let counter_clone = Arc::clone(&counter);
let task = TimeIntervalBuilder::new(work(move || {
let c = Arc::clone(&counter_clone);
async move {
c.fetch_add(1, Ordering::SeqCst);
WorkResult::NeedRetry
}
}))
.intervals_millis([0, 10_000]) .build()?;
beaver.enqueue(task).await?;
tokio::time::sleep(Duration::from_millis(100)).await;
beaver.cancel_all().await?;
tokio::time::sleep(Duration::from_millis(500)).await;
assert_eq!(
counter.load(Ordering::SeqCst),
1,
"Should only execute once before interrupt"
);
Ok(())
}
#[tokio::test]
async fn test_time_interval_cancel_during_work() -> BeaverResult<()> {
let beaver = Beaver::new("test", 256);
let execution_count = Arc::new(AtomicU32::new(0));
let ec_clone = Arc::clone(&execution_count);
let task = TimeIntervalBuilder::new(work(move || {
let ec = Arc::clone(&ec_clone);
async move {
ec.fetch_add(1, Ordering::SeqCst);
tokio::time::sleep(Duration::from_millis(100)).await;
WorkResult::NeedRetry
}
}))
.intervals_millis([0, 0, 0, 0, 0]) .build()?;
beaver.enqueue(task).await?;
tokio::time::sleep(Duration::from_millis(250)).await;
beaver.cancel_all().await?;
tokio::time::sleep(Duration::from_millis(200)).await;
let count = execution_count.load(Ordering::SeqCst);
assert!(
count >= 1 && count < 5,
"Should execute some but not all, got {}",
count
);
Ok(())
}
#[tokio::test]
async fn test_time_interval_single_interval() -> BeaverResult<()> {
let beaver = Beaver::new("test", 256);
let counter = Arc::new(AtomicU32::new(0));
let counter_clone = Arc::clone(&counter);
let task = TimeIntervalBuilder::new(work(move || {
let c = Arc::clone(&counter_clone);
async move {
c.fetch_add(1, Ordering::SeqCst);
WorkResult::NeedRetry
}
}))
.intervals_millis([0]) .build()?;
beaver.enqueue(task).await?;
tokio::time::sleep(Duration::from_millis(100)).await;
assert_eq!(counter.load(Ordering::SeqCst), 1);
Ok(())
}
#[tokio::test]
async fn test_time_interval_all_zero_delays() -> BeaverResult<()> {
let beaver = Beaver::new("test", 256);
let counter = Arc::new(AtomicU32::new(0));
let counter_clone = Arc::clone(&counter);
let task = TimeIntervalBuilder::new(work(move || {
let c = Arc::clone(&counter_clone);
async move {
c.fetch_add(1, Ordering::SeqCst);
WorkResult::NeedRetry
}
}))
.intervals_millis([0, 0, 0, 0, 0])
.build()?;
let start = Instant::now();
beaver.enqueue(task).await?;
tokio::time::sleep(Duration::from_millis(100)).await;
assert_eq!(counter.load(Ordering::SeqCst), 5);
assert!(
start.elapsed() < Duration::from_millis(200),
"Zero delays should complete quickly"
);
Ok(())
}
#[tokio::test]
async fn test_time_interval_long_delay() -> BeaverResult<()> {
let beaver = Beaver::new("test", 256);
let counter = Arc::new(AtomicU32::new(0));
let counter_clone = Arc::clone(&counter);
let task = TimeIntervalBuilder::new(work(move || {
let c = Arc::clone(&counter_clone);
async move {
c.fetch_add(1, Ordering::SeqCst);
WorkResult::NeedRetry
}
}))
.intervals_millis([0, 100_000]) .build()?;
beaver.enqueue(task).await?;
tokio::time::sleep(Duration::from_millis(500)).await;
assert_eq!(
counter.load(Ordering::SeqCst),
1,
"Should only execute first interval within test time"
);
beaver.cancel_all().await?;
Ok(())
}
#[tokio::test]
async fn test_time_interval_builder_chain_order() -> BeaverResult<()> {
let _task1 = TimeIntervalBuilder::new(work(|| async { WorkResult::Done(()) }))
.intervals_millis([0, 1000])
.tag("task1")
.listener(listener(|| {}, || {}))
.build()?;
let _task2 = TimeIntervalBuilder::new(work(|| async { WorkResult::Done(()) }))
.listener(listener(|| {}, || {}))
.tag("task2")
.intervals_millis([0, 1000])
.build()?;
Ok(())
}
#[tokio::test]
async fn test_time_interval_intervals_type_flexibility() -> BeaverResult<()> {
let _task1 = TimeIntervalBuilder::new(work(|| async { WorkResult::Done(()) }))
.intervals_millis([1000, 2000, 3000])
.build()?;
let _task2 = TimeIntervalBuilder::new(work(|| async { WorkResult::Done(()) }))
.intervals_millis(vec![1000, 2000, 3000])
.build()?;
let intervals: Vec<u64> = vec![1000, 2000, 3000];
let _task3 = TimeIntervalBuilder::new(work(|| async { WorkResult::Done(()) }))
.intervals_millis(intervals)
.build()?;
Ok(())
}
#[tokio::test]
async fn test_time_interval_unique_task_id() -> BeaverResult<()> {
let task1 = TimeIntervalBuilder::new(work(|| async { WorkResult::Done(()) }))
.intervals_millis([0])
.build()?;
let task2 = TimeIntervalBuilder::new(work(|| async { WorkResult::Done(()) }))
.intervals_millis([0])
.build()?;
assert_ne!(task1.id(), task2.id());
Ok(())
}
#[tokio::test]
async fn test_http_retry_simulation() -> BeaverResult<()> {
let beaver = Beaver::new("test", 256);
let attempt = Arc::new(AtomicU32::new(0));
let success = Arc::new(AtomicBool::new(false));
let att = Arc::clone(&attempt);
let suc = Arc::clone(&success);
let task = TimeIntervalBuilder::new(work(move || {
let a = Arc::clone(&att);
let s = Arc::clone(&suc);
async move {
let current = a.fetch_add(1, Ordering::SeqCst) + 1;
let status = if current < 3 { 503 } else { 200 };
if status == 200 {
s.store(true, Ordering::SeqCst);
WorkResult::Done(())
} else {
WorkResult::NeedRetry
}
}
}))
.intervals_millis([0, 1000, 2000, 4000]) .tag("http-retry")
.build()?;
beaver.enqueue(task).await?;
tokio::time::sleep(Duration::from_secs(5)).await;
assert_eq!(attempt.load(Ordering::SeqCst), 3, "Should retry 3 times");
assert!(success.load(Ordering::SeqCst), "Should eventually succeed");
Ok(())
}
#[tokio::test]
async fn test_db_reconnection_pattern() -> BeaverResult<()> {
let beaver = Beaver::new("test", 256);
let connected = Arc::new(AtomicBool::new(false));
let connection_attempts = Arc::new(AtomicU32::new(0));
let conn = Arc::clone(&connected);
let attempts = Arc::clone(&connection_attempts);
let task = TimeIntervalBuilder::new(work(move || {
let c = Arc::clone(&conn);
let a = Arc::clone(&attempts);
async move {
let attempt = a.fetch_add(1, Ordering::SeqCst) + 1;
if attempt >= 4 {
c.store(true, Ordering::SeqCst);
WorkResult::Done(())
} else {
WorkResult::NeedRetry
}
}
}))
.intervals_millis([0, 1000, 1000, 2000, 2000, 4000]) .tag("db-reconnect")
.listener(listener(
|| println!("DB connected!"),
|| println!("Reconnection aborted"),
))
.build()?;
beaver.enqueue(task).await?;
tokio::time::sleep(Duration::from_secs(6)).await;
assert!(connected.load(Ordering::SeqCst), "Should reconnect");
assert_eq!(connection_attempts.load(Ordering::SeqCst), 4);
Ok(())
}