From e4e1d24731f9fddd385a64c09d7b9bdc8295ab3d Mon Sep 17 00:00:00 2001 From: Weiko Date: Thu, 30 Jul 2026 14:34:27 +0200 Subject: [PATCH] Prevent overlapping workspace cleanup executions (#23522) ## Context The suspended-workspace cleanup is a long-running scheduled job. Under database or cache pressure, BullMQ can consider an execution stalled and start a replacement on another worker while the original execution is still running. Both executions can then enumerate the same suspended workspaces and run destructive cleanup concurrently. This amplifies the initial slowdown: 1. Multiple cleanup transactions target the same workspace data. 2. Transactions wait on each other's locks. 3. Database connections remain occupied while waiting. 4. Other workers and API requests have fewer connections available. There is a second source of unnecessary lock duration in workspace deletion. The deletion transaction currently starts before field metadata is read from the workspace cache. If that lookup is slow, the transaction stays open during an unrelated cache wait. ## What changed ### Prevent overlapping scheduled cleanups - Acquire a non-blocking PostgreSQL advisory lock before listing suspended workspaces. - Skip the execution when another worker already holds the lock. - Keep the lock on one dedicated PostgreSQL session for the full callback. - Release the lock in all normal and error paths. - Discard the database connection if lock acquisition or release has an ambiguous failure, preventing a session that may still own the lock from returning to the pool. - Encapsulate this lifecycle in `PostgresAdvisoryLockService`, exported by `TypeORMModule`, so other coarse-grained jobs can reuse it without handling acquisition and release themselves. ### Shorten the workspace deletion transaction - Read field metadata and build deletion chunks before starting the transaction. - Pass the precomputed chunks into the transactional deletion loop. - Keep the existing deletion order and SQL behavior unchanged. ## Why a PostgreSQL advisory lock The lock needs to coordinate workers running in different pods. A PostgreSQL session advisory lock provides the required behavior: - It is shared across all workers using the same database. - Acquisition is non-blocking, a duplicate execution can exit immediately. - It has no TTL or renewal heartbeat that could expire during the same event-loop stall that caused BullMQ to recover the job. - PostgreSQL automatically releases it when the owning session or process disappears. This is deliberately scoped to `CleanSuspendedWorkspacesJob`. It prevents overlapping scheduled executions, but it is not an exactly-once mechanism or a global mutex around every workspace-deletion entry point. ## Expected impact - Prevent one slow cleanup execution from becoming several concurrent cleanup executions. - Reduce database lock contention and connection-pool pressure during cleanup. - Avoid holding deletion transaction locks while waiting for workspace-cache data. - Reduce cleanup-related API latency bursts without changing normal cleanup semantics. The advisory lock holds one core database connection for the duration of the scheduled cleanup. This is intentional and bounded to the single lock owner. ## Validation - Focused advisory-lock tests cover successful execution, contention, callback failure, and unsafe connection disposal when unlock fails. - Cleanup-job tests cover both the lock-owner and skipped-execution paths. - Workspace-service coverage verifies that field metadata is loaded before the deletion transaction starts. - `yarn nx typecheck twenty-server` - Oxlint, Prettier, and Oxfmt checks on the changed files --- .../postgres-advisory-lock.service.spec.ts | 212 ++++++++++++++++++ .../typeorm/postgres-advisory-lock.service.ts | 184 +++++++++++++++ .../src/database/typeorm/typeorm.module.ts | 9 +- .../__tests__/workspace.service.spec.ts | 53 ++++- .../workspace/services/workspace.service.ts | 26 ++- .../clean-suspended-workspaces.job.spec.ts | 74 ++++++ .../crons/clean-suspended-workspaces.job.ts | 38 +++- 7 files changed, 566 insertions(+), 30 deletions(-) create mode 100644 packages/twenty-server/src/database/typeorm/postgres-advisory-lock.service.spec.ts create mode 100644 packages/twenty-server/src/database/typeorm/postgres-advisory-lock.service.ts create mode 100644 packages/twenty-server/src/engine/workspace-manager/workspace-cleaner/crons/__tests__/clean-suspended-workspaces.job.spec.ts diff --git a/packages/twenty-server/src/database/typeorm/postgres-advisory-lock.service.spec.ts b/packages/twenty-server/src/database/typeorm/postgres-advisory-lock.service.spec.ts new file mode 100644 index 0000000000..488ff620c5 --- /dev/null +++ b/packages/twenty-server/src/database/typeorm/postgres-advisory-lock.service.spec.ts @@ -0,0 +1,212 @@ +import { Logger } from '@nestjs/common'; + +import { type PoolClient } from 'pg'; +import { type DataSource } from 'typeorm'; +import { type PostgresDriver } from 'typeorm/driver/postgres/PostgresDriver'; + +import { PostgresAdvisoryLockService } from 'src/database/typeorm/postgres-advisory-lock.service'; + +describe('PostgresAdvisoryLockService', () => { + const obtainMasterConnection = jest.fn(); + const dataSource = { + driver: { + obtainMasterConnection, + } as unknown as PostgresDriver, + } as unknown as DataSource; + + const createMockConnection = (isLockAcquired: boolean) => { + let connectionErrorHandler: ((error: Error) => void) | undefined; + const query = jest + .fn() + .mockResolvedValueOnce({ + rows: [{ acquired: isLockAcquired }], + }) + .mockResolvedValueOnce({ + rows: [{ released: true }], + }); + const connection = { + query, + on: jest.fn((_eventName: string, handler: (error: Error) => void) => { + connectionErrorHandler = handler; + }), + removeListener: jest.fn(), + } as unknown as PoolClient; + const release = jest.fn(); + + return { + connection, + emitConnectionError: (error: Error) => connectionErrorHandler?.(error), + query, + release, + }; + }; + + let loggerWarn: jest.SpyInstance; + + beforeEach(() => { + jest.clearAllMocks(); + loggerWarn = jest.spyOn(Logger.prototype, 'warn').mockImplementation(); + }); + + afterEach(() => { + loggerWarn.mockRestore(); + }); + + it('runs the callback while holding the lock', async () => { + const { connection, query, release } = createMockConnection(true); + const callback = jest.fn().mockResolvedValue('result'); + + obtainMasterConnection.mockResolvedValue([connection, release]); + + const result = await new PostgresAdvisoryLockService( + dataSource, + ).tryWithLock('lock-name', callback); + + expect(result).toEqual({ + acquired: true, + value: 'result', + }); + expect(callback).toHaveBeenCalledTimes(1); + expect(query).toHaveBeenCalledTimes(2); + expect(query).toHaveBeenNthCalledWith( + 2, + expect.stringContaining('pg_advisory_unlock'), + ['lock-name'], + ); + expect(query.mock.invocationCallOrder[1]).toBeLessThan( + release.mock.invocationCallOrder[0], + ); + }); + + it('does not run the callback when the lock is held', async () => { + const { connection, query, release } = createMockConnection(false); + const callback = jest.fn(); + + obtainMasterConnection.mockResolvedValue([connection, release]); + + await expect( + new PostgresAdvisoryLockService(dataSource).tryWithLock( + 'lock-name', + callback, + ), + ).resolves.toEqual({ + acquired: false, + }); + + expect(callback).not.toHaveBeenCalled(); + expect(query).toHaveBeenCalledTimes(1); + expect(release).toHaveBeenCalledTimes(1); + }); + + it('discards the connection when acquiring the lock fails', async () => { + const { connection, query, release } = createMockConnection(true); + const callback = jest.fn(); + const acquisitionError = new Error('acquisition failed'); + + query.mockReset().mockRejectedValueOnce(acquisitionError); + obtainMasterConnection.mockResolvedValue([connection, release]); + + await expect( + new PostgresAdvisoryLockService(dataSource).tryWithLock( + 'lock-name', + callback, + ), + ).rejects.toBe(acquisitionError); + + expect(callback).not.toHaveBeenCalled(); + expect(release).toHaveBeenCalledWith(acquisitionError); + }); + + it('releases the lock when the callback fails', async () => { + const { connection, query, release } = createMockConnection(true); + const callbackError = new Error('callback failed'); + + obtainMasterConnection.mockResolvedValue([connection, release]); + + await expect( + new PostgresAdvisoryLockService(dataSource).tryWithLock( + 'lock-name', + async () => { + throw callbackError; + }, + ), + ).rejects.toBe(callbackError); + + expect(query).toHaveBeenCalledTimes(2); + expect(release).toHaveBeenCalledWith(); + }); + + it('discards the connection when releasing the lock fails', async () => { + const { connection, query, release } = createMockConnection(true); + const releaseError = new Error('release failed'); + + query + .mockReset() + .mockResolvedValueOnce({ + rows: [{ acquired: true }], + }) + .mockRejectedValueOnce(releaseError); + obtainMasterConnection.mockResolvedValue([connection, release]); + + await expect( + new PostgresAdvisoryLockService(dataSource).tryWithLock( + 'lock-name', + async () => undefined, + ), + ).rejects.toBe(releaseError); + + expect(release).toHaveBeenCalledWith(releaseError); + }); + + it('preserves the callback error when the connection fails', async () => { + const { connection, emitConnectionError, release } = + createMockConnection(true); + const callbackError = new Error('callback failed'); + const connectionError = new Error('connection failed'); + + obtainMasterConnection.mockResolvedValue([connection, release]); + + await expect( + new PostgresAdvisoryLockService(dataSource).tryWithLock( + 'lock-name', + async () => { + emitConnectionError(connectionError); + throw callbackError; + }, + ), + ).rejects.toBe(callbackError); + + expect(release).toHaveBeenCalledWith(connectionError); + expect(loggerWarn).toHaveBeenCalledWith( + expect.stringContaining(connectionError.message), + ); + }); + + it('preserves the callback error when releasing the lock fails', async () => { + const { connection, query, release } = createMockConnection(true); + const callbackError = new Error('callback failed'); + const releaseError = new Error('release failed'); + + query + .mockReset() + .mockResolvedValueOnce({ + rows: [{ acquired: true }], + }) + .mockRejectedValueOnce(releaseError); + obtainMasterConnection.mockResolvedValue([connection, release]); + + await expect( + new PostgresAdvisoryLockService(dataSource).tryWithLock( + 'lock-name', + async () => { + throw callbackError; + }, + ), + ).rejects.toBe(callbackError); + + expect(release).toHaveBeenCalledWith(releaseError); + expect(loggerWarn).toHaveBeenCalledWith( + expect.stringContaining(releaseError.message), + ); + }); +}); diff --git a/packages/twenty-server/src/database/typeorm/postgres-advisory-lock.service.ts b/packages/twenty-server/src/database/typeorm/postgres-advisory-lock.service.ts new file mode 100644 index 0000000000..60ac61b972 --- /dev/null +++ b/packages/twenty-server/src/database/typeorm/postgres-advisory-lock.service.ts @@ -0,0 +1,184 @@ +import { Injectable, Logger } from '@nestjs/common'; +import { InjectDataSource } from '@nestjs/typeorm'; + +import { type PoolClient } from 'pg'; +import { DataSource } from 'typeorm'; +import { type PostgresDriver } from 'typeorm/driver/postgres/PostgresDriver'; + +const toError = (error: unknown): Error => + error instanceof Error ? error : new Error(String(error)); + +export type PostgresAdvisoryLockResult = + | { + acquired: false; + } + | { + acquired: true; + value: T; + }; + +@Injectable() +export class PostgresAdvisoryLockService { + private readonly logger = new Logger(PostgresAdvisoryLockService.name); + + constructor( + @InjectDataSource() + private readonly coreDataSource: DataSource, + ) {} + + async tryWithLock( + lockName: string, + callback: () => Promise, + ): Promise> { + const lockSession = await PostgresAdvisoryLockSession.tryAcquire( + this.coreDataSource, + lockName, + ); + + if (!lockSession) { + return { + acquired: false, + }; + } + + const callbackResult = await this.captureCallbackResult(callback); + const cleanupError = await lockSession.unlockAndClose(); + + if (callbackResult.status === 'rejected') { + if (cleanupError) { + this.logger.warn( + `Secondary PostgreSQL advisory lock error for "${lockName}" after callback failure: ${cleanupError}`, + ); + } + + throw callbackResult.reason; + } + + if (cleanupError) { + throw cleanupError; + } + + return { + acquired: true, + value: callbackResult.value, + }; + } + + private async captureCallbackResult( + callback: () => Promise, + ): Promise> { + try { + return { + status: 'fulfilled', + value: await callback(), + }; + } catch (reason) { + return { + status: 'rejected', + reason, + }; + } + } +} + +class PostgresAdvisoryLockSession { + private connectionError: Error | undefined; + + private readonly handleConnectionError = (error: Error) => { + this.connectionError ??= error; + }; + + private constructor( + private readonly connection: PoolClient, + private readonly releaseConnection: (error?: Error) => void, + private readonly lockName: string, + ) { + this.connection.on('error', this.handleConnectionError); + } + + static async tryAcquire( + dataSource: DataSource, + lockName: string, + ): Promise { + const [connection, releaseConnection] = (await ( + dataSource.driver as PostgresDriver + ).obtainMasterConnection()) as [PoolClient, (error?: Error) => void]; + const lockSession = new PostgresAdvisoryLockSession( + connection, + releaseConnection, + lockName, + ); + + let isLockAcquired: boolean; + + try { + isLockAcquired = await lockSession.tryAcquireLock(); + } catch (error) { + lockSession.close(toError(error)); + throw error; + } + + if (!isLockAcquired) { + const closeError = lockSession.close(); + + if (closeError) { + throw closeError; + } + + return undefined; + } + + return lockSession; + } + + async unlockAndClose(): Promise { + let unlockError: Error | undefined; + + try { + await this.unlock(); + } catch (error) { + unlockError = toError(error); + } + + return this.close(unlockError); + } + + private async tryAcquireLock(): Promise { + const { + rows: [result], + } = await this.connection.query<{ acquired: boolean }>( + `SELECT pg_try_advisory_lock(hashtextextended($1, 0)) AS "acquired"`, + [this.lockName], + ); + + return result?.acquired === true; + } + + private async unlock(): Promise { + const { + rows: [result], + } = await this.connection.query<{ released: boolean }>( + `SELECT pg_advisory_unlock(hashtextextended($1, 0)) AS "released"`, + [this.lockName], + ); + + if (result?.released !== true) { + throw new Error( + `Could not release PostgreSQL advisory lock ${this.lockName}`, + ); + } + } + + private close(error?: Error): Error | undefined { + this.connection.removeListener('error', this.handleConnectionError); + const errorToRelease = error ?? this.connectionError; + + if (errorToRelease) { + this.releaseConnection(errorToRelease); + } else { + this.releaseConnection(); + } + + return errorToRelease; + } +} diff --git a/packages/twenty-server/src/database/typeorm/typeorm.module.ts b/packages/twenty-server/src/database/typeorm/typeorm.module.ts index 3838a07541..a06d5003cb 100644 --- a/packages/twenty-server/src/database/typeorm/typeorm.module.ts +++ b/packages/twenty-server/src/database/typeorm/typeorm.module.ts @@ -6,6 +6,7 @@ import { DataSource, type DataSourceOptions } from 'typeorm'; import { typeORMCoreModuleOptions } from 'src/database/typeorm/core/core.datasource'; import { DatabaseGaugeService } from 'src/database/typeorm/database-gauge.service'; import { DatabasePoolMetricsService } from 'src/database/typeorm/database-pool-metrics.service'; +import { PostgresAdvisoryLockService } from 'src/database/typeorm/postgres-advisory-lock.service'; import { MetricsModule } from 'src/engine/core-modules/metrics/metrics.module'; import { installUpgradeAwareRepositoryProxy } from 'src/engine/twenty-orm/upgrade-aware/install-upgrade-aware-repository-proxy'; @@ -24,7 +25,11 @@ import { installUpgradeAwareRepositoryProxy } from 'src/engine/twenty-orm/upgrad }), MetricsModule, ], - providers: [DatabasePoolMetricsService, DatabaseGaugeService], - exports: [DatabasePoolMetricsService], + providers: [ + DatabasePoolMetricsService, + DatabaseGaugeService, + PostgresAdvisoryLockService, + ], + exports: [DatabasePoolMetricsService, PostgresAdvisoryLockService], }) export class TypeORMModule {} diff --git a/packages/twenty-server/src/engine/core-modules/workspace/services/__tests__/workspace.service.spec.ts b/packages/twenty-server/src/engine/core-modules/workspace/services/__tests__/workspace.service.spec.ts index e1301ecbd0..8d2636c3bd 100644 --- a/packages/twenty-server/src/engine/core-modules/workspace/services/__tests__/workspace.service.spec.ts +++ b/packages/twenty-server/src/engine/core-modules/workspace/services/__tests__/workspace.service.spec.ts @@ -2,7 +2,7 @@ import { Test, type TestingModule } from '@nestjs/testing'; import { getDataSourceToken, getRepositoryToken } from '@nestjs/typeorm'; import { WorkspaceActivationStatus } from 'twenty-shared/workspace'; -import { IsNull, Not, type Repository } from 'typeorm'; +import { IsNull, Not, type QueryRunner, type Repository } from 'typeorm'; import { BillingSubscriptionService } from 'src/engine/core-modules/billing/services/billing-subscription.service'; import { BillingService } from 'src/engine/core-modules/billing/services/billing.service'; @@ -47,8 +47,21 @@ describe('WorkspaceService', () => { let dnsManagerService: DnsManagerService; let billingSubscriptionService: BillingSubscriptionService; let userWorkspaceService: UserWorkspaceService; + let flatEntityMapsCacheService: WorkspaceManyOrAllFlatEntityMapsCacheService; + let queryRunner: QueryRunner; beforeEach(async () => { + queryRunner = { + connect: jest.fn(), + startTransaction: jest.fn(), + commitTransaction: jest.fn(), + rollbackTransaction: jest.fn(), + release: jest.fn(), + manager: { + delete: jest.fn().mockResolvedValue({ affected: 0 }), + }, + } as unknown as QueryRunner; + const module: TestingModule = await Test.createTestingModule({ providers: [ WorkspaceService, @@ -156,16 +169,7 @@ describe('WorkspaceService', () => { { provide: getDataSourceToken(), useValue: { - createQueryRunner: jest.fn().mockReturnValue({ - connect: jest.fn(), - startTransaction: jest.fn(), - commitTransaction: jest.fn(), - rollbackTransaction: jest.fn(), - release: jest.fn(), - manager: { - delete: jest.fn().mockResolvedValue({ affected: 0 }), - }, - }), + createQueryRunner: jest.fn().mockReturnValue(queryRunner), getRepository: jest.fn().mockReturnValue({ find: jest.fn().mockResolvedValue([]), }), @@ -197,6 +201,10 @@ describe('WorkspaceService', () => { ); userWorkspaceService = module.get(UserWorkspaceService); + flatEntityMapsCacheService = + module.get( + WorkspaceManyOrAllFlatEntityMapsCacheService, + ); }); afterEach(() => { @@ -312,6 +320,29 @@ describe('WorkspaceService', () => { expect(workspaceRepository.softDelete).not.toHaveBeenCalled(); }); + it('should retrieve field metadata before starting the deletion transaction', async () => { + const mockWorkspace = { + id: 'workspace-id', + metadataVersion: 0, + } as WorkspaceEntity; + + jest + .spyOn(workspaceRepository, 'findOne') + .mockResolvedValue(mockWorkspace); + jest.spyOn(userWorkspaceRepository, 'find').mockResolvedValue([]); + + await service.deleteWorkspace(mockWorkspace.id, false); + + const cacheRetrieval = jest.mocked( + flatEntityMapsCacheService.getOrRecomputeManyOrAllFlatEntityMaps, + ); + + expect(cacheRetrieval).toHaveBeenCalledTimes(1); + expect(cacheRetrieval.mock.invocationCallOrder[0]).toBeLessThan( + (queryRunner.startTransaction as jest.Mock).mock.invocationCallOrder[0], + ); + }); + it('should soft delete the workspace', async () => { const mockWorkspace = { id: 'workspace-id', diff --git a/packages/twenty-server/src/engine/core-modules/workspace/services/workspace.service.ts b/packages/twenty-server/src/engine/core-modules/workspace/services/workspace.service.ts index 05d175db04..fa0ab9bdbc 100644 --- a/packages/twenty-server/src/engine/core-modules/workspace/services/workspace.service.ts +++ b/packages/twenty-server/src/engine/core-modules/workspace/services/workspace.service.ts @@ -631,6 +631,9 @@ export class WorkspaceService { private async deleteWorkspaceSyncableMetadataEntities( workspace: WorkspaceEntity, ): Promise { + const fieldMetadataIdChunks = await this.getFieldMetadataIdChunks( + workspace.id, + ); const queryRunner = this.coreDataSource.createQueryRunner(); await queryRunner.connect(); @@ -643,6 +646,7 @@ export class WorkspaceService { const deletedCount = await this.deleteFieldMetadataInChunks( queryRunner, workspace.id, + fieldMetadataIdChunks, ); if (deletedCount > 0) { @@ -679,12 +683,10 @@ export class WorkspaceService { // FieldMetadataEntity has a self-referencing FK (relationTargetFieldMetadataId) // Related fields must be deleted together to avoid constraint violations - private async deleteFieldMetadataInChunks( - queryRunner: QueryRunner, + private async getFieldMetadataIdChunks( workspaceId: string, - ): Promise { + ): Promise { const CHUNK_SIZE = 50; - let totalDeleted = 0; const { flatFieldMetadataMaps } = await this.flatEntityMapsCacheService.getOrRecomputeManyOrAllFlatEntityMaps( @@ -699,7 +701,7 @@ export class WorkspaceService { ).filter(isDefined); if (fields.length === 0) { - return 0; + return []; } const processedIds = new Set(); @@ -734,7 +736,17 @@ export class WorkspaceService { chunks.push(currentChunk); } - for (const [index, chunk] of chunks.entries()) { + return chunks; + } + + private async deleteFieldMetadataInChunks( + queryRunner: QueryRunner, + workspaceId: string, + fieldMetadataIdChunks: string[][], + ): Promise { + let totalDeleted = 0; + + for (const [index, chunk] of fieldMetadataIdChunks.entries()) { const result = await queryRunner.manager .createQueryBuilder() .delete() @@ -747,7 +759,7 @@ export class WorkspaceService { totalDeleted += deletedInChunk; this.logger.log( - `workspace ${workspaceId}: fieldMetadata chunk ${index + 1}/${chunks.length} - deleted ${deletedInChunk} record(s)`, + `workspace ${workspaceId}: fieldMetadata chunk ${index + 1}/${fieldMetadataIdChunks.length} - deleted ${deletedInChunk} record(s)`, ); } diff --git a/packages/twenty-server/src/engine/workspace-manager/workspace-cleaner/crons/__tests__/clean-suspended-workspaces.job.spec.ts b/packages/twenty-server/src/engine/workspace-manager/workspace-cleaner/crons/__tests__/clean-suspended-workspaces.job.spec.ts new file mode 100644 index 0000000000..43363cf151 --- /dev/null +++ b/packages/twenty-server/src/engine/workspace-manager/workspace-cleaner/crons/__tests__/clean-suspended-workspaces.job.spec.ts @@ -0,0 +1,74 @@ +import { type Repository } from 'typeorm'; + +import { type PostgresAdvisoryLockService } from 'src/database/typeorm/postgres-advisory-lock.service'; +import { WorkspaceEntity } from 'src/engine/core-modules/workspace/workspace.entity'; +import { CleanSuspendedWorkspacesJob } from 'src/engine/workspace-manager/workspace-cleaner/crons/clean-suspended-workspaces.job'; +import { type CleanerWorkspaceService } from 'src/engine/workspace-manager/workspace-cleaner/services/cleaner.workspace-service'; + +jest.mock( + 'src/engine/workspace-manager/workspace-cleaner/services/cleaner.workspace-service', + () => ({ + CleanerWorkspaceService: class {}, + }), +); + +describe('CleanSuspendedWorkspacesJob', () => { + const workspaceRepository = { + find: jest.fn(), + }; + const cleanerWorkspaceService = { + batchWarnOrCleanSuspendedWorkspaces: jest.fn(), + }; + const postgresAdvisoryLockService = { + tryWithLock: jest.fn(), + }; + + const createJob = () => + new CleanSuspendedWorkspacesJob( + cleanerWorkspaceService as unknown as CleanerWorkspaceService, + workspaceRepository as unknown as Repository, + postgresAdvisoryLockService as unknown as PostgresAdvisoryLockService, + ); + + beforeEach(() => { + jest.clearAllMocks(); + workspaceRepository.find.mockResolvedValue([{ id: 'workspace-id' }]); + }); + + it('skips cleanup when another execution holds the lock', async () => { + postgresAdvisoryLockService.tryWithLock.mockResolvedValue({ + acquired: false, + }); + + await createJob().handle(); + + expect(workspaceRepository.find).not.toHaveBeenCalled(); + expect( + cleanerWorkspaceService.batchWarnOrCleanSuspendedWorkspaces, + ).not.toHaveBeenCalled(); + }); + + it('cleans suspended workspaces while holding the lock', async () => { + postgresAdvisoryLockService.tryWithLock.mockImplementation( + async (_lockName, callback) => ({ + acquired: true, + value: await callback(), + }), + ); + + await createJob().handle(); + + expect(workspaceRepository.find).toHaveBeenCalledWith({ + select: ['id'], + where: { + activationStatus: 'SUSPENDED', + }, + withDeleted: true, + }); + expect( + cleanerWorkspaceService.batchWarnOrCleanSuspendedWorkspaces, + ).toHaveBeenCalledWith({ + workspaceIds: ['workspace-id'], + }); + }); +}); diff --git a/packages/twenty-server/src/engine/workspace-manager/workspace-cleaner/crons/clean-suspended-workspaces.job.ts b/packages/twenty-server/src/engine/workspace-manager/workspace-cleaner/crons/clean-suspended-workspaces.job.ts index 3fcdde8788..ad5f3efcf6 100644 --- a/packages/twenty-server/src/engine/workspace-manager/workspace-cleaner/crons/clean-suspended-workspaces.job.ts +++ b/packages/twenty-server/src/engine/workspace-manager/workspace-cleaner/crons/clean-suspended-workspaces.job.ts @@ -1,8 +1,10 @@ +import { Logger } from '@nestjs/common'; import { InjectRepository } from '@nestjs/typeorm'; import { Repository } from 'typeorm'; import { WorkspaceActivationStatus } from 'twenty-shared/workspace'; +import { PostgresAdvisoryLockService } from 'src/database/typeorm/postgres-advisory-lock.service'; import { SentryCronMonitor } from 'src/engine/core-modules/cron/sentry-cron-monitor.decorator'; import { Process } from 'src/engine/core-modules/message-queue/decorators/process.decorator'; import { Processor } from 'src/engine/core-modules/message-queue/decorators/processor.decorator'; @@ -11,12 +13,17 @@ import { WorkspaceEntity } from 'src/engine/core-modules/workspace/workspace.ent import { cleanSuspendedWorkspaceCronPattern } from 'src/engine/workspace-manager/workspace-cleaner/crons/clean-suspended-workspaces.cron.pattern'; import { CleanerWorkspaceService } from 'src/engine/workspace-manager/workspace-cleaner/services/cleaner.workspace-service'; +const CLEAN_SUSPENDED_WORKSPACES_LOCK_NAME = 'clean-suspended-workspaces-job'; + @Processor(MessageQueue.cronQueue) export class CleanSuspendedWorkspacesJob { + private readonly logger = new Logger(CleanSuspendedWorkspacesJob.name); + constructor( private readonly cleanerWorkspaceService: CleanerWorkspaceService, @InjectRepository(WorkspaceEntity) private readonly workspaceRepository: Repository, + private readonly postgresAdvisoryLockService: PostgresAdvisoryLockService, ) {} @Process(CleanSuspendedWorkspacesJob.name) @@ -25,16 +32,27 @@ export class CleanSuspendedWorkspacesJob { cleanSuspendedWorkspaceCronPattern, ) async handle(): Promise { - const suspendedWorkspaceIds = await this.workspaceRepository.find({ - select: ['id'], - where: { - activationStatus: WorkspaceActivationStatus.SUSPENDED, - }, - withDeleted: true, - }); + const result = await this.postgresAdvisoryLockService.tryWithLock( + CLEAN_SUSPENDED_WORKSPACES_LOCK_NAME, + async () => { + const suspendedWorkspaceIds = await this.workspaceRepository.find({ + select: ['id'], + where: { + activationStatus: WorkspaceActivationStatus.SUSPENDED, + }, + withDeleted: true, + }); - await this.cleanerWorkspaceService.batchWarnOrCleanSuspendedWorkspaces({ - workspaceIds: suspendedWorkspaceIds.map((workspace) => workspace.id), - }); + await this.cleanerWorkspaceService.batchWarnOrCleanSuspendedWorkspaces({ + workspaceIds: suspendedWorkspaceIds.map((workspace) => workspace.id), + }); + }, + ); + + if (!result.acquired) { + this.logger.log( + 'Skipping suspended workspace cleanup because another execution is running', + ); + } } }