diff --git a/packages/@n8n/db/src/repositories/workflow-publication-outbox.repository.ts b/packages/@n8n/db/src/repositories/workflow-publication-outbox.repository.ts index c4288b42c79..16617073947 100644 --- a/packages/@n8n/db/src/repositories/workflow-publication-outbox.repository.ts +++ b/packages/@n8n/db/src/repositories/workflow-publication-outbox.repository.ts @@ -1,6 +1,7 @@ import { GlobalConfig } from '@n8n/config'; import { Service } from '@n8n/di'; import { DataSource, Repository } from '@n8n/typeorm'; +import type { EntityManager } from '@n8n/typeorm'; import { UnexpectedError } from 'n8n-workflow'; import { @@ -26,25 +27,33 @@ export class WorkflowPublicationOutboxRepository extends Repository { + async enqueue( + workflowId: string, + publishedVersionId: string, + trx?: EntityManager, + ): Promise { if (this.globalConfig.database.type === 'postgresdb') { - await this.enqueueWithPostgresUpsert(workflowId, publishedVersionId); + await this.enqueueWithPostgresUpsert(workflowId, publishedVersionId, trx ?? this.manager); return; } - await this.enqueueWithSqliteUpsert(workflowId, publishedVersionId); + await this.enqueueWithSqliteUpsert(workflowId, publishedVersionId, trx ?? this.manager); } private async enqueueWithPostgresUpsert( workflowId: string, publishedVersionId: string, + trx: EntityManager, ): Promise { const tableName = this.getTableName('workflow_publication_outbox'); // `createdAt`/`updatedAt` carry DB-level defaults, so the insert omits // them; the conflict path bumps `updatedAt` explicitly. - await this.query( + await trx.query( `INSERT INTO ${tableName} ("workflowId", "publishedVersionId", "status") VALUES ($1, $2, '${Status.Pending}') ON CONFLICT ("workflowId", "status") WHERE "status" IN ('${Status.Pending}', '${Status.InProgress}') @@ -56,10 +65,11 @@ export class WorkflowPublicationOutboxRepository extends Repository { const tableName = this.getTableName('workflow_publication_outbox'); - await this.query( + await trx.query( `INSERT INTO ${tableName} ("workflowId", "publishedVersionId", "status") VALUES (?, ?, '${Status.Pending}') ON CONFLICT ("workflowId", "status") WHERE "status" IN ('${Status.Pending}', '${Status.InProgress}') diff --git a/packages/@n8n/db/src/repositories/workflow-publish-history.repository.ts b/packages/@n8n/db/src/repositories/workflow-publish-history.repository.ts index b03857ee7e3..6a62561bb2b 100644 --- a/packages/@n8n/db/src/repositories/workflow-publish-history.repository.ts +++ b/packages/@n8n/db/src/repositories/workflow-publish-history.repository.ts @@ -1,5 +1,6 @@ import { Service } from '@n8n/di'; import { DataSource, Repository } from '@n8n/typeorm'; +import type { EntityManager } from '@n8n/typeorm'; import { WorkflowPublishHistory } from '../entities'; @@ -9,13 +10,17 @@ export class WorkflowPublishHistoryRepository extends Repository) { - await this.insert({ + async addRecord( + { + workflowId, + versionId, + event, + userId, + }: Pick, + trx?: EntityManager, + ) { + const repository = trx ? trx.getRepository(WorkflowPublishHistory) : this; + await repository.insert({ workflowId, versionId, event, diff --git a/packages/cli/src/commands/__tests__/start.test.ts b/packages/cli/src/commands/__tests__/start.test.ts index 7c43d8f1fbb..e13311ffe55 100644 --- a/packages/cli/src/commands/__tests__/start.test.ts +++ b/packages/cli/src/commands/__tests__/start.test.ts @@ -153,6 +153,7 @@ describe('Start - AuthRolesService initialization', () => { cache: { backend: 'memory' }, taskRunners: {}, expressionEngine: { engine: 'legacy', poolSize: 1, maxCodeCacheSize: 1024 }, + workflows: { useWorkflowPublicationService: false }, }; // @ts-expect-error - Accessing protected method for testing start.initCrashJournal = jest.fn().mockResolvedValue(undefined); @@ -206,6 +207,7 @@ describe('Start - AuthRolesService initialization', () => { cache: { backend: 'memory' }, taskRunners: {}, expressionEngine: { engine: 'legacy', poolSize: 1, maxCodeCacheSize: 1024 }, + workflows: { useWorkflowPublicationService: false }, }; await start.init(); @@ -240,6 +242,7 @@ describe('Start - AuthRolesService initialization', () => { cache: { backend: 'memory' }, taskRunners: {}, expressionEngine: { engine: 'legacy', poolSize: 1, maxCodeCacheSize: 1024 }, + workflows: { useWorkflowPublicationService: false }, }; await start.init(); @@ -285,6 +288,7 @@ describe('Start - AuthRolesService initialization', () => { cache: { backend: 'memory' }, taskRunners: {}, expressionEngine: { engine: 'legacy' as const, poolSize: 1, maxCodeCacheSize: 1024 }, + workflows: { useWorkflowPublicationService: false }, }; beforeEach(() => { diff --git a/packages/cli/src/commands/start.ts b/packages/cli/src/commands/start.ts index 6f6082e6971..434111c16f7 100644 --- a/packages/cli/src/commands/start.ts +++ b/packages/cli/src/commands/start.ts @@ -232,6 +232,13 @@ export class Start extends BaseCommand> { await this.initOrchestration(); } + if (this.globalConfig.workflows.useWorkflowPublicationService) { + const { WorkflowPublicationOutboxConsumer } = await import( + '@/workflows/workflow-publication-outbox-consumer' + ); + Container.get(WorkflowPublicationOutboxConsumer).init(); + } + await this.instanceSettings.initialize(Container.get(DeploymentKeyRepository)); await Container.get(JwtService).initialize(Container.get(DeploymentKeyRepository)); await Container.get(BinaryDataConfig).initialize(Container.get(DeploymentKeyRepository)); diff --git a/packages/cli/src/workflows/__tests__/workflow-publication-outbox-consumer.test.ts b/packages/cli/src/workflows/__tests__/workflow-publication-outbox-consumer.test.ts index 3c5a5591850..f3ad7116a85 100644 --- a/packages/cli/src/workflows/__tests__/workflow-publication-outbox-consumer.test.ts +++ b/packages/cli/src/workflows/__tests__/workflow-publication-outbox-consumer.test.ts @@ -13,7 +13,7 @@ import type { } from '@n8n/db'; import type { EntityManager } from '@n8n/typeorm'; import { mock } from 'jest-mock-extended'; -import type { ErrorReporter } from 'n8n-core'; +import type { ErrorReporter, InstanceSettings } from 'n8n-core'; import type { INode } from 'n8n-workflow'; import type { ActivationErrorsService } from '@/activation-errors.service'; @@ -31,6 +31,7 @@ describe('WorkflowPublicationOutboxConsumer', () => { const workflowPublishedVersionRepository = mock(); const activeWorkflowManager = mock(); const activationErrorsService = mock(); + const instanceSettings = mock(); let consumer: WorkflowPublicationOutboxConsumer; @@ -51,6 +52,7 @@ describe('WorkflowPublicationOutboxConsumer', () => { workflowPublishedVersionRepository, activeWorkflowManager, activationErrorsService, + instanceSettings, ); } diff --git a/packages/cli/src/workflows/__tests__/workflow.service.test.ts b/packages/cli/src/workflows/__tests__/workflow.service.test.ts index e1b72b37aaa..228206a8d5b 100644 --- a/packages/cli/src/workflows/__tests__/workflow.service.test.ts +++ b/packages/cli/src/workflows/__tests__/workflow.service.test.ts @@ -1,7 +1,15 @@ import type { LicenseState } from '@n8n/backend-common'; -import type { Project, User, WorkflowRepository, WorkflowPublishHistoryRepository } from '@n8n/db'; +import type { GlobalConfig, WorkflowsConfig } from '@n8n/config'; +import type { + Project, + User, + WorkflowRepository, + WorkflowPublishHistoryRepository, + WorkflowPublicationOutboxRepository, +} from '@n8n/db'; import { WorkflowEntity, WorkflowHistory } from '@n8n/db'; import type { Scope } from '@n8n/permissions'; +import type { EntityManager } from '@n8n/typeorm'; import type { MockProxy } from 'jest-mock-extended'; import { mock } from 'jest-mock-extended'; import type { IConnections, INode } from 'n8n-workflow'; @@ -10,6 +18,7 @@ import type { ActiveWorkflowManager } from '@/active-workflow-manager'; import { BadRequestError } from '@/errors/response-errors/bad-request.error'; import { UnprocessableRequestError } from '@/errors/response-errors/unprocessable.error'; import { WorkflowActivationBadRequestError } from '@/errors/response-errors/workflow-activation-bad-request.error'; +import type { EventService } from '@/events/event.service'; import type { ExternalHooks } from '@/external-hooks'; import type { RedactionEnforcementService } from '@/modules/redaction/redaction-enforcement.service'; import { userHasScopes } from '@/permissions.ee/check-access'; @@ -67,6 +76,7 @@ describe('WorkflowService', () => { mock(), // workflowFinderService mock(), // workflowPublishedVersionRepository mock(), // workflowPublishHistoryRepository + mock(), // outboxRepository Object.assign(mock(), { validateCredentialNodeRestrictions: () => ({ isValid: true }), }), // workflowValidationService @@ -218,6 +228,7 @@ describe('WorkflowService', () => { workflowFinderServiceMock, // workflowFinderService mock(), // workflowPublishedVersionRepository mock(), // workflowPublishHistoryRepository + mock(), // outboxRepository Object.assign(mock(), { validateCredentialNodeRestrictions: () => ({ isValid: true }), }), // workflowValidationService @@ -785,8 +796,11 @@ describe('WorkflowService', () => { let workflowHistoryServiceMock: MockProxy; let workflowRepositoryMock: MockProxy; let workflowPublishHistoryRepositoryMock: MockProxy; + let outboxRepositoryMock: MockProxy; + let globalConfigMock: MockProxy; let activeWorkflowManagerMock: MockProxy; let externalHooksMock: MockProxy; + let eventServiceMock: MockProxy; const WORKFLOW_ID = 'workflow-1'; const PREVIOUS_VERSION_ID = 'v1'; @@ -821,8 +835,13 @@ describe('WorkflowService', () => { workflowHistoryServiceMock = mock(); workflowRepositoryMock = mock(); workflowPublishHistoryRepositoryMock = mock(); + outboxRepositoryMock = mock(); + globalConfigMock = mock({ + workflows: mock({ useWorkflowPublicationService: false }), + }); activeWorkflowManagerMock = mock(); externalHooksMock = mock(); + eventServiceMock = mock(); workflowRepositoryMock.create.mockImplementation( (data) => Object.assign(new WorkflowEntity(), data) as WorkflowEntity, @@ -842,12 +861,13 @@ describe('WorkflowService', () => { mock(), // roleService mock(), // projectService mock(), // executionRepository - mock(), // eventService - mock(), // globalConfig + eventServiceMock, // eventService + globalConfigMock, // globalConfig mock(), // folderRepository workflowFinderServiceMock, // workflowFinderService mock(), // workflowPublishedVersionRepository workflowPublishHistoryRepositoryMock, // workflowPublishHistoryRepository + outboxRepositoryMock, // outboxRepository Object.assign(mock(), { validateCredentialNodeRestrictions: () => ({ isValid: true }), }), // workflowValidationService @@ -941,5 +961,87 @@ describe('WorkflowService', () => { expect(candidate.nodes).toBe(workflow.nodes); expect(candidate.connections).toBe(workflow.connections); }); + + test('with the publication outbox enabled, updates the version, writes history, enqueues and emits events without touching the active workflow manager', async () => { + globalConfigMock.workflows.useWorkflowPublicationService = true; + + const workflow = makeWorkflowEntity({ activeVersionId: PREVIOUS_VERSION_ID }); + const versionToActivate = makeVersionToActivate(); + workflowFinderServiceMock.findWorkflowForUser.mockResolvedValue(workflow); + workflowHistoryServiceMock.getVersion.mockResolvedValue(versionToActivate); + workflowRepositoryMock.findOne.mockResolvedValue(workflow); + externalHooksMock.run.mockResolvedValue(undefined); + + const trx = mock(); + const managerMock = mock(); + (managerMock.transaction as unknown as jest.Mock).mockImplementation( + async (runInTransaction: (entityManager: EntityManager) => Promise) => + await runInTransaction(trx), + ); + Object.defineProperty(workflowRepositoryMock, 'manager', { + value: managerMock, + configurable: true, + }); + + const addToActiveWorkflowManagerSpy = jest.spyOn( + workflowService as never, + '_addToActiveWorkflowManager', + ); + + const user = mock({ id: 'user-1' }); + + await workflowService.activateWorkflow(user, WORKFLOW_ID, { + versionId: TARGET_VERSION_ID, + }); + + // activeVersionId + active are updated inside the transaction + expect(trx.update).toHaveBeenCalledWith( + WorkflowEntity, + { id: WORKFLOW_ID }, + expect.objectContaining({ active: true, activeVersionId: TARGET_VERSION_ID }), + ); + // the outbox record is enqueued in the same transaction + expect(outboxRepositoryMock.enqueue).toHaveBeenCalledWith( + WORKFLOW_ID, + TARGET_VERSION_ID, + trx, + ); + // publish-history records (deactivated for the previous version, activated for the + // target) are written in the same transaction + expect(workflowPublishHistoryRepositoryMock.addRecord).toHaveBeenCalledWith( + expect.objectContaining({ event: 'deactivated', versionId: PREVIOUS_VERSION_ID }), + trx, + ); + expect(workflowPublishHistoryRepositoryMock.addRecord).toHaveBeenCalledWith( + expect.objectContaining({ event: 'activated', versionId: TARGET_VERSION_ID }), + trx, + ); + expect(eventServiceMock.emit).toHaveBeenNthCalledWith(1, 'workflow-deactivated', { + user, + workflowId: WORKFLOW_ID, + workflow, + publicApi: false, + deactivatedVersionId: PREVIOUS_VERSION_ID, + source: 'ui', + }); + expect(eventServiceMock.emit).toHaveBeenNthCalledWith(2, 'workflow-activated', { + user, + workflowId: WORKFLOW_ID, + workflow: expect.objectContaining({ + active: true, + activeVersionId: TARGET_VERSION_ID, + activeVersion: versionToActivate, + nodes: versionToActivate.nodes, + connections: versionToActivate.connections, + }), + publicApi: false, + source: 'ui', + }); + // trigger reapplication is deferred to the consumer + expect(addToActiveWorkflowManagerSpy).not.toHaveBeenCalled(); + expect(activeWorkflowManagerMock.add).not.toHaveBeenCalled(); + expect(activeWorkflowManagerMock.remove).not.toHaveBeenCalled(); + expect(workflowRepositoryMock.update).not.toHaveBeenCalled(); + }); }); }); diff --git a/packages/cli/src/workflows/workflow-publication-outbox-consumer.ts b/packages/cli/src/workflows/workflow-publication-outbox-consumer.ts index bbc3b6e98aa..1562535a9fc 100644 --- a/packages/cli/src/workflows/workflow-publication-outbox-consumer.ts +++ b/packages/cli/src/workflows/workflow-publication-outbox-consumer.ts @@ -12,7 +12,7 @@ import { } from '@n8n/db'; import { OnLeaderStepdown, OnLeaderTakeover, OnShutdown } from '@n8n/decorators'; import { Service } from '@n8n/di'; -import { ErrorReporter } from 'n8n-core'; +import { ErrorReporter, InstanceSettings } from 'n8n-core'; import { UnexpectedError, ensureError } from 'n8n-workflow'; import { ActivationErrorsService } from '@/activation-errors.service'; @@ -48,10 +48,17 @@ export class WorkflowPublicationOutboxConsumer { private readonly workflowPublishedVersionRepository: WorkflowPublishedVersionRepository, private readonly activeWorkflowManager: ActiveWorkflowManager, private readonly activationErrorsService: ActivationErrorsService, + private readonly instanceSettings: InstanceSettings, ) { this.logger = this.logger.scoped('workflow-publication'); } + init() { + if (this.instanceSettings.isLeader) { + this.startPolling(); + } + } + @OnLeaderTakeover() startPolling() { if (!this.workflowsConfig.useWorkflowPublicationService || this.isShuttingDown) return; diff --git a/packages/cli/src/workflows/workflow.service.ts b/packages/cli/src/workflows/workflow.service.ts index 73abd1d8756..3e1928f065f 100644 --- a/packages/cli/src/workflows/workflow.service.ts +++ b/packages/cli/src/workflows/workflow.service.ts @@ -1,15 +1,10 @@ import { UpdateWorkflowHistoryVersionDto } from '@n8n/api-types'; import { LicenseState, Logger } from '@n8n/backend-common'; import { GlobalConfig } from '@n8n/config'; -import type { - User, - ListQueryDb, - WorkflowFolderUnionFull, - WorkflowHistory, - WorkflowEntity, -} from '@n8n/db'; +import type { User, ListQueryDb, WorkflowFolderUnionFull, WorkflowHistory } from '@n8n/db'; import { SharedWorkflow, + WorkflowEntity, ExecutionRepository, FolderRepository, WorkflowTagMappingRepository, @@ -17,6 +12,7 @@ import { WorkflowRepository, WorkflowPublishedVersionRepository, WorkflowPublishHistoryRepository, + WorkflowPublicationOutboxRepository, ProjectRepository, } from '@n8n/db'; import { Container, Service } from '@n8n/di'; @@ -94,6 +90,7 @@ export class WorkflowService { private readonly workflowFinderService: WorkflowFinderService, private readonly workflowPublishedVersionRepository: WorkflowPublishedVersionRepository, private readonly workflowPublishHistoryRepository: WorkflowPublishHistoryRepository, + private readonly outboxRepository: WorkflowPublicationOutboxRepository, private readonly workflowValidationService: WorkflowValidationService, private readonly nodeTypes: NodeTypes, private readonly webhookService: WebhookService, @@ -646,16 +643,6 @@ export class WorkflowService { if (didPublish) { assert(workflow.activeVersionId !== null); - // Temporary: In the future, the workflow publication service will - // manage this mapping. We set it here for now to support incremental - // development and testing. - if (this.globalConfig.workflows.useWorkflowPublicationService) { - await this.workflowPublishedVersionRepository.setPublishedVersion( - workflowId, - workflow.activeVersionId, - ); - } - await this.workflowPublishHistoryRepository.addRecord({ workflowId, versionId: workflow.activeVersionId, @@ -806,51 +793,89 @@ export class WorkflowService { }); } - if (previousActiveVersionId) { - await this.activeWorkflowManager.remove(workflowId); - - this.eventService.emit('workflow-deactivated', { + if (this.globalConfig.workflows.useWorkflowPublicationService) { + await this._publishViaOutbox( user, workflowId, - workflow, + versionIdToActivate, + previousActiveVersionId, + workflow.updatedAt, + ); + + if (previousActiveVersionId) { + this.eventService.emit('workflow-deactivated', { + user, + workflowId, + workflow, + publicApi, + deactivatedVersionId: previousActiveVersionId, + source, + }); + } + + const activatedWorkflow = this.workflowRepository.create({ + ...workflow, + active: true, + activeVersionId: versionIdToActivate, + activeVersion: versionToActivate, + nodes: versionToActivate.nodes, + connections: versionToActivate.connections, + }); + + this.eventService.emit('workflow-activated', { + user, + workflowId, + workflow: activatedWorkflow, publicApi, - deactivatedVersionId: previousActiveVersionId, source, }); - await this.workflowPublishHistoryRepository.addRecord({ - workflowId, - versionId: previousActiveVersionId, - event: 'deactivated', - userId: user.id, + } else { + if (previousActiveVersionId) { + await this.activeWorkflowManager.remove(workflowId); + + this.eventService.emit('workflow-deactivated', { + user, + workflowId, + workflow, + publicApi, + deactivatedVersionId: previousActiveVersionId, + source, + }); + await this.workflowPublishHistoryRepository.addRecord({ + workflowId, + versionId: previousActiveVersionId, + event: 'deactivated', + userId: user.id, + }); + } + + const activationMode = previousActiveVersionId ? 'update' : 'activate'; + + await this.workflowRepository.update(workflowId, { + activeVersionId: versionIdToActivate, + active: true, + // workflow content did not change, so we keep updatedAt as is + updatedAt: workflow.updatedAt, }); + + const workflowForActivation = await this.workflowRepository.findOne({ + where: { id: workflowId }, + relations: ['activeVersion'], + }); + + if (!workflowForActivation) { + throw new NotFoundError(`Workflow with ID "${workflowId}" could not be found.`); + } + + await this._addToActiveWorkflowManager( + user, + workflowId, + workflowForActivation, + activationMode, + { source }, + ); } - const activationMode = previousActiveVersionId ? 'update' : 'activate'; - - await this.workflowRepository.update(workflowId, { - activeVersionId: versionIdToActivate, - active: true, - // workflow content did not change, so we keep updatedAt as is - updatedAt: workflow.updatedAt, - }); - - const workflowForActivation = await this.workflowRepository.findOne({ - where: { id: workflowId }, - relations: ['activeVersion'], - }); - - if (!workflowForActivation) { - throw new NotFoundError(`Workflow with ID "${workflowId}" could not be found.`); - } - - await this._addToActiveWorkflowManager( - user, - workflowId, - workflowForActivation, - activationMode, - { source }, - ); - if (options?.name !== undefined || options?.description !== undefined) { const updateFields: UpdateWorkflowHistoryVersionDto = {}; if (options.name !== undefined) updateFields.name = options.name; @@ -1318,4 +1343,55 @@ export class WorkflowService { throw new WorkflowValidationError(validation.error ?? 'Sub-workflow validation failed'); } } + + /** + * Atomically records the requested version and enqueues an outbox record. + * The publication outbox consumer reapplies the triggers and advances the + * published version asynchronously, so we do not touch the active workflow + * manager here. + */ + private async _publishViaOutbox( + user: User, + workflowId: string, + versionIdToActivate: string, + previousActiveVersionId: string | null, + updatedAt: Date, + ): Promise { + await this.workflowRepository.manager.transaction(async (trx) => { + await trx.update( + WorkflowEntity, + { id: workflowId }, + { + activeVersionId: versionIdToActivate, + active: true, + // workflow content did not change, so we keep updatedAt as is + updatedAt, + }, + ); + + if (previousActiveVersionId) { + await this.workflowPublishHistoryRepository.addRecord( + { + workflowId, + versionId: previousActiveVersionId, + event: 'deactivated', + userId: user.id, + }, + trx, + ); + } + + await this.workflowPublishHistoryRepository.addRecord( + { + workflowId, + versionId: versionIdToActivate, + event: 'activated', + userId: user.id, + }, + trx, + ); + + await this.outboxRepository.enqueue(workflowId, versionIdToActivate, trx); + }); + } } diff --git a/packages/cli/test/integration/database/repositories/workflow-publication-outbox.repository.test.ts b/packages/cli/test/integration/database/repositories/workflow-publication-outbox.repository.test.ts index 0f60609d4b9..3b5171fd4d1 100644 --- a/packages/cli/test/integration/database/repositories/workflow-publication-outbox.repository.test.ts +++ b/packages/cli/test/integration/database/repositories/workflow-publication-outbox.repository.test.ts @@ -45,6 +45,27 @@ describe('WorkflowPublicationOutboxRepository', () => { expect(claimedAgain).toBeNull(); }); + it('enqueues within a provided transaction and is visible once it commits', async () => { + await repository.manager.transaction(async (trx) => { + await repository.enqueue('wf-1', 'v-1', trx); + }); + + const claimed = await repository.claimNextPendingRecord(); + expect(claimed?.workflowId).toBe('wf-1'); + expect(claimed?.publishedVersionId).toBe('v-1'); + }); + + it('discards the enqueued record when the surrounding transaction rolls back', async () => { + await expect( + repository.manager.transaction(async (trx) => { + await repository.enqueue('wf-1', 'v-1', trx); + throw new Error('rollback'); + }), + ).rejects.toThrow('rollback'); + + expect(await repository.claimNextPendingRecord()).toBeNull(); + }); + it('claims pending records in FIFO order', async () => { await repository.enqueue('wf-1', 'v-1'); await repository.enqueue('wf-2', 'v-1'); diff --git a/packages/cli/test/integration/workflows/workflow.service.test.ts b/packages/cli/test/integration/workflows/workflow.service.test.ts index 158ff96eb5b..6e8d8408f42 100644 --- a/packages/cli/test/integration/workflows/workflow.service.test.ts +++ b/packages/cli/test/integration/workflows/workflow.service.test.ts @@ -14,6 +14,7 @@ import { type WorkflowEntity, WorkflowPublishedVersionRepository, WorkflowPublishHistoryRepository, + WorkflowPublicationOutboxRepository, WorkflowRepository, ProjectRepository, } from '@n8n/db'; @@ -44,6 +45,7 @@ let workflowRepository: WorkflowRepository; let workflowService: WorkflowService; let workflowPublishedVersionRepository: WorkflowPublishedVersionRepository; let workflowPublishHistoryRepository: WorkflowPublishHistoryRepository; +let outboxRepository: WorkflowPublicationOutboxRepository; let workflowHistoryService: WorkflowHistoryService; const loggerMock = mock(); const activeWorkflowManager = mockInstance(ActiveWorkflowManager); @@ -60,6 +62,7 @@ beforeAll(async () => { workflowRepository = Container.get(WorkflowRepository); workflowPublishedVersionRepository = Container.get(WorkflowPublishedVersionRepository); workflowPublishHistoryRepository = Container.get(WorkflowPublishHistoryRepository); + outboxRepository = Container.get(WorkflowPublicationOutboxRepository); workflowHistoryService = Container.get(WorkflowHistoryService); workflowService = new WorkflowService( loggerMock, @@ -81,6 +84,7 @@ beforeAll(async () => { Container.get(WorkflowFinderService), workflowPublishedVersionRepository, workflowPublishHistoryRepository, + outboxRepository, workflowValidationService, nodeTypes, webhookServiceMock, @@ -104,6 +108,7 @@ afterEach(async () => { 'SharedWorkflow', 'ProjectRelation', 'WorkflowPublishedVersion', + 'WorkflowPublicationOutbox', 'WorkflowEntity', 'WorkflowHistory', 'WorkflowPublishHistory', @@ -515,7 +520,7 @@ describe('deactivateWorkflow()', () => { }); }); -describe('workflow_published_version table population', () => { +describe('workflow publication outbox', () => { describe('when feature flag is enabled', () => { beforeEach(() => { globalConfig.workflows.useWorkflowPublicationService = true; @@ -525,19 +530,31 @@ describe('workflow_published_version table population', () => { globalConfig.workflows.useWorkflowPublicationService = false; }); - test('should write to workflow_published_version on activation', async () => { + test('should set the active version and enqueue a pending outbox record on activation', async () => { const owner = await createOwner(); const workflow = await createWorkflowWithHistory({}, owner); await workflowService.activateWorkflow(owner, workflow.id); + const updated = await workflowRepository.findOne({ where: { id: workflow.id } }); + expect(updated?.active).toBe(true); + expect(updated?.activeVersionId).toBe(workflow.versionId); + + const outboxRecord = await outboxRepository.findOne({ + where: { workflowId: workflow.id }, + }); + expect(outboxRecord?.publishedVersionId).toBe(workflow.versionId); + expect(outboxRecord?.status).toBe('pending'); + + // Publication is async: the published version is advanced by the + // outbox consumer, not synchronously by the service. const publishedVersion = await workflowPublishedVersionRepository.findOne({ where: { workflowId: workflow.id }, }); - expect(publishedVersion?.publishedVersionId).toBe(workflow.versionId); + expect(publishedVersion).toBeNull(); }); - test('should update workflow_published_version when activating a new version', async () => { + test('should supersede the pending outbox record when activating a new version', async () => { const owner = await createOwner(); const workflow = await createWorkflowWithHistory({}, owner); @@ -550,10 +567,13 @@ describe('workflow_published_version table population', () => { versionId: newVersionId, }); - const publishedVersion = await workflowPublishedVersionRepository.findOne({ - where: { workflowId: workflow.id }, - }); - expect(publishedVersion?.publishedVersionId).toBe(newVersionId); + const outboxRecords = await outboxRepository.find({ where: { workflowId: workflow.id } }); + expect(outboxRecords).toHaveLength(1); + expect(outboxRecords[0].publishedVersionId).toBe(newVersionId); + expect(outboxRecords[0].status).toBe('pending'); + + const updated = await workflowRepository.findOne({ where: { id: workflow.id } }); + expect(updated?.activeVersionId).toBe(newVersionId); }); test('should remove workflow_published_version on deactivation', async () => { @@ -561,6 +581,8 @@ describe('workflow_published_version table population', () => { const workflow = await createWorkflowWithHistory({}, owner); await workflowService.activateWorkflow(owner, workflow.id); + // Simulate the outbox consumer having advanced the published version. + await workflowPublishedVersionRepository.setPublishedVersion(workflow.id, workflow.versionId); await workflowService.deactivateWorkflow(owner, workflow.id); @@ -575,11 +597,8 @@ describe('workflow_published_version table population', () => { const workflow = await createWorkflowWithHistory({}, owner); await workflowService.activateWorkflow(owner, workflow.id); - - const publishedVersionBefore = await workflowPublishedVersionRepository.findOne({ - where: { workflowId: workflow.id }, - }); - expect(publishedVersionBefore).not.toBeNull(); + // Simulate the outbox consumer having advanced the published version. + await workflowPublishedVersionRepository.setPublishedVersion(workflow.id, workflow.versionId); await workflowService.archive(owner, workflow.id); @@ -591,12 +610,15 @@ describe('workflow_published_version table population', () => { }); describe('when feature flag is disabled', () => { - test('should not write to workflow_published_version on activation', async () => { + test('should not enqueue an outbox record on activation', async () => { const owner = await createOwner(); const workflow = await createWorkflowWithHistory({}, owner); await workflowService.activateWorkflow(owner, workflow.id); + const outboxCount = await outboxRepository.count(); + expect(outboxCount).toBe(0); + const publishedVersion = await workflowPublishedVersionRepository.findOne({ where: { workflowId: workflow.id }, }); diff --git a/packages/testing/playwright/tests/e2e/api/workflow-publication-service.spec.ts b/packages/testing/playwright/tests/e2e/api/workflow-publication-service.spec.ts index 5f57a4d049d..9c4dd17ff61 100644 --- a/packages/testing/playwright/tests/e2e/api/workflow-publication-service.spec.ts +++ b/packages/testing/playwright/tests/e2e/api/workflow-publication-service.spec.ts @@ -5,6 +5,9 @@ test.use({ env: { TEST_ISOLATION: 'workflow-publication-service', N8N_USE_WORKFLOW_PUBLICATION_SERVICE: 'true', + // Activation is applied asynchronously by the publication outbox + // consumer, so poll frequently to keep the test fast. + N8N_WORKFLOW_PUBLICATION_OUTBOX_POLL_INTERVAL_MS: '250', }, }, });