Skip to main content

consume

Function consume 

Source
pub async fn consume<J, F, Fut>(
    batch: MessageBatch<String>,
    env: Env,
    perform: F,
) -> Result<()>
where J: DeserializeOwned, F: Fn(Ctx, J) -> Fut, Fut: Future<Output = 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).

  • Ok acknowledges 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), until max_retries = 5 sends 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
}