From cc5545c3332045eb00249ec5c6a1169cbb71536e Mon Sep 17 00:00:00 2001 From: Hampus Date: Tue, 6 Oct 2026 21:38:25 +0200 Subject: [PATCH] fix(api): sync stripe customer email on change (#3243) --- .../admin/services/AdminUserProfileService.ts | 3 + .../tests/AdminUserChangeLogAndFlags.test.ts | 37 ++- fluxer_api/src/api/auth/AuthEmailRevert.ts | 4 +- .../tests/BouncedEmailRecoveryFlow.test.ts | 13 +- .../api/auth/tests/EmailChangeFlow.test.ts | 14 +- .../api/auth/tests/EmailRevertFlow.test.ts | 14 +- fluxer_api/src/api/stripe/StripeCustomer.ts | 119 +++++++++ .../stripe/services/AgeVerificationService.ts | 66 +---- .../stripe/services/StripeCheckoutService.ts | 65 +---- .../tests/StripeCustomerEmailSync.test.ts | 233 ++++++++++++++++++ .../test/msw/handlers/StripeApiHandlers.ts | 14 ++ .../api/user/services/EmailChangeService.ts | 4 +- .../api/user/services/UserAccountService.ts | 2 + fluxer_api/src/api/worker/WorkerLaneConfig.ts | 1 + .../src/api/worker/WorkerTaskRegistry.ts | 2 + .../api/worker/tasks/ReconcileUserPayments.ts | 6 + .../worker/tasks/SyncStripeCustomerEmail.ts | 37 +++ 17 files changed, 509 insertions(+), 125 deletions(-) create mode 100644 fluxer_api/src/api/stripe/StripeCustomer.ts create mode 100644 fluxer_api/src/api/stripe/tests/StripeCustomerEmailSync.test.ts create mode 100644 fluxer_api/src/api/worker/tasks/SyncStripeCustomerEmail.ts diff --git a/fluxer_api/src/api/admin/services/AdminUserProfileService.ts b/fluxer_api/src/api/admin/services/AdminUserProfileService.ts index bfa61a9c6..6ec3f4a3f 100644 --- a/fluxer_api/src/api/admin/services/AdminUserProfileService.ts +++ b/fluxer_api/src/api/admin/services/AdminUserProfileService.ts @@ -12,6 +12,7 @@ import type {EntityAssetService, PreparedAssetUpload} from '@app/api/infrastruct import {Logger} from '@app/api/Logger'; import {getInstanceConfigRepository} from '@app/api/middleware/ServiceSingletons'; import type {User} from '@app/api/models/User'; +import {enqueueStripeCustomerEmailSync} from '@app/api/stripe/StripeCustomer'; import {assertNoDiscriminatorChange, reserveUsername, type UsernameReservation} from '@app/api/user/UniqueUsernames'; import {USERNAME_MODE_DISCRIMINATOR} from '@app/api/user/UserTag'; import {TagAlreadyTakenError} from '@fluxer/errors/src/domains/user/TagAlreadyTakenError'; @@ -218,6 +219,7 @@ export class AdminUserProfileService { users: userRepository, cache: cacheService, contactChangeLog: contactChangeLogService, + worker: workerService, } = this.deps.apiContext.services; const {auditService, updatePropagator} = this.deps; const userId = createUserID(data.user_id); @@ -240,6 +242,7 @@ export class AdminUserProfileService { reason: 'admin_action', actorUserId: adminUserId, }); + await enqueueStripeCustomerEmailSync(workerService, user, updatedUser); await auditService.createAuditLog({ adminUserId, targetType: 'user', diff --git a/fluxer_api/src/api/admin/tests/AdminUserChangeLogAndFlags.test.ts b/fluxer_api/src/api/admin/tests/AdminUserChangeLogAndFlags.test.ts index 4ad4deb6b..4e8847b8f 100644 --- a/fluxer_api/src/api/admin/tests/AdminUserChangeLogAndFlags.test.ts +++ b/fluxer_api/src/api/admin/tests/AdminUserChangeLogAndFlags.test.ts @@ -2,11 +2,12 @@ import {createTestAccount, setUserACLs} from '@app/api/auth/tests/AuthTestUtils'; import {type ApiTestHarness, createApiTestHarness} from '@app/api/test/ApiTestHarness'; +import {NoopWorkerService} from '@app/api/test/NoopWorkerService'; import {HTTP_STATUS, TEST_CREDENTIALS} from '@app/api/test/TestConstants'; import {createBuilder, createBuilderWithoutAuth} from '@app/api/test/TestRequestBuilder'; import {AdminACLs} from '@fluxer/constants/src/AdminACLs'; import {UserFlags} from '@fluxer/constants/src/UserConstants'; -import {afterAll, beforeAll, beforeEach, describe, expect, test} from 'vitest'; +import {afterAll, afterEach, beforeAll, beforeEach, describe, expect, test, vi} from 'vitest'; interface ChangeLogResponse { entries: Array<{ @@ -44,6 +45,9 @@ describe('Admin User Change Log and Flags', () => { beforeEach(async () => { await harness.reset(); }); + afterEach(() => { + vi.restoreAllMocks(); + }); afterAll(async () => { await harness?.shutdown(); }); @@ -176,6 +180,37 @@ describe('Admin User Change Log and Flags', () => { .execute(); }); }); + describe('PATCH /admin/users/{user_id}/email', () => { + test('queues a Stripe customer email sync for users with a Stripe customer', async () => { + const admin = await createTestAccount(harness); + await setUserACLs(harness, admin, [AdminACLs.AUTHENTICATE, AdminACLs.WILDCARD]); + const target = await createTestAccount(harness); + await createBuilderWithoutAuth(harness) + .post(`/test/users/${target.userId}/premium`) + .body({stripe_customer_id: 'cus_admin_email_sync'}) + .expect(HTTP_STATUS.OK) + .execute(); + const addJob = vi.spyOn(NoopWorkerService.prototype, 'addJob'); + await createBuilder(harness, `${admin.token}`) + .patch(`/admin/users/${target.userId}/email`) + .body({email: `admin-changed-${Date.now()}@example.com`}) + .expect(HTTP_STATUS.OK) + .execute(); + expect(addJob).toHaveBeenCalledWith('syncStripeCustomerEmail', {userId: target.userId}); + }); + test('does not queue a Stripe customer email sync for users without a Stripe customer', async () => { + const admin = await createTestAccount(harness); + await setUserACLs(harness, admin, [AdminACLs.AUTHENTICATE, AdminACLs.WILDCARD]); + const target = await createTestAccount(harness); + const addJob = vi.spyOn(NoopWorkerService.prototype, 'addJob'); + await createBuilder(harness, `${admin.token}`) + .patch(`/admin/users/${target.userId}/email`) + .body({email: `admin-changed-${Date.now()}@example.com`}) + .expect(HTTP_STATUS.OK) + .execute(); + expect(addJob).not.toHaveBeenCalledWith('syncStripeCustomerEmail', expect.anything()); + }); + }); describe('PUT /admin/users/{user_id}/email-verification', () => { test('verifying email clears email_bounced', async () => { const admin = await createTestAccount(harness); diff --git a/fluxer_api/src/api/auth/AuthEmailRevert.ts b/fluxer_api/src/api/auth/AuthEmailRevert.ts index 1982df0f2..b0409e330 100644 --- a/fluxer_api/src/api/auth/AuthEmailRevert.ts +++ b/fluxer_api/src/api/auth/AuthEmailRevert.ts @@ -6,6 +6,7 @@ import * as AuthSession from '@app/api/auth/AuthSession'; import * as AuthUtility from '@app/api/auth/AuthUtility'; import {createEmailRevertToken} from '@app/api/BrandedTypes'; import type {User} from '@app/api/models/User'; +import {enqueueStripeCustomerEmailSync} from '@app/api/stripe/StripeCustomer'; import {mapUserToPrivateResponse} from '@app/api/user/UserMappers'; import {ValidationErrorCodes} from '@fluxer/constants/src/ValidationErrorCodes'; import {InputValidationError} from '@fluxer/errors/src/domains/core/InputValidationError'; @@ -44,7 +45,7 @@ export async function revertEmailChange( user_id: string; token: string; }> { - const {users, gateway, contactChangeLog, config} = ctx.services; + const {users, gateway, contactChangeLog, config, worker} = ctx.services; const {token, password, request} = params; const tokenData = await users.getEmailRevertToken(token); if (!tokenData) { @@ -101,5 +102,6 @@ export async function revertEmailChange( reason: 'user_requested', actorUserId: user.id, }); + await enqueueStripeCustomerEmailSync(worker, user, updatedUser); return {user_id: updatedUser.id.toString(), token: authToken}; } diff --git a/fluxer_api/src/api/auth/tests/BouncedEmailRecoveryFlow.test.ts b/fluxer_api/src/api/auth/tests/BouncedEmailRecoveryFlow.test.ts index 2518610ed..908f52af0 100644 --- a/fluxer_api/src/api/auth/tests/BouncedEmailRecoveryFlow.test.ts +++ b/fluxer_api/src/api/auth/tests/BouncedEmailRecoveryFlow.test.ts @@ -9,8 +9,9 @@ import { type TestAccount, } from '@app/api/auth/tests/AuthTestUtils'; import type {ApiTestHarness} from '@app/api/test/ApiTestHarness'; +import {NoopWorkerService} from '@app/api/test/NoopWorkerService'; import {createBuilder, createBuilderWithoutAuth} from '@app/api/test/TestRequestBuilder'; -import {afterAll, beforeAll, beforeEach, describe, expect, it} from 'vitest'; +import {afterAll, afterEach, beforeAll, beforeEach, describe, expect, it, vi} from 'vitest'; interface BouncedEmailRequestNewResponse { ticket: string; @@ -51,12 +52,21 @@ describe('Bounced email recovery flow', () => { await harness.reset(); await clearTestEmails(harness); }); + afterEach(() => { + vi.restoreAllMocks(); + }); afterAll(async () => { await harness?.shutdown(); }); it('allows bounced users to replace email without original-email verification', async () => { const account = await createTestAccount(harness); await markEmailAsBounced(harness, account); + await createBuilderWithoutAuth(harness) + .post(`/test/users/${account.userId}/premium`) + .body({stripe_customer_id: 'cus_bounced_email_sync'}) + .expect(200) + .execute(); + const addJob = vi.spyOn(NoopWorkerService.prototype, 'addJob'); const initialMe = await createBuilder(harness, account.token) .get('/users/@me') .expect(200) @@ -94,6 +104,7 @@ describe('Bounced email recovery flow', () => { const finalMe = await createBuilder(harness, account.token).get('/users/@me').execute(); expect(finalMe.email).toBe(replacementEmail); expect(finalMe.email_bounced).toBe(false); + expect(addJob).toHaveBeenCalledWith('syncStripeCustomerEmail', {userId: account.userId}); }); it('rejects bounced-email recovery for accounts that are not marked as bounced', async () => { const account = await createTestAccount(harness); diff --git a/fluxer_api/src/api/auth/tests/EmailChangeFlow.test.ts b/fluxer_api/src/api/auth/tests/EmailChangeFlow.test.ts index 66550e9ce..8810921db 100644 --- a/fluxer_api/src/api/auth/tests/EmailChangeFlow.test.ts +++ b/fluxer_api/src/api/auth/tests/EmailChangeFlow.test.ts @@ -12,8 +12,9 @@ import { type TestAccount, } from '@app/api/auth/tests/AuthTestUtils'; import type {ApiTestHarness} from '@app/api/test/ApiTestHarness'; +import {NoopWorkerService} from '@app/api/test/NoopWorkerService'; import {createBuilder, createBuilderWithoutAuth} from '@app/api/test/TestRequestBuilder'; -import {afterAll, beforeAll, beforeEach, describe, expect, it} from 'vitest'; +import {afterAll, afterEach, beforeAll, beforeEach, describe, expect, it, vi} from 'vitest'; interface EmailChangeStartResponse { ticket: string; @@ -129,6 +130,9 @@ describe('Email change flow', () => { await harness.reset(); await clearTestEmails(harness); }); + afterEach(() => { + vi.restoreAllMocks(); + }); afterAll(async () => { await harness?.shutdown(); }); @@ -341,13 +345,14 @@ describe('Email change flow', () => { .execute(); expect(updated.email).toBe(newEmail); }); - it('applies email changes for users who have ever purchased', async () => { + it('applies email changes for users who have ever purchased and syncs their Stripe customer', async () => { const account = await createTestAccount(harness); await createBuilderWithoutAuth(harness) .post(`/test/users/${account.userId}/premium`) - .body({has_ever_purchased: true}) + .body({has_ever_purchased: true, stripe_customer_id: 'cus_email_change_sync'}) .expect(200) .execute(); + const addJob = vi.spyOn(NoopWorkerService.prototype, 'addJob'); const startResp = await startEmailChange(harness, account, account.password); let originalProof: string; if (startResp.require_original) { @@ -387,9 +392,11 @@ describe('Email change flow', () => { .execute(); expect(updated.email).toBe(newEmail); expect(updated.has_ever_purchased).toBe(true); + expect(addJob).toHaveBeenCalledWith('syncStripeCustomerEmail', {userId: account.userId}); }); it('applies ordinary claimed email changes', async () => { const account = await createTestAccount(harness); + const addJob = vi.spyOn(NoopWorkerService.prototype, 'addJob'); const startResp = await startEmailChange(harness, account, account.password); const emails = await listTestEmails(harness, {recipient: account.email}); const originalEmail = findLastTestEmail(emails, 'email_change_original'); @@ -425,6 +432,7 @@ describe('Email change flow', () => { .execute(); expect(updated.email).toBe(newEmail); expect(updated.verified).toBe(true); + expect(addJob).not.toHaveBeenCalledWith('syncStripeCustomerEmail', expect.anything()); }); it('requires MFA (not password) for email_token apply when user has TOTP enabled', async () => { const account = await createTestAccount(harness); diff --git a/fluxer_api/src/api/auth/tests/EmailRevertFlow.test.ts b/fluxer_api/src/api/auth/tests/EmailRevertFlow.test.ts index 97538389d..8e2b08fab 100644 --- a/fluxer_api/src/api/auth/tests/EmailRevertFlow.test.ts +++ b/fluxer_api/src/api/auth/tests/EmailRevertFlow.test.ts @@ -10,8 +10,9 @@ import { type TestAccount, } from '@app/api/auth/tests/AuthTestUtils'; import type {ApiTestHarness} from '@app/api/test/ApiTestHarness'; +import {NoopWorkerService} from '@app/api/test/NoopWorkerService'; import {createBuilder, createBuilderWithoutAuth} from '@app/api/test/TestRequestBuilder'; -import {afterAll, beforeAll, beforeEach, describe, expect, it} from 'vitest'; +import {afterAll, afterEach, beforeAll, beforeEach, describe, expect, it, vi} from 'vitest'; interface EmailChangeStartResponse { ticket: string; @@ -130,11 +131,20 @@ describe('Email revert flow', () => { await harness.reset(); await clearTestEmails(harness); }); + afterEach(() => { + vi.restoreAllMocks(); + }); afterAll(async () => { await harness?.shutdown(); }); it('restores original email and clears mfa', async () => { const account = await createTestAccount(harness); + await createBuilderWithoutAuth(harness) + .post(`/test/users/${account.userId}/premium`) + .body({stripe_customer_id: 'cus_email_revert_sync'}) + .expect(200) + .execute(); + const addJob = vi.spyOn(NoopWorkerService.prototype, 'addJob'); const startResp = await startEmailChange(harness, account, account.password); let originalProof: string; if (startResp.require_original) { @@ -178,6 +188,7 @@ describe('Email revert flow', () => { expect(revertEmail?.metadata?.token).toBeDefined(); const revertToken = revertEmail!.metadata!.token!; const newPassword = uniquePassword(); + addJob.mockClear(); const revertResp = await createBuilderWithoutAuth(harness) .post('/auth/email-revert') .body({ @@ -186,6 +197,7 @@ describe('Email revert flow', () => { }) .execute(); expect(revertResp.token.length).toBeGreaterThan(0); + expect(addJob).toHaveBeenCalledWith('syncStripeCustomerEmail', {userId: account.userId}); await createBuilder(harness, account.token).get('/users/@me').expect(401).execute(); const user = await createBuilder(harness, revertResp.token).get('/users/@me').execute(); expect(user.email).toBe(account.email); diff --git a/fluxer_api/src/api/stripe/StripeCustomer.ts b/fluxer_api/src/api/stripe/StripeCustomer.ts new file mode 100644 index 000000000..c6e5e848e --- /dev/null +++ b/fluxer_api/src/api/stripe/StripeCustomer.ts @@ -0,0 +1,119 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later + +import {Logger} from '@app/api/Logger'; +import {getBillingRepository} from '@app/api/middleware/ServiceRegistry'; +import type {User} from '@app/api/models/User'; +import type {IUserRepository} from '@app/api/user/IUserRepository'; +import type {WorkerTaskName} from '@app/api/worker/WorkerLaneConfig'; +import {StripeError} from '@fluxer/errors/src/domains/payment/StripeError'; +import type {ICacheService} from '@pkgs/cache/src/ICacheService'; +import type {IWorkerService} from '@pkgs/worker/src/contracts/IWorkerService'; +import {seconds} from 'itty-time'; +import type Stripe from 'stripe'; + +const CUSTOMER_LOCK_TTL_SECONDS = seconds('30 seconds'); + +interface EnsureStripeCustomerParams { + stripe: Stripe; + user: User; + userRepository: IUserRepository; + cacheService: ICacheService; +} + +export function isStripeResourceMissingError(error: unknown): boolean { + return typeof error === 'object' && error !== null && 'code' in error && error.code === 'resource_missing'; +} + +export async function ensureStripeCustomer({ + stripe, + user, + userRepository, + cacheService, +}: EnsureStripeCustomerParams): Promise { + if (user.stripeCustomerId) { + try { + await syncStripeCustomerEmail(stripe, user); + } catch (error) { + Logger.warn( + {error, userId: user.id, customerId: user.stripeCustomerId}, + 'Failed to sync Stripe customer email, continuing with the stored customer', + ); + } + return user; + } + const lockKey = `stripe_customer_create_lock:${user.id}`; + const lockToken = await cacheService.acquireLock(lockKey, CUSTOMER_LOCK_TTL_SECONDS); + if (!lockToken) { + const freshUser = await userRepository.findUnique(user.id); + if (freshUser?.stripeCustomerId) { + return freshUser; + } + throw new StripeError('Failed to acquire customer creation lock'); + } + try { + const freshUser = await userRepository.findUnique(user.id); + if (freshUser?.stripeCustomerId) { + return freshUser; + } + const customer = await stripe.customers.create({ + email: user.email ?? undefined, + metadata: { + userId: user.id.toString(), + }, + }); + await mirrorStripeCustomer(customer, user); + const updatedUser = await userRepository.patchUpsert(user.id, {stripe_customer_id: customer.id}, user.toRow()); + Logger.debug({userId: user.id, customerId: customer.id}, 'Stripe customer created'); + return updatedUser; + } finally { + try { + const released = await cacheService.releaseLock(lockKey, lockToken); + if (!released) { + Logger.warn({userId: user.id, lockKey}, 'Customer creation lock token no longer matched on release'); + } + } catch (error) { + Logger.error({error, userId: user.id, lockKey}, 'Failed to release customer creation lock'); + } + } +} + +export async function syncStripeCustomerEmail(stripe: Stripe, user: User): Promise { + const customerId = user.stripeCustomerId; + const email = user.email; + if (!customerId || !email) { + return; + } + const mirrored = await getBillingRepository().customers.findById(customerId); + if (mirrored?.deleted || mirrored?.email === email) { + return; + } + const customer = await stripe.customers.update(customerId, {email}); + await mirrorStripeCustomer(customer, user); + Logger.info({userId: user.id, customerId}, 'Synced Stripe customer email with account email'); +} + +export async function enqueueStripeCustomerEmailSync( + workerService: IWorkerService, + oldUser: User, + newUser: User, +): Promise { + if (!newUser.stripeCustomerId || !newUser.email || oldUser.email === newUser.email) { + return; + } + try { + await workerService.addJob('syncStripeCustomerEmail', {userId: newUser.id.toString()}); + } catch (error) { + Logger.warn( + {error, userId: newUser.id, customerId: newUser.stripeCustomerId}, + 'Failed to enqueue Stripe customer email sync, checkout will sync it instead', + ); + } +} + +async function mirrorStripeCustomer(customer: Stripe.Customer, user: User): Promise { + try { + await getBillingRepository().customers.upsertFromStripe(customer, {knownUserId: user.id}); + } catch (mirrorErr) { + Logger.error({mirrorErr, customerId: customer.id}, 'Mirror upsert failed after Stripe write; reconciler will heal'); + } +} diff --git a/fluxer_api/src/api/stripe/services/AgeVerificationService.ts b/fluxer_api/src/api/stripe/services/AgeVerificationService.ts index ff849c9c1..081e36bd0 100644 --- a/fluxer_api/src/api/stripe/services/AgeVerificationService.ts +++ b/fluxer_api/src/api/stripe/services/AgeVerificationService.ts @@ -5,7 +5,7 @@ import {Config} from '@app/api/Config'; import type {IGatewayService} from '@app/api/infrastructure/IGatewayService'; import {Logger} from '@app/api/Logger'; import {getBillingRepository} from '@app/api/middleware/ServiceRegistry'; -import type {User} from '@app/api/models/User'; +import {ensureStripeCustomer} from '@app/api/stripe/StripeCustomer'; import type {IUserRepository} from '@app/api/user/IUserRepository'; import {mapUserToPrivateResponse} from '@app/api/user/UserMappers'; import {UserFlags} from '@fluxer/constants/src/UserConstants'; @@ -14,11 +14,8 @@ import {StripeError} from '@fluxer/errors/src/domains/payment/StripeError'; import {StripePaymentNotAvailableError} from '@fluxer/errors/src/domains/payment/StripePaymentNotAvailableError'; import {UnknownUserError} from '@fluxer/errors/src/domains/user/UnknownUserError'; import type {ICacheService} from '@pkgs/cache/src/ICacheService'; -import {seconds} from 'itty-time'; import type Stripe from 'stripe'; -const CUSTOMER_LOCK_TTL_SECONDS = seconds('30 seconds'); - export class AgeVerificationService { constructor( private stripe: Stripe | null, @@ -38,7 +35,12 @@ export class AgeVerificationService { if (user.flags & UserFlags.AGE_VERIFIED_ADULT) { throw new AgeVerificationAlreadyVerifiedError(); } - const customerUser = await this.ensureStripeCustomer(user); + const customerUser = await ensureStripeCustomer({ + stripe: this.stripe, + user, + userRepository: this.userRepository, + cacheService: this.cacheService, + }); const customerId = customerUser.stripeCustomerId; if (!customerId) { throw new StripeError('Stripe customer id missing after customer setup'); @@ -134,58 +136,4 @@ export class AgeVerificationService { }); Logger.info({userId}, 'Age verification completed successfully'); } - - private async ensureStripeCustomer(user: User): Promise { - if (user.stripeCustomerId) { - return user; - } - if (!this.stripe) { - throw new StripePaymentNotAvailableError(); - } - const lockKey = `stripe_customer_create_lock:${user.id}`; - const lockToken = await this.cacheService.acquireLock(lockKey, CUSTOMER_LOCK_TTL_SECONDS); - if (!lockToken) { - const freshUser = await this.userRepository.findUnique(user.id); - if (freshUser?.stripeCustomerId) { - return freshUser; - } - throw new StripeError('Failed to acquire customer creation lock'); - } - try { - const freshUser = await this.userRepository.findUnique(user.id); - if (freshUser?.stripeCustomerId) { - return freshUser; - } - const customer = await this.stripe.customers.create({ - email: user.email ?? undefined, - metadata: { - userId: user.id.toString(), - }, - }); - try { - await getBillingRepository().customers.upsertFromStripe(customer, {knownUserId: user.id}); - } catch (mirrorErr) { - Logger.error( - {mirrorErr, customerId: customer.id}, - 'Mirror upsert failed after Stripe write; reconciler will heal', - ); - } - const updatedUser = await this.userRepository.patchUpsert( - user.id, - {stripe_customer_id: customer.id}, - user.toRow(), - ); - Logger.debug({userId: user.id, customerId: customer.id}, 'Stripe customer created for age verification'); - return updatedUser; - } finally { - try { - const released = await this.cacheService.releaseLock(lockKey, lockToken); - if (!released) { - Logger.warn({userId: user.id, lockKey}, 'Customer creation lock token no longer matched on release'); - } - } catch (error) { - Logger.error({error, userId: user.id, lockKey}, 'Failed to release customer creation lock'); - } - } - } } diff --git a/fluxer_api/src/api/stripe/services/StripeCheckoutService.ts b/fluxer_api/src/api/stripe/services/StripeCheckoutService.ts index ab73b3cf0..29ea957a2 100644 --- a/fluxer_api/src/api/stripe/services/StripeCheckoutService.ts +++ b/fluxer_api/src/api/stripe/services/StripeCheckoutService.ts @@ -11,6 +11,7 @@ import type {StoreEntitlementService} from '@app/api/store_billing/StoreEntitlem import {getBillingBranding} from '@app/api/stripe/BillingBranding'; import {getEffectiveBillingConfig, isCurrentCatalogPriceId} from '@app/api/stripe/BillingConfigCache'; import type {ProductInfo, ProductRegistry} from '@app/api/stripe/ProductRegistry'; +import {ensureStripeCustomer, isStripeResourceMissingError} from '@app/api/stripe/StripeCustomer'; import {getCachedStripePriceSummary, type StripePriceSummary} from '@app/api/stripe/StripePriceSummaryCache'; import { canProvisionPremiumFromSubscriptionStatus, @@ -42,15 +43,10 @@ import {UnclaimedAccountCannotMakePurchasesError} from '@fluxer/errors/src/domai import {UnknownUserError} from '@fluxer/errors/src/domains/user/UnknownUserError'; import type {CheckoutPaymentMethod} from '@fluxer/schema/src/domains/premium/GiftCodeSchemas'; import type {ICacheService} from '@pkgs/cache/src/ICacheService'; -import {seconds} from 'itty-time'; import type Stripe from 'stripe'; export const EU_WITHDRAWAL_WAIVER_TEXT_VERSION = '2026-04-23'; -function isStripeResourceMissingError(error: unknown): boolean { - return typeof error === 'object' && error !== null && 'code' in error && error.code === 'resource_missing'; -} - type CheckoutSessionCreateParams = Stripe.Checkout.SessionCreateParams; type CheckoutSessionMode = CheckoutSessionCreateParams['mode']; type CheckoutSessionPaymentMethodType = NonNullable[number]; @@ -653,8 +649,6 @@ export class StripeCheckoutService { } } - private static readonly CUSTOMER_LOCK_TTL_SECONDS = seconds('30 seconds'); - private resolveConfiguredPriceIds(countryCode?: string): ResolvedPriceIds { const recurringCurrencyPreferences = getCurrencyPreferences(countryCode); const giftCurrencyPreferences = getGiftCurrencyPreferences(countryCode); @@ -888,59 +882,14 @@ export class StripeCheckoutService { } private async ensureStripeCustomer(existingUser: User): Promise { - const user = await this.clearStaleStripeCustomer(existingUser); - if (user.stripeCustomerId) { - return user; - } if (!this.stripe) { throw new StripePaymentNotAvailableError(); } - const lockKey = `stripe_customer_create_lock:${user.id}`; - const lockToken = await this.cacheService.acquireLock(lockKey, StripeCheckoutService.CUSTOMER_LOCK_TTL_SECONDS); - if (!lockToken) { - const freshUser = await this.userRepository.findUnique(user.id); - if (freshUser?.stripeCustomerId) { - return freshUser; - } - throw new StripeError('Failed to acquire customer creation lock'); - } - try { - const freshUser = await this.userRepository.findUnique(user.id); - if (freshUser?.stripeCustomerId) { - return freshUser; - } - const customer = await this.stripe.customers.create({ - email: user.email ?? undefined, - metadata: { - userId: user.id.toString(), - }, - }); - try { - await getBillingRepository().customers.upsertFromStripe(customer, {knownUserId: user.id}); - } catch (mirrorErr) { - Logger.error( - {mirrorErr, customerId: customer.id}, - 'Mirror upsert failed after Stripe write; reconciler will heal', - ); - } - const updatedUser = await this.userRepository.patchUpsert( - user.id, - { - stripe_customer_id: customer.id, - }, - user.toRow(), - ); - Logger.debug({userId: user.id, customerId: customer.id}, 'Stripe customer created'); - return updatedUser; - } finally { - try { - const released = await this.cacheService.releaseLock(lockKey, lockToken); - if (!released) { - Logger.warn({userId: user.id, lockKey}, 'Customer creation lock token no longer matched on release'); - } - } catch (error) { - Logger.error({error, userId: user.id, lockKey}, 'Failed to release customer creation lock'); - } - } + return ensureStripeCustomer({ + stripe: this.stripe, + user: await this.clearStaleStripeCustomer(existingUser), + userRepository: this.userRepository, + cacheService: this.cacheService, + }); } } diff --git a/fluxer_api/src/api/stripe/tests/StripeCustomerEmailSync.test.ts b/fluxer_api/src/api/stripe/tests/StripeCustomerEmailSync.test.ts new file mode 100644 index 000000000..dc3e4d39a --- /dev/null +++ b/fluxer_api/src/api/stripe/tests/StripeCustomerEmailSync.test.ts @@ -0,0 +1,233 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later + +import {createTestAccount, type TestAccount} from '@app/api/auth/tests/AuthTestUtils'; +import {createUserID} from '@app/api/BrandedTypes'; +import {Config} from '@app/api/Config'; +import {getBillingRepository} from '@app/api/middleware/ServiceRegistry'; +import {getUserRepository} from '@app/api/middleware/ServiceSingletons'; +import {getStripeClient} from '@app/api/stripe/StripeClient'; +import {type ApiTestHarness, createApiTestHarness} from '@app/api/test/ApiTestHarness'; +import {NoopLogger} from '@app/api/test/mocks/NoopLogger'; +import {createStripeApiHandlers, type StripeApiHandlers} from '@app/api/test/msw/handlers/StripeApiHandlers'; +import {server} from '@app/api/test/msw/server'; +import {HTTP_STATUS} from '@app/api/test/TestConstants'; +import {createBuilder, createBuilderWithoutAuth} from '@app/api/test/TestRequestBuilder'; +import {PaymentRepository} from '@app/api/user/repositories/PaymentRepository'; +import reconcileUserPayments from '@app/api/worker/tasks/ReconcileUserPayments'; +import syncStripeCustomerEmail from '@app/api/worker/tasks/SyncStripeCustomerEmail'; +import {clearWorkerDependencies, setWorkerDependenciesForTest} from '@app/api/worker/WorkerContext'; +import type {WorkerTaskHelpers} from '@pkgs/worker/src/contracts/WorkerTask'; +import type Stripe from 'stripe'; +import {afterAll, afterEach, beforeAll, beforeEach, describe, expect, test} from 'vitest'; + +const MOCK_PRICES = { + monthlyUsd: 'price_email_sync_monthly_usd', + yearlyUsd: 'price_email_sync_yearly_usd', + gift1MonthUsd: 'price_email_sync_gift_1_month_usd', + gift1YearUsd: 'price_email_sync_gift_1_year_usd', +}; + +const MOCK_PRICE_SEEDS = { + [MOCK_PRICES.monthlyUsd]: {unit_amount: 499, currency: 'usd', interval: 'month' as const}, + [MOCK_PRICES.yearlyUsd]: {unit_amount: 4999, currency: 'usd', interval: 'year' as const}, +}; + +const STALE_EMAIL = 'previous-address@example.com'; + +function createHelpers(): WorkerTaskHelpers { + return { + logger: new NoopLogger(), + jobId: 0n, + addJob: async () => 0n, + reportProgress: async () => {}, + shouldCancel: async () => false, + setContextLink: async () => {}, + }; +} + +async function currentEmail(account: TestAccount): Promise { + const user = await getUserRepository().findUnique(createUserID(BigInt(account.userId))); + return user!.email!; +} + +async function mirrorCustomer(customerId: string, account: TestAccount, email: string | null): Promise { + await getBillingRepository().customers.upsertFromStripe( + { + id: customerId, + object: 'customer', + email, + created: 1_600_000_000, + livemode: false, + metadata: {userId: account.userId}, + } as unknown as Stripe.Customer, + {knownUserId: BigInt(account.userId)}, + ); +} + +describe('Stripe customer email sync', () => { + let harness: ApiTestHarness; + let stripeHandlers: StripeApiHandlers; + let originalPrices: typeof Config.stripe.prices | undefined; + + async function createCustomerAccount(customerId: string): Promise { + const account = await createTestAccount(harness); + await createBuilderWithoutAuth(harness) + .post(`/test/users/${account.userId}/security-flags`) + .body({email_verified: true}) + .expect(HTTP_STATUS.OK) + .execute(); + await createBuilderWithoutAuth(harness) + .post(`/test/users/${account.userId}/premium`) + .body({stripe_customer_id: customerId}) + .expect(HTTP_STATUS.OK) + .execute(); + return account; + } + + async function runTask(account: TestAccount): Promise { + await syncStripeCustomerEmail({userId: account.userId}, createHelpers()); + } + + function emailUpdatesFor(customerId: string): Array { + return stripeHandlers.spies.updatedCustomers + .filter((update) => update.id === customerId && 'email' in update.params) + .map((update) => update.params.email); + } + + beforeAll(async () => { + originalPrices = Config.stripe.prices; + Config.stripe.prices = MOCK_PRICES; + harness = await createApiTestHarness(); + }); + afterAll(async () => { + await harness.shutdown(); + Config.stripe.prices = originalPrices; + }); + beforeEach(async () => { + await harness.reset(); + Config.stripe.prices = MOCK_PRICES; + stripeHandlers = createStripeApiHandlers({prices: MOCK_PRICE_SEEDS, subscriptionsListEmpty: true}); + server.use(...stripeHandlers.handlers); + const stripe = getStripeClient(); + expect(stripe).not.toBeNull(); + setWorkerDependenciesForTest({userRepository: getUserRepository(), stripe}); + }); + afterEach(() => { + clearWorkerDependencies(); + }); + + describe('worker task', () => { + test('pushes the account email to a customer that still has the old address', async () => { + const account = await createCustomerAccount('cus_sync_stale'); + await mirrorCustomer('cus_sync_stale', account, STALE_EMAIL); + await runTask(account); + const email = await currentEmail(account); + expect(emailUpdatesFor('cus_sync_stale')).toEqual([email]); + expect((await getBillingRepository().customers.findById('cus_sync_stale'))?.email).toBe(email); + }); + + test('pushes the account email when the customer is not mirrored yet', async () => { + const account = await createCustomerAccount('cus_sync_unmirrored'); + await runTask(account); + expect(emailUpdatesFor('cus_sync_unmirrored')).toEqual([await currentEmail(account)]); + }); + + test('does nothing when the customer already has the account email', async () => { + const account = await createCustomerAccount('cus_sync_current'); + await mirrorCustomer('cus_sync_current', account, await currentEmail(account)); + await runTask(account); + expect(stripeHandlers.spies.updatedCustomers).toHaveLength(0); + }); + + test('does nothing for a customer mirrored as deleted', async () => { + const account = await createCustomerAccount('cus_sync_deleted'); + await getBillingRepository().customers.upsertFromStripe( + {id: 'cus_sync_deleted', object: 'customer', deleted: true} as Stripe.DeletedCustomer, + {knownUserId: BigInt(account.userId)}, + ); + await runTask(account); + expect(stripeHandlers.spies.updatedCustomers).toHaveLength(0); + }); + + test('completes without retrying when Stripe no longer has the customer', async () => { + server.use(...createStripeApiHandlers({customerShouldFail: true}).handlers); + const account = await createCustomerAccount('cus_sync_missing'); + await mirrorCustomer('cus_sync_missing', account, STALE_EMAIL); + await expect(runTask(account)).resolves.toBeUndefined(); + expect((await getBillingRepository().customers.findById('cus_sync_missing'))?.email).toBe(STALE_EMAIL); + }); + + test('does nothing for users without a Stripe customer', async () => { + const account = await createTestAccount(harness); + await runTask(account); + expect(stripeHandlers.spies.updatedCustomers).toHaveLength(0); + }); + }); + + describe('payment reconciliation', () => { + test('heals a customer email that drifted before the account email changed', async () => { + setWorkerDependenciesForTest({ + userRepository: getUserRepository(), + paymentRepository: new PaymentRepository(), + stripe: getStripeClient(), + }); + const account = await createCustomerAccount('cus_sync_reconcile'); + await mirrorCustomer('cus_sync_reconcile', account, STALE_EMAIL); + await reconcileUserPayments({userId: account.userId}, createHelpers()); + expect(emailUpdatesFor('cus_sync_reconcile')).toEqual([await currentEmail(account)]); + }); + + test('leaves a customer alone when the email already matches', async () => { + setWorkerDependenciesForTest({ + userRepository: getUserRepository(), + paymentRepository: new PaymentRepository(), + stripe: getStripeClient(), + }); + const account = await createCustomerAccount('cus_sync_reconcile_current'); + await mirrorCustomer('cus_sync_reconcile_current', account, await currentEmail(account)); + await reconcileUserPayments({userId: account.userId}, createHelpers()); + expect(emailUpdatesFor('cus_sync_reconcile_current')).toHaveLength(0); + }); + }); + + describe('checkout', () => { + test('updates a stale customer email before opening checkout with that customer', async () => { + const account = await createCustomerAccount('cus_checkout_stale'); + await mirrorCustomer('cus_checkout_stale', account, STALE_EMAIL); + await createBuilder(harness, account.token) + .post('/stripe/checkout/subscription') + .body({price_id: MOCK_PRICES.monthlyUsd}) + .expect(HTTP_STATUS.OK) + .execute(); + expect(emailUpdatesFor('cus_checkout_stale')).toEqual([await currentEmail(account)]); + expect(stripeHandlers.spies.createdCheckoutSessions).toHaveLength(1); + expect(stripeHandlers.spies.createdCheckoutSessions[0]?.customer).toBe('cus_checkout_stale'); + }); + + test('leaves the customer alone when its email already matches', async () => { + const account = await createCustomerAccount('cus_checkout_current'); + await mirrorCustomer('cus_checkout_current', account, await currentEmail(account)); + await createBuilder(harness, account.token) + .post('/stripe/checkout/subscription') + .body({price_id: MOCK_PRICES.monthlyUsd}) + .expect(HTTP_STATUS.OK) + .execute(); + expect(emailUpdatesFor('cus_checkout_current')).toHaveLength(0); + expect(stripeHandlers.spies.createdCheckoutSessions[0]?.customer).toBe('cus_checkout_current'); + }); + + test('still opens checkout when the email update fails', async () => { + server.use( + ...createStripeApiHandlers({prices: MOCK_PRICE_SEEDS, subscriptionsListEmpty: true, customerShouldFail: true}) + .handlers, + ); + const account = await createCustomerAccount('cus_checkout_update_fails'); + await mirrorCustomer('cus_checkout_update_fails', account, STALE_EMAIL); + await createBuilder(harness, account.token) + .post('/stripe/checkout/subscription') + .body({price_id: MOCK_PRICES.monthlyUsd}) + .expect(HTTP_STATUS.OK) + .execute(); + }); + }); +}); diff --git a/fluxer_api/src/api/test/msw/handlers/StripeApiHandlers.ts b/fluxer_api/src/api/test/msw/handlers/StripeApiHandlers.ts index 12728763b..d1448cba9 100644 --- a/fluxer_api/src/api/test/msw/handlers/StripeApiHandlers.ts +++ b/fluxer_api/src/api/test/msw/handlers/StripeApiHandlers.ts @@ -1565,6 +1565,19 @@ export function createStripeApiHandlers(config: StripeApiMockConfig = {}): Strip const formData = await request.formData(); const updateParams = parseFormDataToObject(formData); spies.updatedCustomers.push({id: customerId, params: updateParams}); + if (config.customerShouldFail) { + return HttpResponse.json( + { + error: { + type: 'invalid_request_error', + message: 'No such customer', + code: 'resource_missing', + param: 'id', + }, + }, + {status: 404}, + ); + } const customer = getCustomer(customerId); const invoiceSettingsUpdate = updateParams.invoice_settings && typeof updateParams.invoice_settings === 'object' @@ -1572,6 +1585,7 @@ export function createStripeApiHandlers(config: StripeApiMockConfig = {}): Strip : null; const updatedCustomer: MockStripeCustomer = { ...customer, + email: typeof updateParams.email === 'string' ? updateParams.email : customer.email, invoice_settings: invoiceSettingsUpdate ? { ...customer.invoice_settings, diff --git a/fluxer_api/src/api/user/services/EmailChangeService.ts b/fluxer_api/src/api/user/services/EmailChangeService.ts index cb1029bc2..d653931ec 100644 --- a/fluxer_api/src/api/user/services/EmailChangeService.ts +++ b/fluxer_api/src/api/user/services/EmailChangeService.ts @@ -5,6 +5,7 @@ import type {ApiContext} from '@app/api/ApiContext'; import * as AuthPassword from '@app/api/auth/AuthPassword'; import {assertEmailNotBlocklisted} from '@app/api/auth/EmailBlocklist'; import type {User} from '@app/api/models/User'; +import {enqueueStripeCustomerEmailSync} from '@app/api/stripe/StripeCustomer'; import type {EmailChangeRepository} from '@app/api/user/repositories/auth/EmailChangeRepository'; import { assertChangeCooldown, @@ -300,7 +301,7 @@ export class EmailChangeService { } async verifyBouncedNew(user: User, ticket: string, code: string): Promise { - const {users} = this.apiContext.services; + const {users, worker} = this.apiContext.services; this.ensureBouncedEmailRecoveryAllowed(user); const row = await getActiveChangeTicketForUser(this.repo, ticket, user.id); if (row.require_original || !row.original_proof) { @@ -314,6 +315,7 @@ export class EmailChangeService { user.toRow(), ); await this.deleteToken(emailToken); + await enqueueStripeCustomerEmailSync(worker, user, updatedUser); return updatedUser; } diff --git a/fluxer_api/src/api/user/services/UserAccountService.ts b/fluxer_api/src/api/user/services/UserAccountService.ts index fe9095701..7b14188dc 100644 --- a/fluxer_api/src/api/user/services/UserAccountService.ts +++ b/fluxer_api/src/api/user/services/UserAccountService.ts @@ -14,6 +14,7 @@ import {Logger} from '@app/api/Logger'; import type {LimitConfigService} from '@app/api/limits/LimitConfigService'; import type {AuthSession} from '@app/api/models/AuthSession'; import type {User} from '@app/api/models/User'; +import {enqueueStripeCustomerEmailSync} from '@app/api/stripe/StripeCustomer'; import type {IUserAccountRepository} from '@app/api/user/repositories/IUserAccountRepository'; import type {IUserChannelRepository} from '@app/api/user/repositories/IUserChannelRepository'; import type {IUserRelationshipRepository} from '@app/api/user/repositories/IUserRelationshipRepository'; @@ -182,6 +183,7 @@ export class UserAccountService { reason: 'user_requested', actorUserId: user.id, }), + () => enqueueStripeCustomerEmailSync(this.apiContext.services.worker, user, updatedUser), async () => { try { await this.profileService.commitAssetChanges(profileResult); diff --git a/fluxer_api/src/api/worker/WorkerLaneConfig.ts b/fluxer_api/src/api/worker/WorkerLaneConfig.ts index c9bfff2db..83b323621 100644 --- a/fluxer_api/src/api/worker/WorkerLaneConfig.ts +++ b/fluxer_api/src/api/worker/WorkerLaneConfig.ts @@ -49,6 +49,7 @@ const LANE_CONFIG = { 'harvestUserData', 'batchGuildAuditLogMessageDeletes', 'reconcileUserPayments', + 'syncStripeCustomerEmail', 'processAppStoreNotification', 'processGooglePlayNotification', 'refreshStorePurchase', diff --git a/fluxer_api/src/api/worker/WorkerTaskRegistry.ts b/fluxer_api/src/api/worker/WorkerTaskRegistry.ts index 8ed99dc12..43d5ba634 100644 --- a/fluxer_api/src/api/worker/WorkerTaskRegistry.ts +++ b/fluxer_api/src/api/worker/WorkerTaskRegistry.ts @@ -49,6 +49,7 @@ import syncCrosspostCopies from '@app/api/worker/tasks/SyncCrosspostCopies'; import syncCrosspostedMessage from '@app/api/worker/tasks/SyncCrosspostedMessage'; import syncDiscoveryIndex from '@app/api/worker/tasks/SyncDiscoveryIndex'; import syncFileShaBlocklists from '@app/api/worker/tasks/SyncFileShaBlocklists'; +import syncStripeCustomerEmail from '@app/api/worker/tasks/SyncStripeCustomerEmail'; import syncUrlBlocklists from '@app/api/worker/tasks/SyncUrlBlocklists'; import userProcessPendingDeletion from '@app/api/worker/tasks/UserProcessPendingDeletion'; import userProcessPendingDeletions from '@app/api/worker/tasks/UserProcessPendingDeletions'; @@ -100,6 +101,7 @@ export const workerTasks: Record = { refreshSearchIndex, removeChannelFollowers, sendSystemDm, + syncStripeCustomerEmail, syncFileShaBlocklists, syncUrlBlocklists, syncDiscoveryIndex, diff --git a/fluxer_api/src/api/worker/tasks/ReconcileUserPayments.ts b/fluxer_api/src/api/worker/tasks/ReconcileUserPayments.ts index 65092550c..600419ba8 100644 --- a/fluxer_api/src/api/worker/tasks/ReconcileUserPayments.ts +++ b/fluxer_api/src/api/worker/tasks/ReconcileUserPayments.ts @@ -6,6 +6,7 @@ import {mapGiftDurationMonthsToFields} from '@app/api/models/GiftCode'; import type {Payment} from '@app/api/models/Payment'; import type {User} from '@app/api/models/User'; import {getProductRegistry} from '@app/api/stripe/ProductRegistry'; +import {syncStripeCustomerEmail} from '@app/api/stripe/StripeCustomer'; import {extractId} from '@app/api/stripe/StripeUtils'; import type {PaymentRepository} from '@app/api/user/repositories/PaymentRepository'; import {mapUserToPrivateResponse} from '@app/api/user/UserMappers'; @@ -332,6 +333,11 @@ const reconcileUserPayments: WorkerTaskHandler = async (payload, helpers) => { if (!user.stripeCustomerId) { return; } + try { + await syncStripeCustomerEmail(stripe, user); + } catch (error) { + Logger.warn({userId: userIdStr, error}, 'Failed to sync Stripe customer email during payment reconciliation'); + } const payments = await paymentRepository.findPaymentsByUserId(userId); let reconciledCount = 0; for (const payment of payments) { diff --git a/fluxer_api/src/api/worker/tasks/SyncStripeCustomerEmail.ts b/fluxer_api/src/api/worker/tasks/SyncStripeCustomerEmail.ts new file mode 100644 index 000000000..9027c2182 --- /dev/null +++ b/fluxer_api/src/api/worker/tasks/SyncStripeCustomerEmail.ts @@ -0,0 +1,37 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later + +import {createUserID} from '@app/api/BrandedTypes'; +import {isStripeResourceMissingError, syncStripeCustomerEmail} from '@app/api/stripe/StripeCustomer'; +import {getWorkerDependencies} from '@app/api/worker/WorkerContext'; +import type {WorkerTaskHandler} from '@pkgs/worker/src/contracts/WorkerTask'; + +const syncStripeCustomerEmailTask: WorkerTaskHandler = async (payload, helpers) => { + const {userRepository, stripe} = getWorkerDependencies(); + if (!stripe) { + helpers.logger.debug('Stripe is disabled, skipping customer email sync'); + return; + } + const userIdStr = payload.userId as string; + if (!userIdStr) { + helpers.logger.warn({payload}, 'Stripe customer email sync task missing userId'); + return; + } + const user = await userRepository.findUnique(createUserID(BigInt(userIdStr))); + if (!user) { + return; + } + try { + await syncStripeCustomerEmail(stripe, user); + } catch (error) { + if (isStripeResourceMissingError(error)) { + helpers.logger.warn( + {userId: userIdStr, customerId: user.stripeCustomerId}, + 'Stripe customer no longer exists, skipping email sync', + ); + return; + } + throw error; + } +}; + +export default syncStripeCustomerEmailTask;