Skip to main content
Omazy Engineering

Omazy Engineering · Background Tasks

How Omazy runs work in the background.

Every email, export, automation, and analytics drain in Omazy rides one platform: Asynq on Redis, in-process with the API. This is how it's built, why we picked it, what it costs, and where it stops.

RFC 2204 // updated 2026-05-10 // live · production

Section 01 · executive tl;dr

The 60-second read.

For C-level and cross-functional readers. The whole platform in five lines.

STACK

One queue platform: Asynq + Redis, in-process with the API binary. No second broker. No new infra cost.

SCALE

Today ~50 tasks/sec; calibrated headroom to ~5k/s before any architectural change is required.

RELIABILITY

25-retry default with exponential backoff. Three priority queues so auth never waits behind analytics.

AUDITABLE

Every task is inspectable in /admin/jobs. Trace context flows from HTTP request through retries to completion.

COST

$0 incremental. Redis is already deployed. Asynq is open-source. Cloudflare Queues + RabbitMQ remain on the shelf, not on the bill.

Section 02 · how it works

One platform, in one diagram.

Producers call a typed enqueue. Tasks land in Redis under one of three priority queues. The Asynq server in the API binary pulls them and runs the registered handler. Retries, scheduling, and observability come for free.

Producers

HTTP handler

auth, vendor, feedback

Domain service

workspace, automations

Cron tick

@every 30s · hourly

jobs.Client.Enqueue(ctx, kind, payload)

Redis · Asynq broker

pending

Waiting for a worker

active

In-flight on a worker thread

scheduled

Delayed (ProcessIn / cron)

retry

Failed, waiting for next attempt

archived

Exhausted retries — operator triage

completed

Done; held for Retention window

3 queues · critical : default : low = 6 : 3 : 1

Asynq server (in-process)

TaskSpec registry

closed catalogue · MountSpecs at boot

Handler · ctx, *Task

retries · trace · idempotent

Scheduler

cron · @every · ProcessIn

concurrency cap · panics → retry

Section 03 · queues

Three priority queues, one worker pool.

The worker thread always picks the highest-priority pending task. The 6:3:1 ratio means analytics work never starves, but auth never waits behind it.

critical weight · 6

Auth, audit, payment hooks. Short deadlines, low retry budgets, never wait behind anything.

default weight · 3

Most work. Vendor email, workspace invitations, automations runner, AI pipeline. Standard retries.

low weight · 1

Analytics drains, reconciliation sweepers, cache warmers, long exports. Generous retry budget, deferred to off-peak.

Section 04 · why this stack

10 reasons we don't bring in a second broker.

