use std::sync::Arc;
use std::time::Duration;
use taquba::object_store::ObjectStore;
use taquba::object_store::memory::InMemory;
use taquba::{MockClock, OpenOptions, Queue, QueueConfig};
pub const QUEUE_PATH: &str = "test";
const WAIT_DEADLINE: Duration = Duration::from_secs(30);
const WAIT_INTERVAL: Duration = Duration::from_millis(20);
pub async fn open_queue(clock: MockClock) -> (Arc<dyn ObjectStore>, Arc<Queue>) {
let store: Arc<dyn ObjectStore> = Arc::new(InMemory::new());
let opts = OpenOptions::default()
.clock(Arc::new(clock))
.default_queue_config(QueueConfig::default().retry_backoff_base(Duration::ZERO))
.reaper_interval(Duration::from_millis(10))
.scheduler_interval(Duration::from_millis(10));
let queue = Queue::open_with_options(store.clone(), QUEUE_PATH, opts)
.await
.unwrap();
(store, Arc::new(queue))
}
pub async fn wait_until<T>(what: &str, mut probe: impl AsyncFnMut() -> Option<T>) -> T {
let deadline = tokio::time::Instant::now() + WAIT_DEADLINE;
loop {
if let Some(value) = probe().await {
return value;
}
assert!(tokio::time::Instant::now() < deadline, "{what}");
tokio::time::sleep(WAIT_INTERVAL).await;
}
}