Skip to content
Background jobs and webhooks

Background jobs and webhooks

How Neatship's background worker runs queued jobs and scheduled jobs (crons), how to add your own, and how Stripe and Clerk webhooks are received, stored, processed, recovered and cleaned up safely.

Some work shouldn't happen while a person waits for a page: syncing with Stripe, applying a change from Clerk, sending a batch of emails. Some work must happen on a schedule: every night, every 10 minutes. Neatship puts that work in a queue, and a separate process, the worker, runs it. This guide explains how the worker runs, how to add your own job and your own scheduled job, and how incoming webhooks are handled without losing or duplicating anything.

The worker

  • It is a separate process with no web server: packages/app-server/src/queue-worker/queue-worker.ts. When it is ready, it logs Queue worker ready.
  • yarn start runs it next to the API server (npx nx run app-server:worker). It builds into its own dist-worker folder, because the API server's watcher rewrites dist.
  • Jobs wait in Redis, managed by BullMQ (a job queue library). Redis must never evict data under memory pressure, so it runs with maxmemory-policy noeviction (the local Docker setup and the native setup both set it).
  • Several workers can run at once. When a worker is asked to stop (a deploy sends SIGTERM), it lets running jobs finish first.
SettingValueWhere
Queueswebhook-queue, billing-queue, cron-queue (scheduled jobs)MESSAGE_QUEUE_NAMES
Jobs of one queue running at the same time, per worker5MESSAGE_QUEUE_WORKER_CONCURRENCY
Retries after the first attempt3 by default, 10 for webhook jobsMESSAGE_QUEUE_DEFAULT_RETRY_LIMIT, WEBHOOK_EVENT_PROCESSING_RETRY_LIMIT
Wait between retriesdoubles each time from 5 seconds: 5 s, 10 s, 20 s, 40 s…MESSAGE_QUEUE_RETRY_BASE_DELAY_IN_MILLISECONDS
Finished jobs kept in Rediscompleted: 4 hours or 1,000 jobs; failed: 7 days or 1,000 jobsMESSAGE_QUEUE_JOB_RETENTION

The constants live in packages/app-server/src/engine/core-modules/message-queue/constants/.

MESSAGE_QUEUE_DRIVER in packages/app-server/.env picks how jobs run:

ValueBehavior
bullmq (default)Jobs wait in Redis until the worker runs them, with retries
syncA job runs at once, inside the request that queues it. For integration tests only; refused when NODE_ENV=production

Warning: If the worker isn't running, jobs wait in Redis and nothing happens: a payment doesn't change the plan, a Clerk change doesn't reach your data. They run as soon as a worker starts.

Add a job

Say the Clients feature needs to send a weekly summary for one workspace.

1. Write the job class

Create jobs/send-client-summary.job.ts in your feature module:

ts
export type SendClientSummaryJobData = { workspaceId: string };

// Queued once per workspace. Safe to run twice: it reads fresh data.
@Processor({ queueName: 'client-queue' })
export class SendClientSummaryJob {
  constructor(private readonly clientSummaryService: ClientSummaryService) {}

  @Process(SendClientSummaryJob.name)
  async handle({ workspaceId }: SendClientSummaryJobData): Promise<void> {
    await this.clientSummaryService.sendSummary(workspaceId);
  }
}
  • @Processor names the queue. @Process(SendClientSummaryJob.name) names the job after its class, so the name can't be misspelled between the code that queues it and the code that runs it.
  • The data type is exported next to the class.

The existing jobs are good models: ApplyWorkspaceMemberLimitJob in packages/app-server/src/engine/core-modules/billing/jobs/, and the two webhook jobs in packages/app-server/src/engine/core-modules/webhook/jobs/.

2. Register it

  1. The queue. Use an existing queue, or add a name to MESSAGE_QUEUE_NAMES (packages/app-server/src/engine/core-modules/message-queue/constants/message-queue-names.constant.ts), here 'client-queue'. A queue groups jobs that are retried and scaled together.
  2. The module. List the job class in the providers of your feature module. The message queue finds every @Processor class there at start-up; there is nothing else to register.
  3. The worker. Import your feature module in packages/app-server/src/queue-worker/queue-worker.module.ts. The worker only loads the modules listed there, so a job in a module that isn't imported never runs.

3. Queue it

Inject the queue with @InjectMessageQueue and add a job:

ts
constructor(
  @InjectMessageQueue('client-queue')
  private readonly clientQueue: MessageQueueService,
) {}

async requestSummary(workspaceId: string): Promise<void> {
  await this.clientQueue.add<SendClientSummaryJobData>(
    SendClientSummaryJob.name,
    { workspaceId },
  );
}

add accepts options: { retryLimit: 5 } changes the number of retries, { delayInMilliseconds: 60000 } delays the first run, and { deduplicationId: 'some-id' } queues a job once even when many requests ask for it (while a job with that id waits, runs or retries, adding another one does nothing).

