Skip to content
DocspackagesDocumentation

@lunora/queue

Cloudflare Queues for Lunora — defineQueue producers, a generated queue() push consumer (or HTTP pull), and the typed ctx.queues surface.

PackagesQueue

@lunora/queue brings Cloudflare Queues to Lunora. Declare a queue with defineQueue in lunora/queues.ts; codegen wires a typed ctx.queues.<name> producer onto mutations and actions and generates the worker queue() push consumer. The config layer reconciles the wrangler queues.producers[] / queues.consumers[] entries from the same definition.

pnpm add @lunora/queue

Scaffold a queue:

vis generate lunora-queue --name=emailQueue

Declare a queue

// lunora/queues.ts
import { defineQueue } from "@lunora/queue";

import { api } from "./_generated/api";

export const emailQueue = defineQueue<{ to: string }>({
    handler: async (_ctx, batch) => {
        for (const message of batch.messages) {
            await message.run(api.email.send, { to: message.body.to });
            message.ack();
        }
    },
    // Push-consumer tuning (optional):
    // maxBatchSize: 10, maxBatchTimeout: 5, maxRetries: 3, deadLetterQueue: "dlq",
});

The handler runs inside a Lunora context: call any query/mutation/action with message.run(api.x.y, args) (the dispatch goes through the same path the scheduler and workflows use). Each message is ack/retry-able for at-least-once processing.

Note the loop does not wrap message.run(...) in a try/catch. Letting the failure propagate is what the consumer needs to isolate a poison message — see below.

Poison messages

message.run(...) is ctx.run(...) pinned to that one message, and the pin is what keeps one bad message from taking its whole batch down.

Cloudflare delivers up to 100 messages per batch and redelivers the whole batch when the handler throws. So a message that can never succeed — a deleted row, a body that fails validation — burns every sibling's retry budget with it and eventually dead-letters the lot.

When a call made through message.run(...) fails deterministically (the dispatched function answered 400, 403, 404 or 422 — a retry would fail identically), the consumer attributes the failure to that message: it is acked and taken out of the queue, every message the handler had not yet decided is explicitly retried, and the batch as a whole resolves. Messages the handler already acked or retryed keep the decision it made. Transient failures (408, 429, 5xx, a timeout) are unchanged — the batch still throws and workerd retries it.

Attribution needs the pin and the throw. A bare ctx.run(api.x.y, args) inside the loop carries no message id, so its failure is unattributable and the whole batch retries — which is why the loop should call message.run. And a try/catch around message.run(...) that swallows the error and calls message.retry() opts the message out entirely: the consumer only attributes a failure the handler let propagate, and an explicit retry() is a decision it never overrides. The poison message then burns its retry budget and dead-letters exactly as if isolation did not exist. If you need to log the failure yourself, re-throw after logging. Isolation does not depend on the dev capture sink; it applies in production too.

Because an isolated message is acked, it never reaches the dead-letter queue. It is recorded with outcome error (and the failure message) in the Studio Queues log, which is where you see it — so keep dev capture on, or log it in a catch that re-throws.

Redeliveries and exactly-once calls

A redelivered message runs its handler again from the top. Every call made through message.run(...) carries a replay-dedup id, <messageId>#<n>, where n counts that message's calls in order. The shard applies each id once, so a redelivery gets the first delivery's results back instead of charging the customer twice. The handler must make its calls in the same order on every delivery. A handler that branches on something non-deterministic should pass its own { dedupId }. A bare ctx.run(...) carries no id and is at-least-once.

The dedup window is 24 hours. The shard prunes dedup rows older than that, so a redelivery that arrives more than 24 hours after the first delivery applies its calls again. Keep retryDelay × maxRetries well under a day, or make the called functions idempotent on their own key.

Slow calls and the retry budget

message.run(...) gives up on a call after 30 seconds, but the shard keeps running it. The message fails, and its redelivery re-issues the same id while the first run is still going. The shard answers that with 409 DISPATCH_IN_PROGRESS instead of running the call twice.

Cloudflare Queues has no retry that does not count against maxRetries, so a decline still costs one attempt. What the consumer bounds is how many it costs. A declined message is not rethrown into an immediate redelivery, which would meet the same decline again. It is retried with a 15-minute delay, the longest the shard's claim can stand. So one slow call costs at most one extra attempt. The price is latency: a declined message always waits the full 15 minutes, even when the action behind it finishes 31 seconds after it started. The consumer cannot see the action finish, so it cannot retry any sooner without risking another decline. The next delivery is then served the finished result from the replay cache, or runs the call if the first run died. A declined message is never acked, so delivery stays at-least-once.

