Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
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
29 changes: 13 additions & 16 deletions apps/web/app/(ee)/api/cron/groups/create-default-links/route.ts
Original file line number Diff line number Diff line change
@@ -1,20 +1,15 @@
import { handleAndReturnErrorResponse } from "@/lib/api/errors";
import { bulkCreateLinks } from "@/lib/api/links";
import { generatePartnerLink } from "@/lib/api/partners/generate-partner-link";
import { extractUtmParams } from "@/lib/api/utm/extract-utm-params";
import { applyGroupUtmToLink } from "@/lib/api/utm/apply-group-utm-to-link";
import { qstash } from "@/lib/cron";
import { verifyQstashSignature } from "@/lib/cron/verify-qstash";
import { loadAppsFlyerParameters } from "@/lib/integrations/appsflyer/apply-parameters";
import { AppsFlyerSettings } from "@/lib/integrations/appsflyer/schema";
import { isAppsFlyerTrackingUrl } from "@/lib/middleware/utils/is-appsflyer-tracking-url";
import { prisma } from "@/lib/prisma";
import { WorkspaceProps } from "@/lib/types";
import {
APP_DOMAIN_WITH_NGROK,
constructURLFromUTMParams,
isFulfilled,
log,
} from "@dub/utils";
import { APP_DOMAIN_WITH_NGROK, isFulfilled, log } from "@dub/utils";
import * as z from "zod/v4";
import { logAndRespond } from "../../utils";
export const dynamic = "force-dynamic";
Expand Down Expand Up @@ -133,8 +128,8 @@ export async function POST(req: Request) {
// Create a new defaultLink for each partner in the group
const processedLinks = (
await Promise.allSettled(
programEnrollments.map(({ partner, ...programEnrollment }) =>
generatePartnerLink({
programEnrollments.map(async ({ partner, ...programEnrollment }) => {
const processedLink = await generatePartnerLink({
workspace: {
id: workspace.id,
plan: workspace.plan as WorkspaceProps["plan"],
Expand All @@ -151,18 +146,20 @@ export async function POST(req: Request) {
},
link: {
domain: defaultLink.domain,
url: constructURLFromUTMParams(
defaultLink.url,
extractUtmParams(group.utmTemplate),
),
...extractUtmParams(group.utmTemplate, { excludeRef: true }),
url: defaultLink.url,
tenantId: programEnrollment.tenantId ?? undefined,
partnerGroupDefaultLinkId: defaultLink.id,
},
userId,
appsFlyerParameters,
}),
),
});

return applyGroupUtmToLink({
link: processedLink,
utmTemplate: group.utmTemplate,
partnerName: partner.name,
});
}),
)
)
.filter(isFulfilled)
Expand Down
25 changes: 13 additions & 12 deletions apps/web/app/(ee)/api/cron/groups/remap-default-links/route.ts
Original file line number Diff line number Diff line change
@@ -1,10 +1,12 @@
import { handleAndReturnErrorResponse } from "@/lib/api/errors";
import { bulkCreateLinks } from "@/lib/api/links";
import { generatePartnerLink } from "@/lib/api/partners/generate-partner-link";
import { applyGroupUtmToLink } from "@/lib/api/utm/apply-group-utm-to-link";
import { qstash } from "@/lib/cron";
import { verifyQstashSignature } from "@/lib/cron/verify-qstash";
import { loadAppsFlyerParameters } from "@/lib/integrations/appsflyer/apply-parameters";
import { AppsFlyerSettings } from "@/lib/integrations/appsflyer/schema";
import { syncGroupUtmJob } from "@/lib/jobs/handlers/sync-group-utm-job";
import { isAppsFlyerTrackingUrl } from "@/lib/middleware/utils/is-appsflyer-tracking-url";
import { prisma } from "@/lib/prisma";
import { WorkspaceProps } from "@/lib/types";
Expand Down Expand Up @@ -158,14 +160,14 @@ export async function POST(req: Request) {

const processedLinks = (
await Promise.allSettled(
linksToCreate.map((link) => {
linksToCreate.map(async (link) => {
const programEnrollment = programEnrollments.find(
(p) => p.partner.id === link.partnerId,
);

const partner = programEnrollment?.partner;

return generatePartnerLink({
const processedLink = await generatePartnerLink({
workspace: {
id: program.workspace.id,
plan: program.workspace.plan as WorkspaceProps["plan"],
Expand All @@ -189,6 +191,12 @@ export async function POST(req: Request) {
userId: userId ?? undefined,
appsFlyerParameters,
});

return applyGroupUtmToLink({
link: processedLink,
utmTemplate: partnerGroup.utmTemplate,
partnerName: partner?.name,
});
}),
)
)
Expand Down Expand Up @@ -251,18 +259,11 @@ export async function POST(req: Request) {
);
}

const syncUtmJob = await qstash.publishJSON({
url: `${APP_DOMAIN_WITH_NGROK}/api/cron/groups/sync-utm`,
body: {
groupId,
partnerIds,
},
await syncGroupUtmJob.dispatch({
groupId,
partnerIds,
});

console.log(
`Scheduled sync-utm job for group ${groupId}: ${prettyPrint(syncUtmJob)}`,
);

const remapDiscountCodesJob = await qstash.publishJSON({
url: `${APP_DOMAIN_WITH_NGROK}/api/cron/groups/remap-discount-codes`,
body: {
Expand Down
158 changes: 17 additions & 141 deletions apps/web/app/(ee)/api/cron/groups/sync-utm/route.ts
Original file line number Diff line number Diff line change
@@ -1,153 +1,29 @@
import { handleAndReturnErrorResponse } from "@/lib/api/errors";
import { linkCache } from "@/lib/api/links/cache";
import { extractUtmParams } from "@/lib/api/utm/extract-utm-params";
import { qstash } from "@/lib/cron";
import { verifyQstashSignature } from "@/lib/cron/verify-qstash";
import { prisma } from "@/lib/prisma";
import {
APP_DOMAIN_WITH_NGROK,
constructURLFromUTMParams,
log,
} from "@dub/utils";
import { withCron } from "@/lib/cron/with-cron";
import { syncGroupUtmJob } from "@/lib/jobs/handlers/sync-group-utm-job";
import * as z from "zod/v4";
import { logAndRespond } from "../../utils";
export const dynamic = "force-dynamic";

const PAGE_SIZE = 50;

const schema = z.object({
groupId: z.string(),
partnerIds: z.array(z.string()).optional(),
startAfterProgramEnrollmentId: z.string().optional(),
});

/**
Syncs the UTM parameter settings for a given group (whether there is a UTM template or not)

This job is triggered when:
1. a UTM template is created for a group
2. a UTM template is updated
3. in groups/remap-default-links cron
*/
// TODO:
// Remove this route after few hours of deployment

// POST /api/cron/groups/sync-utm
export async function POST(req: Request) {
try {
const rawBody = await req.text();
await verifyQstashSignature({ req, rawBody });

const { groupId, partnerIds, startAfterProgramEnrollmentId } = schema.parse(
JSON.parse(rawBody),
);

// Find the UTM template
const group = await prisma.partnerGroup.findUnique({
where: {
id: groupId,
},
include: {
utmTemplate: true,
},
});

if (!group) {
return logAndRespond(
`Group ${groupId} not found for groups/sync-utm cron. Skipping...`,
{
logLevel: "error",
},
);
}

const { utmTemplate } = group;

// Find partners in the group
const programEnrollments = await prisma.programEnrollment.findMany({
where: {
groupId: group.id,
...(partnerIds && {
partnerId: {
in: partnerIds,
},
}),
...(startAfterProgramEnrollmentId && {
id: {
gt: startAfterProgramEnrollmentId,
},
}),
},
take: PAGE_SIZE,
orderBy: {
id: "asc",
},
include: {
links: true,
},
});

if (programEnrollments.length === 0) {
return logAndRespond(`No program enrollments found. Skipping...`);
}

// extract links from program enrollments
const linksToUpdate = programEnrollments.flatMap((enrollment) =>
enrollment.links.map((link) => link),
);
// group links by the same url
const groupedLinksToUpdate = linksToUpdate.reduce(
(acc, link) => {
acc[link.url] = acc[link.url] || [];
acc[link.url].push(link.id);
return acc;
},
{} as Record<string, string[]>,
);

// Update the UTM for each partner links in the group
for (const [url, linkIds] of Object.entries(groupedLinksToUpdate)) {
const payload = {
url: constructURLFromUTMParams(url, extractUtmParams(utmTemplate)),
...extractUtmParams(utmTemplate, { excludeRef: true }),
};

const updatedLinks = await prisma.link.updateMany({
where: {
id: {
in: linkIds,
},
},
data: payload,
});
console.log(
`Updated ${updatedLinks.count} links with URL: ${payload.url}`,
);
}

const redisRes = await linkCache.expireMany(linksToUpdate);
console.log(`Updated Redis cache: ${JSON.stringify(redisRes, null, 2)}`);

if (programEnrollments.length === PAGE_SIZE) {
await qstash.publishJSON({
url: `${APP_DOMAIN_WITH_NGROK}/api/cron/groups/sync-utm`,
method: "POST",
body: {
groupId,
partnerIds,
startAfterProgramEnrollmentId:
programEnrollments[programEnrollments.length - 1].id,
},
});
}

return logAndRespond(
`Finished syncing UTM settings for ${programEnrollments.length} partners in the ${group.name} group (${group.id}).`,
);
} catch (error) {
await log({
message: `Error syncing UTM settings: ${error.message}.`,
type: "errors",
});

return handleAndReturnErrorResponse(error);
}
}
export const POST = withCron(async ({ rawBody }) => {
const { groupId, partnerIds, startAfterProgramEnrollmentId } = schema.parse(
JSON.parse(rawBody),
);

await syncGroupUtmJob.dispatch({
groupId,
partnerIds,
startAfterProgramEnrollmentId,
});

return logAndRespond("OK");
});
43 changes: 27 additions & 16 deletions apps/web/app/(ee)/api/cron/groups/update-default-links/route.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
import { handleAndReturnErrorResponse } from "@/lib/api/errors";
import { linkCache } from "@/lib/api/links/cache";
import { extractUtmParams } from "@/lib/api/utm/extract-utm-params";
import { applyGroupUtmToLink } from "@/lib/api/utm/apply-group-utm-to-link";
import { qstash } from "@/lib/cron";
import { verifyQstashSignature } from "@/lib/cron/verify-qstash";
import {
Expand All @@ -10,11 +10,8 @@ import {
import { AppsFlyerSettings } from "@/lib/integrations/appsflyer/schema";
import { isAppsFlyerTrackingUrl } from "@/lib/middleware/utils/is-appsflyer-tracking-url";
import { prisma } from "@/lib/prisma";
import {
APP_DOMAIN_WITH_NGROK,
constructURLFromUTMParams,
log,
} from "@dub/utils";
import { ProcessedLinkProps } from "@/lib/types";
import { APP_DOMAIN_WITH_NGROK, log } from "@dub/utils";
import { Link } from "@prisma/client";
import * as z from "zod/v4";
import { logAndRespond } from "../../utils";
Expand Down Expand Up @@ -148,10 +145,24 @@ export async function POST(req: Request) {
}[] = [];

for (const defaultPartnerLink of defaultPartnerLinks) {
let url = constructURLFromUTMParams(
defaultLink.url,
extractUtmParams(group.utmTemplate),
);
const utmContext = {
partnerName:
defaultPartnerLink.partner?.name || defaultPartnerLink.key,
partnerLinkKey: defaultPartnerLink.key,
};

const linkWithUtm = applyGroupUtmToLink({
link: {
domain: defaultPartnerLink.domain,
key: defaultPartnerLink.key,
url: defaultLink.url,
projectId: defaultLink.program.workspaceId,
} as ProcessedLinkProps,
utmTemplate: group.utmTemplate,
partnerName: defaultPartnerLink.partner?.name,
});

let url = linkWithUtm.url;

// Inject AppsFlyer parameters with resolved macros
if (
Expand All @@ -161,19 +172,19 @@ export async function POST(req: Request) {
url = applyAppsFlyerParameters({
url,
parameters: appsFlyerParameters,
context: {
partnerName:
defaultPartnerLink.partner?.name || defaultPartnerLink.key,
partnerLinkKey: defaultPartnerLink.key,
},
context: utmContext,
});
}

linksToUpdate.push({
id: defaultPartnerLink.id,
link: {
url,
...extractUtmParams(group.utmTemplate, { excludeRef: true }),
utm_source: linkWithUtm.utm_source ?? null,
utm_medium: linkWithUtm.utm_medium ?? null,
utm_campaign: linkWithUtm.utm_campaign ?? null,
utm_term: linkWithUtm.utm_term ?? null,
utm_content: linkWithUtm.utm_content ?? null,
},
});
}
Expand Down
Loading
Loading