mod common;
use std::time::Duration;
use common::RunningWorker;
use serde::{Deserialize, Serialize};
use steda::{Result, Task, TaskContext};
#[derive(Clone, Debug, Deserialize, Serialize)]
struct PaymentCaptured {
payment_id: String,
order_id: String,
amount_cents: u64,
}
#[derive(Debug, Deserialize, Serialize)]
struct Fulfillment {
order_id: String,
accepted_payment: String,
}
const FULFILL_ORDER: Task<PaymentCaptured, Fulfillment> = Task::new("fulfill-order");
#[tokio::main(flavor = "current_thread")]
async fn main() -> Result<()> {
let steda = common::connect().await?;
let queue = steda.queue("example-webhooks")?;
queue.create().await?;
let worker = queue
.worker()
.task(FULFILL_ORDER, async |payment: PaymentCaptured, _ctx: TaskContext| {
println!(
"processing payment {} for {} (€{}.{:02})",
payment.payment_id,
payment.order_id,
payment.amount_cents / 100,
payment.amount_cents % 100
);
Ok(Fulfillment { order_id: payment.order_id, accepted_payment: payment.payment_id })
})
.build()?;
let webhook = PaymentCaptured {
payment_id: "PAY-1001".to_owned(),
order_id: "ORD-1001".to_owned(),
amount_cents: 8_990,
};
let idempotency_key = common::unique_key("payment-captured")?;
let first =
queue.spawn(FULFILL_ORDER, webhook.clone()).idempotency_key(&idempotency_key).await?;
let duplicate = queue.spawn(FULFILL_ORDER, webhook).idempotency_key(&idempotency_key).await?;
assert!(first.created());
assert!(!duplicate.created());
assert_eq!(first.task_id(), duplicate.task_id());
println!("duplicate webhook reused the same logical task");
let worker = RunningWorker::start(worker);
let fulfillment = first.result_with_timeout(Duration::from_secs(10)).await?;
println!("fulfillment completed once for {}", fulfillment.order_id);
println!("accepted payment: {}", fulfillment.accepted_payment);
worker.stop().await?;
Ok(())
}