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.

<!-- This is an auto-generated description by cubic. -->
<a
href="https://cubic.dev/pr/twentyhq/twenty/pull/23374?utm_source=github"
target="_blank" rel="noopener noreferrer"
data-no-image-dialog="true"><picture><source
media="(prefers-color-scheme: dark)"
srcset="https://www.cubic.dev/buttons/review-in-cubic-dark.svg"><source
media="(prefers-color-scheme: light)"
srcset="https://www.cubic.dev/buttons/review-in-cubic-light.svg"><img
alt="Review in cubic"
src="https://www.cubic.dev/buttons/review-in-cubic-dark.svg"></picture></a>
<!-- End of auto-generated description by cubic. -->
This commit is contained in:
Thomas Trompette
2026-07-27 17:52:15 +02:00
committed by GitHub
parent 3c48e27b2e
commit 2c49c4169c
7 changed files with 246 additions and 115 deletions
@@ -280,6 +280,39 @@ export class WorkflowVersionCoreSyncService {
await this.invalidateAutomatedTriggerMaps(workspaceId);
}
async deleteCoreVersionsByWorkspaceVersionIds(
workspaceId: string,
workflowVersionIds: string[],
): Promise<void> {
if (workflowVersionIds.length === 0) {
return;
}
const coreWorkflowVersionIds =
await this.globalWorkspaceOrmManager.executeInWorkspaceContext(
async () => {
const workflowVersionRepository =
await this.globalWorkspaceOrmManager.getRepository<WorkflowVersionWorkspaceEntity>(
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,
@@ -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,
@@ -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<void> {
const workspace = authContext.workspace;
assertIsDefinedOrThrow(workspace, WorkspaceNotFoundDefaultError);
await this.workflowVersionCoreSyncService.deleteCoreVersionsByWorkspaceVersionIds(
workspace.id,
payload.map((workflowVersion) => workflowVersion.id),
);
}
}
@@ -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<WorkflowVersionWorkspaceEntity>
>,
): Promise<void> {
await this.upsertToCore(
batchEvent.workspaceId,
batchEvent.events.map((event) => event.properties.after),
);
}
@OnDatabaseBatchEvent('workflowVersion', DatabaseEventAction.DELETED)
async handleDeleted(
batchEvent: CustomWorkspaceEventBatch<
ObjectRecordDeleteEvent<WorkflowVersionWorkspaceEntity>
>,
): Promise<void> {
await this.deleteFromCore(
batchEvent.workspaceId,
batchEvent.events
.map((event) => event.properties.before.coreWorkflowVersionId)
.filter(isDefined),
);
}
@OnDatabaseBatchEvent('workflowVersion', DatabaseEventAction.DESTROYED)
async handleDestroyed(
batchEvent: CustomWorkspaceEventBatch<
ObjectRecordDestroyEvent<WorkflowVersionWorkspaceEntity>
>,
): Promise<void> {
await this.deleteFromCore(
batchEvent.workspaceId,
batchEvent.events
.map((event) => event.properties.before.coreWorkflowVersionId)
.filter(isDefined),
);
}
private async upsertToCore(
workspaceId: string | undefined,
workflowVersions: WorkflowVersionWorkspaceEntity[],
): Promise<void> {
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<void> {
if (!isDefined(workspaceId)) {
return;
}
try {
await this.workflowVersionCoreSyncService.deleteFromCore(
workspaceId,
coreWorkflowVersionIds,
);
} catch (error) {
this.exceptionHandlerService.captureExceptions([error], {
workspace: { id: workspaceId },
});
}
}
}
@@ -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 {}
@@ -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,
],
})
@@ -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);
});
});