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 {}