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 <!-- This is an auto-generated description by cubic. --> <a href="https://cubic.dev/pr/twentyhq/twenty/pull/21845?utm_source=github" target="_blank" rel="noopener noreferrer" data-no-image-dialog="true"><picture><source media="(prefers-color-scheme: dark)" srcset="https://www.cubic.dev/buttons/review-in-cubic-dark.svg"><source media="(prefers-color-scheme: light)" srcset="https://www.cubic.dev/buttons/review-in-cubic-light.svg"><img alt="Review in cubic" src="https://www.cubic.dev/buttons/review-in-cubic-dark.svg"></picture></a> <!-- End of auto-generated description by cubic. -->
This commit is contained in:
+29
-41
@@ -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<FlatDeleteLogicFunctionAction>,
|
||||
): Promise<void> {
|
||||
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<LogicFunctionEntity>(
|
||||
@@ -64,13 +54,25 @@ export class DeleteLogicFunctionActionHandlerService extends WorkspaceMigrationR
|
||||
id: flatAction.entityId,
|
||||
workspaceId,
|
||||
});
|
||||
}
|
||||
|
||||
protected override getAfterCommitSideEffects(
|
||||
context: WorkspaceMigrationActionRunnerContext<FlatDeleteLogicFunctionAction>,
|
||||
): 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<void>,
|
||||
): Promise<void> {
|
||||
try {
|
||||
await operation();
|
||||
} catch (error) {
|
||||
this.logger.warn(
|
||||
`Failed to delete ${description}: ${error instanceof Error ? error.message : String(error)}`,
|
||||
DeleteLogicFunctionActionHandlerService.name,
|
||||
);
|
||||
}
|
||||
},
|
||||
];
|
||||
}
|
||||
|
||||
async rollbackForMetadata(): Promise<void> {}
|
||||
|
||||
+14
-1
@@ -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<TMetadataName extends AllMetadataName> =
|
||||
| MetadataToFlatEntityMapsKey<TMetadataName>
|
||||
>;
|
||||
metadataEvents: MetadataEvent[];
|
||||
afterCommitSideEffects: AfterCommitSideEffect[];
|
||||
};
|
||||
|
||||
export abstract class BaseWorkspaceMigrationRunnerActionHandlerService<
|
||||
@@ -134,6 +136,12 @@ export abstract class BaseWorkspaceMigrationRunnerActionHandlerService<
|
||||
return Promise.resolve();
|
||||
}
|
||||
|
||||
protected getAfterCommitSideEffects(
|
||||
_context: WorkspaceMigrationActionRunnerContext<TFlatAction>,
|
||||
): 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(
|
||||
|
||||
+27
-1
@@ -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');
|
||||
|
||||
+4
@@ -0,0 +1,4 @@
|
||||
export type AfterCommitSideEffect = {
|
||||
description: string;
|
||||
run: () => Promise<void>;
|
||||
};
|
||||
Reference in New Issue
Block a user