Getting Started
Core Architecture
Link Engine
Analytics & Attribution
Partners & Affiliates
Third-Party Integrations
Identity & Security
Automation & Messaging
Developer Tools
The following files were used as context for generating this wiki page:
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, apps/web/lib/jobs/index.ts:211-221, apps/web/app/api/jobs/process/%5BjobName%5D/route.ts:10-15, apps/web/prisma/schema/job.prisma:3-15, apps/web/lib/cron/with-cron.ts:23-28
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.
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.
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;
}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
The registry currently maps 25 distinct background job handlers covering domain operations, partner management, discounts, and analytics.
Sources: apps/web/lib/jobs/registry.ts:5-113
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.
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.
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 }),
};
}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.
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
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, apps/web/lib/jobs/index.ts:175-197
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.
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
Sources: apps/web/lib/jobs/send-jobs.ts:10-20
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.
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
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.
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
withCronPeriodic 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.
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
Sources: apps/web/lib/cron/with-cron.ts:39-47
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, apps/web/app/(ee)/api/cron/queue/retry/route.ts:14-16, apps/web/prisma/schema/job.prisma:1-15
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.
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
}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
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().
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);
}
});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.
Job transport routing maps specific job naming conventions to their respective dispatch handlers using a transport array.
Sources: apps/web/lib/jobs/outbox.ts:10-24
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
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, apps/web/lib/cron/enqueue-batch-jobs.ts:17-18
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
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)}`,
);
}
}
}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
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, apps/web/lib/jobs/send-workflows.ts:82-95
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
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
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
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, apps/web/app/(ee)/api/cron/streams/update-click-stats/route.ts:1-11
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.
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:
GET, verifyVercelSignature(req) validates the invocation against Vercel Cron secrets.POST, clonedReq.text() reads the raw body text, and verifyQstashSignature({ req, rawBody }) verifies the Upstash QStash signature.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-75Note
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
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, apps/web/app/(ee)/api/cron/bounties/queue-sync-social-metrics/route.ts:11-49, apps/web/app/(ee)/api/cron/rewards/queue-custom-commissions/route.ts:10-40, apps/web/app/(ee)/api/cron/payouts/aggregate-due-commissions/route.ts:13-68
Sources: apps/web/app/(ee)/api/cron/program-application-reminder/route.ts:7-13, apps/web/app/(ee)/api/cron/bounties/queue-sync-social-metrics/route.ts:11-49, apps/web/app/(ee)/api/cron/rewards/queue-custom-commissions/route.ts:10-40, apps/web/app/(ee)/api/cron/sitemaps/queue/route.ts:9-103, apps/web/app/(ee)/api/cron/trial-emails/route.ts:9-63, apps/web/app/(ee)/api/cron/payouts/aggregate-due-commissions/route.ts:13-68
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, apps/web/app/(ee)/api/cron/streams/update-partner-stats/route.ts:1-370
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, apps/web/app/(ee)/api/cron/streams/update-partner-stats/route.ts:15-343
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, apps/web/app/(ee)/api/cron/streams/update-partner-stats/route.ts:15-15