From 2c49c4169cd58e2bf70d0b7b15ff5549d5feb1fa Mon Sep 17 00:00:00 2001 From: Thomas Trompette Date: Mon, 27 Jul 2026 17:52:15 +0200 Subject: [PATCH] feat(workflow): remove async workflowVersion core dual-write listener (#23374) ## What Removes `WorkflowVersionCoreDualWriteListener` (and its now-empty module), replacing the last async best-effort core writes for workflow versions with synchronous mirrors. Adds one missing synchronous funnel so nothing is left uncovered. ## Why After #23356, the listener's `handleRestored` / `handleDeleted` / `handleDestroyed` handlers are redundant with the synchronous lifecycle mirror, so the async path (which can silently drift on failure) can go. While removing it I found one path the listener was **not** redundant on: **direct `deleteOneWorkflowVersion` (discard draft)** is an allowed operation (discard a DRAFT version that isn't the only version, via the `DISCARD_DRAFT_WORKFLOW` command) and had **no** synchronous post-hook. The async listener was the sole thing deleting its core row. Removing the listener without a replacement would have drifted on every draft discard. So this PR also adds a `workflowVersion.deleteOne` post-hook that mirrors the deletion to core. ## Coverage after this change | version lifecycle path | synchronous coverage | | --- | --- | | `deleteOneWorkflowVersion` (discard draft) | **new** `workflowVersion.deleteOne` post-hook | | delete via workflow cascade | `handleWorkflowSubEntities` -> `deleteCoreVersionsByWorkflowIds` (#23356) | | restore via workflow cascade | `handleWorkflowSubEntities` -> `recreateCoreVersionsByWorkflowId` (#23356) | | destroy via workflow | `workflow.destroy*` post-hooks (#23356) | | `deleteMany` / `destroyOne|Many` / `restoreOne|Many` version | blocked by pre-hooks ("Method not allowed") | ## Notes - The new post-hook re-fetches the version (`withDeleted`) to resolve its `coreWorkflowVersionId`, because the delete post-hook payload only carries the columns the client selected (the delete `RETURNING` set is built from `selectedFieldsResult.select`), so `coreWorkflowVersionId` is not reliably present. - `deleteCoreVersionsByWorkspaceVersionIds` deletes precisely by `coreWorkflowVersionId` (not by `workflowId`), so discarding one draft does not touch the core rows of the workflow's other versions. - The workflow-side `WorkflowCoreSyncModule` listener is intentionally left in place (separate migration track). ## Verification - `nx typecheck twenty-server` green - `nx lint:diff-with-main twenty-server` green - Added integration test `workflow-version-discard-draft-core-mirror`: activate v1, create a draft, discard it, assert only the draft's core row is removed and the active version's core row remains. - Live run on a dev instance still pending. ## Merge gate Per the migration plan, removing the async backstop should land only after the drift cron reports zero drift over a soak period. Review in cubic --- .../workflow-version-core-sync.service.ts | 33 ++++ .../query-hooks/workflow-query-hook.module.ts | 2 + ...flow-version-delete-one.post-query.hook.ts | 35 ++++ ...rkflow-version-core-dual-write.listener.ts | 103 ---------- .../workflow-version-core-sync.module.ts | 10 - .../src/modules/workflow/workflow.module.ts | 2 - ...card-draft-core-mirror.integration-spec.ts | 176 ++++++++++++++++++ 7 files changed, 246 insertions(+), 115 deletions(-) create mode 100644 packages/twenty-server/src/modules/workflow/common/query-hooks/workflow-version-delete-one.post-query.hook.ts delete mode 100644 packages/twenty-server/src/modules/workflow/workflow-version-core-sync/listeners/workflow-version-core-dual-write.listener.ts delete mode 100644 packages/twenty-server/src/modules/workflow/workflow-version-core-sync/workflow-version-core-sync.module.ts create mode 100644 packages/twenty-server/test/integration/graphql/suites/workflow/workflow-version-discard-draft-core-mirror.integration-spec.ts 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); + }); +});