diff --git a/apps/web/app/(ee)/api/cron/partners/merge-accounts/route.ts b/apps/web/app/(ee)/api/cron/partners/merge-accounts/route.ts deleted file mode 100644 index f0d85d09160..00000000000 --- a/apps/web/app/(ee)/api/cron/partners/merge-accounts/route.ts +++ /dev/null @@ -1,488 +0,0 @@ -import { handleAndReturnErrorResponse } from "@/lib/api/errors"; -import { resolveFraudGroups } from "@/lib/api/fraud/resolve-fraud-groups"; -import { linkCache } from "@/lib/api/links/cache"; -import { includeProgramEnrollment } from "@/lib/api/links/include-program-enrollment"; -import { includeTags } from "@/lib/api/links/include-tags"; -import { syncTotalCommissions } from "@/lib/api/partners/sync-total-commissions"; -import { PRISMA_UPDATEMANY_LIMIT } from "@/lib/cron"; -import { verifyQstashSignature } from "@/lib/cron/verify-qstash"; -import { conn } from "@/lib/planetscale"; -import { prisma } from "@/lib/prisma"; -import { storage } from "@/lib/storage"; -import { recordLink } from "@/lib/tinybird"; -import { redis } from "@/lib/upstash"; -import { sendBatchEmail } from "@dub/email"; -import PartnerAccountMerged from "@dub/email/templates/partner-account-merged"; -import { log, prettyPrint, R2_URL } from "@dub/utils"; -import { FraudRuleType } from "@prisma/client"; -import * as z from "zod/v4"; - -export const dynamic = "force-dynamic"; - -const schema = z.object({ - userId: z.string(), - sourceEmail: z.string(), - targetEmail: z.string(), -}); - -const CACHE_KEY_PREFIX = "merge-partner-accounts"; - -// POST /api/cron/partners/merge-accounts -// This route is used to merge a partner account into another account -export async function POST(req: Request) { - let userId: string | null = null; - - try { - const rawBody = await req.text(); - - await verifyQstashSignature({ - req, - rawBody, - }); - - const { - userId: parsedUserId, - sourceEmail, - targetEmail, - } = schema.parse(JSON.parse(rawBody)); - - userId = parsedUserId; - - console.log({ - userId, - sourceEmail, - targetEmail, - }); - - const partnerAccounts = await prisma.partner.findMany({ - where: { - email: { - in: [sourceEmail, targetEmail], - }, - }, - select: { - id: true, - email: true, - image: true, - programs: { - select: { - programId: true, - tenantId: true, - status: true, - }, - }, - users: { - select: { - userId: true, - }, - }, - }, - }); - - if (partnerAccounts.length === 0) { - return new Response("Partner accounts not found."); - } - - const sourceAccount = partnerAccounts.find( - ({ email }) => email?.toLowerCase() === sourceEmail.toLowerCase(), - ); - - const targetAccount = partnerAccounts.find( - ({ email }) => email?.toLowerCase() === targetEmail.toLowerCase(), - ); - - if (!sourceAccount) { - return new Response( - `Partner account with email ${sourceEmail} not found.`, - ); - } - - if (!targetAccount) { - return new Response( - `Partner account with email ${targetEmail} not found.`, - ); - } - - if (sourceAccount.id === targetAccount.id) { - return new Response( - `Source and target partner accounts must be different. Source account: ${sourceAccount.email} (${sourceAccount.id}), Target account: ${targetAccount.email} (${targetAccount.id})`, - ); - } - - const { - id: sourcePartnerId, - users: sourcePartnerUsers, - programs: sourcePartnerEnrollments, - } = sourceAccount; - - const { id: targetPartnerId, programs: targetPartnerEnrollments } = - targetAccount; - - const programIdsToTransfer = sourcePartnerEnrollments.map( - ({ programId }) => programId, - ); - - const updateManyPayload = { - where: { - programId: { - in: programIdsToTransfer, - }, - partnerId: sourcePartnerId, - }, - data: { - partnerId: targetPartnerId, - }, - }; - - if (programIdsToTransfer.length > 0) { - // update links and payouts - const [updatedLinksRes, updatedPayoutsRes] = await Promise.all([ - prisma.link.updateMany(updateManyPayload), - prisma.payout.updateMany(updateManyPayload), - ]); - console.log( - `Updated ${updatedLinksRes.count} links, and ${updatedPayoutsRes.count} payouts`, - ); - - // for commissions / customers, we need to update them in batches of PRISMA_UPDATEMANY_LIMIT cause there can be a lot of them - while (true) { - const { count } = await prisma.commission.updateMany({ - ...updateManyPayload, - limit: PRISMA_UPDATEMANY_LIMIT, - }); - console.log(`Updated ${count} commissions`); - if (count < PRISMA_UPDATEMANY_LIMIT) break; - } - - while (true) { - const { count } = await prisma.customer.updateMany({ - ...updateManyPayload, - limit: PRISMA_UPDATEMANY_LIMIT, - }); - console.log(`Updated ${count} customers`); - if (count < PRISMA_UPDATEMANY_LIMIT) break; - } - - // update discount codes, notification emails, messages, and partner comments - const [ - updatedDiscountCodesRes, - updatedNotificationEmailsRes, - updatedMessagesRes, - updatedPartnerCommentsRes, - ] = await Promise.all([ - prisma.discountCode.updateMany(updateManyPayload), - prisma.notificationEmail.updateMany(updateManyPayload), - prisma.message.updateMany(updateManyPayload), - prisma.partnerComment.updateMany(updateManyPayload), - ]); - console.log( - `Updated ${updatedDiscountCodesRes.count} discount codes, ${updatedNotificationEmailsRes.count} notification emails, ${updatedMessagesRes.count} messages, and ${updatedPartnerCommentsRes.count} partner comments`, - ); - - const updatedLinks = await prisma.link.findMany({ - where: { - programId: { - in: programIdsToTransfer, - }, - partnerId: targetPartnerId, - }, - include: { - ...includeTags, - ...includeProgramEnrollment, - }, - }); - - // only transfer bounty submissions if the target partner has no submissions for the same bounty - const bountySubmissionStats = await prisma.bountySubmission.groupBy({ - by: ["bountyId"], - where: { - partnerId: { - in: [sourcePartnerId, targetPartnerId], - }, - }, - _count: { - partnerId: true, - }, - }); - const bountiesToTransfer = bountySubmissionStats - .filter(({ _count }) => _count.partnerId === 1) - .map(({ bountyId }) => bountyId); - - if (bountiesToTransfer.length > 0) { - const updatedBountySubmissions = - await prisma.bountySubmission.updateMany({ - where: { - bountyId: { in: bountiesToTransfer }, - partnerId: sourcePartnerId, - }, - data: { - partnerId: targetPartnerId, - }, - }); - console.log( - `Transferred ${updatedBountySubmissions.count} bounty submissions`, - ); - } - - const res = await Promise.allSettled([ - // update link metadata in Tinybird - recordLink(updatedLinks), - // expire link cache in Redis - linkCache.expireMany(updatedLinks), - // Sync total commissions for the target partner in each program - ...programIdsToTransfer.map((programId) => - syncTotalCommissions({ - partnerId: targetPartnerId, - programId, - }), - ), - ]); - console.log(prettyPrint(res)); - } - - // Update program enrollments – we do this last to avoid relational updates on Commission, Customer, etc. from running into transaction limits - // First, we start with new enrollments (enrollments that are not duplicate in both source and target) - const newEnrollments = sourcePartnerEnrollments.filter( - ({ programId }) => - !targetPartnerEnrollments.some( - ({ programId: targetProgramId }) => programId === targetProgramId, - ), - ); - if (newEnrollments.length > 0) { - await prisma.programEnrollment.updateMany({ - where: { - programId: { - in: newEnrollments.map(({ programId }) => programId), - }, - partnerId: sourcePartnerId, - }, - data: { - partnerId: targetPartnerId, - }, - }); - } - - // Then, we transfer existing enrollments (enrollments that are duplicate in both source and target) – needs to be handled delicately via a for loop - const existingEnrollments = sourcePartnerEnrollments.filter( - ({ programId }) => - targetPartnerEnrollments.some( - ({ programId: targetProgramId }) => programId === targetProgramId, - ), - ); - - if (existingEnrollments.length > 0) { - for (const sourceEnrollment of existingEnrollments) { - const targetEnrollment = targetPartnerEnrollments.find( - ({ programId }) => programId === sourceEnrollment.programId, - ); - - await prisma.$transaction(async (tx) => { - if ( - targetEnrollment && - sourceEnrollment.status === "approved" && - ["pending", "invited"].includes(targetEnrollment.status) - ) { - await tx.programEnrollment.update({ - where: { - partnerId_programId: { - partnerId: targetPartnerId, - programId: sourceEnrollment.programId, - }, - }, - data: { status: "approved" }, - }); - } - - await tx.programEnrollment.delete({ - where: { - partnerId_programId: { - partnerId: sourcePartnerId, - programId: sourceEnrollment.programId, - }, - }, - }); - - // update target enrollment with source enrollment's tenantId if target enrollment does not have a tenantId - if (sourceEnrollment.tenantId && !targetEnrollment?.tenantId) { - await tx.programEnrollment.update({ - where: { - partnerId_programId: { - partnerId: targetPartnerId, - programId: sourceEnrollment.programId, - }, - }, - data: { - tenantId: sourceEnrollment.tenantId, - }, - }); - } - }); - console.log( - `Deleted old source enrollment for program ${sourceEnrollment.programId}.${sourceEnrollment.tenantId ? ` Since there was a tenantId, we updated the target enrollment with the same tenantId: ${sourceEnrollment.tenantId}` : ""}`, - ); - } - } - - // Remove the user if there are no workspaces left - // TODO: we need to handle deleting multiple users when we allow partners to invite their team members in the future - const sourcePartnerUser = sourcePartnerUsers[0]; - - if (sourcePartnerUser) { - const workspaceCount = await prisma.projectUsers.count({ - where: { - userId: sourcePartnerUser.userId, - }, - }); - - if (workspaceCount === 0) { - try { - const deletedUser = await prisma.user.delete({ - where: { - id: sourcePartnerUser.userId, - }, - select: { - id: true, - email: true, - image: true, - }, - }); - console.log(`Deleted user ${deletedUser.email} (${deletedUser.id})`); - - if (deletedUser.image) { - await storage.delete({ - key: deletedUser.image.replace(`${R2_URL}/`, ""), - }); - } - } catch (error) { - console.error( - `Error deleting user ${sourcePartnerUser.userId}: ${error.message}`, - ); - } - } - } - - const fraudEventsToDelete = await prisma.fraudEvent.findMany({ - where: { - partnerId: sourcePartnerId, - fraudEventGroup: { - type: FraudRuleType.partnerDuplicateAccount, - }, - }, - include: { - fraudEventGroup: { - select: { - id: true, - _count: { - select: { - fraudEvents: true, - }, - }, - }, - }, - }, - }); - - if (fraudEventsToDelete.length > 0) { - await prisma.fraudEvent.deleteMany({ - where: { - id: { in: fraudEventsToDelete.map((e) => e.id) }, - }, - }); - } - - const fraudEventGroupsToResolve = fraudEventsToDelete.filter( - // this is the count pre-deletion the fraud event, so if there are 2 fraud events - // that means post-deletion will leave 1 fraud event in the group (no additional duplicates), hence can be resolved - (e) => e.fraudEventGroup._count.fraudEvents === 2, - ); - - await resolveFraudGroups({ - where: { - OR: [ - { - partnerId: sourcePartnerId, - }, - ...(fraudEventGroupsToResolve.length > 0 - ? [ - { - id: { - in: fraudEventGroupsToResolve.map( - (e) => e.fraudEventGroup.id, - ), - }, - }, - ] - : []), - ], - type: FraudRuleType.partnerDuplicateAccount, - }, - resolutionReason: - "Automatically resolved because partners with duplicate payout methods were merged. No other partners share this payout method.", - }); - - try { - // Finally, delete the partner account - await conn.execute(`DELETE FROM Partner WHERE id = ?`, [sourcePartnerId]); - console.log( - `Deleted partner ${sourceAccount.email} (${sourceAccount.id})`, - ); - - if (sourceAccount.image) { - await storage.delete({ - key: sourceAccount.image.replace(`${R2_URL}/`, ""), - }); - } - } catch (error) { - console.error( - `Error deleting partner ${sourcePartnerId}: ${error.message}`, - ); - } - - // Make sure the cache is cleared - await redis.del(`${CACHE_KEY_PREFIX}:${userId}`); - - const resendBatchEmailRes = await sendBatchEmail( - [ - { - variant: "notifications", - to: sourceEmail, - subject: "Your Dub partner accounts are now merged", - react: PartnerAccountMerged({ - email: sourceEmail, - sourceEmail, - targetEmail, - }), - }, - { - variant: "notifications", - to: targetEmail, - subject: "Your Dub partner accounts are now merged", - react: PartnerAccountMerged({ - email: targetEmail, - sourceEmail, - targetEmail, - }), - }, - ], - { - idempotencyKey: `${CACHE_KEY_PREFIX}/${userId}`, - }, - ); - console.log(prettyPrint(resendBatchEmailRes)); - - return new Response( - `Partner account ${sourceEmail} merged into ${targetEmail}.`, - ); - } catch (error) { - if (userId) { - await redis.del(`${CACHE_KEY_PREFIX}:${userId}`); - } - - await log({ - message: `Error merging partner accounts: ${error.message}`, - type: "alerts", - }); - - return handleAndReturnErrorResponse(error); - } -} diff --git a/apps/web/app/(ee)/api/e2e/trigger-merge-accounts/route.ts b/apps/web/app/(ee)/api/e2e/trigger-merge-accounts/route.ts new file mode 100644 index 00000000000..8b6a7a2465c --- /dev/null +++ b/apps/web/app/(ee)/api/e2e/trigger-merge-accounts/route.ts @@ -0,0 +1,74 @@ +import { DubApiError } from "@/lib/api/errors"; +import { parseRequestBody } from "@/lib/api/utils"; +import { withWorkspace } from "@/lib/auth"; +import { triggerQStashWorkflow } from "@/lib/cron/qstash-workflow"; +import { prisma } from "@/lib/prisma"; +import { ACME_PROGRAM_ID, nanoid } from "@dub/utils"; +import { NextResponse } from "next/server"; +import * as z from "zod/v4"; +import { assertE2EWorkspace } from "../guard"; + +const bodySchema = z.object({ + sourceEmail: z.email(), + targetEmail: z.email(), +}); + +// POST /api/e2e/trigger-merge-accounts +export const POST = withWorkspace( + async ({ req, workspace }) => { + assertE2EWorkspace(workspace); + + const { sourceEmail, targetEmail } = bodySchema.parse( + await parseRequestBody(req), + ); + + if (sourceEmail.toLowerCase() === targetEmail.toLowerCase()) { + throw new DubApiError({ + code: "bad_request", + message: "Source and target emails must be different.", + }); + } + + const partners = await prisma.partner.findMany({ + where: { + email: { in: [sourceEmail, targetEmail] }, + programs: { some: { programId: ACME_PROGRAM_ID } }, + }, + select: { email: true }, + }); + + const enrolledEmails = new Set(partners.map((p) => p.email?.toLowerCase())); + + if ( + !enrolledEmails.has(sourceEmail.toLowerCase()) || + !enrolledEmails.has(targetEmail.toLowerCase()) + ) { + throw new DubApiError({ + code: "bad_request", + message: + "Both partners must exist and be enrolled in the Acme test program.", + }); + } + + const userId = `e2e-merge-${nanoid()}`; + + const res = await triggerQStashWorkflow({ + workflowType: "merge-partner-accounts", + workflowLabel: userId, + body: { + userId, + sourceEmail, + targetEmail, + }, + flowControl: { + key: userId, + parallelism: 1, + }, + }); + + return NextResponse.json(res); + }, + { + requiredPermissions: ["workspaces.write"], + }, +); diff --git a/apps/web/app/(ee)/api/workflows/merge-partner-accounts/route.ts b/apps/web/app/(ee)/api/workflows/merge-partner-accounts/route.ts new file mode 100644 index 00000000000..1b3ad6bb395 --- /dev/null +++ b/apps/web/app/(ee)/api/workflows/merge-partner-accounts/route.ts @@ -0,0 +1,890 @@ +import { resolveFraudGroups } from "@/lib/api/fraud/resolve-fraud-groups"; +import { linkCache } from "@/lib/api/links/cache"; +import { includeProgramEnrollment } from "@/lib/api/links/include-program-enrollment"; +import { includeTags } from "@/lib/api/links/include-tags"; +import { syncTotalCommissions } from "@/lib/api/partners/sync-total-commissions"; +import { logger } from "@/lib/axiom/server"; +import { PRISMA_UPDATEMANY_LIMIT } from "@/lib/cron"; +import { getWorkflowConfig } from "@/lib/cron/qstash-workflow"; +import { conn } from "@/lib/planetscale"; +import { prisma } from "@/lib/prisma"; +import { storage } from "@/lib/storage"; +import { recordLink } from "@/lib/tinybird"; +import { redis } from "@/lib/upstash"; +import { sendBatchEmail } from "@dub/email"; +import PartnerAccountMerged from "@dub/email/templates/partner-account-merged"; +import { log, prettyPrint, R2_URL } from "@dub/utils"; +import { FraudRuleType } from "@prisma/client"; +import { serve } from "@upstash/workflow/nextjs"; +import * as z from "zod/v4"; +import { logAndReturn } from "../../cron/utils"; + +const inputSchema = z.object({ + userId: z.string(), + sourceEmail: z.string(), + targetEmail: z.string(), +}); + +type Input = z.infer; + +const CACHE_KEY_PREFIX = "merge-partner-accounts"; + +/** + * Steps: + * 1. load-merge-plan: resolve + validate accounts, build the ordered list of + * source enrollments to process (overlaps first, then transfers). + * 2. merge-enrollment- (one per enrollment): transfer the enrollment's + * program data to the target and either merge into the existing target + * enrollment (overlap) or move the enrollment over (transfer). + * 3. finalize-transfers: transfer bounty submissions + sync transferred links + * (Tinybird/cache) and total commissions. + * 4. cleanup-source-account: delete the source partner's rewinds, source user, + * duplicate-account fraud events, and finally the source partner itself. + * 5. send-merged-emails: clear the verification cache + notify both accounts. + */ + +// POST /api/workflows/merge-partner-accounts +export const { POST } = serve( + async (context) => { + const { userId, sourceEmail, targetEmail } = context.requestPayload; + + // Step 1: Resolve + validate accounts and build the merge plan + const plan = await context.run("load-merge-plan", async () => { + return await loadMergePlan({ sourceEmail, targetEmail }); + }); + + if (!plan.proceed) { + console.log(`Skipping merge: ${plan.reason}`); + + // Clear the verification cache so the user can cleanly retry (sendTokens + // rejects a new request while this key is still set). + await context.run("clear-cache-after-skip", async () => { + await redis.del(`${CACHE_KEY_PREFIX}:${userId}`); + return logAndReturn({ + outputLog: `Cleared merge cache after skipped merge: ${plan.reason}`, + }); + }); + + return; + } + + const { + sourcePartnerId, + targetPartnerId, + sourceImage, + sourceUserId, + hasRewinds, + orderedSourceEnrollmentIds, + programIdsToTransfer, + } = plan; + + // Step 2: Merge each source enrollment in its own durable step. + // Overlaps are ordered first, then transfers. Each step re-fetches the + // live enrollment so it is safe to retry after a partial run. + for (const enrollmentId of orderedSourceEnrollmentIds) { + await context.run(`merge-enrollment-${enrollmentId}`, async () => { + return await mergeSingleEnrollment({ + enrollmentId, + sourcePartnerId, + targetPartnerId, + }); + }); + } + + // Step 3: Finalize the transfers. Both the bounty transfer and the + // link/commission sync are idempotent post-transfer reconciliation + await context.run("finalize-transfers", async () => { + if (programIdsToTransfer.length === 0) { + return logAndReturn({ outputLog: "No programs to finalize." }); + } + + // Transfer bounty submissions + const { outputLog: bountyLog } = await transferBountySubmissions({ + sourcePartnerId, + targetPartnerId, + }); + + // Sync transferred links (Tinybird + cache) and total commissions + const { outputLog: syncLog } = await syncLinksAndCommissions({ + targetPartnerId, + programIdsToTransfer, + }); + + return logAndReturn({ outputLog: `${bountyLog} | ${syncLog}` }); + }); + + // Step 4: Tear down the source account + await context.run("cleanup-source-account", async () => { + const logs: string[] = []; + + // Delete the source partner's rewinds + if (hasRewinds) { + const deletedRewinds = await prisma.partnerRewind.deleteMany({ + where: { partnerId: sourcePartnerId }, + }); + logs.push(`Deleted ${deletedRewinds.count} partner rewinds`); + } + + // Remove the source user if there are no workspaces left + if (sourceUserId) { + const { outputLog } = await deleteSourceUser({ sourceUserId }); + logs.push(outputLog); + } + + // Clean up duplicate-account fraud events + resolve their groups + const { outputLog: fraudLog } = await cleanupFraudEvents({ + sourcePartnerId, + }); + logs.push(fraudLog); + + // Delete the source partner account (must be last) + const { outputLog: partnerLog } = await deleteSourcePartner({ + sourcePartnerId, + sourceEmail, + sourceImage, + }); + logs.push(partnerLog); + + return logAndReturn({ outputLog: logs.join(" | ") }); + }); + + // Step 5: Clear the verification cache and notify both accounts + await context.run("send-merged-emails", async () => { + await redis.del(`${CACHE_KEY_PREFIX}:${userId}`); + + const resendBatchEmailRes = await sendBatchEmail( + [ + { + variant: "notifications", + to: sourceEmail, + subject: "Your Dub partner accounts are now merged", + react: PartnerAccountMerged({ + email: sourceEmail, + sourceEmail, + targetEmail, + }), + }, + { + variant: "notifications", + to: targetEmail, + subject: "Your Dub partner accounts are now merged", + react: PartnerAccountMerged({ + email: targetEmail, + sourceEmail, + targetEmail, + }), + }, + ], + { + idempotencyKey: `${CACHE_KEY_PREFIX}/${userId}`, + }, + ); + + return logAndReturn({ + outputLog: `Partner account ${sourceEmail} merged into ${targetEmail}. ${prettyPrint(resendBatchEmailRes)}`, + }); + }); + }, + { + initialPayloadParser: (requestPayload) => { + return inputSchema.parse(JSON.parse(requestPayload)); + }, + failureFunction: async ({ + context, + failStatus, + failResponse, + failHeaders, + }) => { + const { userId } = inputSchema.parse(context.requestPayload); + + // Clear the verification cache so the partner can retry from the start + await redis.del(`${CACHE_KEY_PREFIX}:${userId}`); + + const { correlation } = getWorkflowConfig({ + workflowType: "merge-partner-accounts", + body: context.requestPayload, + }); + + await log({ + message: `Error merging partner accounts: ${JSON.stringify(correlation)}, workflowRunId=${context.workflowRunId}, failStatus=${failStatus}, failResponse=${failResponse}. Some enrollments may already be merged (see workflow run for completed steps) - manual cleanup may be required.`, + type: "alerts", + mention: true, + }); + + logger.error("workflow.failed", { + service: "qstash", + event: "workflow.failed", + workflowType: "merge-partner-accounts", + workflowRunId: context.workflowRunId, + failStatus, + failResponse, + failHeaders, + correlation, + }); + + await logger.flush(); + }, + }, +); + +type MergePlan = + | { proceed: false; reason: string } + | { + proceed: true; + sourcePartnerId: string; + targetPartnerId: string; + sourceImage: string | null; + sourceUserId: string | null; + hasRewinds: boolean; + orderedSourceEnrollmentIds: string[]; + programIdsToTransfer: string[]; + }; + +async function loadMergePlan({ + sourceEmail, + targetEmail, +}: { + sourceEmail: string; + targetEmail: string; +}): Promise { + const partnerAccounts = await prisma.partner.findMany({ + where: { + email: { + in: [sourceEmail, targetEmail], + }, + }, + select: { + id: true, + email: true, + image: true, + users: { + select: { + userId: true, + }, + }, + partnerRewinds: true, + }, + }); + + if (partnerAccounts.length === 0) { + return { proceed: false, reason: "Partner accounts not found." }; + } + + const sourceAccount = partnerAccounts.find( + ({ email }) => email?.toLowerCase() === sourceEmail.toLowerCase(), + ); + + const targetAccount = partnerAccounts.find( + ({ email }) => email?.toLowerCase() === targetEmail.toLowerCase(), + ); + + if (!sourceAccount) { + return { + proceed: false, + reason: `Partner account with email ${sourceEmail} not found.`, + }; + } + + if (!targetAccount) { + return { + proceed: false, + reason: `Partner account with email ${targetEmail} not found.`, + }; + } + + if (sourceAccount.id === targetAccount.id) { + return { + proceed: false, + reason: `Source and target partner accounts must be different. Source account: ${sourceAccount.email} (${sourceAccount.id}), Target account: ${targetAccount.email} (${targetAccount.id})`, + }; + } + + const sourcePartnerId = sourceAccount.id; + const targetPartnerId = targetAccount.id; + + const [sourceEnrollments, targetEnrollments] = await Promise.all([ + prisma.programEnrollment.findMany({ + where: { partnerId: sourcePartnerId }, + select: { id: true, programId: true }, + }), + prisma.programEnrollment.findMany({ + where: { partnerId: targetPartnerId }, + select: { programId: true }, + }), + ]); + + const targetProgramIds = new Set( + targetEnrollments.map((enrollment) => enrollment.programId), + ); + + const overlappingEnrollments = sourceEnrollments.filter((enrollment) => + targetProgramIds.has(enrollment.programId), + ); + + const transferEnrollments = sourceEnrollments.filter( + (enrollment) => !targetProgramIds.has(enrollment.programId), + ); + + // Overlaps first, then transfers (preserves the original processing order). + const orderedSourceEnrollmentIds = [ + ...overlappingEnrollments, + ...transferEnrollments, + ].map(({ id }) => id); + + return { + proceed: true, + sourcePartnerId, + targetPartnerId, + sourceImage: sourceAccount.image, + sourceUserId: sourceAccount.users[0]?.userId ?? null, + hasRewinds: sourceAccount.partnerRewinds.length > 0, + orderedSourceEnrollmentIds, + programIdsToTransfer: sourceEnrollments.map(({ programId }) => programId), + }; +} + +async function transferRowsInBatches( + updateBatch: () => Promise, + { + resourceName, + }: { + resourceName: string; + }, +) { + while (true) { + const count = await updateBatch(); + console.log(`Transferred ${count} ${resourceName} in batch`); + if (count < PRISMA_UPDATEMANY_LIMIT) { + break; + } + } +} + +async function transferPartnerProgramData({ + sourcePartnerId, + targetPartnerId, + programId, +}: { + sourcePartnerId: string; + targetPartnerId: string; + programId: string; +}) { + const where = { + programId, + partnerId: sourcePartnerId, + }; + const payload = { + where, + data: { + partnerId: targetPartnerId, + }, + }; + + await Promise.all([ + // High-volume tables: move in batches of PRISMA_UPDATEMANY_LIMIT + transferRowsInBatches( + async () => + ( + await prisma.commission.updateMany({ + ...payload, + limit: PRISMA_UPDATEMANY_LIMIT, + }) + ).count, + { resourceName: "commission" }, + ), + transferRowsInBatches( + async () => + ( + await prisma.link.updateMany({ + ...payload, + limit: PRISMA_UPDATEMANY_LIMIT, + }) + ).count, + { resourceName: "link" }, + ), + transferRowsInBatches( + async () => + ( + await prisma.customer.updateMany({ + ...payload, + limit: PRISMA_UPDATEMANY_LIMIT, + }) + ).count, + { resourceName: "customer" }, + ), + // Low-volume tables: single updateMany is fine + prisma.payout.updateMany(payload), + prisma.discountCode.updateMany(payload), + prisma.notificationEmail.updateMany(payload), + prisma.message.updateMany(payload), + prisma.partnerComment.updateMany(payload), + ]); + + // After payouts are moved onto the target partner, fold any duplicate + // pending payouts for this program into a single payout. + await combinePendingPayouts({ + partnerId: targetPartnerId, + programId, + }); +} + +async function mergeSingleEnrollment({ + enrollmentId, + sourcePartnerId, + targetPartnerId, +}: { + enrollmentId: string; + sourcePartnerId: string; + targetPartnerId: string; +}) { + const sourceEnrollment = await prisma.programEnrollment.findUnique({ + where: { id: enrollmentId }, + }); + + if (!sourceEnrollment) { + return logAndReturn({ + programId: null, + action: "skip", + outputLog: `Enrollment ${enrollmentId} no longer exists, skipping`, + }); + } + + if (sourceEnrollment.partnerId === targetPartnerId) { + return logAndReturn({ + programId: sourceEnrollment.programId, + action: "skip", + outputLog: `Enrollment ${enrollmentId} already on target partner, skipping`, + }); + } + + // Another process could have reassigned this enrollment away from the source + // partner; only the source partner's own enrollments should be merged. + if (sourceEnrollment.partnerId !== sourcePartnerId) { + return logAndReturn({ + programId: sourceEnrollment.programId, + action: "skip", + outputLog: `Enrollment ${enrollmentId} no longer belongs to ${sourcePartnerId} (now ${sourceEnrollment.partnerId}), skipping`, + }); + } + + const { programId } = sourceEnrollment; + + const targetEnrollment = await prisma.programEnrollment.findUnique({ + where: { + partnerId_programId: { + partnerId: targetPartnerId, + programId, + }, + }, + }); + + await transferPartnerProgramData({ + sourcePartnerId, + targetPartnerId, + programId, + }); + + if (targetEnrollment) { + await prisma.$transaction(async (tx) => { + if ( + sourceEnrollment.status === "approved" && + ["pending", "invited"].includes(targetEnrollment.status) + ) { + await tx.programEnrollment.update({ + where: { + partnerId_programId: { + partnerId: targetPartnerId, + programId, + }, + }, + data: { status: "approved" }, + }); + } + + if (sourceEnrollment.applicationId) { + await tx.programEnrollment.updateMany({ + where: { id: sourceEnrollment.id, partnerId: sourcePartnerId }, + data: { applicationId: null }, + }); + } + + await tx.programEnrollment.deleteMany({ + where: { id: sourceEnrollment.id, partnerId: sourcePartnerId }, + }); + + const tenantIdToCopy = + targetEnrollment.tenantId ?? sourceEnrollment.tenantId; + + if (tenantIdToCopy && tenantIdToCopy !== targetEnrollment.tenantId) { + const existingTenantEnrollment = await tx.programEnrollment.findUnique({ + where: { + tenantId_programId: { + tenantId: tenantIdToCopy, + programId, + }, + }, + }); + + if (!existingTenantEnrollment) { + await tx.programEnrollment.update({ + where: { + partnerId_programId: { + partnerId: targetPartnerId, + programId, + }, + }, + data: { tenantId: tenantIdToCopy }, + }); + } + } + }); + + return logAndReturn({ + programId, + action: "overlap", + outputLog: `Merged overlapping enrollment for program ${programId}`, + }); + } + + // Scope the transfer to the source partner so a concurrent reassignment + // can't make us steal another partner's enrollment. + const { count } = await prisma.programEnrollment.updateMany({ + where: { id: sourceEnrollment.id, partnerId: sourcePartnerId }, + data: { partnerId: targetPartnerId }, + }); + + if (count === 0) { + return logAndReturn({ + programId, + action: "skip", + outputLog: `Enrollment ${sourceEnrollment.id} no longer owned by ${sourcePartnerId}, skipping transfer`, + }); + } + + return logAndReturn({ + programId, + action: "transfer", + outputLog: `Transferred enrollment for program ${programId}`, + }); +} + +async function transferBountySubmissions({ + sourcePartnerId, + targetPartnerId, +}: { + sourcePartnerId: string; + targetPartnerId: string; +}) { + const bountySubmissionStats = await prisma.bountySubmission.groupBy({ + by: ["bountyId"], + where: { + partnerId: { + in: [sourcePartnerId, targetPartnerId], + }, + }, + _count: { + partnerId: true, + }, + }); + + // only transfer bounty submissions if the target partner has no submissions for the same bounty + const bountiesToTransfer = bountySubmissionStats + .filter(({ _count }) => _count.partnerId === 1) + .map(({ bountyId }) => bountyId); + + if (bountiesToTransfer.length === 0) { + return logAndReturn({ outputLog: "No bounty submissions to transfer." }); + } + + const updatedBountySubmissions = await prisma.bountySubmission.updateMany({ + where: { + bountyId: { in: bountiesToTransfer }, + partnerId: sourcePartnerId, + }, + data: { + partnerId: targetPartnerId, + }, + }); + + return logAndReturn({ + outputLog: `Transferred ${updatedBountySubmissions.count} bounty submissions`, + }); +} + +async function combinePendingPayouts({ + partnerId, + programId, +}: { + partnerId: string; + programId: string; +}) { + const payoutsToCombine = await prisma.payout.findMany({ + where: { + programId, + partnerId, + status: "pending", + }, + orderBy: { + createdAt: "asc", + }, + }); + + if (payoutsToCombine.length < 2) { + return; + } + + const periodStarts = payoutsToCombine + .map((payout) => payout.periodStart) + .filter((date): date is Date => date !== null); + const periodEnds = payoutsToCombine + .map((payout) => payout.periodEnd) + .filter((date): date is Date => date !== null); + + const periodStart = + periodStarts.length > 0 + ? new Date(Math.min(...periodStarts.map((date) => date.getTime()))) + : payoutsToCombine[0].periodStart; + const periodEnd = + periodEnds.length > 0 + ? new Date(Math.max(...periodEnds.map((date) => date.getTime()))) + : payoutsToCombine[0].periodEnd; + + const totalAmount = payoutsToCombine.reduce( + (sum, payout) => sum + payout.amount, + 0, + ); + + const combinedPayoutId = payoutsToCombine[0].id; + const payoutIdsToDelete = payoutsToCombine.slice(1).map((p) => p.id); + + await prisma.payout.update({ + where: { + id: combinedPayoutId, + }, + data: { + amount: totalAmount, + periodStart, + periodEnd, + }, + }); + + await transferRowsInBatches( + async () => + ( + await prisma.commission.updateMany({ + where: { + payoutId: { + in: payoutIdsToDelete, + }, + }, + data: { + payoutId: combinedPayoutId, + }, + limit: PRISMA_UPDATEMANY_LIMIT, + }) + ).count, + { resourceName: "commission" }, + ); + + const deletedPayouts = await prisma.payout.deleteMany({ + where: { + id: { + in: payoutIdsToDelete, + }, + }, + }); + + return logAndReturn({ + outputLog: `Combined ${payoutsToCombine.length} pending payouts for program ${programId} into ${combinedPayoutId} (deleted ${deletedPayouts.count})`, + }); +} + +async function syncLinksAndCommissions({ + targetPartnerId, + programIdsToTransfer, +}: { + targetPartnerId: string; + programIdsToTransfer: string[]; +}) { + const updatedLinks = await prisma.link.findMany({ + where: { + programId: { + in: programIdsToTransfer, + }, + partnerId: targetPartnerId, + }, + include: { + ...includeTags, + ...includeProgramEnrollment, + }, + }); + + const res = await Promise.allSettled([ + recordLink(updatedLinks), + linkCache.expireMany(updatedLinks), + ...programIdsToTransfer.map((programId) => + syncTotalCommissions({ + partnerId: targetPartnerId, + programId, + }), + ), + ]); + + // Fail the step (so QStash retries it) if any sync rejected. All of these + // ops are idempotent, so re-running the step is safe. + const rejected = res.filter( + (result): result is PromiseRejectedResult => result.status === "rejected", + ); + + if (rejected.length > 0) { + throw new Error( + `Failed to sync links/commissions: ${prettyPrint( + rejected.map(({ reason }) => reason), + )}`, + ); + } + + return logAndReturn({ + outputLog: `Synced ${updatedLinks.length} links and commissions. ${prettyPrint(res)}`, + }); +} + +async function deleteSourceUser({ sourceUserId }: { sourceUserId: string }) { + const workspaceCount = await prisma.projectUsers.count({ + where: { + userId: sourceUserId, + }, + }); + + if (workspaceCount > 0) { + return logAndReturn({ + outputLog: `User ${sourceUserId} still has ${workspaceCount} workspace(s), not deleting.`, + }); + } + + try { + const deletedUser = await prisma.user.delete({ + where: { + id: sourceUserId, + }, + select: { + id: true, + email: true, + image: true, + }, + }); + + if (deletedUser.image) { + await storage.delete({ + key: deletedUser.image.replace(`${R2_URL}/`, ""), + }); + } + + return logAndReturn({ + outputLog: `Deleted user ${deletedUser.email} (${deletedUser.id})`, + }); + } catch (error) { + return logAndReturn({ + outputLog: `Error deleting user ${sourceUserId}: ${error.message}`, + }); + } +} + +async function cleanupFraudEvents({ + sourcePartnerId, +}: { + sourcePartnerId: string; +}) { + const fraudEventsToDelete = await prisma.fraudEvent.findMany({ + where: { + partnerId: sourcePartnerId, + fraudEventGroup: { + type: FraudRuleType.partnerDuplicateAccount, + }, + }, + include: { + fraudEventGroup: { + select: { + id: true, + _count: { + select: { + fraudEvents: true, + }, + }, + }, + }, + }, + }); + + if (fraudEventsToDelete.length > 0) { + await prisma.fraudEvent.deleteMany({ + where: { + id: { in: fraudEventsToDelete.map((e) => e.id) }, + }, + }); + } + + const fraudEventGroupsToResolve = fraudEventsToDelete.filter( + // this is the count pre-deletion the fraud event, so if there are 2 fraud events + // that means post-deletion will leave 1 fraud event in the group (no additional duplicates), hence can be resolved + (e) => e.fraudEventGroup._count.fraudEvents === 2, + ); + + await resolveFraudGroups({ + where: { + OR: [ + { + partnerId: sourcePartnerId, + }, + ...(fraudEventGroupsToResolve.length > 0 + ? [ + { + id: { + in: fraudEventGroupsToResolve.map( + (e) => e.fraudEventGroup.id, + ), + }, + }, + ] + : []), + ], + type: FraudRuleType.partnerDuplicateAccount, + }, + resolutionReason: + "Automatically resolved because partners with duplicate payout methods were merged. No other partners share this payout method.", + }); + + return logAndReturn({ + outputLog: `Deleted ${fraudEventsToDelete.length} duplicate-account fraud events`, + }); +} + +async function deleteSourcePartner({ + sourcePartnerId, + sourceEmail, + sourceImage, +}: { + sourcePartnerId: string; + sourceEmail: string; + sourceImage: string | null; +}) { + await conn.execute(`DELETE FROM Partner WHERE id = ?`, [sourcePartnerId]); + + if (sourceImage) { + try { + await storage.delete({ + key: sourceImage.replace(`${R2_URL}/`, ""), + }); + } catch (error) { + logger.error("partner.image_delete_failed", { + sourcePartnerId, + sourceImage, + error, + }); + } + } + + return logAndReturn({ + outputLog: `Deleted partner ${sourceEmail} (${sourcePartnerId})`, + }); +} diff --git a/apps/web/lib/actions/partners/merge-partner-accounts.ts b/apps/web/lib/actions/partners/merge-partner-accounts.ts index afeba176f78..8fb1315f59b 100644 --- a/apps/web/lib/actions/partners/merge-partner-accounts.ts +++ b/apps/web/lib/actions/partners/merge-partner-accounts.ts @@ -1,13 +1,12 @@ "use server"; import { generateOTP } from "@/lib/auth/utils"; -import { qstash } from "@/lib/cron"; +import { triggerQStashWorkflow } from "@/lib/cron/qstash-workflow"; import { prisma } from "@/lib/prisma"; import { ratelimit, redis } from "@/lib/upstash"; import { emailSchema } from "@/lib/zod/schemas/auth"; import { sendBatchEmail } from "@dub/email"; import VerifyEmailForAccountMerge from "@dub/email/templates/verify-email-for-account-merge"; -import { APP_DOMAIN_WITH_NGROK } from "@dub/utils"; import * as z from "zod/v4"; import { authPartnerActionClient } from "../safe-action"; @@ -341,12 +340,17 @@ const mergeAccounts = async ({ userId }: { userId: string }) => { const { sourceEmail, targetEmail } = accounts; - await qstash.publishJSON({ - url: `${APP_DOMAIN_WITH_NGROK}/api/cron/partners/merge-accounts`, + await triggerQStashWorkflow({ + workflowType: "merge-partner-accounts", + workflowLabel: userId, body: { userId, sourceEmail, targetEmail, }, + flowControl: { + key: userId, + parallelism: 1, + }, }); }; diff --git a/apps/web/lib/cron/qstash-workflow.ts b/apps/web/lib/cron/qstash-workflow.ts index 62a31e427d3..a1023c9d0b9 100644 --- a/apps/web/lib/cron/qstash-workflow.ts +++ b/apps/web/lib/cron/qstash-workflow.ts @@ -1,5 +1,5 @@ import { logger, toErrorFields } from "@/lib/axiom/server"; -import { APP_DOMAIN_WITH_NGROK, pluralize } from "@dub/utils"; +import { APP_DOMAIN, pluralize } from "@dub/utils"; import { FlowControl } from "@upstash/qstash"; import { Client } from "@upstash/workflow"; @@ -14,7 +14,10 @@ const client = new Client({ }), }); -type WorkflowType = "partner-approved" | "create-partner-commission"; +type WorkflowType = + | "partner-approved" + | "create-partner-commission" + | "merge-partner-accounts"; interface QStashWorkflow { workflowType: WorkflowType; @@ -34,7 +37,7 @@ export async function triggerQStashWorkflow( try { const response = await client.trigger( workflows.map((workflow) => ({ - url: `${APP_DOMAIN_WITH_NGROK}/api/workflows/${workflow.workflowType}`, + url: `${APP_DOMAIN}/api/workflows/${workflow.workflowType}`, body: workflow.body, label: workflow.workflowLabel, retries: 5, @@ -107,6 +110,16 @@ export function getWorkflowConfig({ }; } + case "merge-partner-accounts": { + return { + correlation: { + userId: body.userId, + sourceEmail: body.sourceEmail, + targetEmail: body.targetEmail, + }, + }; + } + default: return { correlation: {}, diff --git a/apps/web/tests/partners/applications/approve-reject-partner-application.test.ts b/apps/web/tests/partner-applications/approve-reject-partner-application.test.ts similarity index 94% rename from apps/web/tests/partners/applications/approve-reject-partner-application.test.ts rename to apps/web/tests/partner-applications/approve-reject-partner-application.test.ts index cabfd33347c..21e47d4cf6b 100644 --- a/apps/web/tests/partners/applications/approve-reject-partner-application.test.ts +++ b/apps/web/tests/partner-applications/approve-reject-partner-application.test.ts @@ -1,9 +1,9 @@ import { generateRandomName } from "@/lib/names"; import { Partner } from "@prisma/client"; import { describe, expect, test } from "vitest"; -import { randomPartnerEmail } from "../../utils/helpers"; -import { IntegrationHarness } from "../../utils/integration"; -import { E2E_PARTNER_GROUP } from "../../utils/resource"; +import { randomPartnerEmail } from "../utils/helpers"; +import { IntegrationHarness } from "../utils/integration"; +import { E2E_PARTNER_GROUP } from "../utils/resource"; describe.sequential( "POST /partners/applications/reject and /approve", diff --git a/apps/web/tests/partners/applications/list-partner-applications.test.ts b/apps/web/tests/partner-applications/list-partner-applications.test.ts similarity index 94% rename from apps/web/tests/partners/applications/list-partner-applications.test.ts rename to apps/web/tests/partner-applications/list-partner-applications.test.ts index 63c9524674a..aefd43338c9 100644 --- a/apps/web/tests/partners/applications/list-partner-applications.test.ts +++ b/apps/web/tests/partner-applications/list-partner-applications.test.ts @@ -1,8 +1,8 @@ import { PartnerApplicationProps } from "@/lib/types"; import { PartnerApplicationSchema } from "@/lib/zod/schemas/program-application"; import { describe, expect, test } from "vitest"; -import { IntegrationHarness } from "../../utils/integration"; -import { E2E_PARTNER_GROUP, E2E_PARTNERS } from "../../utils/resource"; +import { IntegrationHarness } from "../utils/integration"; +import { E2E_PARTNER_GROUP, E2E_PARTNERS } from "../utils/resource"; describe.sequential("GET /partners/applications", async () => { const h = new IntegrationHarness(); diff --git a/apps/web/tests/workflows/merge-partner-accounts-workflow.test.ts b/apps/web/tests/workflows/merge-partner-accounts-workflow.test.ts new file mode 100644 index 00000000000..cc4ffa0e2ca --- /dev/null +++ b/apps/web/tests/workflows/merge-partner-accounts-workflow.test.ts @@ -0,0 +1,153 @@ +import { + VITEST_POLL_INTERVAL_MS, + VITEST_TEST_TIMEOUT_MS, +} from "@/lib/constants/misc"; +import { EnrolledPartnerProps } from "@/lib/types"; +import { describe, expect, test } from "vitest"; +import { randomPartnerEmail } from "../utils/helpers"; +import { IntegrationHarness } from "../utils/integration"; +import { E2E_PARTNER_GROUP } from "../utils/resource"; +import { verifyMergeCompleted } from "./utils/verify-merge-completed"; + +describe.sequential("Workflow - MergePartnerAccounts", async () => { + const h = new IntegrationHarness(); + const { http } = await h.init(); + + // Creates a partner enrolled (approved) in the Acme program with a default link. + async function createEnrolledPartner(label: string) { + const { status, data: partner } = await http.post({ + path: "/partners", + body: { + name: `E2E Merge ${label}`, + email: randomPartnerEmail(), + groupId: E2E_PARTNER_GROUP.id, + }, + }); + + expect(status).toEqual(201); + expect(partner.links).not.toBeNull(); + expect(partner.links!.length).toBeGreaterThan(0); + + return partner; + } + + test( + "Overlap merge transfers child data and deletes source", + { timeout: VITEST_TEST_TIMEOUT_MS }, + async () => { + const source = await createEnrolledPartner("source"); + const target = await createEnrolledPartner("target"); + const sourceLinkId = source.links![0].id; + + const { status: triggerStatus, data: triggerRes } = await http.post<{ + workflowRunId?: string; + }>({ + path: "/e2e/trigger-merge-accounts", + body: { sourceEmail: source.email, targetEmail: target.email }, + }); + + expect(triggerStatus).toEqual(200); + expect(triggerRes).not.toBeNull(); + + const merged = await verifyMergeCompleted({ + http, + sourcePartnerId: source.id, + targetPartnerId: target.id, + expectedLinkId: sourceLinkId, + }); + + expect(merged.links!.map((link) => link.id)).toContain(sourceLinkId); + }, + ); + + test( + "Overlap merge upgrades target status from pending to approved", + { timeout: VITEST_TEST_TIMEOUT_MS }, + async () => { + const source = await createEnrolledPartner("upgrade-source"); + const target = await createEnrolledPartner("upgrade-target"); + + const { status: pendingStatus } = await http.post({ + path: "/e2e/partners/pending-program-application", + body: { partnerId: target.id }, + }); + expect(pendingStatus).toEqual(200); + + const { status: triggerStatus } = await http.post({ + path: "/e2e/trigger-merge-accounts", + body: { sourceEmail: source.email, targetEmail: target.email }, + }); + expect(triggerStatus).toEqual(200); + + const startTime = Date.now(); + let lastTargetStatus: string | undefined; + + while (Date.now() - startTime < VITEST_TEST_TIMEOUT_MS) { + const [sourceRes, targetRes] = await Promise.all([ + http.get({ path: `/partners/${source.id}` }), + http.get({ path: `/partners/${target.id}` }), + ]); + + lastTargetStatus = + targetRes.status === 200 ? targetRes.data.status : undefined; + + if (sourceRes.status === 404 && lastTargetStatus === "approved") { + expect(lastTargetStatus).toBe("approved"); + return; + } + + await new Promise((resolve) => + setTimeout(resolve, VITEST_POLL_INTERVAL_MS), + ); + } + + throw new Error( + `Target status was not upgraded to approved within ${VITEST_TEST_TIMEOUT_MS / 1000}s. ` + + `Last seen status: ${lastTargetStatus}`, + ); + }, + ); + + test( + "Repeat merge is rejected once the source is already merged", + { timeout: VITEST_TEST_TIMEOUT_MS }, + async () => { + const source = await createEnrolledPartner("repeat-source"); + const target = await createEnrolledPartner("repeat-target"); + const sourceLinkId = source.links![0].id; + + const { status: firstTrigger } = await http.post({ + path: "/e2e/trigger-merge-accounts", + body: { sourceEmail: source.email, targetEmail: target.email }, + }); + expect(firstTrigger).toEqual(200); + + await verifyMergeCompleted({ + http, + sourcePartnerId: source.id, + targetPartnerId: target.id, + expectedLinkId: sourceLinkId, + }); + + const { data: afterFirst } = await http.get({ + path: `/partners/${target.id}`, + }); + const linkCountAfterFirst = afterFirst.links!.length; + + // The source partner no longer exists, so triggering again is rejected by + // the guard (both partners must be enrolled in the Acme program) - the + // merge can't be double-processed. + const { status: secondTrigger } = await http.post({ + path: "/e2e/trigger-merge-accounts", + body: { sourceEmail: source.email, targetEmail: target.email }, + }); + expect(secondTrigger).toEqual(400); + + const { data: afterSecond } = await http.get({ + path: `/partners/${target.id}`, + }); + + expect(afterSecond.links!.length).toBe(linkCountAfterFirst); + }, + ); +}); diff --git a/apps/web/tests/workflows/utils/verify-merge-completed.ts b/apps/web/tests/workflows/utils/verify-merge-completed.ts new file mode 100644 index 00000000000..491b44bc9b6 --- /dev/null +++ b/apps/web/tests/workflows/utils/verify-merge-completed.ts @@ -0,0 +1,64 @@ +import { + VITEST_POLL_INTERVAL_MS, + VITEST_TEST_TIMEOUT_MS, +} from "@/lib/constants/misc"; +import { EnrolledPartnerProps } from "@/lib/types"; +import { expect } from "vitest"; +import { HttpClient } from "../../utils/http"; + +interface VerifyMergeCompletedProps { + http: HttpClient; + sourcePartnerId: string; + targetPartnerId: string; + // A link id that belonged to the source partner and should end up on the target + expectedLinkId: string; +} + +/** + * Polls until the merge-partner-accounts workflow has finished: + * - the source partner is deleted (GET /partners/:id returns 404), and + * - the target partner now owns the source's moved link. + */ +export const verifyMergeCompleted = async ({ + http, + sourcePartnerId, + targetPartnerId, + expectedLinkId, +}: VerifyMergeCompletedProps) => { + const startTime = Date.now(); + + let lastSourceStatus: number | null = null; + let lastTargetLinkIds: string[] = []; + + while (Date.now() - startTime < VITEST_TEST_TIMEOUT_MS) { + const [sourceRes, targetRes] = await Promise.all([ + http.get<{ error?: unknown }>({ path: `/partners/${sourcePartnerId}` }), + http.get({ path: `/partners/${targetPartnerId}` }), + ]); + + lastSourceStatus = sourceRes.status; + + const sourceDeleted = sourceRes.status === 404; + const targetLinks = + targetRes.status === 200 ? targetRes.data.links ?? [] : []; + lastTargetLinkIds = targetLinks.map((link) => link.id); + const targetOwnsLink = lastTargetLinkIds.includes(expectedLinkId); + + if (sourceDeleted && targetOwnsLink) { + expect(sourceRes.status).toBe(404); + expect(lastTargetLinkIds).toContain(expectedLinkId); + return targetRes.data; + } + + await new Promise((resolve) => + setTimeout(resolve, VITEST_POLL_INTERVAL_MS), + ); + } + + throw new Error( + `Merge did not complete within ${VITEST_TEST_TIMEOUT_MS / 1000} seconds. ` + + `sourcePartnerId: ${sourcePartnerId} (last status: ${lastSourceStatus}), ` + + `targetPartnerId: ${targetPartnerId}, expectedLinkId: ${expectedLinkId}. ` + + `Last seen target link ids: [${lastTargetLinkIds.join(", ")}]`, + ); +};