From 79bf20515d8c83adb298582f8ef470f0ef043de4 Mon Sep 17 00:00:00 2001 From: Thomas Trompette Date: Mon, 27 Jul 2026 16:43:59 +0200 Subject: [PATCH] feat(workflow): mirror version delete/restore/destroy to core transactionally (#23356) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Next step in the workflowVersion -> core soft-ref migration. The transactional mirror (#23243) made **content** writes (create/update) drift-free. This does the same for the **lifecycle** events (delete / restore / destroy), which were still handled only by the async best-effort listener. It's the prerequisite for dropping that listener. ## What changed `delete` and `restore` already soft-delete / restore the workflow's versions inside `handleWorkflowSubEntities` (twenty doesn't cascade soft-deletes, so it does each sub-entity explicitly). So the core delete/recreate just sits next to the existing version write: - **delete** — after `workflowVersionRepository.softDelete({ workflowId })`, `deleteCoreVersionsByWorkflowIds` removes the `core.workflowVersion` rows (`workflowId IN (...)`). (`deactivateVersionOnDelete` no longer re-mirrors the deactivated version — that was recreating the core row it just deleted; it only flips the workspace status to `DEACTIVATED` so a restore comes back deactivated.) - **restore** — after `workflowVersionRepository.restore({ workflowId })`, `recreateCoreVersionsByWorkflowId` re-reads the restored versions and reuses the existing `upsertToCore` (which reuses the stored `coreWorkflowVersionId` soft-ref, so rows come back with their original ids and current status). - **destroy** — version destroy isn't done in `handleWorkflowSubEntities` (it happens via the generic cascade), so there's no existing place to hang the core delete. New `workflow.destroyOne`/`destroyMany` **post**-hooks call `deleteCoreVersionsByWorkflowIds` (batched `IN`) only after the destroy commits, so a rejected destroy can't remove core rows while the workspace versions survive. No new transactional wrappers or raw SQL — the delete/recreate reuse the existing `WorkflowVersionCoreSyncService` methods (`deleteFromCore`-style delete, `upsertToCore`). The async listener stays as an idempotent backstop until the cron soaks zero drift. ## Async listener kept as backstop `handleRestored` / `handleDeleted` / `handleDestroyed` stay for now. Both paths are idempotent (delete-of-deleted is a no-op; upsert converges), so they don't conflict. Those handlers come out in a follow-up once the consistency cron soaks zero drift - which this PR unblocks. ## Verification Lifecycle integration test — creates a workflow, **activates** the version (the active path is where the delete re-mirror bug bit), then asserts the core row: present -> gone after delete -> back after restore (as `DEACTIVATED`, same id) -> gone after destroy. Plus the existing `workflow-resolver` delete/restore suite (regression, since `handleWorkflowSubEntities` is shared). Also verified **live** on a running instance against the real DB: the full active-version lifecycle above, plus a batched `destroyWorkflows` on two workflows removing both core rows in one `IN` delete. `nx typecheck` + oxlint + oxfmt clean. Review in cubic --- .../workflow-version-core-sync.service.ts | 35 ++++ .../workflow-destroy-many.post-query.hook.ts | 35 ++++ .../workflow-destroy-one.post-query.hook.ts | 35 ++++ .../query-hooks/workflow-query-hook.module.ts | 4 + .../workflow-common.workspace-service.ts | 19 +- ...-lifecycle-core-mirror.integration-spec.ts | 168 ++++++++++++++++++ 6 files changed, 287 insertions(+), 9 deletions(-) create mode 100644 packages/twenty-server/src/modules/workflow/common/query-hooks/workflow-destroy-many.post-query.hook.ts create mode 100644 packages/twenty-server/src/modules/workflow/common/query-hooks/workflow-destroy-one.post-query.hook.ts create mode 100644 packages/twenty-server/test/integration/graphql/suites/workflow/workflow-lifecycle-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 156caf9eed..800dc6bbd3 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 @@ -265,6 +265,41 @@ export class WorkflowVersionCoreSyncService { await this.invalidateAutomatedTriggerMaps(workspaceId); } + async deleteCoreVersionsByWorkflowIds( + workspaceId: string, + workflowIds: string[], + ): Promise { + if (workflowIds.length === 0) { + return; + } + + await this.coreWorkflowVersionRepository.delete(workspaceId, { + workflowId: In(workflowIds), + }); + + await this.invalidateAutomatedTriggerMaps(workspaceId); + } + + async recreateCoreVersionsByWorkflowId( + workspaceId: string, + workflowId: string, + ): Promise { + await this.globalWorkspaceOrmManager.executeInWorkspaceContext(async () => { + const workflowVersionRepository = + await this.globalWorkspaceOrmManager.getRepository( + workspaceId, + 'workflowVersion', + { shouldBypassPermissionChecks: true }, + ); + + const versions = await workflowVersionRepository.find({ + where: { workflowId }, + }); + + await this.upsertToCore(workspaceId, versions); + }, buildSystemAuthContext(workspaceId)); + } + private async writeBackCoreVersionIds( workspaceId: string, coreVersionIdByWorkspaceRecordId: Map, diff --git a/packages/twenty-server/src/modules/workflow/common/query-hooks/workflow-destroy-many.post-query.hook.ts b/packages/twenty-server/src/modules/workflow/common/query-hooks/workflow-destroy-many.post-query.hook.ts new file mode 100644 index 0000000000..f7e152d720 --- /dev/null +++ b/packages/twenty-server/src/modules/workflow/common/query-hooks/workflow-destroy-many.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 WorkflowWorkspaceEntity } from 'src/modules/workflow/common/standard-objects/workflow.workspace-entity'; + +@WorkspaceQueryHook({ + key: `workflow.destroyMany`, + type: WorkspaceQueryHookType.POST_HOOK, +}) +export class WorkflowDestroyManyPostQueryHook implements WorkspacePostQueryHookInstance { + constructor( + private readonly workflowVersionCoreSyncService: WorkflowVersionCoreSyncService, + ) {} + + async execute( + authContext: WorkspaceAuthContext, + _objectName: string, + payload: WorkflowWorkspaceEntity[], + ): Promise { + const workspace = authContext.workspace; + + assertIsDefinedOrThrow(workspace, WorkspaceNotFoundDefaultError); + + await this.workflowVersionCoreSyncService.deleteCoreVersionsByWorkflowIds( + workspace.id, + payload.map((workflow) => workflow.id), + ); + } +} diff --git a/packages/twenty-server/src/modules/workflow/common/query-hooks/workflow-destroy-one.post-query.hook.ts b/packages/twenty-server/src/modules/workflow/common/query-hooks/workflow-destroy-one.post-query.hook.ts new file mode 100644 index 0000000000..81cf6c796a --- /dev/null +++ b/packages/twenty-server/src/modules/workflow/common/query-hooks/workflow-destroy-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 WorkflowWorkspaceEntity } from 'src/modules/workflow/common/standard-objects/workflow.workspace-entity'; + +@WorkspaceQueryHook({ + key: `workflow.destroyOne`, + type: WorkspaceQueryHookType.POST_HOOK, +}) +export class WorkflowDestroyOnePostQueryHook implements WorkspacePostQueryHookInstance { + constructor( + private readonly workflowVersionCoreSyncService: WorkflowVersionCoreSyncService, + ) {} + + async execute( + authContext: WorkspaceAuthContext, + _objectName: string, + payload: WorkflowWorkspaceEntity[], + ): Promise { + const workspace = authContext.workspace; + + assertIsDefinedOrThrow(workspace, WorkspaceNotFoundDefaultError); + + await this.workflowVersionCoreSyncService.deleteCoreVersionsByWorkflowIds( + workspace.id, + payload.map((workflow) => workflow.id), + ); + } +} 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 b40c9ff615..f766ed6d86 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 @@ -14,7 +14,9 @@ import { WorkflowCreateOnePostQueryHook } from 'src/modules/workflow/common/quer import { WorkflowCreateOnePreQueryHook } from 'src/modules/workflow/common/query-hooks/workflow-create-one.pre-query.hook'; import { WorkflowDeleteManyPostQueryHook } from 'src/modules/workflow/common/query-hooks/workflow-delete-many.post-query.hook'; import { WorkflowDeleteOnePostQueryHook } from 'src/modules/workflow/common/query-hooks/workflow-delete-one.post-query.hook'; +import { WorkflowDestroyManyPostQueryHook } from 'src/modules/workflow/common/query-hooks/workflow-destroy-many.post-query.hook'; import { WorkflowDestroyManyPreQueryHook } from 'src/modules/workflow/common/query-hooks/workflow-destroy-many.pre-query.hook'; +import { WorkflowDestroyOnePostQueryHook } from 'src/modules/workflow/common/query-hooks/workflow-destroy-one.post-query.hook'; import { WorkflowDestroyOnePreQueryHook } from 'src/modules/workflow/common/query-hooks/workflow-destroy-one.pre-query.hook'; import { WorkflowRestoreManyPostQueryHook } from 'src/modules/workflow/common/query-hooks/workflow-restore-many.post-query.hook'; import { WorkflowRestoreOnePostQueryHook } from 'src/modules/workflow/common/query-hooks/workflow-restore-one.post-query.hook'; @@ -85,6 +87,8 @@ import { WorkflowVersionValidationWorkspaceService } from 'src/modules/workflow/ WorkflowDeleteOnePostQueryHook, WorkflowDestroyOnePreQueryHook, WorkflowDestroyManyPreQueryHook, + WorkflowDestroyOnePostQueryHook, + WorkflowDestroyManyPostQueryHook, ], }) export class WorkflowQueryHookModule {} diff --git a/packages/twenty-server/src/modules/workflow/common/workspace-services/workflow-common.workspace-service.ts b/packages/twenty-server/src/modules/workflow/common/workspace-services/workflow-common.workspace-service.ts index 1c7feab3e7..f27e66f474 100644 --- a/packages/twenty-server/src/modules/workflow/common/workspace-services/workflow-common.workspace-service.ts +++ b/packages/twenty-server/src/modules/workflow/common/workspace-services/workflow-common.workspace-service.ts @@ -312,6 +312,11 @@ export class WorkflowCommonWorkspaceService { workflowId, }); + await this.workflowVersionCoreSyncService.deleteCoreVersionsByWorkflowIds( + workspaceId, + [workflowId], + ); + break; case 'restore': await workflowAutomatedTriggerRepository.restore({ @@ -326,6 +331,11 @@ export class WorkflowCommonWorkspaceService { workflowId, }); + await this.workflowVersionCoreSyncService.recreateCoreVersionsByWorkflowId( + workspaceId, + workflowId, + ); + break; } @@ -414,15 +424,6 @@ export class WorkflowCommonWorkspaceService { undefined, queryRunner.manager, ); - - await this.workflowVersionCoreSyncService.mirrorWorkflowVersionWrite({ - workspaceId, - entityManager: queryRunner.manager, - workflowVersion: { - ...workflowVersion, - status: WorkflowVersionStatus.DEACTIVATED, - }, - }); } } diff --git a/packages/twenty-server/test/integration/graphql/suites/workflow/workflow-lifecycle-core-mirror.integration-spec.ts b/packages/twenty-server/test/integration/graphql/suites/workflow/workflow-lifecycle-core-mirror.integration-spec.ts new file mode 100644 index 0000000000..55fb0083e6 --- /dev/null +++ b/packages/twenty-server/test/integration/graphql/suites/workflow/workflow-lifecycle-core-mirror.integration-spec.ts @@ -0,0 +1,168 @@ +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('workflow lifecycle core mirror with an active version (e2e)', () => { + let workflowId: string; + let alreadyDestroyed = false; + + const coreVersionRowCount = async (): Promise => { + const rows = await global.testDataSource.query( + `SELECT "id" FROM core."workflowVersion" + WHERE "workspaceId" = $1 AND "workflowId" = $2`, + [SEED_APPLE_WORKSPACE_ID, workflowId], + ); + + return rows.length; + }; + + beforeAll(async () => { + const createResponse = await graphql(` + mutation { + createWorkflow(data: { name: "Lifecycle 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 }, + ); + + const workflowVersionId = + getResponse.body.data.workflow.versions.edges[0].node.id; + + await updateWorkflowVersionTrigger({ + workflowVersionId, + 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, + 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 }, + ); + + expect(activateResponse.body.errors).toBeUndefined(); + expect(activateResponse.body.data.activateWorkflowVersion).toBe(true); + }); + + afterAll(async () => { + if (workflowId && !alreadyDestroyed) { + await graphql( + ` + mutation DestroyWorkflow($id: ID!) { + destroyWorkflow(id: $id) { + id + } + } + `, + { id: workflowId }, + ); + } + }); + + it('removes the core version row on delete, recreates it on restore, removes it on destroy', async () => { + // the v1 version created with the workflow is mirrored to core + expect(await coreVersionRowCount()).toBeGreaterThan(0); + + const deleteResponse = await graphql( + ` + mutation DeleteWorkflow($id: ID!) { + deleteWorkflow(id: $id) { + id + } + } + `, + { id: workflowId }, + ); + + expect(deleteResponse.body.errors).toBeUndefined(); + expect(await coreVersionRowCount()).toBe(0); + + const restoreResponse = await graphql( + ` + mutation RestoreWorkflow($id: ID!) { + restoreWorkflow(id: $id) { + id + } + } + `, + { id: workflowId }, + ); + + expect(restoreResponse.body.errors).toBeUndefined(); + expect(await coreVersionRowCount()).toBeGreaterThan(0); + + const destroyResponse = await graphql( + ` + mutation DestroyWorkflow($id: ID!) { + destroyWorkflow(id: $id) { + id + } + } + `, + { id: workflowId }, + ); + + expect(destroyResponse.body.errors).toBeUndefined(); + alreadyDestroyed = true; + expect(await coreVersionRowCount()).toBe(0); + }); +});