pub async fn consume<J, F, Fut>(
batch: MessageBatch<String>,
env: Env,
perform: F,
) -> Result<()>Expand description
Runs a batch of queue messages through the app’s perform; the Worker’s queue entry point.
Each job goes to perform(ctx, job), each email from
deliver_later to
mail::send. Messages run one after the other in one
Worker invocation (10 per batch with the generated max_batch_size, all
sharing its CPU limit: 10 ms on Free).
Okacknowledges the message and logs[ocre jobs] <job> done.- An error that another try cannot fix, like Rails’
discard_on(Error::NotFound,BadRequest,Unauthorized,Forbidden,Invalid,PayloadTooLarge: a record deleted since the job was enqueued, bad input) logs[ocre jobs] <job> discarded, not retried: <error>and acknowledges it. - Any other
Err(Internal,TooManyRequests) logs[ocre jobs] <job> failed, retrying in <n> s: <error>and retries the message after twice the time since it was due, between 30 s and 24 hours (30 s, 1 min, 3 min, 9 min, 27 min), untilmax_retries = 5sends it to the dead-letter queue<app>-jobs-failed. - A message that does not decode (not an Ocre message, or a job renamed
or changed while messages were queued) is logged as
[ocre jobs] dropped message <id>: <reason>and acknowledged, never retried.
Delivery is at-least-once, so a job can run twice: make it safe to repeat.
Free-plan cost: reading and acknowledging a message are 2 Queues operations (of 10,000 a day); each retry is one more read, a dead-lettered message one more write.
§Errors
Never returns Err itself: job errors are logged, then discarded or retried.
§Examples
ocre g job writes the entry point in src/lib.rs and perform in src/jobs/mod.rs:
mod jobs {
use ocre::{Ctx, Result};
use serde::Deserialize;
#[derive(Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum Job {
SendWelcome { user_id: i64 },
}
pub async fn perform(_ctx: Ctx, job: Job) -> Result<()> {
match job {
Job::SendWelcome { user_id } => {
Ok(())
}
}
}
}
#[worker::event(queue)]
async fn queue(batch: worker::MessageBatch<String>, env: worker::Env, _ctx: worker::Context) -> worker::Result<()> {
ocre::jobs::consume(batch, env, jobs::perform).await
}