Fix slow db query issue (#19770)

https://github.com/twentyhq/twenty/pull/19586#discussion_r3074136617
This commit is contained in:
Etienne
2026-04-16 17:37:58 +02:00
committed by GitHub
parent 446a3923f2
commit aecbc89a3f
3 changed files with 94 additions and 83 deletions
@@ -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'),
@@ -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<BillingSubscriptionEntity>,
@InjectRepository(BillingSubscriptionItemEntity)
private readonly billingSubscriptionItemRepository: Repository<BillingSubscriptionItemEntity>,
@InjectRepository(BillingProductEntity)
private readonly billingProductRepository: Repository<BillingProductEntity>,
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<BillingSubscriptionEntity, 'id'>[];
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) {
@@ -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,