4. Follow the three rules

  1. Small payloads, made of ids. The job loads fresh data itself. Data copied into the payload is out of date by the time the job runs.
  2. Safe to run twice, and in parallel with others. A retry after a crash, or two workers, can run the same work twice. Use unique indexes, "insert if absent" writes, or a lock, so the second run changes nothing.
  3. Throw to retry. An error thrown by handle makes the queue retry the job with its backoff. Catch only what you can really handle.

5. Test it

  • In integration tests, the sync driver runs the job inside the request that queues it, so its effects are there when the response comes back. No Redis is needed.
  • In a service unit test, replace the queue with a mock and check what was queued. billing-subscription-sync.service.spec.ts does this.

Add a scheduled job (cron)

A cron is a job that runs on a schedule instead of being queued by a request: every night, every 10 minutes. In Neatship a cron is an ordinary job of the cron-queue, so it gets the same retries and failure reports as any job, and it runs once per tick, however many workers you run.

Say the Clients feature must send its weekly summaries every Monday at 8:00.

1. Write the cron job

Create crons/jobs/send-weekly-summaries.cron.job.ts in your feature module:

ts
// Every Monday at 08:00 UTC.
export const SEND_WEEKLY_SUMMARIES_CRON_PATTERN = '0 8 * * 1';

// Queues one summary per workspace: the slow work runs in many small jobs.
@CronProcessor({ pattern: SEND_WEEKLY_SUMMARIES_CRON_PATTERN })
export class SendWeeklySummariesCronJob {
  constructor(private readonly clientSummaryService: ClientSummaryService) {}

  @Process(SendWeeklySummariesCronJob.name)
  async handle(): Promise<void> {
    await this.clientSummaryService.queueSummaryForEveryWorkspace();
  }
}
  • @CronProcessor replaces @Processor: it puts the job on the cron-queue and gives its schedule.
  • The pattern has 5 fields: minute, hour, day of the month, month, day of the week. It is read in UTC, whatever your server's time zone. */10 * * * * is every 10 minutes; 30 3 * * * is every day at 03:30. crontab.guru explains any pattern.
  • A cron that has work to do for each workspace should queue one small job per workspace (like SendClientSummaryJob above) rather than do everything itself: each workspace then gets its own retries.

2. Register it

List the class in the providers of your feature module, and make sure the worker imports that module (packages/app-server/src/queue-worker/queue-worker.module.ts). That's all: the worker schedules every cron when it starts. Its log says so:

text
Cron SendWeeklySummariesCronJob scheduled: 0 8 * * 1 (UTC).

Change the pattern, or delete the class, and the next start of the worker (your next deploy) updates or removes the schedule. There is no command to run.

3. Test it

In integration tests, nothing is scheduled. Run the cron by adding its job to the cron-queue; the sync driver runs it at once:

ts
await app
  .get<MessageQueueService>(getMessageQueueToken('cron-queue'))
  .add(SendWeeklySummariesCronJob.name, {});

packages/app-server/test/integration/webhook/webhook-inbox-crons.integration-spec.ts tests the two crons of the starter this way.

The crons of the starter

CronWhenWhat it does
WebhookInboxRecoveryCronJobevery 10 minutesQueues again the webhook events nobody processed, and alerts you about the ones that keep failing (see "When an event keeps failing" below)
WebhookInboxRetentionCronJobevery day at 03:30 UTCDeletes the processed webhook events older than WEBHOOK_EVENT_RETENTION_DAYS days (default 30)

Both live in packages/app-server/src/engine/core-modules/webhook/crons/jobs/.

The webhook inbox

A webhook is an HTTP request that a service like Stripe or Clerk sends to your server when something changes on its side. Neatship receives two: POST /webhooks/stripe and POST /webhooks/clerk.

Webhooks have three traps: anyone can send a fake one, the same event can arrive several times, and events can arrive late or out of order. Here is how each request is handled:

  1. The signature is checked first. These routes don't use the session middleware. Their authentication guard is a signature guard (StripeWebhookSignatureGuard, ClerkWebhookSignatureGuard), which checks the provider's signature against the exact bytes received (the raw body), with the secret in STRIPE_WEBHOOK_SECRET or CLERK_WEBHOOK_SIGNING_SECRET. A missing or wrong signature gets a 400 answer and nothing is stored. A provider that is switched off (BILLING_PROVIDER=none, or no Clerk signing secret) gets a 404. Signed requests are then rate-limited per IP address (600 per minute on each route): beyond that, the answer is 429 and the provider sends the event again later.
  2. Unhandled event types are acknowledged and dropped. Only the events the code handles are stored: less noise, and no personal data kept for nothing.
  3. The event is written to the inbox, the webhookEvent table, which has one unique key per provider and event id (Stripe's evt_... id, or Clerk's message id, which stays the same on every retry). A second delivery of the same event creates no second row.
  4. A job is queued with only the inbox row id, and the server answers 200 { received: true } at once. Providers expect a fast answer and retry when they don't get one. While a job waits, runs or retries for an event, queueing that event again adds nothing.
  5. The worker processes the event. It skips an event already processed, counts the attempt, runs the handler, and on success sets processedAt. On failure it records the error in lastError and throws, so the queue retries: 10 retries over about 85 minutes, enough for a short outage. Longer failures: see "When an event keeps failing" below.

