diff --git a/packages/twenty-server/src/engine/core-modules/workflow/services/workflow-version-core-sync.service.ts b/packages/twenty-server/src/engine/core-modules/workflow/services/workflow-version-core-sync.service.ts index 800dc6bbd3..a68cd40cc4 100644 --- a/packages/twenty-server/src/engine/core-modules/workflow/services/workflow-version-core-sync.service.ts +++ b/packages/twenty-server/src/engine/core-modules/workflow/services/workflow-version-core-sync.service.ts @@ -280,6 +280,39 @@ export class WorkflowVersionCoreSyncService { await this.invalidateAutomatedTriggerMaps(workspaceId); } + async deleteCoreVersionsByWorkspaceVersionIds( + workspaceId: string, + workflowVersionIds: string[], + ): Promise { + if (workflowVersionIds.length === 0) { + return; + } + + const coreWorkflowVersionIds = + await this.globalWorkspaceOrmManager.executeInWorkspaceContext( + async () => { + const workflowVersionRepository = + await this.globalWorkspaceOrmManager.getRepository( + workspaceId, + 'workflowVersion', + { shouldBypassPermissionChecks: true }, + ); + + const versions = await workflowVersionRepository.find({ + where: { id: In(workflowVersionIds) }, + withDeleted: true, + }); + + return versions + .map((version) => version.coreWorkflowVersionId) + .filter(isNonEmptyString); + }, + buildSystemAuthContext(workspaceId), + ); + + await this.deleteFromCore(workspaceId, coreWorkflowVersionIds); + } + async recreateCoreVersionsByWorkflowId( workspaceId: string, workflowId: string, diff --git a/packages/twenty-server/src/modules/workflow/common/query-hooks/workflow-query-hook.module.ts b/packages/twenty-server/src/modules/workflow/common/query-hooks/workflow-query-hook.module.ts index f766ed6d86..3af67e4d52 100644 --- a/packages/twenty-server/src/modules/workflow/common/query-hooks/workflow-query-hook.module.ts +++ b/packages/twenty-server/src/modules/workflow/common/query-hooks/workflow-query-hook.module.ts @@ -33,6 +33,7 @@ import { WorkflowUpdateOnePreQueryHook } from 'src/modules/workflow/common/query import { WorkflowVersionCreateManyPreQueryHook } from 'src/modules/workflow/common/query-hooks/workflow-version-create-many.pre-query.hook'; import { WorkflowVersionCreateOnePreQueryHook } from 'src/modules/workflow/common/query-hooks/workflow-version-create-one.pre-query.hook'; import { WorkflowVersionDeleteManyPreQueryHook } from 'src/modules/workflow/common/query-hooks/workflow-version-delete-many.pre-query.hook'; +import { WorkflowVersionDeleteOnePostQueryHook } from 'src/modules/workflow/common/query-hooks/workflow-version-delete-one.post-query.hook'; import { WorkflowVersionDeleteOnePreQueryHook } from 'src/modules/workflow/common/query-hooks/workflow-version-delete-one.pre-query.hook'; import { WorkflowVersionDestroyManyPreQueryHook } from 'src/modules/workflow/common/query-hooks/workflow-version-destroy-many.pre-query.hook'; import { WorkflowVersionDestroyOnePreQueryHook } from 'src/modules/workflow/common/query-hooks/workflow-version-destroy-one.pre-query.hook'; @@ -74,6 +75,7 @@ import { WorkflowVersionValidationWorkspaceService } from 'src/modules/workflow/ WorkflowVersionUpdateOnePreQueryHook, WorkflowVersionUpdateManyPreQueryHook, WorkflowVersionDeleteOnePreQueryHook, + WorkflowVersionDeleteOnePostQueryHook, WorkflowVersionDeleteManyPreQueryHook, WorkflowVersionDestroyOnePreQueryHook, WorkflowVersionDestroyManyPreQueryHook, diff --git a/packages/twenty-server/src/modules/workflow/common/query-hooks/workflow-version-delete-one.post-query.hook.ts b/packages/twenty-server/src/modules/workflow/common/query-hooks/workflow-version-delete-one.post-query.hook.ts new file mode 100644 index 0000000000..c669007c9c --- /dev/null +++ b/packages/twenty-server/src/modules/workflow/common/query-hooks/workflow-version-delete-one.post-query.hook.ts @@ -0,0 +1,35 @@ +import { assertIsDefinedOrThrow } from 'twenty-shared/utils'; + +import { type WorkspacePostQueryHookInstance } from 'src/engine/api/graphql/workspace-query-runner/workspace-query-hook/interfaces/workspace-query-hook.interface'; + +import { WorkspaceQueryHook } from 'src/engine/api/graphql/workspace-query-runner/workspace-query-hook/decorators/workspace-query-hook.decorator'; +import { WorkspaceQueryHookType } from 'src/engine/api/graphql/workspace-query-runner/workspace-query-hook/types/workspace-query-hook.type'; +import { type WorkspaceAuthContext } from 'src/engine/core-modules/auth/types/workspace-auth-context.type'; +import { WorkflowVersionCoreSyncService } from 'src/engine/core-modules/workflow/services/workflow-version-core-sync.service'; +import { WorkspaceNotFoundDefaultError } from 'src/engine/core-modules/workspace/workspace.exception'; +import { type WorkflowVersionWorkspaceEntity } from 'src/modules/workflow/common/standard-objects/workflow-version.workspace-entity'; + +@WorkspaceQueryHook({ + key: `workflowVersion.deleteOne`, + type: WorkspaceQueryHookType.POST_HOOK, +}) +export class WorkflowVersionDeleteOnePostQueryHook implements WorkspacePostQueryHookInstance { + constructor( + private readonly workflowVersionCoreSyncService: WorkflowVersionCoreSyncService, + ) {} + + async execute( + authContext: WorkspaceAuthContext, + _objectName: string, + payload: WorkflowVersionWorkspaceEntity[], + ): Promise { + const workspace = authContext.workspace; + + assertIsDefinedOrThrow(workspace, WorkspaceNotFoundDefaultError); + + await this.workflowVersionCoreSyncService.deleteCoreVersionsByWorkspaceVersionIds( + workspace.id, + payload.map((workflowVersion) => workflowVersion.id), + ); + } +} diff --git a/packages/twenty-server/src/modules/workflow/workflow-version-core-sync/listeners/workflow-version-core-dual-write.listener.ts b/packages/twenty-server/src/modules/workflow/workflow-version-core-sync/listeners/workflow-version-core-dual-write.listener.ts deleted file mode 100644 index 590ac59c6e..0000000000 --- a/packages/twenty-server/src/modules/workflow/workflow-version-core-sync/listeners/workflow-version-core-dual-write.listener.ts +++ /dev/null @@ -1,103 +0,0 @@ -import { Injectable } from '@nestjs/common'; - -import { - type ObjectRecordDeleteEvent, - type ObjectRecordDestroyEvent, - type ObjectRecordRestoreEvent, -} from 'twenty-shared/database-events'; -import { isDefined } from 'twenty-shared/utils'; - -import { OnDatabaseBatchEvent } from 'src/engine/api/graphql/graphql-query-runner/decorators/on-database-batch-event.decorator'; -import { DatabaseEventAction } from 'src/engine/api/graphql/graphql-query-runner/enums/database-event-action'; -import { ExceptionHandlerService } from 'src/engine/core-modules/exception-handler/exception-handler.service'; -import { WorkflowVersionCoreSyncService } from 'src/engine/core-modules/workflow/services/workflow-version-core-sync.service'; -import { type CustomWorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/custom-workspace-batch-event.type'; -import { type WorkflowVersionWorkspaceEntity } from 'src/modules/workflow/common/standard-objects/workflow-version.workspace-entity'; - -@Injectable() -export class WorkflowVersionCoreDualWriteListener { - constructor( - private readonly exceptionHandlerService: ExceptionHandlerService, - private readonly workflowVersionCoreSyncService: WorkflowVersionCoreSyncService, - ) {} - - @OnDatabaseBatchEvent('workflowVersion', DatabaseEventAction.RESTORED) - async handleRestored( - batchEvent: CustomWorkspaceEventBatch< - ObjectRecordRestoreEvent - >, - ): Promise { - await this.upsertToCore( - batchEvent.workspaceId, - batchEvent.events.map((event) => event.properties.after), - ); - } - - @OnDatabaseBatchEvent('workflowVersion', DatabaseEventAction.DELETED) - async handleDeleted( - batchEvent: CustomWorkspaceEventBatch< - ObjectRecordDeleteEvent - >, - ): Promise { - await this.deleteFromCore( - batchEvent.workspaceId, - batchEvent.events - .map((event) => event.properties.before.coreWorkflowVersionId) - .filter(isDefined), - ); - } - - @OnDatabaseBatchEvent('workflowVersion', DatabaseEventAction.DESTROYED) - async handleDestroyed( - batchEvent: CustomWorkspaceEventBatch< - ObjectRecordDestroyEvent - >, - ): Promise { - await this.deleteFromCore( - batchEvent.workspaceId, - batchEvent.events - .map((event) => event.properties.before.coreWorkflowVersionId) - .filter(isDefined), - ); - } - - private async upsertToCore( - workspaceId: string | undefined, - workflowVersions: WorkflowVersionWorkspaceEntity[], - ): Promise { - if (!isDefined(workspaceId)) { - return; - } - - try { - await this.workflowVersionCoreSyncService.upsertToCore( - workspaceId, - workflowVersions, - ); - } catch (error) { - this.exceptionHandlerService.captureExceptions([error], { - workspace: { id: workspaceId }, - }); - } - } - - private async deleteFromCore( - workspaceId: string | undefined, - coreWorkflowVersionIds: string[], - ): Promise { - if (!isDefined(workspaceId)) { - return; - } - - try { - await this.workflowVersionCoreSyncService.deleteFromCore( - workspaceId, - coreWorkflowVersionIds, - ); - } catch (error) { - this.exceptionHandlerService.captureExceptions([error], { - workspace: { id: workspaceId }, - }); - } - } -} diff --git a/packages/twenty-server/src/modules/workflow/workflow-version-core-sync/workflow-version-core-sync.module.ts b/packages/twenty-server/src/modules/workflow/workflow-version-core-sync/workflow-version-core-sync.module.ts deleted file mode 100644 index 19b15565d7..0000000000 --- a/packages/twenty-server/src/modules/workflow/workflow-version-core-sync/workflow-version-core-sync.module.ts +++ /dev/null @@ -1,10 +0,0 @@ -import { Module } from '@nestjs/common'; - -import { WorkflowVersionCoreModule } from 'src/engine/core-modules/workflow/workflow-version-core.module'; -import { WorkflowVersionCoreDualWriteListener } from 'src/modules/workflow/workflow-version-core-sync/listeners/workflow-version-core-dual-write.listener'; - -@Module({ - imports: [WorkflowVersionCoreModule], - providers: [WorkflowVersionCoreDualWriteListener], -}) -export class WorkflowVersionCoreSyncModule {} diff --git a/packages/twenty-server/src/modules/workflow/workflow.module.ts b/packages/twenty-server/src/modules/workflow/workflow.module.ts index 12b79c9117..1cc76b601e 100644 --- a/packages/twenty-server/src/modules/workflow/workflow.module.ts +++ b/packages/twenty-server/src/modules/workflow/workflow.module.ts @@ -3,13 +3,11 @@ import { Module } from '@nestjs/common'; import { WorkflowCoreSyncModule } from 'src/modules/workflow/workflow-core-sync/workflow-core-sync.module'; import { WorkflowStatusModule } from 'src/modules/workflow/workflow-status/workflow-status.module'; import { WorkflowTriggerModule } from 'src/modules/workflow/workflow-trigger/workflow-trigger.module'; -import { WorkflowVersionCoreSyncModule } from 'src/modules/workflow/workflow-version-core-sync/workflow-version-core-sync.module'; @Module({ imports: [ WorkflowTriggerModule, WorkflowStatusModule, - WorkflowVersionCoreSyncModule, WorkflowCoreSyncModule, ], }) diff --git a/packages/twenty-server/test/integration/graphql/suites/workflow/workflow-version-discard-draft-core-mirror.integration-spec.ts b/packages/twenty-server/test/integration/graphql/suites/workflow/workflow-version-discard-draft-core-mirror.integration-spec.ts new file mode 100644 index 0000000000..15da9efcc5 --- /dev/null +++ b/packages/twenty-server/test/integration/graphql/suites/workflow/workflow-version-discard-draft-core-mirror.integration-spec.ts @@ -0,0 +1,176 @@ +import request from 'supertest'; +import { updateWorkflowVersionTrigger } from 'test/integration/graphql/suites/workflow/utils/update-workflow-version-trigger.util'; + +import { SEED_APPLE_WORKSPACE_ID } from 'src/engine/workspace-manager/dev-seeder/core/constants/seeder-workspaces.constant'; + +const client = request(`http://localhost:${APP_PORT}`); + +const graphql = (query: string, variables?: object) => + client + .post('/graphql') + .set('Authorization', `Bearer ${APPLE_JANE_ADMIN_ACCESS_TOKEN}`) + .send({ query, variables }); + +describe('discard draft workflow version core mirror (e2e)', () => { + let workflowId: string; + let firstVersionId: string; + let draftVersionId: string; + + const coreVersions = async (): Promise<{ id: string; status: string }[]> => { + return global.testDataSource.query( + `SELECT "id", "status" FROM core."workflowVersion" + WHERE "workspaceId" = $1 AND "workflowId" = $2`, + [SEED_APPLE_WORKSPACE_ID, workflowId], + ); + }; + + beforeAll(async () => { + const createResponse = await graphql(` + mutation { + createWorkflow(data: { name: "Discard Draft Mirror" }) { + id + } + } + `); + + expect(createResponse.body.errors).toBeUndefined(); + workflowId = createResponse.body.data.createWorkflow.id; + + const getResponse = await graphql( + ` + query GetWorkflow($id: UUID!) { + workflow(filter: { id: { eq: $id } }) { + versions { + edges { + node { + id + } + } + } + } + } + `, + { id: workflowId }, + ); + + firstVersionId = getResponse.body.data.workflow.versions.edges[0].node.id; + + await updateWorkflowVersionTrigger({ + workflowVersionId: firstVersionId, + trigger: { + name: 'Manual Trigger', + type: 'MANUAL', + settings: { outputSchema: {} }, + nextStepIds: [], + position: { x: 0, y: 0 }, + }, + }); + + const stepResponse = await graphql( + ` + mutation CreateWorkflowVersionStep( + $input: CreateWorkflowVersionStepInput! + ) { + createWorkflowVersionStep(input: $input) { + stepsDiff + } + } + `, + { + input: { + workflowVersionId: firstVersionId, + stepType: 'FIND_RECORDS', + parentStepId: 'trigger', + position: { x: 200, y: 0 }, + }, + }, + ); + + expect(stepResponse.body.errors).toBeUndefined(); + + const activateResponse = await graphql( + ` + mutation ActivateWorkflowVersion($workflowVersionId: UUID!) { + activateWorkflowVersion(workflowVersionId: $workflowVersionId) + } + `, + { workflowVersionId: firstVersionId }, + ); + + expect(activateResponse.body.errors).toBeUndefined(); + + const draftResponse = await graphql( + ` + mutation CreateDraft($input: CreateDraftFromWorkflowVersionInput!) { + createDraftFromWorkflowVersion(input: $input) { + id + } + } + `, + { + input: { + workflowId, + workflowVersionIdToCopy: firstVersionId, + }, + }, + ); + + expect(draftResponse.body.errors).toBeUndefined(); + draftVersionId = draftResponse.body.data.createDraftFromWorkflowVersion.id; + }); + + afterAll(async () => { + if (workflowId) { + await graphql( + ` + mutation DestroyWorkflow($id: ID!) { + destroyWorkflow(id: $id) { + id + } + } + `, + { id: workflowId }, + ); + } + }); + + it('removes only the discarded draft core row, keeping the active version', async () => { + // the active version and the new draft are both mirrored to core + const versionsBefore = await coreVersions(); + + expect(versionsBefore.map((version) => version.status).sort()).toEqual([ + 'ACTIVE', + 'DRAFT', + ]); + + const activeCoreId = versionsBefore.find( + (version) => version.status === 'ACTIVE', + )?.id; + const draftCoreId = versionsBefore.find( + (version) => version.status === 'DRAFT', + )?.id; + + expect(activeCoreId).toBeDefined(); + expect(draftCoreId).toBeDefined(); + + const deleteResponse = await graphql( + ` + mutation DeleteWorkflowVersion($id: ID!) { + deleteWorkflowVersion(id: $id) { + id + } + } + `, + { id: draftVersionId }, + ); + + expect(deleteResponse.body.errors).toBeUndefined(); + + const versionsAfter = await coreVersions(); + + // only the discarded draft's core row is removed; the active row remains + expect(versionsAfter).toHaveLength(1); + expect(versionsAfter[0].id).toBe(activeCoreId); + expect(versionsAfter[0].id).not.toBe(draftCoreId); + }); +});