Background work runs on BullMQ over Redis. The queue factory is @perform/queue, and every worker lives in apps/core-api/src/workers/. Workers start in server.ts in the same process as the HTTP server. There is no separate worker deployment.

The queue factory

  • One Redis connection is shared by every queue, every worker and all other Redis use in the app. It is created with maxRetriesPerRequest: null, which BullMQ requires.
  • Every queue and worker gets the key prefix perform.
  • Queue names are constants in PerformQueue.

Queue names

PerformQueue in packages/queue/src/index.ts declares sixteen names. Eleven are used.
Two names no longer describe their job. automation-evaluate runs the task generator, and check-in-reminder runs the check-in form dispatcher, the form reminder engine and the pending form notifier. Renaming a queue would orphan its repeatable job in Redis, which is why the names stayed.
WhatsApp sends, meal analysis and PDF generation happen inside the request today, not on a queue.

Two kinds of worker

Schedulers register a repeatable job on their own queue and process it. They also enqueue one boot job at startup so a deploy does not wait a full interval for the first run. Each job scans the database for work that is due. Event workers process jobs that a request or another job enqueued, with a payload that names one unit of work. Every scheduler follows the same template:
  • The fixed jobId keeps a restart from registering a second repeatable.
  • The long lock stops BullMQ from marking a slow full scan as stalled and running it again in a loop.
  • Every scheduled job is idempotent. A boot run racing a repeat, or two instances, must not double send.
Schedulers also share a transient error filter. Prisma codes P1001, P1002, P1008, P1017, a PrismaClientInitializationError, and messages containing missing lock, could not renew lock, connection terminated or connection closed are logged at warn as “skipped: infrastructure unavailable”. Anything else is an error.

Scheduled jobs

Task generator

Scans every studio’s roster and evaluates each task automation trigger. It creates automatic tasks for newly matched entities and auto-closes tasks whose condition no longer holds. A task carries an autoKey (for example <type>:<entityId>) so the same entity never gets the same task twice. Triggers that have a flow attached do not write tasks directly. The generator hands the newly matched entities to a FlowRunner, which the scheduler builds as startFlowRuns(ctx, args, { enqueue }). That creates AutomationFlowRun rows and enqueues them on flow-run. The result is logged as task generator run with created and closed counts.

Check-in form dispatcher, reminders and pending forms

The worker branches on the job name. A form-reminders job runs only runFormReminderDispatcher: it sends the configured FORM_REMINDERS messages whose send time has arrived, per trainee, in the trainee’s time zone. REMINDER_TICK_MINUTES is 5. An hourly or boot job runs four steps in order:
  1. runUpdateFormDispatcher: assigns the recurring check-in form to each client whose cadence says it is due today. Candidates are active or churn risk clients, not frozen, inside their subscription dates, with an active CHECK_IN update form, in a studio that is not archived. A client is only sent to after SEND_HOUR (8) in their own time zone, and lastUpdateFormSentAt guards against a second send the same day. Each send raises a form sent event.
  2. runFormReminderDispatcher: the same reminder pass as the 5 minute job.
  3. listDueFormSentEvents and raiseFormSent: fires the FORM_SENT automation trigger for assignments that became visible. A failure for one assignment is logged and the loop continues.
  4. notifyPendingForms: pushes a “form waiting” notification for assignments that have not been announced yet, scanning a 14 day window in batches of 200. This is also what announces a form assigned silently through the automation API.

Subscription reconcile

Finds clients with an ACTIVE subscription whose endsOn has passed or a SCHEDULED one whose startedOn has arrived, and runs reconcile(repo, tx, clientId, now) for each inside its own transaction. Reconcile promotes scheduled subscriptions to active, expires lapsed ones and keeps the client’s mirrored plan fields and status in sync. These are trainee subscriptions to a studio’s plan, not the studio’s own Polar subscription.

Plan release sweep

A plan assigned from an intake review is held (PAUSED with Program.releasePending) until the coach sends it. If the coach schedules it, releaseAt is set. The sweep calls programs.releaseDue(now) to activate plans whose time arrived and notify the trainee. An atomic claim in programs.service means overlapping ticks release each plan exactly once. A plan whose activation failed is held again and retried on the next tick. Activation fires the FIRST_PLAN subscription start event.

Flow sweep

The durable half of every wait in an automation flow. It re-enqueues up to 200 WAITING runs whose resumeAt has passed (job id flow-<runId>-sweep, so a double sweep collapses), and marks RUNNING runs untouched for 30 minutes as FAILED with error stalled.

Demo activity

Fabricates recent activity (workouts, nutrition days, metrics, form responses, showcase tasks) for studios in demo mode so a sales demo account always looks alive. It loads demo studios with loadDemoStudios, backfills 2 days by default, and does nothing when no studio is in demo mode.

