Writing a background job
This page explains the job mechanism in server/src/worker.ts so you can add a queue or debug one. The list of existing queues, their crons and what they do is in Background jobs; how to run the worker in production is in Worker and scaling.
One helper declares a queue, its processor and its cron
Section titled “One helper declares a queue, its processor and its cron”worker.ts creates a single PgBoss instance, boss, on DATABASE_URL, logs its error events as worker error, and exports createQueue:
export const downloadCreateArchiveQueue = createQueue<{ downloadId: string }>({ name: 'download/create-archive', processor: (data) => dataSource.transaction(async (em) => { const download = await em.getRepository(Download).findOneBy({ id: data.downloadId }) if (!download) { return } await createDownloadArchive({ em, download }) await mailerDownloadReadyQueue.push({ downloadId: download.id }) }), workerOptions: { batchSize: 1 },})| Option | Meaning |
|---|---|
name |
The pg-boss queue name, by convention <area>/<verb-object> (asset/update-content, mailer/invitation) |
processor |
(data: T) => any, called once per job with the job’s payload |
cron |
Optional. When set, boss.schedule(name, cron) makes pg-boss enqueue an empty job on that schedule; the payload type is then void |
workerOptions |
Optional PgBoss.WorkOptions; the code uses batchSize (10 for asset/update-content, 1 for collection/synchronization and download/create-archive) |
createQueue returns an object with two methods and nothing else:
push(data, { uniqueKey? })callsboss.sendwithretryBackoff: trueandsingletonKey: uniqueKey.bulkPush([{ data, uniqueKey? }])callsboss.insertwith the same two options on every entry, in one round trip.upsertFolderuses it to queue onecollection/synchronizationper affected collection, and the integrity check to re-queue every outdated file.
retryBackoff: true makes pg-boss retry a failed job with exponential backoff, using pg-boss’s default retry limit since none is set here. A singletonKey deduplicates: pg-boss refuses a second job with the same key while one is still queued or active.
Registration happens at module load, activation in startWorker
Section titled “Registration happens at module load, activation in startWorker”createQueue does not talk to the database when called. It pushes an initializer into the module-level workerInitializers array and returns the push / bulkPush pair immediately. startWorker() then runs boss.start() and every initializer in declaration order; each one calls boss.createQueue(name), boss.schedule(name, cron) when a cron is set, and boss.work(name, workerOptions, handler).
Two consequences for a new queue:
- Declare it as an
export constinworker.tsand import that constant where you push (services/asset.tsimportsassetUpdateContentQueueandcollectionSynchronizationQueue;trpc/router/user.tsimports the mailer queues). A queue declared in another module would never be registered, becausestartWorkeronly knows about calls made whileworker.tswas loaded. - Pushing works in any process, worker or not:
pushonly needsbossto be started. The API process startsbossinsidestartWorkerwhenENABLE_WORKER=true; the CLI starts it explicitly withboss.start(). A process that has neither cannot push.
worker.ts imports the services, and services/asset.ts imports worker.ts back. This circular import works because the queue constants are only dereferenced inside functions, after both modules have loaded; keep new code in the same shape and do not call push at module top level.
The handler wraps every job in the same error logging
Section titled “The handler wraps every job in the same error logging”The boss.work handler receives a batch of jobs and runs processor(job.data) for each with Promise.all. A rejected processor is logged as:
error: job {"status":"failed","queue":"<name>","jobId":"<uuid>","error":"<message>","stacktrace":...}then rethrown so pg-boss marks the job failed and schedules the retry. The wrapper reads error.stacktrace, a property Error does not define, so no stack trace reaches the log; log with logger from env.ts inside the processor when you need more than the message. Anything the processor resolves with is ignored.
Use a transaction when a job writes several rows
Section titled “Use a transaction when a job writes several rows”Processors that touch more than one table run inside dataSource.transaction(async (em) => ...) and pass em down to the service (synchronizeCollection(em, id), createDownloadArchive({ em, download })). The services accept an EntityManager for that reason; follow the same signature in a new service so the job can pass its transaction. Processors that do a single repository write (mailer/*, asset/update-content) use dataSource.getRepository directly.
Because the transaction spans the whole processor, an exception rolls back every row the job wrote before pg-boss retries it. Side effects outside Postgres (an S3 upload, an email) are not rolled back; do them last, after the database work, as download/create-archive does by pushing mailer/download-ready at the end.
Cron jobs receive no payload
Section titled “Cron jobs receive no payload”A scheduled queue is createQueue<void> with processor: () => someService(). pg-boss stores the schedule and enqueues a job at each tick; the same processor also runs for jobs pushed by hand with queue.push(undefined). The four crons in the code are listed in Background jobs.
Running a job by hand
Section titled “Running a job by hand”The CLI in server/src/cli.ts is the pattern: it initialises dataSource, calls boss.start(), and defines commands with commander. check-integrity calls the service directly (systemService.integrityCheck()) rather than pushing a job, so it runs in the CLI process and exits with process.exit(0).
To run a processor or push a job during development, add a command following check-integrity:
const pushCmd = new Command('push-update-content')pushCmd.argument('<assetFileId>')pushCmd.action(async (assetFileId: string) => { await assetUpdateContentQueue.push({ assetFileId }) process.exit(0)})program.addCommand(pushCmd)then npm run cli -- push-update-content <id> from server/. The job runs in whichever process has the worker enabled, typically the npm run dev server. To run the processor itself in the CLI process instead, call the service function directly as check-integrity does; dataSource.initialize() has already run by then.
SELECT name, state, count(*) FROM pgboss.job GROUP BY 1, 2; shows what is queued, active, completed or failed.
Checklist for a new queue
Section titled “Checklist for a new queue”export const myQueue = createQueue<Payload>({ name: 'area/verb-object', processor, cron?, workerOptions? })inserver/src/worker.ts.- Import
myQueuewhere it is pushed and callpushorbulkPush, with auniqueKeywhen the same job must not be queued twice. - Add a row to the table in Background jobs;
scripts/check-docs.shfails when a queue name inworker.tsis missing there.