mod common;
use std::time::Duration;
use common::RunningWorker;
use serde::{Deserialize, Serialize};
use steda::{Result, Step, Task, TaskContext};
#[derive(Debug, Deserialize, Serialize)]
struct FulfillOrderInput {
order_id: String,
amount_cents: u64,
}
#[derive(Debug, Deserialize, Serialize)]
struct Reservation {
reservation_id: String,
}
#[derive(Debug, Deserialize, Serialize)]
struct Payment {
payment_id: String,
}
#[derive(Debug, Deserialize, Serialize)]
struct Shipment {
tracking_number: String,
}
#[derive(Debug, Deserialize, Serialize)]
struct FulfillOrderOutput {
order_id: String,
reservation_id: String,
payment_id: String,
tracking_number: String,
}
const FULFILL_ORDER: Task<FulfillOrderInput, FulfillOrderOutput> = Task::new("fulfill-order");
const RESERVE_INVENTORY: Step<Reservation> = Step::new("reserve-inventory");
const CAPTURE_PAYMENT: Step<Payment> = Step::new("capture-payment");
const CREATE_SHIPMENT: Step<Shipment> = Step::new("create-shipment");
async fn fulfill_order(input: FulfillOrderInput, ctx: TaskContext) -> Result<FulfillOrderOutput> {
let reservation = ctx
.step(RESERVE_INVENTORY, async || {
println!("reserving inventory for {}", input.order_id);
Ok(Reservation { reservation_id: "RES-1001".to_owned() })
})
.await?;
let payment = ctx
.step(CAPTURE_PAYMENT, async || {
println!(
"capturing payment of €{}.{:02}",
input.amount_cents / 100,
input.amount_cents % 100
);
Ok(Payment { payment_id: "PAY-1001".to_owned() })
})
.await?;
let shipment = ctx
.step(CREATE_SHIPMENT, async || {
println!("creating shipment");
Ok(Shipment { tracking_number: "TRACK-1001".to_owned() })
})
.await?;
Ok(FulfillOrderOutput {
order_id: input.order_id,
reservation_id: reservation.reservation_id,
payment_id: payment.payment_id,
tracking_number: shipment.tracking_number,
})
}
#[tokio::main(flavor = "current_thread")]
async fn main() -> Result<()> {
let steda = common::connect().await?;
let queue = steda.queue("example-multistep")?;
queue.create().await?;
let worker = queue.worker().task(FULFILL_ORDER, fulfill_order).build()?;
let worker = RunningWorker::start(worker);
let task = queue
.spawn(
FULFILL_ORDER,
FulfillOrderInput { order_id: "ORD-1001".to_owned(), amount_cents: 14_950 },
)
.await?;
let result = task.result_with_timeout(Duration::from_secs(10)).await?;
println!("order {} fulfilled", result.order_id);
println!(" reservation: {}", result.reservation_id);
println!(" payment: {}", result.payment_id);
println!(" tracking: {}", result.tracking_number);
worker.stop().await?;
Ok(())
}