Migrate workflow serverless to logic (#17646)

Migrations commands

---------

Co-authored-by: Charles Bochet <charles@macbook-pro-de-charles-1.home>
This commit is contained in:
Charles Bochet
2026-02-02 19:39:20 +01:00
committed by GitHub
parent fd47d5a1a9
commit cfdcc0065c
9 changed files with 1212 additions and 21 deletions
@@ -0,0 +1,208 @@
On v1.17
"settings": {
"input": {
"logicFunctionId": "fa1d4f4c-b0b6-43cd-8ab9-c87ac3b233d1",
"logicFunctionInput": {
"a": null,
"b": null
}
},
"outputSchema": {
"link": {
"tab": "test",
"icon": "IconVariable",
"label": "Generate Function Output",
"isLeaf": true
},
"_outputSchemaType": "LINK"
}
"settings": {
"input": {
"serverlessFunctionId": "84d3e16c-10bb-4eed-ae15-714f657a2726",
"serverlessFunctionInput": {
"a": null,
"b": null
},
"serverlessFunctionVersion": "draft"
},
"outputSchema": {
"link": {
"tab": "test",
"icon": "IconVariable",
"label": "Generate Function Output",
"isLeaf": true
},
"_outputSchemaType": "LINK"
}
"settings": {
"input": {
"serverlessFunctionId": "cfc3c41f-5520-41f5-a83e-6a56ca6b90ba",
"serverlessFunctionInput": {
"a": null,
"b": null
},
"serverlessFunctionVersion": "1"
},
"outputSchema": {
"link": {
"tab": "test",
"icon": "IconVariable",
"label": "Generate Function Output",
"isLeaf": true
},
"_outputSchemaType": "LINK"
}
"settings": {
"input": {
"serverlessFunctionId": "42b11c8c-0238-4044-a5d3-9cc662066de4",
"serverlessFunctionInput": {
"a": null,
"b": null
},
"serverlessFunctionVersion": "1"
},
"outputSchema": {
"link": {
"tab": "test",
"icon": "IconVariable",
"label": "Generate Function Output",
"isLeaf": true
},
"_outputSchemaType": "LINK"
}
"settings": {
"input": {
"serverlessFunctionId": "42b11c8c-0238-4044-a5d3-9cc662066de4",
"serverlessFunctionInput": {
"a": null,
"b": null
},
"serverlessFunctionVersion": "draft"
},
"outputSchema": {
"link": {
"tab": "test",
"icon": "IconVariable",
"label": "Generate Function Output",
"isLeaf": true
},
"_outputSchemaType": "LINK"
}
"settings": {
"input": {
"serverlessFunctionId": "a7d07444-092f-4a96-8e3c-b379bc49f6e3",
"serverlessFunctionInput": {
"a": null,
"b": null
},
"serverlessFunctionVersion": "draft"
},
"outputSchema": {
"link": {
"tab": "test",
"icon": "IconVariable",
"label": "Generate Function Output",
"isLeaf": true
},
"_outputSchemaType": "LINK"
}
"settings": {
"input": {
"serverlessFunctionId": "a3d42cb4-a5bd-4eb1-b564-5a207fa98dc6",
"serverlessFunctionInput": {
"a": null,
"b": null
},
"serverlessFunctionVersion": "1"
},
"outputSchema": {
"link": {
"tab": "test",
"icon": "IconVariable",
"label": "Generate Function Output",
"isLeaf": true
},
"_outputSchemaType": "LINK"
}
"settings": {
"input": {
"serverlessFunctionId": "8fdff250-0b7d-4d20-ae5f-8defe4ff9d37",
"serverlessFunctionInput": {
"a": null,
"b": null
},
"serverlessFunctionVersion": "1"
},
"outputSchema": {
"link": {
"tab": "test",
"icon": "IconVariable",
"label": "Generate Function Output",
"isLeaf": true
},
"_outputSchemaType": "LINK"
}
"settings": {
"input": {
"serverlessFunctionId": "a3d42cb4-a5bd-4eb1-b564-5a207fa98dc6",
"serverlessFunctionInput": {
"a": null,
"b": null
},
"serverlessFunctionVersion": "1"
},
"outputSchema": {
"link": {
"tab": "test",
"icon": "IconVariable",
"label": "Generate Function Output",
"isLeaf": true
},
"_outputSchemaType": "LINK"
}
"settings": {
"input": {
"serverlessFunctionId": "798f33e3-04ee-453f-94a1-ddc1518e4618",
"serverlessFunctionInput": {
"a": null,
"b": null
},
"serverlessFunctionVersion": "1"
},
"outputSchema": {
"link": {
"tab": "test",
"icon": "IconVariable",
"label": "Generate Function Output",
"isLeaf": true
},
"_outputSchemaType": "LINK"
}
built-function/42b11c8c-0238-4044-a5d3-9cc662066de4/1/index.mjs
built-function/42b11c8c-0238-4044-a5d3-9cc662066de4/draft/index.mjs
built-function/798f33e3-04ee-453f-94a1-ddc1518e4618/1/index.mjs
built-function/798f33e3-04ee-453f-94a1-ddc1518e4618/draft/index.mjs
built-function/84d3e16c-10bb-4eed-ae15-714f657a2726/draft/index.mjs
built-function/8fdff250-0b7d-4d20-ae5f-8defe4ff9d37/1/index.mjs
built-function/8fdff250-0b7d-4d20-ae5f-8defe4ff9d37/draft/index.mjs
built-function/a3d42cb4-a5bd-4eb1-b564-5a207fa98dc6/1/index.mjs
built-function/a3d42cb4-a5bd-4eb1-b564-5a207fa98dc6/draft/index.mjs
built-function/a7d07444-092f-4a96-8e3c-b379bc49f6e3/draft/index.mjs
built-function/cfc3c41f-5520-41f5-a83e-6a56ca6b90ba/1/index.mjs
built-function/cfc3c41f-5520-41f5-a83e-6a56ca6b90ba/draft/index.mjs
built-function/e8a6a1c8-3e63-4d8b-aaa7-d8c932a6b314/draft/index.mjs
serverless-function/42b11c8c-0238-4044-a5d3-9cc662066de4/1/src/index.ts
serverless-function/42b11c8c-0238-4044-a5d3-9cc662066de4/draft/src/index.ts
serverless-function/798f33e3-04ee-453f-94a1-ddc1518e4618/1/src/index.ts
serverless-function/798f33e3-04ee-453f-94a1-ddc1518e4618/draft/src/index.ts
serverless-function/84d3e16c-10bb-4eed-ae15-714f657a2726/draft/src/index.ts
serverless-function/8fdff250-0b7d-4d20-ae5f-8defe4ff9d37/1/src/index.ts
serverless-function/8fdff250-0b7d-4d20-ae5f-8defe4ff9d37/draft/src/index.ts
serverless-function/a3d42cb4-a5bd-4eb1-b564-5a207fa98dc6/1/src/index.ts
serverless-function/a3d42cb4-a5bd-4eb1-b564-5a207fa98dc6/draft/src/index.ts
serverless-function/a7d07444-092f-4a96-8e3c-b379bc49f6e3/draft/src/index.ts
serverless-function/cfc3c41f-5520-41f5-a83e-6a56ca6b90ba/1/src/index.ts
serverless-function/cfc3c41f-5520-41f5-a83e-6a56ca6b90ba/draft/src/index.ts
serverless-function/e8a6a1c8-3e63-4d8b-aaa7-d8c932a6b314/draft/src/index.ts
@@ -12,12 +12,12 @@ import { REACT_APP_SERVER_BASE_URL } from '~/config';
import { useUpdateEffect } from '~/hooks/useUpdateEffect';
import { isMatchingLocation } from '~/utils/isMatchingLocation';
import { ApolloFactory, type Options } from '@/apollo/services/apollo.factory';
import { currentUserWorkspaceState } from '@/auth/states/currentUserWorkspaceState';
import { appVersionState } from '@/client-config/states/appVersionState';
import { useSnackBar } from '@/ui/feedback/snack-bar-manager/hooks/useSnackBar';
import { AppPath } from 'twenty-shared/types';
import { isDefined } from 'twenty-shared/utils';
import { ApolloFactory, type Options } from '@/apollo/services/apollo.factory';
export const useApolloFactory = (options: Partial<Options<any>> = {}) => {
// eslint-disable-next-line twenty/no-state-useref
@@ -6,7 +6,7 @@ import { getObjectTypename } from '@/object-record/cache/utils/getObjectTypename
import { getRecordFromCache } from '@/object-record/cache/utils/getRecordFromCache';
import { getRecordNodeFromRecord } from '@/object-record/cache/utils/getRecordNodeFromRecord';
import { updateRecordFromCache } from '@/object-record/cache/utils/updateRecordFromCache';
import { generateDepthRecordGqlFieldsFromObject } from '@/object-record/graphql/record-gql-fields/utils/generateDepthRecordGqlFieldsFromObject';
import { generateDepthRecordGqlFieldsFromRecord } from '@/object-record/graphql/record-gql-fields/utils/generateDepthRecordGqlFieldsFromRecord';
import { useObjectPermissions } from '@/object-record/hooks/useObjectPermissions';
import { useRefetchAggregateQueriesForObjectMetadataItem } from '@/object-record/hooks/useRefetchAggregateQueriesForObjectMetadataItem';
import { useUpsertRecordsInStore } from '@/object-record/record-store/hooks/useUpsertRecordsInStore';
@@ -34,12 +34,6 @@ export const useTriggerOptimisticEffectFromSseUpdateEvents = () => {
objectRecordEvents: ObjectRecordEvent[];
objectMetadataItem: ObjectMetadataItem;
}) => {
const recordGqlFields = generateDepthRecordGqlFieldsFromObject({
objectMetadataItem,
objectMetadataItems,
depth: 1,
});
const updateEvents = objectRecordEvents.filter((objectRecordEvent) => {
return objectRecordEvent.action === DatabaseEventAction.UPDATED;
});
@@ -53,6 +47,26 @@ export const useTriggerOptimisticEffectFromSseUpdateEvents = () => {
upsertRecordsInStore({ partialRecords: [updatedRecord] });
const computedOptimisticRecord = {
...computeOptimisticRecordFromInput({
cache: apolloCoreClient.cache,
objectMetadataItem,
objectMetadataItems,
recordInput: updatedRecord,
objectPermissionsByObjectMetadataId,
currentWorkspaceMember: null,
}),
id: updatedRecord.id,
__typename: getObjectTypename(objectMetadataItem.nameSingular),
};
const recordGqlFields = generateDepthRecordGqlFieldsFromRecord({
objectMetadataItem,
objectMetadataItems,
record: computedOptimisticRecord,
depth: 0,
});
const cachedRecord = getRecordFromCache({
cache: apolloCoreClient.cache,
objectMetadataItem,
@@ -77,19 +91,6 @@ export const useTriggerOptimisticEffectFromSseUpdateEvents = () => {
continue;
}
const computedOptimisticRecord = {
...computeOptimisticRecordFromInput({
cache: apolloCoreClient.cache,
objectMetadataItem,
objectMetadataItems,
recordInput: updatedRecord,
objectPermissionsByObjectMetadataId,
currentWorkspaceMember: null,
}),
id: updatedRecord.id,
__typename: getObjectTypename(objectMetadataItem.nameSingular),
};
updateRecordFromCache({
objectMetadataItems,
objectMetadataItem,
@@ -0,0 +1,309 @@
import { Logger } from '@nestjs/common';
import { InjectRepository } from '@nestjs/typeorm';
import * as fs from 'fs/promises';
import { dirname, join } from 'path';
import { isObject } from '@sniptt/guards';
import { Command } from 'nest-commander';
import { FileFolder, type Sources } from 'twenty-shared/types';
import { isDefined } from 'twenty-shared/utils';
import { In, Repository } from 'typeorm';
import { ActiveOrSuspendedWorkspacesMigrationCommandRunner } from 'src/database/commands/command-runners/active-or-suspended-workspaces-migration.command-runner';
import { RunOnWorkspaceArgs } from 'src/database/commands/command-runners/workspaces-migration.command-runner';
import {
collectCodeStepMigrationTargets,
migrateWorkflowCodeStepsWithMapping,
type CodeStepMigrationTarget,
type ServerlessToLogicFunctionMapping,
} from 'src/database/commands/upgrade-version-command/1-17/utils/migrate-workflow-code-step.util';
import { ApplicationService } from 'src/engine/core-modules/application/services/application.service';
import { FileStorageService } from 'src/engine/core-modules/file-storage/file-storage.service';
import { WorkspaceEntity } from 'src/engine/core-modules/workspace/workspace.entity';
import { DataSourceService } from 'src/engine/metadata-modules/data-source/data-source.service';
import { LogicFunctionEntity } from 'src/engine/metadata-modules/logic-function/logic-function.entity';
import { LogicFunctionService } from 'src/engine/metadata-modules/logic-function/services/logic-function.service';
import { GlobalWorkspaceOrmManager } from 'src/engine/twenty-orm/global-workspace-datasource/global-workspace-orm.manager';
import { WorkflowVersionStatus } from 'src/modules/workflow/common/standard-objects/workflow-version.workspace-entity';
const OLD_BUILT_FOLDER = 'built-function';
const OLD_SOURCE_FOLDER = 'serverless-function';
const NEW_WORKFLOW_RESOURCE_PREFIX = 'workflow';
@Command({
name: 'upgrade:1-17:migrate-workflow-code-steps',
description:
'Migrate workflow code steps from v1.16 (serverless) to v1.17 (logic function): create new logic function per (serverlessFunctionId, version), move files, update steps, delete old logic function.',
})
export class MigrateWorkflowCodeStepsCommand extends ActiveOrSuspendedWorkspacesMigrationCommandRunner {
protected readonly logger = new Logger(MigrateWorkflowCodeStepsCommand.name);
constructor(
@InjectRepository(WorkspaceEntity)
protected readonly workspaceRepository: Repository<WorkspaceEntity>,
@InjectRepository(LogicFunctionEntity)
private readonly logicFunctionRepository: Repository<LogicFunctionEntity>,
protected readonly globalWorkspaceOrmManager: GlobalWorkspaceOrmManager,
protected readonly dataSourceService: DataSourceService,
private readonly fileStorageService: FileStorageService,
private readonly applicationService: ApplicationService,
private readonly logicFunctionService: LogicFunctionService,
) {
super(workspaceRepository, globalWorkspaceOrmManager, dataSourceService);
}
override async runOnWorkspace({
workspaceId,
options,
}: RunOnWorkspaceArgs): Promise<void> {
const isDryRun = options.dryRun ?? false;
this.logger.log(
`Running MigrateWorkflowCodeStepsCommand for workspace ${workspaceId}`,
);
const workflowVersionRepository =
await this.globalWorkspaceOrmManager.getRepository(
workspaceId,
'workflowVersion',
{ shouldBypassPermissionChecks: true },
);
const workflowVersions = await workflowVersionRepository.find({
select: ['id', 'steps', 'status'],
where: {
status: In([WorkflowVersionStatus.DRAFT, WorkflowVersionStatus.ACTIVE]),
},
});
const allTargets = new Map<string, CodeStepMigrationTarget>();
for (const version of workflowVersions) {
const targets = collectCodeStepMigrationTargets(version.steps);
for (const target of targets) {
const key = `${target.serverlessFunctionId}:${target.serverlessFunctionVersion}`;
if (!allTargets.has(key)) {
allTargets.set(key, target);
}
}
}
const targetsList = Array.from(allTargets.values());
if (targetsList.length === 0) {
this.logger.log(`No code steps to migrate in workspace ${workspaceId}`);
return;
}
const mapping = new Map<string, string>();
if (!isDryRun) {
for (const target of targetsList) {
const newLogicFunctionId =
await this.createLogicFunctionAndMigrateFiles(workspaceId, target);
if (isDefined(newLogicFunctionId)) {
const key = `${target.serverlessFunctionId}:${target.serverlessFunctionVersion}`;
mapping.set(key, newLogicFunctionId);
}
}
const serverlessToLogicMapping: ServerlessToLogicFunctionMapping = (
serverlessFunctionId,
serverlessFunctionVersion,
) => {
const key = `${serverlessFunctionId}:${serverlessFunctionVersion}`;
const newId = mapping.get(key);
if (!isDefined(newId)) {
throw new Error(
`Missing mapping for ${serverlessFunctionId}:${serverlessFunctionVersion}`,
);
}
return newId;
};
for (const version of workflowVersions) {
const { migratedSteps, hasChanges } =
migrateWorkflowCodeStepsWithMapping(
version.steps,
serverlessToLogicMapping,
);
if (!hasChanges) {
continue;
}
await workflowVersionRepository.update(
{ id: version.id },
{ steps: migratedSteps },
);
this.logger.log(
`Migrated workflow version ${version.id} in workspace ${workspaceId}`,
);
}
} else {
this.logger.log(
`[DRY RUN] Would create ${targetsList.length} new logic function(s), migrate files, update workflow steps, and delete ${new Set(targetsList.map((t) => t.serverlessFunctionId)).size} old logic function(s) in workspace ${workspaceId}`,
);
}
}
private async getApplicationUniversalIdentifier(
applicationId: string,
): Promise<string | null> {
const application = await this.applicationService.findById(applicationId);
return application?.universalIdentifier ?? null;
}
private async createLogicFunctionAndMigrateFiles(
workspaceId: string,
target: CodeStepMigrationTarget,
): Promise<string | null> {
const { serverlessFunctionId, serverlessFunctionVersion } = target;
const version = serverlessFunctionVersion ?? 'draft';
const oldLogicFunction = await this.logicFunctionRepository.findOne({
where: { id: serverlessFunctionId, workspaceId },
});
if (!isDefined(oldLogicFunction)) {
this.logger.warn(
`Logic function ${serverlessFunctionId} not found in workspace ${workspaceId}, skipping`,
);
return null;
}
const tempRoot = await this.migrateFilesFromOldPathToTemp(
workspaceId,
serverlessFunctionId,
version,
);
const newFlatLogicFunction = await this.logicFunctionService.createOne({
input: {
name: oldLogicFunction.name,
description: oldLogicFunction.description ?? undefined,
timeoutSeconds: oldLogicFunction.timeoutSeconds ?? 300,
logicFunctionLayerId: oldLogicFunction.logicFunctionLayerId,
toolInputSchema: oldLogicFunction.toolInputSchema ?? undefined,
isTool: oldLogicFunction.isTool ?? false,
},
workspaceId,
applicationId: oldLogicFunction.applicationId,
});
const newLogicFunctionId = newFlatLogicFunction.id;
const applicationUniversalIdentifier =
await this.getApplicationUniversalIdentifier(
newFlatLogicFunction.applicationId,
);
if (isDefined(applicationUniversalIdentifier)) {
await this.uploadTempToNewPath(
workspaceId,
applicationUniversalIdentifier,
newLogicFunctionId,
tempRoot,
);
}
await fs.rm(tempRoot, { recursive: true, force: true });
this.logger.log(
`Created logic function ${newLogicFunctionId} (from ${serverlessFunctionId}/${version}) and migrated files in workspace ${workspaceId}`,
);
return newLogicFunctionId;
}
private async migrateFilesFromOldPathToTemp(
workspaceId: string,
serverlessFunctionId: string,
version: string,
): Promise<string> {
const workspacePrefix = `workspace-${workspaceId}`;
const oldPaths = {
built: `${workspacePrefix}/${OLD_BUILT_FOLDER}/${serverlessFunctionId}/${version}`,
source: `${workspacePrefix}/${OLD_SOURCE_FOLDER}/${serverlessFunctionId}/${version}`,
};
const tempRoot = await fs.mkdtemp(
`/tmp/twenty-migrate-code-step-${workspaceId}-${serverlessFunctionId}-${version}-`,
);
const builtTempDir = join(tempRoot, 'built');
const sourceTempDir = join(tempRoot, 'source');
await fs.mkdir(builtTempDir, { recursive: true });
await fs.mkdir(sourceTempDir, { recursive: true });
const builtSources = await this.fileStorageService.readFolder(
oldPaths.built,
);
await this.writeSourcesToLocalFolder(builtSources as Sources, builtTempDir);
const sourceSources = await this.fileStorageService.readFolder(
oldPaths.source,
);
const flattened =
(sourceSources.src as Sources) ?? (sourceSources as Sources);
await this.writeSourcesToLocalFolder(flattened, sourceTempDir);
return tempRoot;
}
private async uploadTempToNewPath(
workspaceId: string,
applicationUniversalIdentifier: string,
newLogicFunctionId: string,
tempRoot: string,
): Promise<void> {
const resourcePath = `${NEW_WORKFLOW_RESOURCE_PREFIX}/${newLogicFunctionId}/src`;
const builtTempDir = join(tempRoot, 'built');
const sourceTempDir = join(tempRoot, 'source');
await this.fileStorageService.uploadFolder_v2({
workspaceId,
applicationUniversalIdentifier,
fileFolder: FileFolder.BuiltLogicFunction,
resourcePath,
localPath: builtTempDir,
});
await this.fileStorageService.uploadFolder_v2({
workspaceId,
applicationUniversalIdentifier,
fileFolder: FileFolder.Source,
resourcePath,
localPath: sourceTempDir,
});
}
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);
}
}
}
@@ -0,0 +1,529 @@
import { Logger } from '@nestjs/common';
import { InjectRepository } from '@nestjs/typeorm';
import { v4 as uuidv4 } from 'uuid';
import { Command } from 'nest-commander';
import { Repository } from 'typeorm';
import { ActiveOrSuspendedWorkspacesMigrationCommandRunner } from 'src/database/commands/command-runners/active-or-suspended-workspaces-migration.command-runner';
import { RunOnWorkspaceArgs } from 'src/database/commands/command-runners/workspaces-migration.command-runner';
import { ApplicationService } from 'src/engine/core-modules/application/services/application.service';
import { FileStorageService } from 'src/engine/core-modules/file-storage/file-storage.service';
import { LogicFunctionLayerService } from 'src/engine/core-modules/logic-function/logic-function-layer/services/logic-function-layer.service';
import { RecordPositionService } from 'src/engine/core-modules/record-position/services/record-position.service';
import { WorkspaceEntity } from 'src/engine/core-modules/workspace/workspace.entity';
import { DataSourceService } from 'src/engine/metadata-modules/data-source/data-source.service';
import {
DEFAULT_BUILT_HANDLER_PATH,
DEFAULT_HANDLER_NAME,
DEFAULT_SOURCE_HANDLER_PATH,
LogicFunctionEntity,
LogicFunctionRuntime,
} from 'src/engine/metadata-modules/logic-function/logic-function.entity';
import { GlobalWorkspaceOrmManager } from 'src/engine/twenty-orm/global-workspace-datasource/global-workspace-orm.manager';
import { WorkflowVersionStatus } from 'src/modules/workflow/common/standard-objects/workflow-version.workspace-entity';
import { WorkflowStatus } from 'src/modules/workflow/common/standard-objects/workflow.workspace-entity';
import { WorkflowActionType } from 'src/modules/workflow/workflow-executor/workflow-actions/types/workflow-action-type.enum';
import { WorkflowTriggerType } from 'src/modules/workflow/workflow-trigger/types/workflow-trigger.type';
const OLD_BUILT_FOLDER = 'built-function';
const OLD_SOURCE_FOLDER = 'serverless-function';
const SEED_VERSION_DRAFT = 'draft';
const SEED_VERSION_PUBLISHED = '1';
const OUTPUT_SCHEMA_LINK = {
link: {
tab: 'test',
icon: 'IconVariable',
label: 'Generate Function Output',
isLeaf: true,
},
_outputSchemaType: 'LINK',
} as const;
const OUTPUT_SCHEMA_MESSAGE = {
message: {
type: 'string',
label: 'message',
value: 'Hello, input: null and null',
isLeaf: true,
},
} as const;
@Command({
name: 'upgrade:1-17:seed-workflow-v1-16',
description:
'[Temporary] Clean existing workflow runs, workflow versions, workflows, logic functions and old file storage, then seed 3 scenarios: (1) draft+active with LINK outputSchema, (2) draft-only with message outputSchema, (3) draft+active with mixed outputSchema (message on draft, LINK on active). For testing the 1-17 migrate-workflow-code-steps command.',
})
export class SeedWorkflowV1_16Command extends ActiveOrSuspendedWorkspacesMigrationCommandRunner {
protected readonly logger = new Logger(SeedWorkflowV1_16Command.name);
constructor(
@InjectRepository(WorkspaceEntity)
protected readonly workspaceRepository: Repository<WorkspaceEntity>,
@InjectRepository(LogicFunctionEntity)
private readonly logicFunctionRepository: Repository<LogicFunctionEntity>,
protected readonly globalWorkspaceOrmManager: GlobalWorkspaceOrmManager,
protected readonly dataSourceService: DataSourceService,
private readonly applicationService: ApplicationService,
private readonly logicFunctionLayerService: LogicFunctionLayerService,
private readonly fileStorageService: FileStorageService,
private readonly recordPositionService: RecordPositionService,
) {
super(workspaceRepository, globalWorkspaceOrmManager, dataSourceService);
}
override async runOnWorkspace({
workspaceId,
}: RunOnWorkspaceArgs): Promise<void> {
this.logger.log(`Seeding workflow v1.16 data for workspace ${workspaceId}`);
await this.cleanWorkflowsAndOldFileStorage(workspaceId);
const workflowRepository =
await this.globalWorkspaceOrmManager.getRepository(
workspaceId,
'workflow',
{
shouldBypassPermissionChecks: true,
},
);
const workflowVersionRepository =
await this.globalWorkspaceOrmManager.getRepository(
workspaceId,
'workflowVersion',
{ shouldBypassPermissionChecks: true },
);
await this.seedScenarioDraftAndActiveLink(
workspaceId,
workflowRepository,
workflowVersionRepository,
);
await this.seedScenarioDraftOnlyMessage(
workspaceId,
workflowRepository,
workflowVersionRepository,
);
await this.seedScenarioDraftAndActiveMixedOutputSchema(
workspaceId,
workflowRepository,
workflowVersionRepository,
);
this.logger.log(
`Seeded 3 workflows (draft+active LINK, draft-only message, draft+active mixed outputSchema) in workspace ${workspaceId}. Run upgrade:1-17:migrate-workflow-code-steps to test migration.`,
);
}
private buildCodeStep(
logicFunctionId: string,
version: string,
outputSchema: object,
stepName: string,
) {
return {
id: uuidv4(),
name: stepName,
type: WorkflowActionType.CODE,
settings: {
input: {
serverlessFunctionId: logicFunctionId,
serverlessFunctionInput: {},
serverlessFunctionVersion: version,
},
outputSchema,
},
valid: true,
};
}
private buildTrigger(): {
name: string;
type: string;
settings: { outputSchema: object };
} {
return {
name: 'trigger',
type: WorkflowTriggerType.MANUAL,
settings: { outputSchema: {} },
};
}
private async seedScenarioDraftAndActiveLink(
workspaceId: string,
workflowRepository: Awaited<
ReturnType<GlobalWorkspaceOrmManager['getRepository']>
>,
workflowVersionRepository: Awaited<
ReturnType<GlobalWorkspaceOrmManager['getRepository']>
>,
): Promise<void> {
const logicFunctionId = await this.insertLogicFunctionRow(
workspaceId,
'Seed (draft+active LINK)',
);
await this.writeOldFormatFiles(logicFunctionId, SEED_VERSION_DRAFT);
await this.writeOldFormatFiles(logicFunctionId, SEED_VERSION_PUBLISHED);
const workflowId = uuidv4();
const draftVersionId = uuidv4();
const activeVersionId = uuidv4();
const workflowPosition =
await this.recordPositionService.buildRecordPosition({
value: 'first',
objectMetadata: {
isCustom: false,
nameSingular: 'workflow',
},
workspaceId,
});
const draftVersionPosition =
await this.recordPositionService.buildRecordPosition({
value: 'first',
objectMetadata: {
isCustom: false,
nameSingular: 'workflowVersion',
},
workspaceId,
});
const activeVersionPosition =
await this.recordPositionService.buildRecordPosition({
value: 'last',
objectMetadata: {
isCustom: false,
nameSingular: 'workflowVersion',
},
workspaceId,
});
await workflowRepository.insert({
id: workflowId,
name: 'Seed draft+active (LINK)',
statuses: [WorkflowStatus.DRAFT],
position: workflowPosition,
});
const trigger = this.buildTrigger();
const draftSteps = [
this.buildCodeStep(
logicFunctionId,
SEED_VERSION_DRAFT,
OUTPUT_SCHEMA_LINK,
'Code step (draft)',
),
];
const activeSteps = [
this.buildCodeStep(
logicFunctionId,
SEED_VERSION_PUBLISHED,
OUTPUT_SCHEMA_LINK,
'Code step (v1)',
),
];
await workflowVersionRepository.insert({
id: draftVersionId,
workflowId,
name: 'v1',
status: WorkflowVersionStatus.DRAFT,
trigger,
steps: draftSteps,
position: draftVersionPosition,
});
await workflowVersionRepository.insert({
id: activeVersionId,
workflowId,
name: 'v2',
status: WorkflowVersionStatus.ACTIVE,
trigger,
steps: activeSteps,
position: activeVersionPosition,
});
await workflowRepository.update(workflowId, {
lastPublishedVersionId: activeVersionId,
statuses: [WorkflowStatus.ACTIVE],
});
}
private async seedScenarioDraftOnlyMessage(
workspaceId: string,
workflowRepository: Awaited<
ReturnType<GlobalWorkspaceOrmManager['getRepository']>
>,
workflowVersionRepository: Awaited<
ReturnType<GlobalWorkspaceOrmManager['getRepository']>
>,
): Promise<void> {
const logicFunctionId = await this.insertLogicFunctionRow(
workspaceId,
'Seed (draft-only message)',
);
await this.writeOldFormatFiles(logicFunctionId, SEED_VERSION_DRAFT);
const workflowId = uuidv4();
const draftVersionId = uuidv4();
const workflowPosition =
await this.recordPositionService.buildRecordPosition({
value: 'first',
objectMetadata: {
isCustom: false,
nameSingular: 'workflow',
},
workspaceId,
});
const draftVersionPosition =
await this.recordPositionService.buildRecordPosition({
value: 'first',
objectMetadata: {
isCustom: false,
nameSingular: 'workflowVersion',
},
workspaceId,
});
await workflowRepository.insert({
id: workflowId,
name: 'Seed draft-only (message)',
statuses: [WorkflowStatus.DRAFT],
position: workflowPosition,
});
const trigger = this.buildTrigger();
const draftSteps = [
this.buildCodeStep(
logicFunctionId,
SEED_VERSION_DRAFT,
OUTPUT_SCHEMA_MESSAGE,
'Code step (draft)',
),
];
await workflowVersionRepository.insert({
id: draftVersionId,
workflowId,
name: 'v1',
status: WorkflowVersionStatus.DRAFT,
trigger,
steps: draftSteps,
position: draftVersionPosition,
});
}
private async seedScenarioDraftAndActiveMixedOutputSchema(
workspaceId: string,
workflowRepository: Awaited<
ReturnType<GlobalWorkspaceOrmManager['getRepository']>
>,
workflowVersionRepository: Awaited<
ReturnType<GlobalWorkspaceOrmManager['getRepository']>
>,
): Promise<void> {
const logicFunctionId = await this.insertLogicFunctionRow(
workspaceId,
'Seed (draft+active mixed)',
);
await this.writeOldFormatFiles(logicFunctionId, SEED_VERSION_DRAFT);
await this.writeOldFormatFiles(logicFunctionId, SEED_VERSION_PUBLISHED);
const workflowId = uuidv4();
const draftVersionId = uuidv4();
const activeVersionId = uuidv4();
const workflowPosition =
await this.recordPositionService.buildRecordPosition({
value: 'first',
objectMetadata: {
isCustom: false,
nameSingular: 'workflow',
},
workspaceId,
});
const draftVersionPosition =
await this.recordPositionService.buildRecordPosition({
value: 'first',
objectMetadata: {
isCustom: false,
nameSingular: 'workflowVersion',
},
workspaceId,
});
const activeVersionPosition =
await this.recordPositionService.buildRecordPosition({
value: 'last',
objectMetadata: {
isCustom: false,
nameSingular: 'workflowVersion',
},
workspaceId,
});
await workflowRepository.insert({
id: workflowId,
name: 'Seed draft+active (mixed)',
statuses: [WorkflowStatus.DRAFT],
position: workflowPosition,
});
const trigger = this.buildTrigger();
const draftSteps = [
this.buildCodeStep(
logicFunctionId,
SEED_VERSION_DRAFT,
OUTPUT_SCHEMA_MESSAGE,
'Code step (draft, message)',
),
];
const activeSteps = [
this.buildCodeStep(
logicFunctionId,
SEED_VERSION_PUBLISHED,
OUTPUT_SCHEMA_LINK,
'Code step (v1, LINK)',
),
];
await workflowVersionRepository.insert({
id: draftVersionId,
workflowId,
name: 'v1',
status: WorkflowVersionStatus.DRAFT,
trigger,
steps: draftSteps,
position: draftVersionPosition,
});
await workflowVersionRepository.insert({
id: activeVersionId,
workflowId,
name: 'v2',
status: WorkflowVersionStatus.ACTIVE,
trigger,
steps: activeSteps,
position: activeVersionPosition,
});
await workflowRepository.update(workflowId, {
lastPublishedVersionId: activeVersionId,
statuses: [WorkflowStatus.ACTIVE],
});
}
private async insertLogicFunctionRow(
workspaceId: string,
name: string = 'Seed code step (v1.16)',
): Promise<string> {
const { id: logicFunctionLayerId } =
await this.logicFunctionLayerService.createCommonLayerIfNotExist(
workspaceId,
);
const { workspaceCustomFlatApplication } =
await this.applicationService.findWorkspaceTwentyStandardAndCustomApplicationOrThrow(
{ workspaceId },
);
const applicationId = workspaceCustomFlatApplication.id;
const id = uuidv4();
const universalIdentifier = uuidv4();
const now = new Date();
await this.logicFunctionRepository.insert({
id,
workspaceId,
universalIdentifier,
applicationId,
name,
description: 'Temporary logic function for 1.17 migration testing',
sourceHandlerPath: DEFAULT_SOURCE_HANDLER_PATH,
builtHandlerPath: DEFAULT_BUILT_HANDLER_PATH,
handlerName: DEFAULT_HANDLER_NAME,
runtime: LogicFunctionRuntime.NODE22,
timeoutSeconds: 300,
checksum: null,
toolInputSchema: null,
isTool: false,
logicFunctionLayerId,
cronTriggerSettings: null,
databaseEventTriggerSettings: null,
httpRouteTriggerSettings: null,
createdAt: now,
updatedAt: now,
deletedAt: null,
});
return id;
}
private async cleanWorkflowsAndOldFileStorage(
workspaceId: string,
): Promise<void> {
const workflowRunRepository =
await this.globalWorkspaceOrmManager.getRepository(
workspaceId,
'workflowRun',
{ shouldBypassPermissionChecks: true },
);
const workflowVersionRepository =
await this.globalWorkspaceOrmManager.getRepository(
workspaceId,
'workflowVersion',
{ shouldBypassPermissionChecks: true },
);
const workflowRepository =
await this.globalWorkspaceOrmManager.getRepository(
workspaceId,
'workflow',
{
shouldBypassPermissionChecks: true,
},
);
const deletedRuns = await workflowRunRepository.delete({});
const deletedVersions = await workflowVersionRepository.delete({});
const deletedWorkflows = await workflowRepository.delete({});
const deletedLogicFunctions = await this.logicFunctionRepository.delete({
workspaceId,
});
this.logger.log(
`Cleaned workspace ${workspaceId}: ${deletedRuns.affected ?? 0} workflow run(s), ${deletedVersions.affected ?? 0} workflow version(s), ${deletedWorkflows.affected ?? 0} workflow(s), ${deletedLogicFunctions.affected ?? 0} logic function(s)`,
);
try {
await this.fileStorageService.delete({ folderPath: OLD_BUILT_FOLDER });
await this.fileStorageService.delete({ folderPath: OLD_SOURCE_FOLDER });
this.logger.log(
`Cleaned old file storage: ${OLD_BUILT_FOLDER}, ${OLD_SOURCE_FOLDER}`,
);
} catch (error) {
this.logger.warn(
`Old file storage cleanup skipped (folders may not exist): ${error instanceof Error ? error.message : String(error)}`,
);
}
}
private async writeOldFormatFiles(
logicFunctionId: string,
version: string,
): Promise<void> {
const builtFolder = `${OLD_BUILT_FOLDER}/${logicFunctionId}/${version}`;
const sourceFolder = `${OLD_SOURCE_FOLDER}/${logicFunctionId}/${version}`;
const builtSources = {
'index.mjs':
'export default async function main() { return { message: "ok" }; }\n',
};
const sourceSources = {
src: {
'index.ts':
'export default async function main(): Promise<{ message: string }> {\n return { message: "ok" };\n}\n',
},
};
await this.fileStorageService.writeFolder(builtSources, builtFolder);
await this.fileStorageService.writeFolder(sourceSources, sourceFolder);
}
}
@@ -6,14 +6,21 @@ import { IdentifyWebhookMetadataCommand } from 'src/database/commands/upgrade-ve
import { MakeWebhookUniversalIdentifierAndApplicationIdNotNullableMigrationCommand } from 'src/database/commands/upgrade-version-command/1-17/1-17-make-webhook-universal-identifier-and-application-id-not-nullable-migration.command';
import { MigrateAttachmentToMorphRelationsCommand } from 'src/database/commands/upgrade-version-command/1-17/1-17-migrate-attachment-to-morph-relations.command';
import { MigrateSendEmailRecipientsCommand } from 'src/database/commands/upgrade-version-command/1-17/1-17-migrate-send-email-recipients.command';
import { MigrateWorkflowCodeStepsCommand } from 'src/database/commands/upgrade-version-command/1-17/1-17-migrate-workflow-code-steps.command';
import { SeedWorkflowV1_16Command } from 'src/database/commands/upgrade-version-command/1-17/1-17-seed-workflow-v1-16.command';
import { ApplicationModule } from 'src/engine/core-modules/application/application.module';
import { FeatureFlagEntity } from 'src/engine/core-modules/feature-flag/feature-flag.entity';
import { FeatureFlagModule } from 'src/engine/core-modules/feature-flag/feature-flag.module';
import { FileStorageModule } from 'src/engine/core-modules/file-storage/file-storage.module';
import { FileEntity } from 'src/engine/core-modules/file/entities/file.entity';
import { CoreLogicFunctionLayerModule } from 'src/engine/core-modules/logic-function/logic-function-layer/logic-function-layer.module';
import { RecordPositionModule } from 'src/engine/core-modules/record-position/record-position.module';
import { WorkspaceEntity } from 'src/engine/core-modules/workspace/workspace.entity';
import { DataSourceModule } from 'src/engine/metadata-modules/data-source/data-source.module';
import { FieldMetadataEntity } from 'src/engine/metadata-modules/field-metadata/field-metadata.entity';
import { FieldMetadataModule } from 'src/engine/metadata-modules/field-metadata/field-metadata.module';
import { LogicFunctionEntity } from 'src/engine/metadata-modules/logic-function/logic-function.entity';
import { LogicFunctionModule } from 'src/engine/metadata-modules/logic-function/logic-function.module';
import { ObjectMetadataEntity } from 'src/engine/metadata-modules/object-metadata/object-metadata.entity';
import { ObjectMetadataModule } from 'src/engine/metadata-modules/object-metadata/object-metadata.module';
import { WebhookEntity } from 'src/engine/metadata-modules/webhook/entities/webhook.entity';
@@ -33,15 +40,20 @@ import { AttachmentWorkspaceEntity } from 'src/modules/attachment/standard-objec
AttachmentWorkspaceEntity,
WebhookEntity,
FileEntity,
LogicFunctionEntity,
]),
DataSourceModule,
WorkspaceCacheStorageModule,
WorkspaceMetadataVersionModule,
FeatureFlagModule,
FileStorageModule.forRoot(),
WorkspaceCacheModule,
FieldMetadataModule,
ObjectMetadataModule,
ApplicationModule,
CoreLogicFunctionLayerModule,
LogicFunctionModule,
RecordPositionModule,
GlobalWorkspaceDataSourceModule,
],
providers: [
@@ -50,6 +62,8 @@ import { AttachmentWorkspaceEntity } from 'src/modules/attachment/standard-objec
MakeWebhookUniversalIdentifierAndApplicationIdNotNullableMigrationCommand,
DeleteFileRecordsCommand,
MigrateSendEmailRecipientsCommand,
MigrateWorkflowCodeStepsCommand,
SeedWorkflowV1_16Command,
],
exports: [
MigrateAttachmentToMorphRelationsCommand,
@@ -57,6 +71,8 @@ import { AttachmentWorkspaceEntity } from 'src/modules/attachment/standard-objec
MakeWebhookUniversalIdentifierAndApplicationIdNotNullableMigrationCommand,
DeleteFileRecordsCommand,
MigrateSendEmailRecipientsCommand,
MigrateWorkflowCodeStepsCommand,
SeedWorkflowV1_16Command,
],
})
export class V1_17_UpgradeVersionCommandModule {}
@@ -0,0 +1,125 @@
import { isDefined } from 'twenty-shared/utils';
import { WorkflowActionType } from 'src/modules/workflow/workflow-executor/workflow-actions/types/workflow-action.type';
type LegacyCodeStepInput = {
serverlessFunctionId: string;
serverlessFunctionVersion?: string;
serverlessFunctionInput?: Record<string, unknown>;
};
type MigratedCodeStepInput = {
logicFunctionId: string;
logicFunctionInput?: Record<string, unknown>;
};
type WorkflowStep = {
id: string;
type: string;
settings: {
input: LegacyCodeStepInput | MigratedCodeStepInput;
outputSchema?: unknown;
};
};
export const needsCodeStepMigration = (
input: LegacyCodeStepInput | MigratedCodeStepInput,
): input is LegacyCodeStepInput => {
return (
'serverlessFunctionId' in input && isDefined(input.serverlessFunctionId)
);
};
export type CodeStepMigrationTarget = {
serverlessFunctionId: string;
serverlessFunctionVersion: string;
};
export type ServerlessToLogicFunctionMapping = (
serverlessFunctionId: string,
serverlessFunctionVersion: string,
) => string;
export const collectCodeStepMigrationTargets = (
steps: unknown,
): CodeStepMigrationTarget[] => {
if (!isDefined(steps) || !Array.isArray(steps) || steps.length === 0) {
return [];
}
const typedSteps = steps as WorkflowStep[];
const seen = new Set<string>();
const targets: CodeStepMigrationTarget[] = [];
for (const step of typedSteps) {
if (step.type !== WorkflowActionType.CODE) {
continue;
}
const input = step.settings.input as
| LegacyCodeStepInput
| MigratedCodeStepInput;
if (!needsCodeStepMigration(input)) {
continue;
}
const version = input.serverlessFunctionVersion ?? 'draft';
const key = `${input.serverlessFunctionId}:${version}`;
if (seen.has(key)) {
continue;
}
seen.add(key);
targets.push({
serverlessFunctionId: input.serverlessFunctionId,
serverlessFunctionVersion: version,
});
}
return targets;
};
export const migrateWorkflowCodeStepsWithMapping = (
steps: unknown,
mapping: ServerlessToLogicFunctionMapping,
): { migratedSteps: WorkflowStep[]; hasChanges: boolean } => {
if (!isDefined(steps) || !Array.isArray(steps) || steps.length === 0) {
return { migratedSteps: [], hasChanges: false };
}
const typedSteps = steps as WorkflowStep[];
let hasChanges = false;
const migratedSteps = typedSteps.map((step) => {
if (step.type !== WorkflowActionType.CODE) {
return step;
}
const input = step.settings.input as
| LegacyCodeStepInput
| MigratedCodeStepInput;
if (!needsCodeStepMigration(input)) {
return step;
}
hasChanges = true;
const version = input.serverlessFunctionVersion ?? 'draft';
const newLogicFunctionId = mapping(input.serverlessFunctionId, version);
return {
...step,
settings: {
...step.settings,
input: {
logicFunctionId: newLogicFunctionId,
logicFunctionInput: input.serverlessFunctionInput ?? {},
} as MigratedCodeStepInput,
},
};
});
return { migratedSteps, hasChanges };
};
@@ -13,6 +13,7 @@ import { DeleteFileRecordsCommand } from 'src/database/commands/upgrade-version-
import { IdentifyWebhookMetadataCommand } from 'src/database/commands/upgrade-version-command/1-17/1-17-identify-webhook-metadata.command';
import { MakeWebhookUniversalIdentifierAndApplicationIdNotNullableMigrationCommand } from 'src/database/commands/upgrade-version-command/1-17/1-17-make-webhook-universal-identifier-and-application-id-not-nullable-migration.command';
import { MigrateAttachmentToMorphRelationsCommand } from 'src/database/commands/upgrade-version-command/1-17/1-17-migrate-attachment-to-morph-relations.command';
import { MigrateWorkflowCodeStepsCommand } from 'src/database/commands/upgrade-version-command/1-17/1-17-migrate-workflow-code-steps.command';
import { TwentyConfigService } from 'src/engine/core-modules/twenty-config/twenty-config.service';
import { WorkspaceEntity } from 'src/engine/core-modules/workspace/workspace.entity';
import { DataSourceService } from 'src/engine/metadata-modules/data-source/data-source.service';
@@ -37,6 +38,7 @@ export class UpgradeCommand extends UpgradeCommandRunner {
protected readonly migrateAttachmentToMorphRelationsCommand: MigrateAttachmentToMorphRelationsCommand,
protected readonly identifyWebhookMetadataCommand: IdentifyWebhookMetadataCommand,
protected readonly makeWebhookUniversalIdentifierAndApplicationIdNotNullableMigrationCommand: MakeWebhookUniversalIdentifierAndApplicationIdNotNullableMigrationCommand,
protected readonly migrateWorkflowCodeStepsCommand: MigrateWorkflowCodeStepsCommand,
) {
super(
workspaceRepository,
@@ -53,6 +55,7 @@ export class UpgradeCommand extends UpgradeCommandRunner {
this.identifyWebhookMetadataCommand,
this
.makeWebhookUniversalIdentifierAndApplicationIdNotNullableMigrationCommand,
this.migrateWorkflowCodeStepsCommand,
this.deleteFileRecordsCommand,
];