From 7dd253f3432f5c1b37afc5f50bd00703b6eb8bf1 Mon Sep 17 00:00:00 2001 From: Kiran K Date: Tue, 18 Aug 2026 19:14:02 +0530 Subject: [PATCH 01/13] Queue due campaigns from a cron instead of per-campaign QStash schedules --- .../(ee)/api/campaigns/[campaignId]/route.ts | 23 --- .../api/cron/campaigns/broadcast/route.ts | 45 ++--- .../cron/campaigns/queue-scheduled/route.ts | 138 +++++++++++++ .../pause-campaigns-on-plan-downgrade.ts | 44 +--- .../lib/api/campaigns/schedule-campaigns.ts | 189 ------------------ apps/web/lib/zod/schemas/workflows.ts | 3 - apps/web/playwright/assert-local-database.ts | 6 +- apps/web/prisma/schema/campaign.prisma | 21 +- apps/web/scripts/dev/data.json | 7 + apps/web/scripts/dev/seed.ts | 68 ++++--- apps/web/vercel.json | 4 + 11 files changed, 221 insertions(+), 327 deletions(-) create mode 100644 apps/web/app/(ee)/api/cron/campaigns/queue-scheduled/route.ts delete mode 100644 apps/web/lib/api/campaigns/schedule-campaigns.ts diff --git a/apps/web/app/(ee)/api/campaigns/[campaignId]/route.ts b/apps/web/app/(ee)/api/campaigns/[campaignId]/route.ts index a5ded7839a0..eaa8feb6f0c 100644 --- a/apps/web/app/(ee)/api/campaigns/[campaignId]/route.ts +++ b/apps/web/app/(ee)/api/campaigns/[campaignId]/route.ts @@ -1,8 +1,4 @@ import { getCampaignOrThrow } from "@/lib/api/campaigns/get-campaign-or-throw"; -import { - deleteCampaignSchedule, - scheduleCampaign, -} from "@/lib/api/campaigns/schedule-campaigns"; import { campaignEligibilityIncludes, transformCampaign, @@ -21,7 +17,6 @@ import { } from "@/lib/zod/schemas/campaigns"; import { arrayEqual, pluck } from "@dub/utils"; import { PartnerGroup } from "@prisma/client"; -import { waitUntil } from "@vercel/functions"; import { NextResponse } from "next/server"; // GET /api/campaigns/[campaignId] - get an email campaign @@ -191,13 +186,6 @@ export const PATCH = withWorkspace( }); }); - waitUntil( - scheduleCampaign({ - campaign, - updatedCampaign, - }), - ); - return NextResponse.json( CampaignSchema.parse(transformCampaign(updatedCampaign)), ); @@ -217,15 +205,6 @@ export const DELETE = withWorkspace( const campaign = await getCampaignOrThrow({ programId, campaignId, - include: { - workflow: { - select: { - id: true, - actions: true, - triggerConditions: true, - }, - }, - }, }); await prisma.$transaction(async (tx) => { @@ -244,8 +223,6 @@ export const DELETE = withWorkspace( } }); - waitUntil(deleteCampaignSchedule(campaign)); - return NextResponse.json({ id: campaignId }); }, { diff --git a/apps/web/app/(ee)/api/cron/campaigns/broadcast/route.ts b/apps/web/app/(ee)/api/cron/campaigns/broadcast/route.ts index ced6b7dd481..d4dcfa19491 100644 --- a/apps/web/app/(ee)/api/cron/campaigns/broadcast/route.ts +++ b/apps/web/app/(ee)/api/cron/campaigns/broadcast/route.ts @@ -15,7 +15,6 @@ import CampaignEmail from "@dub/email/templates/campaign-email"; import { APP_DOMAIN_WITH_NGROK, chunk, log, pluck } from "@dub/utils"; import { NotificationEmailType } from "@prisma/client"; import { differenceInMinutes } from "date-fns"; -import { headers } from "next/headers"; import * as z from "zod/v4"; import { logAndRespond } from "../../utils"; @@ -102,20 +101,6 @@ export async function POST(req: Request) { } } - // This is a safety check to ensure the campaign broadcast is not "initiated" multiple times - const headersList = await headers(); - const upstashMessageId = headersList.get("Upstash-Message-Id"); - - if ( - !startingAfter && // First run - campaign.qstashMessageId && - upstashMessageId !== campaign.qstashMessageId - ) { - return logAndRespond( - `Campaign ${campaignId} broadcast was skipped because it is not the current message being processed.`, - ); - } - const program = campaign.program; // TODO: We should make the from address required. There are existing campaign without from address @@ -126,19 +111,23 @@ export async function POST(req: Request) { }); } - // Mark the campaign as sending (if it's in scheduled status) - if (campaign.status === "scheduled") { - try { - await prisma.campaign.update({ - where: { - id: campaignId, - }, - data: { - status: "sending", - }, - }); - } catch (error) { - // + // Claim the first run so leftover delayed messages / scanner retries + // cannot start a second broadcast. + if (!startingAfter) { + const claimed = await prisma.campaign.updateMany({ + where: { + id: campaignId, + status: "scheduled", + }, + data: { + status: "sending", + }, + }); + + if (claimed.count === 0) { + return logAndRespond( + `Campaign ${campaignId} broadcast already initiated. Skipping...`, + ); } } diff --git a/apps/web/app/(ee)/api/cron/campaigns/queue-scheduled/route.ts b/apps/web/app/(ee)/api/cron/campaigns/queue-scheduled/route.ts new file mode 100644 index 00000000000..9008e80021a --- /dev/null +++ b/apps/web/app/(ee)/api/cron/campaigns/queue-scheduled/route.ts @@ -0,0 +1,138 @@ +import { isScheduledWorkflow } from "@/lib/api/workflows/utils"; +import { CRON_BATCH_SIZE, qstash } from "@/lib/cron"; +import { withCron } from "@/lib/cron/with-cron"; +import { prisma } from "@/lib/prisma"; +import { APP_DOMAIN_WITH_NGROK } from "@dub/utils"; +import { CampaignStatus, CampaignType } from "@prisma/client"; +import { logAndRespond } from "../../utils"; + +export const dynamic = "force-dynamic"; +export const maxDuration = 600; + +// GET /api/cron/campaigns/queue-scheduled +// Fans out due marketing broadcasts and (on the 12h tick) scheduled transactional workflows. +export const GET = withCron(async () => { + const now = new Date(); + + await Promise.all([ + queueTransactionalCampaigns(now), + queueMarketingCampaigns(now), + ]); + + return logAndRespond("Finished the campaigns queueing process."); +}); + +async function queueTransactionalCampaigns(now: Date) { + // Matches the 12h enrollment window in executeSendCampaignWorkflow. + if (now.getUTCMinutes() !== 0 || now.getUTCHours() % 12 !== 0) { + console.log("[Transactional] Not the 12h tick, skipping campaigns."); + return; + } + + let queued = 0; + let page = 0; + + while (true) { + const campaigns = await prisma.campaign.findMany({ + where: { + type: CampaignType.transactional, + status: CampaignStatus.active, + workflow: { + disabledAt: null, + }, + }, + select: { + workflow: { + select: { + id: true, + triggerConditions: true, + actions: true, + }, + }, + }, + take: CRON_BATCH_SIZE, + skip: page * CRON_BATCH_SIZE, + orderBy: { + id: "asc", + }, + }); + + if (campaigns.length === 0) { + console.log("[Transactional] No more campaigns to queue."); + break; + } + + const scheduledWorkflows = campaigns.flatMap((campaign) => + campaign.workflow && isScheduledWorkflow(campaign.workflow) + ? [campaign.workflow] + : [], + ); + + if (scheduledWorkflows.length > 0) { + await qstash.batchJSON( + scheduledWorkflows.map((workflow) => ({ + url: `${APP_DOMAIN_WITH_NGROK}/api/cron/workflows/${workflow.id}`, + deduplicationId: workflow.id, + flowControl: { + key: "execute-scheduled-workflow", + parallelism: 10, + }, + body: {}, + })), + ); + + queued += scheduledWorkflows.length; + } + + page++; + } + + console.log(`[Transactional] Queued ${queued} campaigns.`); +} + +async function queueMarketingCampaigns(now: Date) { + let queued = 0; + let page = 0; + + while (true) { + const campaigns = await prisma.campaign.findMany({ + where: { + type: CampaignType.marketing, + status: CampaignStatus.scheduled, + OR: [{ scheduledAt: null }, { scheduledAt: { lte: now } }], + }, + select: { + id: true, + }, + take: CRON_BATCH_SIZE, + skip: page * CRON_BATCH_SIZE, + orderBy: { + id: "asc", + }, + }); + + if (campaigns.length === 0) { + console.log("[Marketing] No more campaigns to queue."); + break; + } + + await qstash.batchJSON( + campaigns.map((campaign) => ({ + url: `${APP_DOMAIN_WITH_NGROK}/api/cron/campaigns/broadcast`, + deduplicationId: campaign.id, + flowControl: { + key: "broadcast-marketing-campaign", + parallelism: 1, + }, + body: { + campaignId: campaign.id, + }, + })), + ); + + queued += campaigns.length; + page++; + } + + console.log(`[Marketing] Queued ${queued} campaigns.`); +} diff --git a/apps/web/lib/api/campaigns/pause-campaigns-on-plan-downgrade.ts b/apps/web/lib/api/campaigns/pause-campaigns-on-plan-downgrade.ts index e42c8cc74f6..89cf02a91ec 100644 --- a/apps/web/lib/api/campaigns/pause-campaigns-on-plan-downgrade.ts +++ b/apps/web/lib/api/campaigns/pause-campaigns-on-plan-downgrade.ts @@ -1,5 +1,3 @@ -import { scheduleTransactionalCampaign } from "@/lib/api/campaigns/schedule-campaigns"; -import { qstash } from "@/lib/cron"; import { prisma } from "@/lib/prisma"; import { CampaignStatus, CampaignType } from "@prisma/client"; @@ -8,7 +6,7 @@ export async function pauseOrCancelCampaignsForProgramOnPlanDowngrade({ }: { programId: string; }): Promise { - const marketingCampaigns = await prisma.campaign.findMany({ + await prisma.campaign.updateMany({ where: { programId, type: CampaignType.marketing, @@ -16,46 +14,22 @@ export async function pauseOrCancelCampaignsForProgramOnPlanDowngrade({ in: [CampaignStatus.scheduled, CampaignStatus.sending], }, }, + data: { + status: CampaignStatus.canceled, + }, }); - for (const campaign of marketingCampaigns) { - let qstashDeleteSucceeded = !campaign.qstashMessageId; - - if (campaign.qstashMessageId) { - try { - await qstash.messages.cancel(campaign.qstashMessageId); - qstashDeleteSucceeded = true; - } catch (error) { - console.warn( - `Failed to delete QStash message ${campaign.qstashMessageId} for campaign ${campaign.id}:`, - error, - ); - } - } - - await prisma.campaign.update({ - where: { id: campaign.id }, - data: { - status: CampaignStatus.canceled, - ...(qstashDeleteSucceeded ? { qstashMessageId: null } : {}), - }, - }); - } - const transactionalCampaigns = await prisma.campaign.findMany({ where: { programId, type: CampaignType.transactional, status: CampaignStatus.active, }, - include: { - workflow: true, - }, }); for (const campaign of transactionalCampaigns) { try { - const updatedCampaign = await prisma.$transaction(async (tx) => { + await prisma.$transaction(async (tx) => { if (campaign.workflowId) { await tx.workflow.update({ where: { id: campaign.workflowId }, @@ -63,17 +37,11 @@ export async function pauseOrCancelCampaignsForProgramOnPlanDowngrade({ }); } - return tx.campaign.update({ + await tx.campaign.update({ where: { id: campaign.id }, data: { status: CampaignStatus.paused }, - include: { workflow: true }, }); }); - - await scheduleTransactionalCampaign({ - campaign, - updatedCampaign, - }); } catch (error) { console.warn( `Failed to pause transactional campaign ${campaign.id} on plan downgrade:`, diff --git a/apps/web/lib/api/campaigns/schedule-campaigns.ts b/apps/web/lib/api/campaigns/schedule-campaigns.ts deleted file mode 100644 index a8748beef00..00000000000 --- a/apps/web/lib/api/campaigns/schedule-campaigns.ts +++ /dev/null @@ -1,189 +0,0 @@ -import { qstash } from "@/lib/cron"; -import { prisma } from "@/lib/prisma"; -import { PARTNER_ENROLLED_WORKFLOW_CRON } from "@/lib/zod/schemas/workflows"; -import { APP_DOMAIN_WITH_NGROK, log } from "@dub/utils"; -import { Campaign, CampaignType, Workflow } from "@prisma/client"; -import { isScheduledWorkflow } from "../workflows/utils"; - -type ScheduleCampaignProps = { - campaign: Campaign; - updatedCampaign: Campaign & { - workflow: Workflow | null; - }; -}; - -export const scheduleCampaign = async ({ - campaign, - updatedCampaign, -}: ScheduleCampaignProps) => { - if (campaign.type == CampaignType.marketing) { - return scheduleMarketingCampaign({ - campaign, - updatedCampaign, - }); - } - - if (campaign.type == CampaignType.transactional) { - return scheduleTransactionalCampaign({ - campaign, - updatedCampaign, - }); - } -}; - -export const deleteCampaignSchedule = async ( - campaign: Pick & { - workflow: Pick | null; - }, -) => { - if (campaign.type == CampaignType.marketing && campaign.qstashMessageId) { - try { - await qstash.messages.cancel(campaign.qstashMessageId); - } catch (error) { - console.warn( - `Failed to delete QStash message ${campaign.qstashMessageId}:`, - error, - ); - } - return; - } - - if ( - campaign.type == CampaignType.transactional && - campaign.workflow && - isScheduledWorkflow(campaign.workflow) - ) { - try { - await qstash.schedules.delete(campaign.workflow.id); - } catch (error) { - console.warn( - `Failed to delete QStash schedule ${campaign.workflow.id}:`, - error, - ); - } - } -}; - -// Schedule a marketing campaign -const scheduleMarketingCampaign = async ({ - campaign, - updatedCampaign, -}: ScheduleCampaignProps) => { - if (updatedCampaign.status === "draft") { - return; - } - - const scheduleChanged = - campaign.scheduledAt?.getTime() !== updatedCampaign.scheduledAt?.getTime(); - - const statusChanged = - (campaign.status === "draft" && updatedCampaign.status === "scheduled") || - (campaign.status === "scheduled" && updatedCampaign.status === "canceled"); - - if (!statusChanged && !scheduleChanged) { - return; - } - - let qstashMessageId = updatedCampaign.qstashMessageId; - - // Delete the existing message - if (campaign.qstashMessageId) { - try { - await qstash.messages.cancel(campaign.qstashMessageId); - qstashMessageId = null; - } catch (error) { - console.warn( - `Failed to delete QStash message ${campaign.qstashMessageId}:`, - error, - ); - } - } - - // Queue a new message - if (updatedCampaign.status === "scheduled") { - const notBefore = updatedCampaign.scheduledAt - ? Math.floor(updatedCampaign.scheduledAt.getTime() / 1000) - : null; - - try { - const response = await qstash.publishJSON({ - url: `${APP_DOMAIN_WITH_NGROK}/api/cron/campaigns/broadcast`, - method: "POST", - ...(notBefore && { notBefore }), - body: { - campaignId: campaign.id, - }, - }); - - qstashMessageId = response.messageId; - } catch (error) { - console.warn( - `Failed to queue QStash message for campaign ${campaign.id}:`, - error, - ); - } - } - - await prisma.campaign.update({ - where: { - id: campaign.id, - }, - data: { - qstashMessageId, - }, - }); -}; - -// Schedule a transactional campaign -export const scheduleTransactionalCampaign = async ({ - campaign, - updatedCampaign, -}: ScheduleCampaignProps) => { - if ( - !updatedCampaign.workflow || - !isScheduledWorkflow(updatedCampaign.workflow) - ) { - console.log(`No workflow found for campaign ${campaign.id}. Skipping...`); - return; - } - - const shouldSchedule = - (campaign.status === "draft" || campaign.status === "paused") && - updatedCampaign.status === "active"; - - if (shouldSchedule) { - try { - return await qstash.schedules.create({ - destination: `${APP_DOMAIN_WITH_NGROK}/api/cron/workflows/${updatedCampaign.workflow.id}`, - cron: PARTNER_ENROLLED_WORKFLOW_CRON, - scheduleId: updatedCampaign.workflow.id, - }); - } catch (error) { - // should never happen, but just in case - const errorMessage = `Failed to create QStash schedule ${updatedCampaign.workflow.id}: ${error}`; - console.warn(errorMessage); - await log({ - type: "errors", - message: errorMessage, - }); - } - return; - } - - const shouldDeleteSchedule = - campaign.status === "active" && updatedCampaign.status === "paused"; - - if (shouldDeleteSchedule) { - try { - return await qstash.schedules.delete(updatedCampaign.workflow.id); - } catch (error) { - // should never happen, but just in case - const errorMessage = `Failed to delete QStash schedule ${updatedCampaign.workflow.id}: ${error}`; - console.warn(errorMessage); - await log({ - type: "errors", - message: errorMessage, - }); - } - } -}; diff --git a/apps/web/lib/zod/schemas/workflows.ts b/apps/web/lib/zod/schemas/workflows.ts index c6daae06265..e41a238c3a7 100644 --- a/apps/web/lib/zod/schemas/workflows.ts +++ b/apps/web/lib/zod/schemas/workflows.ts @@ -5,9 +5,6 @@ import { } from "@/lib/api/workflows/operator-definitions"; import * as z from "zod/v4"; -// Cron for scheduled workflows that use partnerEnrolledDays conditions -export const PARTNER_ENROLLED_WORKFLOW_CRON = "0 */12 * * *"; // every 12 hours - export enum WORKFLOW_ACTION_TYPES { AwardBounty = "awardBounty", SendCampaign = "sendCampaign", diff --git a/apps/web/playwright/assert-local-database.ts b/apps/web/playwright/assert-local-database.ts index 897f6279ba8..7359c68faa6 100644 --- a/apps/web/playwright/assert-local-database.ts +++ b/apps/web/playwright/assert-local-database.ts @@ -16,7 +16,7 @@ function assertUrl( if (!value) { if (required) { throw new Error( - `${name} is not set. Playwright tests require a local database. ${LOCAL_URL_HINT}`, + `${name} is not set. A local database is required. ${LOCAL_URL_HINT}`, ); } return; @@ -27,13 +27,13 @@ function assertUrl( hostname = hostnameFromDatabaseUrl(value); } catch { throw new Error( - `${name} is not a valid URL. Playwright tests require a local database. ${LOCAL_URL_HINT}`, + `${name} is not a valid URL. A local database is required. ${LOCAL_URL_HINT}`, ); } if (!LOCAL_HOSTNAMES.has(hostname)) { throw new Error( - `Refusing to run Playwright tests: ${name} host "${hostname}" is not a local database. Only localhost / 127.0.0.1 / ::1 are allowed. ${LOCAL_URL_HINT}`, + `Refusing to proceed: ${name} host "${hostname}" is not a local database. Only localhost / 127.0.0.1 / ::1 are allowed. ${LOCAL_URL_HINT}`, ); } } diff --git a/apps/web/prisma/schema/campaign.prisma b/apps/web/prisma/schema/campaign.prisma index 7f94a4e0602..32192e4535a 100644 --- a/apps/web/prisma/schema/campaign.prisma +++ b/apps/web/prisma/schema/campaign.prisma @@ -18,28 +18,29 @@ enum CampaignStatus { } model Campaign { - id String @id + id String @id programId String - workflowId String? @unique + workflowId String? @unique userId String - qstashMessageId String? @unique + qstashMessageId String? @unique type CampaignType - status CampaignStatus @default(draft) + status CampaignStatus @default(draft) name String subject String - preview String? @db.Text + preview String? @db.Text from String? - bodyJson Json @db.Json + bodyJson Json @db.Json scheduledAt DateTime? - createdAt DateTime @default(now()) - updatedAt DateTime @updatedAt - program Program @relation(fields: [programId], references: [id]) - workflow Workflow? @relation(fields: [workflowId], references: [id], onDelete: Cascade) + createdAt DateTime @default(now()) + updatedAt DateTime @updatedAt + program Program @relation(fields: [programId], references: [id]) + workflow Workflow? @relation(fields: [workflowId], references: [id], onDelete: Cascade) groups CampaignGroup[] partnerTags CampaignPartnerTag[] emails NotificationEmail[] @@index(programId) + @@index([type, status, scheduledAt]) } model CampaignGroup { diff --git a/apps/web/scripts/dev/data.json b/apps/web/scripts/dev/data.json index d245191c346..ecf61cdc755 100644 --- a/apps/web/scripts/dev/data.json +++ b/apps/web/scripts/dev/data.json @@ -60,6 +60,13 @@ "verified": true } ], + "emailDomains": [ + { + "id": "dom_1KETZ919F83ZJH6A80EMAILDOM", + "slug": "getacme.link", + "status": "verified" + } + ], "folders": [ { "id": "fold_1K2J9DRWPPJ2F1RX53N92TSGB", diff --git a/apps/web/scripts/dev/seed.ts b/apps/web/scripts/dev/seed.ts index 4683c622ba5..9fafe6d266f 100644 --- a/apps/web/scripts/dev/seed.ts +++ b/apps/web/scripts/dev/seed.ts @@ -3,6 +3,7 @@ import { hashPassword } from "@/lib/auth/password"; import { prisma } from "@/lib/prisma"; import { Domain, + EmailDomain, Folder, Integration, Partner, @@ -18,7 +19,7 @@ import { import "dotenv-flow/config"; import fs from "fs"; import path from "path"; -import readline from "readline"; +import { assertLocalDatabaseEnv } from "../../playwright/assert-local-database"; type Workspace = Pick< Project, @@ -48,6 +49,8 @@ type Workspace = Pick< type DomainSeed = Pick; +type EmailDomainSeed = Pick; + type FolderSeed = Pick; type RewardSeed = Pick< @@ -118,6 +121,7 @@ type SeedData = { workspace: Workspace; users: WorkspaceUser[]; domains: DomainSeed[]; + emailDomains: EmailDomainSeed[]; folders: FolderSeed[]; rewards: RewardSeed[]; groups: GroupSeed[]; @@ -225,6 +229,33 @@ const createDomains = async (data: SeedData) => { console.log(`Created ${count} domains`); }; +// Create email domains +const createEmailDomains = async (data: SeedData) => { + const { emailDomains, workspace, program } = data; + + if (!emailDomains || emailDomains.length === 0) { + console.log("No email domains to insert"); + return; + } + + if (!program) { + console.log("Program is required to create email domains"); + return; + } + + const { count } = await prisma.emailDomain.createMany({ + data: emailDomains.map((emailDomain) => ({ + id: emailDomain.id, + slug: emailDomain.slug, + status: emailDomain.status, + workspaceId: workspace.id, + programId: program.id, + })), + }); + + console.log(`Created ${count} email domains`); +}; + // Create folders const createFolders = async (data: SeedData) => { const { folders, workspace } = data; @@ -566,44 +597,14 @@ const truncate = async () => { console.log("Database truncated successfully"); }; -// Ask for confirmation - requires typing "YES DELETE DATA" -const askConfirmation = (question: string): Promise => { - const rl = readline.createInterface({ - input: process.stdin, - output: process.stdout, - }); - - return new Promise((resolve) => { - rl.question( - `${question}\nType "YES DELETE DATA" to confirm: `, - (answer) => { - rl.close(); - resolve(answer === "YES DELETE DATA"); - }, - ); - }); -}; - async function main() { + assertLocalDatabaseEnv(); + // Check for --truncate flag // process.argv[0] = node, process.argv[1] = script path, process.argv[2+] = arguments const shouldTruncate = process.argv.slice(2).includes("--truncate"); if (shouldTruncate) { - console.log( - "\n⚠️ WARNING: This will delete ALL data from the database.\n", - ); - console.log("⚠️ Make sure you are NOT on production database!\n"); - const confirmed = await askConfirmation( - "Are you sure you want to delete ALL data from the database?", - ); - - if (!confirmed) { - console.log("\nTruncate canceled. Exiting..."); - process.exit(0); - } - - console.log("\n"); await truncate(); console.log("\n"); } @@ -615,6 +616,7 @@ async function main() { await createDomains(data); await createFolders(data); await createProgram(data); + await createEmailDomains(data); await createRewards(data); await createGroups(data); await createPartners(data); diff --git a/apps/web/vercel.json b/apps/web/vercel.json index 35588c8ed5d..e59dfb5eda2 100644 --- a/apps/web/vercel.json +++ b/apps/web/vercel.json @@ -91,6 +91,10 @@ { "path": "/api/cron/queue/retry", "schedule": "* * * * *" + }, + { + "path": "/api/cron/campaigns/queue-scheduled", + "schedule": "* * * * *" } ], "functions": { From c703c4db589436826ba6303d4ffc1a7811b30651 Mon Sep 17 00:00:00 2001 From: Kiran K Date: Tue, 18 Aug 2026 21:34:55 +0530 Subject: [PATCH 02/13] Publish due marketing campaigns immediately and retry failed first-run claims. --- .../(ee)/api/campaigns/[campaignId]/route.ts | 26 ++++++++- .../api/cron/campaigns/broadcast/route.ts | 20 ++++--- .../campaigns/marketing-campaign-broadcast.ts | 35 ++++++++++++ apps/web/prisma/schema/campaign.prisma | 2 +- .../delete-campaign-qstash-schedules.ts | 54 +++++++++++++++++++ 5 files changed, 129 insertions(+), 8 deletions(-) create mode 100644 apps/web/lib/api/campaigns/marketing-campaign-broadcast.ts create mode 100644 apps/web/scripts/migrations/delete-campaign-qstash-schedules.ts diff --git a/apps/web/app/(ee)/api/campaigns/[campaignId]/route.ts b/apps/web/app/(ee)/api/campaigns/[campaignId]/route.ts index eaa8feb6f0c..d4714d2fe2b 100644 --- a/apps/web/app/(ee)/api/campaigns/[campaignId]/route.ts +++ b/apps/web/app/(ee)/api/campaigns/[campaignId]/route.ts @@ -1,4 +1,5 @@ import { getCampaignOrThrow } from "@/lib/api/campaigns/get-campaign-or-throw"; +import { shouldEnqueueDueMarketingBroadcast } from "@/lib/api/campaigns/marketing-campaign-broadcast"; import { campaignEligibilityIncludes, transformCampaign, @@ -10,13 +11,15 @@ import { getDefaultProgramIdOrThrow } from "@/lib/api/programs/get-default-progr import { parseRequestBody } from "@/lib/api/utils"; import { validateWorkflowConditions } from "@/lib/api/workflows/validate-workflow-conditions"; import { withWorkspace } from "@/lib/auth"; +import { qstash } from "@/lib/cron"; import { prisma } from "@/lib/prisma"; import { CampaignSchema, updateCampaignSchema, } from "@/lib/zod/schemas/campaigns"; -import { arrayEqual, pluck } from "@dub/utils"; +import { APP_DOMAIN_WITH_NGROK, arrayEqual, pluck } from "@dub/utils"; import { PartnerGroup } from "@prisma/client"; +import { waitUntil } from "@vercel/functions"; import { NextResponse } from "next/server"; // GET /api/campaigns/[campaignId] - get an email campaign @@ -186,6 +189,27 @@ export const PATCH = withWorkspace( }); }); + if ( + shouldEnqueueDueMarketingBroadcast({ + previous: campaign, + next: updatedCampaign, + }) + ) { + waitUntil( + qstash.publishJSON({ + url: `${APP_DOMAIN_WITH_NGROK}/api/cron/campaigns/broadcast`, + deduplicationId: campaignId, + flowControl: { + key: "broadcast-marketing-campaign", + parallelism: 1, + }, + body: { + campaignId, + }, + }), + ); + } + return NextResponse.json( CampaignSchema.parse(transformCampaign(updatedCampaign)), ); diff --git a/apps/web/app/(ee)/api/cron/campaigns/broadcast/route.ts b/apps/web/app/(ee)/api/cron/campaigns/broadcast/route.ts index d4dcfa19491..e23a4904589 100644 --- a/apps/web/app/(ee)/api/cron/campaigns/broadcast/route.ts +++ b/apps/web/app/(ee)/api/cron/campaigns/broadcast/route.ts @@ -13,7 +13,7 @@ import { ACTIVE_ENROLLMENT_STATUSES } from "@/lib/zod/schemas/partners"; import { sendBatchEmail } from "@dub/email"; import CampaignEmail from "@dub/email/templates/campaign-email"; import { APP_DOMAIN_WITH_NGROK, chunk, log, pluck } from "@dub/utils"; -import { NotificationEmailType } from "@prisma/client"; +import { CampaignStatus, NotificationEmailType } from "@prisma/client"; import { differenceInMinutes } from "date-fns"; import * as z from "zod/v4"; import { logAndRespond } from "../../utils"; @@ -112,22 +112,30 @@ export async function POST(req: Request) { } // Claim the first run so leftover delayed messages / scanner retries - // cannot start a second broadcast. + // cannot start a second broadcast. A QStash retry of the message that + // already claimed (Upstash-Retried > 0) is allowed to continue. if (!startingAfter) { const claimed = await prisma.campaign.updateMany({ where: { id: campaignId, - status: "scheduled", + status: CampaignStatus.scheduled, }, data: { - status: "sending", + status: CampaignStatus.sending, }, }); if (claimed.count === 0) { - return logAndRespond( - `Campaign ${campaignId} broadcast already initiated. Skipping...`, + const retried = Number.parseInt( + req.headers.get("Upstash-Retried") ?? "0", + 10, ); + + if (!Number.isFinite(retried) || retried < 1) { + return logAndRespond( + `Campaign ${campaignId} broadcast already initiated. Skipping...`, + ); + } } } diff --git a/apps/web/lib/api/campaigns/marketing-campaign-broadcast.ts b/apps/web/lib/api/campaigns/marketing-campaign-broadcast.ts new file mode 100644 index 00000000000..795cb2fb6d5 --- /dev/null +++ b/apps/web/lib/api/campaigns/marketing-campaign-broadcast.ts @@ -0,0 +1,35 @@ +import { Campaign, CampaignStatus, CampaignType } from "@prisma/client"; + +type MarketingBroadcastCampaign = Pick< + Campaign, + "type" | "status" | "scheduledAt" +>; + +export function isDueMarketingCampaign({ + campaign, + now = new Date(), +}: { + campaign: MarketingBroadcastCampaign; + now?: Date; +}) { + return ( + campaign.type === CampaignType.marketing && + campaign.status === CampaignStatus.scheduled && + (!campaign.scheduledAt || campaign.scheduledAt <= now) + ); +} + +export function shouldEnqueueDueMarketingBroadcast({ + previous, + next, + now = new Date(), +}: { + previous: MarketingBroadcastCampaign; + next: MarketingBroadcastCampaign; + now?: Date; +}) { + return ( + isDueMarketingCampaign({ campaign: next, now }) && + !isDueMarketingCampaign({ campaign: previous, now }) + ); +} diff --git a/apps/web/prisma/schema/campaign.prisma b/apps/web/prisma/schema/campaign.prisma index 32192e4535a..92cc6907c10 100644 --- a/apps/web/prisma/schema/campaign.prisma +++ b/apps/web/prisma/schema/campaign.prisma @@ -22,7 +22,7 @@ model Campaign { programId String workflowId String? @unique userId String - qstashMessageId String? @unique + qstashMessageId String? @unique // TODO: remove this column type CampaignType status CampaignStatus @default(draft) name String diff --git a/apps/web/scripts/migrations/delete-campaign-qstash-schedules.ts b/apps/web/scripts/migrations/delete-campaign-qstash-schedules.ts new file mode 100644 index 00000000000..f054a62a7ea --- /dev/null +++ b/apps/web/scripts/migrations/delete-campaign-qstash-schedules.ts @@ -0,0 +1,54 @@ +import "dotenv-flow/config"; + +import { qstash } from "@/lib/cron"; +import { prisma } from "@/lib/prisma"; +import { CampaignType } from "@prisma/client"; + +// Remove the existing schedules from QStash +async function main() { + const campaigns = await prisma.campaign.findMany({ + where: { + type: CampaignType.transactional, + workflowId: { + not: null, + }, + }, + select: { + id: true, + workflowId: true, + status: true, + }, + }); + + console.table(campaigns); + + if (campaigns.length === 0) { + console.log("No transactional campaigns with a workflowId."); + return; + } + + let deleted = 0; + let skipped = 0; + let failed = 0; + + for (const { id, workflowId } of campaigns) { + if (!workflowId) { + continue; + } + + try { + await qstash.schedules.delete(workflowId); + deleted++; + } catch (error) { + console.error( + `Failed to delete schedule ${workflowId} (campaign ${id}):`, + error, + ); + failed++; + } + } + + console.log(`Done. deleted=${deleted} skipped=${skipped} failed=${failed}`); +} + +main(); From d5c1ca0b6cabe58b4d71845cfaa390772930e85b Mon Sep 17 00:00:00 2001 From: Kiran K Date: Wed, 19 Aug 2026 12:47:40 +0530 Subject: [PATCH 03/13] Isolate campaign queueing failures and cancel invalid-from broadcasts. --- .../(ee)/api/campaigns/[campaignId]/route.ts | 2 +- .../api/cron/campaigns/broadcast/route.ts | 70 ++++++++++++++++--- .../cron/campaigns/queue-scheduled/route.ts | 68 +++++++++++++----- .../pause-campaigns-on-plan-downgrade.ts | 50 ++++++++----- apps/web/lib/cron/enqueue-batch-jobs.ts | 2 +- packages/utils/src/functions/promises.ts | 10 +-- 6 files changed, 149 insertions(+), 53 deletions(-) diff --git a/apps/web/app/(ee)/api/campaigns/[campaignId]/route.ts b/apps/web/app/(ee)/api/campaigns/[campaignId]/route.ts index d4714d2fe2b..a5d6d3302d4 100644 --- a/apps/web/app/(ee)/api/campaigns/[campaignId]/route.ts +++ b/apps/web/app/(ee)/api/campaigns/[campaignId]/route.ts @@ -200,7 +200,7 @@ export const PATCH = withWorkspace( url: `${APP_DOMAIN_WITH_NGROK}/api/cron/campaigns/broadcast`, deduplicationId: campaignId, flowControl: { - key: "broadcast-marketing-campaign", + key: `broadcast-marketing-campaign-${campaignId}`, parallelism: 1, }, body: { diff --git a/apps/web/app/(ee)/api/cron/campaigns/broadcast/route.ts b/apps/web/app/(ee)/api/cron/campaigns/broadcast/route.ts index e23a4904589..c2540b84fd4 100644 --- a/apps/web/app/(ee)/api/cron/campaigns/broadcast/route.ts +++ b/apps/web/app/(ee)/api/cron/campaigns/broadcast/route.ts @@ -3,7 +3,7 @@ import { renderCampaignEmailHTML } from "@/lib/api/campaigns/render-campaign-ema import { campaignEligibilityIncludes } from "@/lib/api/campaigns/transform-campaign"; import { validateCampaignFromAddress } from "@/lib/api/campaigns/validate-campaign"; import { createId } from "@/lib/api/create-id"; -import { handleAndReturnErrorResponse } from "@/lib/api/errors"; +import { DubApiError, handleAndReturnErrorResponse } from "@/lib/api/errors"; import { qstash } from "@/lib/cron"; import { verifyQstashSignature } from "@/lib/cron/verify-qstash"; import { resolveCampaignFromAddress } from "@/lib/email/parse-campaign-from-address"; @@ -13,7 +13,12 @@ import { ACTIVE_ENROLLMENT_STATUSES } from "@/lib/zod/schemas/partners"; import { sendBatchEmail } from "@dub/email"; import CampaignEmail from "@dub/email/templates/campaign-email"; import { APP_DOMAIN_WITH_NGROK, chunk, log, pluck } from "@dub/utils"; -import { CampaignStatus, NotificationEmailType } from "@prisma/client"; +import { + Campaign, + CampaignStatus, + EmailDomain, + NotificationEmailType, +} from "@prisma/client"; import { differenceInMinutes } from "date-fns"; import * as z from "zod/v4"; import { logAndRespond } from "../../utils"; @@ -103,14 +108,6 @@ export async function POST(req: Request) { const program = campaign.program; - // TODO: We should make the from address required. There are existing campaign without from address - if (campaign.from) { - validateCampaignFromAddress({ - campaign, - emailDomains: program.emailDomains, - }); - } - // Claim the first run so leftover delayed messages / scanner retries // cannot start a second broadcast. A QStash retry of the message that // already claimed (Upstash-Retried > 0) is allowed to continue. @@ -139,6 +136,15 @@ export async function POST(req: Request) { } } + const invalidFromResponse = await cancelCampaignIfInvalidFromAddress({ + campaign, + emailDomains: program.emailDomains, + }); + + if (invalidFromResponse) { + return invalidFromResponse; + } + const campaignGroupIds = pluck(campaign.groups, "groupId"); const campaignPartnerTagIds = pluck(campaign.partnerTags, "partnerTagId"); @@ -373,3 +379,47 @@ export async function POST(req: Request) { return handleAndReturnErrorResponse(error); } } + +async function cancelCampaignIfInvalidFromAddress({ + campaign, + emailDomains, +}: { + campaign: Pick; + emailDomains: Pick[]; +}) { + if (!campaign.from) { + return; + } + + try { + validateCampaignFromAddress({ + campaign, + emailDomains, + }); + } catch (error) { + if (!(error instanceof DubApiError)) { + throw error; + } + + await prisma.campaign.updateMany({ + where: { + id: campaign.id, + status: { + in: [CampaignStatus.scheduled, CampaignStatus.sending], + }, + }, + data: { + status: CampaignStatus.canceled, + }, + }); + + await log({ + type: "errors", + message: `Campaign ${campaign.id} canceled: ${error.message}`, + }); + + return logAndRespond( + `Campaign ${campaign.id} canceled: invalid from address.`, + ); + } +} diff --git a/apps/web/app/(ee)/api/cron/campaigns/queue-scheduled/route.ts b/apps/web/app/(ee)/api/cron/campaigns/queue-scheduled/route.ts index 9008e80021a..6f3d4ebf069 100644 --- a/apps/web/app/(ee)/api/cron/campaigns/queue-scheduled/route.ts +++ b/apps/web/app/(ee)/api/cron/campaigns/queue-scheduled/route.ts @@ -1,8 +1,14 @@ import { isScheduledWorkflow } from "@/lib/api/workflows/utils"; -import { CRON_BATCH_SIZE, qstash } from "@/lib/cron"; +import { CRON_BATCH_SIZE } from "@/lib/cron"; +import { enqueueBatchJobs } from "@/lib/cron/enqueue-batch-jobs"; import { withCron } from "@/lib/cron/with-cron"; import { prisma } from "@/lib/prisma"; -import { APP_DOMAIN_WITH_NGROK } from "@dub/utils"; +import { + APP_DOMAIN_WITH_NGROK, + isRejected, + log, + serializeError, +} from "@dub/utils"; import { CampaignStatus, CampaignType } from "@prisma/client"; import { logAndRespond } from "../../utils"; @@ -14,19 +20,46 @@ export const maxDuration = 600; export const GET = withCron(async () => { const now = new Date(); - await Promise.all([ + const [transactional, marketing] = await Promise.allSettled([ queueTransactionalCampaigns(now), queueMarketingCampaigns(now), ]); - return logAndRespond("Finished the campaigns queueing process."); + const failures: string[] = []; + + if (isRejected(transactional)) { + failures.push(`transactional: ${serializeError(transactional.reason)}`); + } + + if (isRejected(marketing)) { + failures.push(`marketing: ${serializeError(marketing.reason)}`); + } + + if (failures.length > 0) { + const message = `Campaign queueing partially failed: ${failures.join("; ")}`; + await log({ type: "errors", message }); + return logAndRespond(message, { logLevel: "error" }); + } + + const transactionalQueued = + transactional.status === "fulfilled" ? transactional.value : 0; + const marketingQueued = + marketing.status === "fulfilled" ? marketing.value : 0; + + if (transactionalQueued + marketingQueued === 0) { + return logAndRespond("No campaigns to queue."); + } + + return logAndRespond( + `Queued ${marketingQueued} marketing and ${transactionalQueued} transactional campaigns.`, + ); }); async function queueTransactionalCampaigns(now: Date) { // Matches the 12h enrollment window in executeSendCampaignWorkflow. - if (now.getUTCMinutes() !== 0 || now.getUTCHours() % 12 !== 0) { - console.log("[Transactional] Not the 12h tick, skipping campaigns."); - return; + // 5-minute window absorbs Vercel cron jitter; QStash dedup (10 min) collapses extra publishes. + if (!isTransactionalTick(now)) { + return 0; } let queued = 0; @@ -58,7 +91,6 @@ async function queueTransactionalCampaigns(now: Date) { }); if (campaigns.length === 0) { - console.log("[Transactional] No more campaigns to queue."); break; } @@ -69,10 +101,11 @@ async function queueTransactionalCampaigns(now: Date) { ); if (scheduledWorkflows.length > 0) { - await qstash.batchJSON( + await enqueueBatchJobs( scheduledWorkflows.map((workflow) => ({ url: `${APP_DOMAIN_WITH_NGROK}/api/cron/workflows/${workflow.id}`, deduplicationId: workflow.id, + label: "execute-scheduled-workflow", flowControl: { key: "execute-scheduled-workflow", parallelism: 10, @@ -87,41 +120,42 @@ async function queueTransactionalCampaigns(now: Date) { page++; } - console.log(`[Transactional] Queued ${queued} campaigns.`); + return queued; } async function queueMarketingCampaigns(now: Date) { let queued = 0; - let page = 0; + let lastCampaignId: string | undefined; while (true) { const campaigns = await prisma.campaign.findMany({ where: { type: CampaignType.marketing, + // Do not reclaim `sending` campaigns; failures Slack-alert and we resume them manually. status: CampaignStatus.scheduled, OR: [{ scheduledAt: null }, { scheduledAt: { lte: now } }], + ...(lastCampaignId && { id: { gt: lastCampaignId } }), }, select: { id: true, }, take: CRON_BATCH_SIZE, - skip: page * CRON_BATCH_SIZE, orderBy: { id: "asc", }, }); if (campaigns.length === 0) { - console.log("[Marketing] No more campaigns to queue."); break; } - await qstash.batchJSON( + await enqueueBatchJobs( campaigns.map((campaign) => ({ url: `${APP_DOMAIN_WITH_NGROK}/api/cron/campaigns/broadcast`, deduplicationId: campaign.id, + label: "broadcast-marketing-campaign", flowControl: { - key: "broadcast-marketing-campaign", + key: `broadcast-marketing-campaign-${campaign.id}`, parallelism: 1, }, body: { @@ -131,8 +165,8 @@ async function queueMarketingCampaigns(now: Date) { ); queued += campaigns.length; - page++; + lastCampaignId = campaigns[campaigns.length - 1].id; } - console.log(`[Marketing] Queued ${queued} campaigns.`); + return queued; } diff --git a/apps/web/lib/api/campaigns/pause-campaigns-on-plan-downgrade.ts b/apps/web/lib/api/campaigns/pause-campaigns-on-plan-downgrade.ts index 89cf02a91ec..889988229d7 100644 --- a/apps/web/lib/api/campaigns/pause-campaigns-on-plan-downgrade.ts +++ b/apps/web/lib/api/campaigns/pause-campaigns-on-plan-downgrade.ts @@ -6,6 +6,7 @@ export async function pauseOrCancelCampaignsForProgramOnPlanDowngrade({ }: { programId: string; }): Promise { + // Cancel marketing campaigns await prisma.campaign.updateMany({ where: { programId, @@ -25,28 +26,39 @@ export async function pauseOrCancelCampaignsForProgramOnPlanDowngrade({ type: CampaignType.transactional, status: CampaignStatus.active, }, + select: { + workflowId: true, + }, }); - for (const campaign of transactionalCampaigns) { - try { - await prisma.$transaction(async (tx) => { - if (campaign.workflowId) { - await tx.workflow.update({ - where: { id: campaign.workflowId }, - data: { disabledAt: new Date() }, - }); - } + const workflowIds = transactionalCampaigns.flatMap((campaign) => + campaign.workflowId ? [campaign.workflowId] : [], + ); - await tx.campaign.update({ - where: { id: campaign.id }, - data: { status: CampaignStatus.paused }, - }); + await prisma.$transaction(async (tx) => { + if (workflowIds.length > 0) { + await tx.workflow.updateMany({ + where: { + id: { + in: workflowIds, + }, + disabledAt: null, + }, + data: { + disabledAt: new Date(), + }, }); - } catch (error) { - console.warn( - `Failed to pause transactional campaign ${campaign.id} on plan downgrade:`, - error, - ); } - } + + await tx.campaign.updateMany({ + where: { + programId, + type: CampaignType.transactional, + status: CampaignStatus.active, + }, + data: { + status: CampaignStatus.paused, + }, + }); + }); } diff --git a/apps/web/lib/cron/enqueue-batch-jobs.ts b/apps/web/lib/cron/enqueue-batch-jobs.ts index 00c69a65556..f6c8a8e259f 100644 --- a/apps/web/lib/cron/enqueue-batch-jobs.ts +++ b/apps/web/lib/cron/enqueue-batch-jobs.ts @@ -3,7 +3,7 @@ import type { PublishBatchRequest } from "@upstash/qstash"; import { qstash } from "."; type EnqueueBatchJobsProps = PublishBatchRequest & { - queueName: + queueName?: | "ban-partner" | "send-partner-summary" | "create-discount-code" diff --git a/packages/utils/src/functions/promises.ts b/packages/utils/src/functions/promises.ts index adf3842c648..91799818ef0 100644 --- a/packages/utils/src/functions/promises.ts +++ b/packages/utils/src/functions/promises.ts @@ -6,6 +6,10 @@ export const isRejected = ( p: PromiseSettledResult, ): p is PromiseRejectedResult => p.status === "rejected"; +export function serializeError(error: unknown) { + return error instanceof Error ? error.message : String(error); +} + export function logPromiseResults( results: PromiseSettledResult[], options?: { @@ -26,11 +30,7 @@ export function logPromiseResults( console.log(`${label}${id} succeeded.`); } else { failureCount++; - const reason = - result.reason instanceof Error - ? result.reason.message - : String(result.reason); - console.error(`${label}${id} failed: ${reason}`); + console.error(`${label}${id} failed: ${serializeError(result.reason)}`); } } From 723005ce5a7147448530c6f49e39c7c1de84b55f Mon Sep 17 00:00:00 2001 From: Kiran K Date: Wed, 19 Aug 2026 12:56:37 +0530 Subject: [PATCH 04/13] Gate transactional campaign queueing to the 12h UTC window and mark stuck sending campaigns as sent. --- .../cron/campaigns/queue-scheduled/route.ts | 6 +++ .../delete-campaign-qstash-schedules.ts | 51 ++++++++++++++++++- 2 files changed, 55 insertions(+), 2 deletions(-) diff --git a/apps/web/app/(ee)/api/cron/campaigns/queue-scheduled/route.ts b/apps/web/app/(ee)/api/cron/campaigns/queue-scheduled/route.ts index 6f3d4ebf069..99cd26be481 100644 --- a/apps/web/app/(ee)/api/cron/campaigns/queue-scheduled/route.ts +++ b/apps/web/app/(ee)/api/cron/campaigns/queue-scheduled/route.ts @@ -55,6 +55,12 @@ export const GET = withCron(async () => { ); }); +// First 5 minutes of 00:00/12:00 UTC. QStash dedup lasts 10 minutes, so this +// absorbs Vercel cron jitter without leaking a second publish after expiry. +function isTransactionalTick(now: Date) { + return now.getUTCHours() % 12 === 0 && now.getUTCMinutes() < 5; +} + async function queueTransactionalCampaigns(now: Date) { // Matches the 12h enrollment window in executeSendCampaignWorkflow. // 5-minute window absorbs Vercel cron jitter; QStash dedup (10 min) collapses extra publishes. diff --git a/apps/web/scripts/migrations/delete-campaign-qstash-schedules.ts b/apps/web/scripts/migrations/delete-campaign-qstash-schedules.ts index f054a62a7ea..3eeed52a2c6 100644 --- a/apps/web/scripts/migrations/delete-campaign-qstash-schedules.ts +++ b/apps/web/scripts/migrations/delete-campaign-qstash-schedules.ts @@ -1,11 +1,17 @@ +// @ts-ignore import "dotenv-flow/config"; import { qstash } from "@/lib/cron"; import { prisma } from "@/lib/prisma"; -import { CampaignType } from "@prisma/client"; +import { CampaignStatus, CampaignType } from "@prisma/client"; -// Remove the existing schedules from QStash async function main() { + await deleteTransactionalCampaignQstashSchedules(); + await markStuckSendingCampaignsAsSent(); +} + +// Remove the existing schedules from QStash +async function deleteTransactionalCampaignQstashSchedules() { const campaigns = await prisma.campaign.findMany({ where: { type: CampaignType.transactional, @@ -51,4 +57,45 @@ async function main() { console.log(`Done. deleted=${deleted} skipped=${skipped} failed=${failed}`); } +// Marketing campaigns that started sending in mid-March 2026 and never +// flipped to `sent` after the next QStash batch failed (~665 / ~673 emails already +// delivered). Mark these campaigns as `sent` so they leave the in-progress UI and a stray QStash retry cannot continue. +async function markStuckSendingCampaignsAsSent() { + const stuckSendingCampaignIds = [ + "cmp_1KKSXA2G25J2KBNGZQERMAPEX", + "cmp_1KKZ458FAK2PY9W9677TT3VS2", + ]; + + const campaigns = await prisma.campaign.findMany({ + where: { + id: { + in: stuckSendingCampaignIds, + }, + }, + select: { + id: true, + name: true, + status: true, + scheduledAt: true, + updatedAt: true, + }, + }); + + console.table(campaigns); + + const result = await prisma.campaign.updateMany({ + where: { + id: { + in: stuckSendingCampaignIds, + }, + status: CampaignStatus.sending, + }, + data: { + status: CampaignStatus.sent, + }, + }); + + console.log(`Marked ${result.count} stuck sending campaigns as sent.`); +} + main(); From 5153b7668e69c7b66ff7732dbe2827aadcda7e13 Mon Sep 17 00:00:00 2001 From: Kiran K Date: Wed, 19 Aug 2026 12:59:26 +0530 Subject: [PATCH 05/13] Identify in-flight broadcasts by QStash message id instead of retry count. --- .../app/(ee)/api/cron/campaigns/broadcast/route.ts | 14 ++++++-------- apps/web/prisma/schema/campaign.prisma | 2 +- 2 files changed, 7 insertions(+), 9 deletions(-) diff --git a/apps/web/app/(ee)/api/cron/campaigns/broadcast/route.ts b/apps/web/app/(ee)/api/cron/campaigns/broadcast/route.ts index c2540b84fd4..fdc7a75cb32 100644 --- a/apps/web/app/(ee)/api/cron/campaigns/broadcast/route.ts +++ b/apps/web/app/(ee)/api/cron/campaigns/broadcast/route.ts @@ -109,9 +109,11 @@ export async function POST(req: Request) { const program = campaign.program; // Claim the first run so leftover delayed messages / scanner retries - // cannot start a second broadcast. A QStash retry of the message that - // already claimed (Upstash-Retried > 0) is allowed to continue. + // cannot start a second broadcast. A QStash retry of the claiming + // message (matching qstashMessageId) is allowed to continue. if (!startingAfter) { + const messageId = req.headers.get("Upstash-Message-Id"); + const claimed = await prisma.campaign.updateMany({ where: { id: campaignId, @@ -119,16 +121,12 @@ export async function POST(req: Request) { }, data: { status: CampaignStatus.sending, + qstashMessageId: messageId, }, }); if (claimed.count === 0) { - const retried = Number.parseInt( - req.headers.get("Upstash-Retried") ?? "0", - 10, - ); - - if (!Number.isFinite(retried) || retried < 1) { + if (!messageId || campaign.qstashMessageId !== messageId) { return logAndRespond( `Campaign ${campaignId} broadcast already initiated. Skipping...`, ); diff --git a/apps/web/prisma/schema/campaign.prisma b/apps/web/prisma/schema/campaign.prisma index 92cc6907c10..5b61f7c49a2 100644 --- a/apps/web/prisma/schema/campaign.prisma +++ b/apps/web/prisma/schema/campaign.prisma @@ -22,7 +22,7 @@ model Campaign { programId String workflowId String? @unique userId String - qstashMessageId String? @unique // TODO: remove this column + qstashMessageId String? @unique // Claiming QStash message id for the in-flight broadcast type CampaignType status CampaignStatus @default(draft) name String From 1bd1522c111119089e4ebcb61cd920ab301411a3 Mon Sep 17 00:00:00 2001 From: Kiran K Date: Wed, 19 Aug 2026 12:59:32 +0530 Subject: [PATCH 06/13] Update campaign.prisma --- apps/web/prisma/schema/campaign.prisma | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/apps/web/prisma/schema/campaign.prisma b/apps/web/prisma/schema/campaign.prisma index 5b61f7c49a2..32192e4535a 100644 --- a/apps/web/prisma/schema/campaign.prisma +++ b/apps/web/prisma/schema/campaign.prisma @@ -22,7 +22,7 @@ model Campaign { programId String workflowId String? @unique userId String - qstashMessageId String? @unique // Claiming QStash message id for the in-flight broadcast + qstashMessageId String? @unique type CampaignType status CampaignStatus @default(draft) name String From b91022bf73743740f94fafa4646a2784359014c4 Mon Sep 17 00:00:00 2001 From: Kiran K Date: Wed, 19 Aug 2026 13:07:11 +0530 Subject: [PATCH 07/13] Update route.ts --- .../app/(ee)/api/cron/campaigns/queue-scheduled/route.ts | 7 ++++--- 1 file changed, 4 insertions(+), 3 deletions(-) diff --git a/apps/web/app/(ee)/api/cron/campaigns/queue-scheduled/route.ts b/apps/web/app/(ee)/api/cron/campaigns/queue-scheduled/route.ts index 99cd26be481..753a7deed4c 100644 --- a/apps/web/app/(ee)/api/cron/campaigns/queue-scheduled/route.ts +++ b/apps/web/app/(ee)/api/cron/campaigns/queue-scheduled/route.ts @@ -69,7 +69,7 @@ async function queueTransactionalCampaigns(now: Date) { } let queued = 0; - let page = 0; + let lastCampaignId: string | undefined; while (true) { const campaigns = await prisma.campaign.findMany({ @@ -79,8 +79,10 @@ async function queueTransactionalCampaigns(now: Date) { workflow: { disabledAt: null, }, + ...(lastCampaignId && { id: { gt: lastCampaignId } }), }, select: { + id: true, workflow: { select: { id: true, @@ -90,7 +92,6 @@ async function queueTransactionalCampaigns(now: Date) { }, }, take: CRON_BATCH_SIZE, - skip: page * CRON_BATCH_SIZE, orderBy: { id: "asc", }, @@ -123,7 +124,7 @@ async function queueTransactionalCampaigns(now: Date) { queued += scheduledWorkflows.length; } - page++; + lastCampaignId = campaigns[campaigns.length - 1].id; } return queued; From 745bd57704950186a6b2baf6205ff5e423df2e92 Mon Sep 17 00:00:00 2001 From: Kiran K Date: Wed, 19 Aug 2026 13:35:38 +0530 Subject: [PATCH 08/13] Drop QStash dedup ids on marketing broadcasts so a later reschedule is not collapsed into the earlier publish. --- apps/web/app/(ee)/api/campaigns/[campaignId]/route.ts | 1 - apps/web/app/(ee)/api/cron/campaigns/queue-scheduled/route.ts | 1 - 2 files changed, 2 deletions(-) diff --git a/apps/web/app/(ee)/api/campaigns/[campaignId]/route.ts b/apps/web/app/(ee)/api/campaigns/[campaignId]/route.ts index a5d6d3302d4..288cf9d3d6b 100644 --- a/apps/web/app/(ee)/api/campaigns/[campaignId]/route.ts +++ b/apps/web/app/(ee)/api/campaigns/[campaignId]/route.ts @@ -198,7 +198,6 @@ export const PATCH = withWorkspace( waitUntil( qstash.publishJSON({ url: `${APP_DOMAIN_WITH_NGROK}/api/cron/campaigns/broadcast`, - deduplicationId: campaignId, flowControl: { key: `broadcast-marketing-campaign-${campaignId}`, parallelism: 1, diff --git a/apps/web/app/(ee)/api/cron/campaigns/queue-scheduled/route.ts b/apps/web/app/(ee)/api/cron/campaigns/queue-scheduled/route.ts index 753a7deed4c..f9f93be5b09 100644 --- a/apps/web/app/(ee)/api/cron/campaigns/queue-scheduled/route.ts +++ b/apps/web/app/(ee)/api/cron/campaigns/queue-scheduled/route.ts @@ -159,7 +159,6 @@ async function queueMarketingCampaigns(now: Date) { await enqueueBatchJobs( campaigns.map((campaign) => ({ url: `${APP_DOMAIN_WITH_NGROK}/api/cron/campaigns/broadcast`, - deduplicationId: campaign.id, label: "broadcast-marketing-campaign", flowControl: { key: `broadcast-marketing-campaign-${campaign.id}`, From 639ba9e0f13dcfe065579cad16796c1765cdbd13 Mon Sep 17 00:00:00 2001 From: Kiran K Date: Wed, 19 Aug 2026 14:13:44 +0530 Subject: [PATCH 09/13] Persist numeric campaign logic edits by marking the form dirty. --- .../transactional-campaign-logic.tsx | 126 +++++++----------- 1 file changed, 51 insertions(+), 75 deletions(-) diff --git a/apps/web/app/app.dub.co/(dashboard)/[slug]/(ee)/program/campaigns/[campaignId]/transactional-campaign-logic.tsx b/apps/web/app/app.dub.co/(dashboard)/[slug]/(ee)/program/campaigns/[campaignId]/transactional-campaign-logic.tsx index 93b48192dd0..7004dda4912 100644 --- a/apps/web/app/app.dub.co/(dashboard)/[slug]/(ee)/program/campaigns/[campaignId]/transactional-campaign-logic.tsx +++ b/apps/web/app/app.dub.co/(dashboard)/[slug]/(ee)/program/campaigns/[campaignId]/transactional-campaign-logic.tsx @@ -9,17 +9,16 @@ import { type SendCampaignAttributeKey, } from "@/lib/api/workflows/send-campaign/schema"; import { satisfiesExclusiveAttributeRules } from "@/lib/api/workflows/utils"; -import { handleMoneyInputChange, handleMoneyKeyDown } from "@/lib/form-utils"; import { DurationPopoverContent } from "@/ui/shared/duration-popover-content"; import { InlineBadgePopover, - InlineBadgePopoverContext, + InlineBadgePopoverAmountInput, InlineBadgePopoverMenu, } from "@/ui/shared/inline-badge-popover"; import { Button } from "@dub/ui"; import { Xmark } from "@dub/ui/icons"; -import { cn, currencyFormatter, pluralize } from "@dub/utils"; -import { useContext, useEffect, useMemo, useRef } from "react"; +import { currencyFormatter, pluralize } from "@dub/utils"; +import { ChangeEvent, useEffect, useMemo, useRef } from "react"; import { Controller, useFieldArray } from "react-hook-form"; import { useCampaignFormContext } from "./campaign-form-context"; @@ -152,8 +151,11 @@ function ConditionRow({ setValue( `triggerConditions.${index}.value`, attribute === "partnerJoined" ? 0 : (null as any), + { shouldDirty: true }, ); - setValue(`triggerConditions.${index}.operator`, "gte"); + setValue(`triggerConditions.${index}.operator`, "gte", { + shouldDirty: true, + }); } prevAttributeRef.current = attribute; @@ -162,7 +164,7 @@ function ConditionRow({ // Ensure partnerJoined always has value 0 useEffect(() => { if (attribute === "partnerJoined" && value !== 0) { - setValue(`triggerConditions.${index}.value`, 0); + setValue(`triggerConditions.${index}.value`, 0, { shouldDirty: true }); } }, [attribute, value, index, setValue]); @@ -235,7 +237,7 @@ function ConditionRow({ {config.inputType === "dropdown" ? ( ) : ( - + )} )} @@ -299,83 +301,57 @@ function DropdownValueInput({ function ValueInput({ index, config, - value, }: { index: number; config: { inputType?: string }; - value: number | null | undefined; }) { - const { watch, setValue } = useCampaignFormContext(); - const { setIsOpen } = useContext(InlineBadgePopoverContext); - - const storedValue = watch(`triggerConditions.${index}.value`); - + const { control } = useCampaignFormContext(); const isCurrency = config.inputType === "currency"; - const displayValue = - isCurrency && storedValue ? storedValue / 100 : storedValue; - - const hasValue = value !== null && value !== undefined; - return ( - -
- {isCurrency && ( - - $ - - )} - { - const nextValue = e.target.value; - if (nextValue === "") { - setValue(`triggerConditions.${index}.value`, null as any); - } else { - const numValue = +nextValue; - setValue( - `triggerConditions.${index}.value`, - isCurrency ? Math.round(numValue * 100) : numValue, - ); - } + { + const storedValue = field.value; + const displayValue = + isCurrency && storedValue ? storedValue / 100 : storedValue; + const hasValue = storedValue !== null && storedValue !== undefined; - if (isCurrency) { - handleMoneyInputChange(e); - } - }} - onKeyDown={(e) => { - if (e.key === "Enter") { - e.preventDefault(); - setIsOpen(false); - return; + return ( + + ) => { + const nextValue = e.target.value; - if (isCurrency) { - handleMoneyKeyDown(e); - } - }} - /> - {isCurrency && ( - - USD - - )} -
-
+ if (nextValue === "") { + field.onChange(null as unknown as number); + return; + } + + const numValue = +nextValue; + field.onChange( + isCurrency ? Math.round(numValue * 100) : numValue, + ); + }} + onBlur={field.onBlur} + /> + + ); + }} + /> ); } From c8d629d6d1a6bc912beaf9621b1c4716445d2999 Mon Sep 17 00:00:00 2001 From: Kiran K Date: Wed, 19 Aug 2026 14:29:24 +0530 Subject: [PATCH 10/13] Skip claiming broadcasts that have been rescheduled into the future. --- apps/web/app/(ee)/api/cron/campaigns/broadcast/route.ts | 1 + 1 file changed, 1 insertion(+) diff --git a/apps/web/app/(ee)/api/cron/campaigns/broadcast/route.ts b/apps/web/app/(ee)/api/cron/campaigns/broadcast/route.ts index fdc7a75cb32..9cae722a3b5 100644 --- a/apps/web/app/(ee)/api/cron/campaigns/broadcast/route.ts +++ b/apps/web/app/(ee)/api/cron/campaigns/broadcast/route.ts @@ -118,6 +118,7 @@ export async function POST(req: Request) { where: { id: campaignId, status: CampaignStatus.scheduled, + OR: [{ scheduledAt: null }, { scheduledAt: { lte: new Date() } }], }, data: { status: CampaignStatus.sending, From 88a781e815c769bd7d2a1d924f12a24862042ffa Mon Sep 17 00:00:00 2001 From: Marcus Farrell Date: Wed, 19 Aug 2026 20:32:41 -0700 Subject: [PATCH 11/13] Partner interest empty state --- .../(dashboard)/profile/about-you-form.tsx | 34 +++++++++++++++++-- 1 file changed, 31 insertions(+), 3 deletions(-) diff --git a/apps/web/app/(ee)/partners.dub.co/(dashboard)/profile/about-you-form.tsx b/apps/web/app/(ee)/partners.dub.co/(dashboard)/profile/about-you-form.tsx index 8111a245215..e2f9d321bb4 100644 --- a/apps/web/app/(ee)/partners.dub.co/(dashboard)/profile/about-you-form.tsx +++ b/apps/web/app/(ee)/partners.dub.co/(dashboard)/profile/about-you-form.tsx @@ -9,6 +9,7 @@ import { PartnerProps } from "@/lib/types"; import { MAX_PARTNER_DESCRIPTION_LENGTH } from "@/lib/zod/schemas/partners"; import { MaxCharactersCounter } from "@/ui/shared/max-characters-counter"; import { Button, RadioGroup, RadioGroupItem, useEnterSubmit } from "@dub/ui"; +import { Plus } from "@dub/ui/icons"; import { cn } from "@dub/utils"; import { IndustryInterest, MonthlyTraffic } from "@prisma/client"; import { useAction } from "next-safe-action/hooks"; @@ -151,12 +152,39 @@ export function AboutYouForm({ partner }: { partner?: PartnerProps }) { )) : [...Array(3)].map((_, idx) => ( -
setShowIndustryInterestsModal(true)} className={cn( - "border-border-subtle h-11 w-32 rounded-full border border-dashed bg-white", + "relative flex h-11 w-32 items-center justify-center rounded-full bg-white", + !disabled && + "transition-colors hover:bg-neutral-50/60", + disabled && "cursor-not-allowed", )} - /> + > + + + ))}