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).
This commit is contained in:
+5
-1
@@ -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 {}
|
||||
|
||||
+78
@@ -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<void> {
|
||||
let workspaceWorkflowVersions: WorkflowVersionWorkspaceEntity[];
|
||||
|
||||
try {
|
||||
workspaceWorkflowVersions =
|
||||
await this.globalWorkspaceOrmManager.executeInWorkspaceContext(
|
||||
async () => {
|
||||
const workflowVersionRepository =
|
||||
await this.globalWorkspaceOrmManager.getRepository<WorkflowVersionWorkspaceEntity>(
|
||||
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}`,
|
||||
);
|
||||
}
|
||||
}
|
||||
+70
@@ -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<WorkflowVersionEntity>,
|
||||
private readonly workspaceCacheService: WorkspaceCacheService,
|
||||
) {}
|
||||
|
||||
async upsertToCore(
|
||||
workspaceId: string,
|
||||
workflowVersions: WorkflowVersionWorkspaceEntity[],
|
||||
): Promise<void> {
|
||||
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<void> {
|
||||
if (workflowVersionIds.length === 0) {
|
||||
return;
|
||||
}
|
||||
|
||||
await this.workflowVersionRepository.delete(workspaceId, {
|
||||
id: In(workflowVersionIds),
|
||||
});
|
||||
|
||||
await this.invalidateAutomatedTriggerMaps(workspaceId);
|
||||
}
|
||||
|
||||
private async invalidateAutomatedTriggerMaps(
|
||||
workspaceId: string,
|
||||
): Promise<void> {
|
||||
await this.workspaceCacheService.invalidateAndRecompute(workspaceId, [
|
||||
'workflowAutomatedTriggerMaps',
|
||||
]);
|
||||
}
|
||||
}
|
||||
+12
-2
@@ -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 {}
|
||||
|
||||
+125
@@ -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<WorkflowVersionWorkspaceEntity>
|
||||
>,
|
||||
): Promise<void> {
|
||||
await this.upsertToCore(
|
||||
batchEvent.workspaceId,
|
||||
batchEvent.events.map((event) => event.properties.after),
|
||||
);
|
||||
}
|
||||
|
||||
@OnDatabaseBatchEvent('workflowVersion', DatabaseEventAction.UPDATED)
|
||||
async handleUpdated(
|
||||
batchEvent: CustomWorkspaceEventBatch<
|
||||
ObjectRecordUpdateEvent<WorkflowVersionWorkspaceEntity>
|
||||
>,
|
||||
): Promise<void> {
|
||||
await this.upsertToCore(
|
||||
batchEvent.workspaceId,
|
||||
batchEvent.events.map((event) => event.properties.after),
|
||||
);
|
||||
}
|
||||
|
||||
@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.id),
|
||||
);
|
||||
}
|
||||
|
||||
@OnDatabaseBatchEvent('workflowVersion', DatabaseEventAction.DESTROYED)
|
||||
async handleDestroyed(
|
||||
batchEvent: CustomWorkspaceEventBatch<
|
||||
ObjectRecordDestroyEvent<WorkflowVersionWorkspaceEntity>
|
||||
>,
|
||||
): Promise<void> {
|
||||
await this.deleteFromCore(
|
||||
batchEvent.workspaceId,
|
||||
batchEvent.events.map((event) => event.properties.before.id),
|
||||
);
|
||||
}
|
||||
|
||||
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,
|
||||
workflowVersionIds: string[],
|
||||
): Promise<void> {
|
||||
if (!isDefined(workspaceId)) {
|
||||
return;
|
||||
}
|
||||
|
||||
try {
|
||||
await this.workflowVersionCoreSyncService.deleteFromCore(
|
||||
workspaceId,
|
||||
workflowVersionIds,
|
||||
);
|
||||
} catch (error) {
|
||||
this.exceptionHandlerService.captureExceptions([error], {
|
||||
workspace: { id: workspaceId },
|
||||
});
|
||||
}
|
||||
}
|
||||
}
|
||||
+10
@@ -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 {}
|
||||
@@ -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 {}
|
||||
|
||||
Reference in New Issue
Block a user