Event workers

Flow run

Each job calls advanceRun(ctx, runId) in modules/task-automations/flow-runtime.ts, which walks one AutomationFlowRun through its node tree:
  • One run row exists per studio, trigger and entity. The unique constraint dedups starts, so an entity that already ran a trigger’s flow never runs it again.
  • position is a work stack of frames. A wait node parks the run as WAITING with resumeAt and enqueues a delayed job. The delayed job and the hourly sweep race to resume it, and the loser no-ops on the entry guard.
  • Every effect node inserts a FlowNodeExecution row before acting, so a retry or a resume that re-walks old ground skips nodes that already ran.
  • The flow definition is re-read from the automation row on every advance. A coach’s edit applies to the remaining steps, and a deleted branch ends the run.
When a job exhausts its attempts the worker marks a run still in RUNNING as FAILED with the flow run was interrupted. Without that, a job killed outside the process would leave the run in RUNNING forever. Redis is optional for readiness, so a delayed job can be lost. Postgres resumeAt is the source of truth and the sweep is what makes waits durable.

Automation hook delivery

Delivers webhooks off the request path with deliverHook (10 second timeout). The trainee’s form submission is already saved when a delivery is queued, so nothing here may block or fail it. A receiver that answers 4xx, or a hook that is gone or switched off, is dropped and not retried. Other failures throw so the queue retries.

Push receipt check

Reads Expo’s delivery receipts for pushes it accepted. See Notifications and push. The enqueue is deliberately not awaited: a coach’s send must never wait on Redis.

Agent reply

Turns a debounced burst of WhatsApp messages from a coach into one reply by calling service.replyToPending. Concurrency is 1 so two replies to the same coach never interleave. There are no retries because the service sends its own “try again” fallback and a retry would answer twice.

AutoFit import

Runs one migration from AutoFit stage by stage: PULLING, WRITING_LIBRARY, IMPORTING, WRITING_PROGRAMS, DONE. A run only walks the stages its include flags need. Progress and counters live on the AutofitImport row, which the web wizard polls. Safety properties built into the worker:
  • Every write goes straight through Prisma with no notifier, so an import can never send a WhatsApp message or a push.
  • Raw pulls are encrypted before they are written to R2 and the snapshot prefix is deleted when the run is done.
  • Mirrored media URLs are treated as hostile: https only, public hosts only, a streamed size cap, and the URL re-validated after redirects.
  • Rows are keyed by the AutoFit id they carry, so a re-run only updates rows this import wrote.
  • A stage aborts after 25 individual client errors.
If the job dies outside the process, the failed handler marks a run still in a running status as FAILED with the import process was interrupted and records the stage.

In-process jobs that are not on BullMQ

Two AI features run long work without a queue:
  • Plan import (modules/plan-import/plan-import.service.ts)
  • Forms AI (modules/forms-ai/forms-ai.service.ts)
analyze answers 202 at once and continues with setImmediate. Progress is written to one Redis key (perform:plan-import:<jobId> or perform:forms-ai:<jobId>) with a one hour TTL, and the client polls GET /analyze/:jobId. The stored state carries the owning studioId and userId, and a mismatch answers 404. A job id is not a capability. A restart mid job loses it, and the client sees the key expire.

Retention

Completed jobs are removed (removeOnComplete: true). Failed jobs are kept up to a count: 50 for most queues, 100 for hook deliveries.

Adding a queue

1

Name it

Add a constant to PerformQueue in packages/queue/src/index.ts and rebuild the package (pnpm --filter @perform/queue build), or core-api type checks against a stale dist.
2

Write the worker

Create workers/<name>-worker.ts exporting start<Name>Worker(ctx: AppContext) that returns the BullMQ worker. If something else enqueues, export a producer too, create<Name>Queue(ctx), with the job options in one place.
3

Make the job idempotent

A job can run twice. Use a deterministic jobId to collapse duplicates, a unique constraint to dedup starts, or a row claimed with an atomic update.
4

Keep durable state in Postgres

A delayed job is only as durable as Redis. If losing it would lose work, store the due time in a column and add a sweep.
5

Start and stop it

Call the start function in server.ts and add worker.close() to both the shutdown handler and the SIGUSR2 handler.
6

Handle the out of process death

If the job moves a row through running states, add a failed handler that marks the row failed when attempts are exhausted.

Running a job by hand

backend/scripts/run-task-generator.ts and apps/core-api/src/scripts/run-demo-activity.ts call the job logic directly without the queue. The job functions take a PrismaClient and optional now, which makes them callable from a script or a test.