diff --git a/apps/web/app/(ee)/api/stripe/integration/webhook/checkout-session-completed.ts b/apps/web/app/(ee)/api/stripe/integration/webhook/checkout-session-completed.ts index 5f4e2bc11dc..89f5be373cf 100644 --- a/apps/web/app/(ee)/api/stripe/integration/webhook/checkout-session-completed.ts +++ b/apps/web/app/(ee)/api/stripe/integration/webhook/checkout-session-completed.ts @@ -1,5 +1,9 @@ import { convertCurrency } from "@/lib/analytics/convert-currency"; import { isFirstConversion } from "@/lib/analytics/is-first-conversion"; +import { + invoiceDedupeKey, + legacyStripeInvoiceDedupeKey, +} from "@/lib/api/conversions/invoice-idempotency"; import { createId } from "@/lib/api/create-id"; import { getOrCreateCustomer } from "@/lib/api/customers/get-or-create-customer"; import { includeTags } from "@/lib/api/links/include-tags"; @@ -334,9 +338,24 @@ export async function checkoutSessionCompleted({ } if (invoiceId) { + const legacyRecord = await redis.get( + legacyStripeInvoiceDedupeKey(invoiceId), + ); + + if (legacyRecord) { + console.info( + "[checkout.session.completed] Skipping already processed invoice (legacy key).", + invoiceId, + ); + + return { + response: `Invoice with ID ${invoiceId} already processed, skipping...`, + }; + } + // Skip if invoice id is already processed const ok = await redis.set( - `trackSale:stripe:invoiceId:${invoiceId}`, // here we assume that Stripe's invoice ID is unique across all customers + invoiceDedupeKey(workspace.id, invoiceId), { timestamp: new Date().toISOString(), dubCustomerExternalId, diff --git a/apps/web/app/(ee)/api/stripe/integration/webhook/invoice-paid.ts b/apps/web/app/(ee)/api/stripe/integration/webhook/invoice-paid.ts index 9cae9a50370..d9d9fef2dd6 100644 --- a/apps/web/app/(ee)/api/stripe/integration/webhook/invoice-paid.ts +++ b/apps/web/app/(ee)/api/stripe/integration/webhook/invoice-paid.ts @@ -1,5 +1,9 @@ import { convertCurrency } from "@/lib/analytics/convert-currency"; import { isFirstConversion } from "@/lib/analytics/is-first-conversion"; +import { + invoiceDedupeKey, + legacyStripeInvoiceDedupeKey, +} from "@/lib/api/conversions/invoice-idempotency"; import { includeTags } from "@/lib/api/links/include-tags"; import { syncPartnerLinksStats } from "@/lib/api/partners/sync-partner-links-stats"; import { executeWorkflows } from "@/lib/api/workflows/execute-workflows"; @@ -135,9 +139,21 @@ export async function invoicePaid({ ? invoice.total_excluding_tax : invoice.amount_paid; + const legacyRecord = await redis.get(legacyStripeInvoiceDedupeKey(invoiceId)); + + if (legacyRecord) { + console.info( + "[invoice.paid] Skipping already processed invoice (legacy key).", + invoiceId, + ); + return { + response: `Invoice with ID ${invoiceId} already processed, skipping...`, + }; + } + // Skip if invoice id is already processed const ok = await redis.set( - `trackSale:stripe:invoiceId:${invoiceId}`, // here we assume that Stripe's invoice ID is unique across all customers + invoiceDedupeKey(workspace.id, invoiceId), { timestamp: new Date().toISOString(), dubCustomerExternalId: customer.externalId, diff --git a/apps/web/lib/api/conversions/invoice-idempotency.ts b/apps/web/lib/api/conversions/invoice-idempotency.ts new file mode 100644 index 00000000000..a79f1bff608 --- /dev/null +++ b/apps/web/lib/api/conversions/invoice-idempotency.ts @@ -0,0 +1,9 @@ +// TODO: remove after 2026-08-10 (10 days after rollout) once the transition +// window has elapsed and no keys written under the old Stripe-only format +// remain (they carry a 7-day TTL). + +export const invoiceDedupeKey = (workspaceId: string, invoiceId: string) => + `trackSale:${workspaceId}:invoiceId:${invoiceId}`; + +export const legacyStripeInvoiceDedupeKey = (invoiceId: string) => + `trackSale:stripe:invoiceId:${invoiceId}`; diff --git a/apps/web/lib/api/conversions/track-sale.ts b/apps/web/lib/api/conversions/track-sale.ts index b22e04de093..323bb457f44 100644 --- a/apps/web/lib/api/conversions/track-sale.ts +++ b/apps/web/lib/api/conversions/track-sale.ts @@ -32,6 +32,7 @@ import * as z from "zod/v4"; import { createId } from "../create-id"; import { syncPartnerLinksStats } from "../partners/sync-partner-links-stats"; import { executeWorkflows } from "../workflows/execute-workflows"; +import { invoiceDedupeKey } from "./invoice-idempotency"; type TrackSaleParams = z.input & { workspace: Pick; @@ -62,8 +63,9 @@ export const trackSale = async ({ // Return idempotent response if invoiceId is already processed if (invoiceId) { const cachedResponse = await redis.get( - `trackSale:${workspace.id}:invoiceId:${invoiceId}`, + invoiceDedupeKey(workspace.id, invoiceId), ); + if (cachedResponse) { return cachedResponse; } @@ -707,13 +709,9 @@ const _trackSale = async ({ if (invoiceId) { waitUntil( - redis.set( - `trackSale:${workspace.id}:invoiceId:${invoiceId}`, - trackSaleResponse, - { - ex: 60 * 60 * 24 * 7, // cache for 1 week - }, - ), + redis.set(invoiceDedupeKey(workspace.id, invoiceId), trackSaleResponse, { + ex: 60 * 60 * 24 * 7, // cache for 1 week + }), ); }