From ca90a9358f6b7f71ef4bc6d0dbcd28b4ee191419 Mon Sep 17 00:00:00 2001 From: Thomas Trompette Date: Wed, 8 Jul 2026 18:45:57 +0200 Subject: [PATCH] feat(workflow): backfill workspace workflowVersion into core (phase A) (#22663) ## workflowVersion -> core, Phase A Follows #21674 (Phase 0, merged). Base: `main`. Populates core `workflowVersion` and keeps it in sync with the workspace object, so a later phase can switch reads to core. Reads stay on the workspace object in this PR. ### 1. Backfill (upgrade command) `BackfillWorkflowVersionToCoreCommand`, a `@RegisteredWorkspaceCommand('2.20.0', ...)`. Per workspace, reads all workspace `workflowVersion` records and upserts them into core, preserving ids (idempotent), dry-run aware. ### 2. Dual-write (always on, not flag-gated) `WorkflowVersionCoreDualWriteListener` hooks `@OnDatabaseBatchEvent('workflowVersion', CREATED/UPDATED/DELETED)` (same mechanism as the existing workflow-version status listener) and mirrors every mutation into core. Sync failures are logged, never break the user's write; drift is repaired by re-running the backfill command. Dual-write is deliberately not behind a flag: reading from core (next phase) is only safe if core has been continuously in sync since the backfill. An always-on mirror makes "core is fresh" an invariant, so the read switch becomes a plain flag flip. The cost is one extra upsert on infrequent workflowVersion writes. `IS_WORKFLOW_VERSION_IN_CORE_ENABLED` is reserved for the read switch (Phase B): dispatch from the `workflowAutomatedTriggerMaps` cache, runner and builder reading trigger/steps from core. Until then it gates nothing. Both the backfill and the listener go through a single `WorkflowVersionCoreSyncService` (`upsertToCore`/`deleteFromCore`): the workspace-to-core mapping (`trigger` -> `triggers[]`, plus `steps`, `status`, `workflowId`) and `workflowAutomatedTriggerMaps` invalidation live in one place. ### Rollout plan (following phases) - **B, read switch (flag per workspace):** reads move to core; writes keep flowing workspace -> listener -> core. Rollback = flip the flag back, workspace never stopped being source of truth. - **C, contract (code change):** write paths write trigger/steps to core directly; workspace `workflowVersion` stays as a thin shell (nav/relations/search) but drops the trigger/steps columns; listener and flag removed. ### Not in this PR - Reconciliation tooling beyond re-running the backfill. - The read switch (Phase B). --- .../2-20-upgrade-version-command.module.ts | 6 +- ...ckfill-workflow-version-to-core.command.ts | 78 +++++++++++ .../workflow-version-core-sync.service.ts | 70 ++++++++++ .../workflow/workflow-version-core.module.ts | 14 +- ...rkflow-version-core-dual-write.listener.ts | 125 ++++++++++++++++++ .../workflow-version-core-sync.module.ts | 10 ++ .../src/modules/workflow/workflow.module.ts | 7 +- 7 files changed, 306 insertions(+), 4 deletions(-) create mode 100644 packages/twenty-server/src/database/commands/upgrade-version-command/2-20/2-20-workspace-command-1783526282685-backfill-workflow-version-to-core.command.ts create mode 100644 packages/twenty-server/src/engine/core-modules/workflow/services/workflow-version-core-sync.service.ts create mode 100644 packages/twenty-server/src/modules/workflow/workflow-version-core-sync/listeners/workflow-version-core-dual-write.listener.ts create mode 100644 packages/twenty-server/src/modules/workflow/workflow-version-core-sync/workflow-version-core-sync.module.ts diff --git a/packages/twenty-server/src/database/commands/upgrade-version-command/2-20/2-20-upgrade-version-command.module.ts b/packages/twenty-server/src/database/commands/upgrade-version-command/2-20/2-20-upgrade-version-command.module.ts index 26e67effdb..da9a87b326 100644 --- a/packages/twenty-server/src/database/commands/upgrade-version-command/2-20/2-20-upgrade-version-command.module.ts +++ b/packages/twenty-server/src/database/commands/upgrade-version-command/2-20/2-20-upgrade-version-command.module.ts @@ -1,10 +1,12 @@ import { Module } from '@nestjs/common'; import { WorkspaceIteratorModule } from 'src/database/commands/command-runners/workspace-iterator.module'; +import { BackfillActorSourceEnumValuesCommand } from 'src/database/commands/upgrade-version-command/2-20/2-20-workspace-command-1783499671542-backfill-actor-source-enum-values.command'; +import { BackfillWorkflowVersionToCoreCommand } from 'src/database/commands/upgrade-version-command/2-20/2-20-workspace-command-1783526282685-backfill-workflow-version-to-core.command'; import { AddMessageCampaignStatFieldsCommand } from 'src/database/commands/upgrade-version-command/2-20/2-20-workspace-command-1783525261000-add-message-campaign-stat-fields.command'; import { CreateMessageListViewCommand } from 'src/database/commands/upgrade-version-command/2-20/2-20-workspace-command-1783525261001-create-message-list-view.command'; -import { BackfillActorSourceEnumValuesCommand } from 'src/database/commands/upgrade-version-command/2-20/2-20-workspace-command-1783499671542-backfill-actor-source-enum-values.command'; import { ApplicationModule } from 'src/engine/core-modules/application/application.module'; +import { WorkflowVersionCoreModule } from 'src/engine/core-modules/workflow/workflow-version-core.module'; import { WorkspaceSchemaManagerModule } from 'src/engine/twenty-orm/workspace-schema-manager/workspace-schema-manager.module'; import { WorkspaceCacheModule } from 'src/engine/workspace-cache/workspace-cache.module'; import { WorkspaceMigrationModule } from 'src/engine/workspace-manager/workspace-migration/workspace-migration.module'; @@ -16,11 +18,13 @@ import { WorkspaceMigrationModule } from 'src/engine/workspace-manager/workspace WorkspaceCacheModule, WorkspaceMigrationModule, WorkspaceSchemaManagerModule, + WorkflowVersionCoreModule, ], providers: [ AddMessageCampaignStatFieldsCommand, CreateMessageListViewCommand, BackfillActorSourceEnumValuesCommand, + BackfillWorkflowVersionToCoreCommand, ], }) export class V2_20_UpgradeVersionCommandModule {} diff --git a/packages/twenty-server/src/database/commands/upgrade-version-command/2-20/2-20-workspace-command-1783526282685-backfill-workflow-version-to-core.command.ts b/packages/twenty-server/src/database/commands/upgrade-version-command/2-20/2-20-workspace-command-1783526282685-backfill-workflow-version-to-core.command.ts new file mode 100644 index 0000000000..d2ddf51fb0 --- /dev/null +++ b/packages/twenty-server/src/database/commands/upgrade-version-command/2-20/2-20-workspace-command-1783526282685-backfill-workflow-version-to-core.command.ts @@ -0,0 +1,78 @@ +import { Command } from 'nest-commander'; +import { EntityMetadataNotFoundError } from 'typeorm/error/EntityMetadataNotFoundError'; + +import { ActiveOrSuspendedWorkspaceCommandRunner } from 'src/database/commands/command-runners/active-or-suspended-workspace.command-runner'; +import { WorkspaceIteratorService } from 'src/database/commands/command-runners/workspace-iterator.service'; +import { type RunOnWorkspaceArgs } from 'src/database/commands/command-runners/workspace.command-runner'; +import { RegisteredWorkspaceCommand } from 'src/engine/core-modules/upgrade/decorators/registered-workspace-command.decorator'; +import { WorkflowVersionCoreSyncService } from 'src/engine/core-modules/workflow/services/workflow-version-core-sync.service'; +import { buildSystemAuthContext } from 'src/engine/twenty-orm/utils/build-system-auth-context.util'; +import { GlobalWorkspaceOrmManager } from 'src/engine/twenty-orm/global-workspace-datasource/global-workspace-orm.manager'; +import { type WorkflowVersionWorkspaceEntity } from 'src/modules/workflow/common/standard-objects/workflow-version.workspace-entity'; + +@RegisteredWorkspaceCommand('2.20.0', 1783526282685) +@Command({ + name: 'upgrade:2-20:backfill-workflow-version-to-core', + description: + 'Copy each workspace workflowVersion (trigger, steps, status, workflowId) into the core workflowVersion table, preserving ids', +}) +export class BackfillWorkflowVersionToCoreCommand extends ActiveOrSuspendedWorkspaceCommandRunner { + constructor( + protected readonly workspaceIteratorService: WorkspaceIteratorService, + private readonly globalWorkspaceOrmManager: GlobalWorkspaceOrmManager, + private readonly workflowVersionCoreSyncService: WorkflowVersionCoreSyncService, + ) { + super(workspaceIteratorService); + } + + override async runOnWorkspace({ + workspaceId, + options, + }: RunOnWorkspaceArgs): Promise { + let workspaceWorkflowVersions: WorkflowVersionWorkspaceEntity[]; + + try { + workspaceWorkflowVersions = + await this.globalWorkspaceOrmManager.executeInWorkspaceContext( + async () => { + const workflowVersionRepository = + await this.globalWorkspaceOrmManager.getRepository( + workspaceId, + 'workflowVersion', + { shouldBypassPermissionChecks: true }, + ); + + return workflowVersionRepository.find(); + }, + buildSystemAuthContext(workspaceId), + ); + } catch (error) { + if (error instanceof EntityMetadataNotFoundError) { + this.logger.log( + `workflowVersion object does not exist for workspace ${workspaceId}, skipping`, + ); + + return; + } + + throw error; + } + + if (options.dryRun === true) { + this.logger.log( + `[DRY RUN] Would upsert ${workspaceWorkflowVersions.length} workflowVersion row(s) into core for workspace ${workspaceId}`, + ); + + return; + } + + await this.workflowVersionCoreSyncService.upsertToCore( + workspaceId, + workspaceWorkflowVersions, + ); + + this.logger.log( + `Backfilled ${workspaceWorkflowVersions.length} workflowVersion row(s) into core for workspace ${workspaceId}`, + ); + } +} 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 new file mode 100644 index 0000000000..04a91e2a84 --- /dev/null +++ b/packages/twenty-server/src/engine/core-modules/workflow/services/workflow-version-core-sync.service.ts @@ -0,0 +1,70 @@ +import { Injectable } from '@nestjs/common'; + +import { isDefined } from 'twenty-shared/utils'; +import { In } from 'typeorm'; + +import { + WorkflowVersionEntity, + WorkflowVersionStatus, +} from 'src/engine/core-modules/workflow/entities/workflow-version.entity'; +import { InjectWorkspaceScopedRepository } from 'src/engine/twenty-orm/workspace-scoped-repository/inject-workspace-scoped-repository.decorator'; +import { WorkspaceScopedRepository } from 'src/engine/twenty-orm/workspace-scoped-repository/workspace-scoped-repository'; +import { WorkspaceCacheService } from 'src/engine/workspace-cache/services/workspace-cache.service'; +import { type WorkflowVersionWorkspaceEntity } from 'src/modules/workflow/common/standard-objects/workflow-version.workspace-entity'; + +@Injectable() +export class WorkflowVersionCoreSyncService { + constructor( + @InjectWorkspaceScopedRepository(WorkflowVersionEntity) + private readonly workflowVersionRepository: WorkspaceScopedRepository, + private readonly workspaceCacheService: WorkspaceCacheService, + ) {} + + async upsertToCore( + workspaceId: string, + workflowVersions: WorkflowVersionWorkspaceEntity[], + ): Promise { + if (workflowVersions.length === 0) { + return; + } + + await this.workflowVersionRepository.upsert( + workspaceId, + workflowVersions.map((workflowVersion) => ({ + id: workflowVersion.id, + workflowId: workflowVersion.workflowId, + triggers: isDefined(workflowVersion.trigger) + ? [workflowVersion.trigger] + : null, + steps: workflowVersion.steps ?? null, + status: workflowVersion.status as unknown as WorkflowVersionStatus, + })), + ['id'], + ); + + await this.invalidateAutomatedTriggerMaps(workspaceId); + } + + async deleteFromCore( + workspaceId: string, + workflowVersionIds: string[], + ): Promise { + if (workflowVersionIds.length === 0) { + return; + } + + await this.workflowVersionRepository.delete(workspaceId, { + id: In(workflowVersionIds), + }); + + await this.invalidateAutomatedTriggerMaps(workspaceId); + } + + private async invalidateAutomatedTriggerMaps( + workspaceId: string, + ): Promise { + await this.workspaceCacheService.invalidateAndRecompute(workspaceId, [ + 'workflowAutomatedTriggerMaps', + ]); + } +} diff --git a/packages/twenty-server/src/engine/core-modules/workflow/workflow-version-core.module.ts b/packages/twenty-server/src/engine/core-modules/workflow/workflow-version-core.module.ts index 48132b5d0d..9849724cad 100644 --- a/packages/twenty-server/src/engine/core-modules/workflow/workflow-version-core.module.ts +++ b/packages/twenty-server/src/engine/core-modules/workflow/workflow-version-core.module.ts @@ -2,15 +2,25 @@ import { Module } from '@nestjs/common'; import { TypeOrmModule } from '@nestjs/typeorm'; import { WorkflowVersionEntity } from 'src/engine/core-modules/workflow/entities/workflow-version.entity'; +import { WorkflowVersionCoreSyncService } from 'src/engine/core-modules/workflow/services/workflow-version-core-sync.service'; import { WorkspaceWorkflowAutomatedTriggerMapCacheService } from 'src/engine/core-modules/workflow/services/workspace-workflow-automated-trigger-map-cache.service'; import { provideWorkspaceScopedRepository } from 'src/engine/twenty-orm/workspace-scoped-repository/provide-workspace-scoped-repository'; +import { WorkspaceCacheModule } from 'src/engine/workspace-cache/workspace-cache.module'; @Module({ - imports: [TypeOrmModule.forFeature([WorkflowVersionEntity])], + imports: [ + TypeOrmModule.forFeature([WorkflowVersionEntity]), + WorkspaceCacheModule, + ], providers: [ WorkspaceWorkflowAutomatedTriggerMapCacheService, + WorkflowVersionCoreSyncService, provideWorkspaceScopedRepository(WorkflowVersionEntity), ], - exports: [TypeOrmModule, WorkspaceWorkflowAutomatedTriggerMapCacheService], + exports: [ + TypeOrmModule, + WorkspaceWorkflowAutomatedTriggerMapCacheService, + WorkflowVersionCoreSyncService, + ], }) export class WorkflowVersionCoreModule {} 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 new file mode 100644 index 0000000000..d1d085ba4d --- /dev/null +++ b/packages/twenty-server/src/modules/workflow/workflow-version-core-sync/listeners/workflow-version-core-dual-write.listener.ts @@ -0,0 +1,125 @@ +import { Injectable } from '@nestjs/common'; + +import { + type ObjectRecordCreateEvent, + type ObjectRecordDeleteEvent, + type ObjectRecordDestroyEvent, + type ObjectRecordRestoreEvent, + type ObjectRecordUpdateEvent, +} 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.CREATED) + async handleCreated( + batchEvent: CustomWorkspaceEventBatch< + ObjectRecordCreateEvent + >, + ): Promise { + await this.upsertToCore( + batchEvent.workspaceId, + batchEvent.events.map((event) => event.properties.after), + ); + } + + @OnDatabaseBatchEvent('workflowVersion', DatabaseEventAction.UPDATED) + async handleUpdated( + batchEvent: CustomWorkspaceEventBatch< + ObjectRecordUpdateEvent + >, + ): Promise { + await this.upsertToCore( + batchEvent.workspaceId, + batchEvent.events.map((event) => event.properties.after), + ); + } + + @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.id), + ); + } + + @OnDatabaseBatchEvent('workflowVersion', DatabaseEventAction.DESTROYED) + async handleDestroyed( + batchEvent: CustomWorkspaceEventBatch< + ObjectRecordDestroyEvent + >, + ): Promise { + await this.deleteFromCore( + batchEvent.workspaceId, + batchEvent.events.map((event) => event.properties.before.id), + ); + } + + 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, + workflowVersionIds: string[], + ): Promise { + if (!isDefined(workspaceId)) { + return; + } + + try { + await this.workflowVersionCoreSyncService.deleteFromCore( + workspaceId, + workflowVersionIds, + ); + } 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 new file mode 100644 index 0000000000..19b15565d7 --- /dev/null +++ b/packages/twenty-server/src/modules/workflow/workflow-version-core-sync/workflow-version-core-sync.module.ts @@ -0,0 +1,10 @@ +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 08456972b3..485874e4e6 100644 --- a/packages/twenty-server/src/modules/workflow/workflow.module.ts +++ b/packages/twenty-server/src/modules/workflow/workflow.module.ts @@ -2,8 +2,13 @@ import { Module } from '@nestjs/common'; 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], + imports: [ + WorkflowTriggerModule, + WorkflowStatusModule, + WorkflowVersionCoreSyncModule, + ], }) export class WorkflowModule {}