From ee5b3a65f428db47055fb492286a8a90dec86a5c Mon Sep 17 00:00:00 2001 From: Paul Rastoin <45004772+prastoin@users.noreply.github.com> Date: Fri, 19 Jun 2026 17:50:13 +0200 Subject: [PATCH] Workspace migration post transaction commit side effect (#21845) # Introduction Avoid side effect in transaction to reduce lock duration As discussed we're going to introduce in migration runner side effect later through jobs and metadata boolean state tracker in db Review in cubic --- ...e-logic-function-action-handler.service.ts | 70 ++++++++----------- ...runner-action-handler-service.interface.ts | 15 +++- .../workspace-migration-runner.service.ts | 28 +++++++- .../types/after-commit-side-effect.type.ts | 4 ++ 4 files changed, 74 insertions(+), 43 deletions(-) create mode 100644 packages/twenty-server/src/engine/workspace-manager/workspace-migration/workspace-migration-runner/types/after-commit-side-effect.type.ts diff --git a/packages/twenty-server/src/engine/workspace-manager/workspace-migration/workspace-migration-runner/action-handlers/logic-function/services/delete-logic-function-action-handler.service.ts b/packages/twenty-server/src/engine/workspace-manager/workspace-migration/workspace-migration-runner/action-handlers/logic-function/services/delete-logic-function-action-handler.service.ts index b25da9ebf1..ca7ac4dc71 100644 --- a/packages/twenty-server/src/engine/workspace-manager/workspace-migration/workspace-migration-runner/action-handlers/logic-function/services/delete-logic-function-action-handler.service.ts +++ b/packages/twenty-server/src/engine/workspace-manager/workspace-migration/workspace-migration-runner/action-handlers/logic-function/services/delete-logic-function-action-handler.service.ts @@ -8,13 +8,14 @@ import { FileStorageService } from 'src/engine/core-modules/file-storage/file-st import { LOGIC_FUNCTION_DRIVER_FACTORY_TOKEN } from 'src/engine/core-modules/logic-function/logic-function-drivers/constants/logic-function-driver-factory.token'; import { findFlatEntityByIdInFlatEntityMapsOrThrow } from 'src/engine/metadata-modules/flat-entity/utils/find-flat-entity-by-id-in-flat-entity-maps-or-throw.util'; import { LogicFunctionEntity } from 'src/engine/metadata-modules/logic-function/logic-function.entity'; +import { getLogicFunctionSubfolderForFromSource } from 'src/engine/metadata-modules/logic-function/utils/get-logic-function-subfolder-for-from-source'; import type { LogicFunctionDriverFactory } from 'src/engine/core-modules/logic-function/logic-function-drivers/logic-function-driver.factory'; -import { getLogicFunctionSubfolderForFromSource } from 'src/engine/metadata-modules/logic-function/utils/get-logic-function-subfolder-for-from-source'; import { FlatDeleteLogicFunctionAction, UniversalDeleteLogicFunctionAction, } from 'src/engine/workspace-manager/workspace-migration/workspace-migration-builder/builders/logic-function/types/workspace-migration-logic-function-action.type'; +import { type AfterCommitSideEffect } from 'src/engine/workspace-manager/workspace-migration/workspace-migration-runner/types/after-commit-side-effect.type'; import { WorkspaceMigrationActionRunnerArgs, WorkspaceMigrationActionRunnerContext, @@ -42,18 +43,7 @@ export class DeleteLogicFunctionActionHandlerService extends WorkspaceMigrationR async executeForMetadata( context: WorkspaceMigrationActionRunnerContext, ): Promise { - const { - flatAction, - queryRunner, - workspaceId, - allFlatEntityMaps, - flatApplication, - } = context; - - const flatLogicFunction = findFlatEntityByIdInFlatEntityMapsOrThrow({ - flatEntityMaps: allFlatEntityMaps.flatLogicFunctionMaps, - flatEntityId: flatAction.entityId, - }); + const { flatAction, queryRunner, workspaceId } = context; const logicFunctionRepository = queryRunner.manager.getRepository( @@ -64,13 +54,25 @@ export class DeleteLogicFunctionActionHandlerService extends WorkspaceMigrationR id: flatAction.entityId, workspaceId, }); + } + + protected override getAfterCommitSideEffects( + context: WorkspaceMigrationActionRunnerContext, + ): AfterCommitSideEffect[] { + const { flatAction, workspaceId, allFlatEntityMaps, flatApplication } = + context; + + const flatLogicFunction = findFlatEntityByIdInFlatEntityMapsOrThrow({ + flatEntityMaps: allFlatEntityMaps.flatLogicFunctionMaps, + flatEntityId: flatAction.entityId, + }); const applicationUniversalIdentifier = flatApplication.universalIdentifier; - await Promise.all([ - this.deleteBestEffort( - `source folder for logic function ${flatLogicFunction.id}`, - () => + return [ + { + description: `source folder for logic function ${flatLogicFunction.id} (universalIdentifier=${flatLogicFunction.universalIdentifier}, applicationUniversalIdentifier=${applicationUniversalIdentifier})`, + run: () => this.fileStorageService.deleteFolder({ workspaceId, applicationUniversalIdentifier, @@ -79,39 +81,25 @@ export class DeleteLogicFunctionActionHandlerService extends WorkspaceMigrationR flatLogicFunction.id, ), }), - ), - this.deleteBestEffort( - `built handler for logic function ${flatLogicFunction.id}`, - () => + }, + { + description: `built handler for logic function ${flatLogicFunction.id} (universalIdentifier=${flatLogicFunction.universalIdentifier}, applicationUniversalIdentifier=${applicationUniversalIdentifier})`, + run: () => this.fileStorageService.deleteFile({ workspaceId, applicationUniversalIdentifier, fileFolder: FileFolder.BuiltLogicFunction, resourcePath: flatLogicFunction.builtHandlerPath, }), - ), - this.deleteBestEffort( - `runtime resource for logic function ${flatLogicFunction.id}`, - () => + }, + { + description: `runtime resource for logic function ${flatLogicFunction.id} (universalIdentifier=${flatLogicFunction.universalIdentifier}, applicationUniversalIdentifier=${applicationUniversalIdentifier})`, + run: () => this.logicFunctionDriverFactory .getCurrentDriver() .delete(flatLogicFunction), - ), - ]); - } - - private async deleteBestEffort( - description: string, - operation: () => Promise, - ): Promise { - try { - await operation(); - } catch (error) { - this.logger.warn( - `Failed to delete ${description}: ${error instanceof Error ? error.message : String(error)}`, - DeleteLogicFunctionActionHandlerService.name, - ); - } + }, + ]; } async rollbackForMetadata(): Promise {} diff --git a/packages/twenty-server/src/engine/workspace-manager/workspace-migration/workspace-migration-runner/interfaces/workspace-migration-runner-action-handler-service.interface.ts b/packages/twenty-server/src/engine/workspace-manager/workspace-migration/workspace-migration-runner/interfaces/workspace-migration-runner-action-handler-service.interface.ts index b2434500b9..ccc0635e70 100644 --- a/packages/twenty-server/src/engine/workspace-manager/workspace-migration/workspace-migration-runner/interfaces/workspace-migration-runner-action-handler-service.interface.ts +++ b/packages/twenty-server/src/engine/workspace-manager/workspace-migration/workspace-migration-runner/interfaces/workspace-migration-runner-action-handler-service.interface.ts @@ -26,6 +26,7 @@ import { WorkspaceMigrationRunnerException, WorkspaceMigrationRunnerExceptionCode, } from 'src/engine/workspace-manager/workspace-migration/workspace-migration-runner/exceptions/workspace-migration-runner.exception'; +import { type AfterCommitSideEffect } from 'src/engine/workspace-manager/workspace-migration/workspace-migration-runner/types/after-commit-side-effect.type'; import { type MetadataEvent } from 'src/engine/workspace-manager/workspace-migration/workspace-migration-runner/types/metadata-event'; import { WorkspaceMigrationActionRunnerContext, @@ -54,6 +55,7 @@ export type ActionHandlerExecuteResult = | MetadataToFlatEntityMapsKey >; metadataEvents: MetadataEvent[]; + afterCommitSideEffects: AfterCommitSideEffect[]; }; export abstract class BaseWorkspaceMigrationRunnerActionHandlerService< @@ -134,6 +136,12 @@ export abstract class BaseWorkspaceMigrationRunnerActionHandlerService< return Promise.resolve(); } + protected getAfterCommitSideEffects( + _context: WorkspaceMigrationActionRunnerContext, + ): AfterCommitSideEffect[] { + return []; + } + private optimisticallyApplyActionOnAllFlatEntityMaps({ flatAction, allFlatEntityMaps, @@ -278,13 +286,18 @@ export abstract class BaseWorkspaceMigrationRunnerActionHandlerService< allFlatEntityMaps: context.allFlatEntityMaps, }); + const afterCommitSideEffects = this.getAfterCommitSideEffects({ + ...context, + flatAction, + }); + const partialOptimisticCache = this.optimisticallyApplyActionOnAllFlatEntityMaps({ flatAction, allFlatEntityMaps: context.allFlatEntityMaps, }); - return { partialOptimisticCache, metadataEvents }; + return { partialOptimisticCache, metadataEvents, afterCommitSideEffects }; } async rollback( diff --git a/packages/twenty-server/src/engine/workspace-manager/workspace-migration/workspace-migration-runner/services/workspace-migration-runner.service.ts b/packages/twenty-server/src/engine/workspace-manager/workspace-migration/workspace-migration-runner/services/workspace-migration-runner.service.ts index 62476baf5f..de4902d64b 100644 --- a/packages/twenty-server/src/engine/workspace-manager/workspace-migration/workspace-migration-runner/services/workspace-migration-runner.service.ts +++ b/packages/twenty-server/src/engine/workspace-manager/workspace-migration/workspace-migration-runner/services/workspace-migration-runner.service.ts @@ -23,6 +23,7 @@ import { WorkspaceMigrationRunnerExceptionCode, } from 'src/engine/workspace-manager/workspace-migration/workspace-migration-runner/exceptions/workspace-migration-runner.exception'; import { WorkspaceMigrationRunnerActionHandlerRegistryService } from 'src/engine/workspace-manager/workspace-migration/workspace-migration-runner/registry/workspace-migration-runner-action-handler-registry.service'; +import { type AfterCommitSideEffect } from 'src/engine/workspace-manager/workspace-migration/workspace-migration-runner/types/after-commit-side-effect.type'; import { type MetadataEvent } from 'src/engine/workspace-manager/workspace-migration/workspace-migration-runner/types/metadata-event'; @Injectable() @@ -283,6 +284,7 @@ export class WorkspaceMigrationRunnerService { await queryRunner.startTransaction(); const allMetadataEvents: MetadataEvent[] = []; + const allAfterCommitSideEffects: AfterCommitSideEffect[] = []; const transactionStart = performance.now(); let slowestActionMs = 0; @@ -294,7 +296,11 @@ export class WorkspaceMigrationRunnerService { for (const action of actions) { const actionStart = performance.now(); - const { partialOptimisticCache, metadataEvents } = + const { + partialOptimisticCache, + metadataEvents, + afterCommitSideEffects, + } = await this.workspaceMigrationRunnerActionHandlerRegistry.executeActionHandler( { action, @@ -330,6 +336,7 @@ export class WorkspaceMigrationRunnerService { } as typeof allFlatEntityMaps; allMetadataEvents.push(...metadataEvents); + allAfterCommitSideEffects.push(...afterCommitSideEffects); } const commitStart = performance.now(); @@ -432,6 +439,25 @@ export class WorkspaceMigrationRunnerService { 'Runner', ); + const sideEffectResults = await Promise.allSettled( + allAfterCommitSideEffects.map((sideEffect) => + Promise.resolve().then(() => sideEffect.run()), + ), + ); + + sideEffectResults.forEach((result, index) => { + if (result.status === 'rejected') { + this.logger.warn( + `After-commit side effect failed (${allAfterCommitSideEffects[index].description}): ${ + result.reason instanceof Error + ? result.reason.message + : String(result.reason) + }`, + 'Runner', + ); + } + }); + const hasSchemaMetadataChanged = allFlatEntityMapsKeys.includes('flatObjectMetadataMaps') || allFlatEntityMapsKeys.includes('flatFieldMetadataMaps'); diff --git a/packages/twenty-server/src/engine/workspace-manager/workspace-migration/workspace-migration-runner/types/after-commit-side-effect.type.ts b/packages/twenty-server/src/engine/workspace-manager/workspace-migration/workspace-migration-runner/types/after-commit-side-effect.type.ts new file mode 100644 index 0000000000..822357f2bc --- /dev/null +++ b/packages/twenty-server/src/engine/workspace-manager/workspace-migration/workspace-migration-runner/types/after-commit-side-effect.type.ts @@ -0,0 +1,4 @@ +export type AfterCommitSideEffect = { + description: string; + run: () => Promise; +};