From 1bb3ed47d5aae3c2bf0ab0fe0a712cb541c0284c Mon Sep 17 00:00:00 2001 From: Pierre Brisorgueil Date: Mon, 28 Sep 2026 21:23:10 +0200 Subject: [PATCH 01/15] fix(billing): redeliver checkout webhook on Stripe retrieve failure stripe.subscriptions.retrieve failing in the checkout.session.completed handler returned silently, so the event was recorded as processed while the subscription row stayed unlinked on the free plan. Throw instead so withIdempotency lets Stripe redeliver the event. Claude-Session: https://claude.ai/code/session_01EqkUXh5nmxbVBTvhgM6zAo --- .../services/billing.webhook.service.js | 16 +++++++++---- .../billing.webhook.checkout.unit.tests.js | 24 ++++++++++++------- 2 files changed, 26 insertions(+), 14 deletions(-) diff --git a/modules/billing/services/billing.webhook.service.js b/modules/billing/services/billing.webhook.service.js index fa54b2ea9..3acdef204 100644 --- a/modules/billing/services/billing.webhook.service.js +++ b/modules/billing/services/billing.webhook.service.js @@ -245,9 +245,10 @@ const handleCheckoutCompleted = async (session, event) => { } // Fetch real status from Stripe — never assume 'active' (could be 'trialing', 'incomplete', etc.) - // On retrieval failure we abort: persisting 'active' on a failed fetch would silently - // misclassify trialing/incomplete subscriptions and bypass dunning. The next webhook - // (subscription.updated or invoice.payment_failed) will reconcile the correct status. + // On retrieval failure we throw (see catch below): persisting 'active' on a failed + // fetch would silently misclassify trialing/incomplete subscriptions and bypass + // dunning, and swallowing the failure would let the event be recorded as processed + // while the subscription row stays unlinked (#4151). Throwing lets Stripe redeliver. const stripe = getStripe(); if (!stripe) { logger.error('[billing.webhook] checkout.session.completed — Stripe not configured, aborting', { @@ -266,11 +267,16 @@ const handleCheckoutCompleted = async (session, event) => { return; } } catch (err) { - logger.error('[billing.webhook] checkout.session.completed — subscription retrieve failed, aborting to avoid stale active assumption', { + // Throw (do NOT swallow) — a silent return here still lets withIdempotency record + // the event as processed while the subscription row stays unlinked on the free + // plan. Throwing propagates through withIdempotency's catch (attempts persists, + // the claim doc is never deleted on failure), so the controller returns a 5xx and + // Stripe redelivers the event instead (#4151). + logger.error('[billing.webhook] checkout.session.completed — subscription retrieve failed, will retry via Stripe redelivery', { stripeSubscriptionId, error: err?.message ?? String(err), }); - return; + throw err; } const existing = await SubscriptionRepository.findByOrganization(organizationId); diff --git a/modules/billing/tests/billing.webhook.checkout.unit.tests.js b/modules/billing/tests/billing.webhook.checkout.unit.tests.js index 1f6ed5664..6e62260d9 100644 --- a/modules/billing/tests/billing.webhook.checkout.unit.tests.js +++ b/modules/billing/tests/billing.webhook.checkout.unit.tests.js @@ -429,17 +429,23 @@ describe('Billing webhook checkout unit tests:', () => { ); }); - test('should abort without persisting when stripe.subscriptions.retrieve throws', async () => { + test('should throw and not persist when stripe.subscriptions.retrieve fails, so Stripe redelivers the event', async () => { + // #4151: a silent `return` here recorded the event as processed while the + // subscription row stayed unlinked on the free plan. Throwing propagates through + // withIdempotency (attempts persist, no delete-on-failure) so the controller + // returns 5xx and Stripe redelivers instead. mockStripeInstance.subscriptions.retrieve.mockRejectedValue(new Error('Stripe API error')); - await BillingWebhookService.handleCheckoutCompleted( - { - customer: 'cus_123', - subscription: 'sub_456', - metadata: { organizationId: orgId, plan: 'pro' }, - }, - checkoutEvent, - ); + await expect( + BillingWebhookService.handleCheckoutCompleted( + { + customer: 'cus_123', + subscription: 'sub_456', + metadata: { organizationId: orgId, plan: 'pro' }, + }, + checkoutEvent, + ), + ).rejects.toThrow('Stripe API error'); // Should not persist anything — aborting to avoid stale 'active' assumption expect(mockSubscriptionRepository.updateIfEventNewer).not.toHaveBeenCalled(); From d89970fbab79565bece696e7c9d8dbe8150bb943 Mon Sep 17 00:00:00 2001 From: Pierre Brisorgueil Date: Mon, 28 Sep 2026 21:24:04 +0200 Subject: [PATCH 02/15] fix(billing): invoice webhooks no longer resurrect a canceled subscription A late invoice.payment_failed or invoice.payment_succeeded event could bring a canceled subscription back to past_due/active. Both handlers now return early once the stored subscription status is 'canceled'. Claude-Session: https://claude.ai/code/session_01EqkUXh5nmxbVBTvhgM6zAo --- .../services/billing.webhook.service.js | 6 +++++ ...billing.webhook.subscription.unit.tests.js | 23 +++++++++++++++++++ 2 files changed, 29 insertions(+) diff --git a/modules/billing/services/billing.webhook.service.js b/modules/billing/services/billing.webhook.service.js index 3acdef204..276a8cd76 100644 --- a/modules/billing/services/billing.webhook.service.js +++ b/modules/billing/services/billing.webhook.service.js @@ -645,6 +645,9 @@ const handleInvoicePaymentFailed = async (invoice, event) => { const existing = await SubscriptionRepository.findByStripeSubscriptionId(stripeSubscriptionId); if (!existing) return; + // A late/out-of-order invoice event must not bring a canceled subscription back to + // past_due — the subscription no longer exists on Stripe (#4151). + if (existing.status === 'canceled') return; const fields = { status: 'past_due' }; @@ -691,6 +694,9 @@ const handleInvoicePaymentSucceeded = async (invoice, event) => { const existing = await SubscriptionRepository.findByStripeSubscriptionId(stripeSubscriptionId); if (!existing) return; + // A late invoice.payment_succeeded (e.g. a delayed Stripe retry queued before the + // cancellation) must not resurrect a canceled subscription to 'active' (#4151). + if (existing.status === 'canceled') return; // Always advance the invoice-family marker (lastInvoiceEventCreatedAt / lastInvoiceEventId) // so stale replays of older invoice events are correctly rejected by the ordering guard. diff --git a/modules/billing/tests/billing.webhook.subscription.unit.tests.js b/modules/billing/tests/billing.webhook.subscription.unit.tests.js index 2ed3069a1..bd7b6d7b3 100644 --- a/modules/billing/tests/billing.webhook.subscription.unit.tests.js +++ b/modules/billing/tests/billing.webhook.subscription.unit.tests.js @@ -526,6 +526,18 @@ describe('Billing webhook subscription unit tests:', () => { expect(mockSubscriptionRepository.updateIfEventNewer).not.toHaveBeenCalled(); }); + test('#4151 — should return early without writing when subscription is already canceled', async () => { + // A late invoice.payment_succeeded (e.g. a delayed retry Stripe queued before the + // cancellation) must not bring a canceled subscription back to 'active'. + const existing = { _id: subId, organization: orgId, status: 'canceled', pastDueSince: null }; + mockSubscriptionRepository.findByStripeSubscriptionId.mockResolvedValue(existing); + + await BillingWebhookService.handleInvoicePaymentSucceeded({ subscription: 'sub_456' }, makeEvent()); + + expect(mockSubscriptionRepository.updateIfEventNewer).not.toHaveBeenCalled(); + expect(mockStripe.subscriptions.retrieve).not.toHaveBeenCalled(); + }); + test('should log info when event is stale (V5 P1 #1 ordering guard)', async () => { const existing = { _id: subId, @@ -800,5 +812,16 @@ describe('Billing webhook subscription unit tests:', () => { expect(mockSubscriptionRepository.updateIfEventNewer).not.toHaveBeenCalled(); }); + + test('#4151 — should return early without writing when subscription is already canceled', async () => { + // A late invoice.payment_failed for a canceled subscription must not touch it — + // there is nothing to mark past_due on a subscription that no longer exists on Stripe. + const existing = { _id: subId, organization: orgId, status: 'canceled', pastDueSince: null }; + mockSubscriptionRepository.findByStripeSubscriptionId.mockResolvedValue(existing); + + await BillingWebhookService.handleInvoicePaymentFailed({ subscription: 'sub_456' }, makeEvent()); + + expect(mockSubscriptionRepository.updateIfEventNewer).not.toHaveBeenCalled(); + }); }); }); From 19b0b37bbe9075531d44b9908c3d62cfbd562e49 Mon Sep 17 00:00:00 2001 From: Pierre Brisorgueil Date: Mon, 28 Sep 2026 21:26:25 +0200 Subject: [PATCH 03/15] fix(billing): meter bills the free plan on a fail-closed subscription status incrementMeter kept billing against the paid-plan quota snapshot even when the admission gate had already routed the same subscription status to the free plan, letting the paid quota drain for free. The fail-closed status list is now a single shared constant, imported by both the gate and the meter, so they cannot drift apart again. Claude-Session: https://claude.ai/code/session_01EqkUXh5nmxbVBTvhgM6zAo --- modules/billing/lib/constants.js | 10 ++++ .../billing.subscription.repository.js | 13 +++-- .../billing/services/billing.quota.service.js | 7 +-- .../billing/services/billing.usage.service.js | 10 +++- ...ling.subscription.repository.unit.tests.js | 9 ++-- .../tests/billing.usage.service.unit.tests.js | 53 +++++++++++++++++++ 6 files changed, 89 insertions(+), 13 deletions(-) diff --git a/modules/billing/lib/constants.js b/modules/billing/lib/constants.js index 31a7a4068..ec791bba1 100644 --- a/modules/billing/lib/constants.js +++ b/modules/billing/lib/constants.js @@ -3,3 +3,13 @@ * Any status not in this list is treated as inactive (falls back to free plan). */ export const activeStatuses = ['active', 'trialing']; + +/** + * Stripe subscription statuses that fail closed — billed and gated as if the + * organization were on the default (free) plan. Shared by billing.quota.service.js + * (the admission gate) and billing.usage.service.js (the meter) so the two + * enforcement paths can never drift apart: duplicating this list in both places let + * the gate treat a fail-closed status as free while the meter kept consuming the + * paid quota (#4151). + */ +export const failClosedStatuses = ['paused', 'unpaid', 'incomplete_expired', 'incomplete', 'canceled']; diff --git a/modules/billing/repositories/billing.subscription.repository.js b/modules/billing/repositories/billing.subscription.repository.js index 4f9aa141a..9ce318951 100644 --- a/modules/billing/repositories/billing.subscription.repository.js +++ b/modules/billing/repositories/billing.subscription.repository.js @@ -112,16 +112,19 @@ const findByStripeSubscriptionId = (stripeSubscriptionId) => { /** * @function findPlan - * @description Lean lookup that returns only the `plan` field for a given organization. - * Used on hot paths (meter attribution, weekly reset) where only the plan - * identifier is needed — avoids the full populate overhead of findByOrganization. + * @description Lean lookup that returns the `plan` and `status` fields for a given + * organization. Used on hot paths (meter attribution, weekly reset) where + * only these two fields are needed — avoids the full populate overhead of + * findByOrganization. `status` lets callers detect a fail-closed + * subscription (see `failClosedStatuses` in `../lib/constants.js`) without + * a second query (#4151). * @param {String} organizationId - The organization ID. - * @returns {Promise<{plan: string}|null>} A lean plain object with just `plan`, or null. + * @returns {Promise<{plan: string, status: string}|null>} A lean plain object, or null. */ // biome-ignore lint/correctness/useQwikValidLexicalScope: false positive — Node.js repository, not Qwik const findPlan = (organizationId) => { if (!mongoose.Types.ObjectId.isValid(organizationId)) return null; - return Subscription.findOne({ organization: organizationId }, { plan: 1 }).lean().exec(); + return Subscription.findOne({ organization: organizationId }, { plan: 1, status: 1 }).lean().exec(); }; /** diff --git a/modules/billing/services/billing.quota.service.js b/modules/billing/services/billing.quota.service.js index 1bf3381d5..eab8243af 100644 --- a/modules/billing/services/billing.quota.service.js +++ b/modules/billing/services/billing.quota.service.js @@ -23,7 +23,7 @@ import BillingUsageService from './billing.usage.service.js'; import BillingExtraBalanceRepository from '../repositories/billing.extraBalance.repository.js'; import BillingPlanService from './billing.plan.service.js'; -import { activeStatuses } from '../lib/constants.js'; +import { activeStatuses, failClosedStatuses } from '../lib/constants.js'; import { getDefaultPlanId, getGracePeriodDays } from '../lib/billing.constants.js'; import config from '../../../config/index.js'; import AppError from '../../../lib/helpers/AppError.js'; @@ -75,8 +75,9 @@ async function assertCanExecute({ orgId, organization, user, resource, action }) if (config.billing?.meterMode === true) { const subscription = await SubscriptionRepository.findByOrganization(orgIdStr); - // Fail-closed statuses → route to free plan quota - const failClosedStatuses = ['paused', 'unpaid', 'incomplete_expired', 'incomplete', 'canceled']; + // Fail-closed statuses → route to free plan quota. Shared with + // billing.usage.service.js's incrementMeter so the gate and the meter can never + // diverge on which statuses are treated as free (#4151). if (subscription && failClosedStatuses.includes(subscription.status)) { const planId = getDefaultPlanId(); const freePlan = BillingPlanService.getActivePlan(planId); diff --git a/modules/billing/services/billing.usage.service.js b/modules/billing/services/billing.usage.service.js index 6d48bacd3..9108a417f 100644 --- a/modules/billing/services/billing.usage.service.js +++ b/modules/billing/services/billing.usage.service.js @@ -10,6 +10,7 @@ import BillingExtraService from './billing.extra.service.js'; import billingEvents from '../lib/events.js'; import { currentWeekKey } from '../lib/billing.isoWeek.js'; import { getAlertThresholdPercents, getDefaultPlanId } from '../lib/billing.constants.js'; +import { failClosedStatuses } from '../lib/constants.js'; /** * Compute the current month string in YYYY-MM format. @@ -115,9 +116,14 @@ const incrementMeter = async (organizationId, units, breakdown, idempotencyKey) const weekKey = currentWeekKey(); const monthKey = currentMonth(); - // Fetch active plan for quota snapshot — config-static, no DB read + // Fetch active plan for quota snapshot — config-static, no DB read. + // Fail-closed statuses (paused/unpaid/incomplete/incomplete_expired/canceled) bill + // against the default (free) plan — same list and same fallback the admission gate + // uses (billing.quota.service.js) — otherwise the gate treats the org as free while + // the meter keeps consuming the stale paid-plan quota snapshot (#4151). const subscription = await BillingSubscriptionRepository.findPlan(organizationId); - const planId = subscription?.plan ?? getDefaultPlanId(); + const isFailClosed = subscription && failClosedStatuses.includes(subscription.status); + const planId = isFailClosed ? getDefaultPlanId() : (subscription?.plan ?? getDefaultPlanId()); const activePlan = BillingPlanService.getActivePlan(planId); const meterQuota = activePlan?.meterQuota ?? 0; const planVersion = activePlan?.version ?? null; diff --git a/modules/billing/tests/billing.subscription.repository.unit.tests.js b/modules/billing/tests/billing.subscription.repository.unit.tests.js index d6e4f0f47..1630c588b 100644 --- a/modules/billing/tests/billing.subscription.repository.unit.tests.js +++ b/modules/billing/tests/billing.subscription.repository.unit.tests.js @@ -226,8 +226,11 @@ describe('BillingSubscriptionRepository unit tests:', () => { // ── findPlan ────────────────────────────────────────────────────────────── describe('findPlan', () => { - test('returns plan field for a valid organizationId', async () => { - const planDoc = { _id: subId, plan: 'pro' }; + // #4151: widened to also project `status` so callers (incrementMeter) can detect a + // fail-closed subscription without a second query — same fields the admission gate + // (billing.quota.service.js) already reads. + test('returns plan and status fields for a valid organizationId', async () => { + const planDoc = { _id: subId, plan: 'pro', status: 'active' }; const execMock = jest.fn().mockResolvedValue(planDoc); const leanMock = jest.fn().mockReturnValue({ exec: execMock }); mockModel.findOne.mockReturnValue({ lean: leanMock }); @@ -236,7 +239,7 @@ describe('BillingSubscriptionRepository unit tests:', () => { expect(mockModel.findOne).toHaveBeenCalledWith( { organization: orgId }, - { plan: 1 }, + { plan: 1, status: 1 }, ); expect(leanMock).toHaveBeenCalled(); expect(result).toEqual(planDoc); diff --git a/modules/billing/tests/billing.usage.service.unit.tests.js b/modules/billing/tests/billing.usage.service.unit.tests.js index 9ed16310f..92b606ea9 100644 --- a/modules/billing/tests/billing.usage.service.unit.tests.js +++ b/modules/billing/tests/billing.usage.service.unit.tests.js @@ -243,7 +243,60 @@ describe('BillingUsageService — meter extensions unit tests:', () => { expect(mockPlanService.getActivePlan).toHaveBeenCalledWith('free'); }); + }); + + /** + * #4151 — billing.quota.service.js (the admission gate) already treats a fail-closed + * subscription status (paused/unpaid/incomplete/incomplete_expired/canceled) as the + * free plan. incrementMeter must bill the SAME way, or the gate lets the org through + * on free-plan grounds while the meter keeps consuming the paid quota snapshot. + */ + describe('incrementMeter — fail-closed statuses bill against the default (free) plan', () => { + test('unpaid subscription bills the default plan quota, not the stale paid-plan snapshot', async () => { + mockConfig.billing.defaultPlan = 'free'; + mockSubscriptionRepository.findPlan.mockResolvedValue({ plan: 'pro', status: 'unpaid' }); + mockPlanService.getActivePlan.mockImplementation((planId) => + makePlan({ planId, meterQuota: planId === 'free' ? 0 : 500000 })); + mockUsageRepository.incrementMeter.mockResolvedValue(makeUsageDoc({ meterUsed: 10, meterQuota: 0 })); + + const result = await BillingUsageService.incrementMeter(orgId, 10, { scrap: 10 }, 'hist_failclosed'); + + expect(mockPlanService.getActivePlan).toHaveBeenCalledWith('free'); + expect(mockPlanService.getActivePlan).not.toHaveBeenCalledWith('pro'); + expect(result.meterQuota).toBe(0); + // Free-plan behavior: every unit is billed to extras, none absorbed by a quota + // the org no longer has access to. + expect(result.extrasConsumed).toBe(10); + }); + + test('every status in the shared fail-closed list routes to the default plan (reuse, not a second literal list)', async () => { + mockConfig.billing.defaultPlan = 'free'; + mockPlanService.getActivePlan.mockImplementation((planId) => + makePlan({ planId, meterQuota: planId === 'free' ? 0 : 500000 })); + mockUsageRepository.incrementMeter.mockResolvedValue(makeUsageDoc({ meterUsed: 1, meterQuota: 0 })); + + const { failClosedStatuses } = await import('../lib/constants.js'); + expect(failClosedStatuses.length).toBeGreaterThan(0); + + for (const status of failClosedStatuses) { + mockSubscriptionRepository.findPlan.mockResolvedValue({ plan: 'enterprise', status }); + await BillingUsageService.incrementMeter(orgId, 1, {}, `hist_${status}`); + expect(mockPlanService.getActivePlan).toHaveBeenLastCalledWith('free'); + } + }); + + test('active subscription still bills its own paid plan (fail-closed guard does not over-fire)', async () => { + mockSubscriptionRepository.findPlan.mockResolvedValue({ plan: 'pro', status: 'active' }); + mockPlanService.getActivePlan.mockReturnValue(makePlan()); + mockUsageRepository.incrementMeter.mockResolvedValue(makeUsageDoc({ meterUsed: 10 })); + + await BillingUsageService.incrementMeter(orgId, 10, {}, 'hist_active'); + + expect(mockPlanService.getActivePlan).toHaveBeenCalledWith('pro'); + }); + }); + describe('incrementMeter — replay', () => { test('should return applied=false and fetch existing doc on replay', async () => { mockSubscriptionRepository.findPlan.mockResolvedValue({ plan: 'pro' }); mockPlanService.getActivePlan.mockReturnValue(makePlan()); From 73e72d34e3d5e6b7bd1608cbdd219eaae0045c2c Mon Sep 17 00:00:00 2001 From: Pierre Brisorgueil Date: Mon, 28 Sep 2026 21:27:10 +0200 Subject: [PATCH 04/15] fix(billing): retry the extras debit on a transient DB error A transient write failure on the extras debit fell straight through the overflow branch with no retry, so a paying org's overage was never charged. Debit is idempotent by refId, so wrapping it in the module's existing retryWithBackoff cannot double-charge on the retried write. Claude-Session: https://claude.ai/code/session_01EqkUXh5nmxbVBTvhgM6zAo --- .../billing/services/billing.extra.service.js | 6 ++++- .../tests/billing.extra.service.unit.tests.js | 23 +++++++++++++++++++ 2 files changed, 28 insertions(+), 1 deletion(-) diff --git a/modules/billing/services/billing.extra.service.js b/modules/billing/services/billing.extra.service.js index b7d17b509..d1c1714b4 100644 --- a/modules/billing/services/billing.extra.service.js +++ b/modules/billing/services/billing.extra.service.js @@ -6,6 +6,7 @@ import logger from '../../../lib/services/logger.js'; import billingEvents from '../lib/events.js'; import BillingExtraBalanceRepository from '../repositories/billing.extraBalance.repository.js'; import { SENTINEL_PENDING } from '../lib/billing.constants.js'; +import { retryWithBackoff } from '../lib/billing.retry.js'; /** * @function creditPack @@ -40,6 +41,9 @@ const creditPack = async (orgId, packId, stripeSessionId) => { * @function debit * @description Debit meter units from the extra balance. * Returns applied=false if balance is insufficient or refId already used. + * Wrapped in retryWithBackoff: the write is idempotent by refId (a retry + * cannot double-charge), so a transient DB error must not fall through and + * silently make the overflow free (#4151). * * @param {string} orgId - The organization ObjectId (string). * @param {number} units - Meter units to debit (must be > 0). @@ -48,7 +52,7 @@ const creditPack = async (orgId, packId, stripeSessionId) => { */ // biome-ignore lint/correctness/useQwikValidLexicalScope: false positive — Node.js service, not Qwik const debit = (orgId, units, refId) => - BillingExtraBalanceRepository.debit(orgId, units, refId); + retryWithBackoff(() => BillingExtraBalanceRepository.debit(orgId, units, refId)); /** * @function getOrgBalanceContext diff --git a/modules/billing/tests/billing.extra.service.unit.tests.js b/modules/billing/tests/billing.extra.service.unit.tests.js index 15434908d..e294aaadd 100644 --- a/modules/billing/tests/billing.extra.service.unit.tests.js +++ b/modules/billing/tests/billing.extra.service.unit.tests.js @@ -143,6 +143,29 @@ describe('BillingExtraService unit tests:', () => { expect(r1.applied).toBe(true); expect(r2.applied).toBe(false); }); + + // #4151 — a transient DB error must not make the overflow free: debit is idempotent + // by refId (see 'debit same refId twice' above), so retrying it cannot double-charge. + test('retries a transient repository failure and succeeds without double-crediting', async () => { + const doc = makeDoc({ cachedBalance: 300000 }); + mockRepository.debit + .mockRejectedValueOnce(new Error('Mongo topology closed')) + .mockResolvedValueOnce({ doc, applied: true }); + + const result = await BillingExtraService.debit(orgId, 100000, 'ref_transient'); + + expect(mockRepository.debit).toHaveBeenCalledTimes(2); + expect(mockRepository.debit).toHaveBeenNthCalledWith(1, orgId, 100000, 'ref_transient'); + expect(mockRepository.debit).toHaveBeenNthCalledWith(2, orgId, 100000, 'ref_transient'); + expect(result.applied).toBe(true); + expect(result.doc).toBe(doc); + }); + + test('gives up after repeated transient failures and surfaces the error (never silently frees the overflow)', async () => { + mockRepository.debit.mockRejectedValue(new Error('Mongo topology closed')); + + await expect(BillingExtraService.debit(orgId, 100000, 'ref_down')).rejects.toThrow('Mongo topology closed'); + }); }); describe('getOrgBalanceContext', () => { From 6425ecfccb82192563fe87c373803c68f018416f Mon Sep 17 00:00:00 2001 From: Pierre Brisorgueil Date: Mon, 28 Sep 2026 21:28:25 +0200 Subject: [PATCH 05/15] fix(billing): creditGrant idempotency guard checks every org's ledger The idempotency check before a grant write only excluded a matching refId within the target org's own ledger, so a retry whose target org changed between attempts (e.g. an org merge/reassignment) was credited twice under the same refId. The guard now checks across all orgs before writing. Claude-Session: https://claude.ai/code/session_01EqkUXh5nmxbVBTvhgM6zAo --- .../billing.extraBalance.repository.js | 12 +++++++++++- .../tests/billing.extraBalance.unit.tests.js | 17 +++++++++++++++++ 2 files changed, 28 insertions(+), 1 deletion(-) diff --git a/modules/billing/repositories/billing.extraBalance.repository.js b/modules/billing/repositories/billing.extraBalance.repository.js index fbbbde948..7b18f0b49 100644 --- a/modules/billing/repositories/billing.extraBalance.repository.js +++ b/modules/billing/repositories/billing.extraBalance.repository.js @@ -169,7 +169,11 @@ const debit = async (orgId, amount, refId) => { * @description Atomically credit extra meter units for a non-Stripe grant (e.g. signup free * tier, referral grant). Idempotent: if a ledger entry with the same refId * already exists, the update is a no-op and applied=false is returned. - * 2-step pattern aligned with creditPack: + * Idempotency is checked ACROSS ALL orgs, not just this one: a retry whose + * target org changed between attempts (e.g. an org merge/reassignment) must + * not be credited twice under the same refId (#4151). + * 3-step pattern aligned with creditPack: + * Step 0 — cross-org existence check (any org's ledger, not just this one). * Step 1 — ensure doc exists (atomic getOrCreate, no-op on replay). * Step 2 — idempotency-guarded credit (no upsert). * Idempotency key: `options.refId` when supplied (#3842 referral grants — @@ -212,6 +216,12 @@ const creditGrant = async (orgId, amount, source, { refId = null, expiresAt = nu ...(expiresAt ? { expiresAt } : {}), }; + // Step 0: cross-org idempotency guard. The step-2 filter below only excludes a + // matching refId within THIS org's own ledger, so it misses a replay whose target + // org changed between attempts. Check across every org's ledger before writing. + const alreadyGranted = await BillingExtraBalance().exists({ 'ledger.refId': idempotencyKey }); + if (alreadyGranted) return { doc: null, applied: false, reason: 'duplicate_grant' }; + // Step 1: ensure the document exists (atomic getOrCreate, no-op if already present). await getOrCreate(orgId); diff --git a/modules/billing/tests/billing.extraBalance.unit.tests.js b/modules/billing/tests/billing.extraBalance.unit.tests.js index 6238a5f11..5463b5ca1 100644 --- a/modules/billing/tests/billing.extraBalance.unit.tests.js +++ b/modules/billing/tests/billing.extraBalance.unit.tests.js @@ -1184,6 +1184,23 @@ describe('BillingExtraBalance unit tests:', () => { expect(result).toEqual({ doc: null, applied: false }); expect(mockModel.findOneAndUpdate).not.toHaveBeenCalled(); }); + + // #4151 — the idempotency guard in step 2 below only excludes a matching refId + // WITHIN this org's own ledger ('ledger.refId': { $ne }), so a retry whose target + // org changed between attempts (e.g. an org merge/reassignment) is credited twice. + test('returns applied:false duplicate_grant when the refId was already granted to a DIFFERENT org', async () => { + // Cross-org existence check finds a prior grant under this refId in any org. + mockModel.exists.mockResolvedValue({ _id: 'some-other-orgs-balance-doc' }); + + const result = await BillingExtraBalanceRepository.creditGrant(orgId, 500, 'referral', { + refId: 'referral:64b2f0000000000000000001:referrer', + }); + + expect(result).toEqual({ doc: null, applied: false, reason: 'duplicate_grant' }); + expect(mockModel.exists).toHaveBeenCalledWith({ 'ledger.refId': 'referral:64b2f0000000000000000001:referrer' }); + // Must short-circuit before any write. + expect(mockModel.findOneAndUpdate).not.toHaveBeenCalled(); + }); }); }); }); From 245b3371a6afd70293281128d2f101f7789e7c0f Mon Sep 17 00:00:00 2001 From: Pierre Brisorgueil Date: Mon, 28 Sep 2026 21:31:43 +0200 Subject: [PATCH 06/15] fix(auth): close the email-verification concurrent-provisioning race verifyEmail read the user by token, then wrote it in a separate step, so two concurrent requests for the same token could both pass the read check and both provision a workspace. UserService now exposes an atomic consumeEmailVerificationToken (one findOneAndUpdate on the still- unexpired token) that verifyEmail consumes instead. Claude-Session: https://claude.ai/code/session_01EqkUXh5nmxbVBTvhgM6zAo --- modules/auth/controllers/auth.controller.js | 17 ++- .../auth.verifyEmail.grant.unit.tests.js | 13 ++- .../auth.verifyEmail.signup-org.unit.tests.js | 9 +- .../users/repositories/users.repository.js | 29 +++++ modules/users/services/users.service.js | 13 +++ ...ationToken.concurrent.integration.tests.js | 109 ++++++++++++++++++ 6 files changed, 173 insertions(+), 17 deletions(-) create mode 100644 modules/users/tests/users.consumeEmailVerificationToken.concurrent.integration.tests.js diff --git a/modules/auth/controllers/auth.controller.js b/modules/auth/controllers/auth.controller.js index 2b775ae21..729e6cf10 100644 --- a/modules/auth/controllers/auth.controller.js +++ b/modules/auth/controllers/auth.controller.js @@ -750,18 +750,15 @@ const getConfig = async (req, res) => { */ const verifyEmail = async (req, res) => { try { - const user = await UserService.getBrut({ emailVerificationToken: req.params.token }); - const isExpired = !user?.emailVerificationExpires || Number(user.emailVerificationExpires) < Date.now(); - if (!user || !user.email || isExpired) { + // Atomic read+write (#4151): a separate getBrut() read then update() write let two + // concurrent requests for the same token both pass the read check and both + // provision a workspace. consumeEmailVerificationToken does both in one + // findOneAndUpdate — only the first concurrent caller can match the still-unexpired + // token; the second gets null and falls into the same 400 as an invalid token. + const user = await UserService.consumeEmailVerificationToken(req.params.token); + if (!user) { return responses.error(res, 400, 'Bad Request', 'Email verification token is invalid or has expired.')(); } - await UserService.update(user, { - emailVerified: true, - emailVerificationToken: null, - emailVerificationExpires: null, - }, 'recover'); - // Mark verified on the local object so handleSignupOrganization sees emailVerified=true - user.emailVerified = true; // Post-verification org setup — provision org/grant if not yet done (best-effort). // handleSignupOrganization is idempotent: if the org already exists it converges diff --git a/modules/auth/tests/auth.verifyEmail.grant.unit.tests.js b/modules/auth/tests/auth.verifyEmail.grant.unit.tests.js index 76d16ce11..9b6cae9b5 100644 --- a/modules/auth/tests/auth.verifyEmail.grant.unit.tests.js +++ b/modules/auth/tests/auth.verifyEmail.grant.unit.tests.js @@ -32,8 +32,15 @@ describe('verifyEmail — org provisioning after email verification (non-fatal): jest.resetModules(); mockUserService = { - getBrut: jest.fn().mockResolvedValue({ ...fakeUser }), - update: jest.fn().mockResolvedValue({}), + // #4151: verifyEmail consumes the token via one atomic findOneAndUpdate instead + // of a separate getBrut() read + update() write — the resolved doc already + // carries emailVerified: true (set by the atomic write), not mutated locally. + consumeEmailVerificationToken: jest.fn().mockResolvedValue({ + ...fakeUser, + emailVerified: true, + emailVerificationToken: null, + emailVerificationExpires: null, + }), }; mockOrganizationsService = { @@ -162,7 +169,7 @@ describe('verifyEmail — org provisioning after email verification (non-fatal): }); test('invalid/expired token — handleSignupOrganization NOT called', async () => { - mockUserService.getBrut.mockResolvedValue(null); + mockUserService.consumeEmailVerificationToken.mockResolvedValue(null); const req = makeReq('bad_token'); const res = makeRes(); diff --git a/modules/auth/tests/auth.verifyEmail.signup-org.unit.tests.js b/modules/auth/tests/auth.verifyEmail.signup-org.unit.tests.js index 7c83edc14..465fbd77d 100644 --- a/modules/auth/tests/auth.verifyEmail.signup-org.unit.tests.js +++ b/modules/auth/tests/auth.verifyEmail.signup-org.unit.tests.js @@ -22,14 +22,15 @@ describe('auth.controller verifyEmail — handleSignupOrganization wiring:', () mockUserService = { create: jest.fn(), - getBrut: jest.fn().mockResolvedValue({ + // #4151: verifyEmail consumes the token via one atomic findOneAndUpdate instead + // of a separate getBrut() read + update() write — the resolved doc already + // carries emailVerified: true (set by the atomic write), not mutated locally. + consumeEmailVerificationToken: jest.fn().mockResolvedValue({ _id: 'user_001', id: 'user_001', email: 'user@example.com', - emailVerificationToken: 'tok', - emailVerificationExpires: Date.now() + 3600000, + emailVerified: true, }), - update: jest.fn().mockResolvedValue({}), remove: jest.fn(), search: jest.fn(), count: jest.fn().mockResolvedValue(0), diff --git a/modules/users/repositories/users.repository.js b/modules/users/repositories/users.repository.js index b33b689c9..4652925d6 100644 --- a/modules/users/repositories/users.repository.js +++ b/modules/users/repositories/users.repository.js @@ -107,6 +107,34 @@ const update = (user) => { return new User(user).save(); }; +/** + * @desc Atomically verify an email address: only a document whose + * emailVerificationToken matches AND emailVerificationExpires is still in the + * future is updated, in the SAME findOneAndUpdate that reads it — closing the + * read-then-write race where two concurrent requests for the same token both + * pass a separate read check and both provision a workspace. A second + * concurrent call, or a replay after the token was already consumed, matches + * no document and returns null. + * @param {String} token - The raw emailVerificationToken from the verification link. + * @returns {Object|null} the updated user document, or null when the token is + * missing, unknown, expired, or already consumed. + */ +const consumeEmailVerificationToken = (token) => { + if (!token) return Promise.resolve(null); + return User.findOneAndUpdate( + { + emailVerificationToken: token, + emailVerificationExpires: { $gt: Date.now() }, + }, + { + emailVerified: true, + emailVerificationToken: null, + emailVerificationExpires: null, + }, + { returnDocument: 'after', runValidators: true }, + ).exec(); +}; + /** * @desc Function to remove a user from db by id or email * @param {Object} user @@ -252,6 +280,7 @@ export default { get, search, update, + consumeEmailVerificationToken, remove, stats, count, diff --git a/modules/users/services/users.service.js b/modules/users/services/users.service.js index 49564f653..4017ed852 100644 --- a/modules/users/services/users.service.js +++ b/modules/users/services/users.service.js @@ -148,6 +148,18 @@ const update = async (user, body, option) => { return removeSensitive(result); }; +/** + * @desc Atomically consume an email-verification token: verifies the email and clears + * the token in one write (see UserRepository.consumeEmailVerificationToken), + * instead of a separate read-then-write that lets two concurrent requests for + * the same token both pass the read check. + * @param {String} token - The raw emailVerificationToken from the verification link. + * @returns {Promise} full ("brut", unsanitized — same convention as + * getBrut) user document, or null when the token could not be atomically consumed + * (unknown, expired, or already used). + */ +const consumeEmailVerificationToken = (token) => UserRepository.consumeEmailVerificationToken(token); + /** * @desc Function to ask repository to sign terms for current user * @param {Object} user - original user document @@ -288,6 +300,7 @@ export default { get, getBrut, update, + consumeEmailVerificationToken, terms, remove, stats, diff --git a/modules/users/tests/users.consumeEmailVerificationToken.concurrent.integration.tests.js b/modules/users/tests/users.consumeEmailVerificationToken.concurrent.integration.tests.js new file mode 100644 index 000000000..ef1878b76 --- /dev/null +++ b/modules/users/tests/users.consumeEmailVerificationToken.concurrent.integration.tests.js @@ -0,0 +1,109 @@ +/** + * Module dependencies. + */ +import path from 'path'; +import mongoose from 'mongoose'; +import { beforeAll, afterEach, afterAll, describe, test, expect } from '@jest/globals'; + +import { bootstrap } from '../../../lib/app.js'; + +/** + * #4151 — verifyEmail used to read the user (getBrut) then write it (update) in two + * separate steps. Two concurrent requests for the SAME token both passed the read + * check before either write landed, so both provisioned an organization/grant for the + * same signup. UserService.consumeEmailVerificationToken closes the race with one + * atomic findOneAndUpdate on the still-unexpired token: only the first concurrent + * caller can match it, the second gets null. + * + * Uses the real bootstrapped app + real User model against the test Mongo — the + * atomicity guarantee this test proves cannot be faked with mocks (see + * feedback_stale_evidence_worse_than_none). + */ +describe('UserService.consumeEmailVerificationToken — concurrent verification race (#4151):', () => { + let UserService; + let User; + + beforeAll(async () => { + await bootstrap(); + UserService = (await import(path.resolve('./modules/users/services/users.service.js'))).default; + User = mongoose.model('User'); + }); + + afterEach(async () => { + await User.deleteMany({ email: { $regex: /^verify-race-/ } }).exec(); + }); + + afterAll(async () => { + await User.deleteMany({ email: { $regex: /^verify-race-/ } }).exec(); + }); + + test('two concurrent requests for the same token: exactly one succeeds, the other gets null', async () => { + const token = `tok_${Date.now()}_${Math.random().toString(36).slice(2)}`; + const created = await User.create({ + firstName: 'Race', + lastName: 'Condition', + email: 'verify-race-1@example.com', + provider: 'local', + emailVerified: false, + emailVerificationToken: token, + emailVerificationExpires: Date.now() + 3600000, + }); + + const [r1, r2] = await Promise.all([ + UserService.consumeEmailVerificationToken(token), + UserService.consumeEmailVerificationToken(token), + ]); + + const results = [r1, r2]; + const succeeded = results.filter(Boolean); + const failed = results.filter((r) => r === null); + + // Exactly ONE of the two concurrent requests may consume the token — the atomic + // findOneAndUpdate's filter (unexpired token) can only ever match once, so the + // second concurrent caller cannot re-provision a second workspace for this user. + expect(succeeded).toHaveLength(1); + expect(failed).toHaveLength(1); + expect(String(succeeded[0]._id)).toBe(String(created._id)); + expect(succeeded[0].emailVerified).toBe(true); + + const dbUser = await User.findById(created._id).lean(); + expect(dbUser.emailVerified).toBe(true); + expect(dbUser.emailVerificationToken).toBeNull(); + expect(dbUser.emailVerificationExpires).toBeNull(); + }); + + test('an expired token is never consumed', async () => { + const token = `tok_expired_${Date.now()}`; + await User.create({ + firstName: 'Expired', + lastName: 'Token', + email: 'verify-race-2@example.com', + provider: 'local', + emailVerified: false, + emailVerificationToken: token, + emailVerificationExpires: Date.now() - 1000, + }); + + const result = await UserService.consumeEmailVerificationToken(token); + expect(result).toBeNull(); + }); + + test('a replay after the token was already consumed returns null (single-use)', async () => { + const token = `tok_replay_${Date.now()}`; + await User.create({ + firstName: 'Replay', + lastName: 'Guard', + email: 'verify-race-3@example.com', + provider: 'local', + emailVerified: false, + emailVerificationToken: token, + emailVerificationExpires: Date.now() + 3600000, + }); + + const first = await UserService.consumeEmailVerificationToken(token); + const second = await UserService.consumeEmailVerificationToken(token); + + expect(first).not.toBeNull(); + expect(second).toBeNull(); + }); +}); From a5469716de8efe453f2593449472724f2f8c1d4c Mon Sep 17 00:00:00 2001 From: Pierre Brisorgueil Date: Mon, 28 Sep 2026 21:31:55 +0200 Subject: [PATCH 07/15] docs(billing): runbook line for the customer-paid-no-credit symptom Names the existing dead-letter replay endpoint for the specific symptom of a customer who paid but never received credit. Claude-Session: https://claude.ai/code/session_01EqkUXh5nmxbVBTvhgM6zAo --- modules/billing/RUNBOOKS.md | 2 ++ 1 file changed, 2 insertions(+) diff --git a/modules/billing/RUNBOOKS.md b/modules/billing/RUNBOOKS.md index d1aec7301..5644425b9 100644 --- a/modules/billing/RUNBOOKS.md +++ b/modules/billing/RUNBOOKS.md @@ -56,6 +56,8 @@ Operational runbooks for the billing module. Each runbook references real endpoi **Context**: Stripe webhook events that fail processing 5+ times (or where the idempotency guard fires on a poisoned payload) are marked `deadLetter: true` in `processedStripeEvents`. They accumulate and must be reviewed manually — partial TTL index excludes them from auto-expiry. +**Common trigger:** customer paid but no credit → `POST /api/admin/billing/webhook/replay {eventId}` (an event lost mid-handler is replayed by hand). + **Steps**: 1. List all dead-letter events: From ba1f95612d645f5daf87c234a20ad7e727e1a038 Mon Sep 17 00:00:00 2001 From: Pierre Brisorgueil Date: Mon, 28 Sep 2026 21:38:24 +0200 Subject: [PATCH 08/15] fix(auth): consumeEmailVerificationToken rejects an account with no email The atomic filter dropped the old controller-level !user.email guard. Restore the equivalent check in the filter itself so a token belonging to an emailless account still falls into the same invalid-token 400. Claude-Session: https://claude.ai/code/session_01EqkUXh5nmxbVBTvhgM6zAo --- .../users/repositories/users.repository.js | 16 ++++++++------ ...ationToken.concurrent.integration.tests.js | 22 +++++++++++++++++++ 2 files changed, 31 insertions(+), 7 deletions(-) diff --git a/modules/users/repositories/users.repository.js b/modules/users/repositories/users.repository.js index 4652925d6..f24f4445d 100644 --- a/modules/users/repositories/users.repository.js +++ b/modules/users/repositories/users.repository.js @@ -109,15 +109,16 @@ const update = (user) => { /** * @desc Atomically verify an email address: only a document whose - * emailVerificationToken matches AND emailVerificationExpires is still in the - * future is updated, in the SAME findOneAndUpdate that reads it — closing the - * read-then-write race where two concurrent requests for the same token both - * pass a separate read check and both provision a workspace. A second - * concurrent call, or a replay after the token was already consumed, matches - * no document and returns null. + * emailVerificationToken matches, emailVerificationExpires is still in the + * future, AND email is a non-empty string is updated, in the SAME + * findOneAndUpdate that reads it — closing the read-then-write race where two + * concurrent requests for the same token both pass a separate read check and + * both provision a workspace. A second concurrent call, a replay after the + * token was already consumed, or a doc with no email (mirrors the previous + * controller-level `!user.email` guard), matches no document and returns null. * @param {String} token - The raw emailVerificationToken from the verification link. * @returns {Object|null} the updated user document, or null when the token is - * missing, unknown, expired, or already consumed. + * missing, unknown, expired, already consumed, or the account has no email. */ const consumeEmailVerificationToken = (token) => { if (!token) return Promise.resolve(null); @@ -125,6 +126,7 @@ const consumeEmailVerificationToken = (token) => { { emailVerificationToken: token, emailVerificationExpires: { $gt: Date.now() }, + email: { $exists: true, $nin: [null, ''] }, }, { emailVerified: true, diff --git a/modules/users/tests/users.consumeEmailVerificationToken.concurrent.integration.tests.js b/modules/users/tests/users.consumeEmailVerificationToken.concurrent.integration.tests.js index ef1878b76..00b1f43ed 100644 --- a/modules/users/tests/users.consumeEmailVerificationToken.concurrent.integration.tests.js +++ b/modules/users/tests/users.consumeEmailVerificationToken.concurrent.integration.tests.js @@ -72,6 +72,28 @@ describe('UserService.consumeEmailVerificationToken — concurrent verification expect(dbUser.emailVerificationExpires).toBeNull(); }); + test('a user with no email is never verified (mirrors the old !user.email guard)', async () => { + const token = `tok_noemail_${Date.now()}`; + // Schema allows a missing email (no `required`); the old controller code checked + // `!user.email` after a separate read — the atomic filter must reject the same case. + const created = await User.create({ + firstName: 'No', + lastName: 'Email', + provider: 'local', + emailVerified: false, + emailVerificationToken: token, + emailVerificationExpires: Date.now() + 3600000, + }); + + try { + const result = await UserService.consumeEmailVerificationToken(token); + expect(result).toBeNull(); + } finally { + // Not covered by the afterEach's email-regex cleanup (this doc has no email). + await User.deleteOne({ _id: created._id }).exec(); + } + }); + test('an expired token is never consumed', async () => { const token = `tok_expired_${Date.now()}`; await User.create({ From 3341ff1c6eb7f73730db392a7ab4a4f0e5a8e2ee Mon Sep 17 00:00:00 2001 From: Pierre Brisorgueil Date: Mon, 28 Sep 2026 21:38:49 +0200 Subject: [PATCH 09/15] docs(errors): log the four non-obvious #4151 fixes Checkout webhook silent return, fail-closed meter/gate list drift, creditGrant's per-org-only idempotency guard, and the verifyEmail read-then-write race. Claude-Session: https://claude.ai/code/session_01EqkUXh5nmxbVBTvhgM6zAo --- ERRORS.md | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/ERRORS.md b/ERRORS.md index bcf8c7f91..333c9bfa8 100644 --- a/ERRORS.md +++ b/ERRORS.md @@ -49,3 +49,7 @@ Use this file as a compact memory of recurring AI mistakes. - [2026-09-25] billing/stripe: `billing.plans.service.js fetchPlansFromStripe` fell back to the raw Stripe product id (`product.metadata?.planId || product.id`) when a product carried no `planId` metadata -> ANY active Stripe product (a one-time pack, a recurring product sold outside the plans catalogue via a Payment Link) advertised itself as a public plan via `GET /api/billing/plans`, often with a null price id; fix = filter to `product.metadata?.planId` truthy BEFORE mapping, drop the id fallback entirely — a product now needs `metadata.planId` to be listed. Left every OTHER `metadata?.planId ||` fallback untouched (`billing.planResolver.js`, `billing.webhook.service.js`, `billing.admin.service.js`) since those resolve an EXISTING subscription's plan, not the public catalogue, and must keep working for a recurring product sold outside it; see pierreb-devkit/Node#4113 - [2026-09-25] auth: `oauthCallback` never provisioned an org, and review of the first fix caught that gating on `!user.currentOrganization` would also fire for an EXISTING org-less user (removed from org, pending join) on every login -> gate on `info.created` instead (set by `checkOAuthUserProfile`'s create branch, relayed via passport's verify-callback `info`), so only a genuine new signup provisions. `oauthCallback` also needed an outer try/catch (passport invokes it fire-and-forget) and a `headersSent` guard before any fallback redirect; see pierreb-devkit/Node#4115 - [2026-09-25] billing: computing a percent-of-grant level as `(1 - threshold/100) * signupGrant` (e.g. `500 * (1 - 80/100)`) lands on `99.99999999999997`, not `100`, due to float imprecision -> silently misses an exact `post === level` boundary crossing; use `signupGrant * (100 - threshold) / 100` instead, which is exact at common values; see pierreb-devkit/Node#4117 +- [2026-09-28] billing/stripe: `checkout.session.completed`'s subscription-retrieve `catch` returned silently on failure -> the event was recorded as processed by the idempotency wrapper while the subscription row stayed unlinked on the free plan, with no retry; fix = throw instead of return, so the claim persists (attempts increments, the doc is never deleted on failure) and Stripe redelivers; see pierreb-devkit/Node#4151 +- [2026-09-28] billing/meter: `incrementMeter` always billed against the subscription's own plan quota, while the admission gate (`billing.quota.service.js`) already treated a fail-closed status (`paused`/`unpaid`/`incomplete`/`incomplete_expired`/`canceled`) as the free plan -> a fail-closed org kept consuming its stale paid-plan quota on the meter even though the gate let it through on free-plan grounds; the fail-closed status list is now one shared constant imported by both the gate and the meter — a literal duplicated in two enforcement paths drifts by design, not by mistake; see pierreb-devkit/Node#4151 +- [2026-09-28] billing/extras: `creditGrant`'s idempotency guard excluded a matching `refId` only within the TARGET org's own ledger -> a retry whose target org changed between attempts (e.g. an org merge/reassignment) was credited twice under the same `refId`; fix = a cross-org existence check on the `refId` before the `getOrCreate`/write steps; see pierreb-devkit/Node#4151 +- [2026-09-28] auth: `verifyEmail` read the user by token (`getBrut`) then wrote it in a separate step (`update`) -> two concurrent requests for the same token both passed the read check before either write landed, so both provisioned an organization/grant for the same signup; fix = one atomic `findOneAndUpdate` (`UserService.consumeEmailVerificationToken`) that verifies and clears the token together, so only the first concurrent caller can match the still-unexpired token; see pierreb-devkit/Node#4151 From b95a1c3cafacedfb3da34d925ac0edc8f7cbf870 Mon Sep 17 00:00:00 2001 From: Pierre Brisorgueil Date: Tue, 29 Sep 2026 09:27:29 +0200 Subject: [PATCH 10/15] fix(billing): invoice write re-checks canceled status atomically MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The #4151 guard read existing.status before the write but updateIfEventNewer only filtered on the event-ordering markers, not status — a customer.subscription.deleted that committed between the read and the write could still have its canceled row flipped back to past_due or active by a late invoice event. updateIfEventNewer gains an optional extraMatch param, ANDed into the same findOneAndUpdate filter as the event-ordering guard. The two invoice handlers pass status: { $ne: 'canceled' }, so the exclusion is re-checked at write time instead of relying solely on the earlier read. The read-side check stays as a fast path. Addresses a CodeRabbit review comment on #4155. Claude-Session: https://claude.ai/code/session_01EqkUXh5nmxbVBTvhgM6zAo --- .../billing.subscription.repository.js | 17 ++++++++-- .../services/billing.webhook.service.js | 18 +++++++++-- .../tests/billing.service.unit.tests.js | 1 + ...iption.repository.per-family.unit.tests.js | 31 +++++++++++++++++++ .../billing.webhook.integration.tests.js | 1 + ...billing.webhook.subscription.unit.tests.js | 7 +++++ 6 files changed, 69 insertions(+), 6 deletions(-) diff --git a/modules/billing/repositories/billing.subscription.repository.js b/modules/billing/repositories/billing.subscription.repository.js index 9ce318951..6b444582b 100644 --- a/modules/billing/repositories/billing.subscription.repository.js +++ b/modules/billing/repositories/billing.subscription.repository.js @@ -227,16 +227,26 @@ const markUnpaid = (id, threshold) => { * 'subscription' family (customer.subscription.*) or the 'invoice' family * (invoice.*). Same-second cross-family deliveries no longer cancel each other. * - * Returns null when the guard prevents the write (stale event). + * Returns null when the guard prevents the write (stale event, or — when + * `extraMatch` excludes it — a status the write must not apply over, e.g. a + * canceled subscription). * @param {string} id - The subscription ObjectId (string). * @param {number} eventCreatedAt - Stripe event.created Unix timestamp (seconds). * @param {string} eventId - Stripe event.id (e.g. evt_xxx) — tiebreaker for same-second delivery. * @param {Object} fields - Fields to $set on the document. * @param {'subscription'|'invoice'} [family='subscription'] - Event family to scope the ordering guard. - * @returns {Promise} Updated doc, or null if id invalid or event is stale. + * @param {Object} [extraMatch={}] - Additional top-level filter fields ANDed into the + * query alongside the event-ordering guard, evaluated atomically in the same + * findOneAndUpdate. Used by the invoice handlers to require + * `status: { $ne: 'canceled' }` so a cancellation that commits between the + * caller's own status read and this write still wins — the write's own filter + * is re-checked at write time, closing the read-then-write race the read-side + * check alone cannot (#4155). + * @returns {Promise} Updated doc, or null if id invalid, event is stale, or + * `extraMatch` no longer matches. */ // biome-ignore lint/correctness/useQwikValidLexicalScope: false positive — Node.js repository, not Qwik -const updateIfEventNewer = (id, eventCreatedAt, eventId, fields, family = 'subscription') => { +const updateIfEventNewer = (id, eventCreatedAt, eventId, fields, family = 'subscription', extraMatch = {}) => { if (!id || !mongoose.Types.ObjectId.isValid(id)) return null; const createdAtField = family === 'invoice' ? 'lastInvoiceEventCreatedAt' : 'lastSubscriptionEventCreatedAt'; @@ -245,6 +255,7 @@ const updateIfEventNewer = (id, eventCreatedAt, eventId, fields, family = 'subsc return Subscription.findOneAndUpdate( { _id: id, + ...extraMatch, $or: [ { [createdAtField]: { $exists: false } }, { [createdAtField]: null }, diff --git a/modules/billing/services/billing.webhook.service.js b/modules/billing/services/billing.webhook.service.js index 276a8cd76..b18adf54e 100644 --- a/modules/billing/services/billing.webhook.service.js +++ b/modules/billing/services/billing.webhook.service.js @@ -646,7 +646,9 @@ const handleInvoicePaymentFailed = async (invoice, event) => { const existing = await SubscriptionRepository.findByStripeSubscriptionId(stripeSubscriptionId); if (!existing) return; // A late/out-of-order invoice event must not bring a canceled subscription back to - // past_due — the subscription no longer exists on Stripe (#4151). + // past_due — the subscription no longer exists on Stripe (#4151). Fast-path only: a + // concurrent cancellation can still commit after this read, so the write below + // re-asserts `status: { $ne: 'canceled' }` atomically (#4155). if (existing.status === 'canceled') return; const fields = { status: 'past_due' }; @@ -656,7 +658,14 @@ const handleInvoicePaymentFailed = async (invoice, event) => { fields.pastDueSince = new Date(); } - const updated = await SubscriptionRepository.updateIfEventNewer(String(existing._id), event.created, event.id, fields, 'invoice'); + const updated = await SubscriptionRepository.updateIfEventNewer( + String(existing._id), + event.created, + event.id, + fields, + 'invoice', + { status: { $ne: 'canceled' } }, + ); if (!updated) { logger.info('[billing.webhook] skipped stale event', { eventId: event.id, type: event.type }); return; @@ -695,7 +704,9 @@ const handleInvoicePaymentSucceeded = async (invoice, event) => { const existing = await SubscriptionRepository.findByStripeSubscriptionId(stripeSubscriptionId); if (!existing) return; // A late invoice.payment_succeeded (e.g. a delayed Stripe retry queued before the - // cancellation) must not resurrect a canceled subscription to 'active' (#4151). + // cancellation) must not resurrect a canceled subscription to 'active' (#4151). Fast-path + // only: a concurrent cancellation can still commit after this read, so the write below + // re-asserts `status: { $ne: 'canceled' }` atomically (#4155). if (existing.status === 'canceled') return; // Always advance the invoice-family marker (lastInvoiceEventCreatedAt / lastInvoiceEventId) @@ -737,6 +748,7 @@ const handleInvoicePaymentSucceeded = async (invoice, event) => { event.id, fields, 'invoice', + { status: { $ne: 'canceled' } }, ); if (!updated) { logger.info('[billing.webhook] skipped stale event', { eventId: event.id, type: event.type }); diff --git a/modules/billing/tests/billing.service.unit.tests.js b/modules/billing/tests/billing.service.unit.tests.js index f738a4c2e..1316ab901 100644 --- a/modules/billing/tests/billing.service.unit.tests.js +++ b/modules/billing/tests/billing.service.unit.tests.js @@ -368,6 +368,7 @@ describe('Billing webhook service unit tests:', () => { 'evt_fail', expect.objectContaining({ status: 'past_due', pastDueSince: expect.any(Date) }), 'invoice', + { status: { $ne: 'canceled' } }, ); }); diff --git a/modules/billing/tests/billing.subscription.repository.per-family.unit.tests.js b/modules/billing/tests/billing.subscription.repository.per-family.unit.tests.js index 5108d2514..d80c21356 100644 --- a/modules/billing/tests/billing.subscription.repository.per-family.unit.tests.js +++ b/modules/billing/tests/billing.subscription.repository.per-family.unit.tests.js @@ -91,4 +91,35 @@ describe('updateIfEventNewer — per-family guard:', () => { expect(update.$set.lastSubscriptionEventCreatedAt).toBe(300); expect(update.$set.lastInvoiceEventCreatedAt).toBeUndefined(); }); + + // #4155 — extraMatch: an invoice-family write must not resurrect a subscription a + // concurrent cancellation already moved to 'canceled' between the caller's own status + // read and this write. The exclusion has to live in the SAME findOneAndUpdate filter as + // the event-ordering guard so both are checked atomically at write time. + test('extraMatch is merged into the top-level filter alongside the event-ordering guard', async () => { + await SubscriptionRepository.updateIfEventNewer( + '507f1f77bcf86cd799439011', + 400, + 'evt_4', + { status: 'active' }, + 'invoice', + { status: { $ne: 'canceled' } }, + ); + + expect(mockModel.findOneAndUpdate).toHaveBeenCalledTimes(1); + const [filter] = mockModel.findOneAndUpdate.mock.calls[0]; + + expect(filter.status).toEqual({ $ne: 'canceled' }); + // The event-ordering guard still applies — extraMatch adds to it, not replaces it. + expect(Array.isArray(filter.$or)).toBe(true); + }); + + test('extraMatch defaults to {} — no unrelated top-level filter field is added', async () => { + await SubscriptionRepository.updateIfEventNewer('507f1f77bcf86cd799439011', 500, 'evt_5', { plan: 'pro' }, 'subscription'); + + expect(mockModel.findOneAndUpdate).toHaveBeenCalledTimes(1); + const [filter] = mockModel.findOneAndUpdate.mock.calls[0]; + + expect(filter.status).toBeUndefined(); + }); }); diff --git a/modules/billing/tests/billing.webhook.integration.tests.js b/modules/billing/tests/billing.webhook.integration.tests.js index 44ec15c1e..b2cccebc5 100644 --- a/modules/billing/tests/billing.webhook.integration.tests.js +++ b/modules/billing/tests/billing.webhook.integration.tests.js @@ -494,6 +494,7 @@ describe('Billing webhook integration tests:', () => { 'evt_failed', expect.objectContaining({ status: 'past_due' }), 'invoice', + { status: { $ne: 'canceled' } }, ); }); diff --git a/modules/billing/tests/billing.webhook.subscription.unit.tests.js b/modules/billing/tests/billing.webhook.subscription.unit.tests.js index bd7b6d7b3..6deed916b 100644 --- a/modules/billing/tests/billing.webhook.subscription.unit.tests.js +++ b/modules/billing/tests/billing.webhook.subscription.unit.tests.js @@ -462,6 +462,7 @@ describe('Billing webhook subscription unit tests:', () => { 'evt_succeeded', expect.objectContaining({ pastDueSince: null, status: 'active' }), 'invoice', + { status: { $ne: 'canceled' } }, ); }); @@ -485,6 +486,7 @@ describe('Billing webhook subscription unit tests:', () => { 'evt_succeeded', {}, 'invoice', + { status: { $ne: 'canceled' } }, ); }); @@ -509,6 +511,7 @@ describe('Billing webhook subscription unit tests:', () => { 'evt_succeeded', expect.objectContaining({ pastDueSince: null, status: 'active' }), 'invoice', + { status: { $ne: 'canceled' } }, ); }); @@ -584,6 +587,7 @@ describe('Billing webhook subscription unit tests:', () => { 'evt_succeeded', expect.objectContaining({ plan: 'pro', status: 'active', pastDueSince: null }), 'invoice', + { status: { $ne: 'canceled' } }, ); expect(mockOrganizationRepository.setPlan).toHaveBeenCalledWith(orgId, 'pro'); }); @@ -610,6 +614,7 @@ describe('Billing webhook subscription unit tests:', () => { 'evt_succeeded', expect.not.objectContaining({ plan: expect.anything() }), 'invoice', + { status: { $ne: 'canceled' } }, ); }); @@ -728,6 +733,7 @@ describe('Billing webhook subscription unit tests:', () => { 'evt_succeeded', expect.objectContaining({ plan: 'pro', status: 'active', pastDueSince: null }), 'invoice', + { status: { $ne: 'canceled' } }, ); expect(mockOrganizationRepository.setPlan).toHaveBeenCalledWith(orgId, 'pro'); expect(mockEvents.emit).toHaveBeenCalledWith( @@ -752,6 +758,7 @@ describe('Billing webhook subscription unit tests:', () => { 'evt_failed', expect.objectContaining({ status: 'past_due' }), 'invoice', + { status: { $ne: 'canceled' } }, ); }); From 70c97abf8757f7b4a62f093326c12f9cfc0e8d2b Mon Sep 17 00:00:00 2001 From: Pierre Brisorgueil Date: Tue, 29 Sep 2026 09:27:40 +0200 Subject: [PATCH 11/15] fix(billing): creditGrant serializes cross-org idempotency-key claims The #4151 cross-org exists() check closed the sequential-retry case but was still a check-then-write race for genuine concurrency: two creditGrant calls sharing one idempotencyKey but resolving to DIFFERENT orgs (e.g. an in-process grant listener racing a reconcile sweep for the same event) could both observe "not yet granted" before either wrote, each crediting its own org. creditGrant now wraps its existing 3-step sequence (cross-org exists check, getOrCreate, org-scoped guarded write) in a per-idempotencyKey lock from the existing distributedLock service (a unique-_id claim, already used for cron mutual exclusion). Losing the lock is treated like any other idempotent replay: { applied: false, reason: 'duplicate_grant' }. The org-scoped ledger guard remains the durable dedup once the lock is released. Addresses a CodeRabbit review comment on #4155. Claude-Session: https://claude.ai/code/session_01EqkUXh5nmxbVBTvhgM6zAo --- .../billing.extraBalance.repository.js | 99 ++++++++++++++----- ...ance.creditGrant.race.integration.tests.js | 89 +++++++++++++++++ .../tests/billing.extraBalance.unit.tests.js | 39 +++++++- 3 files changed, 199 insertions(+), 28 deletions(-) create mode 100644 modules/billing/tests/billing.extraBalance.creditGrant.race.integration.tests.js diff --git a/modules/billing/repositories/billing.extraBalance.repository.js b/modules/billing/repositories/billing.extraBalance.repository.js index 7b18f0b49..329576d21 100644 --- a/modules/billing/repositories/billing.extraBalance.repository.js +++ b/modules/billing/repositories/billing.extraBalance.repository.js @@ -1,6 +1,7 @@ /** * Module dependencies */ +import { randomUUID } from 'node:crypto'; import mongoose from 'mongoose'; import AppError from '../../../lib/helpers/AppError.js'; import BillingExtraBalanceSchema from '../models/billing.extraBalance.schema.js'; @@ -164,18 +165,49 @@ const debit = async (orgId, amount, refId) => { return { doc: null, applied: false, reason: 'duplicate_step' }; }; +/** + * Lock TTL for the creditGrant critical section (existence check + getOrCreate + guarded + * write). Generous relative to the operation's normal sub-100ms cost so a slow write never + * loses its own lock mid-flight; short enough that a crashed holder does not block a retry + * for long — the caller's own retry (Stripe redelivery) or the referral reconcile cron + * picks it back up well within this window. + */ +const GRANT_LOCK_TTL_MS = 15_000; + /** * @function creditGrant * @description Atomically credit extra meter units for a non-Stripe grant (e.g. signup free * tier, referral grant). Idempotent: if a ledger entry with the same refId * already exists, the update is a no-op and applied=false is returned. - * Idempotency is checked ACROSS ALL orgs, not just this one: a retry whose + * Idempotency is enforced ACROSS ALL orgs, not just this one: a retry whose * target org changed between attempts (e.g. an org merge/reassignment) must * not be credited twice under the same refId (#4151). - * 3-step pattern aligned with creditPack: + * + * A cross-org check-then-write is not atomic by itself — two concurrent calls + * sharing the same idempotencyKey but resolving to DIFFERENT orgs (e.g. the + * in-process referral listener racing the reconcile cron backfill for the same + * invitation) can both observe "not yet granted" before either writes, and both + * would credit their own org (#4155, CodeRabbit). This is closed with a + * database-enforced claim: `lib/services/distributedLock.js` upserts a lock + * document on a unique `_id`, so only ONE concurrent call can hold the key at a + * time — the same primitive the billing crons already use for pod-level mutual + * exclusion, here scoped per-idempotencyKey instead of per-cron. + * + * Sequence while the lock is held: * Step 0 — cross-org existence check (any org's ledger, not just this one). * Step 1 — ensure doc exists (atomic getOrCreate, no-op on replay). * Step 2 — idempotency-guarded credit (no upsert). + * The lock only serializes this sequence for a given key; the org-scoped + * `'ledger.refId': { $ne: key }` guard in Step 2 remains the durable dedup once + * the lock is released (a later replay of the SAME org+key still resolves + * correctly without needing the lock to still be held). + * Lock not acquired (another call for this exact key is mid-flight) → treated + * like any other idempotent replay: `{ applied: false, reason: 'duplicate_grant' }`. + * A request that loses this race is not lost — referral grants are re-derived + * and retried by crons/billing.referralReconcile.js; the signup-grant synthetic + * key (`-`) embeds orgId and so can never collide across orgs by + * construction, making this path relevant to the referral (explicit refId) case. + * * Idempotency key: `options.refId` when supplied (#3842 referral grants — * several grants per org, one per invitation, e.g. * `referral::referrer`), otherwise the synthetic per-org key @@ -216,31 +248,46 @@ const creditGrant = async (orgId, amount, source, { refId = null, expiresAt = nu ...(expiresAt ? { expiresAt } : {}), }; - // Step 0: cross-org idempotency guard. The step-2 filter below only excludes a - // matching refId within THIS org's own ledger, so it misses a replay whose target - // org changed between attempts. Check across every org's ledger before writing. - const alreadyGranted = await BillingExtraBalance().exists({ 'ledger.refId': idempotencyKey }); - if (alreadyGranted) return { doc: null, applied: false, reason: 'duplicate_grant' }; - - // Step 1: ensure the document exists (atomic getOrCreate, no-op if already present). - await getOrCreate(orgId); - - // Step 2: idempotency-guarded credit (no upsert — doc is guaranteed to exist after step 1). - const doc = await BillingExtraBalance().findOneAndUpdate( - { - organization: orgId, - 'ledger.refId': { $ne: idempotencyKey }, - }, - { - $push: { ledger: entry }, - $inc: { cachedBalance: amount }, - $set: { cachedBalanceAt: new Date() }, - }, - { returnDocument: 'after' }, - ); + // Database-enforced claim on the idempotencyKey (see JSDoc): serializes the cross-org + // check + write below against any other creditGrant call for this exact key. Dynamic + // import — this is the only function in the repository that needs the lock's own + // CronLock model, so no other test in this file has to account for a second + // mongoose.model() registration it never exercises. + const { acquireLock, releaseLock } = await import('../../../lib/services/distributedLock.js'); + const lockName = `billing.creditGrant.${idempotencyKey}`; + const lockHolder = `${process.env.HOSTNAME ?? 'unknown'}:${randomUUID()}`; + const locked = await acquireLock({ name: lockName, ttlMs: GRANT_LOCK_TTL_MS, holder: lockHolder }); + if (!locked) return { doc: null, applied: false, reason: 'duplicate_grant' }; + + try { + // Step 0: cross-org idempotency guard. The step-2 filter below only excludes a + // matching refId within THIS org's own ledger, so it misses a replay whose target + // org changed between attempts. Check across every org's ledger before writing. + const alreadyGranted = await BillingExtraBalance().exists({ 'ledger.refId': idempotencyKey }); + if (alreadyGranted) return { doc: null, applied: false, reason: 'duplicate_grant' }; + + // Step 1: ensure the document exists (atomic getOrCreate, no-op if already present). + await getOrCreate(orgId); + + // Step 2: idempotency-guarded credit (no upsert — doc is guaranteed to exist after step 1). + const doc = await BillingExtraBalance().findOneAndUpdate( + { + organization: orgId, + 'ledger.refId': { $ne: idempotencyKey }, + }, + { + $push: { ledger: entry }, + $inc: { cachedBalance: amount }, + $set: { cachedBalanceAt: new Date() }, + }, + { returnDocument: 'after' }, + ); - if (doc) return { doc, applied: true }; - return { doc: null, applied: false, reason: 'duplicate_grant' }; + if (doc) return { doc, applied: true }; + return { doc: null, applied: false, reason: 'duplicate_grant' }; + } finally { + await releaseLock({ name: lockName, holder: lockHolder }); + } }; /** diff --git a/modules/billing/tests/billing.extraBalance.creditGrant.race.integration.tests.js b/modules/billing/tests/billing.extraBalance.creditGrant.race.integration.tests.js new file mode 100644 index 000000000..5d34e25ca --- /dev/null +++ b/modules/billing/tests/billing.extraBalance.creditGrant.race.integration.tests.js @@ -0,0 +1,89 @@ +/** + * Module dependencies. + */ +import mongoose from 'mongoose'; +import { describe, beforeAll, beforeEach, afterAll, test, expect } from '@jest/globals'; + +import mongooseService from '../../../lib/services/mongoose.js'; + +/** + * Integration tests for creditGrant cross-organization idempotency-key race closure (#4155). + * + * A cross-org check-then-write is not atomic by itself: two concurrent creditGrant calls + * sharing the same idempotencyKey but targeting DIFFERENT orgs (e.g. the in-process + * referral listener racing the reconcile cron backfill for the same invitation) could + * both observe "not yet granted" before either writes, and both would credit their own + * org — double-crediting the same logical grant. Fixed by serializing the cross-org + * check + write behind a database-enforced lock (lib/services/distributedLock.js) keyed + * per idempotencyKey. These tests run against REAL MongoDB — the lock's unique-`_id` + * claim is the property under test, which a mocked model cannot demonstrate. + */ +describe('BillingExtraBalanceRepository.creditGrant — cross-org race closure integration tests:', () => { + let BillingExtraBalance; + let BillingExtraBalanceRepository; + let CronLock; + + const orgA = new mongoose.Types.ObjectId().toString(); + const orgB = new mongoose.Types.ObjectId().toString(); + + beforeAll(async () => { + await mongooseService.loadModels(); + await mongooseService.connect(); + + BillingExtraBalance = mongoose.model('BillingExtraBalance'); + BillingExtraBalanceRepository = (await import('../repositories/billing.extraBalance.repository.js')).default; + ({ CronLock } = await import('../../../lib/services/distributedLock.js')); + }); + + beforeEach(async () => { + await Promise.all([ + BillingExtraBalance.deleteMany({ organization: { $in: [orgA, orgB] } }), + CronLock.deleteMany({ _id: /^billing\.creditGrant\./ }), + ]); + }); + + afterAll(async () => { + await mongooseService.disconnect(); + }); + + test('concurrent creditGrant calls sharing one refId across TWO orgs → exactly ONE applied:true, ONE ledger entry total', async () => { + const sharedKey = `referral:${new mongoose.Types.ObjectId()}:referrer`; + + const [first, second] = await Promise.all([ + BillingExtraBalanceRepository.creditGrant(orgA, 1000, 'referral', { refId: sharedKey }), + BillingExtraBalanceRepository.creditGrant(orgB, 1000, 'referral', { refId: sharedKey }), + ]); + + // Exactly one org wins the grant; the loser surfaces the idempotent no-op — never both. + const applied = [first, second].filter((r) => r.applied === true); + expect(applied).toHaveLength(1); + const duplicate = [first, second].filter((r) => r.applied === false); + expect(duplicate).toHaveLength(1); + expect(duplicate[0].reason).toBe('duplicate_grant'); + + // The key was credited to exactly ONE of the two orgs, never both. + const [docA, docB] = await Promise.all([ + BillingExtraBalance.findOne({ organization: orgA }).lean(), + BillingExtraBalance.findOne({ organization: orgB }).lean(), + ]); + const entriesA = (docA?.ledger ?? []).filter((e) => e.refId === sharedKey); + const entriesB = (docB?.ledger ?? []).filter((e) => e.refId === sharedKey); + expect(entriesA.length + entriesB.length).toBe(1); + }, 15000); + + test('a losing call releases the lock — a later replay for the SAME org+key still resolves idempotently', async () => { + const sharedKey = `referral:${new mongoose.Types.ObjectId()}:referrer`; + + const [first] = await Promise.all([ + BillingExtraBalanceRepository.creditGrant(orgA, 1000, 'referral', { refId: sharedKey }), + BillingExtraBalanceRepository.creditGrant(orgB, 1000, 'referral', { refId: sharedKey }), + ]); + const winnerOrg = first.applied ? orgA : orgB; + + // Sequential replay against the org that actually won — must stay a no-op (lock is + // released after each call; this is the durable Step 0/2 dedup doing its job, not + // the lock still being held). + const replay = await BillingExtraBalanceRepository.creditGrant(winnerOrg, 1000, 'referral', { refId: sharedKey }); + expect(replay).toMatchObject({ applied: false, reason: 'duplicate_grant' }); + }, 15000); +}); diff --git a/modules/billing/tests/billing.extraBalance.unit.tests.js b/modules/billing/tests/billing.extraBalance.unit.tests.js index 5463b5ca1..dad1c0cff 100644 --- a/modules/billing/tests/billing.extraBalance.unit.tests.js +++ b/modules/billing/tests/billing.extraBalance.unit.tests.js @@ -146,6 +146,7 @@ describe('BillingExtraBalance unit tests:', () => { describe('Repository', () => { let BillingExtraBalanceRepository; let mockModel; + let mockLockModel; const orgId = '507f1f77bcf86cd799439011'; /** @@ -172,9 +173,29 @@ describe('BillingExtraBalance unit tests:', () => { exists: jest.fn(), }; + // #4155 — creditGrant dynamically imports lib/services/distributedLock.js, which + // registers its OWN 'CronLock' model. Keyed by name so that model does not share + // mockModel's findOneAndUpdate call queue with the ExtraBalance step1/step2 + // sequencing used throughout this describe block. Defaults to "lock always + // acquired" (echoes the holder it was given, matching acquireLock's success + // check) and "release always succeeds" — individual tests override + // mockLockModel.findOneAndUpdate to simulate lock contention. + mockLockModel = { + findOneAndUpdate: jest.fn((_filter, update) => Promise.resolve({ holder: update?.$set?.holder })), + deleteOne: jest.fn().mockResolvedValue({ deletedCount: 1 }), + }; + jest.unstable_mockModule('mongoose', () => ({ default: { - model: jest.fn(() => mockModel), + model: jest.fn((name) => (name === 'CronLock' ? mockLockModel : mockModel)), + models: {}, + // distributedLock.js does `new mongoose.Schema(...).index(...)` at import time — + // never actually executed by mongoose here, just needs to not throw. + Schema: class MockSchema { + index() { + return this; + } + }, Types: { ObjectId: { isValid: jest.fn(() => true), @@ -1321,6 +1342,7 @@ describe('Referral grant extensions (#3842):', () => { describe('repository', () => { let BillingExtraBalanceRepository; let mockModel; + let mockLockModel; /** * @param {Object[]} rows - Rows the aggregation should resolve with. @@ -1340,9 +1362,22 @@ describe('Referral grant extensions (#3842):', () => { exists: jest.fn(), aggregate: jest.fn(), }; + // #4155 — see the identical comment in the outer 'Repository' describe's beforeEach: + // creditGrant's dynamic import of distributedLock.js needs its own 'CronLock' model, + // kept off mockModel's findOneAndUpdate call queue. + mockLockModel = { + findOneAndUpdate: jest.fn((_filter, update) => Promise.resolve({ holder: update?.$set?.holder })), + deleteOne: jest.fn().mockResolvedValue({ deletedCount: 1 }), + }; jest.unstable_mockModule('mongoose', () => ({ default: { - model: jest.fn(() => mockModel), + model: jest.fn((name) => (name === 'CronLock' ? mockLockModel : mockModel)), + models: {}, + Schema: class MockSchema { + index() { + return this; + } + }, Types: { ObjectId: { isValid: jest.fn(() => true) } }, }, })); From 45f28aa7c293686f15b4e3c4b567640c00bf3186 Mon Sep 17 00:00:00 2001 From: Pierre Brisorgueil Date: Tue, 29 Sep 2026 09:45:12 +0200 Subject: [PATCH 12/15] fix(billing): creditGrant claims its idempotency key on a durable unique index, not a lease MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The previous fix (70c97abf) serialized the cross-org check + write behind a TTL lock from lib/services/distributedLock.js. CodeRabbit's follow-up review was right that a lease does not close the race, only narrow it: if the first grant's write ever outran the TTL, a second call could acquire the expired lock and pass the same cross-org check, crediting a different org. Replaces the lock with a genuinely durable claim: BillingGrantClaim is a brand-new collection (billing.grantClaim.model.mongoose.js) with a unique index on `key` — new collection, so the index builds against no pre-existing data, no migration needed, and MongoDB enforces it permanently. BillingGrantClaimRepository.tryClaim (mirrors the existing ProcessedStripeEvent claim-by-unique-index pattern) inserts a claim before crediting; a conflict from a DIFFERENT org is the cross-org double-grant this exists to prevent, rejected immediately. A conflict from the SAME org (a replay, or this exact call retrying after a crash between claiming and writing) falls through to the existing per-org ledger guard, which is already idempotent on its own. The claim is never rolled back or expired — nothing to release, nothing to expire into a race. The original Step 0 cross-org ledger existence check stays as a legacy backstop for entries written before this claim mechanism existed. Addresses CodeRabbit's follow-up on the same #4155 review comment. Claude-Session: https://claude.ai/code/session_01EqkUXh5nmxbVBTvhgM6zAo --- .../billing.grantClaim.model.mongoose.js | 65 ++++++++++ .../billing.extraBalance.repository.js | 121 ++++++++---------- .../billing.grantClaim.repository.js | 43 +++++++ ...ance.creditGrant.race.integration.tests.js | 39 ++++-- .../tests/billing.extraBalance.unit.tests.js | 116 ++++++++++++----- ...illing.grantClaim.repository.unit.tests.js | 92 +++++++++++++ 6 files changed, 364 insertions(+), 112 deletions(-) create mode 100644 modules/billing/models/billing.grantClaim.model.mongoose.js create mode 100644 modules/billing/repositories/billing.grantClaim.repository.js create mode 100644 modules/billing/tests/billing.grantClaim.repository.unit.tests.js diff --git a/modules/billing/models/billing.grantClaim.model.mongoose.js b/modules/billing/models/billing.grantClaim.model.mongoose.js new file mode 100644 index 000000000..b4ec157bc --- /dev/null +++ b/modules/billing/models/billing.grantClaim.model.mongoose.js @@ -0,0 +1,65 @@ +/** + * Module dependencies + */ +import mongoose from 'mongoose'; + +const Schema = mongoose.Schema; + +/** + * BillingGrantClaim Data Model Mongoose + * + * Durable, database-enforced claim on a BillingExtraBalanceRepository.creditGrant + * idempotencyKey (#4155). The ledger (`modules/billing/models/billing.extraBalance.model.mongoose.js`) + * is an embedded array per organization — a unique index on it cannot enforce uniqueness + * ACROSS organizations without a migration against already-deployed data. This is a + * separate, brand-new collection instead: the unique index on `key` is built fresh (no + * pre-existing data, no migration), and MongoDB enforces it permanently — unlike a + * TTL/lease-based lock, there is no window where an expired holder can still write. + * + * The claim is never deleted or expired. `organization` records who claimed it, so a + * later call with the SAME key can tell "my own retry/replay after a crash" (same org — + * fall through to the ordinary per-org ledger guard, which is idempotent on its own) from + * "a different org already holds this key" (the cross-org double-grant this exists to + * prevent — rejected as a duplicate). + */ +const GrantClaimMongoose = new Schema( + { + key: { + type: String, + required: true, + unique: true, + trim: true, + }, + organization: { + type: Schema.ObjectId, + ref: 'Organization', + required: true, + }, + at: { + type: Date, + required: true, + default: () => new Date(), + }, + }, + { + timestamps: false, + }, +); + +/** + * Returns the hex string representation of the document ObjectId. + * @returns {string} Hex string of the ObjectId. + */ +function addID() { + return this._id.toHexString(); +} + +/** + * Model configuration + */ +GrantClaimMongoose.virtual('id').get(addID); +GrantClaimMongoose.set('toJSON', { + virtuals: true, +}); + +mongoose.model('BillingGrantClaim', GrantClaimMongoose); diff --git a/modules/billing/repositories/billing.extraBalance.repository.js b/modules/billing/repositories/billing.extraBalance.repository.js index 329576d21..8e9b18553 100644 --- a/modules/billing/repositories/billing.extraBalance.repository.js +++ b/modules/billing/repositories/billing.extraBalance.repository.js @@ -1,9 +1,9 @@ /** * Module dependencies */ -import { randomUUID } from 'node:crypto'; import mongoose from 'mongoose'; import AppError from '../../../lib/helpers/AppError.js'; +import BillingGrantClaimRepository from './billing.grantClaim.repository.js'; import BillingExtraBalanceSchema from '../models/billing.extraBalance.schema.js'; import { computeExpiryRemovals } from '../lib/billing.packExpiry.js'; @@ -165,15 +165,6 @@ const debit = async (orgId, amount, refId) => { return { doc: null, applied: false, reason: 'duplicate_step' }; }; -/** - * Lock TTL for the creditGrant critical section (existence check + getOrCreate + guarded - * write). Generous relative to the operation's normal sub-100ms cost so a slow write never - * loses its own lock mid-flight; short enough that a crashed holder does not block a retry - * for long — the caller's own retry (Stripe redelivery) or the referral reconcile cron - * picks it back up well within this window. - */ -const GRANT_LOCK_TTL_MS = 15_000; - /** * @function creditGrant * @description Atomically credit extra meter units for a non-Stripe grant (e.g. signup free @@ -187,31 +178,38 @@ const GRANT_LOCK_TTL_MS = 15_000; * sharing the same idempotencyKey but resolving to DIFFERENT orgs (e.g. the * in-process referral listener racing the reconcile cron backfill for the same * invitation) can both observe "not yet granted" before either writes, and both - * would credit their own org (#4155, CodeRabbit). This is closed with a - * database-enforced claim: `lib/services/distributedLock.js` upserts a lock - * document on a unique `_id`, so only ONE concurrent call can hold the key at a - * time — the same primitive the billing crons already use for pod-level mutual - * exclusion, here scoped per-idempotencyKey instead of per-cron. + * would credit their own org. A TTL/lease-based lock does not close this either + * — it only narrows the window to "the write took longer than the lease" (#4155, + * CodeRabbit, both rounds). Closed instead with a genuinely durable claim: + * `BillingGrantClaimRepository.tryClaim` inserts into a brand-new collection + * with a unique index on `key` — new collection, so the index is built fresh + * with no pre-existing data to migrate, and MongoDB enforces it permanently. + * There is no lease to expire and nothing to release. * - * Sequence while the lock is held: + * Sequence: + * Step -1 — claim the key (BillingGrantClaimRepository.tryClaim). A claim + * held by a DIFFERENT org is the cross-org conflict this exists to catch + * → duplicate_grant immediately. A claim held by THIS SAME org (a replay, + * or this exact call retrying after a crash between claiming and writing) + * falls through — the per-org ledger guard in Step 2 is idempotent on its + * own and decides the outcome from there. The claim itself is never rolled + * back: it is a permanent "this key belongs to this org" record, not a + * lock — a same-org retry after ANY failure (crash, transient DB error) + * re-enters here and proceeds exactly like a fresh call whose claim already + * exists. * Step 0 — cross-org existence check (any org's ledger, not just this one). + * Legacy backstop: entries written before this claim mechanism existed have + * no claim doc, so a legacy key retried from a reassigned org would pass + * Step -1 as a "fresh" claim without this check. * Step 1 — ensure doc exists (atomic getOrCreate, no-op on replay). * Step 2 — idempotency-guarded credit (no upsert). - * The lock only serializes this sequence for a given key; the org-scoped - * `'ledger.refId': { $ne: key }` guard in Step 2 remains the durable dedup once - * the lock is released (a later replay of the SAME org+key still resolves - * correctly without needing the lock to still be held). - * Lock not acquired (another call for this exact key is mid-flight) → treated - * like any other idempotent replay: `{ applied: false, reason: 'duplicate_grant' }`. - * A request that loses this race is not lost — referral grants are re-derived - * and retried by crons/billing.referralReconcile.js; the signup-grant synthetic - * key (`-`) embeds orgId and so can never collide across orgs by - * construction, making this path relevant to the referral (explicit refId) case. * * Idempotency key: `options.refId` when supplied (#3842 referral grants — * several grants per org, one per invitation, e.g. * `referral::referrer`), otherwise the synthetic per-org key - * `-` (signup grant — one per org). + * `-` (signup grant — one per org; embeds orgId, so it can + * never collide across orgs by construction — the cross-org path is only ever + * reachable via an explicit `refId`, i.e. referral grants). * `options.expiresAt` mirrors the creditPack expiry mechanism: the entry is * swept by crons/billing.extrasExpiration.js once past its expiry. * No stripeSessionId required. @@ -248,46 +246,39 @@ const creditGrant = async (orgId, amount, source, { refId = null, expiresAt = nu ...(expiresAt ? { expiresAt } : {}), }; - // Database-enforced claim on the idempotencyKey (see JSDoc): serializes the cross-org - // check + write below against any other creditGrant call for this exact key. Dynamic - // import — this is the only function in the repository that needs the lock's own - // CronLock model, so no other test in this file has to account for a second - // mongoose.model() registration it never exercises. - const { acquireLock, releaseLock } = await import('../../../lib/services/distributedLock.js'); - const lockName = `billing.creditGrant.${idempotencyKey}`; - const lockHolder = `${process.env.HOSTNAME ?? 'unknown'}:${randomUUID()}`; - const locked = await acquireLock({ name: lockName, ttlMs: GRANT_LOCK_TTL_MS, holder: lockHolder }); - if (!locked) return { doc: null, applied: false, reason: 'duplicate_grant' }; - - try { - // Step 0: cross-org idempotency guard. The step-2 filter below only excludes a - // matching refId within THIS org's own ledger, so it misses a replay whose target - // org changed between attempts. Check across every org's ledger before writing. - const alreadyGranted = await BillingExtraBalance().exists({ 'ledger.refId': idempotencyKey }); - if (alreadyGranted) return { doc: null, applied: false, reason: 'duplicate_grant' }; - - // Step 1: ensure the document exists (atomic getOrCreate, no-op if already present). - await getOrCreate(orgId); - - // Step 2: idempotency-guarded credit (no upsert — doc is guaranteed to exist after step 1). - const doc = await BillingExtraBalance().findOneAndUpdate( - { - organization: orgId, - 'ledger.refId': { $ne: idempotencyKey }, - }, - { - $push: { ledger: entry }, - $inc: { cachedBalance: amount }, - $set: { cachedBalanceAt: new Date() }, - }, - { returnDocument: 'after' }, - ); - - if (doc) return { doc, applied: true }; + // Step -1: durable cross-org claim (see JSDoc). A different org already owns this key → + // reject now. This org already owns it (fresh or a replay) → fall through to the + // existing per-org guards, which decide the outcome idempotently on their own. + const claim = await BillingGrantClaimRepository.tryClaim(idempotencyKey, orgId); + if (!claim.claimed && claim.ownerOrgId !== orgId) { return { doc: null, applied: false, reason: 'duplicate_grant' }; - } finally { - await releaseLock({ name: lockName, holder: lockHolder }); } + + // Step 0: legacy backstop — cross-org existence check for ledger entries written before + // the claim mechanism existed (no claim doc for those). The step-2 filter below only + // excludes a matching refId within THIS org's own ledger, so it misses those. + const alreadyGranted = await BillingExtraBalance().exists({ 'ledger.refId': idempotencyKey }); + if (alreadyGranted) return { doc: null, applied: false, reason: 'duplicate_grant' }; + + // Step 1: ensure the document exists (atomic getOrCreate, no-op if already present). + await getOrCreate(orgId); + + // Step 2: idempotency-guarded credit (no upsert — doc is guaranteed to exist after step 1). + const doc = await BillingExtraBalance().findOneAndUpdate( + { + organization: orgId, + 'ledger.refId': { $ne: idempotencyKey }, + }, + { + $push: { ledger: entry }, + $inc: { cachedBalance: amount }, + $set: { cachedBalanceAt: new Date() }, + }, + { returnDocument: 'after' }, + ); + + if (doc) return { doc, applied: true }; + return { doc: null, applied: false, reason: 'duplicate_grant' }; }; /** diff --git a/modules/billing/repositories/billing.grantClaim.repository.js b/modules/billing/repositories/billing.grantClaim.repository.js new file mode 100644 index 000000000..d457bd834 --- /dev/null +++ b/modules/billing/repositories/billing.grantClaim.repository.js @@ -0,0 +1,43 @@ +/** + * Module dependencies + */ +import mongoose from 'mongoose'; +import { isDuplicateKeyError } from '../lib/billing.errors.js'; + +/** + * Lazily resolves the BillingGrantClaim Mongoose model. + * Deferred to keep unit tests importable before model registration. + * @returns {import('mongoose').Model} The registered BillingGrantClaim model. + */ +// biome-ignore lint/correctness/useQwikValidLexicalScope: false positive — Node.js repository, not Qwik +const BillingGrantClaim = () => mongoose.model('BillingGrantClaim'); + +/** + * @function tryClaim + * @description Atomically claim a creditGrant idempotencyKey using the unique index on + * `key` (see models/billing.grantClaim.model.mongoose.js for why this is a + * separate collection rather than a unique index on the ledger). No rollback, + * no TTL — the claim is permanent once inserted. + * @param {string} key - The creditGrant idempotencyKey to claim. + * @param {string} orgId - The organization ObjectId (string) attempting the claim. + * @returns {Promise<{claimed: boolean, ownerOrgId?: string|null}>} `claimed: true` on a + * fresh claim. `claimed: false` with `ownerOrgId` set to the claim holder's + * org (string) when the key is already claimed — by this same org (a + * same-org retry/replay) or a different one (the cross-org conflict #4155 + * exists to catch). + */ +// biome-ignore lint/correctness/useQwikValidLexicalScope: false positive — Node.js repository, not Qwik +const tryClaim = async (key, orgId) => { + try { + await BillingGrantClaim().create({ key, organization: orgId, at: new Date() }); + return { claimed: true }; + } catch (err) { + if (!isDuplicateKeyError(err)) throw err; + const existing = await BillingGrantClaim().findOne({ key }).lean(); + return { claimed: false, ownerOrgId: existing ? String(existing.organization) : null }; + } +}; + +export default { + tryClaim, +}; diff --git a/modules/billing/tests/billing.extraBalance.creditGrant.race.integration.tests.js b/modules/billing/tests/billing.extraBalance.creditGrant.race.integration.tests.js index 5d34e25ca..fd63c0ed6 100644 --- a/modules/billing/tests/billing.extraBalance.creditGrant.race.integration.tests.js +++ b/modules/billing/tests/billing.extraBalance.creditGrant.race.integration.tests.js @@ -13,15 +13,17 @@ import mongooseService from '../../../lib/services/mongoose.js'; * sharing the same idempotencyKey but targeting DIFFERENT orgs (e.g. the in-process * referral listener racing the reconcile cron backfill for the same invitation) could * both observe "not yet granted" before either writes, and both would credit their own - * org — double-crediting the same logical grant. Fixed by serializing the cross-org - * check + write behind a database-enforced lock (lib/services/distributedLock.js) keyed - * per idempotencyKey. These tests run against REAL MongoDB — the lock's unique-`_id` - * claim is the property under test, which a mocked model cannot demonstrate. + * org — double-crediting the same logical grant. A TTL/lease-based lock does not close + * this either — it only narrows the window to "the write took longer than the lease". + * Fixed with a genuinely durable claim (BillingGrantClaimRepository.tryClaim): a brand-new + * collection, unique index on `key`, no TTL, nothing to release. These tests run against + * REAL MongoDB — the unique-index claim is the property under test, which a mocked model + * cannot demonstrate. */ describe('BillingExtraBalanceRepository.creditGrant — cross-org race closure integration tests:', () => { let BillingExtraBalance; + let BillingGrantClaim; let BillingExtraBalanceRepository; - let CronLock; const orgA = new mongoose.Types.ObjectId().toString(); const orgB = new mongoose.Types.ObjectId().toString(); @@ -31,14 +33,17 @@ describe('BillingExtraBalanceRepository.creditGrant — cross-org race closure i await mongooseService.connect(); BillingExtraBalance = mongoose.model('BillingExtraBalance'); + BillingGrantClaim = mongoose.model('BillingGrantClaim'); + // The unique index on `key` is what makes the claim atomic — build it before the first + // concurrent test runs, otherwise the first run can race the (async) index build. + await BillingGrantClaim.syncIndexes(); BillingExtraBalanceRepository = (await import('../repositories/billing.extraBalance.repository.js')).default; - ({ CronLock } = await import('../../../lib/services/distributedLock.js')); }); beforeEach(async () => { await Promise.all([ BillingExtraBalance.deleteMany({ organization: { $in: [orgA, orgB] } }), - CronLock.deleteMany({ _id: /^billing\.creditGrant\./ }), + BillingGrantClaim.deleteMany({ organization: { $in: [orgA, orgB] } }), ]); }); @@ -46,7 +51,7 @@ describe('BillingExtraBalanceRepository.creditGrant — cross-org race closure i await mongooseService.disconnect(); }); - test('concurrent creditGrant calls sharing one refId across TWO orgs → exactly ONE applied:true, ONE ledger entry total', async () => { + test('concurrent creditGrant calls sharing one refId across TWO orgs → exactly ONE applied:true, ONE ledger entry total, ONE claim', async () => { const sharedKey = `referral:${new mongoose.Types.ObjectId()}:referrer`; const [first, second] = await Promise.all([ @@ -69,9 +74,15 @@ describe('BillingExtraBalanceRepository.creditGrant — cross-org race closure i const entriesA = (docA?.ledger ?? []).filter((e) => e.refId === sharedKey); const entriesB = (docB?.ledger ?? []).filter((e) => e.refId === sharedKey); expect(entriesA.length + entriesB.length).toBe(1); + + // Exactly one claim document for the key, owned by whichever org actually won. + const claims = await BillingGrantClaim.find({ key: sharedKey }).lean(); + expect(claims).toHaveLength(1); + const winnerOrg = entriesA.length === 1 ? orgA : orgB; + expect(String(claims[0].organization)).toBe(winnerOrg); }, 15000); - test('a losing call releases the lock — a later replay for the SAME org+key still resolves idempotently', async () => { + test('a losing call leaves the claim permanently owned by the winner — a later replay for the SAME org+key still resolves idempotently', async () => { const sharedKey = `referral:${new mongoose.Types.ObjectId()}:referrer`; const [first] = await Promise.all([ @@ -80,10 +91,14 @@ describe('BillingExtraBalanceRepository.creditGrant — cross-org race closure i ]); const winnerOrg = first.applied ? orgA : orgB; - // Sequential replay against the org that actually won — must stay a no-op (lock is - // released after each call; this is the durable Step 0/2 dedup doing its job, not - // the lock still being held). + // Sequential replay against the org that actually won — must stay a no-op. No lock to + // release: the claim is permanent, the per-org ledger guard is the durable dedup. const replay = await BillingExtraBalanceRepository.creditGrant(winnerOrg, 1000, 'referral', { refId: sharedKey }); expect(replay).toMatchObject({ applied: false, reason: 'duplicate_grant' }); + + // A retry from the LOSING org must still be rejected — the claim does not expire. + const loserOrg = winnerOrg === orgA ? orgB : orgA; + const loserRetry = await BillingExtraBalanceRepository.creditGrant(loserOrg, 1000, 'referral', { refId: sharedKey }); + expect(loserRetry).toMatchObject({ applied: false, reason: 'duplicate_grant' }); }, 15000); }); diff --git a/modules/billing/tests/billing.extraBalance.unit.tests.js b/modules/billing/tests/billing.extraBalance.unit.tests.js index dad1c0cff..8223743ba 100644 --- a/modules/billing/tests/billing.extraBalance.unit.tests.js +++ b/modules/billing/tests/billing.extraBalance.unit.tests.js @@ -146,7 +146,7 @@ describe('BillingExtraBalance unit tests:', () => { describe('Repository', () => { let BillingExtraBalanceRepository; let mockModel; - let mockLockModel; + let mockClaimModel; const orgId = '507f1f77bcf86cd799439011'; /** @@ -173,29 +173,20 @@ describe('BillingExtraBalance unit tests:', () => { exists: jest.fn(), }; - // #4155 — creditGrant dynamically imports lib/services/distributedLock.js, which - // registers its OWN 'CronLock' model. Keyed by name so that model does not share - // mockModel's findOneAndUpdate call queue with the ExtraBalance step1/step2 - // sequencing used throughout this describe block. Defaults to "lock always - // acquired" (echoes the holder it was given, matching acquireLock's success - // check) and "release always succeeds" — individual tests override - // mockLockModel.findOneAndUpdate to simulate lock contention. - mockLockModel = { - findOneAndUpdate: jest.fn((_filter, update) => Promise.resolve({ holder: update?.$set?.holder })), - deleteOne: jest.fn().mockResolvedValue({ deletedCount: 1 }), + // #4155 — creditGrant claims its idempotencyKey via BillingGrantClaimRepository, + // which resolves its OWN 'BillingGrantClaim' model. Keyed by name so that model does + // not share mockModel's findOneAndUpdate call queue with the ExtraBalance + // step1/step2 sequencing used throughout this describe block. Defaults to "claim + // always succeeds" (fresh key) — individual tests override mockClaimModel.create to + // simulate an existing claim (same-org replay or a different org's conflict). + mockClaimModel = { + create: jest.fn().mockResolvedValue({}), + findOne: jest.fn(() => ({ lean: jest.fn().mockResolvedValue(null) })), }; jest.unstable_mockModule('mongoose', () => ({ default: { - model: jest.fn((name) => (name === 'CronLock' ? mockLockModel : mockModel)), - models: {}, - // distributedLock.js does `new mongoose.Schema(...).index(...)` at import time — - // never actually executed by mongoose here, just needs to not throw. - Schema: class MockSchema { - index() { - return this; - } - }, + model: jest.fn((name) => (name === 'BillingGrantClaim' ? mockClaimModel : mockModel)), Types: { ObjectId: { isValid: jest.fn(() => true), @@ -1209,8 +1200,12 @@ describe('BillingExtraBalance unit tests:', () => { // #4151 — the idempotency guard in step 2 below only excludes a matching refId // WITHIN this org's own ledger ('ledger.refId': { $ne }), so a retry whose target // org changed between attempts (e.g. an org merge/reassignment) is credited twice. - test('returns applied:false duplicate_grant when the refId was already granted to a DIFFERENT org', async () => { - // Cross-org existence check finds a prior grant under this refId in any org. + // Legacy backstop (step 0, after the #4155 claim below): a ledger entry written + // before the claim mechanism existed has no claim doc, so the claim itself looks + // "fresh" — this cross-org existence check is what still catches it. + test('returns applied:false duplicate_grant when the refId already exists in a DIFFERENT org\'s ledger (legacy, pre-claim data)', async () => { + // Claim succeeds (fresh key, default mock) — no claim doc for this pre-existing entry. + // Cross-org existence check then finds the prior grant under this refId in any org. mockModel.exists.mockResolvedValue({ _id: 'some-other-orgs-balance-doc' }); const result = await BillingExtraBalanceRepository.creditGrant(orgId, 500, 'referral', { @@ -1222,6 +1217,63 @@ describe('BillingExtraBalance unit tests:', () => { // Must short-circuit before any write. expect(mockModel.findOneAndUpdate).not.toHaveBeenCalled(); }); + + // #4155 — the durable claim itself: two concurrent creditGrant calls sharing one + // idempotencyKey but targeting DIFFERENT orgs must not both credit. The claim's + // unique index on `key` is what makes this atomic (not a lease/lock), so the second + // caller's create() rejects with a duplicate-key error before either org is touched. + describe('creditGrant — durable claim (#4155):', () => { + test('a claim already held by a DIFFERENT org → duplicate_grant, before any ledger read or write', async () => { + const otherOrgId = '507f1f77bcf86cd799439022'; + mockClaimModel.create.mockRejectedValue(Object.assign(new Error('E11000 duplicate key'), { code: 11000 })); + mockClaimModel.findOne.mockReturnValue({ lean: jest.fn().mockResolvedValue({ organization: otherOrgId }) }); + + const result = await BillingExtraBalanceRepository.creditGrant(orgId, 500, 'referral', { + refId: 'referral:64b2f0000000000000000002:referrer', + }); + + expect(result).toEqual({ doc: null, applied: false, reason: 'duplicate_grant' }); + expect(mockClaimModel.create).toHaveBeenCalledWith( + expect.objectContaining({ key: 'referral:64b2f0000000000000000002:referrer', organization: orgId }), + ); + // Rejected at the claim — never reaches the ledger at all. + expect(mockModel.exists).not.toHaveBeenCalled(); + expect(mockModel.findOneAndUpdate).not.toHaveBeenCalled(); + }); + + test('a claim already held by the SAME org (replay, or retry after a crash) falls through to the ordinary per-org guard', async () => { + mockClaimModel.create.mockRejectedValue(Object.assign(new Error('E11000 duplicate key'), { code: 11000 })); + mockClaimModel.findOne.mockReturnValue({ lean: jest.fn().mockResolvedValue({ organization: orgId }) }); + // Step 1: getOrCreate; Step 2: per-org guard says already credited. + mockModel.findOneAndUpdate + .mockResolvedValueOnce(makeDoc()) + .mockResolvedValueOnce(null); + + const result = await BillingExtraBalanceRepository.creditGrant(orgId, 500, 'referral', { + refId: 'referral:64b2f0000000000000000003:referrer', + }); + + expect(result).toEqual({ doc: null, applied: false, reason: 'duplicate_grant' }); + // Did NOT short-circuit — the same-org branch reaches the ordinary guards. + expect(mockModel.findOneAndUpdate).toHaveBeenCalledTimes(2); + }); + + test('a fresh claim proceeds normally and credits the org', async () => { + const updatedDoc = makeDoc({ cachedBalance: 500 }); + mockModel.findOneAndUpdate + .mockResolvedValueOnce(makeDoc()) + .mockResolvedValueOnce(updatedDoc); + + const result = await BillingExtraBalanceRepository.creditGrant(orgId, 500, 'referral', { + refId: 'referral:64b2f0000000000000000004:referrer', + }); + + expect(result.applied).toBe(true); + expect(mockClaimModel.create).toHaveBeenCalledWith( + expect.objectContaining({ key: 'referral:64b2f0000000000000000004:referrer', organization: orgId }), + ); + }); + }); }); }); }); @@ -1342,7 +1394,7 @@ describe('Referral grant extensions (#3842):', () => { describe('repository', () => { let BillingExtraBalanceRepository; let mockModel; - let mockLockModel; + let mockClaimModel; /** * @param {Object[]} rows - Rows the aggregation should resolve with. @@ -1363,21 +1415,15 @@ describe('Referral grant extensions (#3842):', () => { aggregate: jest.fn(), }; // #4155 — see the identical comment in the outer 'Repository' describe's beforeEach: - // creditGrant's dynamic import of distributedLock.js needs its own 'CronLock' model, - // kept off mockModel's findOneAndUpdate call queue. - mockLockModel = { - findOneAndUpdate: jest.fn((_filter, update) => Promise.resolve({ holder: update?.$set?.holder })), - deleteOne: jest.fn().mockResolvedValue({ deletedCount: 1 }), + // creditGrant claims its idempotencyKey via its own 'BillingGrantClaim' model, kept + // off mockModel's findOneAndUpdate call queue. + mockClaimModel = { + create: jest.fn().mockResolvedValue({}), + findOne: jest.fn(() => ({ lean: jest.fn().mockResolvedValue(null) })), }; jest.unstable_mockModule('mongoose', () => ({ default: { - model: jest.fn((name) => (name === 'CronLock' ? mockLockModel : mockModel)), - models: {}, - Schema: class MockSchema { - index() { - return this; - } - }, + model: jest.fn((name) => (name === 'BillingGrantClaim' ? mockClaimModel : mockModel)), Types: { ObjectId: { isValid: jest.fn(() => true) } }, }, })); diff --git a/modules/billing/tests/billing.grantClaim.repository.unit.tests.js b/modules/billing/tests/billing.grantClaim.repository.unit.tests.js new file mode 100644 index 000000000..ab6509ca1 --- /dev/null +++ b/modules/billing/tests/billing.grantClaim.repository.unit.tests.js @@ -0,0 +1,92 @@ +/** + * Module dependencies. + */ +import { jest, describe, test, beforeEach, afterEach, expect } from '@jest/globals'; + +/** + * Unit tests for billing.grantClaim.repository.js (#4155). + */ +describe('BillingGrantClaimRepository unit tests:', () => { + let BillingGrantClaimRepository; + let mockModel; + + const orgId = '507f1f77bcf86cd799439011'; + const otherOrgId = '507f1f77bcf86cd799439022'; + + const makeE11000 = () => { + const err = new Error('E11000 duplicate key error'); + err.code = 11000; + return err; + }; + + beforeEach(async () => { + jest.resetModules(); + + mockModel = { + create: jest.fn(), + findOne: jest.fn(), + }; + + jest.unstable_mockModule('mongoose', () => ({ + default: { + model: jest.fn(() => mockModel), + }, + })); + + const mod = await import('../repositories/billing.grantClaim.repository.js'); + BillingGrantClaimRepository = mod.default; + }); + + afterEach(() => { + jest.restoreAllMocks(); + }); + + describe('tryClaim', () => { + test('fresh key → { claimed: true }, no lookup needed', async () => { + mockModel.create.mockResolvedValue({ key: 'k1', organization: orgId }); + + const result = await BillingGrantClaimRepository.tryClaim('k1', orgId); + + expect(result).toEqual({ claimed: true }); + expect(mockModel.create).toHaveBeenCalledWith( + expect.objectContaining({ key: 'k1', organization: orgId }), + ); + expect(mockModel.findOne).not.toHaveBeenCalled(); + }); + + test('duplicate key already owned by the SAME org → { claimed: false, ownerOrgId: orgId }', async () => { + mockModel.create.mockRejectedValueOnce(makeE11000()); + mockModel.findOne.mockReturnValue({ lean: jest.fn().mockResolvedValue({ key: 'k1', organization: orgId }) }); + + const result = await BillingGrantClaimRepository.tryClaim('k1', orgId); + + expect(result).toEqual({ claimed: false, ownerOrgId: orgId }); + expect(mockModel.findOne).toHaveBeenCalledWith({ key: 'k1' }); + }); + + test('duplicate key already owned by a DIFFERENT org → { claimed: false, ownerOrgId: otherOrgId }', async () => { + mockModel.create.mockRejectedValueOnce(makeE11000()); + mockModel.findOne.mockReturnValue({ lean: jest.fn().mockResolvedValue({ key: 'k1', organization: otherOrgId }) }); + + const result = await BillingGrantClaimRepository.tryClaim('k1', orgId); + + expect(result).toEqual({ claimed: false, ownerOrgId: otherOrgId }); + }); + + test('duplicate key but lookup finds nothing (race window: claim deleted between insert-fail and lookup) → ownerOrgId: null', async () => { + mockModel.create.mockRejectedValueOnce(makeE11000()); + mockModel.findOne.mockReturnValue({ lean: jest.fn().mockResolvedValue(null) }); + + const result = await BillingGrantClaimRepository.tryClaim('k1', orgId); + + expect(result).toEqual({ claimed: false, ownerOrgId: null }); + }); + + test('a non-duplicate-key error propagates (never swallowed as a claim conflict)', async () => { + mockModel.create.mockRejectedValueOnce(new Error('connection reset')); + + await expect(BillingGrantClaimRepository.tryClaim('k1', orgId)).rejects.toThrow('connection reset'); + expect(mockModel.findOne).not.toHaveBeenCalled(); + }); + }); +}); From 7a9d294df9a8accb195af3aa46e355ded6336615 Mon Sep 17 00:00:00 2001 From: Pierre Brisorgueil Date: Tue, 29 Sep 2026 10:45:14 +0200 Subject: [PATCH 13/15] fix(billing): debit does not retry deterministic repository failures retryWithBackoff's default shouldRetry retries every thrown error. BillingExtraBalanceRepository.debit's own validation errors (invalid argument: bad amount/refId) and its ORGANIZATION_NOT_FOUND AppError never succeed on retry, so the default wastes 2 extra DB round trips (~600ms) before surfacing the identical error to the caller. Passes a shouldRetry predicate that excludes those two cases; transient failures (the case retryWithBackoff exists for, #4151) still retry as before. Addresses a CodeRabbit finding from the #4155 review's full-review pass. Claude-Session: https://claude.ai/code/session_01EqkUXh5nmxbVBTvhgM6zAo --- .../billing/services/billing.extra.service.js | 11 ++++++++++- .../tests/billing.extra.service.unit.tests.js | 19 +++++++++++++++++++ 2 files changed, 29 insertions(+), 1 deletion(-) diff --git a/modules/billing/services/billing.extra.service.js b/modules/billing/services/billing.extra.service.js index d1c1714b4..ec50b1933 100644 --- a/modules/billing/services/billing.extra.service.js +++ b/modules/billing/services/billing.extra.service.js @@ -44,6 +44,11 @@ const creditPack = async (orgId, packId, stripeSessionId) => { * Wrapped in retryWithBackoff: the write is idempotent by refId (a retry * cannot double-charge), so a transient DB error must not fall through and * silently make the overflow free (#4151). + * `shouldRetry` excludes the repository's deterministic failures — an + * `invalid argument: ...` (bad amount/refId) or the `ORGANIZATION_NOT_FOUND` + * AppError never succeeds on retry, so retrying it only burns the backoff + * budget (2 extra DB round trips, ~600ms) before surfacing the identical + * error. * * @param {string} orgId - The organization ObjectId (string). * @param {number} units - Meter units to debit (must be > 0). @@ -52,7 +57,11 @@ const creditPack = async (orgId, packId, stripeSessionId) => { */ // biome-ignore lint/correctness/useQwikValidLexicalScope: false positive — Node.js service, not Qwik const debit = (orgId, units, refId) => - retryWithBackoff(() => BillingExtraBalanceRepository.debit(orgId, units, refId)); + retryWithBackoff(() => BillingExtraBalanceRepository.debit(orgId, units, refId), { + shouldRetry: (err) => + !(typeof err?.message === 'string' && err.message.startsWith('invalid argument')) + && err?.code !== 'ORGANIZATION_NOT_FOUND', + }); /** * @function getOrgBalanceContext diff --git a/modules/billing/tests/billing.extra.service.unit.tests.js b/modules/billing/tests/billing.extra.service.unit.tests.js index e294aaadd..e207de877 100644 --- a/modules/billing/tests/billing.extra.service.unit.tests.js +++ b/modules/billing/tests/billing.extra.service.unit.tests.js @@ -166,6 +166,25 @@ describe('BillingExtraService unit tests:', () => { await expect(BillingExtraService.debit(orgId, 100000, 'ref_down')).rejects.toThrow('Mongo topology closed'); }); + + // CodeRabbit follow-up on #4155's PR: retryWithBackoff's default shouldRetry retries + // every error, including the repository's own deterministic validation/existence + // failures — those never succeed on retry, so retrying them only wastes 2 DB round + // trips (~600ms) before surfacing the identical error. + test('does NOT retry an "invalid argument" error from the repository — fails fast on attempt 1', async () => { + mockRepository.debit.mockRejectedValue(new Error('invalid argument: amount must be a positive finite number')); + + await expect(BillingExtraService.debit(orgId, -1, 'ref_bad_amount')).rejects.toThrow('invalid argument'); + expect(mockRepository.debit).toHaveBeenCalledTimes(1); + }); + + test('does NOT retry an ORGANIZATION_NOT_FOUND AppError from the repository — fails fast on attempt 1', async () => { + const err = Object.assign(new Error('Organization not found: ghost-org'), { code: 'ORGANIZATION_NOT_FOUND', status: 404 }); + mockRepository.debit.mockRejectedValue(err); + + await expect(BillingExtraService.debit('ghost-org', 100, 'ref_ghost')).rejects.toThrow('Organization not found'); + expect(mockRepository.debit).toHaveBeenCalledTimes(1); + }); }); describe('getOrgBalanceContext', () => { From 415c17f5601845d37e8d2dfa5f156942ec616767 Mon Sep 17 00:00:00 2001 From: Pierre Brisorgueil Date: Tue, 29 Sep 2026 11:13:23 +0200 Subject: [PATCH 14/15] fix(billing): checkout.session.completed throws instead of returning when Stripe is unconfigured The retrieval-catch a few lines below already documents why a silent return is wrong here: withIdempotency records the event as processed even though nothing happened, so Stripe never redelivers it once Stripe configuration is restored. The `!stripe` guard above it had the identical failure mode but still returned. Throws now, matching the retrieval-catch's own pattern. Addresses an outside-diff-range CodeRabbit finding from the #4155 review's full-review pass on billing.webhook.service.js:253-257. Claude-Session: https://claude.ai/code/session_01EqkUXh5nmxbVBTvhgM6zAo --- .../services/billing.webhook.service.js | 9 ++++++-- .../billing.webhook.checkout.unit.tests.js | 23 +++++++++++-------- 2 files changed, 21 insertions(+), 11 deletions(-) diff --git a/modules/billing/services/billing.webhook.service.js b/modules/billing/services/billing.webhook.service.js index b18adf54e..7a98ebae0 100644 --- a/modules/billing/services/billing.webhook.service.js +++ b/modules/billing/services/billing.webhook.service.js @@ -251,10 +251,15 @@ const handleCheckoutCompleted = async (session, event) => { // while the subscription row stays unlinked (#4151). Throwing lets Stripe redeliver. const stripe = getStripe(); if (!stripe) { - logger.error('[billing.webhook] checkout.session.completed — Stripe not configured, aborting', { + // Throw (do NOT return) — same reasoning as the retrieval catch below: a silent + // return here still lets withIdempotency record the event as processed while the + // subscription row stays unlinked, even though nothing was actually done. Throwing + // lets an event received during a Stripe-configuration outage retry once + // configuration is restored (#4155). + logger.error('[billing.webhook] checkout.session.completed — Stripe not configured, will retry via Stripe redelivery', { stripeSubscriptionId, }); - return; + throw new Error('Stripe not configured'); } let realStatus; try { diff --git a/modules/billing/tests/billing.webhook.checkout.unit.tests.js b/modules/billing/tests/billing.webhook.checkout.unit.tests.js index 6e62260d9..eca84e633 100644 --- a/modules/billing/tests/billing.webhook.checkout.unit.tests.js +++ b/modules/billing/tests/billing.webhook.checkout.unit.tests.js @@ -482,7 +482,10 @@ describe('Billing webhook checkout unit tests:', () => { expect(mockSubscriptionRepository.create).not.toHaveBeenCalled(); }); - test('should abort without querying when Stripe is not configured (getStripe returns null)', async () => { + // #4155 — CodeRabbit: a silent return here let withIdempotency record the event as + // processed while the subscription row stayed unlinked. Must throw so Stripe + // redelivers once configuration is restored (mirrors the retrieval-catch below it). + test('throws (does not silently return) when Stripe is not configured (getStripe returns null) — lets Stripe redeliver', async () => { // The billing.webhook.checkout.unit.tests.js mocks stripe.js at module level. // To test the getStripe()=null branch, we reload the module with a null-returning mock. jest.resetModules(); @@ -516,14 +519,16 @@ describe('Billing webhook checkout unit tests:', () => { const mod2 = await import('../services/billing.webhook.service.js'); const svc2 = mod2.default; - await svc2.handleCheckoutCompleted( - { - customer: 'cus_123', - subscription: 'sub_456', - metadata: { organizationId: orgId, plan: 'pro' }, - }, - checkoutEvent, - ); + await expect( + svc2.handleCheckoutCompleted( + { + customer: 'cus_123', + subscription: 'sub_456', + metadata: { organizationId: orgId, plan: 'pro' }, + }, + checkoutEvent, + ), + ).rejects.toThrow('Stripe not configured'); expect(mockSubscriptionRepository.updateIfEventNewer).not.toHaveBeenCalled(); expect(mockSubscriptionRepository.create).not.toHaveBeenCalled(); From 51f9e7d55d472d0e099b67bf81f67bd2c4a6cdf8 Mon Sep 17 00:00:00 2001 From: Pierre Brisorgueil Date: Tue, 29 Sep 2026 11:52:39 +0200 Subject: [PATCH 15/15] =?UTF-8?q?revert(billing):=20drop=20the=20cross-org?= =?UTF-8?q?=20grant=20idempotency=20change=20=E2=80=94=20accepted=20as=20n?= =?UTF-8?q?egligible?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Claude-Session: https://claude.ai/code/session_01EqkUXh5nmxbVBTvhgM6zAo --- ERRORS.md | 1 - .../billing.grantClaim.model.mongoose.js | 65 ----------- .../billing.extraBalance.repository.js | 52 +-------- .../billing.grantClaim.repository.js | 43 -------- ...ance.creditGrant.race.integration.tests.js | 104 ------------------ .../tests/billing.extraBalance.unit.tests.js | 102 +---------------- ...illing.grantClaim.repository.unit.tests.js | 92 ---------------- 7 files changed, 4 insertions(+), 455 deletions(-) delete mode 100644 modules/billing/models/billing.grantClaim.model.mongoose.js delete mode 100644 modules/billing/repositories/billing.grantClaim.repository.js delete mode 100644 modules/billing/tests/billing.extraBalance.creditGrant.race.integration.tests.js delete mode 100644 modules/billing/tests/billing.grantClaim.repository.unit.tests.js diff --git a/ERRORS.md b/ERRORS.md index 333c9bfa8..c52d78c22 100644 --- a/ERRORS.md +++ b/ERRORS.md @@ -51,5 +51,4 @@ Use this file as a compact memory of recurring AI mistakes. - [2026-09-25] billing: computing a percent-of-grant level as `(1 - threshold/100) * signupGrant` (e.g. `500 * (1 - 80/100)`) lands on `99.99999999999997`, not `100`, due to float imprecision -> silently misses an exact `post === level` boundary crossing; use `signupGrant * (100 - threshold) / 100` instead, which is exact at common values; see pierreb-devkit/Node#4117 - [2026-09-28] billing/stripe: `checkout.session.completed`'s subscription-retrieve `catch` returned silently on failure -> the event was recorded as processed by the idempotency wrapper while the subscription row stayed unlinked on the free plan, with no retry; fix = throw instead of return, so the claim persists (attempts increments, the doc is never deleted on failure) and Stripe redelivers; see pierreb-devkit/Node#4151 - [2026-09-28] billing/meter: `incrementMeter` always billed against the subscription's own plan quota, while the admission gate (`billing.quota.service.js`) already treated a fail-closed status (`paused`/`unpaid`/`incomplete`/`incomplete_expired`/`canceled`) as the free plan -> a fail-closed org kept consuming its stale paid-plan quota on the meter even though the gate let it through on free-plan grounds; the fail-closed status list is now one shared constant imported by both the gate and the meter — a literal duplicated in two enforcement paths drifts by design, not by mistake; see pierreb-devkit/Node#4151 -- [2026-09-28] billing/extras: `creditGrant`'s idempotency guard excluded a matching `refId` only within the TARGET org's own ledger -> a retry whose target org changed between attempts (e.g. an org merge/reassignment) was credited twice under the same `refId`; fix = a cross-org existence check on the `refId` before the `getOrCreate`/write steps; see pierreb-devkit/Node#4151 - [2026-09-28] auth: `verifyEmail` read the user by token (`getBrut`) then wrote it in a separate step (`update`) -> two concurrent requests for the same token both passed the read check before either write landed, so both provisioned an organization/grant for the same signup; fix = one atomic `findOneAndUpdate` (`UserService.consumeEmailVerificationToken`) that verifies and clears the token together, so only the first concurrent caller can match the still-unexpired token; see pierreb-devkit/Node#4151 diff --git a/modules/billing/models/billing.grantClaim.model.mongoose.js b/modules/billing/models/billing.grantClaim.model.mongoose.js deleted file mode 100644 index b4ec157bc..000000000 --- a/modules/billing/models/billing.grantClaim.model.mongoose.js +++ /dev/null @@ -1,65 +0,0 @@ -/** - * Module dependencies - */ -import mongoose from 'mongoose'; - -const Schema = mongoose.Schema; - -/** - * BillingGrantClaim Data Model Mongoose - * - * Durable, database-enforced claim on a BillingExtraBalanceRepository.creditGrant - * idempotencyKey (#4155). The ledger (`modules/billing/models/billing.extraBalance.model.mongoose.js`) - * is an embedded array per organization — a unique index on it cannot enforce uniqueness - * ACROSS organizations without a migration against already-deployed data. This is a - * separate, brand-new collection instead: the unique index on `key` is built fresh (no - * pre-existing data, no migration), and MongoDB enforces it permanently — unlike a - * TTL/lease-based lock, there is no window where an expired holder can still write. - * - * The claim is never deleted or expired. `organization` records who claimed it, so a - * later call with the SAME key can tell "my own retry/replay after a crash" (same org — - * fall through to the ordinary per-org ledger guard, which is idempotent on its own) from - * "a different org already holds this key" (the cross-org double-grant this exists to - * prevent — rejected as a duplicate). - */ -const GrantClaimMongoose = new Schema( - { - key: { - type: String, - required: true, - unique: true, - trim: true, - }, - organization: { - type: Schema.ObjectId, - ref: 'Organization', - required: true, - }, - at: { - type: Date, - required: true, - default: () => new Date(), - }, - }, - { - timestamps: false, - }, -); - -/** - * Returns the hex string representation of the document ObjectId. - * @returns {string} Hex string of the ObjectId. - */ -function addID() { - return this._id.toHexString(); -} - -/** - * Model configuration - */ -GrantClaimMongoose.virtual('id').get(addID); -GrantClaimMongoose.set('toJSON', { - virtuals: true, -}); - -mongoose.model('BillingGrantClaim', GrantClaimMongoose); diff --git a/modules/billing/repositories/billing.extraBalance.repository.js b/modules/billing/repositories/billing.extraBalance.repository.js index 8e9b18553..fbbbde948 100644 --- a/modules/billing/repositories/billing.extraBalance.repository.js +++ b/modules/billing/repositories/billing.extraBalance.repository.js @@ -3,7 +3,6 @@ */ import mongoose from 'mongoose'; import AppError from '../../../lib/helpers/AppError.js'; -import BillingGrantClaimRepository from './billing.grantClaim.repository.js'; import BillingExtraBalanceSchema from '../models/billing.extraBalance.schema.js'; import { computeExpiryRemovals } from '../lib/billing.packExpiry.js'; @@ -170,46 +169,13 @@ const debit = async (orgId, amount, refId) => { * @description Atomically credit extra meter units for a non-Stripe grant (e.g. signup free * tier, referral grant). Idempotent: if a ledger entry with the same refId * already exists, the update is a no-op and applied=false is returned. - * Idempotency is enforced ACROSS ALL orgs, not just this one: a retry whose - * target org changed between attempts (e.g. an org merge/reassignment) must - * not be credited twice under the same refId (#4151). - * - * A cross-org check-then-write is not atomic by itself — two concurrent calls - * sharing the same idempotencyKey but resolving to DIFFERENT orgs (e.g. the - * in-process referral listener racing the reconcile cron backfill for the same - * invitation) can both observe "not yet granted" before either writes, and both - * would credit their own org. A TTL/lease-based lock does not close this either - * — it only narrows the window to "the write took longer than the lease" (#4155, - * CodeRabbit, both rounds). Closed instead with a genuinely durable claim: - * `BillingGrantClaimRepository.tryClaim` inserts into a brand-new collection - * with a unique index on `key` — new collection, so the index is built fresh - * with no pre-existing data to migrate, and MongoDB enforces it permanently. - * There is no lease to expire and nothing to release. - * - * Sequence: - * Step -1 — claim the key (BillingGrantClaimRepository.tryClaim). A claim - * held by a DIFFERENT org is the cross-org conflict this exists to catch - * → duplicate_grant immediately. A claim held by THIS SAME org (a replay, - * or this exact call retrying after a crash between claiming and writing) - * falls through — the per-org ledger guard in Step 2 is idempotent on its - * own and decides the outcome from there. The claim itself is never rolled - * back: it is a permanent "this key belongs to this org" record, not a - * lock — a same-org retry after ANY failure (crash, transient DB error) - * re-enters here and proceeds exactly like a fresh call whose claim already - * exists. - * Step 0 — cross-org existence check (any org's ledger, not just this one). - * Legacy backstop: entries written before this claim mechanism existed have - * no claim doc, so a legacy key retried from a reassigned org would pass - * Step -1 as a "fresh" claim without this check. + * 2-step pattern aligned with creditPack: * Step 1 — ensure doc exists (atomic getOrCreate, no-op on replay). * Step 2 — idempotency-guarded credit (no upsert). - * * Idempotency key: `options.refId` when supplied (#3842 referral grants — * several grants per org, one per invitation, e.g. * `referral::referrer`), otherwise the synthetic per-org key - * `-` (signup grant — one per org; embeds orgId, so it can - * never collide across orgs by construction — the cross-org path is only ever - * reachable via an explicit `refId`, i.e. referral grants). + * `-` (signup grant — one per org). * `options.expiresAt` mirrors the creditPack expiry mechanism: the entry is * swept by crons/billing.extrasExpiration.js once past its expiry. * No stripeSessionId required. @@ -246,20 +212,6 @@ const creditGrant = async (orgId, amount, source, { refId = null, expiresAt = nu ...(expiresAt ? { expiresAt } : {}), }; - // Step -1: durable cross-org claim (see JSDoc). A different org already owns this key → - // reject now. This org already owns it (fresh or a replay) → fall through to the - // existing per-org guards, which decide the outcome idempotently on their own. - const claim = await BillingGrantClaimRepository.tryClaim(idempotencyKey, orgId); - if (!claim.claimed && claim.ownerOrgId !== orgId) { - return { doc: null, applied: false, reason: 'duplicate_grant' }; - } - - // Step 0: legacy backstop — cross-org existence check for ledger entries written before - // the claim mechanism existed (no claim doc for those). The step-2 filter below only - // excludes a matching refId within THIS org's own ledger, so it misses those. - const alreadyGranted = await BillingExtraBalance().exists({ 'ledger.refId': idempotencyKey }); - if (alreadyGranted) return { doc: null, applied: false, reason: 'duplicate_grant' }; - // Step 1: ensure the document exists (atomic getOrCreate, no-op if already present). await getOrCreate(orgId); diff --git a/modules/billing/repositories/billing.grantClaim.repository.js b/modules/billing/repositories/billing.grantClaim.repository.js deleted file mode 100644 index d457bd834..000000000 --- a/modules/billing/repositories/billing.grantClaim.repository.js +++ /dev/null @@ -1,43 +0,0 @@ -/** - * Module dependencies - */ -import mongoose from 'mongoose'; -import { isDuplicateKeyError } from '../lib/billing.errors.js'; - -/** - * Lazily resolves the BillingGrantClaim Mongoose model. - * Deferred to keep unit tests importable before model registration. - * @returns {import('mongoose').Model} The registered BillingGrantClaim model. - */ -// biome-ignore lint/correctness/useQwikValidLexicalScope: false positive — Node.js repository, not Qwik -const BillingGrantClaim = () => mongoose.model('BillingGrantClaim'); - -/** - * @function tryClaim - * @description Atomically claim a creditGrant idempotencyKey using the unique index on - * `key` (see models/billing.grantClaim.model.mongoose.js for why this is a - * separate collection rather than a unique index on the ledger). No rollback, - * no TTL — the claim is permanent once inserted. - * @param {string} key - The creditGrant idempotencyKey to claim. - * @param {string} orgId - The organization ObjectId (string) attempting the claim. - * @returns {Promise<{claimed: boolean, ownerOrgId?: string|null}>} `claimed: true` on a - * fresh claim. `claimed: false` with `ownerOrgId` set to the claim holder's - * org (string) when the key is already claimed — by this same org (a - * same-org retry/replay) or a different one (the cross-org conflict #4155 - * exists to catch). - */ -// biome-ignore lint/correctness/useQwikValidLexicalScope: false positive — Node.js repository, not Qwik -const tryClaim = async (key, orgId) => { - try { - await BillingGrantClaim().create({ key, organization: orgId, at: new Date() }); - return { claimed: true }; - } catch (err) { - if (!isDuplicateKeyError(err)) throw err; - const existing = await BillingGrantClaim().findOne({ key }).lean(); - return { claimed: false, ownerOrgId: existing ? String(existing.organization) : null }; - } -}; - -export default { - tryClaim, -}; diff --git a/modules/billing/tests/billing.extraBalance.creditGrant.race.integration.tests.js b/modules/billing/tests/billing.extraBalance.creditGrant.race.integration.tests.js deleted file mode 100644 index fd63c0ed6..000000000 --- a/modules/billing/tests/billing.extraBalance.creditGrant.race.integration.tests.js +++ /dev/null @@ -1,104 +0,0 @@ -/** - * Module dependencies. - */ -import mongoose from 'mongoose'; -import { describe, beforeAll, beforeEach, afterAll, test, expect } from '@jest/globals'; - -import mongooseService from '../../../lib/services/mongoose.js'; - -/** - * Integration tests for creditGrant cross-organization idempotency-key race closure (#4155). - * - * A cross-org check-then-write is not atomic by itself: two concurrent creditGrant calls - * sharing the same idempotencyKey but targeting DIFFERENT orgs (e.g. the in-process - * referral listener racing the reconcile cron backfill for the same invitation) could - * both observe "not yet granted" before either writes, and both would credit their own - * org — double-crediting the same logical grant. A TTL/lease-based lock does not close - * this either — it only narrows the window to "the write took longer than the lease". - * Fixed with a genuinely durable claim (BillingGrantClaimRepository.tryClaim): a brand-new - * collection, unique index on `key`, no TTL, nothing to release. These tests run against - * REAL MongoDB — the unique-index claim is the property under test, which a mocked model - * cannot demonstrate. - */ -describe('BillingExtraBalanceRepository.creditGrant — cross-org race closure integration tests:', () => { - let BillingExtraBalance; - let BillingGrantClaim; - let BillingExtraBalanceRepository; - - const orgA = new mongoose.Types.ObjectId().toString(); - const orgB = new mongoose.Types.ObjectId().toString(); - - beforeAll(async () => { - await mongooseService.loadModels(); - await mongooseService.connect(); - - BillingExtraBalance = mongoose.model('BillingExtraBalance'); - BillingGrantClaim = mongoose.model('BillingGrantClaim'); - // The unique index on `key` is what makes the claim atomic — build it before the first - // concurrent test runs, otherwise the first run can race the (async) index build. - await BillingGrantClaim.syncIndexes(); - BillingExtraBalanceRepository = (await import('../repositories/billing.extraBalance.repository.js')).default; - }); - - beforeEach(async () => { - await Promise.all([ - BillingExtraBalance.deleteMany({ organization: { $in: [orgA, orgB] } }), - BillingGrantClaim.deleteMany({ organization: { $in: [orgA, orgB] } }), - ]); - }); - - afterAll(async () => { - await mongooseService.disconnect(); - }); - - test('concurrent creditGrant calls sharing one refId across TWO orgs → exactly ONE applied:true, ONE ledger entry total, ONE claim', async () => { - const sharedKey = `referral:${new mongoose.Types.ObjectId()}:referrer`; - - const [first, second] = await Promise.all([ - BillingExtraBalanceRepository.creditGrant(orgA, 1000, 'referral', { refId: sharedKey }), - BillingExtraBalanceRepository.creditGrant(orgB, 1000, 'referral', { refId: sharedKey }), - ]); - - // Exactly one org wins the grant; the loser surfaces the idempotent no-op — never both. - const applied = [first, second].filter((r) => r.applied === true); - expect(applied).toHaveLength(1); - const duplicate = [first, second].filter((r) => r.applied === false); - expect(duplicate).toHaveLength(1); - expect(duplicate[0].reason).toBe('duplicate_grant'); - - // The key was credited to exactly ONE of the two orgs, never both. - const [docA, docB] = await Promise.all([ - BillingExtraBalance.findOne({ organization: orgA }).lean(), - BillingExtraBalance.findOne({ organization: orgB }).lean(), - ]); - const entriesA = (docA?.ledger ?? []).filter((e) => e.refId === sharedKey); - const entriesB = (docB?.ledger ?? []).filter((e) => e.refId === sharedKey); - expect(entriesA.length + entriesB.length).toBe(1); - - // Exactly one claim document for the key, owned by whichever org actually won. - const claims = await BillingGrantClaim.find({ key: sharedKey }).lean(); - expect(claims).toHaveLength(1); - const winnerOrg = entriesA.length === 1 ? orgA : orgB; - expect(String(claims[0].organization)).toBe(winnerOrg); - }, 15000); - - test('a losing call leaves the claim permanently owned by the winner — a later replay for the SAME org+key still resolves idempotently', async () => { - const sharedKey = `referral:${new mongoose.Types.ObjectId()}:referrer`; - - const [first] = await Promise.all([ - BillingExtraBalanceRepository.creditGrant(orgA, 1000, 'referral', { refId: sharedKey }), - BillingExtraBalanceRepository.creditGrant(orgB, 1000, 'referral', { refId: sharedKey }), - ]); - const winnerOrg = first.applied ? orgA : orgB; - - // Sequential replay against the org that actually won — must stay a no-op. No lock to - // release: the claim is permanent, the per-org ledger guard is the durable dedup. - const replay = await BillingExtraBalanceRepository.creditGrant(winnerOrg, 1000, 'referral', { refId: sharedKey }); - expect(replay).toMatchObject({ applied: false, reason: 'duplicate_grant' }); - - // A retry from the LOSING org must still be rejected — the claim does not expire. - const loserOrg = winnerOrg === orgA ? orgB : orgA; - const loserRetry = await BillingExtraBalanceRepository.creditGrant(loserOrg, 1000, 'referral', { refId: sharedKey }); - expect(loserRetry).toMatchObject({ applied: false, reason: 'duplicate_grant' }); - }, 15000); -}); diff --git a/modules/billing/tests/billing.extraBalance.unit.tests.js b/modules/billing/tests/billing.extraBalance.unit.tests.js index 8223743ba..6238a5f11 100644 --- a/modules/billing/tests/billing.extraBalance.unit.tests.js +++ b/modules/billing/tests/billing.extraBalance.unit.tests.js @@ -146,7 +146,6 @@ describe('BillingExtraBalance unit tests:', () => { describe('Repository', () => { let BillingExtraBalanceRepository; let mockModel; - let mockClaimModel; const orgId = '507f1f77bcf86cd799439011'; /** @@ -173,20 +172,9 @@ describe('BillingExtraBalance unit tests:', () => { exists: jest.fn(), }; - // #4155 — creditGrant claims its idempotencyKey via BillingGrantClaimRepository, - // which resolves its OWN 'BillingGrantClaim' model. Keyed by name so that model does - // not share mockModel's findOneAndUpdate call queue with the ExtraBalance - // step1/step2 sequencing used throughout this describe block. Defaults to "claim - // always succeeds" (fresh key) — individual tests override mockClaimModel.create to - // simulate an existing claim (same-org replay or a different org's conflict). - mockClaimModel = { - create: jest.fn().mockResolvedValue({}), - findOne: jest.fn(() => ({ lean: jest.fn().mockResolvedValue(null) })), - }; - jest.unstable_mockModule('mongoose', () => ({ default: { - model: jest.fn((name) => (name === 'BillingGrantClaim' ? mockClaimModel : mockModel)), + model: jest.fn(() => mockModel), Types: { ObjectId: { isValid: jest.fn(() => true), @@ -1196,84 +1184,6 @@ describe('BillingExtraBalance unit tests:', () => { expect(result).toEqual({ doc: null, applied: false }); expect(mockModel.findOneAndUpdate).not.toHaveBeenCalled(); }); - - // #4151 — the idempotency guard in step 2 below only excludes a matching refId - // WITHIN this org's own ledger ('ledger.refId': { $ne }), so a retry whose target - // org changed between attempts (e.g. an org merge/reassignment) is credited twice. - // Legacy backstop (step 0, after the #4155 claim below): a ledger entry written - // before the claim mechanism existed has no claim doc, so the claim itself looks - // "fresh" — this cross-org existence check is what still catches it. - test('returns applied:false duplicate_grant when the refId already exists in a DIFFERENT org\'s ledger (legacy, pre-claim data)', async () => { - // Claim succeeds (fresh key, default mock) — no claim doc for this pre-existing entry. - // Cross-org existence check then finds the prior grant under this refId in any org. - mockModel.exists.mockResolvedValue({ _id: 'some-other-orgs-balance-doc' }); - - const result = await BillingExtraBalanceRepository.creditGrant(orgId, 500, 'referral', { - refId: 'referral:64b2f0000000000000000001:referrer', - }); - - expect(result).toEqual({ doc: null, applied: false, reason: 'duplicate_grant' }); - expect(mockModel.exists).toHaveBeenCalledWith({ 'ledger.refId': 'referral:64b2f0000000000000000001:referrer' }); - // Must short-circuit before any write. - expect(mockModel.findOneAndUpdate).not.toHaveBeenCalled(); - }); - - // #4155 — the durable claim itself: two concurrent creditGrant calls sharing one - // idempotencyKey but targeting DIFFERENT orgs must not both credit. The claim's - // unique index on `key` is what makes this atomic (not a lease/lock), so the second - // caller's create() rejects with a duplicate-key error before either org is touched. - describe('creditGrant — durable claim (#4155):', () => { - test('a claim already held by a DIFFERENT org → duplicate_grant, before any ledger read or write', async () => { - const otherOrgId = '507f1f77bcf86cd799439022'; - mockClaimModel.create.mockRejectedValue(Object.assign(new Error('E11000 duplicate key'), { code: 11000 })); - mockClaimModel.findOne.mockReturnValue({ lean: jest.fn().mockResolvedValue({ organization: otherOrgId }) }); - - const result = await BillingExtraBalanceRepository.creditGrant(orgId, 500, 'referral', { - refId: 'referral:64b2f0000000000000000002:referrer', - }); - - expect(result).toEqual({ doc: null, applied: false, reason: 'duplicate_grant' }); - expect(mockClaimModel.create).toHaveBeenCalledWith( - expect.objectContaining({ key: 'referral:64b2f0000000000000000002:referrer', organization: orgId }), - ); - // Rejected at the claim — never reaches the ledger at all. - expect(mockModel.exists).not.toHaveBeenCalled(); - expect(mockModel.findOneAndUpdate).not.toHaveBeenCalled(); - }); - - test('a claim already held by the SAME org (replay, or retry after a crash) falls through to the ordinary per-org guard', async () => { - mockClaimModel.create.mockRejectedValue(Object.assign(new Error('E11000 duplicate key'), { code: 11000 })); - mockClaimModel.findOne.mockReturnValue({ lean: jest.fn().mockResolvedValue({ organization: orgId }) }); - // Step 1: getOrCreate; Step 2: per-org guard says already credited. - mockModel.findOneAndUpdate - .mockResolvedValueOnce(makeDoc()) - .mockResolvedValueOnce(null); - - const result = await BillingExtraBalanceRepository.creditGrant(orgId, 500, 'referral', { - refId: 'referral:64b2f0000000000000000003:referrer', - }); - - expect(result).toEqual({ doc: null, applied: false, reason: 'duplicate_grant' }); - // Did NOT short-circuit — the same-org branch reaches the ordinary guards. - expect(mockModel.findOneAndUpdate).toHaveBeenCalledTimes(2); - }); - - test('a fresh claim proceeds normally and credits the org', async () => { - const updatedDoc = makeDoc({ cachedBalance: 500 }); - mockModel.findOneAndUpdate - .mockResolvedValueOnce(makeDoc()) - .mockResolvedValueOnce(updatedDoc); - - const result = await BillingExtraBalanceRepository.creditGrant(orgId, 500, 'referral', { - refId: 'referral:64b2f0000000000000000004:referrer', - }); - - expect(result.applied).toBe(true); - expect(mockClaimModel.create).toHaveBeenCalledWith( - expect.objectContaining({ key: 'referral:64b2f0000000000000000004:referrer', organization: orgId }), - ); - }); - }); }); }); }); @@ -1394,7 +1304,6 @@ describe('Referral grant extensions (#3842):', () => { describe('repository', () => { let BillingExtraBalanceRepository; let mockModel; - let mockClaimModel; /** * @param {Object[]} rows - Rows the aggregation should resolve with. @@ -1414,16 +1323,9 @@ describe('Referral grant extensions (#3842):', () => { exists: jest.fn(), aggregate: jest.fn(), }; - // #4155 — see the identical comment in the outer 'Repository' describe's beforeEach: - // creditGrant claims its idempotencyKey via its own 'BillingGrantClaim' model, kept - // off mockModel's findOneAndUpdate call queue. - mockClaimModel = { - create: jest.fn().mockResolvedValue({}), - findOne: jest.fn(() => ({ lean: jest.fn().mockResolvedValue(null) })), - }; jest.unstable_mockModule('mongoose', () => ({ default: { - model: jest.fn((name) => (name === 'BillingGrantClaim' ? mockClaimModel : mockModel)), + model: jest.fn(() => mockModel), Types: { ObjectId: { isValid: jest.fn(() => true) } }, }, })); diff --git a/modules/billing/tests/billing.grantClaim.repository.unit.tests.js b/modules/billing/tests/billing.grantClaim.repository.unit.tests.js deleted file mode 100644 index ab6509ca1..000000000 --- a/modules/billing/tests/billing.grantClaim.repository.unit.tests.js +++ /dev/null @@ -1,92 +0,0 @@ -/** - * Module dependencies. - */ -import { jest, describe, test, beforeEach, afterEach, expect } from '@jest/globals'; - -/** - * Unit tests for billing.grantClaim.repository.js (#4155). - */ -describe('BillingGrantClaimRepository unit tests:', () => { - let BillingGrantClaimRepository; - let mockModel; - - const orgId = '507f1f77bcf86cd799439011'; - const otherOrgId = '507f1f77bcf86cd799439022'; - - const makeE11000 = () => { - const err = new Error('E11000 duplicate key error'); - err.code = 11000; - return err; - }; - - beforeEach(async () => { - jest.resetModules(); - - mockModel = { - create: jest.fn(), - findOne: jest.fn(), - }; - - jest.unstable_mockModule('mongoose', () => ({ - default: { - model: jest.fn(() => mockModel), - }, - })); - - const mod = await import('../repositories/billing.grantClaim.repository.js'); - BillingGrantClaimRepository = mod.default; - }); - - afterEach(() => { - jest.restoreAllMocks(); - }); - - describe('tryClaim', () => { - test('fresh key → { claimed: true }, no lookup needed', async () => { - mockModel.create.mockResolvedValue({ key: 'k1', organization: orgId }); - - const result = await BillingGrantClaimRepository.tryClaim('k1', orgId); - - expect(result).toEqual({ claimed: true }); - expect(mockModel.create).toHaveBeenCalledWith( - expect.objectContaining({ key: 'k1', organization: orgId }), - ); - expect(mockModel.findOne).not.toHaveBeenCalled(); - }); - - test('duplicate key already owned by the SAME org → { claimed: false, ownerOrgId: orgId }', async () => { - mockModel.create.mockRejectedValueOnce(makeE11000()); - mockModel.findOne.mockReturnValue({ lean: jest.fn().mockResolvedValue({ key: 'k1', organization: orgId }) }); - - const result = await BillingGrantClaimRepository.tryClaim('k1', orgId); - - expect(result).toEqual({ claimed: false, ownerOrgId: orgId }); - expect(mockModel.findOne).toHaveBeenCalledWith({ key: 'k1' }); - }); - - test('duplicate key already owned by a DIFFERENT org → { claimed: false, ownerOrgId: otherOrgId }', async () => { - mockModel.create.mockRejectedValueOnce(makeE11000()); - mockModel.findOne.mockReturnValue({ lean: jest.fn().mockResolvedValue({ key: 'k1', organization: otherOrgId }) }); - - const result = await BillingGrantClaimRepository.tryClaim('k1', orgId); - - expect(result).toEqual({ claimed: false, ownerOrgId: otherOrgId }); - }); - - test('duplicate key but lookup finds nothing (race window: claim deleted between insert-fail and lookup) → ownerOrgId: null', async () => { - mockModel.create.mockRejectedValueOnce(makeE11000()); - mockModel.findOne.mockReturnValue({ lean: jest.fn().mockResolvedValue(null) }); - - const result = await BillingGrantClaimRepository.tryClaim('k1', orgId); - - expect(result).toEqual({ claimed: false, ownerOrgId: null }); - }); - - test('a non-duplicate-key error propagates (never swallowed as a claim conflict)', async () => { - mockModel.create.mockRejectedValueOnce(new Error('connection reset')); - - await expect(BillingGrantClaimRepository.tryClaim('k1', orgId)).rejects.toThrow('connection reset'); - expect(mockModel.findOne).not.toHaveBeenCalled(); - }); - }); -});