From aecbc89a3fa8f9386d4f2f34d9529133f418220b Mon Sep 17 00:00:00 2001 From: Etienne <45695613+etiennejouan@users.noreply.github.com> Date: Thu, 16 Apr 2026 17:37:58 +0200 Subject: [PATCH] Fix slow db query issue (#19770) https://github.com/twentyhq/twenty/pull/19586#discussion_r3074136617 --- .../__test__/enforce-usage-cap.job.spec.ts | 96 +++++++++---------- .../billing/crons/enforce-usage-cap.job.ts | 75 +++++++++------ .../core-modules/message-queue/jobs.module.ts | 6 +- 3 files changed, 94 insertions(+), 83 deletions(-) diff --git a/packages/twenty-server/src/engine/core-modules/billing/crons/__test__/enforce-usage-cap.job.spec.ts b/packages/twenty-server/src/engine/core-modules/billing/crons/__test__/enforce-usage-cap.job.spec.ts index 4f37f1ac1ce..87e4c992b62 100644 --- a/packages/twenty-server/src/engine/core-modules/billing/crons/__test__/enforce-usage-cap.job.spec.ts +++ b/packages/twenty-server/src/engine/core-modules/billing/crons/__test__/enforce-usage-cap.job.spec.ts @@ -6,16 +6,19 @@ import { getRepositoryToken } from '@nestjs/typeorm'; import { In } from 'typeorm'; import { EnforceUsageCapJob } from 'src/engine/core-modules/billing/crons/enforce-usage-cap.job'; +import { BillingProductEntity } from 'src/engine/core-modules/billing/entities/billing-product.entity'; import { BillingSubscriptionItemEntity } from 'src/engine/core-modules/billing/entities/billing-subscription-item.entity'; import { BillingSubscriptionEntity } from 'src/engine/core-modules/billing/entities/billing-subscription.entity'; import { BillingProductKey } from 'src/engine/core-modules/billing/enums/billing-product-key.enum'; import { BillingUsageCapService } from 'src/engine/core-modules/billing/services/billing-usage-cap.service'; import { TwentyConfigService } from 'src/engine/core-modules/twenty-config/twenty-config.service'; +const METERED_STRIPE_PRODUCT_ID = 'prod_metered'; +const METERED_STRIPE_PRICE_ID = 'price_metered'; + describe('EnforceUsageCapJob', () => { let job: EnforceUsageCapJob; - let idQueryMock: jest.Mock; - let fullQueryMock: jest.Mock; + let billingSubscriptionFindMock: jest.Mock; let billingSubscriptionItemRepository: jest.Mocked<{ update: jest.Mock; }>; @@ -44,54 +47,40 @@ describe('EnforceUsageCapJob', () => { { id: itemId, hasReachedCurrentPeriodCap, - billingProduct: { - metadata: { - productKey: BillingProductKey.WORKFLOW_NODE_EXECUTION, - }, - }, + stripeProductId: METERED_STRIPE_PRODUCT_ID, + stripePriceId: METERED_STRIPE_PRICE_ID, }, ], }) as unknown as BillingSubscriptionEntity; beforeEach(async () => { - idQueryMock = jest.fn().mockResolvedValue([]); - fullQueryMock = jest.fn().mockResolvedValue([]); - - const idQueryBuilderMock = { - select: jest.fn().mockReturnThis(), - innerJoin: jest.fn().mockReturnThis(), - where: jest.fn().mockReturnThis(), - andWhere: jest.fn().mockReturnThis(), - orderBy: jest.fn().mockReturnThis(), - limit: jest.fn().mockReturnThis(), - offset: jest.fn().mockReturnThis(), - getRawMany: idQueryMock, - }; - - const fullQueryBuilderMock = { - innerJoinAndSelect: jest.fn().mockReturnThis(), - leftJoinAndSelect: jest.fn().mockReturnThis(), - where: jest.fn().mockReturnThis(), - orderBy: jest.fn().mockReturnThis(), - getMany: fullQueryMock, - }; - - const createQueryBuilder = jest - .fn() - .mockReturnValueOnce(idQueryBuilderMock) - .mockReturnValueOnce(fullQueryBuilderMock); + billingSubscriptionFindMock = jest.fn().mockResolvedValue([]); const module: TestingModule = await Test.createTestingModule({ providers: [ EnforceUsageCapJob, { provide: getRepositoryToken(BillingSubscriptionEntity), - useValue: { createQueryBuilder }, + useValue: { find: billingSubscriptionFindMock }, }, { provide: getRepositoryToken(BillingSubscriptionItemEntity), useValue: { update: jest.fn() }, }, + { + provide: getRepositoryToken(BillingProductEntity), + useValue: { + find: jest.fn().mockResolvedValue([ + { + stripeProductId: METERED_STRIPE_PRODUCT_ID, + metadata: { + productKey: BillingProductKey.WORKFLOW_NODE_EXECUTION, + }, + billingPrices: [], + }, + ]), + }, + }, { provide: BillingUsageCapService, useValue: { @@ -136,7 +125,7 @@ describe('EnforceUsageCapJob', () => { await job.handle(); - expect(idQueryMock).not.toHaveBeenCalled(); + expect(billingSubscriptionFindMock).not.toHaveBeenCalled(); }); it('no-ops when ClickHouse is not configured', async () => { @@ -145,15 +134,16 @@ describe('EnforceUsageCapJob', () => { await job.handle(); - expect(idQueryMock).not.toHaveBeenCalled(); + expect(billingSubscriptionFindMock).not.toHaveBeenCalled(); }); it('skips transitions in shadow mode (flag off)', async () => { mockConfig({ BILLING_USAGE_CAP_CLICKHOUSE_ENABLED: false }); const sub = buildSubscription({ hasReachedCurrentPeriodCap: false }); - idQueryMock.mockResolvedValueOnce([{ subscription_id: 'sub_123' }]); - fullQueryMock.mockResolvedValueOnce([sub]); + billingSubscriptionFindMock + .mockResolvedValueOnce([{ id: 'sub_123' }]) + .mockResolvedValueOnce([sub]); billingUsageCapService.getBatchPeriodCreditsUsed.mockResolvedValue( new Map([['workspace_123', 2_000_000]]), @@ -186,8 +176,9 @@ describe('EnforceUsageCapJob', () => { hasReachedCurrentPeriodCap: false, }); - idQueryMock.mockResolvedValueOnce([{ subscription_id: 'sub_123' }]); - fullQueryMock.mockResolvedValueOnce([sub]); + billingSubscriptionFindMock + .mockResolvedValueOnce([{ id: 'sub_123' }]) + .mockResolvedValueOnce([sub]); billingUsageCapService.getBatchPeriodCreditsUsed.mockResolvedValue( new Map([['workspace_123', 2_000_000]]), @@ -223,8 +214,9 @@ describe('EnforceUsageCapJob', () => { hasReachedCurrentPeriodCap: true, }); - idQueryMock.mockResolvedValueOnce([{ subscription_id: 'sub_123' }]); - fullQueryMock.mockResolvedValueOnce([sub]); + billingSubscriptionFindMock + .mockResolvedValueOnce([{ id: 'sub_123' }]) + .mockResolvedValueOnce([sub]); billingUsageCapService.getBatchPeriodCreditsUsed.mockResolvedValue( new Map([['workspace_123', 500_000]]), @@ -257,8 +249,9 @@ describe('EnforceUsageCapJob', () => { mockConfig({ BILLING_USAGE_CAP_CLICKHOUSE_ENABLED: true }); const sub = buildSubscription({ hasReachedCurrentPeriodCap: true }); - idQueryMock.mockResolvedValueOnce([{ subscription_id: 'sub_123' }]); - fullQueryMock.mockResolvedValueOnce([sub]); + billingSubscriptionFindMock + .mockResolvedValueOnce([{ id: 'sub_123' }]) + .mockResolvedValueOnce([sub]); billingUsageCapService.getBatchPeriodCreditsUsed.mockResolvedValue( new Map([['workspace_123', 2_000_000]]), @@ -288,8 +281,9 @@ describe('EnforceUsageCapJob', () => { mockConfig({ BILLING_USAGE_CAP_CLICKHOUSE_ENABLED: true }); const sub = buildSubscription(); - idQueryMock.mockResolvedValueOnce([{ subscription_id: 'sub_123' }]); - fullQueryMock.mockResolvedValueOnce([sub]); + billingSubscriptionFindMock + .mockResolvedValueOnce([{ id: 'sub_123' }]) + .mockResolvedValueOnce([sub]); billingUsageCapService.getBatchPeriodCreditsUsed.mockResolvedValue( new Map(), @@ -315,8 +309,9 @@ describe('EnforceUsageCapJob', () => { creditBalanceMicro: 300_000, }); - idQueryMock.mockResolvedValueOnce([{ subscription_id: 'sub_123' }]); - fullQueryMock.mockResolvedValueOnce([sub]); + billingSubscriptionFindMock + .mockResolvedValueOnce([{ id: 'sub_123' }]) + .mockResolvedValueOnce([sub]); billingUsageCapService.getBatchPeriodCreditsUsed.mockResolvedValue( new Map(), @@ -336,8 +331,9 @@ describe('EnforceUsageCapJob', () => { mockConfig({ BILLING_USAGE_CAP_CLICKHOUSE_ENABLED: true }); const sub = buildSubscription({ hasReachedCurrentPeriodCap: true }); - idQueryMock.mockResolvedValueOnce([{ subscription_id: 'sub_123' }]); - fullQueryMock.mockResolvedValueOnce([sub]); + billingSubscriptionFindMock + .mockResolvedValueOnce([{ id: 'sub_123' }]) + .mockResolvedValueOnce([sub]); billingUsageCapService.getBatchPeriodCreditsUsed.mockRejectedValue( new Error('clickhouse exploded'), diff --git a/packages/twenty-server/src/engine/core-modules/billing/crons/enforce-usage-cap.job.ts b/packages/twenty-server/src/engine/core-modules/billing/crons/enforce-usage-cap.job.ts index 53ece0ed54a..0b5bfd5ab9b 100644 --- a/packages/twenty-server/src/engine/core-modules/billing/crons/enforce-usage-cap.job.ts +++ b/packages/twenty-server/src/engine/core-modules/billing/crons/enforce-usage-cap.job.ts @@ -3,9 +3,10 @@ import { Injectable, Logger } from '@nestjs/common'; import { InjectRepository } from '@nestjs/typeorm'; -import { In, Repository } from 'typeorm'; +import { In, IsNull, Repository } from 'typeorm'; import { enforceUsageCapCronPattern } from 'src/engine/core-modules/billing/crons/enforce-usage-cap.cron.pattern'; +import { BillingProductEntity } from 'src/engine/core-modules/billing/entities/billing-product.entity'; import { BillingSubscriptionItemEntity } from 'src/engine/core-modules/billing/entities/billing-subscription-item.entity'; import { BillingSubscriptionEntity } from 'src/engine/core-modules/billing/entities/billing-subscription.entity'; import { BillingProductKey } from 'src/engine/core-modules/billing/enums/billing-product-key.enum'; @@ -29,6 +30,8 @@ export class EnforceUsageCapJob { private readonly billingSubscriptionRepository: Repository, @InjectRepository(BillingSubscriptionItemEntity) private readonly billingSubscriptionItemRepository: Repository, + @InjectRepository(BillingProductEntity) + private readonly billingProductRepository: Repository, private readonly billingUsageCapService: BillingUsageCapService, private readonly twentyConfigService: TwentyConfigService, ) {} @@ -57,46 +60,48 @@ export class EnforceUsageCapJob { let errors = 0; let offset = 0; + const allProducts = await this.billingProductRepository.find({ + relations: { billingPrices: true }, + }); + const productByStripeProductId = new Map( + allProducts.map((product) => [product.stripeProductId, product]), + ); + let batch: BillingSubscriptionEntity[]; - let idRows: { subscription_id: string }[]; + let idRows: Pick[]; do { - idRows = await this.billingSubscriptionRepository - .createQueryBuilder('subscription') - .select('subscription.id') - .innerJoin('subscription.workspace', 'workspace') - .where('subscription.status IN (:...statuses)', { - statuses: [ + idRows = await this.billingSubscriptionRepository.find({ + select: { id: true }, + relations: { + workspace: true, + }, + where: { + status: In([ SubscriptionStatus.Active, SubscriptionStatus.Trialing, SubscriptionStatus.PastDue, - ], - }) - .andWhere('workspace.suspendedAt IS NULL') - .orderBy('subscription.id', 'ASC') - .limit(BATCH_SIZE) - .offset(offset) - .getRawMany<{ subscription_id: string }>(); + ]), + workspace: { suspendedAt: IsNull() }, + }, + order: { id: 'ASC' }, + take: BATCH_SIZE, + skip: offset, + }); if (idRows.length === 0) { break; } - const ids = idRows.map((row) => row.subscription_id); + const ids = idRows.map((row) => row.id); - batch = await this.billingSubscriptionRepository - .createQueryBuilder('subscription') - .innerJoinAndSelect( - 'subscription.billingSubscriptionItems', - 'item', - 'item.billingSubscriptionId = subscription.id', - ) - .innerJoinAndSelect('item.billingProduct', 'product') - .leftJoinAndSelect('product.billingPrices', 'price') - .innerJoinAndSelect('subscription.billingCustomer', 'customer') - .where('subscription.id IN (:...ids)', { ids }) - .orderBy('subscription.id', 'ASC') - .getMany(); + batch = await this.billingSubscriptionRepository.find({ + where: { id: In(ids) }, + relations: { + billingCustomer: true, + billingSubscriptionItems: true, + }, + }); if (batch.length === 0) { break; @@ -140,6 +145,14 @@ export class EnforceUsageCapJob { subscription.billingCustomer.creditBalanceMicro, ); } + + for (const item of subscription.billingSubscriptionItems ?? []) { + const product = productByStripeProductId.get(item.stripeProductId); + + if (product) { + item.billingProduct = product; + } + } } const evaluations = this.billingUsageCapService.evaluateCapBatch( @@ -166,8 +179,8 @@ export class EnforceUsageCapJob { const meteredItem = subscription.billingSubscriptionItems.find( (item) => - item.billingProduct?.metadata?.productKey === - BillingProductKey.WORKFLOW_NODE_EXECUTION, + productByStripeProductId.get(item.stripeProductId)?.metadata + ?.productKey === BillingProductKey.WORKFLOW_NODE_EXECUTION, ); if (!meteredItem) { diff --git a/packages/twenty-server/src/engine/core-modules/message-queue/jobs.module.ts b/packages/twenty-server/src/engine/core-modules/message-queue/jobs.module.ts index 2ef8f03a95d..5713e199ab1 100644 --- a/packages/twenty-server/src/engine/core-modules/message-queue/jobs.module.ts +++ b/packages/twenty-server/src/engine/core-modules/message-queue/jobs.module.ts @@ -7,6 +7,7 @@ import { AuditJobModule } from 'src/engine/core-modules/audit/jobs/audit-job.mod import { AuthModule } from 'src/engine/core-modules/auth/auth.module'; import { BillingModule } from 'src/engine/core-modules/billing/billing.module'; import { EnforceUsageCapJob } from 'src/engine/core-modules/billing/crons/enforce-usage-cap.job'; +import { BillingProductEntity } from 'src/engine/core-modules/billing/entities/billing-product.entity'; import { BillingSubscriptionItemEntity } from 'src/engine/core-modules/billing/entities/billing-subscription-item.entity'; import { BillingSubscriptionEntity } from 'src/engine/core-modules/billing/entities/billing-subscription.entity'; import { UpdateSubscriptionQuantityJob } from 'src/engine/core-modules/billing/jobs/update-subscription-quantity.job'; @@ -14,6 +15,8 @@ import { StripeModule } from 'src/engine/core-modules/billing/stripe/stripe.modu import { EmailSenderJob } from 'src/engine/core-modules/email/email-sender.job'; import { EmailModule } from 'src/engine/core-modules/email/email.module'; import { EnterpriseModule } from 'src/engine/core-modules/enterprise/enterprise.module'; +import { GenerateSdkClientJob } from 'src/engine/core-modules/sdk-client/jobs/generate-sdk-client.job'; +import { SdkClientModule } from 'src/engine/core-modules/sdk-client/sdk-client.module'; import { UserWorkspaceModule } from 'src/engine/core-modules/user-workspace/user-workspace.module'; import { UpdateWorkspaceMemberEmailJob } from 'src/engine/core-modules/user/jobs/update-workspace-member-email.job'; import { UserVarsModule } from 'src/engine/core-modules/user/user-vars/user-vars.module'; @@ -30,8 +33,6 @@ import { WebhookJobModule } from 'src/engine/metadata-modules/webhook/jobs/webho import { SubscriptionsModule } from 'src/engine/subscriptions/subscriptions.module'; import { CleanOnboardingWorkspacesJob } from 'src/engine/workspace-manager/workspace-cleaner/crons/clean-onboarding-workspaces.job'; import { CleanSuspendedWorkspacesJob } from 'src/engine/workspace-manager/workspace-cleaner/crons/clean-suspended-workspaces.job'; -import { GenerateSdkClientJob } from 'src/engine/core-modules/sdk-client/jobs/generate-sdk-client.job'; -import { SdkClientModule } from 'src/engine/core-modules/sdk-client/sdk-client.module'; import { CleanWorkspaceDeletionWarningUserVarsJob } from 'src/engine/workspace-manager/workspace-cleaner/jobs/clean-workspace-deletion-warning-user-vars.job'; import { WorkspaceCleanerModule } from 'src/engine/workspace-manager/workspace-cleaner/workspace-cleaner.module'; import { CalendarEventParticipantManagerModule } from 'src/modules/calendar/calendar-event-participant-manager/calendar-event-participant-manager.module'; @@ -48,6 +49,7 @@ import { WorkflowModule } from 'src/modules/workflow/workflow.module'; WorkspaceEntity, BillingSubscriptionEntity, BillingSubscriptionItemEntity, + BillingProductEntity, ]), ObjectMetadataModule, TypeORMModule,