Keyboard shortcuts

Press ← or → to navigate between chapters

Press S or / to search in the book

Press ? to show this help

Press Esc to hide this help

Background jobs and schedules

This guide moves slow or retryable work out of requests with background jobs on Cloudflare Queues (ocre g job, perform_later, ocre::jobs::enqueue_all), explains retries, discarding, the dead-letter queue and idempotency, gives urgent jobs their own queue, and runs tasks on a timetable with Cron Triggers (ocre g schedule "every day at 3am", ocre schedules), all within the Workers free plan.

Before you start

How jobs run

The app’s Worker is both the producer and the consumer of its queue, <app>-jobs, bound as JOBS (plus any named queue):

handler ── ocre::jobs::enqueue ──> <app>-jobs queue ──> queue event (src/lib.rs)
                                                          └─> ocre::jobs::consume ──> jobs::perform (src/jobs/mod.rs) ──> SendWelcome::perform

A handler enqueues a job and answers as soon as Cloudflare stored the message. Moments later Cloudflare calls the Worker’s queue event with a batch of messages; Ocre decodes each one into the app’s Job enum and calls perform. Dispatch is a plain match in app code, not a registry.

Generate a job

The arguments are fields, with the same types as models (integer is i64, string is String…):

ocre g job SendWelcome user_id:integer
  create  src/jobs/send_welcome.rs
  create  src/jobs/mod.rs
  update  src/lib.rs
  update  cloudflare.config.ts

Next:
  enqueue it from a handler: jobs::SendWelcome { user_id }.perform_later(&ctx).await?
  ocre dev (jobs run locally; look for `[ocre jobs]` lines in the output)
  ocre deploy creates the queue shop-jobs and its dead-letter queue

The first job wires everything; later ones only create their file and add a variant to src/jobs/mod.rs.

FileWhat the generator writes
src/jobs/send_welcome.rspub struct SendWelcome { pub user_id: i64 } (serde), fn perform_later(self, ctx) (sends it to the queue) and async fn perform(self, ctx: &Ctx) -> Result<()>, which does nothing yet
src/jobs/mod.rsThe Job enum (one variant per job) and perform, which matches on it; keep the // ocre:jobs, // ocre:job-variants and // ocre:job-dispatch markers
src/lib.rsmod jobs; and the queue entry point (first job only)
cloudflare.config.tsThe JOBS producer binding and the consumer trigger (first job only)

src/jobs/mod.rs after the first job:

// ocre:jobs
pub mod send_welcome;
pub use send_welcome::SendWelcome;

/// Every job of the app. A queue message holds one as JSON:
/// `{"send_welcome": {"user_id": 1}}`. Renaming a variant or changing its
/// fields makes messages already queued undecodable (they are logged and
/// dropped), so change jobs when the queue is empty or add a new variant.
#[derive(Debug, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum Job {
    // ocre:job-variants
    SendWelcome(SendWelcome),
}

/// Runs one job; called by `ocre::jobs::consume` for each queue message.
/// `Ok` acknowledges it; `Err(Error::NotFound)` and other 4xx errors drop
/// it (logged); other errors retry it later. Code here runs around every
/// job, like Rails' `around_perform`.
pub async fn perform(ctx: Ctx, job: Job) -> Result<()> {
    match job {
        // ocre:job-dispatch
        Job::SendWelcome(job) => job.perform(&ctx).await,
    }
}

The entry point added to src/lib.rs:

/// Background jobs from the `JOBS` queue, run by `jobs::perform` (src/jobs/mod.rs).
#[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
}

And the queue in cloudflare.config.ts (for an app named shop), the producer after // ocre:env and the consumer after // ocre:triggers:

// in worker.env
// Background jobs (`ocre g job`): the Worker sends jobs to this queue and runs
// them (src/jobs/). ...
JOBS: bindings.queue({ name: "shop-jobs" }),