The decision isn't "Asynq beats RabbitMQ." It's that Omazy's workloads — email, exports, automations, analytics drains — fit comfortably inside what Asynq + Redis already provides. Adding a second broker means adding ops cost without adding capability.

  • ·01

    Zero new infrastructure

    Redis is already deployed via Docker on the existing droplet. Asynq rides on the same connection. No second cluster to operate, monitor, or pay for.

  • ·02

    Type-safe producers

    A closed catalogue of TaskKind constants. Producers compile against the registry; typos are rejected at runtime. No silent dropped tasks.

  • ·03

    Three priority queues

    critical (auth, audit) · default (most work) · low (analytics, reconciliation). Weight ratio 6:3:1 — the worker thread always picks the highest-priority pending task.

  • ·04

    Retries + backoff baked in

    Default 25 retries with exponential backoff per task. Permanent failures opt out via SkipRetry. Per-spec overrides for short-lived work.

  • ·05

    Scheduled tasks + cron

    Five-field cron syntax or @every shorthand. One-shot delays via ProcessIn. Already powers the analytics drain (every 30s) and custom-domain checks (hourly).

  • ·06

    Trace context propagation

    OpenTelemetry span IDs ride on the task payload. A single trace covers HTTP → enqueue → handler → retries → completion. No grepping logs.

  • ·07

    Boot-time validation

    Unknown TaskKind, duplicate registration, or a nil handler — the server panics at start. Programmer errors fail loud, not at 3 a.m.

  • ·08

    Admin inspector, no extra UI

    GET /api/v1/admin/jobs/* exposes queue depth, per-task state, retry, delete. Audited. Zero dashboard to deploy.

  • ·09

    Driver-abstraction friendly

    Workers consume Go interfaces (EmailSender, StorageDriver). Swapping Resend → SES touches one factory case, not the workers.

  • ·10

    Test-friendly NopClient

    jobs.NopClient() makes Enqueue a no-op. Unit tests exercise handlers in isolation. Redis is not a test dependency.

Section 05 · adding a task

Three files, in this order.

The boot-time validator panics if any step is skipped. There is no missing-fourth-step. By design.

01 internal/jobs/registry.go

Register the TaskKind

Append a const to the closed catalogue. Naming is <domain>:<verb>. Renames are forbidden — the wire format is what producers persisted to Redis.

TaskFeedbackEmail TaskKind = "feedback:email"
02 internal/<domain>/worker.go

Write the handler

A function with the (ctx, *asynq.Task) error signature. Use injected interfaces (EmailSender, StorageDriver) — never branch on vendor name. Idempotent: Asynq retries.

func (w *Worker) HandleEmail(ctx context.Context, t *asynq.Task) error { ... }
03 cmd/server/wiring/<domain>.go

Wire it at boot

Register a TaskSpec with queue, retry budget, retention, deadline, and the handler. MountSpecs walks the registry on Server.Start. No deployment changes.

jobs.RegisterSpec(jobs.TaskSpec{ Kind: jobs.TaskFeedbackEmail, Queue: jobs.QueueLow, ... })

Section 06 · monitor

Five endpoints, no separate dashboard.

Queue depth, per-task state, retry, delete — all over plain HTTP, all audited, all already running on production. Wire alerts directly off these; don't build a queue UI.

GET /admin/jobs/queues All queues with depth + per-state counts + memory usage
GET /admin/jobs?queue=&state= Paginated list filtered by queue + state
GET /admin/jobs/:queue/:id Full task info: payload, last error, retry count, next-process timestamp
POST /admin/jobs/:queue/:id/retry Force-re-run a task (audited)
DELETE /admin/jobs/:queue/:id Remove a task (audited)

// what to alert on

  • active + pending on critical > 100 sustained 5 min · auth/audit backing up
  • archived count non-zero and growing · a handler is permanently failing — read last_err
  • retry count for one task > 10 · vendor degraded or handler bug; manual review
  • Redis memory > 70% of host · cap or evict; tighten Retention

Section 07 · debug

When something didn't happen.

Walk this in order. Most reports of "X didn't run" are upstream of the queue, not inside it.

  1. Was the task ever enqueued? Check the producer's logs or /admin/jobs?state=completed for the task ID. If it's missing, the producer didn't fire — fix the upstream code path, not the queue.
  2. Did it run? /admin/jobs/:queue/:id shows state + completed_at. archived = exhausted retries; read last_err.
  3. Stuck pending? Pending for minutes means the worker isn't pulling. Check redis: healthy on /health and the asynq server started line in boot logs.
  4. Did the side effect occur? Container up ≠ feature works. For email, verify it landed in the inbox; for storage, retrieve the object. Logs are necessary, not sufficient.

// common failure modes

Symptom Root cause Fix
jobs.Enqueue: unknown TaskKind Producer using a string not in the catalogue Add to registry.go const block + allKnownKinds map
Boot panic: TaskKind already registered Two domains claim the same kind Pick one owner; refactor the second to its own kind
Boot panic: nil Handler TaskSpec.Handler field missing Set Handler: w.HandleX on the spec
Tasks pile up in pending Worker crashed; EMAIL_DRIVER=logged in prod; Redis down Check container health + boot logs; verify driver env var
Task succeeded but feature broken Handler returned nil silently; SkipRetry hid an error Read the handler's success path and the trace; never swallow errors
Same task ran twice Asynq retried after a successful run that exceeded Deadline Make handlers idempotent (always); raise Deadline if needed
Archived count climbing Permanent failure (invalid recipient, deleted parent row) Read last_err; if non-retryable, return SkipRetry early
Lost the trace ID Producer used context.Background() instead of request ctx Always thread the request ctx through the call chain

Section 08 · gaps

Honest limitations.

None are blockers. Each has a workaround that fits inside Asynq + Redis. Listed so we don't kid ourselves.

No standalone worker process

impact Asynq runs in the API binary; a runaway handler can starve API request handling at the concurrency cap.

workaround Conservative concurrency cap; isolate hot tasks to the low queue; split into a dedicated container only when CPU contention is measurable.

No native outbox helper

impact Each domain re-implements the transactional drain pattern.

workaround Three call sites is below the abstraction threshold. Promote to internal/outbox when the fourth lands.

No replay / time-travel

impact Once a task is past its Retention window, it's gone. Can't rerun "everything from yesterday."

workaround For audit-critical tasks, log payload to Postgres separately. Not the queue's job.

Polyglot consumers unsupported

impact Asynq is Go-only. A Python ML worker can't consume the same queue.

workaround All workers live in the middleware today. If a non-Go consumer ever lands, that's a NATS conversation, not Asynq.

Single-region Redis

impact Cross-region failover requires Redis replication we don't have today.

workaround Out of scope until a second region is on the roadmap.

Retention vs Redis memory

impact Long retention on high-volume tasks bloats Redis.

workaround Set Retention per spec (24h is plenty for admin debugging). Don't retain analytics writes.

No streaming / partial results

impact Long exports can't stream progress through the queue.

workaround Handler writes progress to a Postgres row; UI polls that row, not the queue. Pattern in TaskWorkspaceExportBuild.

No CF Workers integration

impact A CF Worker that wants to enqueue must call middleware over HTTPS.

workaround Acceptable today — both CF apps are stateless proxies. Add an internal POST /enqueue if volume grows.

Manual dead-letter triage

impact Archived tasks sit until an operator retries or deletes via admin UI.

workaround Acceptable while volume is low. When archived count justifies it, add a cron that pages on threshold breach.

// anti-patterns

  • Enqueueing inside a Postgres transaction (rollback fires ghost tasks — use an outbox).
  • Branching on vendor name in a handler (violates the driver-abstraction rule).
  • context.Background() in handlers (loses trace context — thread the request ctx).
  • String-typed asynq.NewTask without a TaskKind (typos drop tasks silently).
  • Long retention on high-volume tasks (bloats Redis; retention is for debugging, not auditing).
  • A handler that's not idempotent (double-charge, double-email, double-row on retry).
  • New callers of EnqueueRaw (deprecated; review will reject it).
Companion · load matrix

The honest yardstick.

Adopting Asynq + Redis is the chosen path. This section is the matrix that backs that call: how to load-test the chosen stack, and how it compares against the two alternatives that always come up — "should we use Cloudflare Queues?" / "should we move to RabbitMQ?"

Workloads · the three scenarios that matter

W1

Email fanout

10k tasks in a 60s burst · 1 KB payload · ~200 ms outbound HTTPS

Steady producer → I/O-bound consumer. Tests retry under vendor latency.

W2

Report generation

100 tasks · 4 KB payload · 30-300s of Postgres reads + S3 upload

Long-running tasks. Tests deadline, concurrency, progress reporting.

W3

High-volume analytics

1M tasks over 1 hour (~280/s) · 0.5 KB payload · in-memory + ClickHouse insert

Sustained throughput. Tests broker memory + scheduling fairness.

Calibration ranges per stack

Numbers below are what to expect on the harness above. If your run is wildly outside the band, a config is wrong (not the broker). Cells read across the three workloads.

Asynq + Redis

adopted
W1 (10k burst)W2 (long-running)W3 (sustained 280/s)
Producer p99 enqueue 5–15 ms5–15 ms5–20 ms
End-to-end p99 1–3 shandler + ≤ 1 s1–5 s
Sustained throughput 200–1000/sn/a500–2000/s · redis-bound past ~5k/s
Broker memory ~150 MBminimal1–3 GB depending on Retention
Cost (incremental) $0$0$0
Disqualifies if Redis OOMsConcurrency starves API> 5k/s sustained

Cloudflare Queues

edge-only
W1 (10k burst)W2 (long-running)W3 (sustained 280/s)
Producer p99 enqueue 50–150 mssamesame
End-to-end p99 2–5 sdisqualified · 30s CPU cap2–8 s
Sustained throughput up to 5k/sn/aup to 5k/s
Broker memory n/a (managed)n/an/a
Cost ~$8 for the burstn/a~$60–90/mo at 280/s
Disqualifies if Producer not on CFHandler > 30sNeed priority queues

RabbitMQ

shelf
W1 (10k burst)W2 (long-running)W3 (sustained 280/s)
Producer p99 enqueue 1–5 ms1–5 ms1–5 ms
End-to-end p99 0.5–2 shandler + ≤ 1 s0.5–3 s
Sustained throughput 5k–20k/sn/atrivial
Broker memory ~200–500 MBminimal1–2 GB
Cost $0 self-host · $50–100/mo CloudAMQPsamesame
Disqualifies if No polyglot / routing needsamesame

Decision criteria

When each stack is the right answer.

Default for everything Omazy does today

→ Asynq + Redis

Producer runs inside a Cloudflare Worker

→ Cloudflare Queues

Polyglot consumers · > 5k/s sustained · complex routing

→ RabbitMQ

W1 burst with simple retry semantics

→ Asynq + Redis (zero infra cost)

W2 long-running > 30 s

→ Asynq + Redis OR RabbitMQ (CF disqualified)

W3 sustained ≥ 500/s

→ Asynq + Redis (only at ≥ 5k/s does Rabbit win)

Cost model · ~5M tasks/month

What each option actually costs.

Marginal-cost differences are pennies until ~50M tasks/month. Ops cost dominates the decision: Asynq adds zero ops; RabbitMQ adds non-trivial ops; CF Queues adds vendor dependency. Pick on ops, not cents.

Stack Marginal monthly Notes
Asynq + Redis $0 Already deployed. Concurrency raised in-process.
Cloudflare Queues ~$5–15 $0.40/M ops × ~10M ops at 5M tasks/mo.
RabbitMQ — self-hosted $0 incremental Runs alongside Postgres/Redis. Adds ~512 MB RAM and an Erlang process to operate.
RabbitMQ — managed $50–100 CloudAMQP "Tiger" or equivalent. HA, mirrored, monitored. Worth it the day you go multi-region; not before.

Re-run triggers

When these numbers go stale.

Update the matrix whenever any of the following change by ≥ 2x. Until then, the calibration ranges above are still load-bearing.

  • Sustained tasks/sec at peak goes from ~50/s today to > 100/s sustained.
  • Redis memory at peak goes from ~200 MB today to > 500 MB.
  • A new TaskKind lands with W2-shaped runtime (> 30 s deadline).
  • A second region is added to the deployment topology.
  • A non-Go consumer joins the stack.

// sources: docs/rfc2204-background-tasks.md · docs/background-tasks-load-testing.md

// edit a source doc → open a PR → it ships on merge