Refactor workflow-logic-function-interaction (#17699)

# Refactor workflow–logic function interaction

## Why

Workflow code steps and standalone logic functions shared the same build
layer and DB layer, which blurred two use cases: code steps belong to a
workflow version; standalone functions are deployable units. That made
workflow code steps harder to own and evolve.

## Goal

Treat code steps as **workflow-owned**: build and run them in workflow
context, and expose workflow-scoped APIs so the editor can load, test,
and save code step source without going through the generic
logic-function layer.
This commit is contained in:
Charles Bochet
2026-02-05 01:46:55 +01:00
committed by GitHub
parent 752359335e
commit 67a98f77e3
77 changed files with 1750 additions and 1336 deletions
@@ -1,40 +1,27 @@
import { Injectable } from '@nestjs/common';
import { InjectRepository } from '@nestjs/typeorm';
import crypto from 'crypto';
import { Repository } from 'typeorm';
import { FileFolder } from 'twenty-shared/types';
import { isDefined } from 'twenty-shared/utils';
import { WorkspaceMigrationRunnerActionHandler } from 'src/engine/workspace-manager/workspace-migration/workspace-migration-runner/interfaces/workspace-migration-runner-action-handler-service.interface';
import { ApplicationEntity } from 'src/engine/core-modules/application/application.entity';
import { FileStorageService } from 'src/engine/core-modules/file-storage/file-storage.service';
import { getSeedProjectFiles } from 'src/engine/core-modules/logic-function/logic-function-drivers/utils/get-seed-project-files';
import { getLogicFunctionBaseFolderPath } from 'src/engine/core-modules/logic-function/logic-function-build/utils/get-logic-function-base-folder-path.util';
import { LogicFunctionEntity } from 'src/engine/metadata-modules/logic-function/logic-function.entity';
import {
LogicFunctionException,
LogicFunctionExceptionCode,
} from 'src/engine/metadata-modules/logic-function/logic-function.exception';
import { FlatLogicFunction } from 'src/engine/metadata-modules/logic-function/types/flat-logic-function.type';
import { FlatCreateLogicFunctionAction } from 'src/engine/workspace-manager/workspace-migration/workspace-migration-builder/builders/logic-function/types/workspace-migration-logic-function-action.type';
import {
WorkspaceMigrationActionRunnerArgs,
WorkspaceMigrationActionRunnerContext,
} from 'src/engine/workspace-manager/workspace-migration/workspace-migration-runner/types/workspace-migration-action-runner-args.type';
import { FlatCreateLogicFunctionAction } from 'src/engine/workspace-manager/workspace-migration/workspace-migration-builder/builders/logic-function/types/workspace-migration-logic-function-action.type';
@Injectable()
export class CreateLogicFunctionActionHandlerService extends WorkspaceMigrationRunnerActionHandler(
'create',
'logicFunction',
) {
constructor(
private readonly fileStorageService: FileStorageService,
@InjectRepository(ApplicationEntity)
private readonly applicationRepository: Repository<ApplicationEntity>,
) {
constructor(private readonly fileStorageService: FileStorageService) {
super();
}
@@ -47,18 +34,30 @@ export class CreateLogicFunctionActionHandlerService extends WorkspaceMigrationR
async executeForMetadata(
context: WorkspaceMigrationActionRunnerContext<FlatCreateLogicFunctionAction>,
): Promise<void> {
const { flatAction, queryRunner, workspaceId } = context;
const { flatAction, queryRunner, workspaceId, flatApplication } = context;
const { flatEntity: logicFunction } = flatAction;
const applicationUniversalIdentifier =
await this.getApplicationUniversalIdentifier(logicFunction.applicationId);
const applicationUniversalIdentifier = flatApplication.universalIdentifier;
let seedChecksum: string | undefined;
if (!isDefined(logicFunction.checksum)) {
seedChecksum = await this.seedLogicFunctionFiles(
logicFunction,
const [sourceExists, builtExists] = await Promise.all([
this.fileStorageService.checkFileExists_v2({
workspaceId,
applicationUniversalIdentifier,
fileFolder: FileFolder.Source,
resourcePath: logicFunction.sourceHandlerPath,
}),
this.fileStorageService.checkFileExists_v2({
workspaceId,
applicationUniversalIdentifier,
fileFolder: FileFolder.BuiltLogicFunction,
resourcePath: logicFunction.builtHandlerPath,
}),
]);
if (!sourceExists || !builtExists) {
throw new LogicFunctionException(
`Logic function source or built file missing before create (source: ${sourceExists}, built: ${builtExists})`,
LogicFunctionExceptionCode.LOGIC_FUNCTION_CREATE_FAILED,
);
}
@@ -70,103 +69,16 @@ export class CreateLogicFunctionActionHandlerService extends WorkspaceMigrationR
await logicFunctionRepository.insert({
...logicFunction,
workspaceId,
checksum: seedChecksum ?? logicFunction.checksum,
});
}
private async getApplicationUniversalIdentifier(
applicationId: string,
): Promise<string> {
const application = await this.applicationRepository.findOne({
where: { id: applicationId },
select: ['universalIdentifier'],
});
if (!isDefined(application)) {
throw new LogicFunctionException(
`Application with id ${applicationId} not found`,
LogicFunctionExceptionCode.LOGIC_FUNCTION_NOT_FOUND,
);
}
return application.universalIdentifier;
}
private async seedLogicFunctionFiles(
logicFunction: FlatLogicFunction,
applicationUniversalIdentifier: string,
): Promise<string> {
const seedProjectFiles = await getSeedProjectFiles;
const sourceFiles = seedProjectFiles.filter((file) =>
file.name.endsWith('index.ts'),
);
const builtFiles = seedProjectFiles.filter((file) =>
file.name.endsWith('.mjs'),
);
if (sourceFiles.length !== 1 || builtFiles.length !== 1) {
throw new LogicFunctionException(
'Seed project should have one index.ts file and one index.mjs file',
LogicFunctionExceptionCode.LOGIC_FUNCTION_CREATE_FAILED,
);
}
const sourceFile = sourceFiles[0];
const builtFile = builtFiles[0];
await this.fileStorageService.writeFile_v2({
workspaceId: logicFunction.workspaceId,
applicationUniversalIdentifier,
fileFolder: FileFolder.Source,
resourcePath: logicFunction.sourceHandlerPath,
sourceFile: sourceFile.content,
mimeType: 'application/typescript',
settings: {
isTemporaryFile: false,
toDelete: false,
},
});
await this.fileStorageService.writeFile_v2({
workspaceId: logicFunction.workspaceId,
applicationUniversalIdentifier,
fileFolder: FileFolder.BuiltLogicFunction,
resourcePath: logicFunction.builtHandlerPath,
sourceFile: builtFile.content,
mimeType: 'application/javascript',
settings: {
isTemporaryFile: false,
toDelete: false,
},
});
return crypto.createHash('md5').update(builtFile.content).digest('hex');
}
async rollbackForMetadata(
context: Omit<
_context: Omit<
WorkspaceMigrationActionRunnerArgs<FlatCreateLogicFunctionAction>,
'queryRunner'
>,
): Promise<void> {
const { action } = context;
const applicationUniversalIdentifier =
await this.getApplicationUniversalIdentifier(
action.flatEntity.applicationId,
);
const baseFolderPath = getLogicFunctionBaseFolderPath(
action.flatEntity.sourceHandlerPath,
);
await this.fileStorageService.delete_v2({
workspaceId: action.flatEntity.workspaceId,
applicationUniversalIdentifier,
fileFolder: FileFolder.Source,
resourcePath: baseFolderPath,
});
// Nothing to rollback for now
return;
}
}
@@ -7,7 +7,6 @@ import { WorkspaceMigrationRunnerActionHandler } from 'src/engine/workspace-mana
import { FileStorageService } from 'src/engine/core-modules/file-storage/file-storage.service';
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 { getLogicFunctionBaseFolderPath } from 'src/engine/core-modules/logic-function/logic-function-build/utils/get-logic-function-base-folder-path.util';
import {
FlatDeleteLogicFunctionAction,
UniversalDeleteLogicFunctionAction,
@@ -35,7 +34,13 @@ export class DeleteLogicFunctionActionHandlerService extends WorkspaceMigrationR
async executeForMetadata(
context: WorkspaceMigrationActionRunnerContext<FlatDeleteLogicFunctionAction>,
): Promise<void> {
const { flatAction, queryRunner, workspaceId, allFlatEntityMaps } = context;
const {
flatAction,
queryRunner,
workspaceId,
allFlatEntityMaps,
flatApplication,
} = context;
const flatLogicFunction = findFlatEntityByIdInFlatEntityMapsOrThrow({
flatEntityMaps: allFlatEntityMaps.flatLogicFunctionMaps,
@@ -52,18 +57,14 @@ export class DeleteLogicFunctionActionHandlerService extends WorkspaceMigrationR
workspaceId,
});
const sourceBaseFolderPath = getLogicFunctionBaseFolderPath(
flatLogicFunction.sourceHandlerPath,
);
const builtBaseFolderPath = getLogicFunctionBaseFolderPath(
flatLogicFunction.builtHandlerPath,
);
const applicationUniversalIdentifier = flatApplication.universalIdentifier;
const sourceFolderPath = `workspace-${workspaceId}/${FileFolder.Source}/${sourceBaseFolderPath}`;
const builtFolderPath = `workspace-${workspaceId}/${FileFolder.BuiltLogicFunction}/${builtBaseFolderPath}`;
await this.fileStorageService.delete({ folderPath: sourceFolderPath });
await this.fileStorageService.delete({ folderPath: builtFolderPath });
await this.fileStorageService.delete_v2({
workspaceId,
applicationUniversalIdentifier,
fileFolder: FileFolder.BuiltLogicFunction,
resourcePath: flatLogicFunction.builtHandlerPath,
});
}
async rollbackForMetadata(): Promise<void> {}
@@ -1,21 +1,10 @@
import { Injectable } from '@nestjs/common';
import { InjectRepository } from '@nestjs/typeorm';
import { promises as fs } from 'fs';
import { dirname, join } from 'path';
import { isObject } from '@sniptt/guards';
import { FileFolder, Sources } from 'twenty-shared/types';
import { isDefined } from 'twenty-shared/utils';
import { Repository } from 'typeorm';
import { FileFolder } from 'twenty-shared/types';
import { WorkspaceMigrationRunnerActionHandler } from 'src/engine/workspace-manager/workspace-migration/workspace-migration-runner/interfaces/workspace-migration-runner-action-handler-service.interface';
import { ApplicationEntity } from 'src/engine/core-modules/application/application.entity';
import { FileStorageService } from 'src/engine/core-modules/file-storage/file-storage.service';
import { LogicFunctionBuildService } from 'src/engine/core-modules/logic-function/logic-function-build/services/logic-function-build.service';
import { getLogicFunctionBaseFolderPath } from 'src/engine/core-modules/logic-function/logic-function-build/utils/get-logic-function-base-folder-path.util';
import { LambdaBuildDirectoryManager } from 'src/engine/core-modules/logic-function/logic-function-drivers/utils/lambda-build-directory-manager';
import { LogicFunctionExecutorService } from 'src/engine/core-modules/logic-function/logic-function-executor/services/logic-function-executor.service';
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';
@@ -37,9 +26,6 @@ export class UpdateLogicFunctionActionHandlerService extends WorkspaceMigrationR
) {
constructor(
private readonly logicFunctionExecutorService: LogicFunctionExecutorService,
private readonly functionBuildService: LogicFunctionBuildService,
@InjectRepository(ApplicationEntity)
private readonly applicationRepository: Repository<ApplicationEntity>,
private readonly fileStorageService: FileStorageService,
) {
super();
@@ -54,8 +40,8 @@ export class UpdateLogicFunctionActionHandlerService extends WorkspaceMigrationR
async executeForMetadata(
context: WorkspaceMigrationActionRunnerContext<FlatUpdateLogicFunctionAction>,
): Promise<void> {
const { flatAction, queryRunner, workspaceId } = context;
const { entityId, code, update } = flatAction;
const { flatAction, queryRunner, flatApplication, workspaceId } = context;
const { entityId, update } = flatAction;
const logicFunctionRepository =
queryRunner.manager.getRepository<LogicFunctionEntity>(
@@ -69,111 +55,42 @@ export class UpdateLogicFunctionActionHandlerService extends WorkspaceMigrationR
flatEntityMaps: context.allFlatEntityMaps.flatLogicFunctionMaps,
});
if (isDefined(update.checksum) && isDefined(code)) {
await this.handleChecksumUpdate({
flatLogicFunction,
code,
});
}
if (update.deletedAt !== undefined) {
await this.handleDeletedAtUpdate({
const applicationUniversalIdentifier = flatApplication.universalIdentifier;
if (update.checksum) {
await this.verifySourceAndBuiltFilesExist({
flatLogicFunction,
applicationUniversalIdentifier,
});
}
}
async handleDeletedAtUpdate({
private async verifySourceAndBuiltFilesExist({
flatLogicFunction,
applicationUniversalIdentifier,
}: {
flatLogicFunction: FlatLogicFunction;
}) {
await this.logicFunctionExecutorService.delete(flatLogicFunction);
}
private async getApplicationUniversalIdentifier(
applicationId: string,
): Promise<string> {
const application = await this.applicationRepository.findOne({
where: { id: applicationId },
select: ['universalIdentifier'],
});
if (!isDefined(application)) {
throw new LogicFunctionException(
`Application with id ${applicationId} not found`,
LogicFunctionExceptionCode.LOGIC_FUNCTION_NOT_FOUND,
);
}
return application.universalIdentifier;
}
private async writeSourcesToLocalFolder(
sources: Sources,
localPath: string,
): Promise<void> {
for (const key of Object.keys(sources)) {
const filePath = join(localPath, key);
const value = sources[key];
if (isObject(value)) {
await this.writeSourcesToLocalFolder(value as Sources, filePath);
continue;
}
await fs.mkdir(dirname(filePath), { recursive: true });
await fs.writeFile(filePath, value);
}
}
async handleChecksumUpdate({
flatLogicFunction,
code,
}: {
flatLogicFunction: FlatLogicFunction;
code: Sources;
}) {
const applicationUniversalIdentifier =
await this.getApplicationUniversalIdentifier(
flatLogicFunction.applicationId,
);
const lambdaBuildDirectoryManager = new LambdaBuildDirectoryManager();
try {
const { sourceTemporaryDir } = await lambdaBuildDirectoryManager.init();
await this.writeSourcesToLocalFolder(code, sourceTemporaryDir);
const baseFolderPath = getLogicFunctionBaseFolderPath(
flatLogicFunction.sourceHandlerPath,
);
await this.fileStorageService.uploadFolder_v2({
applicationUniversalIdentifier: string;
}): Promise<void> {
const [sourceExists, builtExists] = await Promise.all([
this.fileStorageService.checkFileExists_v2({
workspaceId: flatLogicFunction.workspaceId,
applicationUniversalIdentifier,
fileFolder: FileFolder.Source,
resourcePath: baseFolderPath,
localPath: sourceTemporaryDir,
});
} catch (error) {
this.logger.log(
'workspace-migration-runner',
`Error updating logic function ${flatLogicFunction.id}: ${error.message}`,
);
} finally {
await lambdaBuildDirectoryManager.clean();
}
try {
await this.functionBuildService.buildAndUpload({
flatLogicFunction,
resourcePath: flatLogicFunction.sourceHandlerPath,
}),
this.fileStorageService.checkFileExists_v2({
workspaceId: flatLogicFunction.workspaceId,
applicationUniversalIdentifier,
});
} catch (error) {
// TODO: logic should be improved
this.logger.log(
'workspace-migration-runner',
`Error building and uploading logic function ${flatLogicFunction.id}: ${error.message}`,
fileFolder: FileFolder.BuiltLogicFunction,
resourcePath: flatLogicFunction.builtHandlerPath,
}),
]);
if (!sourceExists || !builtExists) {
throw new LogicFunctionException(
`Logic function source or built file missing before update (source: ${sourceExists}, built: ${builtExists})`,
LogicFunctionExceptionCode.LOGIC_FUNCTION_NOT_READY,
);
}
}
@@ -155,6 +155,30 @@ export class WorkspaceMigrationRunnerService {
this.logger.time('Runner', 'Total execution');
this.logger.time('Runner', 'Initial cache retrieval');
const queryRunner = this.coreDataSource.createQueryRunner();
const actionMetadataNames = [
...new Set(actions.flatMap((action) => action.metadataName)),
];
const actionsMetadataAndRelatedMetadataNames: AllMetadataName[] = [
...new Set([
...actionMetadataNames,
...actionMetadataNames.flatMap(getMetadataRelatedMetadataNames),
]),
];
const allFlatEntityMapsKeys = actionsMetadataAndRelatedMetadataNames.map(
getMetadataFlatEntityMapsKey,
);
let allFlatEntityMaps =
await this.flatEntityMapsCacheService.getOrRecomputeManyOrAllFlatEntityMaps<
typeof allFlatEntityMapsKeys
>({
workspaceId,
flatMapsKeys: allFlatEntityMapsKeys,
});
this.logger.timeEnd('Runner', 'Initial cache retrieval');
const { flatApplicationMaps } =
await this.workspaceCacheService.getOrRecompute(workspaceId, [
'flatApplicationMaps',
@@ -175,35 +199,12 @@ export class WorkspaceMigrationRunnerService {
});
}
const queryRunner = this.coreDataSource.createQueryRunner();
const actionMetadataNames = [
...new Set(actions.flatMap((action) => action.metadataName)),
];
const actionsMetadataAndRelatedMetadataNames: AllMetadataName[] = [
...new Set([
...actionMetadataNames,
...actionMetadataNames.flatMap(getMetadataRelatedMetadataNames),
]),
];
const allFlatEntityMapsKeys = actionsMetadataAndRelatedMetadataNames.map(
getMetadataFlatEntityMapsKey,
);
let allFlatEntityMaps =
await this.flatEntityMapsCacheService.getOrRecomputeManyOrAllFlatEntityMaps<
typeof allFlatEntityMapsKeys,
false
>({
workspaceId,
flatMapsKeys: allFlatEntityMapsKeys,
});
this.logger.timeEnd('Runner', 'Initial cache retrieval');
this.logger.time('Runner', 'Transaction execution');
try {
await queryRunner.connect();
await queryRunner.startTransaction();
for (const action of actions) {
const result =
await this.workspaceMigrationRunnerActionHandlerRegistry.executeActionHandler(
@@ -222,7 +223,7 @@ export class WorkspaceMigrationRunnerService {
allFlatEntityMaps = {
...allFlatEntityMaps,
...result,
};
} as typeof allFlatEntityMaps;
}
await queryRunner.commitTransaction();