---
title: "Background Jobs and Queues"
description: "Background Jobs and Queues manages asynchronous task execution, scheduled background routines, and reliable message transport across the platform. It solves the challenge of handling long-running o..."
last_updated: "2026-10-05T05:07:35.174967+00:00"
canonical_url: "https://www.doc0.dev/docs/934e554a-e6a1-476f-bb2f-23e62d86c3fd/technical/core-architecture/background-jobs-and-queues"
---

<details>
<summary>Relevant source files</summary>

The following files were used as context for generating this wiki page:

- [apps/web/lib/jobs/outbox.ts](https://github.com/blade47/dub/blob/HEAD/apps/web/lib/jobs/outbox.ts)
- [apps/web/lib/jobs/index.ts](https://github.com/blade47/dub/blob/HEAD/apps/web/lib/jobs/index.ts)
- [apps/web/app/api/jobs/process/jobName/route.ts](https://github.com/blade47/dub/blob/HEAD/apps/web/app/api/jobs/process/%5BjobName%5D/route.ts)
- [apps/web/app/ee/api/cron/campaigns/broadcast/route.ts](https://github.com/blade47/dub/blob/HEAD/apps/web/app/(ee)/api/cron/campaigns/broadcast/route.ts)
- [apps/web/app/ee/api/cron/campaigns/queue-scheduled/route.ts](https://github.com/blade47/dub/blob/HEAD/apps/web/app/(ee)/api/cron/campaigns/queue-scheduled/route.ts)
- [apps/web/lib/jobs/send-jobs.ts](https://github.com/blade47/dub/blob/HEAD/apps/web/lib/jobs/send-jobs.ts)
- [apps/web/app/ee/api/cron/framer/backfill-leads-batch/route.ts](https://github.com/blade47/dub/blob/HEAD/apps/web/app/(ee)/api/cron/framer/backfill-leads-batch/route.ts)
- [apps/web/lib/cron/enqueue-batch-jobs.ts](https://github.com/blade47/dub/blob/HEAD/apps/web/lib/cron/enqueue-batch-jobs.ts)
- [apps/web/app/ee/api/cron/program-application-reminder/route.ts](https://github.com/blade47/dub/blob/HEAD/apps/web/app/(ee)/api/cron/program-application-reminder/route.ts)
- [apps/web/app/ee/api/cron/streams/update-click-stats/route.ts](https://github.com/blade47/dub/blob/HEAD/apps/web/app/(ee)/api/cron/streams/update-click-stats/route.ts)
- [apps/web/lib/api/rewards/queue-reward-processing.ts](https://github.com/blade47/dub/blob/HEAD/apps/web/lib/api/rewards/queue-reward-processing.ts)
- [apps/web/app/ee/api/cron/rewards/queue-custom-commissions/route.ts](https://github.com/blade47/dub/blob/HEAD/apps/web/app/(ee)/api/cron/rewards/queue-custom-commissions/route.ts)
- [apps/web/app/ee/api/cron/queue/retry/route.ts](https://github.com/blade47/dub/blob/HEAD/apps/web/app/(ee)/api/cron/queue/retry/route.ts)
- [apps/web/app/ee/api/cron/sitemaps/queue/route.ts](https://github.com/blade47/dub/blob/HEAD/apps/web/app/(ee)/api/cron/sitemaps/queue/route.ts)
- [apps/web/app/ee/api/cron/bounties/queue-sync-social-metrics/route.ts](https://github.com/blade47/dub/blob/HEAD/apps/web/app/(ee)/api/cron/bounties/queue-sync-social-metrics/route.ts)
- [apps/web/prisma/schema/job.prisma](https://github.com/blade47/dub/blob/HEAD/apps/web/prisma/schema/job.prisma)
- [apps/web/lib/cron/with-cron.ts](https://github.com/blade47/dub/blob/HEAD/apps/web/lib/cron/with-cron.ts)
- [apps/web/app/ee/api/cron/trial-emails/route.ts](https://github.com/blade47/dub/blob/HEAD/apps/web/app/(ee)/api/cron/trial-emails/route.ts)
- [apps/web/app/ee/api/cron/streams/update-partner-stats/route.ts](https://github.com/blade47/dub/blob/HEAD/apps/web/app/(ee)/api/cron/streams/update-partner-stats/route.ts)
- [apps/web/lib/jobs/registry.ts](https://github.com/blade47/dub/blob/HEAD/apps/web/lib/jobs/registry.ts)
- [apps/web/app/ee/api/cron/payouts/aggregate-due-commissions/route.ts](https://github.com/blade47/dub/blob/HEAD/apps/web/app/(ee)/api/cron/payouts/aggregate-due-commissions/route.ts)
- [apps/web/lib/jobs/publish-workflows.ts](https://github.com/blade47/dub/blob/HEAD/apps/web/lib/jobs/publish-workflows.ts)
- [apps/web/lib/actions/partners/trigger-aggregate-due-commissions.ts](https://github.com/blade47/dub/blob/HEAD/apps/web/lib/actions/partners/trigger-aggregate-due-commissions.ts)
- [apps/web/lib/postback/postback-adapters.ts](https://github.com/blade47/dub/blob/HEAD/apps/web/lib/postback/postback-adapters.ts)
- [apps/web/lib/cron/index.ts](https://github.com/blade47/dub/blob/HEAD/apps/web/lib/cron/index.ts)
- [apps/web/lib/jobs/send-workflows.ts](https://github.com/blade47/dub/blob/HEAD/apps/web/lib/jobs/send-workflows.ts)
- [apps/web/lib/partnerstack/importer.ts](https://github.com/blade47/dub/blob/HEAD/apps/web/lib/partnerstack/importer.ts)
- [apps/web/lib/dub.ts](https://github.com/blade47/dub/blob/HEAD/apps/web/lib/dub.ts)
- [apps/web/app/ee/api/stripe/connect/webhook/balance-available.ts](https://github.com/blade47/dub/blob/HEAD/apps/web/app/(ee)/api/stripe/connect/webhook/balance-available.ts)
- [apps/web/scripts/dev/test-partner-referrals.ts](https://github.com/blade47/dub/blob/HEAD/apps/web/scripts/dev/test-partner-referrals.ts)
</details>

## Overview

Background Jobs and Queues manages asynchronous task execution, scheduled background routines, and reliable message transport across the platform. It solves the challenge of handling long-running or deferred operations reliably by integrating a typed job definition framework with Upstash QStash, a transactional Prisma-backed outbox fallback for failed dispatches, and secure webhook execution endpoints. By decoupling heavy workloads from synchronous request lifecycles, the system ensures robust error handling, concurrency control, and automated retry mechanics for periodic crons and distributed workflows.

Sources: [apps/web/lib/jobs/outbox.ts:3-15](https://github.com/blade47/dub/blob/HEAD/apps/web/lib/jobs/outbox.ts#L3-L15), [apps/web/lib/jobs/index.ts:211-221](https://github.com/blade47/dub/blob/HEAD/apps/web/lib/jobs/index.ts#L211-L221), [apps/web/app/api/jobs/process/%5BjobName%5D/route.ts:10-15](https://github.com/blade47/dub/blob/HEAD/apps/web/app/api/jobs/process/%5BjobName%5D/route.ts#L10-L15), [apps/web/prisma/schema/job.prisma:3-15](https://github.com/blade47/dub/blob/HEAD/apps/web/prisma/schema/job.prisma#L3-L15), [apps/web/lib/cron/with-cron.ts:23-28](https://github.com/blade47/dub/blob/HEAD/apps/web/lib/cron/with-cron.ts#L23-L28)

## Job Definition and Registry System

### Overview

The job definition and registry system provides a typed public API for creating background jobs and mapping them within a centralized registry. Using Zod for payload validation, developers define jobs with strict type safety, automatic dispatch helpers, and optional configuration defaults. The registry uses static dynamic imports to map job names to their respective handlers, enabling webpack code-splitting and efficient runtime loading with in-memory caching.

Sources: [apps/web/lib/jobs/index.ts:211-248](https://github.com/blade47/dub/blob/HEAD/apps/web/lib/jobs/index.ts#L211-L248), [apps/web/lib/jobs/registry.ts:4-136](https://github.com/blade47/dub/blob/HEAD/apps/web/lib/jobs/registry.ts#L4-L136)

### Job Loading and Caching Lifecycle

The registry maps registered job names to loader functions that dynamically import job handler modules. When a job is requested, the system inspects an in-memory cache before triggering the loader.

The loading process follows this exact sequence: `loadJob()` checks `jobCache.get(name)` → if cached, returns the job immediately → otherwise looks up the loader in `jobLoaders` → invokes `loader()` to dynamically import the module → validates that `job.name === name` to prevent mismatches → stores the definition in `jobCache` and returns it.

```typescript
export async function loadJob(
  name: string,
): Promise<JobDefinition | undefined> {
  const cached = jobCache.get(name);
  if (cached) return cached;

  const loader = jobLoaders[name as keyof typeof jobLoaders];
  if (!loader) return undefined;

  const job = await loader();
  if (job.name !== name) {
    throw new Error(`Job name mismatch: ${job.name} !== ${name}`);
  }

  jobCache.set(name, job);
  return job;
}
```

Sources: [apps/web/lib/jobs/registry.ts:118-134](https://github.com/blade47/dub/blob/HEAD/apps/web/lib/jobs/registry.ts#L118-L134)

> [!NOTE]
> Static imports within `jobLoaders` ensure that webpack splits each background job handler into a separate bundle chunk, preventing unnecessary code bloat in core server runtimes.
> Sources: [apps/web/lib/jobs/registry.ts:4-114](https://github.com/blade47/dub/blob/HEAD/apps/web/lib/jobs/registry.ts#L4-L114)

### Registered Job Handlers

The registry currently maps 25 distinct background job handlers covering domain operations, partner management, discounts, and analytics.

| Job Name | Module Import Path |
| :--- | :--- |
| `folder-deleted-job` | `./handlers/folder-deleted-job` |
| `partner-tag-deleted-job` | `./handlers/partner-tag-deleted-job` |
| `unban-partner-job` | `./handlers/unban-partner-job` |
| `link-tag-deleted-job` | `./handlers/link-tag-deleted-job` |
| `domain-deleted-job` | `./handlers/domain-deleted-job` |
| `default-link-deleted-job` | `./handlers/default-link-deleted-job` |
| `create-tremendous-campaign-job` | `./handlers/create-tremendous-campaign-job` |
| `sync-group-utm-job` | `./handlers/sync-group-utm-job` |
| `partner-search-sync-job` | `./handlers/partner-search-sync-job` |
| `process-shopify-order-job` | `./handlers/process-shopify-order-job` |
| `welcome-user-job` | `./handlers/welcome-user-job` |
| `auto-approve-partner-job` | `./handlers/auto-approve-partner-job` |
| `auto-reject-partner-job` | `./handlers/auto-reject-partner-job` |
| `queue-partner-program-summary-job` | `./handlers/queue-partner-program-summary-job` |
| `send-partner-program-summary-job` | `./handlers/send-partner-program-summary-job` |
| `send-connect-payout-reminders-job` | `./handlers/send-connect-payout-reminders-job` |
| `create-custom-commission-job` | `./handlers/create-custom-commission-job` |
| `invalidate-links-for-discounts-job` | `./handlers/invalidate-links-for-discounts-job` |
| `remap-discount-code-job` | `./handlers/remap-discount-code-job` |
| `attach-discount-job` | `./handlers/attach-discount-job` |
| `create-discount-code-for-link-job` | `./handlers/create-discount-code-for-link-job` |
| `publish-discount-codes-creation-job` | `./handlers/publish-discount-codes-creation-job` |
| `aggregate-clicks-job` | `./handlers/aggregate-clicks-job` |
| `process-partner-group-change-job` | `./handlers/process-partner-group-change-job` |
| `program-application-reminder-job` | `./handlers/program-application-reminder-job` |

Sources: [apps/web/lib/jobs/registry.ts:5-113](https://github.com/blade47/dub/blob/HEAD/apps/web/lib/jobs/registry.ts#L5-L113)

## Dispatching and QStash Transport

### Overview

The transport and dispatch layer controls how job payloads are serialized, wrapped into request envelopes, batched, and published to Upstash QStash. It governs retry logic, error handling, fallback deferral to outbox storage, and request construction for both single and batch dispatches.

Sources: [apps/web/lib/jobs/index.ts:35-209](https://github.com/blade47/dub/blob/HEAD/apps/web/lib/jobs/index.ts#L35-L209), [apps/web/lib/jobs/send-jobs.ts:70-104](https://github.com/blade47/dub/blob/HEAD/apps/web/lib/jobs/send-jobs.ts#L70-L104)

### QStash Transport and Request Building

Job dispatch inputs are transformed into QStash-compatible publish payloads via helper utilities that configure endpoint URLs, request bodies, delays, deduplication headers, and flow control.

```typescript
export function buildQStashJobRequest(
  { name, payload, options }: DispatchJobInput,
  opts?: {
    dispatchedAt?: string;
    batch?: boolean;
    notBefore?: number;
  },
) {
  const envelope: JobEnvelope = {
    name,
    payload,
    dispatchedAt: opts?.dispatchedAt ?? new Date().toISOString(),
  };

  const notBefore = opts?.notBefore ?? options?.notBefore;
  const deduplicationId = buildJobDeduplicationId(
    name,
    options?.deduplicationId,
  );

  return {
    url: getJobsEndpointUrl(name),
    body: envelope,
    label: buildJobLabel(name, options?.label),
    ...(options?.delay &&
      opts?.notBefore === undefined && {
        delay: options.delay,
      }),
    ...(notBefore && { notBefore }),
    ...(deduplicationId && { deduplicationId }),
    ...(options?.retries !== undefined && { retries: options.retries }),
    ...(options?.flowControl && { flowControl: options.flowControl }),
    ...(opts?.batch && options?.queue && { queueName: options.queue }),
  };
}
```

Sources: [apps/web/lib/jobs/send-jobs.ts:70-104](https://github.com/blade47/dub/blob/HEAD/apps/web/lib/jobs/send-jobs.ts#L70-L104)

To handle transient network errors when communicating with QStash, publish attempts are wrapped in an exponential backoff retry utility supporting up to 3 retry attempts.

```typescript
async function withQStashRetry<T>(fn: () => Promise<T>): Promise<T> {
  for (let attempt = 0; attempt <= QSTASH_PUBLISH_MAX_RETRIES; attempt++) {
    try {
      return await fn();
    } catch (error) {
      if (attempt < QSTASH_PUBLISH_MAX_RETRIES) {
        await sleep(1000 * Math.pow(2, attempt));
        continue;
      }

      throw error;
    }
  }

  throw new Error("Failed to publish to QStash.");
}
```

Sources: [apps/web/lib/jobs/index.ts:35-50](https://github.com/blade47/dub/blob/HEAD/apps/web/lib/jobs/index.ts#L35-L50)

> [!WARNING]
> If all retry attempts fail during `withQStashRetry`, the error propagates to the dispatch loop, which catches the failure, logs `jobs.publish_failed`, and invokes `deferJobs()` to persist the affected jobs into the database outbox table.
> Sources: [apps/web/lib/jobs/index.ts:143-165](https://github.com/blade47/dub/blob/HEAD/apps/web/lib/jobs/index.ts#L143-L165), [apps/web/lib/jobs/index.ts:175-197](https://github.com/blade47/dub/blob/HEAD/apps/web/lib/jobs/index.ts#L175-L197)

### Dispatch Flow and Batch Processing

The dispatch pipeline processes collections of job inputs by splitting them into chunks, publishing them to QStash, validating response structures, and falling back to outbox persistence if publishing fails.

The end-to-end dispatch execution walkthrough follows this exact sequence: `dispatchJobs()` checks if inputs are empty → iterates over chunks using `chunk(inputs, QSTASH_BATCH_CHUNK_SIZE)` → calls `publishJobsToQStash(inputChunk)` → inspects each response using `isPublishSuccess(response)` → if successful, increments published count and records message ID → if publishing fails for any item in a chunk, routes failed inputs to `deferJobs()` which calls `persistBackgroundJobs(inputs)` → returns a `DispatchBatchResult` containing counts and individual results.

```typescript
async function publishJobsToQStash(inputs: DispatchJobInput[]) {
  if (inputs.length === 0) {
    return [];
  }

  if (inputs.length === 1) {
    const input = inputs[0];
    const request = buildQStashJobRequest(input);

    const response = await withQStashRetry(async () => {
      if (input.options?.queue) {
        return qstash
          .queue({ queueName: input.options.queue })
          .enqueueJSON(request);
      }

      return qstash.publishJSON(request);
    });

    return [response];
  }

  const requests = inputs.map((input) =>
    buildQStashJobRequest(input, { batch: true }),
  );

  return withQStashRetry(() => qstash.batchJSON(requests));
}
```

Sources: [apps/web/lib/jobs/index.ts:61-88](https://github.com/blade47/dub/blob/HEAD/apps/web/lib/jobs/index.ts#L61-L88)

| Dispatch Option Property | Type | Purpose |
| :--- | :--- | :--- |
| `delay` | `number \| string` | Delays job execution by a relative duration in seconds. |
| `notBefore` | `number` | Schedules job execution at an absolute Unix timestamp. |
| `deduplicationId` | `string` | Ensures unique message queuing by deduplication identifier suffix. |
| `retries` | `number` | Configures the maximum number of delivery retries on failure. |
| `flowControl` | `FlowControl` | Applies rate limiting and concurrency controls to queue items. |
| `label` | `string` | Attaches a custom grouping label to the QStash request. |
| `queue` | `string` | Routes the job through a specific named QStash queue. |

Sources: [apps/web/lib/jobs/send-jobs.ts:10-20](https://github.com/blade47/dub/blob/HEAD/apps/web/lib/jobs/send-jobs.ts#L10-L20)

## Execution Webhooks and Signature Verification

### Overview

Background job execution relies on a shared worker endpoint located at `/api/jobs/process/[jobName]` that handles incoming QStash webhook dispatches, verifies request authenticity, validates payloads against job envelopes, and invokes the underlying job definition.

Sources: [apps/web/app/api/jobs/process/%5BjobName%5D/route.ts:10-106](https://github.com/blade47/dub/blob/HEAD/apps/web/app/api/jobs/process/%5BjobName%5D/route.ts#L10-L106)

> [!NOTE]
> The worker execution route sets `maxDuration = 600` to allow up to 10 minutes for long-running background tasks.
> Sources: [apps/web/app/api/jobs/process/%5BjobName%5D/route.ts:8-8](https://github.com/blade47/dub/blob/HEAD/apps/web/app/api/jobs/process/%5BjobName%5D/route.ts#L8-L8)

### Call-Chain Execution Walkthrough

The request processing pipeline follows a rigorous validation and execution sequence before invoking handler logic. The end-to-end execution path proceeds as follows: 

`POST` handler wrapped in `withAxiomBodyLog` → extracts `jobName` from route `params` → clones incoming request and extracts raw text body → `verifyQstashSignature({ req, rawBody })` validates cryptographic headers → `JSON.parse(rawBody)` parses the payload → `jobEnvelopeSchema.safeParse(parsedBody)` validates the job envelope → compares URL `jobName` against envelope `name` → `loadJob(jobName)` retrieves the job definition from the registry → `job.execute(envelope.data.payload)` runs the handler.

Sources: [apps/web/app/api/jobs/process/%5BjobName%5D/route.ts:11-75](https://github.com/blade47/dub/blob/HEAD/apps/web/app/api/jobs/process/%5BjobName%5D/route.ts#L11-L75)

> [!WARNING]
> If a payload fails Zod validation (`z.ZodError`), the worker catches the error and explicitly returns a `200` HTTP status code rather than a `500` error. This guarantees that QStash treats the malformed payload as permanently invalid and halts further automatic retries.
> Sources: [apps/web/app/api/jobs/process/%5BjobName%5D/route.ts:77-88](https://github.com/blade47/dub/blob/HEAD/apps/web/app/api/jobs/process/%5BjobName%5D/route.ts#L77-L88)

### Cron and Webhook Authentication via `withCron`

Periodic crons and webhook triggers share a common wrapper utility `withCron` that enforces signature authentication based on the HTTP method before passing control to the downstream handler.

```typescript
export const withCron = (handler: WithCronHandler) => {
  return withAxiomBodyLog(
    async (
      req,
      { params: initialParams }: { params: Promise<Record<string, string>> },
    ) => {
      const clonedReq = req.clone();
      const params = (await initialParams) || {};
      const searchParams = getSearchParams(req.url);

      try {
        let rawBody: string | undefined;

        if (req.method === "GET") {
          await verifyVercelSignature(req);
        } else if (req.method === "POST") {
          rawBody = await clonedReq.text();
          await verifyQstashSignature({ req, rawBody });
        } else {
          throw new Error(`Unsupported HTTP method: ${req.method}`);
        }

        return await handler({
          req: clonedReq,
          searchParams,
          params,
          rawBody: rawBody ?? "",
        });
      } catch (error) {
        console.error(error);
        const errorMessage =
          error instanceof Error ? error.message : String(error);

        logger.error(errorMessage, error);
        await logger.flush();

        const statusCode =
          error instanceof DubApiError
            ? ErrorCodes[error.code]
            : ErrorCodes.internal_server_error;

        return logAndRespond(errorMessage, { status: statusCode });
      }
    },
  );
};
```

Sources: [apps/web/lib/cron/with-cron.ts:23-76](https://github.com/blade47/dub/blob/HEAD/apps/web/lib/cron/with-cron.ts#L23-L76)

| HTTP Method | Authentication Mechanism | Target Provider / Source |
| :--- | :--- | :--- |
| `GET` | `verifyVercelSignature(req)` | Vercel Cron |
| `POST` | `verifyQstashSignature({ req, rawBody })` | Upstash QStash |

Sources: [apps/web/lib/cron/with-cron.ts:39-47](https://github.com/blade47/dub/blob/HEAD/apps/web/lib/cron/with-cron.ts#L39-L47)

## Transactional Outbox and Failure Recovery

### Overview

When background jobs or workflows fail to publish to Upstash QStash at dispatch time due to network partitions, rate limits, or downstream outages, the system relies on a transactional outbox pattern backed by Prisma and PostgreSQL. Unsent jobs are persisted to the database via `persistBackgroundJobs()` or `persistFailedJobs()`, preserving their payload, options, scheduling time, and attempt counts. A dedicated retry cron endpoint at `/api/cron/queue/retry` regularly polls for due jobs, republishes them through appropriate transport layers, and settles their database state.

Sources: [apps/web/lib/jobs/outbox.ts:45-122](https://github.com/blade47/dub/blob/HEAD/apps/web/lib/jobs/outbox.ts#L45-L122), [apps/web/app/(ee)/api/cron/queue/retry/route.ts:14-16](https://github.com/blade47/dub/blob/HEAD/apps/web/app/(ee)/api/cron/queue/retry/route.ts#L14-L16), [apps/web/prisma/schema/job.prisma:1-15](https://github.com/blade47/dub/blob/HEAD/apps/web/prisma/schema/job.prisma#L1-L15)

### Prisma Job Schema and Indexing

The `Job` model stores un-dispatched work items. It uses a compound index on `scheduledAt` and `attempts` to allow the retry cron to efficiently locate due jobs that have not exceeded the maximum attempt threshold.

```prisma
model Job {
  id          String   @id
  name        String
  payload     Json
  options     Json? // dispatch options replayed verbatim: deduplicationId, retries, queue, flowControl, label
  scheduledAt DateTime @default(now()) // when we should publish the job
  attempts    Int      @default(0)
  lastError   String?  @db.Text
  createdAt   DateTime @default(now())
  updatedAt   DateTime @updatedAt

  @@index([scheduledAt, attempts]) // retry cron: due jobs under the attempt cap
}
```

Sources: [apps/web/prisma/schema/job.prisma:3-15](https://github.com/blade47/dub/blob/HEAD/apps/web/prisma/schema/job.prisma#L3-L15)

> [!NOTE]
> The first failed publish attempt initializes `attempts` to 1 via `toJobCreateInput()`, ensuring that the retry budget aligns with `MAX_JOB_ATTEMPTS`.
> Sources: [apps/web/lib/jobs/outbox.ts:26-43](https://github.com/blade47/dub/blob/HEAD/apps/web/lib/jobs/outbox.ts#L26-L43)

### Retry Cron Polling Flow and Distributed Locking

The retry cron endpoint executes every minute via Vercel Cron. To prevent concurrent executions from overlapping during long-running batch republish cycles, the handler acquires a Redis distributed lock (`lock:queue-retry`) with a 600-second TTL matching the maximum cron duration before invoking `publishPendingJobs()`.

```typescript
export const GET = withCron(async () => {
  const acquired = await redis.set(LOCK_KEY, "1", {
    nx: true,
    ex: LOCK_TTL_SECONDS,
  });

  if (!acquired) {
    return logAndRespond(
      "[queue-retry] Another run is in progress. Skipping...",
    );
  }

  try {
    const { attempted, published, failed } = await publishPendingJobs();

    if (attempted === 0) {
      return logAndRespond("No background jobs to retry.");
    }

    if (failed === 0) {
      return logAndRespond(
        `Republished ${published} background jobs to QStash.`,
      );
    }

    return logAndRespond(
      `Republished ${published} background jobs to QStash; failed to republish ${failed}.`,
    );
  } finally {
    await redis.del(LOCK_KEY);
  }
});
```

Sources: [apps/web/app/(ee)/api/cron/queue/retry/route.ts:8-47](https://github.com/blade47/dub/blob/HEAD/apps/web/app/(ee)/api/cron/queue/retry/route.ts#L8-L47)

Sources: [apps/web/app/(ee)/api/cron/queue/retry/route.ts:8-12](https://github.com/blade47/dub/blob/HEAD/apps/web/app/(ee)/api/cron/queue/retry/route.ts#L8-L12)

### Call-Chain Execution Walkthrough

The background retry mechanism processes pending outbox entries through a structured multi-step retrieval and settlement pipeline. The full execution sequence proceeds as follows:

`GET` request hits `/api/cron/queue/retry` wrapped in `withCron` → acquires Redis distributed lock `lock:queue-retry` with `nx: true` and 600s TTL → `publishPendingJobs()` queries Prisma for due records where `scheduledAt <= now()` and `attempts < MAX_JOB_ATTEMPTS` ordered by `createdAt` asc with a take limit of `MAX_JOBS_PER_BATCH` → matches jobs against transport definitions (`isWorkflowName` or `isDefineJobName`) → dispatches batches via `triggerWorkflows` or `sendJobs` → `settlePublishResults({ results, jobs })` evaluates outcomes → successfully published job IDs are deleted via `prisma.job.deleteMany()` → failed publish results group error messages and update records via `prisma.job.updateMany()` incrementing `attempts` and setting `lastError` → if any job reaches `MAX_JOB_ATTEMPTS`, an error event `jobs.retry_exhausted` is logged to Axiom.

Sources: [apps/web/lib/jobs/outbox.ts:125-220](https://github.com/blade47/dub/blob/HEAD/apps/web/lib/jobs/outbox.ts#L125-L220), [apps/web/app/(ee)/api/cron/queue/retry/route.ts:16-46](https://github.com/blade47/dub/blob/HEAD/apps/web/app/(ee)/api/cron/queue/retry/route.ts#L16-L46)

### Outbox Transport Configuration and Error Handling

Job transport routing maps specific job naming conventions to their respective dispatch handlers using a transport array.

| Matcher Function | Transport Handler | Purpose |
| :--- | :--- | :--- |
| `isWorkflowName` | `triggerWorkflows` | Handles workflow execution steps and triggers |
| `isDefineJobName` | `sendJobs` | Handles standard typed background jobs |

Sources: [apps/web/lib/jobs/outbox.ts:10-24](https://github.com/blade47/dub/blob/HEAD/apps/web/lib/jobs/outbox.ts#L10-L24)

> [!CAUTION]
> When `persistFailedJobs()` encounters a database error while attempting to record failed dispatches, it catches the exception and logs it via Axiom (`jobs.dispatch_lost`) rather than rethrowing, swallowing the error to prevent dispatch request failures from crashing upstream API routes.
> Sources: [apps/web/lib/jobs/outbox.ts:45-84](https://github.com/blade47/dub/blob/HEAD/apps/web/lib/jobs/outbox.ts#L45-L84)

## Workflow Orchestration and Batch Enqueuing

### Overview

Workflow orchestration and batch enqueuing provide mechanisms for triggering complex workflow steps and pacing batch enqueue operations. The system coordinates scheduled campaign fan-out, transactional workflow intervals, and Upstash QStash batch transmissions.
Sources: [apps/web/app/(ee)/api/cron/campaigns/queue-scheduled/route.ts:18-26](https://github.com/blade47/dub/blob/HEAD/apps/web/app/(ee)/api/cron/campaigns/queue-scheduled/route.ts#L18-L26), [apps/web/lib/cron/enqueue-batch-jobs.ts:17-18](https://github.com/blade47/dub/blob/HEAD/apps/web/lib/cron/enqueue-batch-jobs.ts#L17-L18)

### Batch Enqueue and Retry Mechanism

The `enqueueBatchJobs` function handles batch delivery to QStash with built-in retry logic and exponential backoff.
Sources: [apps/web/lib/cron/enqueue-batch-jobs.ts:17-18](https://github.com/blade47/dub/blob/HEAD/apps/web/lib/cron/enqueue-batch-jobs.ts#L17-L18)

```typescript
export async function enqueueBatchJobs(jobs: EnqueueBatchJobsProps[]) {
  const maxRetries = 3;

  for (let attempt = 0; attempt <= maxRetries; attempt++) {
    try {
      return await qstash.batchJSON(jobs);
    } catch (error) {
      if (attempt < maxRetries) {
        await sleep(1000 * Math.pow(2, attempt));
        continue;
      }

      await log({
        message: `[enqueueBatchJobs] Failed to enqueue batch jobs: ${JSON.stringify(error, null, 2)}`,
        type: "errors",
        mention: true,
      });

      throw new Error(
        `Failed to enqueue batch jobs: ${JSON.stringify(error, null, 2)}`,
      );
    }
  }
}
```
Sources: [apps/web/lib/cron/enqueue-batch-jobs.ts:18-41](https://github.com/blade47/dub/blob/HEAD/apps/web/lib/cron/enqueue-batch-jobs.ts#L18-L41)

> [!WARNING]
> If all retry attempts fail in `enqueueBatchJobs`, an error is logged with mention tagging via `log()` before throwing an unhandled `Error`, which can disrupt upstream campaign queueing cron executions.
> Sources: [apps/web/lib/cron/enqueue-batch-jobs.ts:30-39](https://github.com/blade47/dub/blob/HEAD/apps/web/lib/cron/enqueue-batch-jobs.ts#L30-L39)

### Workflow Dispatch and Validation

Workflows are validated against strict kebab-case naming rules and mapped to explicit API endpoints before dispatching requests via the QStash workflow client.
Sources: [apps/web/lib/jobs/send-workflows.ts:20-38](https://github.com/blade47/dub/blob/HEAD/apps/web/lib/jobs/send-workflows.ts#L20-L38), [apps/web/lib/jobs/send-workflows.ts:82-95](https://github.com/blade47/dub/blob/HEAD/apps/web/lib/jobs/send-workflows.ts#L82-L95)

| Workflow Name | Target Path |
| :--- | :--- |
| `partner-approved-workflow` | `/api/workflows/partner-approved` |
| `merge-partner-accounts-workflow` | `/api/workflows/merge-partner-accounts` |
| `create-partner-commission-workflow` | `/api/workflows/create-partner-commission` |
| `reattribute-customer-workflow` | `/api/workflows/reattribute-customer` |
| `detach-discount-workflow` | `/api/workflows/detach-discount` |

Sources: [apps/web/lib/jobs/send-workflows.ts:27-34](https://github.com/blade47/dub/blob/HEAD/apps/web/lib/jobs/send-workflows.ts#L27-L34)

### Scheduled Campaign Queueing Call-Chain

The scheduled campaign queueing cron endpoint processes transactional and marketing campaigns through a structured paging and batching execution sequence.
Sources: [apps/web/app/(ee)/api/cron/campaigns/queue-scheduled/route.ts:20-26](https://github.com/blade47/dub/blob/HEAD/apps/web/app/(ee)/api/cron/campaigns/queue-scheduled/route.ts#L20-L26)

`GET` request hits `/api/cron/campaigns/queue-scheduled` wrapped in `withCron` → `Promise.allSettled()` invokes `queueTransactionalCampaigns(now)` and `queueMarketingCampaigns(now)` in parallel → `queueTransactionalCampaigns` verifies `isTransactionalTick(now)` (checking if UTC hour % 12 === 0 and UTC minute < 5) → queries `prisma.campaign.findMany` with pagination cursor `lastCampaignId` and `CRON_BATCH_SIZE` taking active transactional campaigns with enabled workflows → filters scheduled workflows via `isScheduledWorkflow` → dispatches batches via `enqueueBatchJobs()` pointing to `/api/cron/workflows/${workflow.id}` with parallelism set to 10 → `queueMarketingCampaigns` queries scheduled marketing campaigns where `scheduledAt <= now` → dispatches batches via `enqueueBatchJobs()` pointing to `/api/cron/campaigns/broadcast` with parallelism set to 1 → error handling aggregates rejections and logs failure via `log()` if any.
Sources: [apps/web/app/(ee)/api/cron/campaigns/queue-scheduled/route.ts:20-178](https://github.com/blade47/dub/blob/HEAD/apps/web/app/(ee)/api/cron/campaigns/queue-scheduled/route.ts#L20-L178)

> [!TIP]
> The 5-minute transaction tick window (`isTransactionalTick`) absorbs Vercel cron jitter, while QStash deduplication lasting 10 minutes absorbs duplicate publishes without leaking extra triggers after expiry.
> Sources: [apps/web/app/(ee)/api/cron/campaigns/queue-scheduled/route.ts:58-69](https://github.com/blade47/dub/blob/HEAD/apps/web/app/(ee)/api/cron/campaigns/queue-scheduled/route.ts#L58-L69)

## Periodic Crons and Stream Processing

### Overview

Periodic cron jobs and Redis stream consumers automate scheduled background operations, batch event processing, and periodic maintenance tasks across Dub's enterprise tier. These routines are centrally protected by authentication wrappers that verify invocation signatures from external schedulers and message brokers.
Sources: [apps/web/lib/cron/with-cron.ts:23-76](https://github.com/blade47/dub/blob/HEAD/apps/web/lib/cron/with-cron.ts#L23-L76), [apps/web/app/(ee)/api/cron/streams/update-click-stats/route.ts:1-11](https://github.com/blade47/dub/blob/HEAD/apps/web/app/(ee)/api/cron/streams/update-click-stats/route.ts#L1-L11)

### Cron Endpoint Authentication

Periodic crons and webhook triggers share a common wrapper utility `withCron` that enforces signature authentication based on the HTTP method before passing control to the downstream handler.

#### Call-Chain Execution Walkthrough

`POST`/`GET` request received at cron route → wrapped handler executes via `withAxiomBodyLog` → request is cloned (`req.clone()`) to allow body inspection while preserving the original stream for Axiom logging → `params` and `searchParams` are extracted from the URL → request method branch evaluates:
- If `GET`, `verifyVercelSignature(req)` validates the invocation against Vercel Cron secrets.
- If `POST`, `clonedReq.text()` reads the raw body text, and `verifyQstashSignature({ req, rawBody })` verifies the Upstash QStash signature.
- If any other method is encountered, an unsupported method error is thrown.
→ The inner cron task handler executes with structured arguments (`req`, `searchParams`, `params`, `rawBody`) → Any caught errors are logged to Axiom via `logger.error()` and flushed before translating the error into an appropriate HTTP status code via `ErrorCodes` and returning a standardized error response.
Sources: [apps/web/lib/cron/with-cron.ts:24-75](https://github.com/blade47/dub/blob/HEAD/apps/web/lib/cron/with-cron.ts#L24-L75)

> [!NOTE]
> `withCron` clones the incoming request early in the middleware lifecycle so individual cron handlers can read request bodies without mutating or consuming the stream required by Axiom audit logging.
> Sources: [apps/web/lib/cron/with-cron.ts:29-31](https://github.com/blade47/dub/blob/HEAD/apps/web/lib/cron/with-cron.ts#L29-L31)

### Scheduled Tasks and Fan-Out Crons

Periodic background endpoints handle tasks ranging from program application reminders and social metrics queueing to reward processing and commission aggregation.
Sources: [apps/web/app/(ee)/api/cron/program-application-reminder/route.ts:7-13](https://github.com/blade47/dub/blob/HEAD/apps/web/app/(ee)/api/cron/program-application-reminder/route.ts#L7-L13), [apps/web/app/(ee)/api/cron/bounties/queue-sync-social-metrics/route.ts:11-49](https://github.com/blade47/dub/blob/HEAD/apps/web/app/(ee)/api/cron/bounties/queue-sync-social-metrics/route.ts#L11-L49), [apps/web/app/(ee)/api/cron/rewards/queue-custom-commissions/route.ts:10-40](https://github.com/blade47/dub/blob/HEAD/apps/web/app/(ee)/api/cron/rewards/queue-custom-commissions/route.ts#L10-L40), [apps/web/app/(ee)/api/cron/payouts/aggregate-due-commissions/route.ts:13-68](https://github.com/blade47/dub/blob/HEAD/apps/web/app/(ee)/api/cron/payouts/aggregate-due-commissions/route.ts#L13-L68)

| Cron Route | HTTP Method | Max Duration | Purpose |
| :--- | :--- | :--- | :--- |
| `/api/cron/program-application-reminder` | `POST` | Default | Sends reminders to unverified program applicants |
| `/api/cron/bounties/queue-sync-social-metrics` | `GET` | Default | Queues social metric synchronization batches for submission bounties |
| `/api/cron/rewards/queue-custom-commissions` | `GET` | 600s | Fans out custom commission generation jobs for due rewards |
| `/api/cron/sitemaps/queue` | `GET`/`POST` | Default | Paging loop that queues workspace sitemap imports via QStash |
| `/api/cron/trial-emails` | `GET`/`POST` | Default | Executes paid-plan trial marketing email sequences recursively |
| `/api/cron/payouts/aggregate-due-commissions` | `GET` | 600s | Aggregates eligible pending commissions into partner payouts |

Sources: [apps/web/app/(ee)/api/cron/program-application-reminder/route.ts:7-13](https://github.com/blade47/dub/blob/HEAD/apps/web/app/(ee)/api/cron/program-application-reminder/route.ts#L7-L13), [apps/web/app/(ee)/api/cron/bounties/queue-sync-social-metrics/route.ts:11-49](https://github.com/blade47/dub/blob/HEAD/apps/web/app/(ee)/api/cron/bounties/queue-sync-social-metrics/route.ts#L11-L49), [apps/web/app/(ee)/api/cron/rewards/queue-custom-commissions/route.ts:10-40](https://github.com/blade47/dub/blob/HEAD/apps/web/app/(ee)/api/cron/rewards/queue-custom-commissions/route.ts#L10-L40), [apps/web/app/(ee)/api/cron/sitemaps/queue/route.ts:9-103](https://github.com/blade47/dub/blob/HEAD/apps/web/app/(ee)/api/cron/sitemaps/queue/route.ts#L9-L103), [apps/web/app/(ee)/api/cron/trial-emails/route.ts:9-63](https://github.com/blade47/dub/blob/HEAD/apps/web/app/(ee)/api/cron/trial-emails/route.ts#L9-L63), [apps/web/app/(ee)/api/cron/payouts/aggregate-due-commissions/route.ts:13-68](https://github.com/blade47/dub/blob/HEAD/apps/web/app/(ee)/api/cron/payouts/aggregate-due-commissions/route.ts#L13-L68)

### Redis Stream Event Consumers

High-frequency telemetry such as click events and partner activity logs are ingested into Redis streams and processed in batches by dedicated stream consumers.
Sources: [apps/web/app/(ee)/api/cron/streams/update-click-stats/route.ts:1-170](https://github.com/blade47/dub/blob/HEAD/apps/web/app/(ee)/api/cron/streams/update-click-stats/route.ts#L1-L170), [apps/web/app/(ee)/api/cron/streams/update-partner-stats/route.ts:1-370](https://github.com/blade47/dub/blob/HEAD/apps/web/app/(ee)/api/cron/streams/update-partner-stats/route.ts#L1-L370)

The click stats consumer pulls entries up to a defined batch size, aggregates metrics by link, workspace, and program enrollment, and commits database writes across parallel sub-batches. Similarly, the partner activity stream consumer groups events by program-partner pairs, queries grouped link and commission statistics in parallel from Prisma, computes derived performance metrics (such as net revenue, earnings per click, and consistency scores), and flushes updates via raw PlanetScale SQL execution.
Sources: [apps/web/app/(ee)/api/cron/streams/update-click-stats/route.ts:15-170](https://github.com/blade47/dub/blob/HEAD/apps/web/app/(ee)/api/cron/streams/update-click-stats/route.ts#L15-L170), [apps/web/app/(ee)/api/cron/streams/update-partner-stats/route.ts:15-343](https://github.com/blade47/dub/blob/HEAD/apps/web/app/(ee)/api/cron/streams/update-partner-stats/route.ts#L15-L343)

> [!CAUTION]
> Stream processors define concurrency locks and batch limits (such as `BATCH_SIZE = 10_000` for clicks and `BATCH_SIZE = 6000` for partner activities) to prevent overlapping invocations from double-counting metrics or exhausting database connection pools.
> Sources: [apps/web/app/(ee)/api/cron/streams/update-click-stats/route.ts:15-20](https://github.com/blade47/dub/blob/HEAD/apps/web/app/(ee)/api/cron/streams/update-click-stats/route.ts#L15-L20), [apps/web/app/(ee)/api/cron/streams/update-partner-stats/route.ts:15-15](https://github.com/blade47/dub/blob/HEAD/apps/web/app/(ee)/api/cron/streams/update-partner-stats/route.ts#L15-L15)

## Related

- [Payout Processing](https://www.doc0.dev/docs/934e554a-e6a1-476f-bb2f-23e62d86c3fd/technical/affiliate-platform/payout-processing)
- [Campaign Broadcaster](https://www.doc0.dev/docs/934e554a-e6a1-476f-bb2f-23e62d86c3fd/technical/automation-and-communications/campaign-broadcaster)


## Sitemap

See the full [sitemap](https://www.doc0.dev/docs/934e554a-e6a1-476f-bb2f-23e62d86c3fd/llms.txt) for all pages in this wiki.