// in worker.triggers
// Runs the jobs of shop-jobs: up to 10 messages per run, waiting at most 5 s
// to fill a batch. Each run is one Worker request with 10 ms of CPU on the free
// plan: lower maxBatchSize for CPU-heavy jobs. ...
triggers.queue({ name: "shop-jobs", deadLetterQueue: "shop-jobs-failed", maxBatchSize: 10, maxBatchTimeout: 5, maxRetries: 5 }),

Write perform

A job’s fields are its arguments, stored in the queue message as JSON (128 KB at most): pass ids and small values, and load records in perform. The record may have changed or disappeared since the job was enqueued, so decide what that means. This job sends the welcome email of the User mailer to a user:

// src/jobs/send_welcome.rs
use ocre::{Ctx, OptionExt, Result};
use serde::{Deserialize, Serialize};

use crate::{mailers, models::user};

/// The job's arguments: an id, not the user record.
#[derive(Debug, Serialize, Deserialize)]
pub struct SendWelcome {
    pub user_id: i64,
}

impl SendWelcome {
    /// `Err` retries the job later; `Ok` acknowledges it.
    pub async fn perform(self, ctx: &Ctx) -> Result<()> {
        // Deleted since the job was enqueued: `NotFound` drops the job (logged), no retry.
        let user = user::find(ctx, self.user_id).await?.or_404()?;
        ocre::mail::send(ctx, mailers::user::welcome(&user.email)?).await
    }
}

Returning Ok(()) acknowledges the message. An error that another try cannot fix is dropped with a log line, like Rails’ discard_on: Error::NotFound (here, a user deleted since), BadRequest, Unauthorized, Forbidden, Invalid and PayloadTooLarge. Any other Err (a failed D1 query, a provider error from send: Error::Internal, or TooManyRequests) retries the job later. See Errors: retry or discard.

