use busybeaver::{
listener, listener_with_error, work, Beaver, BeaverResult, FixedCountBuilder,
FixedCountProgress, RuntimeError, WorkResult,
};
use std::sync::atomic::{AtomicBool, AtomicU32, Ordering};
use std::sync::Arc;
use std::time::Duration;
#[tokio::test]
async fn test_basic_fixed_count_task() -> BeaverResult<()> {
let beaver = Beaver::new("test", 256);
let counter = Arc::new(AtomicU32::new(0));
let counter_clone = Arc::clone(&counter);
let task = FixedCountBuilder::new(work(move || {
let c = Arc::clone(&counter_clone);
async move {
c.fetch_add(1, Ordering::SeqCst);
WorkResult::NeedRetry
}
}))
.count(5)
.build()?;
beaver.enqueue(task).await?;
tokio::time::sleep(Duration::from_millis(200)).await;
assert_eq!(
counter.load(Ordering::SeqCst),
5,
"Should execute exactly 5 times"
);
Ok(())
}
#[tokio::test]
async fn test_fixed_count_one() -> BeaverResult<()> {
let beaver = Beaver::new("test", 256);
let counter = Arc::new(AtomicU32::new(0));
let counter_clone = Arc::clone(&counter);
let task = FixedCountBuilder::new(work(move || {
let c = Arc::clone(&counter_clone);
async move {
c.fetch_add(1, Ordering::SeqCst);
WorkResult::NeedRetry
}
}))
.count(1)
.build()?;
beaver.enqueue(task).await?;
tokio::time::sleep(Duration::from_millis(100)).await;
assert_eq!(
counter.load(Ordering::SeqCst),
1,
"count=1 should execute exactly once"
);
Ok(())
}
#[tokio::test]
async fn test_fixed_count_zero_normalized_to_one() -> BeaverResult<()> {
let beaver = Beaver::new("test", 256);
let counter = Arc::new(AtomicU32::new(0));
let counter_clone = Arc::clone(&counter);
let task = FixedCountBuilder::new(work(move || {
let c = Arc::clone(&counter_clone);
async move {
c.fetch_add(1, Ordering::SeqCst);
WorkResult::NeedRetry
}
}))
.count(0) .build()?;
beaver.enqueue(task).await?;
tokio::time::sleep(Duration::from_millis(100)).await;
assert_eq!(
counter.load(Ordering::SeqCst),
1,
"count=0 should be normalized to 1"
);
Ok(())
}
#[tokio::test]
async fn test_fixed_count_stops_early_on_done() -> BeaverResult<()> {
let beaver = Beaver::new("test", 256);
let counter = Arc::new(AtomicU32::new(0));
let counter_clone = Arc::clone(&counter);
let task = FixedCountBuilder::new(work(move || {
let c = Arc::clone(&counter_clone);
async move {
let count = c.fetch_add(1, Ordering::SeqCst) + 1;
if count == 3 {
WorkResult::Done(()) } else {
WorkResult::NeedRetry
}
}
}))
.count(10) .build()?;
beaver.enqueue(task).await?;
tokio::time::sleep(Duration::from_millis(200)).await;
assert_eq!(
counter.load(Ordering::SeqCst),
3,
"Should stop early when Done is returned"
);
Ok(())
}
#[tokio::test]
async fn test_fixed_count_default_count() -> BeaverResult<()> {
let beaver = Beaver::new("test", 256);
let counter = Arc::new(AtomicU32::new(0));
let counter_clone = Arc::clone(&counter);
let task = FixedCountBuilder::new(work(move || {
let c = Arc::clone(&counter_clone);
async move {
c.fetch_add(1, Ordering::SeqCst);
WorkResult::NeedRetry
}
}))
.build()?;
beaver.enqueue(task).await?;
tokio::time::sleep(Duration::from_millis(200)).await;
assert_eq!(
counter.load(Ordering::SeqCst),
3,
"Default count should be 3"
);
Ok(())
}
#[tokio::test]
async fn test_fixed_count_large_count() -> BeaverResult<()> {
let beaver = Beaver::new("test", 256);
let counter = Arc::new(AtomicU32::new(0));
let counter_clone = Arc::clone(&counter);
let task = FixedCountBuilder::new(work(move || {
let c = Arc::clone(&counter_clone);
async move {
c.fetch_add(1, Ordering::SeqCst);
WorkResult::NeedRetry
}
}))
.count(100)
.build()?;
beaver.enqueue(task).await?;
tokio::time::sleep(Duration::from_millis(500)).await;
assert_eq!(
counter.load(Ordering::SeqCst),
100,
"Should execute 100 times"
);
Ok(())
}
#[tokio::test]
async fn test_fixed_count_progress_callback() -> BeaverResult<()> {
let beaver = Beaver::new("test", 256);
let progress_log = Arc::new(std::sync::Mutex::new(Vec::new()));
let progress_clone = Arc::clone(&progress_log);
let progress_fn: Arc<dyn FixedCountProgress> =
Arc::new(move |current: u32, total: u32, tag: &str| {
let mut log = progress_clone.lock().unwrap();
log.push((current, total, tag.to_string()));
});
let task = FixedCountBuilder::new(work(|| async { WorkResult::NeedRetry }))
.count(5)
.tag("progress-test")
.progress(progress_fn)
.build()?;
beaver.enqueue(task).await?;
tokio::time::sleep(Duration::from_millis(200)).await;
let log = progress_log.lock().unwrap();
assert_eq!(log.len(), 5, "Progress should be called 5 times");
assert_eq!(log[0], (1, 5, "progress-test".to_string()));
assert_eq!(log[1], (2, 5, "progress-test".to_string()));
assert_eq!(log[2], (3, 5, "progress-test".to_string()));
assert_eq!(log[3], (4, 5, "progress-test".to_string()));
assert_eq!(log[4], (5, 5, "progress-test".to_string()));
Ok(())
}
#[tokio::test]
async fn test_fixed_count_progress_with_tag() -> BeaverResult<()> {
let beaver = Beaver::new("test", 256);
let received_tag = Arc::new(std::sync::Mutex::new(String::new()));
let tag_clone = Arc::clone(&received_tag);
let progress_fn: Arc<dyn FixedCountProgress> = Arc::new(move |_: u32, _: u32, tag: &str| {
let mut t = tag_clone.lock().unwrap();
*t = tag.to_string();
});
let task = FixedCountBuilder::new(work(|| async { WorkResult::Done(()) }))
.count(1)
.tag("my-special-tag")
.progress(progress_fn)
.build()?;
beaver.enqueue(task).await?;
tokio::time::sleep(Duration::from_millis(100)).await;
assert_eq!(
*received_tag.lock().unwrap(),
"my-special-tag",
"Progress should receive correct tag"
);
Ok(())
}
#[tokio::test]
async fn test_fixed_count_progress_empty_tag() -> BeaverResult<()> {
let beaver = Beaver::new("test", 256);
let received_tag = Arc::new(std::sync::Mutex::new(String::from("placeholder")));
let tag_clone = Arc::clone(&received_tag);
let progress_fn: Arc<dyn FixedCountProgress> = Arc::new(move |_: u32, _: u32, tag: &str| {
let mut t = tag_clone.lock().unwrap();
*t = tag.to_string();
});
let task = FixedCountBuilder::new(work(|| async { WorkResult::Done(()) }))
.count(1)
.progress(progress_fn)
.build()?;
beaver.enqueue(task).await?;
tokio::time::sleep(Duration::from_millis(100)).await;
assert_eq!(
*received_tag.lock().unwrap(),
"",
"Empty tag should be passed as empty string"
);
Ok(())
}
#[tokio::test]
async fn test_fixed_count_on_error_retries_exhausted() -> BeaverResult<()> {
let beaver = Beaver::new("test", 256);
let completed = Arc::new(AtomicBool::new(false));
let exhausted = Arc::new(AtomicBool::new(false));
let completed_clone = Arc::clone(&completed);
let exhausted_clone = Arc::clone(&exhausted);
let task = FixedCountBuilder::new(work(|| async { WorkResult::NeedRetry }))
.count(3)
.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_fixed_count_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 counter = Arc::new(AtomicU32::new(0));
let counter_clone = Arc::clone(&counter);
let task = FixedCountBuilder::new(work(move || {
let c = Arc::clone(&counter_clone);
async move {
c.fetch_add(1, Ordering::SeqCst);
WorkResult::Done(()) }
}))
.count(5)
.listener(listener(
move || {
completed_clone.store(true, Ordering::SeqCst);
},
|| {},
))
.build()?;
beaver.enqueue(task).await?;
tokio::time::sleep(Duration::from_millis(200)).await;
assert_eq!(counter.load(Ordering::SeqCst), 1);
assert!(
completed.load(Ordering::SeqCst),
"on_complete should be called when Done is returned"
);
Ok(())
}
#[tokio::test]
async fn test_fixed_count_on_interrupt() -> BeaverResult<()> {
let beaver = Beaver::new("test", 256);
let interrupted = Arc::new(AtomicBool::new(false));
let interrupted_clone = Arc::clone(&interrupted);
let task = FixedCountBuilder::new(work(|| async {
tokio::time::sleep(Duration::from_millis(100)).await;
WorkResult::NeedRetry
}))
.count(100)
.listener(listener(
|| {},
move || {
interrupted_clone.store(true, Ordering::SeqCst);
},
))
.build()?;
beaver.enqueue(task).await?;
tokio::time::sleep(Duration::from_millis(50)).await;
beaver.cancel_all().await?;
tokio::time::sleep(Duration::from_millis(200)).await;
assert!(
interrupted.load(Ordering::SeqCst),
"on_interrupt should be called"
);
Ok(())
}
#[tokio::test]
async fn test_fixed_count_with_tag() -> BeaverResult<()> {
let task = FixedCountBuilder::new(work(|| async { WorkResult::Done(()) }))
.count(3)
.tag("my-fixed-task")
.build()?;
assert_eq!(task.tag(), "my-fixed-task");
Ok(())
}
#[tokio::test]
async fn test_fixed_count_without_tag() -> BeaverResult<()> {
let task = FixedCountBuilder::new(work(|| async { WorkResult::Done(()) }))
.count(3)
.build()?;
assert_eq!(task.tag(), "");
Ok(())
}
#[tokio::test]
async fn test_fixed_count_interrupt_mid_execution() -> BeaverResult<()> {
let beaver = Beaver::new("test", 256);
let counter = Arc::new(AtomicU32::new(0));
let counter_clone = Arc::clone(&counter);
let task = FixedCountBuilder::new(work(move || {
let c = Arc::clone(&counter_clone);
async move {
c.fetch_add(1, Ordering::SeqCst);
tokio::time::sleep(Duration::from_millis(100)).await;
WorkResult::NeedRetry
}
}))
.count(10)
.build()?;
beaver.enqueue(task).await?;
tokio::time::sleep(Duration::from_millis(50)).await;
beaver.cancel_all().await?;
tokio::time::sleep(Duration::from_millis(500)).await;
let count = counter.load(Ordering::SeqCst);
assert!(
count < 10,
"Should be interrupted before completing all executions"
);
Ok(())
}
#[tokio::test]
async fn test_fixed_count_cancel_stops_execution() -> BeaverResult<()> {
let beaver = Beaver::new("test_fixed_count_cancel_stops_execution", 256);
let execution_count = Arc::new(AtomicU32::new(0));
let interrupted = Arc::new(AtomicBool::new(false));
let ec = Arc::clone(&execution_count);
let int = Arc::clone(&interrupted);
let task = FixedCountBuilder::new(work(move || {
let c = Arc::clone(&ec);
async move {
c.fetch_add(1, Ordering::SeqCst);
tokio::time::sleep(Duration::from_millis(100)).await;
WorkResult::NeedRetry
}
}))
.count(100) .listener(listener(
|| {},
move || {
int.store(true, Ordering::SeqCst);
},
))
.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 < 100,
"Task should not complete all 100 executions, got {}",
count
);
assert!(
interrupted.load(Ordering::SeqCst),
"on_interrupt should be called"
);
Ok(())
}
#[tokio::test]
async fn test_fixed_count_builder_chain_order() -> BeaverResult<()> {
let progress_fn: Arc<dyn FixedCountProgress> = Arc::new(|_: u32, _: u32, _: &str| {});
let _task1 = FixedCountBuilder::new(work(|| async { WorkResult::Done(()) }))
.count(5)
.tag("task1")
.progress(Arc::clone(&progress_fn))
.listener(listener(|| {}, || {}))
.build()?;
let _task2 = FixedCountBuilder::new(work(|| async { WorkResult::Done(()) }))
.listener(listener(|| {}, || {}))
.progress(Arc::clone(&progress_fn))
.tag("task2")
.count(5)
.build()?;
Ok(())
}
#[tokio::test]
async fn test_fixed_count_unique_task_id() -> BeaverResult<()> {
let task1 = FixedCountBuilder::new(work(|| async { WorkResult::Done(()) }))
.count(3)
.build()?;
let task2 = FixedCountBuilder::new(work(|| async { WorkResult::Done(()) }))
.count(3)
.build()?;
assert_ne!(task1.id(), task2.id());
Ok(())
}
#[tokio::test]
async fn test_fixed_count_async_work() -> BeaverResult<()> {
let beaver = Beaver::new("test", 256);
let results = Arc::new(std::sync::Mutex::new(Vec::new()));
let results_clone = Arc::clone(&results);
let task = FixedCountBuilder::new(work(move || {
let r = Arc::clone(&results_clone);
async move {
tokio::time::sleep(Duration::from_millis(10)).await;
r.lock().unwrap().push(42);
WorkResult::NeedRetry
}
}))
.count(3)
.build()?;
beaver.enqueue(task).await?;
tokio::time::sleep(Duration::from_millis(200)).await;
assert_eq!(results.lock().unwrap().len(), 3);
Ok(())
}
#[tokio::test]
async fn test_multiple_fixed_count_tasks() -> BeaverResult<()> {
let beaver = Beaver::new("test", 256);
let total = Arc::new(AtomicU32::new(0));
for i in 1..=3 {
let t = Arc::clone(&total);
let task = FixedCountBuilder::new(work(move || {
let total = Arc::clone(&t);
async move {
total.fetch_add(1, Ordering::SeqCst);
WorkResult::NeedRetry
}
}))
.count(i as u32) .build()?;
beaver
.enqueue_on_new_thread(task, format!("dam-{}", i), 256, false)
.await?;
}
tokio::time::sleep(Duration::from_millis(300)).await;
assert_eq!(total.load(Ordering::SeqCst), 6);
Ok(())
}
#[tokio::test]
async fn test_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 = FixedCountBuilder::new(work(move || {
let a = Arc::clone(&att);
let s = Arc::clone(&suc);
async move {
let current = a.fetch_add(1, Ordering::SeqCst) + 1;
if current >= 3 {
s.store(true, Ordering::SeqCst);
WorkResult::Done(())
} else {
WorkResult::NeedRetry
}
}
}))
.count(5) .build()?;
beaver.enqueue(task).await?;
tokio::time::sleep(Duration::from_millis(200)).await;
assert_eq!(attempt.load(Ordering::SeqCst), 3);
assert!(success.load(Ordering::SeqCst));
Ok(())
}