use std::{
error::Error,
sync::atomic::{AtomicUsize, Ordering},
time::Duration,
};
use queuey::prelude::*;
const DEFAULT_URL: &str = "amqp://guest:guest@localhost:5672/%2f";
const EXPECTED: usize = 3;
const DEADLINE: Duration = Duration::from_secs(30);
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Queues)]
#[queues(prefix = "aq-example")]
enum AppQueues {
#[queue(
prefetch = 8,
retry(
max_attempts = 3,
backoff = "exponential",
base = "500ms",
max = "5s",
jitter = false
)
)]
Emails,
}
#[derive(Debug, Serialize, Deserialize, Job)]
#[job(queue = AppQueues::Emails)]
struct SendEmail {
to: String,
flaky_for: u32,
}
struct Settled(AtomicUsize);
impl Settled {
fn bump(&self) {
self.0.fetch_add(1, Ordering::SeqCst);
}
fn count(&self) -> usize {
self.0.load(Ordering::SeqCst)
}
}
struct EmailHandler {
settled: Arc<Settled>,
}
#[async_trait]
impl JobHandler for EmailHandler {
type Job = SendEmail;
async fn handle(&self, job: SendEmail, ctx: JobContext) -> Result<(), JobError> {
tracing::info!(
to = %job.to,
attempt = ctx.attempt,
max_attempts = ctx.max_attempts,
age = ?ctx.age,
"handling email"
);
if ctx.attempt <= job.flaky_for {
return Err(JobError::retryable_msg("smtp connection reset"));
}
self.settled.bump();
Ok(())
}
}
#[tokio::main]
async fn main() -> Result<(), Box<dyn Error>> {
tracing_subscriber::fmt()
.with_env_filter(
tracing_subscriber::EnvFilter::try_from_default_env().unwrap_or_else(|_| "info".into()),
)
.init();
let url = std::env::var("AMQP_URL").unwrap_or_else(|_| DEFAULT_URL.to_owned());
tracing::info!(%url, "connecting");
let backend = Arc::new(RabbitMqBackend::connect(&url).await?);
let settled = Arc::new(Settled(AtomicUsize::new(0)));
let producer = Producer::<AppQueues, _>::new(backend.clone()).await?;
let worker = Worker::<AppQueues, _>::builder(backend)
.close_backend_on_shutdown(true)
.handler(EmailHandler {
settled: settled.clone(),
})
.build()
.await?;
let handle = worker.handle();
let running = tokio::spawn(worker.run());
for (to, flaky_for) in [
("ada@example.com", 0),
("grace@example.com", 0),
("flaky@example.com", 2),
] {
let id = producer
.enqueue(&SendEmail {
to: to.to_owned(),
flaky_for,
})
.await?;
tracing::info!(%id, to, flaky_for, "enqueued");
}
let all_settled = async {
while settled.count() < EXPECTED {
tokio::time::sleep(Duration::from_millis(50)).await;
}
};
let timed_out = tokio::select! {
() = all_settled => false,
_ = tokio::time::sleep(DEADLINE) => true,
_ = tokio::signal::ctrl_c() => {
tracing::warn!("interrupted; shutting down");
false
}
};
tracing::info!(settled = settled.count(), "shutting down");
handle.shutdown();
running.await??;
if timed_out {
return Err(format!(
"only {}/{EXPECTED} jobs settled within {DEADLINE:?}",
settled.count()
)
.into());
}
tracing::info!("done");
Ok(())
}