@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.Two kinds of worker
Schedulers register a repeatable job on their own queue and process it. They also enqueue oneboot 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
jobIdkeeps 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.
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:
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 activeCHECK_INupdate form, in a studio that is not archived. A client is only sent to afterSEND_HOUR(8) in their own time zone, andlastUpdateFormSentAtguards against a second send the same day. Each send raises a form sent event.runFormReminderDispatcher: the same reminder pass as the 5 minute job.listDueFormSentEventsandraiseFormSent: fires theFORM_SENTautomation trigger for assignments that became visible. A failure for one assignment is logged and the loop continues.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.
positionis a work stack of frames. A wait node parks the run asWAITINGwithresumeAtand 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
FlowNodeExecutionrow 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.
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.
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.