OPCStack uses Cloudflare Queues for async work and Cloudflare Cron Triggers for periodic jobs. Both enter the same Worker deployment through src/index.ts, then dispatch to backend handlers.
Cloudflare Worker
|
+-- queue(batch, env, ctx)
| |
| +-- src/backend/consumers/index.ts
|
+-- scheduled(controller, env, ctx)
|
+-- src/backend/jobs/index.ts
Keep queues and cron boring. A queue message should identify durable state. It should not carry the whole job payload.
Runtime Model
src/index.ts exports three Worker entrypoints:
| Entrypoint | Handler | Purpose |
|---|---|---|
fetch |
API or SvelteKit SSR | User requests and webhooks |
queue |
handleQueue |
Cloudflare Queue batches |
scheduled |
handleScheduled |
Cron Trigger events |
Queues are configured by name. Cron jobs are configured by exact cron expression. Unknown queues or unknown cron expressions are skipped.
Config
Queue and cron config lives in .env.dev and .env.prod.
QUEUE_NAMES=image-generate;tts-generate;video-generate
QUEUE_MAX_CONCURRENCY=
CRONS=*/10 * * * *
Rules:
| Key | Format | Meaning |
|---|---|---|
QUEUE_NAMES |
Semicolon-separated names | Queue resources and Worker producer/consumer bindings |
QUEUE_MAX_CONCURRENCY |
Empty or integer 1..250 |
Optional Cloudflare consumer max concurrency |
CRONS |
Semicolon-separated cron expressions | Worker Cron Triggers |
Empty QUEUE_NAMES means no queue bindings. Empty CRONS means no cron triggers.
prepare-cloudflare.mjs deduplicates queue names after trimming whitespace. It does not validate queue names beyond binding-name conversion, so keep names lowercase and hyphenated.
Generated Bindings
prepare-cloudflare.mjs converts queue names to Worker bindings:
| Queue name | Binding |
|---|---|
image-generate |
Q_IMAGE_GENERATE |
tts-generate |
Q_TTS_GENERATE |
video-generate |
Q_VIDEO_GENERATE |
task-check |
Q_TASK_CHECK |
Binding rule:
Q_<QUEUE_NAME_UPPER_WITH_NON_ALNUM_AS_UNDERSCORE>
Generated wrangler.jsonc contains both producers and consumers:
{
"queues": {
"producers": [
{
"binding": "Q_IMAGE_GENERATE",
"queue": "image-generate"
}
],
"consumers": [
{
"queue": "image-generate"
}
]
}
}
When QUEUE_MAX_CONCURRENCY=1, each consumer entry also gets max_concurrency: 1.
Queue Dispatch
All queue dispatch starts at src/backend/consumers/index.ts.
export async function handleQueue(
batch: MessageBatch<unknown>,
env: Env,
ctx: ExecutionContext
): Promise<void> {
const handler = queueHandlers[batch.queue]
if (!handler) {
return
}
await handler(batch, env, ctx)
}
Registered handlers:
| Queue | Handler |
|---|---|
image-generate |
handleAIImageQueue |
tts-generate |
handleAITTSQueue |
video-generate |
handleAIVideoQueue |
Do not add branching inside src/index.ts. Add a handler to queueHandlers.
Queue Message Rules
Queue messages should be small and durable-state based.
Current AI queues use this shape:
{
taskId: string
userId: string
}
This is intentional:
taskIdpoints to the row that contains prompt, provider, model, references, and output optionsuserIdlets the consumer reopen the correct Tenant Shard DB through Meta DB- retry is idempotent because the task row is the source of truth
Do not put prompts, API keys, provider configs, R2 options, or full user payloads into queue messages.
Existing AI Queues
Image, TTS, and video async tasks are tenant-owned. They are stored in Tenant Shard DB tables:
| Queue | Table | Task creator |
|---|---|---|
image-generate |
ai_image_tasks |
createAIImageTask |
tts-generate |
ai_tts_tasks |
createAITTSTask, createAITTSSourceTask |
video-generate |
ai_video_tasks |
createAIVideoTask |
Flow:
request handler or service
|
+-- insert processing task row in Tenant Shard DB
|
+-- env.Q_*.send({ taskId, userId })
|
+-- consumer opens user's Tenant Shard DB
|
+-- load task row
|
+-- call AI provider
|
+-- update task row to completed or failed
The consumer opens the user's DB with:
const metaDb = getMetaDb(env.META_DB)
const tenant = await createTenantShardAccess(metaDb, env).openUserDb(userId)
That lookup is required because queue events do not have request middleware context.
Retry and Ack Rules
Consumers must explicitly finish each message:
| Action | Meaning |
|---|---|
message.ack() |
Message is done and should not retry |
message.retry({ delaySeconds }) |
Message should be retried later |
Current AI retry behavior:
| Queue | Max attempts | Delay |
|---|---|---|
image-generate |
3 | 10s, 30s |
tts-generate |
3 | 10s, 30s |
video-generate |
3 failure attempts | 10s, 30s |
Video has one extra rule: if the remote SeedDance provider task is still running, the consumer retries the queue message after 30 seconds without incrementing attemptCount.
Consumers ack missing or non-processing tasks. That is not swallowing a bug. It is idempotency: a duplicate message should not re-run a completed task.
Cron Dispatch
All Cron Trigger dispatch starts at src/backend/jobs/index.ts.
export async function handleScheduled(
controller: ScheduledController,
env: Env,
ctx: ExecutionContext
): Promise<void> {
const handler = scheduledHandlers[controller.cron]
if (!handler) {
return
}
await handler(controller, env, ctx)
}
The key is the exact cron expression string from Cloudflare. If CRONS contains */10 * * * *, the handler map must also use */10 * * * *.
Existing Cron Jobs
Current registered job:
| Cron | Job |
|---|---|
*/10 * * * * |
Expire credits and clean old credit transactions on every active tenant shard |
The job:
- Opens Meta DB
- Lists active tenant shard DBs
- Runs
CreditsService.expire({ limit: 20 })on each shard - Runs
CreditsService.cleanupTransactions({ limit: 100 })on each shard - Logs structured job results
The cron reads Credits transaction retention and AI task retention from Meta D1 once per execution. Invalid or missing dynamic configuration fails the job instead of silently using an ENV fallback.
Add A Queue
Add a queue only when work must continue outside the request lifecycle.
Steps:
- Add the queue name to
QUEUE_NAMES - Run
pnpm prepare:cloudflare:dev - Use the generated binding, for example
env.Q_TASK_CHECK.send(...) - Add a message type near the domain that owns the task
- Add a consumer handler under
src/backend/consumers/ - Register it in
queueHandlers - Add unit tests for dispatch, ack, retry, and idempotent duplicate handling
Minimal message shape:
export interface TaskCheckQueueMessage {
taskId: string
userId: string
}
Send:
await env.Q_TASK_CHECK.send({
taskId,
userId
})
Do not create a generic queue framework. The map in queueHandlers is enough.
Add A Cron
Add cron only when the job is truly periodic. If the work is user-triggered or provider-triggered, use a queue or webhook.
Steps:
- Add the cron expression to
CRONS - Add a handler in
scheduledHandlers - Keep the handler short
- For long work, enqueue tasks instead of doing everything inside the scheduled event
- Add unit tests for registered and unknown cron behavior
Example:
export const scheduledHandlers: Record<string, ScheduledJobHandler> = {
'*/10 * * * *': async (controller, env): Promise<void> => {
// existing credits job
},
'0 0 * * *': async (_controller, _env): Promise<void> => {
// daily maintenance
}
}
Cron format:
* * * * *
| | | | |
| | | | +-- day of week
| | | +---- month
| | +------ day of month
| +-------- hour
+---------- minute
Local Testing
For unit tests, call handlers directly:
await handleQueue(batch, env, ctx)
await handleScheduled(controller, env, ctx)
For local scheduled events:
pnpm prepare:cloudflare:dev
pnpm exec wrangler dev --test-scheduled --env-file .wrangler/runtime-secrets.env
For generated bindings and types:
pnpm prepare:cloudflare:dev
pnpm exec wrangler types --config .wrangler/wrangler.types.jsonc --env-file .wrangler/runtime-secrets.env --strict-vars false
Do not run remote E2E to test resource creation. Remote E2E should only call HTTP APIs against an already deployed environment.
Common Mistakes
Adding queue logic to src/index.ts
Keep src/index.ts as the Worker entrypoint only. Register queue handlers in src/backend/consumers/index.ts.
Using comma-separated config
QUEUE_NAMES and CRONS use semicolons.
Putting full job payloads in messages
Put the durable row in D1 and send { taskId, userId }.
Forgetting the generated binding name
task-check becomes Q_TASK_CHECK. Use the generated Env type instead of hand-writing env interfaces.
Doing long work inside cron
Cron should coordinate. Queue should execute long or retryable jobs.
Changing the existing */10 * * * * cron casually
That expression drives credit expiry and cleanup. Changing it changes real business timing.