If the first queueing failed (Redis was down), the recovery cron queues the event again within 10 minutes, and so does the provider's next delivery of the same event, because the inbox row still has no processedAt.

Inbox columnMeaning
providerstripe or clerk
eventIdthe provider's id for the event
typefor example customer.subscription.updated
payloadthe verified event, as sent
receivedAtwhen it arrived
processedAtwhen a job processed it successfully; empty until then
attemptshow many times a job tried
lastErrorthe error of the last failed attempt; cleared on success
abandonedAtwhen the recovery cron set the event aside after too many failed attempts; empty otherwise

Idempotency and ordering

Idempotent means that doing something twice gives the same result as doing it once. Every handler is written so that duplicates, delays and disorder do no harm:

  • Duplicates are stopped twice: by the inbox's unique key, and by the processedAt check in the job.
  • Stripe events only tell the job which customer changed (or which dispute, refund or fraud warning). The job doesn't trust the event's content: it asks Stripe for that customer's subscriptions (or that object) as they are now, and writes that. A late or repeated event therefore writes the same, current state. The workspace row is locked during the sync, so two syncs of one workspace run one after the other.
  • Clerk events carry an updated_at time. A change is applied only if it is newer than the last one copied (clerkUpdatedAt on users, workspaces and members). Deletions are soft deletes: the row is marked deleted, not erased.

Apply the same ideas to your own jobs: read the current state instead of trusting a message, and make every write safe to repeat.

Look inside

The worker's output (in the yarn start terminal) logs each failed job, with its name, id, queue and attempt number, and, at start, each cron it scheduled. In production these lines are JSON, and every line written during a job names the job (jobName, jobId): see "Logs" in Deploy.

To list the latest webhook events on your machine (a read-only query):

bash
bash packages/app-utils/dev-compose.sh exec postgres \
  psql -U app -d default \
  -c 'SELECT provider, type, "receivedAt", "processedAt", attempts, "abandonedAt", "lastError" FROM core."webhookEvent" ORDER BY "receivedAt" DESC LIMIT 20'

If you run Postgres without Docker, use psql -h localhost -p 5433 -U app -d default -c '...' with the same query.

To produce Stripe events locally, keep stripe listen running (see Billing) and use Checkout, the customer portal, or stripe trigger <event>. An event about a Stripe customer that isn't linked to any workspace (typical with stripe trigger) is processed and logged as "nothing to sync".

When an event keeps failing

Nothing is lost when processing fails. Here is what happens to an event whose processing fails, for example because Stripe was unreachable for an hour or because of a bug in a handler:

  1. Its job retries 10 times over about 85 minutes. Every failed attempt is reported as a server fault (in the logs, and in Sentry when error alerts are on: see Deploy).

  2. The recovery cron queues it again. Every 10 minutes, it looks for events still unprocessed 15 minutes after they arrived, and queues them again. This also catches an event whose first queueing failed (Redis was down) or whose job Redis lost. An event whose job is still waiting or retrying is left to that job: the same event is never processed twice at the same time.

  3. After 33 failed attempts (three rounds, about 5 hours), the event is set aside: the cron stops retrying it, fills abandonedAt, and sends one alert naming the event and its last error. Retrying can't help anymore: something must be fixed.

  4. You fix the cause and deploy. Then give the set-aside events their attempts back:

    bash
    # On the production server
    bash packages/app-utils/prod-compose.sh run --rm webhook-retry-abandoned
    # On your machine
    npx nx run app-server:webhook:retry-abandoned

    The command prints how many events it reopened. The worker queues them again within 10 minutes; follow them in the worker's logs. Running it again changes nothing. Resending one event from the Stripe or Clerk dashboard does the same for that event.

Stripe events are safe to process late: the job reads the customer's current state from Stripe, so an old event never undoes a newer change.

Old events are deleted

The inbox holds personal data (names and emails in Clerk events, customer ids in Stripe events), and nothing reads an event after it was processed. Every day at 03:30 UTC, the retention cron deletes the processed events older than WEBHOOK_EVENT_RETENTION_DAYS days (default 30, in packages/app-server/.env or .env.production). Unprocessed events, set aside or not, are never deleted: they still need you.