Remove all saves from workflows (#14283)

- replace save by insert
- if the insert output is needed, cast the generatedMaps
- remove transactions from trigger services. Doing it manually would be
complex and not reliable
This commit is contained in:
Thomas Trompette
2025-09-03 15:41:03 +02:00
committed by GitHub
parent 6ee41f33c6
commit 7ac7f52510
8 changed files with 55 additions and 149 deletions
@@ -47,18 +47,16 @@ export class WorkflowCreateManyPostQueryHook
workspaceId: workspace.id,
});
const workflowVersionsToCreate = payload.map((workflow) => {
return workflowVersionRepository.create({
workflowId: workflow.id,
status: WorkflowVersionStatus.DRAFT,
name: 'v1',
position,
});
});
const workflowVersionsToCreate = payload.map((workflow) => ({
workflowId: workflow.id,
status: WorkflowVersionStatus.DRAFT,
name: 'v1',
position,
}));
await Promise.all(
workflowVersionsToCreate.map((workflowVersion) => {
return workflowVersionRepository.save(workflowVersion);
return workflowVersionRepository.insert(workflowVersion);
}),
);
}
@@ -49,13 +49,11 @@ export class WorkflowCreateOnePostQueryHook
workspaceId: workspace.id,
});
const workflowVersionToCreate = workflowVersionRepository.create({
await workflowVersionRepository.insert({
workflowId: workflow.id,
status: WorkflowVersionStatus.DRAFT,
name: 'v1',
position,
});
await workflowVersionRepository.save(workflowVersionToCreate);
}
}
@@ -84,12 +84,15 @@ export class WorkflowVersionWorkspaceService {
workspaceId,
});
draftWorkflowVersion = await workflowVersionRepository.save({
const insertResult = await workflowVersionRepository.insert({
workflowId,
name: `v${workflowVersionsCount + 1}`,
status: WorkflowVersionStatus.DRAFT,
position,
});
draftWorkflowVersion = insertResult
.generatedMaps[0] as WorkflowVersionWorkspaceEntity;
}
assertWorkflowVersionIsDraft(draftWorkflowVersion);
@@ -103,7 +103,7 @@ export class CreateRecordWorkflowAction implements WorkflowAction {
objectMetadataMapItem: objectMetadataItemWithFieldsMaps,
});
const objectRecord = await repository.save({
const insertResult = await repository.insert({
...transformedObjectRecord,
position,
createdBy: {
@@ -112,8 +112,10 @@ export class CreateRecordWorkflowAction implements WorkflowAction {
},
});
const [createdRecord] = insertResult.generatedMaps;
return {
result: objectRecord,
result: createdRecord,
};
}
}
@@ -115,7 +115,7 @@ export class WorkflowRunWorkspaceService {
? parseInt(workflowRunCountMatch[1], 10)
: 0;
const workflowRun = workflowRunRepository.create({
const workflowRun = {
id: workflowRunId ?? v4(),
name: `#${workflowRunCount + 1} - ${workflow.name}`,
workflowVersionId,
@@ -124,9 +124,8 @@ export class WorkflowRunWorkspaceService {
status,
position,
state: initState,
enqueuedAt:
status === WorkflowRunStatus.ENQUEUED ? new Date().toISOString() : null,
});
enqueuedAt: status === WorkflowRunStatus.ENQUEUED ? new Date() : null,
};
await workflowRunRepository.insert(workflowRun);
@@ -418,16 +418,16 @@ This is the most efficient way for AI to create workflows as it handles all the
workspaceId,
});
const workflow = workflowRepository.create({
const workflow = {
id: uuidv4(),
name,
statuses: [WorkflowStatus.DRAFT],
position: workflowPosition,
});
};
const savedWorkflow = await workflowRepository.save(workflow);
await workflowRepository.insert(workflow);
return savedWorkflow.id;
return workflow.id;
}
private async createWorkflowVersion({
@@ -460,7 +460,7 @@ This is the most efficient way for AI to create workflows as it handles all the
workspaceId,
});
const workflowVersion = workflowVersionRepository.create({
const workflowVersion = {
id: uuidv4(),
workflowId,
name: 'v1',
@@ -468,12 +468,11 @@ This is the most efficient way for AI to create workflows as it handles all the
trigger,
steps,
position: versionPosition,
});
};
const savedWorkflowVersion =
await workflowVersionRepository.save(workflowVersion);
await workflowVersionRepository.insert(workflowVersion);
return savedWorkflowVersion.id;
return workflowVersion.id;
}
private async updateWorkflowStatus({
@@ -1,6 +1,5 @@
import { Injectable } from '@nestjs/common';
import { type WorkspaceEntityManager } from 'src/engine/twenty-orm/entity-manager/workspace-entity-manager';
import { TwentyORMManager } from 'src/engine/twenty-orm/twenty-orm.manager';
import {
type AutomatedTriggerType,
@@ -14,12 +13,10 @@ export class AutomatedTriggerWorkspaceService {
async addAutomatedTrigger({
workflowId,
manager,
type,
settings,
}: {
workflowId: string;
manager: WorkspaceEntityManager;
type: AutomatedTriggerType;
settings: AutomatedTriggerSettings;
}) {
@@ -28,31 +25,19 @@ export class AutomatedTriggerWorkspaceService {
'workflowAutomatedTrigger',
);
const workflowAutomatedTrigger = workflowAutomatedTriggerRepository.create({
await workflowAutomatedTriggerRepository.insert({
type,
settings,
workflowId,
});
await workflowAutomatedTriggerRepository.save(
workflowAutomatedTrigger,
{},
manager,
);
}
async deleteAutomatedTrigger({
workflowId,
manager,
}: {
workflowId: string;
manager: WorkspaceEntityManager;
}) {
async deleteAutomatedTrigger({ workflowId }: { workflowId: string }) {
const workflowAutomatedTriggerRepository =
await this.twentyORMManager.getRepository<WorkflowAutomatedTriggerWorkspaceEntity>(
'workflowAutomatedTrigger',
);
await workflowAutomatedTriggerRepository.delete({ workflowId }, manager);
await workflowAutomatedTriggerRepository.delete({ workflowId });
}
}
@@ -3,7 +3,6 @@ import { Injectable } from '@nestjs/common';
import { t } from '@lingui/core/macro';
import { type ActorMetadata } from 'src/engine/metadata-modules/field-metadata/composite-types/actor.composite-type';
import { type WorkspaceEntityManager } from 'src/engine/twenty-orm/entity-manager/workspace-entity-manager';
import { ScopedWorkspaceContextFactory } from 'src/engine/twenty-orm/factories/scoped-workspace-context.factory';
import { type WorkspaceRepository } from 'src/engine/twenty-orm/repository/workspace.repository';
import { TwentyORMGlobalManager } from 'src/engine/twenty-orm/twenty-orm-global.manager';
@@ -117,84 +116,30 @@ export class WorkflowTriggerWorkspaceService {
assertVersionCanBeActivated(workflowVersion, workflow);
const workspaceDataSource =
await this.twentyORMGlobalManager.getDataSourceForWorkspace({
workspaceId: this.getWorkspaceId(),
});
const queryRunner = workspaceDataSource.createQueryRunner();
await this.performActivationSteps(
workflow,
workflowVersion,
workflowRepository,
workflowVersionRepository,
);
await queryRunner.connect();
await queryRunner.startTransaction();
const manager = queryRunner.manager;
try {
await this.performActivationSteps(
workflow,
workflowVersion,
workflowRepository,
workflowVersionRepository,
manager,
);
await queryRunner.commitTransaction();
return true;
} catch (error) {
if (queryRunner.isTransactionActive) {
try {
await queryRunner.rollbackTransaction();
} catch (error) {
// eslint-disable-next-line no-console
console.trace(`Failed to rollback transaction: ${error.message}`);
}
}
throw error;
} finally {
await queryRunner.release();
}
return true;
}
async deactivateWorkflowVersion(workflowVersionId: string) {
const workspaceDataSource =
await this.twentyORMGlobalManager.getDataSourceForWorkspace({
workspaceId: this.getWorkspaceId(),
});
const queryRunner = workspaceDataSource.createQueryRunner();
await queryRunner.connect();
await queryRunner.startTransaction();
try {
const workflowVersionRepository =
await this.twentyORMGlobalManager.getRepositoryForWorkspace<WorkflowVersionWorkspaceEntity>(
this.getWorkspaceId(),
'workflowVersion',
{ shouldBypassPermissionChecks: true },
);
await this.performDeactivationSteps(
workflowVersionId,
workflowVersionRepository,
queryRunner.manager,
const workflowVersionRepository =
await this.twentyORMGlobalManager.getRepositoryForWorkspace<WorkflowVersionWorkspaceEntity>(
this.getWorkspaceId(),
'workflowVersion',
{ shouldBypassPermissionChecks: true },
);
await queryRunner.commitTransaction();
await this.performDeactivationSteps(
workflowVersionId,
workflowVersionRepository,
);
return true;
} catch (error) {
if (queryRunner.isTransactionActive) {
try {
await queryRunner.rollbackTransaction();
} catch (error) {
// eslint-disable-next-line no-console
console.trace(`Failed to rollback transaction: ${error.message}`);
}
}
throw error;
} finally {
await queryRunner.release();
}
return true;
}
private async performActivationSteps(
@@ -202,7 +147,6 @@ export class WorkflowTriggerWorkspaceService {
workflowVersion: WorkflowVersionWorkspaceEntity,
workflowRepository: WorkspaceRepository<WorkflowWorkspaceEntity>,
workflowVersionRepository: WorkspaceRepository<WorkflowVersionWorkspaceEntity>,
manager: WorkspaceEntityManager,
) {
if (
workflow.lastPublishedVersionId &&
@@ -211,7 +155,6 @@ export class WorkflowTriggerWorkspaceService {
await this.performDeactivationSteps(
workflow.lastPublishedVersionId,
workflowVersionRepository,
manager,
);
}
@@ -220,22 +163,19 @@ export class WorkflowTriggerWorkspaceService {
workflowVersion.id,
workflowRepository,
workflowVersionRepository,
manager,
);
await this.setActiveVersionStatus(
workflowVersion,
workflowVersionRepository,
manager,
);
await this.enableTrigger(workflowVersion, manager);
await this.enableTrigger(workflowVersion);
}
private async performDeactivationSteps(
workflowVersionId: string,
workflowVersionRepository: WorkspaceRepository<WorkflowVersionWorkspaceEntity>,
manager: WorkspaceEntityManager,
) {
const workflowVersionNullable = await workflowVersionRepository.findOne({
where: { id: workflowVersionId },
@@ -253,26 +193,21 @@ export class WorkflowTriggerWorkspaceService {
await this.setDeactivatedVersionStatus(
workflowVersion,
workflowVersionRepository,
manager,
);
await this.disableTrigger(workflowVersion, manager);
await this.disableTrigger(workflowVersion);
}
private async setActiveVersionStatus(
workflowVersion: WorkflowVersionWorkspaceEntity,
workflowVersionRepository: WorkspaceRepository<WorkflowVersionWorkspaceEntity>,
manager: WorkspaceEntityManager,
) {
const activeWorkflowVersions = await workflowVersionRepository.find(
{
where: {
workflowId: workflowVersion.workflowId,
status: WorkflowVersionStatus.ACTIVE,
},
const activeWorkflowVersions = await workflowVersionRepository.find({
where: {
workflowId: workflowVersion.workflowId,
status: WorkflowVersionStatus.ACTIVE,
},
manager,
);
});
if (activeWorkflowVersions.length > 0) {
throw new WorkflowTriggerException(
@@ -287,7 +222,6 @@ export class WorkflowTriggerWorkspaceService {
await workflowVersionRepository.update(
{ id: workflowVersion.id },
{ status: WorkflowVersionStatus.ACTIVE },
manager,
);
await this.emitStatusUpdateEvents(
@@ -300,7 +234,6 @@ export class WorkflowTriggerWorkspaceService {
private async setDeactivatedVersionStatus(
workflowVersion: WorkflowVersionWorkspaceEntity,
workflowVersionRepository: WorkspaceRepository<WorkflowVersionWorkspaceEntity>,
manager: WorkspaceEntityManager,
) {
if (workflowVersion.status !== WorkflowVersionStatus.ACTIVE) {
throw new WorkflowTriggerException(
@@ -315,7 +248,6 @@ export class WorkflowTriggerWorkspaceService {
await workflowVersionRepository.update(
{ id: workflowVersion.id },
{ status: WorkflowVersionStatus.DEACTIVATED },
manager,
);
await this.emitStatusUpdateEvents(
@@ -330,7 +262,6 @@ export class WorkflowTriggerWorkspaceService {
newPublishedVersionId: string,
workflowRepository: WorkspaceRepository<WorkflowWorkspaceEntity>,
workflowVersionRepository: WorkspaceRepository<WorkflowVersionWorkspaceEntity>,
manager: WorkspaceEntityManager,
) {
if (workflow.lastPublishedVersionId === newPublishedVersionId) {
return;
@@ -340,21 +271,16 @@ export class WorkflowTriggerWorkspaceService {
await workflowVersionRepository.update(
{ id: workflow.lastPublishedVersionId },
{ status: WorkflowVersionStatus.ARCHIVED },
manager,
);
}
await workflowRepository.update(
{ id: workflow.id },
{ lastPublishedVersionId: newPublishedVersionId },
manager,
);
}
private async enableTrigger(
workflowVersion: WorkflowVersionWorkspaceEntity,
manager: WorkspaceEntityManager,
) {
private async enableTrigger(workflowVersion: WorkflowVersionWorkspaceEntity) {
assertWorkflowVersionTriggerIsDefined(workflowVersion);
switch (workflowVersion.trigger.type) {
@@ -369,7 +295,6 @@ export class WorkflowTriggerWorkspaceService {
workflowId: workflowVersion.workflowId,
type: AutomatedTriggerType.DATABASE_EVENT,
settings,
manager,
});
return;
@@ -381,7 +306,6 @@ export class WorkflowTriggerWorkspaceService {
workflowId: workflowVersion.workflowId,
type: AutomatedTriggerType.CRON,
settings: { pattern },
manager,
});
return;
@@ -394,7 +318,6 @@ export class WorkflowTriggerWorkspaceService {
private async disableTrigger(
workflowVersion: WorkflowVersionWorkspaceEntity,
manager: WorkspaceEntityManager,
) {
assertWorkflowVersionTriggerIsDefined(workflowVersion);
@@ -403,7 +326,6 @@ export class WorkflowTriggerWorkspaceService {
case WorkflowTriggerType.CRON:
await this.automatedTriggerWorkspaceService.deleteAutomatedTrigger({
workflowId: workflowVersion.workflowId,
manager,
});
return;