Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
21 commits
Select commit Hold shift + click to select a range
7dd253f
Queue due campaigns from a cron instead of per-campaign QStash schedules
devkiran Aug 18, 2026
c703c4d
Publish due marketing campaigns immediately and retry failed first-ru…
devkiran Aug 18, 2026
fbe0d4e
Merge branch 'main' into queue-scheduled-campaigns
steven-tey Aug 18, 2026
cb4c29c
Merge branch 'main' into queue-scheduled-campaigns
devkiran Aug 19, 2026
d5c1ca0
Isolate campaign queueing failures and cancel invalid-from broadcasts.
devkiran Aug 19, 2026
723005c
Gate transactional campaign queueing to the 12h UTC window and mark s…
devkiran Aug 19, 2026
5153b76
Identify in-flight broadcasts by QStash message id instead of retry c…
devkiran Aug 19, 2026
1bd1522
Update campaign.prisma
devkiran Aug 19, 2026
34c50e8
Merge branch 'queue-scheduled-campaigns' of https://github.com/dubinc…
devkiran Aug 19, 2026
5abf7fe
Merge branch 'main' into queue-scheduled-campaigns
devkiran Aug 19, 2026
b91022b
Update route.ts
devkiran Aug 19, 2026
e762909
Merge branch 'queue-scheduled-campaigns' of https://github.com/dubinc…
devkiran Aug 19, 2026
745bd57
Drop QStash dedup ids on marketing broadcasts so a later reschedule i…
devkiran Aug 19, 2026
639ba9e
Persist numeric campaign logic edits by marking the form dirty.
devkiran Aug 19, 2026
c8d629d
Skip claiming broadcasts that have been rescheduled into the future.
devkiran Aug 19, 2026
88a781e
Partner interest empty state
marcusljf Aug 20, 2026
4128db5
Add aria-label
marcusljf Aug 20, 2026
860f32c
Merge pull request #4370 from dubinc/partner-interests
steven-tey Aug 20, 2026
eff8f6d
Merge branch 'main' into queue-scheduled-campaigns
steven-tey Aug 20, 2026
7d2299e
Merge pull request #4351 from dubinc/queue-scheduled-campaigns
steven-tey Aug 20, 2026
4308d00
remove tremendous flag
steven-tey Aug 20, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
44 changes: 22 additions & 22 deletions apps/web/app/(ee)/api/campaigns/[campaignId]/route.ts
Original file line number Diff line number Diff line change
@@ -1,8 +1,5 @@
import { getCampaignOrThrow } from "@/lib/api/campaigns/get-campaign-or-throw";
import {
deleteCampaignSchedule,
scheduleCampaign,
} from "@/lib/api/campaigns/schedule-campaigns";
import { shouldEnqueueDueMarketingBroadcast } from "@/lib/api/campaigns/marketing-campaign-broadcast";
import {
campaignEligibilityIncludes,
transformCampaign,
Expand All @@ -14,12 +11,13 @@ 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";
Expand Down Expand Up @@ -191,12 +189,25 @@ export const PATCH = withWorkspace(
});
});

waitUntil(
scheduleCampaign({
campaign,
updatedCampaign,
}),
);
if (
shouldEnqueueDueMarketingBroadcast({
previous: campaign,
next: updatedCampaign,
})
) {
waitUntil(
qstash.publishJSON({
url: `${APP_DOMAIN_WITH_NGROK}/api/cron/campaigns/broadcast`,
flowControl: {
key: `broadcast-marketing-campaign-${campaignId}`,
parallelism: 1,
},
body: {
campaignId,
},
}),
);
}

return NextResponse.json(
CampaignSchema.parse(transformCampaign(updatedCampaign)),
Expand All @@ -217,15 +228,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) => {
Expand All @@ -244,8 +246,6 @@ export const DELETE = withWorkspace(
}
});

waitUntil(deleteCampaignSchedule(campaign));

return NextResponse.json({ id: campaignId });
},
{
Expand Down
118 changes: 82 additions & 36 deletions apps/web/app/(ee)/api/cron/campaigns/broadcast/route.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand All @@ -13,9 +13,13 @@ 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 {
Campaign,
CampaignStatus,
EmailDomain,
NotificationEmailType,
} from "@prisma/client";
import { differenceInMinutes } from "date-fns";
import { headers } from "next/headers";
import * as z from "zod/v4";
import { logAndRespond } from "../../utils";

Expand Down Expand Up @@ -102,46 +106,44 @@ 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
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 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,
status: CampaignStatus.scheduled,
OR: [{ scheduledAt: null }, { scheduledAt: { lte: new Date() } }],
},
data: {
status: CampaignStatus.sending,
qstashMessageId: messageId,
},
});
}

// 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) {
//
if (claimed.count === 0) {
if (!messageId || campaign.qstashMessageId !== messageId) {
return logAndRespond(
`Campaign ${campaignId} broadcast already initiated. Skipping...`,
);
}
}
}

const invalidFromResponse = await cancelCampaignIfInvalidFromAddress({
campaign,
emailDomains: program.emailDomains,
});

if (invalidFromResponse) {
return invalidFromResponse;
}

const campaignGroupIds = pluck(campaign.groups, "groupId");
const campaignPartnerTagIds = pluck(campaign.partnerTags, "partnerTagId");

Expand Down Expand Up @@ -376,3 +378,47 @@ export async function POST(req: Request) {
return handleAndReturnErrorResponse(error);
}
}

async function cancelCampaignIfInvalidFromAddress({
campaign,
emailDomains,
}: {
campaign: Pick<Campaign, "id" | "from" | "programId">;
emailDomains: Pick<EmailDomain, "slug" | "status">[];
}) {
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.`,
);
}
}
Loading
Loading