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.
Section 01 · executive tl;dr
The 60-second read.
For C-level and cross-functional readers. The whole platform in five lines.
One queue platform: Asynq + Redis, in-process with the API binary. No second broker. No new infra cost.
Today ~50 tasks/sec; calibrated headroom to ~5k/s before any architectural change is required.
25-retry default with exponential backoff. Three priority queues so auth never waits behind analytics.
Every task is inspectable in /admin/jobs. Trace context flows from HTTP request through retries to completion.
$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.
Auth, audit, payment hooks. Short deadlines, low retry budgets, never wait behind anything.
Most work. Vendor email, workspace invitations, automations runner, AI pipeline. Standard retries.
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.
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" 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 { ... } 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.
/admin/jobs/queues All queues with depth + per-state counts + memory usage /admin/jobs?queue=&state= Paginated list filtered by queue + state /admin/jobs/:queue/:id Full task info: payload, last error, retry count, next-process timestamp /admin/jobs/:queue/:id/retry Force-re-run a task (audited) /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.
- 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.
- Did it run? /admin/jobs/:queue/:id shows state + completed_at. archived = exhausted retries; read last_err.
- 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.
- 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).
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
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.
Report generation
100 tasks · 4 KB payload · 30-300s of Postgres reads + S3 upload
Long-running tasks. Tests deadline, concurrency, progress reporting.
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 ms | 5–15 ms | 5–20 ms |
| End-to-end p99 | 1–3 s | handler + ≤ 1 s | 1–5 s |
| Sustained throughput | 200–1000/s | n/a | 500–2000/s · redis-bound past ~5k/s |
| Broker memory | ~150 MB | minimal | 1–3 GB depending on Retention |
| Cost (incremental) | $0 | $0 | $0 |
| Disqualifies if | Redis OOMs | Concurrency starves API | > 5k/s sustained |
Cloudflare Queues
edge-only| W1 (10k burst) | W2 (long-running) | W3 (sustained 280/s) | |
|---|---|---|---|
| Producer p99 enqueue | 50–150 ms | same | same |
| End-to-end p99 | 2–5 s | disqualified · 30s CPU cap | 2–8 s |
| Sustained throughput | up to 5k/s | n/a | up to 5k/s |
| Broker memory | n/a (managed) | n/a | n/a |
| Cost | ~$8 for the burst | n/a | ~$60–90/mo at 280/s |
| Disqualifies if | Producer not on CF | Handler > 30s | Need priority queues |
RabbitMQ
shelf| W1 (10k burst) | W2 (long-running) | W3 (sustained 280/s) | |
|---|---|---|---|
| Producer p99 enqueue | 1–5 ms | 1–5 ms | 1–5 ms |
| End-to-end p99 | 0.5–2 s | handler + ≤ 1 s | 0.5–3 s |
| Sustained throughput | 5k–20k/s | n/a | trivial |
| Broker memory | ~200–500 MB | minimal | 1–2 GB |
| Cost | $0 self-host · $50–100/mo CloudAMQP | same | same |
| Disqualifies if | No polyglot / routing need | same | same |
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