Jobs run outside any request: there is no session, no signed-in user and no request URL. Pass what perform needs (an id, a locale, an absolute base URL) as fields. Records travel as ids, like Rails’ GlobalID arguments; dates as strings or Unix times; any other type works if it implements serde’s Serialize and Deserialize (derive them, or use #[serde(with = "...")]), which replaces Rails’ custom argument serializers.

Enqueue a job

job.perform_later(&ctx).await? (generated in each job) sends the job now and returns as soon as Cloudflare stored it, like Rails’ perform_later. Under it are the functions of ocre::jobs, for any job of the Job enum:

CallDoes
ocre::jobs::enqueue(&ctx, &job)Sends the job to the default queue, due now
ocre::jobs::enqueue_in(&ctx, &job, delay)Due after delay (whole seconds, 24 hours at most), like Rails’ set(wait:)
ocre::jobs::enqueue_all(&ctx, &jobs)Many jobs, in one call per 100 (see Enqueue many jobs)
ocre::jobs::queue(&ctx, "urgent").enqueue(&job)The same on a named queue (also enqueue_in, enqueue_all)
job.perform(&ctx).awaitRuns it now, in the request, like Rails’ perform_now

This module lets a signed-in user ask for the welcome email again; register it with mod welcome; under // ocre:modules and .merge(welcome::routes()) under // ocre:routes in src/lib.rs:

// src/welcome.rs
use std::time::Duration;

use axum::{Router, extract::State, response::Redirect, routing::post};
use ocre::{Ctx, Result, Session};

use crate::{
    auth::CurrentUser,
    jobs::{Job, SendWelcome},
};

pub fn routes() -> Router<Ctx> {
    Router::new().route("/welcome", post(resend)).route("/welcome/later", post(resend_later))
}

/// Answers as soon as Cloudflare stored the message; the email goes out moments later.
async fn resend(State(ctx): State<Ctx>, session: Session, CurrentUser(user): CurrentUser) -> Result<Redirect> {
    SendWelcome { user_id: user.id }.perform_later(&ctx).await?;
    session.flash("notice", "The welcome email is on its way.")?;
    Ok(Redirect::to("/account"))
}

/// Same job, run in 10 minutes (24 hours at most).
async fn resend_later(State(ctx): State<Ctx>, CurrentUser(user): CurrentUser) -> Result<Redirect> {
    let job = Job::SendWelcome(SendWelcome { user_id: user.id });
    ocre::jobs::enqueue_in(&ctx, &job, Duration::from_secs(600)).await?;
    Ok(Redirect::to("/account"))
}

With a signed-in cookie jar:

curl -s -c jar.txt -b jar.txt -X POST http://localhost:8787/welcome -o /dev/null -w '%{http_code} %{redirect_url}\n'
303 http://localhost:8787/account

ocre dev runs the queue in-process; within about 5 seconds (max_batch_timeout) the job runs and its email is printed:

[ocre mail] not sent (MAIL_ADAPTER = "log")
From: shop <noreply@example.com>
To: ada@example.com
Subject: Welcome
...
[ocre mail] end
[ocre jobs] send_welcome done
...
[wrangler:info] QUEUE shop-jobs 2/3 (13ms)

QUEUE shop-jobs 2/3 is the local server’s summary of one consumer run (a wrangler line: cf dev runs wrangler): 2 of the 3 messages of that batch were acknowledged (this job and a deliver_later email; the third was the failing job of Errors: retry or discard).

perform_later, enqueue, enqueue_in and enqueue_all fail with Error::Internal (500, the log names the fix) when the job does not serialize, the message is over 128 KB, the delay is over 24 hours, the queue’s binding (JOBS, JOBS_URGENT…) is missing from cloudflare.config.ts, or Queues refuses the message. For later work than 24 hours, enqueue from a scheduled task or store the due time in D1.

The message

Each job is one queue message, JSON text:

{"at": 1790656502, "job": {"send_welcome": {"user_id": 1}}}

at is the Unix time the job is due (now, or now plus the delay); job is the serde form of the Job enum, whose variant names are snake_case. An email from deliver_later is {"at": ..., "mail": {...}} on the same queue.

Enqueue many jobs

ocre::jobs::enqueue_all(&ctx, &jobs).await? is Rails’ perform_all_later: it sends a list of jobs (any variants of Job) with one sendBatch call per 100 messages (256 KB), instead of one call per job. A Worker invocation may only make a limited number of calls to bindings, so use it whenever a handler or a scheduled task enqueues more than a few jobs. Each call is atomic, the list as a whole is not: if the second call fails, the first 100 jobs are queued. Every job still costs 3 Queues operations.

use ocre::{Ctx, Result, params};
use serde::Deserialize;

use crate::jobs::{Job, SendWelcome};

#[derive(Deserialize)]
struct UserId {
    id: i64,
}

/// One `sendBatch` call for up to 100 users.
pub async fn welcome_everyone(ctx: &Ctx) -> Result<()> {
    let users: Vec<UserId> = ctx.db()?.all("SELECT id FROM users ORDER BY id LIMIT ?1", params![100]).await?;
    let jobs: Vec<Job> = users.into_iter().map(|u| Job::SendWelcome(SendWelcome { user_id: u.id })).collect();
    ocre::jobs::enqueue_all(ctx, &jobs).await
}

Urgent jobs: named queues

Cloudflare Queues has no priorities: every queue delivers its messages roughly in order, in batches. Ocre gives urgent work its own queue instead, like Rails’ queue_as and Loco’s named queues: a sign-in code must not wait behind a thousand digest emails.

ocre g job SendCode user_id:integer --queue urgent
  create  src/jobs/send_code.rs
  update  cloudflare.config.ts
  update  src/jobs/mod.rs

Next:
  enqueue it from a handler: jobs::SendCode { user_id }.perform_later(&ctx).await?
  ocre dev (jobs run locally; look for `[ocre jobs]` lines in the output)
  ocre deploy creates the queue shop-jobs-urgent and its dead-letter queue

The generated perform_later sends to that queue (ocre::jobs::queue(ctx, "urgent").enqueue(...)), and cloudflare.config.ts gets the producer JOBS_URGENT and a consumer that waits at most 1 second to fill a batch:

// in worker.env
JOBS_URGENT: bindings.queue({ name: "shop-jobs-urgent" }),
// in worker.triggers
triggers.queue({ name: "shop-jobs-urgent", deadLetterQueue: "shop-jobs-urgent-failed", maxBatchSize: 10, maxBatchTimeout: 1, maxRetries: 5 }),

Every queue is consumed by the same queue event and the same perform: the queue changes when a job runs, not how. Queue names are lowercase letters, digits and -; default is the JOBS queue. Queues are free to create; a message costs the same 3 operations on any queue. Loco’s worker tags (a process that only runs some jobs) have no equivalent: there are no worker processes, and a queue per kind of work gives the same isolation.

Run code around jobs: callbacks

Rails’ before_perform, around_perform and after_perform are code around the match in perform of src/jobs/mod.rs, which runs every job; before_enqueue and after_enqueue are code in a job’s perform_later. Returning early halts, like throw :abort:

pub async fn perform(ctx: Ctx, job: Job) -> Result<()> {
    let started = ocre::now();
    let result = match job {
        // ocre:job-dispatch
        Job::SendWelcome(job) => job.perform(&ctx).await,
    };
    worker::console_log!("job took {} s", ocre::now() - started);
    result
}

Handle an error of one job there too, like Rails’ rescue_from: Job::ImportFeed(job) => job.perform(&ctx).await.or_else(|err| ...).

Enqueue after the database commits

D1 has no transaction that stays open across awaits: a group of writes is one ctx.db()?.batch(statements).await?, committed when the call returns. Enqueue after it, and a job never sees data that was rolled back, which is what Rails’ enqueue_after_transaction_commit ensures. If enqueuing fails after the commit, the handler returns the error (500) with the data saved; make the job something a scheduled task can catch up on, or accept the rare miss.

Limit concurrency

Cloudflare runs several consumer invocations of a queue in parallel when messages pile up. Rails’ limits_concurrency has two Ocre forms:

  • For a whole queue, add maxConcurrency: 1 to its triggers.queue({ ... }) entry in cloudflare.config.ts: one batch at a time. Combine it with a named queue for the jobs that must not overlap.

  • For jobs sharing a key (one import per account, one GPU task per user), generate the job with --lock: ocre g job ImportCsv account_id:integer --lock account_id. The job takes the lock import_csv:<account_id> (a row in job_locks) before running and releases it after; a second run for the same account is enqueued again every 30 s until the lock is free, without using up its retries. A run that dies loses the lock after an hour. In code of your own, ocre::jobs::lock(&db, key, owner, ttl_seconds) returns whether owner got the lock (the owner holding it gets it again, extended) and ocre::jobs::unlock(&db, key, owner) releases it:

    let db = ctx.db()?;
    let key = format!("import:{}", self.account_id);
    if !ocre::jobs::lock(&db, &key, &self.run, 3600).await? {
        return Err(Error::internal("another import of this account is running")); // retried later
    }
    let result = self.import(ctx).await;
    ocre::jobs::unlock(&db, &key, &self.run).await?;
    result

    Each lock and unlock is one D1 row written (100,000 a day on the free plan); the table is ocre::jobs::LOCKS_TABLE_SQL.

Long jobs: continue in steps

Work made of distinct stages (download, split, process, notify) is a job with steps: ocre g job ProcessVideo video_id:integer --steps fetch,split,upscale,merge,notify --lock video_id. Each step is its own queue message and its own invocation, with its own 10 ms of CPU and 50 subrequests: perform runs the current step’s method, then enqueues the job again with the next step. A step that fails is retried on its own, from that step (the steps before it do not run again), and the run keeps its id (run) through retries, so with --lock no second run for the same video starts until the last step is done. Steps must be safe to repeat, like every job.

Within one step, or in a job of one piece, a loop over many rows uses a budget. A queue batch has the limits of a request: 10 ms of CPU on the free plan, and 50 D1 queries and 50 subrequests (fetch) per invocation; past those, calls fail. A job over many rows does a slice, then enqueues itself with a cursor for the rest, like Rails’ ActiveJob::Continuable. ocre::jobs::Budget counts the calls the job may still make, and run_steps runs steps while the budget covers them:

use ocre::jobs::{Budget, Step, run_steps};
use ocre::{Ctx, Query, Result, bulk};

impl Reindex {
    pub async fn perform(self, ctx: &Ctx) -> Result<()> {
        // 50 queries per invocation, 1 kept to enqueue the rest.
        let budget = Budget::new(Budget::FREE_D1_QUERIES - 1);
        // A step reads a page and writes it back: 2 queries.
        let rest = run_steps(&budget, 2, self.after_id, |after_id| async move {
            let db = ctx.db()?;
            let page: Vec<Row> = Query::table("tracks").gt("id", after_id).order_asc("id").limit(200).all(&db).await?;
            let Some(last) = page.last().map(|row| row.id) else { return Ok(Step::Done) };
            let update = bulk::update("tracks", "id", &["slug"], &slugs(&page), true)?;
            db.execute(&update.sql, update.params).await?;
            Ok(Step::Next(last))
        })
        .await?;
        if let Some(after_id) = rest {
            ocre::jobs::enqueue(ctx, &Job::Reindex(Reindex { after_id })).await?; // the next invocation
        }
        Ok(())
    }
}

run_steps returns Some(cursor) when the budget ran out before a step said Step::Done. budget.take(n) spends calls outside steps (a fetch to another service, a lookup before the loop). Ocre does not count the calls itself: give each step the cost it really has, and keep a margin. A retried job starts again from the cursor it was enqueued with, so steps must be safe to repeat. Writing many rows per query (ocre::bulk) keeps the number of steps down.

Errors: retry or discard

What perform returns decides what happens to the message:

perform returnsOcreLike Rails
Ok(())Acknowledges it
Err(Error::NotFound), BadRequest, Unauthorized, Forbidden, Invalid, PayloadTooLargeLogs discarded, not retried and acknowledges it: another try would fail the same waydiscard_on, ActiveJob::DeserializationError
Any other Err (Internal, TooManyRequests)Retries it with a growing delayretry_on

A retried job comes back after twice the time since it was due, 30 seconds at least. For a job processed right away that gives 30 s, 1 min, 3 min, 9 min and 27 min. After maxRetries: 5 retries, Cloudflare moves the message to the dead-letter queue <app>-jobs-failed, where it stays 24 hours; inspect it in the dashboard (Queues > <app>-jobs-failed), which lists its messages. The last retry comes about 40 minutes after the job was due. The number of retries is per queue (maxRetries in cloudflare.config.ts): put jobs that need another policy on their own queue.

To change the policy of one job, map its errors in perform: Err(Error::internal(..)) to retry what Ocre would discard, Ok(()) (with a log line) to drop what it would retry.

A job that always fails (its perform returns Err(ocre::Error::internal(format!("{} answered 503", self.url)))), as logged by ocre dev:

✘ [ERROR] [ocre jobs] import_feed failed, retrying in 30 s: internal error: https://example.com/feed.xml answered 503
...
✘ [ERROR] [ocre jobs] import_feed failed, retrying in 80 s: internal error: https://example.com/feed.xml answered 503

The second delay is 80 s, not 60 s: the retry was processed 40 s after the job was due (30 s of delay plus the batch wait), and the next delay is twice that.

The log lines, all prefixed [ocre jobs]:

LineMeaning
[ocre jobs] <job> doneperform returned Ok; the message is acknowledged
[ocre jobs] <job> failed, retrying in <n> s: <error>perform returned an error worth retrying; the message comes back after n seconds
[ocre jobs] <job> discarded, not retried: <error>perform returned a 4xx error; the message is acknowledged
[ocre jobs] mail doneA deliver_later email was sent
[ocre jobs] dropped message <id>: <reason>The message did not decode; it is acknowledged and never retried

Changing jobs safely

A message is dropped when it is not an Ocre message, or when its job no longer matches the Job enum: a renamed variant, a renamed field, a new field without a default, a field whose type changed. (Fields that old messages have and the struct no longer has are ignored.) Messages wait in the queue for up to 24 hours, so:

  • Add a new variant (SendWelcomeV2) instead of changing one while messages may be queued, and remove the old one a day later.
  • New fields can be added safely with #[serde(default)]: old messages decode with the default.
  • Never rename a variant while its messages may be queued.

Make jobs safe to repeat

Queues deliver at least once: a job can run twice, for example when a Worker is evicted after perform succeeded but before the acknowledgement, or when a batch is retried. Write perform so a second run does no harm:

  • Prefer statements that are naturally idempotent: UPDATE ... SET status = 'sent', INSERT ... ON CONFLICT DO NOTHING, DELETE.

  • Record that the work was done, and check first. For example, with a welcomed_at column on users:

    let claimed = ctx.db()?
        .execute("UPDATE users SET welcomed_at = datetime('now') WHERE id = ?1 AND welcomed_at IS NULL", params![self.user_id])
        .await?;
    if claimed == 0 {
        return Ok(()); // already welcomed
    }

    Claiming before sending means an email whose sending fails after the claim is not retried; claiming after sending means a crash in between sends it twice. Pick the failure you prefer for each job.

  • Calls to other APIs: pass an idempotency key when the API accepts one (a value stored in the job’s fields, so every run sends the same one).

Schedules

A scheduled task runs on the deployed Worker at times given in plain English or as a cron expression, in UTC:

ocre g schedule nightly_cleanup "every day at 3am"
  create  src/schedules/nightly_cleanup.rs
  create  src/schedules/mod.rs
  update  cloudflare.config.ts
  update  src/lib.rs

Next:
  ocre dev, then: ocre schedules run nightly_cleanup
  ocre deploy (Cron Triggers only fire on the deployed Worker; this one runs at `0 3 * * *`, UTC)
FileWhat the generator writes
src/schedules/nightly_cleanup.rspub async fn run(ctx: &Ctx) -> Result<()>, which does nothing yet
src/schedules/mod.rsrun(ctx, cron), which matches the expression that fired to a task; keep the // ocre:schedules and // ocre:schedule-dispatch markers
cloudflare.config.tsThe expression, triggers.scheduled({ schedule: "0 3 * * *" }), after // ocre:triggers
src/lib.rsmod schedules; and the scheduled entry point (first schedule only)

The entry point calls ocre::jobs::cron, which runs the task and logs the result:

/// Cron Triggers (`triggers.scheduled` in cloudflare.config.ts), run by `schedules::run` (src/schedules/mod.rs).
#[worker::event(scheduled)]
async fn scheduled(event: worker::ScheduledEvent, env: worker::Env, _ctx: worker::ScheduleContext) {
    ocre::jobs::cron(event, env, schedules::run).await
}

When: English or cron

The generator turns a plain-English phrase into Cloudflare’s five cron fields (minute, hour, day of month, month, day of week) and keeps the phrase in the task’s comment, like Loco’s English schedules. Times are UTC, 24-hour (16:30) or with am/pm; midnight and noon work too.

PhraseCronRuns
every minute* * * * *Every minute (the shortest interval Cron Triggers allow)
every 15 minutes*/15 * * * *Every 15 minutes
every hour, hourly0 * * * *At the start of every hour
every 6 hours0 */6 * * *At 00:00, 06:00, 12:00 and 18:00
every day at 3am, daily at 03:00, at 3am0 3 * * *Every day at 03:00
midnight on tuesdays0 0 * * TUETuesdays at midnight
every monday and friday at 9:3030 9 * * MON,FRIMondays and Fridays at 09:30
every weekday at 6pm0 18 * * MON-FRIMonday to Friday at 18:00
every weekend at noon0 12 * * SAT,SUNSaturdays and Sundays at 12:00
weekly, every week0 0 * * SUNSundays at midnight
monthly, every month0 0 1 * *The first day of each month, at midnight
every month at 2am0 2 1 * *The first day of each month, at 02:00

Anything else is refused with both forms in the hint; seconds (every 15 seconds) are refused because Cron Triggers run at most once a minute. A cron expression is accepted as is: letters, digits and * , - / #, with spaces normalized; Cloudflare validates the rest on deploy (see its supported expressions).

One task per expression: ocre g schedule refuses an expression already in [triggers] crons (run the new work from the existing task, or pick another minute). The free plan allows 5 Cron Triggers per account, across all Workers; past 5 in one app, the generator adds a warning to its Next: steps. Run several tasks from one cron when you need more.

List the schedules

ocre schedules reads [triggers] crons and the dispatch of src/schedules/mod.rs, without building the app:

ocre schedules
CRON (UTC)   TASK
0 3 * * *    src/schedules/nightly_cleanup.rs
0 9 * * MON  src/schedules/weekly_digest.rs

A cron without a task shows (no task: fails when it fires) and a Next: step. With --json, the list is in schedules: [{"cron": "0 3 * * *", "task": "nightly_cleanup"}, ...].

Keep tasks short: enqueue jobs

A scheduled run has the same CPU limit as a request (10 ms on the free plan), and a failed run is only logged ([ocre cron] <cron> failed: <error>), never retried: the task runs again at its next time. So a task should do a few queries and enqueue one job per item; the jobs do the slow work, with retries. This task deletes expired sign-in tokens and enqueues a welcome email for each user who signed up in the last day:

// src/schedules/nightly_cleanup.rs
use ocre::{Ctx, Result, params};
use serde::Deserialize;

use crate::jobs::SendWelcome;

#[derive(Deserialize)]
struct UserId {
    id: i64,
}

/// Runs at 03:00 UTC. Two queries and a few enqueues fit in 10 ms of CPU;
/// the per-user work (sending email) runs in the jobs.
pub async fn run(ctx: &Ctx) -> Result<()> {
    let db = ctx.db()?;
    let expired = db.execute("DELETE FROM auth_tokens WHERE expires_at <= datetime('now')", vec![]).await?;
    worker::console_log!("nightly_cleanup: {expired} expired tokens deleted");

    let new_users: Vec<UserId> = db
        .all("SELECT id FROM users WHERE created_at >= datetime('now', '-1 day') ORDER BY id LIMIT ?1", params![20])
        .await?;
    for UserId { id } in new_users {
        SendWelcome { user_id: id }.perform_later(ctx).await?;
    }
    Ok(())
}

The LIMIT bounds the work of one run; each perform_later is one Queues write. For more than a few jobs, send them with one enqueue_all.

Run a task locally

Cron Triggers do not fire in ocre dev. While it runs, fire one task from another terminal, like Loco’s scheduler --name:

ocre schedules run nightly_cleanup
  fired nightly_cleanup (0 3 * * *); its `[ocre cron]` line is in the `ocre dev` output

It calls the dev server’s local endpoint with the task’s expression (--port if ocre dev is not on 8787); the same with curl, spaces as +:

curl 'http://localhost:8787/cdn-cgi/local/scheduled?cron=0+3+*+*+*'

The ocre dev output, with the task above and two users created that day (their jobs run a few seconds later):

nightly_cleanup: 0 expired tokens deleted
[ocre cron] 0 3 * * * done
...
[ocre mail] not sent (MAIL_ADAPTER = "log")
From: shop <noreply@example.com>
To: ada@example.com
Subject: Welcome
...
[ocre mail] end
[ocre jobs] send_welcome done
[ocre mail] not sent (MAIL_ADAPTER = "log")
From: shop <noreply@example.com>
To: grace@example.com
Subject: Welcome
...
[ocre mail] end
[ocre jobs] send_welcome done
[wrangler:info] QUEUE shop-jobs 2/2 (5ms)

An expression with no task in src/schedules/mod.rs fails, and the error names the fix:

curl 'http://localhost:8787/cdn-cgi/local/scheduled?cron=*/5+*+*+*+*'
✘ [ERROR] [ocre cron] */5 * * * * failed: internal error: no scheduled task for cron `*/5 * * * *`. Fix: add it to the match in src/schedules/mod.rs, or remove its `triggers.scheduled(...)` entry from cloudflare.config.ts

One-off tasks

Rails’ rake tasks and Loco’s cargo loco task run app code from a terminal. A Worker has no terminal: its code only runs inside workerd, for a request, a queue batch, a cron or an email. Ocre covers the uses of tasks without opening a remote-execution endpoint:

NeedOcre
Change data once (backfill a column, fix rows)A data migration (ocre g migration, then SQL in the file), or ocre sql "UPDATE ..." --remote
Seed dataocre db seed
Recurring workocre g schedule; ocre schedules run <task> runs it on demand in ocre dev
App code once in production (send a batch of emails, reindex)A job enqueued from a handler that only an admin can call (ocre g auth’s CurrentUser, plus your own admin check); it runs with retries, in steps if long

Test jobs

A job is a plain struct: build it in a unit test and check what it holds, or test the functions perform calls. perform itself needs a Ctx, which only exists inside workerd, so run it end to end in ocre dev: enqueue through a request, then read the [ocre jobs] lines, the database (ocre sql), or the emails it sent at http://localhost:8787/ocre/dev/mailers/sent.json (see Email). In a handler, job.perform(&ctx).await runs a job inline, Rails’ perform_now, with the same code as the queue.

Request tests (ocre test --e2e) check them with ocre::testing::Client::jobs(), Rails’ assert_enqueued_with and assert_performed_jobs: it reads GET /ocre/dev/jobs.json, which the first ocre g job merges into routes() (debug builds only, a 404 after ocre deploy), and lists the last 50 jobs enqueued (queue, job as JSON, name()) and the last 50 runs (job, outcome: done, discarded or retried):

let mut client = ocre::testing::Client::new();
client.post("/signups", &[("email", "ada@example.com")]);
assert_eq!(client.jobs().enqueued.last().unwrap().name(), Some("send_welcome"));
ocre::testing::eventually(|| client.jobs().performed.iter().any(|run| run.job == "send_welcome" && run.outcome == "done").then_some(()));

Deploy

ocre deploy handles the queues and crons of cloudflare.config.ts:

  • Before deploying, it lists the account’s queues (cf queues list) and creates (cf queues create) every queue named there that is missing (producers, consumers and dead-letter queues), because a consumer of a missing queue fails the deploy. For each queue it creates, it prints Created queue <name> on Cloudflare (for example Created queue shop-jobs on Cloudflare and Created queue shop-jobs-failed on Cloudflare); with --json they are listed in provisioned as "queue shop-jobs".
  • cf deploy then registers the consumer and the triggers.scheduled crons. Crons fire only on the deployed Worker.

Follow the deployed Worker’s job and cron lines in Workers Logs (dashboard: your Worker > Logs; search [ocre jobs] or [ocre cron]). See Deployment for the rest of the deploy.

Free-plan budget

Limits of the Workers Free plan (September 2026) and what Ocre does about them:

LimitValueWhat Ocre does
Queues operations10,000 a day; a message costs 3 (write, read, delete), each retry 1 more read, a dead-lettered message 1 more writeOne message per job; about 3,300 jobs a day, whatever the queue; discarded jobs are not retried
Retention24 hours on Free (not configurable)Retries stop long before: the last one comes after about 40 minutes
Message size128 KB; 100 messages and 256 KB per sendBatchenqueue refuses larger jobs with an error naming the fix (pass ids); enqueue_all splits lists into batches
QueuesNamed queues cost nothing to createocre g job --queue <name> adds one, with its own consumer
Delay24 hours, on send and on retryenqueue_in refuses longer delays
BatchesUp to 100 messages, 60 s waitmax_batch_size = 10, max_batch_timeout = 5: one consumer run (one Worker invocation) per 10 jobs
CPU10 ms per invocation, for requests, cron runs and consumer batchesJobs should be I/O (D1, mail, fetch); the jobs of a batch share one invocation, so lower max_batch_size for CPU-heavy jobs
Cron Triggers5 per accountocre g schedule warns past 5 in the app; run several tasks from one cron

See also