A decline on the message's last delivery cannot be retried: the broker would dead-letter or drop it with the call still running. The consumer sends a copy back onto the same queue instead, delayed by the same 15 minutes, and acks the original only once that send succeeded. The copy gets a fresh maxRetries budget. Its handler sees the original message.id and message.body, so its message.run(...) calls carry the same dedup ids and are served the declined call's result rather than running it again. The copy goes through the queue's own producer binding, which codegen always declares.

The copy's body is an envelope under the reserved key "$lunora.requeued$": the original id, the original body as the JSON text of its wire encoding (so a bigint, Date or bytes survives), and an HMAC-SHA256 over both keyed by LUNORA_ADMIN_TOKEN. The consumer unwraps only an envelope whose MAC verifies. Any other body, including one that copies the envelope's shape, is delivered as-is under the broker's own id, so a body forwarded from outside cannot pick the dedup ids its calls carry. ctx.queues.<name>.send and sendBatch refuse a body carrying the reserved key. A dead-lettered copy sits in the DLQ as that JSON envelope; the original body is the body string inside it.

A copy is never copied again. If a copy's own last delivery is declined too, or the copy cannot be sent (the send fails, the envelope would exceed the 128 KB message limit, or the body has no wire encoding, such as a class instance), the message is retried anyway, which takes it to the deadLetterQueue or drops it if the queue has none, and the consumer logs an error. Give a queue that calls slow actions a maxRetries of at least 2, plus a deadLetterQueue.

Enqueue from a mutation or action

ctx.queues.<exportName> is a typed producer on Mutation and Action contexts (enqueue is a side effect, so it's excluded from the deterministic query context):

import { mutation } from "./_generated/server";
import { v } from "@lunora/values";

export const invite = mutation.input({ email: v.string() }).mutation(async ({ ctx, args: { email } }) => {
    await ctx.queues.emailQueue.send({ to: email });
    // or in one call:
    // await ctx.queues.emailQueue.sendBatch([{ body: { to: a } }, { body: { to: b } }]);
});

Push vs pull consumers

By default a queue is a push consumer: Cloudflare delivers batches to the worker's generated queue() handler. To expose the queue to an external HTTP pull consumer instead, omit the handler and set mode: "pull":

export const reports = defineQueue({ mode: "pull", name: "reports" });

lunora dev / lunora deploy reconcile the wrangler config for you: every queue gets a queues.producers[] entry; push queues add a worker consumer, pull queues add a type: "http_pull" consumer (with any batch/retry/dead-letter tuning carried through). The tuning is kept in step after that too: adding a deadLetterQueue or changing maxRetries on an existing queue updates its consumer on the next dev/deploy. Only the options defineQueue sets are written. A consumer field you set by hand and defineQueue leaves unset is not touched.

Removing an option from defineQueue removes it from the consumer on the next dev/deploy. To tell a field it manages from one you set by hand, reconcile records what defineQueue declared on each pass in your package.json, under lunora.queueTuning, next to lunora.crons. Commit it. When an option is dropped, it removes the field only while the field still holds the recorded value. A field you changed by hand to another value is kept, and reconcile warns about it. A hand-set value that already equalled the declared one counts as declared: dropping the option removes it too. A project upgrading from before the record existed starts with nothing recorded, so an option removed in that same change stays deployed once. Delete it from wrangler.jsonc by hand.

lunora deploy --env <name> also retunes the consumers that env.<name> already declares. Each consumer is matched to its defineQueue export through that block's own queues.producers[] binding, since environment queue names usually carry a suffix. It falls back to the declared queue name only when the block maps that queue's binding nowhere. A consumer matched by neither, such as one for a suffixed queue another worker produces, is skipped with a warning naming it. It never adds a producer or consumer to the block, and it never writes dead_letter_queue there, because that value is a queue name and the declared one belongs to the top level. When defineQueue declares a deadLetterQueue and the env consumer has none, the deploy warns. Add it by hand under that environment's queue name.

Studio

Declared queues show in the Studio Queues page (export name, deployed queue name, consumer mode, producer binding, and dead-letter queue), refreshed on every codegen run.