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; +};