@lunora/scheduler ships a SchedulerDO you mount once per app. It stores
pending invocations sorted by scheduled time and fires them via HTTP on the DO
alarm. On top of that it gives you deferred jobs (runAfter/runAt), code-first
cron jobs, and two workpool flavours for bounded-concurrency background work.
Deferred jobs
Inside a mutation or action, schedule a function by path through ctx.scheduler
(codegen wires this surface):
import { mutation, v } from "@/lunora/_generated/server";
export const requestExport = mutation.input({ workspaceId: v.id("workspaces") }).mutation(async ({ ctx, args: { workspaceId } }) => {
const id = await ctx.scheduler.runAfter(5 * 60_000, "exports:run", { workspaceId });
return id;
});ctx.scheduler exposes runAfter, runAt, cancel, get, and list.
runAfter/runAt take the function path as a string (e.g. "exports:run")
and return the job id.
Under the hood ctx.scheduler.runAfter POSTs to the SchedulerDO's /schedule
endpoint with functionPath, args and scheduledFor. The origin it is called
back on is NOT on that request: the DO reads it from its own
env.LUNORA_ORIGIN_URL, because a caller-supplied dispatch target would be an
SSRF vector. Set that var (in vars for local dev, as a secret in production) or
every schedule is refused with ORIGIN_NOT_CONFIGURED.
Standalone scheduler
Outside a Lunora function (a plain Worker, a test), build the client yourself
with createScheduler. This surface takes a typed function reference (from
_generated/api) instead of a path:
import { createScheduler } from "@lunora/scheduler";
import { api } from "@/lunora/_generated/api";
const scheduler = createScheduler({
namespace: env.SCHEDULER, // SchedulerDO binding
});
const id = await scheduler.runAfter(5 * 60_000, api.email.sendReminder, { userId: "u-1" });
await scheduler.runAt(new Date("2026-06-01T12:00:00Z"), api.cleanup.run, { older: 30 });
await scheduler.cancel(id);runAfter/runAt resolve the job id — the same string cancel, get and
ctx.scheduler deal in. Read the instant it will fire at back off
scheduler.get(id).
Pass instanceName to isolate a tenant's jobs into a separate DO instance.
Per-job RunOptions cover shardKey (routing hint) and retry, a retry
policy with maxAttempts (default 5 — retries after the first dispatch, so
six deliveries in all), backoff ("exponential" | "linear"),
baseMs (default 30_000), and an optional maxMs ceiling. On exhaustion the job
is parked under a dead-letter key rather than dropped.
Cron jobs
Declare recurring jobs code-first in lunora/crons.ts with cronJobs(). Codegen
discovers the file, compiles each schedule to a cron expression, emits the
wrangler.jsonc triggers.crons array, and wires the runtime dispatch map, so you
never edit wrangler by hand:
import { cronJobs } from "@lunora/scheduler";
import { internal } from "@/lunora/_generated/api";
const crons = cronJobs();
crons.interval("clear presence", { minutes: 30 }, internal.presence.clear, {});
crons.hourly("sweep sessions", { minuteUTC: 17 }, internal.presence.sweep, {});
crons.daily("send digest", { hourUTC: 9, minuteUTC: 0 }, internal.email.digest, {});
crons.weekly("weekly report", { dayOfWeek: "monday", hourUTC: 8, minuteUTC: 0 }, internal.reports.weekly, {});
crons.monthly("monthly invoice", { day: 1, hourUTC: 0, minuteUTC: 0 }, internal.billing.invoice, {});
crons.cron("custom", "0 * * * *", internal.foo.bar, {});
export default crons;Each method takes a unique name, a schedule, a target, and optional args.
Names must be unique within one cronJobs() registry. The target may be a
function reference (a one-shot dispatch) or a durable workflow reference
(workflows.<name>). A workflow target starts a fresh workflow instance on each
fire, and its args are type-checked against the workflow's params.
interval takes exactly one of {minutes | hours}. { seconds } is rejected at definition time: Cloudflare Cron Triggers have a one-minute floor,
and the 6-field expression it would compile to is refused by wrangler deploy. For sub-minute recurrence use ctx.scheduler.runAfter/runAt (via a
workpool for bounded concurrency); { minutes: 1 } is the fastest cron-native cadence. hourly, daily, weekly, and monthly schedule at a fixed UTC
wall-clock time. Use .cron for the raw 5- or 6-field grammar when the ergonomic forms don't fit.
Migrating from Convex: { hours: 24 } is not "once a day". An interval compiles to a cron step within its field's period, so interval.hours is capped
at 23 and must divide 24: */24 in a 0-23 field is not a recurrence. { days: 1 } is not a unit either. Use crons.daily(name, { hourUTC, minuteUTC }, …) instead.
Prefer hourly/daily over interval generally: they let you place a job off the boundary, which is how you stop a dozen jobs from all firing at :00.
Imperative cron triggers
If you'd rather emit a wrangler.jsonc fragment yourself, createCronTrigger
returns the snippet and dispatcher metadata for a single recurring function:
import { createCronTrigger } from "@lunora/scheduler";
import { internal } from "@/lunora/_generated/api";
const trigger = createCronTrigger({
schedule: "0 3 * * *",
fn: internal.cleanup.cleanupOldMessages,
});
// trigger.crons → ["0 3 * * *"]
// trigger.wranglerJsonc → the JSON snippet to paste under triggers.crons
// trigger.dispatcher → { functionPath, args }Both surfaces validate the expression eagerly via cron-parser, so a malformed
schedule throws at authoring time. isValidCronExpression /
assertValidCronExpression are exported if you need the check standalone.
Workpools
For bounded-concurrency background work, createWorkpool builds a named logical
pool inside the same SchedulerDO, with no extra binding. The DO dispatches at most
maxConcurrency of the pool's jobs at once and queues the rest durably, draining
as the runtime reports completions:
import { createWorkpool } from "@lunora/scheduler";
import { internal } from "@/lunora/_generated/api";
const pool = createWorkpool({
namespace: env.SCHEDULER,
maxConcurrency: 5,
name: "stripe-sync",
});
const { id } = await pool.enqueue(internal.stripe.sync, { invoiceId }, { retry: { maxAttempts: 3 } });
const { inFlight, queued, maxConcurrency } = await pool.status();
await pool.cancel(id);What bounds a drain
Two platform ceilings sit above maxConcurrency, and both apply to the drain as
a whole rather than to one job:
- Six concurrent dispatches. A Durable Object may have at most six
connections simultaneously waiting for response headers, so one alarm drain
keeps six jobs in flight. A
maxConcurrencyabove six is not wrong, just not reachable in a single drain — the surplus goes out on the following alarms. - A 15-minute alarm. An alarm handler has 15 minutes of wall time, and the whole drain shares it. Work that does not fit is deferred to the next alarm, not dropped: an unreached job keeps both its record and its index entry.
Both matter more than they look, because the DO's dispatch only returns once the dispatched function has finished running — a dispatch's latency is the whole job's latency, not a kick's.
Long-running jobs belong on Queues
That is the case to move off this pool. A job that runs for minutes holds one of the six slots and one slice of the 15-minute budget for its whole duration, and it does so for every other job in the app, not just its own pool.
createQueueWorkpool + createQueueConsumer give each batch a fresh Worker
invocation with its own connection budget and its own clock, and the consumer's
max_concurrency is not capped at six. The dispatcher allows each job five
minutes by default (httpDispatcher({ timeoutMs })).
You give up the hard concurrency cap, per-job cancellation, and per-job status —
that is the whole trade. Reach for the DO-backed createWorkpool when you need
those and your jobs are short; reach for Queues when the jobs are long.
Concurrency, retries, and dead-lettering are configured on the consumer in
wrangler.jsonc (max_concurrency, max_retries, dead_letter_queue) rather
than in code:
import { createQueueConsumer, createQueueWorkpool, httpDispatcher } from "@lunora/scheduler";
import { internal } from "@/lunora/_generated/api";
// Producer (inside an action / Worker).
const queue = createQueueWorkpool({ queue: env.JOBS });
await queue.enqueue(internal.images.optimize, { key: "u-1/avatar.png" });
// Consumer (your Worker's queue() handler).
export const queueHandler = createQueueConsumer({
dispatch: httpDispatcher({ originUrl: "https://app.acme.test", adminToken: env.LUNORA_ADMIN_TOKEN }),
// The same producer: a job whose earlier delivery is still running when
// its last delivery arrives is re-enqueued, delayed, instead of dropped.
// The copy carries a MAC keyed by `secret`, so no other body can claim
// the id it dispatches under.
requeue: { queue: env.JOBS, secret: env.LUNORA_ADMIN_TOKEN },
});Neither workpool is for multi-step orchestration. Reach for Cloudflare Workflows (@lunora/workflow) when you need durable step.do / step.sleep /
step.waitForEvent.
Dispatch contract
When the alarm fires the DO POSTs to
${env.LUNORA_ORIGIN_URL}/_lunora/scheduler/dispatch with
{ functionPath, args, shardKey, scheduledFor, id }. The runtime unwraps
that and runs the function on the same code path as an RPC call — with one
difference: the dispatch is server-initiated, so it carries no end-user
identity. ctx.auth.userId is null and ctx.ip is undefined, which is also
what lets a job target an internal function. Pass anything the job needs to
know about a user in its args; never gate a scheduled function on ctx.auth.
A cron trigger takes a different route: the worker's scheduled() handler
dispatches it straight to the shard, without touching the SchedulerDO. So the
retry ladder, backoff, and dead-letter park above apply to runAfter/runAt
jobs only — a failed cron tick is logged as CRON_JOB_FAILED and not retried.
Delivery guarantee
A scheduled job is delivered at least once. That is the contract the retry ladder and the dead-letter park are built on, and it means a job can run more than once. There are two ways it happens, and they are not equally forgiving.
A retry after a failed attempt is the ordinary one. The attempt is over
before the next begins, so the record id — which rides to the shard as the
replay-dedup key, and to Cloudflare Workflows as the instance id — finds the
previous attempt's mark already in place. A mutation short-circuits to its
cached result, a workflow re-attaches to the instance it already created, and
an action short-circuits too.
A re-delivery that overlaps a live attempt is the hard one, and it comes from exactly one source: the Durable Object that owns the schedule being evicted or crashing while a dispatch is open. Every target kind now survives it.
| Target | Overlapping re-delivery |
|---|---|
mutation | Runs once. The dedup read holds the shard's single-writer gate. |
workflow / agent | Runs once. The record id is the instance id; a second create is a no-op. |
action | Runs once. The shard declines the second delivery while the first is in flight¹. |
action used to be the exception. It is dispatched without the shard's
single-writer gate — taking it would let any caller freeze a whole shard for the
length of an action's outbound I/O — and its dedup row is written only once the
handler returns, so a second delivery arriving mid-flight found nothing to dedup
against and ran the handler again, alongside the first.
The shard now claims the dedup key before running the handler and replaces
the claim with the result afterwards. A delivery that arrives while the claim is
live is answered 409 DISPATCH_IN_PROGRESS and the handler is not entered a
second time.
¹ For up to 15 minutes — see the staleness ceiling below.
The decline is temporary, never terminal. A 409 is not a success, so the scheduler keeps the record and re-arms it — the job is still delivered at least once. What it no longer does is run twice at the same time. The retry that follows lands after the first attempt has settled and is served that attempt's cached result.
A decline does not spend the retry budget. The scheduler re-arms a
DISPATCH_IN_PROGRESS decline without counting it as an attempt, at the backoff
the next real attempt would get and never sooner than 30 seconds. A decline says
the first attempt is still doing the work; charging it would let a slow action
use up maxAttempts while running fine, and if that attempt then failed the job
would be parked in /dead having failed only once.
Staleness. A claim is released when its dispatch finishes — success or throw — and it cannot outlive the Durable Object instance that holds it: if that instance is gone, so is the handler it was running, and a successor finds no claim and runs the job. What a claim can outlive is its handler hanging while the instance stays up (an outbound call with no timeout, a promise nothing resolves). So a claim older than 15 minutes — the same bound as the lease below — is treated as stale, and the next delivery runs.
That ceiling is also where the guarantee ends. An action still legitimately running after 15 minutes is no longer protected: a re-delivery past that point runs alongside it.
Two reservations still bound how soon a lost job is retried, and neither is a correctness gap:
- A claimed job stays reserved by the scheduler for 15 minutes before another instance may fire it — the platform's own ceiling on how long a dispatch this scheduler started can still be open (an alarm invocation, and the fetch it owns, are torn down at that mark). A job whose instance genuinely died waits that out. It is delayed, never dropped.
- A job whose handler hangs is declined until its claim goes stale, then runs again and is charged as usual if that attempt fails. The dead-letter is reached only by attempts that actually ran and failed.
Queue consumers and workflow steps that call an action through ctx.run treat
the 409 as retryable too, but they do count it against their own retry budget
(max_retries for a queue, the step's retries for a workflow). Give an action
that can overlap its own re-delivery enough budget — or a retry delay — to
outlast one run.
A scheduled action that can run longer than 15 minutes must still be idempotent, or mark its own completion — a row it writes first and checks on
entry, a provider-side idempotency key. Past that point the shard's claim is stale and an overlapping re-delivery runs. Every action also has to be
idempotent against an ordinary sequential retry whose first attempt's response was lost before its dedup row was written, which is the plain
at-least-once case every scheduled job has. For genuinely long work, createQueueWorkpool + createQueueConsumer remains the better fit: it bounds each
job at five minutes and reports completion explicitly.