From ef2f21fed3894ca9f64fa4a6f2e891affe6a0322 Mon Sep 17 00:00:00 2001 From: Tomi Turtiainen <10324676+tomi@users.noreply.github.com> Date: Wed, 10 Jun 2026 17:57:59 +0300 Subject: [PATCH] feat(core): Add trigger diffing to workflow publication (no-changelog) (#31971) --- .../__tests__/active-workflow-manager.test.ts | 45 ++- packages/cli/src/active-workflow-manager.ts | 208 +++++++++++- packages/cli/src/webhooks/webhook.service.ts | 10 + .../workflows/__tests__/trigger-diff.test.ts | 75 +++++ ...rkflow-publication-outbox-consumer.test.ts | 313 +++++++++++------- packages/cli/src/workflows/trigger-diff.ts | 48 +++ .../workflow-publication-outbox-consumer.ts | 140 +++++--- ...rkflow-publication-outbox-consumer.test.ts | 164 +++++++++ .../active-workflow-triggers.test.ts | 149 ++++++++- .../__tests__/scheduled-task-manager.test.ts | 61 ++++ .../active-workflow-triggers.ts | 112 ++++++- .../scheduled-task-manager.ts | 33 ++ .../workflow-active-triggers-state.ts | 15 + packages/workflow/src/workflow-diff.ts | 2 +- 14 files changed, 1181 insertions(+), 194 deletions(-) create mode 100644 packages/cli/src/workflows/__tests__/trigger-diff.test.ts create mode 100644 packages/cli/src/workflows/trigger-diff.ts create mode 100644 packages/cli/test/integration/workflows/workflow-publication-outbox-consumer.test.ts diff --git a/packages/cli/src/__tests__/active-workflow-manager.test.ts b/packages/cli/src/__tests__/active-workflow-manager.test.ts index b3b8d94372e..319e1f5e11d 100644 --- a/packages/cli/src/__tests__/active-workflow-manager.test.ts +++ b/packages/cli/src/__tests__/active-workflow-manager.test.ts @@ -67,6 +67,49 @@ describe('ActiveWorkflowManager', () => { ); }); + describe('getEnabledTriggerNodes', () => { + function node(id: string, type: string, overrides: Partial = {}): INode { + return { + id, + name: id, + type, + typeVersion: 1, + position: [0, 0], + parameters: {}, + ...overrides, + }; + } + + beforeEach(() => { + const description = { properties: [] }; + nodeTypes.getByNameAndVersion.mockImplementation((type: string) => { + if (type === 'trigger') return { description, trigger: jest.fn() } as never; + if (type === 'poll') return { description, poll: jest.fn() } as never; + if (type === 'webhook') return { description, webhook: jest.fn() } as never; + return { description } as never; + }); + }); + + test('returns enabled trigger, poll and webhook nodes, excluding regular and disabled nodes', () => { + const result = activeWorkflowManager.getEnabledTriggerNodes({ + nodes: [ + node('t', 'trigger'), + node('p', 'poll'), + node('w', 'webhook'), + node('regular', 'n8n-nodes-base.set'), + node('disabled', 'trigger', { disabled: true }), + ], + connections: {}, + }); + + expect(result.map((n) => n.id).sort()).toEqual(['p', 't', 'w']); + }); + + test('returns an empty array when the version is null', () => { + expect(activeWorkflowManager.getEnabledTriggerNodes(null)).toEqual([]); + }); + }); + describe('shouldAddWebhooks', () => { describe('if leader', () => { beforeAll(() => { @@ -694,7 +737,7 @@ describe('ActiveWorkflowManager', () => { } as unknown as ConstructorParameters[4], ); - await realActiveWorkflowTriggers.add( + await realActiveWorkflowTriggers.addAllTriggers( 'wf-1', workflow, additionalData, diff --git a/packages/cli/src/active-workflow-manager.ts b/packages/cli/src/active-workflow-manager.ts index 1a7b9145eb2..3c364f0f49f 100644 --- a/packages/cli/src/active-workflow-manager.ts +++ b/packages/cli/src/active-workflow-manager.ts @@ -24,6 +24,7 @@ import { } from 'n8n-core'; import type { ExecutionError, + IConnections, IDeferredPromise, IExecuteResponsePromiseData, INode, @@ -159,10 +160,17 @@ export class ActiveWorkflowManager { additionalData: IWorkflowExecuteAdditionalData, mode: WorkflowExecuteMode, activation: WorkflowActivateMode, + nodeIds?: Set, ) { - const webhooks = WebhookHelpers.getWorkflowWebhooks(workflow, additionalData, undefined, true); + let webhooks = WebhookHelpers.getWorkflowWebhooks(workflow, additionalData, undefined, true); let path = ''; + if (nodeIds) { + webhooks = webhooks.filter((webhookData) => + nodeIds.has(workflow.getNode(webhookData.node)?.id ?? ''), + ); + } + if (webhooks.length === 0) return false; for (const webhookData of webhooks) { @@ -276,13 +284,29 @@ export class ActiveWorkflowManager { settings: workflowData.settings, }); - const mode = 'internal'; - const additionalData = await WorkflowExecuteAdditionalData.getBase({ workflowId: workflow.id, workflowSettings: workflowData.settings, }); + await this.deregisterWebhooks(workflow, additionalData); + + await this.webhookService.deleteWorkflowWebhooks(workflowId); + } + + /** + * Deregisters a workflow's webhooks from external services and persists any + * resulting static data. When `nodeIds` is given, only the webhooks of those + * nodes are deregistered. Returns the names of the nodes whose webhooks were + * deregistered. + */ + private async deregisterWebhooks( + workflow: Workflow, + additionalData: IWorkflowExecuteAdditionalData, + nodeIds?: Set, + ) { + const removedNodeNames: string[] = []; + await workflow.expression.acquireIsolate(); try { const webhooks = WebhookHelpers.getWorkflowWebhooks( @@ -293,7 +317,11 @@ export class ActiveWorkflowManager { ); for (const webhookData of webhooks) { - await this.webhookService.deleteWebhook(workflow, webhookData, mode, 'update'); + if (nodeIds && !nodeIds.has(workflow.getNode(webhookData.node)?.id ?? '')) { + continue; + } + await this.webhookService.deleteWebhook(workflow, webhookData, 'internal', 'update'); + removedNodeNames.push(webhookData.node); } } finally { await workflow.expression.releaseIsolate(); @@ -301,7 +329,7 @@ export class ActiveWorkflowManager { await this.workflowStaticDataService.saveStaticData(workflow); - await this.webhookService.deleteWorkflowWebhooks(workflowId); + return removedNodeNames; } /** @@ -838,6 +866,162 @@ export class ActiveWorkflowManager { return added; } + /** + * Returns the enabled trigger-like nodes (active, poll, schedule and webhook + * triggers) of a workflow version. Disabled nodes are excluded, so the result + * is the set of nodes that actually drive trigger registration. Used to + * compute the trigger-level diff during publication. + */ + getEnabledTriggerNodes(version: { nodes: INode[]; connections: IConnections } | null): INode[] { + if (!version) return []; + + const workflow = new Workflow({ + id: 'trigger-diff', + name: 'trigger-diff', + nodes: version.nodes, + connections: version.connections, + active: false, + nodeTypes: this.nodeTypes, + }); + + return workflow.queryNodes( + (nodeType) => !!nodeType.trigger || !!nodeType.poll || !!nodeType.webhook, + ); + } + + /** + * Registers only the given trigger nodes (webhook and non-webhook) of the + * given workflow version, leaving any other already-active triggers + * untouched. The "add" side of a publication trigger diff; runs on the leader + * after the published version has been advanced. + */ + async addTriggerNodes( + dbWorkflow: WorkflowEntity, + version: { nodes: INode[]; connections: IConnections }, + nodeIds: Set, + ) { + const { nodes, connections } = version; + dbWorkflow.nodes = nodes; + dbWorkflow.connections = connections; + + const workflow = new Workflow({ + id: dbWorkflow.id, + name: dbWorkflow.name, + nodes, + connections, + active: true, + nodeTypes: this.nodeTypes, + staticData: dbWorkflow.staticData, + settings: dbWorkflow.settings, + }); + + const additionalData = await WorkflowExecuteAdditionalData.getBase({ + workflowId: workflow.id, + workflowSettings: dbWorkflow.settings, + }); + + let triggerCount = 0; + await workflow.expression.acquireIsolate(); + try { + if (this.shouldAddWebhooks('update')) { + await this.addWebhooks(workflow, additionalData, 'trigger', 'update', nodeIds); + } + + if (this.shouldAddNonWebhookTriggers()) { + const resolveWorkflowData = this.workflowsConfig.useWorkflowPublicationService + ? async () => await this.loadPublishedWorkflowData(dbWorkflow) + : async () => dbWorkflow as IWorkflowBase; + + await this.addNonWebhookTriggers(dbWorkflow, workflow, { + activationMode: 'update', + executionMode: 'trigger', + additionalData, + resolveWorkflowData, + nodeIds, + }); + } + + triggerCount = this.countTriggers(workflow, additionalData); + } finally { + await workflow.expression.releaseIsolate(); + } + + await Promise.all([ + this.workflowRepository.updateWorkflowTriggerCount(workflow.id, triggerCount), + this.workflowStaticDataService.saveStaticData(workflow), + ]); + } + + /** + * Recomputes the persisted trigger count for a workflow version without + * registering any triggers. Used when publication only removes triggers. + */ + async updateWorkflowTriggerCount( + dbWorkflow: WorkflowEntity, + version: { nodes: INode[]; connections: IConnections }, + ) { + const workflow = new Workflow({ + id: dbWorkflow.id, + name: dbWorkflow.name, + nodes: version.nodes, + connections: version.connections, + active: true, + nodeTypes: this.nodeTypes, + staticData: dbWorkflow.staticData, + settings: dbWorkflow.settings, + }); + + const additionalData = await WorkflowExecuteAdditionalData.getBase({ + workflowId: workflow.id, + workflowSettings: dbWorkflow.settings, + }); + + let triggerCount = 0; + await workflow.expression.acquireIsolate(); + try { + triggerCount = this.countTriggers(workflow, additionalData); + } finally { + await workflow.expression.releaseIsolate(); + } + + await this.workflowRepository.updateWorkflowTriggerCount(workflow.id, triggerCount); + } + + /** + * Deregisters only the given trigger nodes (webhook and non-webhook) of the + * given workflow version, leaving the rest active. The "remove" side of a + * publication trigger diff; the caller passes the currently published version + * so the right webhooks are deregistered. + */ + async removeTriggerNodes( + dbWorkflow: WorkflowEntity, + version: { nodes: INode[]; connections: IConnections }, + nodeIds: Set, + ) { + if (nodeIds.size === 0) return; + + const workflow = new Workflow({ + id: dbWorkflow.id, + name: dbWorkflow.name, + nodes: version.nodes, + connections: version.connections, + active: true, + nodeTypes: this.nodeTypes, + staticData: dbWorkflow.staticData, + settings: dbWorkflow.settings, + }); + + const additionalData = await WorkflowExecuteAdditionalData.getBase({ + workflowId: workflow.id, + workflowSettings: dbWorkflow.settings, + }); + + const removedNodeNames = await this.deregisterWebhooks(workflow, additionalData, nodeIds); + await this.webhookService.deleteWorkflowWebhooksForNodes(dbWorkflow.id, removedNodeNames); + + await this.activeWorkflowTriggers.removeTriggers(dbWorkflow.id, nodeIds); + } + @OnPubSubEvent('display-workflow-activation', { instanceType: 'main' }) handleDisplayWorkflowActivation({ workflowId, @@ -1130,11 +1314,13 @@ export class ActiveWorkflowManager { executionMode, additionalData, resolveWorkflowData, + nodeIds, }: { activationMode: WorkflowActivateMode; executionMode: WorkflowExecuteMode; additionalData: IWorkflowExecuteAdditionalData; resolveWorkflowData: () => Promise; + nodeIds?: Set; }, ) { const getTriggerFunctions = this.getExecuteTriggerFunctions( @@ -1153,13 +1339,21 @@ export class ActiveWorkflowManager { resolveWorkflowData, ); - if (workflow.getTriggerNodes().length === 0 && workflow.getPollNodes().length === 0) { + const triggerAndPollNodeIds = [...workflow.getTriggerNodes(), ...workflow.getPollNodes()].map( + (node) => node.id, + ); + const nodeIdsToAdd = nodeIds + ? triggerAndPollNodeIds.filter((id) => nodeIds.has(id)) + : triggerAndPollNodeIds; + + if (nodeIdsToAdd.length === 0) { return false; } - await this.activeWorkflowTriggers.add( + await this.activeWorkflowTriggers.addTriggers( workflow.id, workflow, + nodeIdsToAdd, additionalData, executionMode, activationMode, diff --git a/packages/cli/src/webhooks/webhook.service.ts b/packages/cli/src/webhooks/webhook.service.ts index 9cf1dafd651..cc7268e602f 100644 --- a/packages/cli/src/webhooks/webhook.service.ts +++ b/packages/cli/src/webhooks/webhook.service.ts @@ -176,6 +176,16 @@ export class WebhookService { return await this.deleteWebhooks(webhooks); } + /** Delete the webhooks registered for the given nodes of a workflow. */ + async deleteWorkflowWebhooksForNodes(workflowId: string, nodeNames: string[]) { + if (nodeNames.length === 0) return; + + const webhooks = await this.webhookRepository.findBy({ workflowId }); + const toDelete = webhooks.filter((webhook) => nodeNames.includes(webhook.node)); + + return await this.deleteWebhooks(toDelete); + } + private async deleteWebhooks(webhooks: WebhookEntity[]) { void this.cacheService.deleteMany(webhooks.map((w) => w.cacheKey)); diff --git a/packages/cli/src/workflows/__tests__/trigger-diff.test.ts b/packages/cli/src/workflows/__tests__/trigger-diff.test.ts new file mode 100644 index 00000000000..6c5c5a05e7b --- /dev/null +++ b/packages/cli/src/workflows/__tests__/trigger-diff.test.ts @@ -0,0 +1,75 @@ +import type { INode } from 'n8n-workflow'; + +import { computeTriggerDiff } from '@/workflows/trigger-diff'; + +describe('computeTriggerDiff', () => { + function makeNode(id: string, overrides: Partial = {}): INode { + return { + id, + name: id, + type: 'n8n-nodes-base.scheduleTrigger', + typeVersion: 1, + position: [0, 0], + parameters: {}, + ...overrides, + }; + } + + test('returns empty diff when trigger sets are identical', () => { + const a = makeNode('a'); + + const diff = computeTriggerDiff([a], [{ ...a }]); + + expect(diff).toEqual({ toAdd: new Set(), toRemove: new Set() }); + }); + + test('detects added triggers', () => { + const diff = computeTriggerDiff([makeNode('a')], [makeNode('a'), makeNode('b')]); + + expect(diff).toEqual({ toAdd: new Set(['b']), toRemove: new Set() }); + }); + + test('detects removed triggers', () => { + const diff = computeTriggerDiff([makeNode('a'), makeNode('b')], [makeNode('a')]); + + expect(diff).toEqual({ toAdd: new Set(), toRemove: new Set(['b']) }); + }); + + test('treats a parameter change as a modification (remove-then-add)', () => { + const before = makeNode('a', { parameters: { interval: 1 } }); + const after = makeNode('a', { parameters: { interval: 5 } }); + + const diff = computeTriggerDiff([before], [after]); + + expect(diff).toEqual({ toAdd: new Set(['a']), toRemove: new Set(['a']) }); + }); + + test('treats a typeVersion change as a modification', () => { + const diff = computeTriggerDiff( + [makeNode('a', { typeVersion: 1 })], + [makeNode('a', { typeVersion: 2 })], + ); + + expect(diff).toEqual({ toAdd: new Set(['a']), toRemove: new Set(['a']) }); + }); + + test('handles a mix of added, removed, modified and unchanged triggers', () => { + const unchanged = makeNode('unchanged'); + const removed = makeNode('removed'); + const modifiedBefore = makeNode('modified', { parameters: { value: 1 } }); + const modifiedAfter = makeNode('modified', { parameters: { value: 2 } }); + const added = makeNode('added'); + + const diff = computeTriggerDiff( + [unchanged, removed, modifiedBefore], + [{ ...unchanged }, modifiedAfter, added], + ); + + expect(Array.from(diff.toAdd).sort()).toEqual(['added', 'modified']); + expect(Array.from(diff.toRemove).sort()).toEqual(['modified', 'removed']); + }); + + test('returns empty diff for two empty trigger sets', () => { + expect(computeTriggerDiff([], [])).toEqual({ toAdd: new Set(), toRemove: new Set() }); + }); +}); 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 b3a0e5096b0..3c5a5591850 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 @@ -3,13 +3,18 @@ import type { WorkflowsConfig } from '@n8n/config'; import { WorkflowPublishedVersion } from '@n8n/db'; import type { WorkflowEntity, + WorkflowHistory, WorkflowPublicationOutbox, WorkflowPublicationOutboxRepository, + WorkflowPublishedVersion as WorkflowPublishedVersionEntity, + WorkflowPublishedVersionRepository, + WorkflowHistoryRepository, WorkflowRepository, } from '@n8n/db'; import type { EntityManager } from '@n8n/typeorm'; import { mock } from 'jest-mock-extended'; import type { ErrorReporter } from 'n8n-core'; +import type { INode } from 'n8n-workflow'; import type { ActivationErrorsService } from '@/activation-errors.service'; import type { ActiveWorkflowManager } from '@/active-workflow-manager'; @@ -22,9 +27,10 @@ describe('WorkflowPublicationOutboxConsumer', () => { const errorReporter = mock(); const outboxRepository = mock(); const workflowRepository = mock(); + const workflowHistoryRepository = mock(); + const workflowPublishedVersionRepository = mock(); const activeWorkflowManager = mock(); const activationErrorsService = mock(); - const entityManager = mock(); let consumer: WorkflowPublicationOutboxConsumer; @@ -41,6 +47,8 @@ describe('WorkflowPublicationOutboxConsumer', () => { errorReporter, outboxRepository, workflowRepository, + workflowHistoryRepository, + workflowPublishedVersionRepository, activeWorkflowManager, activationErrorsService, ); @@ -65,26 +73,68 @@ describe('WorkflowPublicationOutboxConsumer', () => { return { id: 'wf-1', active: true, - activeVersionId: 'v-1', + activeVersionId: 'v-2', ...overrides, } as WorkflowEntity; } + function makeVersion(versionId: string): WorkflowHistory { + return { + versionId, + workflowId: 'wf-1', + nodes: [], + connections: {}, + } as unknown as WorkflowHistory; + } + + function triggerNode(id: string, overrides: Partial = {}): INode { + return { + id, + name: id, + type: 'n8n-nodes-base.scheduleTrigger', + typeVersion: 1, + position: [0, 0], + parameters: {}, + ...overrides, + }; + } + + /** The `workflow_published_version` mapping read by `resolveVersions`. */ + function makePublishedVersion( + publishedVersion: WorkflowHistory | null, + ): WorkflowPublishedVersionEntity { + return { + workflowId: 'wf-1', + publishedVersionId: publishedVersion?.versionId ?? 'v-1', + publishedVersion, + } as unknown as WorkflowPublishedVersionEntity; + } + + const newVersion = makeVersion('v-2'); + const oldVersion = makeVersion('v-1'); + + /** Drives the trigger diff: first call returns old triggers, second returns new. */ + function setTriggerSets(oldTriggers: INode[], newTriggers: INode[]) { + activeWorkflowManager.getEnabledTriggerNodes + .mockReturnValueOnce(oldTriggers) + .mockReturnValueOnce(newTriggers); + } + /** Pin all tracked mock functions so jest-mock-extended Proxy returns stable instances. */ function setupDefaultMocks() { - workflowRepository.findById.mockResolvedValue(null); + workflowRepository.findOneBy.mockResolvedValue(makeWorkflow({ activeVersionId: 'v-1' })); + workflowPublishedVersionRepository.findOne.mockResolvedValue(makePublishedVersion(oldVersion)); workflowRepository.update.mockResolvedValue({} as never); - activeWorkflowManager.remove.mockResolvedValue(undefined); - activeWorkflowManager.clearWebhooks.mockResolvedValue(undefined); - activeWorkflowManager.removeActivationError.mockResolvedValue(undefined); - activeWorkflowManager.removeNonWebhookTriggers.mockResolvedValue(undefined); - activeWorkflowManager.add.mockResolvedValue({ webhooks: true, triggersAndPollers: true }); + workflowHistoryRepository.findOneBy.mockResolvedValue(newVersion); + activeWorkflowManager.getEnabledTriggerNodes.mockReturnValue([]); + activeWorkflowManager.addTriggerNodes.mockResolvedValue(undefined); + activeWorkflowManager.removeTriggerNodes.mockResolvedValue(undefined); + activeWorkflowManager.updateWorkflowTriggerCount.mockResolvedValue(undefined); outboxRepository.claimNextPendingRecord.mockResolvedValue(null); outboxRepository.markCompleted.mockResolvedValue(undefined); outboxRepository.markFailed.mockResolvedValue(undefined); - entityManager.upsert.mockResolvedValue({} as never); Object.defineProperty(outboxRepository, 'manager', { - value: entityManager, + value: mock({ upsert: jest.fn() }), writable: true, }); activationErrorsService.register.mockResolvedValue(undefined); @@ -169,131 +219,164 @@ describe('WorkflowPublicationOutboxConsumer', () => { }); describe('processRecord', () => { - test('removes old triggers, updates version, adds new triggers, finalizes', async () => { - const record = makeRecord(); - workflowRepository.findById.mockResolvedValue(makeWorkflow()); + test('marks completed when workflow not found', async () => { + workflowRepository.findOneBy.mockResolvedValue(null); - const callOrder: string[] = []; - activeWorkflowManager.clearWebhooks.mockImplementation(async () => { - callOrder.push('clearWebhooks'); - return await Promise.resolve(); - }); - activeWorkflowManager.removeNonWebhookTriggers.mockImplementation(async () => { - callOrder.push('removeNonWebhookTriggers'); - return await Promise.resolve(); - }); - entityManager.upsert.mockImplementation(async () => { - callOrder.push('advanceVersion'); - return await Promise.resolve({} as never); - }); - activeWorkflowManager.add.mockImplementation(async () => { - callOrder.push('add'); - return await Promise.resolve({ webhooks: true, triggersAndPollers: true }); - }); + await consumer.processRecord(makeRecord()); - await consumer.processRecord(record); + expect(outboxRepository.markCompleted).toHaveBeenCalledWith(1); + expect(activeWorkflowManager.getEnabledTriggerNodes).not.toHaveBeenCalled(); + expect(activeWorkflowManager.addTriggerNodes).not.toHaveBeenCalled(); + expect(activeWorkflowManager.removeTriggerNodes).not.toHaveBeenCalled(); + }); - expect(activeWorkflowManager.clearWebhooks).toHaveBeenCalledWith('wf-1'); - expect(activeWorkflowManager.removeActivationError).toHaveBeenCalledWith('wf-1'); - expect(activeWorkflowManager.removeNonWebhookTriggers).toHaveBeenCalledWith('wf-1'); - expect(workflowRepository.update).not.toHaveBeenCalledWith('wf-1', { - activeVersionId: 'v-2', - }); - expect(entityManager.upsert).toHaveBeenCalledWith( + test('marks completed when workflow is no longer active', async () => { + workflowRepository.findOneBy.mockResolvedValue( + makeWorkflow({ active: false, activeVersionId: null }), + ); + + await consumer.processRecord(makeRecord()); + + expect(outboxRepository.markCompleted).toHaveBeenCalledWith(1); + expect(outboxRepository.manager.upsert).not.toHaveBeenCalled(); + expect(activeWorkflowManager.getEnabledTriggerNodes).not.toHaveBeenCalled(); + expect(activeWorkflowManager.addTriggerNodes).not.toHaveBeenCalled(); + expect(activeWorkflowManager.removeTriggerNodes).not.toHaveBeenCalled(); + }); + + test('marks failed when the published version is not found', async () => { + workflowHistoryRepository.findOneBy.mockResolvedValue(null); + + await consumer.processRecord(makeRecord()); + + expect(outboxRepository.markFailed).toHaveBeenCalledWith(1, 'Published version not found'); + expect(activeWorkflowManager.getEnabledTriggerNodes).not.toHaveBeenCalled(); + expect(activeWorkflowManager.addTriggerNodes).not.toHaveBeenCalled(); + }); + + test('advances the published version and finalizes when no triggers changed', async () => { + const trigger = triggerNode('a'); + setTriggerSets([trigger], [{ ...trigger }]); + + await consumer.processRecord(makeRecord()); + + expect(outboxRepository.manager.upsert).toHaveBeenCalledWith( WorkflowPublishedVersion, { workflowId: 'wf-1', publishedVersionId: 'v-2' }, ['workflowId'], ); - expect(activeWorkflowManager.add).toHaveBeenCalledWith('wf-1', 'update', undefined, { - shouldPublish: false, - }); + expect(activeWorkflowManager.removeTriggerNodes).not.toHaveBeenCalled(); + expect(activeWorkflowManager.addTriggerNodes).not.toHaveBeenCalled(); expect(outboxRepository.markCompleted).toHaveBeenCalledWith(1); expect(activationErrorsService.deregister).toHaveBeenCalledWith('wf-1'); - expect(callOrder).toEqual([ - 'clearWebhooks', - 'removeNonWebhookTriggers', - 'advanceVersion', - 'add', - ]); }); - test('marks completed when workflow not found', async () => { - const record = makeRecord(); - workflowRepository.findById.mockResolvedValue(null); + test('registers only added triggers', async () => { + setTriggerSets([triggerNode('a')], [triggerNode('a'), triggerNode('b')]); - await consumer.processRecord(record); + await consumer.processRecord(makeRecord()); + expect(activeWorkflowManager.removeTriggerNodes).not.toHaveBeenCalled(); + expect(activeWorkflowManager.addTriggerNodes).toHaveBeenCalledWith( + expect.objectContaining({ id: 'wf-1' }), + newVersion, + new Set(['b']), + ); + expect(outboxRepository.manager.upsert).toHaveBeenCalled(); expect(outboxRepository.markCompleted).toHaveBeenCalledWith(1); - expect(activeWorkflowManager.clearWebhooks).not.toHaveBeenCalled(); - expect(activeWorkflowManager.removeNonWebhookTriggers).not.toHaveBeenCalled(); - expect(activeWorkflowManager.add).not.toHaveBeenCalled(); }); - test('marks completed when workflow is not active', async () => { - const record = makeRecord(); - workflowRepository.findById.mockResolvedValue( - makeWorkflow({ active: false, activeVersionId: null }), + test('deregisters only removed triggers', async () => { + setTriggerSets([triggerNode('a'), triggerNode('b')], [triggerNode('a')]); + + await consumer.processRecord(makeRecord()); + + expect(activeWorkflowManager.removeTriggerNodes).toHaveBeenCalledWith( + expect.objectContaining({ id: 'wf-1' }), + oldVersion, + new Set(['b']), + ); + expect(activeWorkflowManager.addTriggerNodes).not.toHaveBeenCalled(); + expect(activeWorkflowManager.updateWorkflowTriggerCount).toHaveBeenCalledWith( + expect.objectContaining({ id: 'wf-1' }), + newVersion, + ); + expect(outboxRepository.markCompleted).toHaveBeenCalledWith(1); + }); + + test('reapplies modified triggers as remove-then-add, advancing in between', async () => { + setTriggerSets( + [triggerNode('a', { parameters: { interval: 1 } })], + [triggerNode('a', { parameters: { interval: 5 } })], ); - await consumer.processRecord(record); - - expect(outboxRepository.markCompleted).toHaveBeenCalledWith(1); - expect(activeWorkflowManager.clearWebhooks).not.toHaveBeenCalled(); - expect(activeWorkflowManager.removeNonWebhookTriggers).not.toHaveBeenCalled(); - }); - - test('finalizes without trigger reapply when version already matches', async () => { - const record = makeRecord({ publishedVersionId: 'v-1' }); - workflowRepository.findById.mockResolvedValue(makeWorkflow({ activeVersionId: 'v-1' })); - - await consumer.processRecord(record); - - expect(outboxRepository.markCompleted).toHaveBeenCalledWith(1); - expect(entityManager.upsert).toHaveBeenCalledWith( - WorkflowPublishedVersion, - { workflowId: 'wf-1', publishedVersionId: 'v-1' }, - ['workflowId'], - ); - expect(activeWorkflowManager.clearWebhooks).not.toHaveBeenCalled(); - expect(activeWorkflowManager.removeNonWebhookTriggers).not.toHaveBeenCalled(); - expect(activeWorkflowManager.add).not.toHaveBeenCalled(); - }); - - test('marks failed when non-webhook trigger teardown throws', async () => { - const record = makeRecord(); - workflowRepository.findById.mockResolvedValue(makeWorkflow()); - activeWorkflowManager.removeNonWebhookTriggers.mockRejectedValue( - new Error('trigger cleanup failed'), - ); - - await consumer.processRecord(record); - - expect(outboxRepository.markFailed).toHaveBeenCalledWith(1, 'trigger cleanup failed'); - expect(entityManager.upsert).not.toHaveBeenCalled(); - expect(activeWorkflowManager.add).not.toHaveBeenCalled(); - }); - - test('rolls back and marks failed when registerNewTriggers throws', async () => { - const record = makeRecord(); - workflowRepository.findById.mockResolvedValue(makeWorkflow()); - activeWorkflowManager.add.mockRejectedValue(new Error('trigger registration failed')); - - await consumer.processRecord(record); - - expect(activeWorkflowManager.clearWebhooks).toHaveBeenCalledWith('wf-1'); - expect(activeWorkflowManager.removeNonWebhookTriggers).toHaveBeenCalledWith('wf-1'); - expect(activeWorkflowManager.add).toHaveBeenCalled(); - - // Rollback: deactivate workflow - expect(workflowRepository.update).toHaveBeenCalledWith('wf-1', { - active: false, - activeVersionId: null, + const callOrder: string[] = []; + activeWorkflowManager.removeTriggerNodes.mockImplementation(async () => { + callOrder.push('remove'); }); - expect(activationErrorsService.register).toHaveBeenCalledWith( - 'wf-1', - 'trigger registration failed', + (outboxRepository.manager.upsert as jest.Mock).mockImplementation(async () => { + callOrder.push('advance'); + return await Promise.resolve({} as never); + }); + activeWorkflowManager.addTriggerNodes.mockImplementation(async () => { + callOrder.push('add'); + }); + + await consumer.processRecord(makeRecord()); + + expect(activeWorkflowManager.removeTriggerNodes).toHaveBeenCalledWith( + expect.objectContaining({ id: 'wf-1' }), + oldVersion, + new Set(['a']), ); - expect(outboxRepository.markFailed).toHaveBeenCalledWith(1, 'trigger registration failed'); + expect(activeWorkflowManager.addTriggerNodes).toHaveBeenCalledWith( + expect.objectContaining({ id: 'wf-1' }), + newVersion, + new Set(['a']), + ); + expect(callOrder).toEqual(['remove', 'advance', 'add']); + expect(outboxRepository.markCompleted).toHaveBeenCalledWith(1); + }); + + test('propagates without advancing when removing triggers throws', async () => { + setTriggerSets([triggerNode('a'), triggerNode('b')], [triggerNode('a')]); + activeWorkflowManager.removeTriggerNodes.mockRejectedValue(new Error('teardown failed')); + + // Teardown failures bubble up to the poll loop (tryProcessRecord), which + // marks the record failed. + await expect(consumer.processRecord(makeRecord())).rejects.toThrow('teardown failed'); + + expect(outboxRepository.manager.upsert).not.toHaveBeenCalled(); + expect(activeWorkflowManager.addTriggerNodes).not.toHaveBeenCalled(); + }); + + test('throws unimplemented when adding triggers throws', async () => { + setTriggerSets([triggerNode('a')], [triggerNode('a'), triggerNode('b')]); + activeWorkflowManager.addTriggerNodes.mockRejectedValue(new Error('registration failed')); + + await expect(consumer.processRecord(makeRecord())).rejects.toThrow( + 'Workflow publication trigger activation failure handling is not implemented yet', + ); + + expect(outboxRepository.manager.upsert).toHaveBeenCalled(); + expect(workflowRepository.update).not.toHaveBeenCalled(); + expect(activationErrorsService.register).not.toHaveBeenCalled(); + expect(outboxRepository.markFailed).not.toHaveBeenCalled(); + }); + + test('treats a first publication (no published-version mapping yet) as all-added', async () => { + workflowPublishedVersionRepository.findOne.mockResolvedValue(null); + setTriggerSets([], [triggerNode('a')]); + + await consumer.processRecord(makeRecord()); + + expect(activeWorkflowManager.removeTriggerNodes).not.toHaveBeenCalled(); + expect(activeWorkflowManager.addTriggerNodes).toHaveBeenCalledWith( + expect.objectContaining({ id: 'wf-1' }), + newVersion, + new Set(['a']), + ); + expect(outboxRepository.markCompleted).toHaveBeenCalledWith(1); }); }); }); diff --git a/packages/cli/src/workflows/trigger-diff.ts b/packages/cli/src/workflows/trigger-diff.ts new file mode 100644 index 00000000000..a016ea54c65 --- /dev/null +++ b/packages/cli/src/workflows/trigger-diff.ts @@ -0,0 +1,48 @@ +import type { INode } from 'n8n-workflow'; +import { compareWorkflowsNodes, NodeDiffStatus } from 'n8n-workflow'; + +/** + * The trigger nodes that need to be deregistered (`toRemove`) and registered + * (`toAdd`) to move a workflow's active triggers from `oldTriggerNodes` to + * `newTriggerNodes`. A modified trigger appears in both lists (remove-then-add); + * an unchanged trigger appears in neither. + */ +export interface TriggerDiff { + toAdd: Set; + toRemove: Set; +} + +/** + * Computes the trigger-level diff between two versions of a workflow. Both + * inputs must already be filtered to the enabled trigger-like nodes of their + * version, so a disabled trigger is treated as absent: enabling it yields an + * add, disabling it yields a remove. + */ +export function computeTriggerDiff( + oldTriggerNodes: INode[], + newTriggerNodes: INode[], +): TriggerDiff { + const diff = compareWorkflowsNodes(oldTriggerNodes, newTriggerNodes); + + const toAdd: Set = new Set(); + const toRemove: Set = new Set(); + + for (const [nodeId, { status }] of diff) { + switch (status) { + case NodeDiffStatus.Added: + toAdd.add(nodeId); + break; + case NodeDiffStatus.Deleted: + toRemove.add(nodeId); + break; + case NodeDiffStatus.Modified: + toRemove.add(nodeId); + toAdd.add(nodeId); + break; + case NodeDiffStatus.Eq: + break; + } + } + + return { toAdd, toRemove }; +} diff --git a/packages/cli/src/workflows/workflow-publication-outbox-consumer.ts b/packages/cli/src/workflows/workflow-publication-outbox-consumer.ts index 4ef53d6a353..bbc3b6e98aa 100644 --- a/packages/cli/src/workflows/workflow-publication-outbox-consumer.ts +++ b/packages/cli/src/workflows/workflow-publication-outbox-consumer.ts @@ -1,27 +1,34 @@ import { Logger } from '@n8n/backend-common'; import { WorkflowsConfig } from '@n8n/config'; import { + WorkflowEntity, + WorkflowHistory, + WorkflowHistoryRepository, WorkflowPublicationOutbox, WorkflowPublicationOutboxRepository, WorkflowPublishedVersion, + WorkflowPublishedVersionRepository, WorkflowRepository, } from '@n8n/db'; import { OnLeaderStepdown, OnLeaderTakeover, OnShutdown } from '@n8n/decorators'; import { Service } from '@n8n/di'; import { ErrorReporter } from 'n8n-core'; -import { ensureError } from 'n8n-workflow'; +import { UnexpectedError, ensureError } from 'n8n-workflow'; import { ActivationErrorsService } from '@/activation-errors.service'; import { ActiveWorkflowManager } from '@/active-workflow-manager'; +import { computeTriggerDiff } from '@/workflows/trigger-diff'; /** * Consumes the workflow publication outbox on the leader instance. It polls for - * pending records and, for each one, reapplies the workflow's triggers to match - * the published version: tearing down the old triggers, advancing - * `activeVersionId`, and registering the new triggers. Records for deleted or - * no-longer-active workflows are short-circuited to completed, and activation - * failures are recorded against the workflow and marked failed without halting - * the poll loop. + * pending records and, for each one, reconciles the workflow's triggers to match + * the version being published. It computes a trigger-level diff between the + * currently published version and the new version, then applies only the + * necessary operations: removing deleted triggers, adding new ones, and + * re-applying modified ones (remove-then-add) while leaving unchanged triggers + * running. Records for deleted or no-longer-active workflows are short-circuited + * to completed, and activation failures are recorded against the workflow and + * marked failed without halting the poll loop. */ @Service() export class WorkflowPublicationOutboxConsumer { @@ -37,6 +44,8 @@ export class WorkflowPublicationOutboxConsumer { private readonly errorReporter: ErrorReporter, private readonly outboxRepository: WorkflowPublicationOutboxRepository, private readonly workflowRepository: WorkflowRepository, + private readonly workflowHistoryRepository: WorkflowHistoryRepository, + private readonly workflowPublishedVersionRepository: WorkflowPublishedVersionRepository, private readonly activeWorkflowManager: ActiveWorkflowManager, private readonly activationErrorsService: ActivationErrorsService, ) { @@ -125,14 +134,13 @@ export class WorkflowPublicationOutboxConsumer { async processRecord(record: WorkflowPublicationOutbox) { const { workflowId, publishedVersionId } = record; - - const workflow = await this.workflowRepository.findById(workflowId); + const { workflow, oldVersion, newVersion } = await this.resolveVersions(record); if (!workflow) { this.logger.warn('Workflow not found, marking outbox record as completed', { workflowId, outboxId: record.id, }); - await this.outboxRepository.markCompleted(record.id); + await this.finalizePublication(record); return; } @@ -141,26 +149,51 @@ export class WorkflowPublicationOutboxConsumer { workflowId, outboxId: record.id, }); - await this.outboxRepository.markCompleted(record.id); + await this.finalizePublication(record); return; } - if (workflow.activeVersionId === publishedVersionId) { + if (!newVersion) { + this.logger.warn('Published version not found, marking outbox record as completed', { + workflowId, + publishedVersionId, + outboxId: record.id, + }); + await this.outboxRepository.markFailed(record.id, 'Published version not found'); + return; + } + + const { toAdd, toRemove } = computeTriggerDiff( + this.activeWorkflowManager.getEnabledTriggerNodes(oldVersion), + this.activeWorkflowManager.getEnabledTriggerNodes(newVersion), + ); + + // No trigger changed: advance the published version and finish. Unchanged + // triggers keep running and re-read the new version on their next fire. + if (toAdd.size === 0 && toRemove.size === 0) { await this.advancePublishedVersion(record); await this.finalizePublication(record); return; } + // Must happen BEFORE advancing the version, using the currently + // published version so the right webhooks are deregistered. + if (toRemove.size > 0 && oldVersion) { + await this.activeWorkflowManager.removeTriggerNodes(workflow, oldVersion, toRemove); + } + + await this.advancePublishedVersion(record); + try { - // Must happen BEFORE advancing the version because clearWebhooks() - // reads activeVersion from DB. - await this.tearDownOldTriggers(record); + if (toAdd.size > 0) { + await this.activeWorkflowManager.addTriggerNodes(workflow, newVersion, toAdd); + } else if (toRemove.size > 0) { + await this.activeWorkflowManager.updateWorkflowTriggerCount(workflow, newVersion); + } + } catch (e) { + const error = ensureError(e); - await this.advancePublishedVersion(record); - - await this.registerNewTriggers(record); - } catch (error) { - await this.outboxRepository.markFailed(record.id, ensureError(error).message); + await this.handleTriggerActivationFailure(error, workflowId, record.id); return; } @@ -173,38 +206,31 @@ export class WorkflowPublicationOutboxConsumer { }); } - private async tearDownOldTriggers(record: WorkflowPublicationOutbox) { - // This try-catch reflects the old behaviour in the ActiveWorkflowManager. - // We probably want to revisit this. - try { - await this.activeWorkflowManager.clearWebhooks(record.workflowId); - } catch (error) { - this.errorReporter.error(error, { - shouldBeLogged: true, - tags: { - workflowId: record.workflowId, - outboxId: record.id, - }, - }); - } + /** + * Loads the workflow and the two versions whose triggers are diffed: the + * version being published (`newVersion`, null if its history row no longer + * exists) and the currently published version (`oldVersion`, null on a first + * publication). The workflow is loaded independently of the published-version + * mapping so a first publication (no mapping row yet) still resolves it. + */ + private async resolveVersions(record: WorkflowPublicationOutbox): Promise<{ + workflow: WorkflowEntity | null; + oldVersion: WorkflowHistory | null; + newVersion: WorkflowHistory | null; + }> { + const [workflow, currentlyPublishedVersion, newVersion] = await Promise.all([ + this.workflowRepository.findOneBy({ id: record.workflowId }), + this.workflowPublishedVersionRepository.findOne({ + where: { workflowId: record.workflowId }, + relations: { publishedVersion: true }, + loadEagerRelations: false, + }), + this.workflowHistoryRepository.findOneBy({ versionId: record.publishedVersionId }), + ]); - await this.activeWorkflowManager.removeActivationError(record.workflowId); - await this.activeWorkflowManager.removeNonWebhookTriggers(record.workflowId); - } + const oldVersion = currentlyPublishedVersion?.publishedVersion ?? null; - private async registerNewTriggers(record: WorkflowPublicationOutbox) { - try { - await this.activeWorkflowManager.add(record.workflowId, 'update', undefined, { - shouldPublish: false, - }); - } catch (error) { - await this.workflowRepository.update(record.workflowId, { - active: false, - activeVersionId: null, - }); - await this.activationErrorsService.register(record.workflowId, ensureError(error).message); - throw error; - } + return { workflow, oldVersion, newVersion }; } /** @@ -228,4 +254,18 @@ export class WorkflowPublicationOutboxConsumer { await this.outboxRepository.markCompleted(record.id); await this.activationErrorsService.deregister(record.workflowId); } + + private async handleTriggerActivationFailure( + error: Error, + workflowId: string, + recordId: number, + ): Promise { + throw new UnexpectedError( + 'Workflow publication trigger activation failure handling is not implemented yet', + { + cause: error, + extra: { workflowId, outboxId: recordId }, + }, + ); + } } diff --git a/packages/cli/test/integration/workflows/workflow-publication-outbox-consumer.test.ts b/packages/cli/test/integration/workflows/workflow-publication-outbox-consumer.test.ts new file mode 100644 index 00000000000..d22fe3d65f2 --- /dev/null +++ b/packages/cli/test/integration/workflows/workflow-publication-outbox-consumer.test.ts @@ -0,0 +1,164 @@ +import { + createWorkflowWithHistory, + mockInstance, + setActiveVersion, + testDb, +} from '@n8n/backend-test-utils'; +import { WorkflowPublicationOutboxRepository, WorkflowPublishedVersionRepository } from '@n8n/db'; +import { Container } from '@n8n/di'; +import { ActiveWorkflowTriggers, ExternalSecretsProxy, InstanceSettings } from 'n8n-core'; +import { ScheduleTrigger } from 'n8n-nodes-base/nodes/Schedule/ScheduleTrigger.node'; +import type { INode, INodeTypeData } from 'n8n-workflow'; +import { v4 as uuid } from 'uuid'; + +import { ActiveExecutions } from '@/active-executions'; +import { ActiveWorkflowManager } from '@/active-workflow-manager'; +import { ExecutionService } from '@/executions/execution.service'; +import { ExternalHooks } from '@/external-hooks'; +import { Push } from '@/push'; +import { OwnershipService } from '@/services/ownership.service'; +import { WorkflowPublicationOutboxConsumer } from '@/workflows/workflow-publication-outbox-consumer'; +import { WorkflowService } from '@/workflows/workflow.service'; + +import { createOwner } from '../shared/db/users'; +import { createWorkflowHistoryItem } from '../shared/db/workflow-history'; +import * as utils from '../shared/utils/'; + +// Peripheral services with side effects we don't exercise here; the webhook +// service is left real so non-webhook (schedule) triggers enumerate to zero +// webhooks correctly. +mockInstance(ActiveExecutions); +mockInstance(Push); +mockInstance(ExternalSecretsProxy); +mockInstance(ExecutionService); +mockInstance(WorkflowService); +mockInstance(OwnershipService); +mockInstance(ExternalHooks); + +let consumer: WorkflowPublicationOutboxConsumer; +let activeWorkflowManager: ActiveWorkflowManager; +let activeWorkflowTriggers: ActiveWorkflowTriggers; +let outboxRepository: WorkflowPublicationOutboxRepository; +let publishedVersionRepository: WorkflowPublishedVersionRepository; + +const scheduleNode = (suffix: string): INode => ({ + id: `node-${suffix}`, + name: `Schedule ${suffix}`, + type: 'n8n-nodes-base.scheduleTrigger', + typeVersion: 1, + position: [0, 0], + parameters: {}, +}); + +beforeAll(async () => { + await testDb.init(); + + const nodes: INodeTypeData = { + 'n8n-nodes-base.scheduleTrigger': { type: new ScheduleTrigger(), sourcePath: '' }, + }; + await utils.initNodeTypes(nodes); + + Container.get(InstanceSettings).markAsLeader(); + + consumer = Container.get(WorkflowPublicationOutboxConsumer); + activeWorkflowManager = Container.get(ActiveWorkflowManager); + activeWorkflowTriggers = Container.get(ActiveWorkflowTriggers); + outboxRepository = Container.get(WorkflowPublicationOutboxRepository); + publishedVersionRepository = Container.get(WorkflowPublishedVersionRepository); +}); + +afterEach(async () => { + await activeWorkflowManager.removeAll(); + // Delete WorkflowPublishedVersion first: it references WorkflowHistory with + // onDelete RESTRICT, and deleting WorkflowEntity cascades into WorkflowHistory. + await testDb.truncate([ + 'WorkflowPublishedVersion', + 'WorkflowPublicationOutbox', + 'WorkflowPublishHistory', + 'WorkflowEntity', + 'WorkflowHistory', + ]); +}); + +afterAll(async () => { + await testDb.terminate(); +}); + +describe('WorkflowPublicationOutboxConsumer (integration)', () => { + test('applies only the trigger diff, leaving the unchanged trigger registered', async () => { + const owner = await createOwner(); + + const unchanged = scheduleNode('unchanged'); + const removed = scheduleNode('removed'); + const added = scheduleNode('added'); + + // Currently active version runs `unchanged` + `removed`. + const workflow = await createWorkflowWithHistory( + { active: true, nodes: [unchanged, removed] }, + owner, + ); + await setActiveVersion(workflow.id, workflow.versionId); + await publishedVersionRepository.setPublishedVersion(workflow.id, workflow.versionId); + await activeWorkflowManager.add(workflow.id, 'activate'); + + expect(activeWorkflowTriggers.get(workflow.id)?.has(unchanged.id)).toBe(true); + expect(activeWorkflowTriggers.get(workflow.id)?.has(removed.id)).toBe(true); + + // New version drops `removed`, keeps `unchanged`, adds `added`. + const newVersionId = uuid(); + await createWorkflowHistoryItem(workflow.id, { + versionId: newVersionId, + nodes: [unchanged, added], + connections: {}, + }); + + await outboxRepository.enqueue(workflow.id, newVersionId); + const record = await outboxRepository.claimNextPendingRecord(); + expect(record).not.toBeNull(); + + await consumer.processRecord(record!); + + // Surgical in-memory result: unchanged kept, removed gone, added registered. + const state = activeWorkflowTriggers.get(workflow.id); + expect(state?.has(unchanged.id)).toBe(true); + expect(state?.has(removed.id)).toBe(false); + expect(state?.has(added.id)).toBe(true); + + // Canonical published version advanced and the record completed. + const published = await publishedVersionRepository.getPublishedVersionWithRelations( + workflow.id, + ); + expect(published?.publishedVersionId).toBe(newVersionId); + expect(await outboxRepository.claimNextPendingRecord()).toBeNull(); + }); + + test('does no trigger work when only non-trigger content changed', async () => { + const owner = await createOwner(); + + const trigger = scheduleNode('only'); + const workflow = await createWorkflowWithHistory({ active: true, nodes: [trigger] }, owner); + await setActiveVersion(workflow.id, workflow.versionId); + await publishedVersionRepository.setPublishedVersion(workflow.id, workflow.versionId); + await activeWorkflowManager.add(workflow.id, 'activate'); + + // New version keeps the same trigger (a non-trigger node could have changed). + const newVersionId = uuid(); + await createWorkflowHistoryItem(workflow.id, { + versionId: newVersionId, + nodes: [trigger], + connections: {}, + }); + + await outboxRepository.enqueue(workflow.id, newVersionId); + const record = await outboxRepository.claimNextPendingRecord(); + + await consumer.processRecord(record!); + + expect(activeWorkflowTriggers.get(workflow.id)?.has(trigger.id)).toBe(true); + const published = await publishedVersionRepository.getPublishedVersionWithRelations( + workflow.id, + ); + expect(published?.publishedVersionId).toBe(newVersionId); + expect(await outboxRepository.claimNextPendingRecord()).toBeNull(); + }); +}); diff --git a/packages/core/src/execution-engine/__tests__/active-workflow-triggers.test.ts b/packages/core/src/execution-engine/__tests__/active-workflow-triggers.test.ts index 842565fc284..8f98555716c 100644 --- a/packages/core/src/execution-engine/__tests__/active-workflow-triggers.test.ts +++ b/packages/core/src/execution-engine/__tests__/active-workflow-triggers.test.ts @@ -103,7 +103,7 @@ describe('ActiveWorkflowTriggers', () => { getPollFunctions.mockReturnValue(pollFunctions); } - return await activeWorkflowTriggers.add( + return await activeWorkflowTriggers.addAllTriggers( workflowId, workflow, additionalData, @@ -114,7 +114,7 @@ describe('ActiveWorkflowTriggers', () => { ); }; - describe('add()', () => { + describe('addAllTriggers()', () => { describe('should activate workflow', () => { it('with trigger function nodes', async () => { await addWorkflow({ triggerNodes: [triggerNode] }); @@ -225,7 +225,7 @@ describe('ActiveWorkflowTriggers', () => { .mockRejectedValueOnce(new Error('Trigger activation failed')); await expect( - activeWorkflowTriggers.add( + activeWorkflowTriggers.addAllTriggers( workflowId, workflow, additionalData, @@ -255,7 +255,7 @@ describe('ActiveWorkflowTriggers', () => { ); await expect( - activeWorkflowTriggers.add( + activeWorkflowTriggers.addAllTriggers( workflowId, workflow, additionalData, @@ -281,7 +281,7 @@ describe('ActiveWorkflowTriggers', () => { .mockRejectedValueOnce(new Error('Trigger activation failed')); await expect( - activeWorkflowTriggers.add( + activeWorkflowTriggers.addAllTriggers( workflowId, workflow, additionalData, @@ -692,6 +692,137 @@ describe('ActiveWorkflowTriggers', () => { }); }); + describe('addTriggers()', () => { + const triggerNodeA = mock({ id: 'a' }); + const triggerNodeB = mock({ id: 'b' }); + + const addTriggers = async (nodeIds: string[]) => + await activeWorkflowTriggers.addTriggers( + workflowId, + workflow, + nodeIds, + additionalData, + mode, + activation, + getTriggerFunctions, + getPollFunctions, + ); + + it('registers only the requested trigger nodes', async () => { + workflow.getTriggerNodes.mockReturnValue([triggerNodeA, triggerNodeB]); + workflow.getPollNodes.mockReturnValue([]); + triggersAndPollers.runTriggerFunction.mockResolvedValue(triggerResponse); + + await addTriggers(['a']); + + expect(triggersAndPollers.runTriggerFunction).toHaveBeenCalledTimes(1); + expect(triggersAndPollers.runTriggerFunction).toHaveBeenCalledWith( + workflow, + triggerNodeA, + getTriggerFunctions, + additionalData, + mode, + activation, + ); + expect(activeWorkflowTriggers.isActive(workflowId)).toBe(true); + }); + + it('merges newly added triggers into an already-active workflow', async () => { + workflow.getTriggerNodes.mockReturnValue([triggerNodeA, triggerNodeB]); + workflow.getPollNodes.mockReturnValue([]); + triggersAndPollers.runTriggerFunction.mockResolvedValue(triggerResponse); + + await addTriggers(['a']); + await addTriggers(['b']); + + expect(triggersAndPollers.runTriggerFunction).toHaveBeenCalledTimes(2); + + await activeWorkflowTriggers.removeTriggers(workflowId, new Set(['a', 'b'])); + expect(activeWorkflowTriggers.isActive(workflowId)).toBe(false); + }); + }); + + describe('removeTriggers()', () => { + const triggerNodeA = mock({ id: 'a' }); + const triggerNodeB = mock({ id: 'b' }); + const pollNodeP = mock({ id: 'p' }); + const responseA = mock(); + const responseB = mock(); + + const addTriggerNodesAB = async () => { + workflow.getTriggerNodes.mockReturnValue([triggerNodeA, triggerNodeB]); + workflow.getPollNodes.mockReturnValue([]); + triggersAndPollers.runTriggerFunction.mockImplementation(async (_workflow, node) => + node.id === 'a' ? responseA : responseB, + ); + + await activeWorkflowTriggers.addTriggers( + workflowId, + workflow, + ['a', 'b'], + additionalData, + mode, + activation, + getTriggerFunctions, + getPollFunctions, + ); + }; + + it('closes only the targeted trigger and deregisters its cron, leaving others active', async () => { + await addTriggerNodesAB(); + + await activeWorkflowTriggers.removeTriggers(workflowId, new Set(['a'])); + + expect(responseA.closeFunction).toHaveBeenCalled(); + expect(responseB.closeFunction).not.toHaveBeenCalled(); + expect(scheduledTaskManager.deregisterCron).toHaveBeenCalledWith(workflowId, 'a'); + expect(activeWorkflowTriggers.isActive(workflowId)).toBe(true); + }); + + it('drops the workflow once its last trigger is removed and no crons remain', async () => { + await addTriggerNodesAB(); + scheduledTaskManager.hasCrons.mockReturnValue(false); + + await activeWorkflowTriggers.removeTriggers(workflowId, new Set(['a', 'b'])); + + expect(activeWorkflowTriggers.isActive(workflowId)).toBe(false); + }); + + it('keeps the workflow active while it still has registered crons', async () => { + await addTriggerNodesAB(); + scheduledTaskManager.hasCrons.mockReturnValue(true); + + await activeWorkflowTriggers.removeTriggers(workflowId, new Set(['a', 'b'])); + + expect(activeWorkflowTriggers.isActive(workflowId)).toBe(true); + }); + + it('deregisters crons for a poll node with no trigger response', async () => { + workflow.getTriggerNodes.mockReturnValue([]); + workflow.getPollNodes.mockReturnValue([pollNodeP]); + getPollFunctions.mockReturnValue(pollFunctions); + pollFunctions.getNodeParameter + .calledWith('pollTimes') + .mockReturnValue({ item: [{ mode: 'everyMinute' }] }); + triggersAndPollers.runPollFunction.mockResolvedValue(null); + + await activeWorkflowTriggers.addTriggers( + workflowId, + workflow, + ['p'], + additionalData, + mode, + activation, + getTriggerFunctions, + getPollFunctions, + ); + + await activeWorkflowTriggers.removeTriggers(workflowId, new Set(['p'])); + + expect(scheduledTaskManager.deregisterCron).toHaveBeenCalledWith(workflowId, 'p'); + }); + }); + describe('ScheduledTaskManager cron cleanup', () => { const hourly = '0 * * * *' as CronExpression; let realLogger: ReturnType>; @@ -779,7 +910,7 @@ describe('ActiveWorkflowTriggers', () => { .mockRejectedValueOnce(new Error('Trigger activation failed')); await expect( - activeWorkflowTriggersReal.add( + activeWorkflowTriggersReal.addAllTriggers( workflowId, workflow, additionalData, @@ -796,7 +927,7 @@ describe('ActiveWorkflowTriggers', () => { it('should leave no registered cron when a later poll node fails activation', async () => { // First poll node registers its cron, second fails its test poll → the // registered cron must be torn down. The cron is keyed by workflow.id, so it - // must match the id passed to add(). + // must match the id passed to addAllTriggers(). workflow.id = workflowId; workflow.getTriggerNodes.mockReturnValue([]); workflow.getPollNodes.mockReturnValue([mock(), mock()]); @@ -809,7 +940,7 @@ describe('ActiveWorkflowTriggers', () => { .mockRejectedValueOnce(new Error('Failed to activate poll trigger')); await expect( - activeWorkflowTriggersReal.add( + activeWorkflowTriggersReal.addAllTriggers( workflowId, workflow, additionalData, @@ -841,7 +972,7 @@ describe('ActiveWorkflowTriggers', () => { return triggerResponse; }); - await activeWorkflowTriggersReal.add( + await activeWorkflowTriggersReal.addAllTriggers( workflowId, workflow, additionalData, diff --git a/packages/core/src/execution-engine/__tests__/scheduled-task-manager.test.ts b/packages/core/src/execution-engine/__tests__/scheduled-task-manager.test.ts index 4587d558015..f1faace2a93 100644 --- a/packages/core/src/execution-engine/__tests__/scheduled-task-manager.test.ts +++ b/packages/core/src/execution-engine/__tests__/scheduled-task-manager.test.ts @@ -160,6 +160,67 @@ describe('ScheduledTaskManager', () => { expect(onTick).not.toHaveBeenCalled(); }); + it('should deregister CronJobs for a single node, leaving other nodes intact', () => { + const nodeA = 'node-a'; + const nodeB = 'node-b'; + + scheduledTaskManager.registerCron( + { + workflowId: workflow.id, + nodeId: nodeA, + timezone: workflow.timezone, + expression: everyMinute, + }, + onTick, + ); + scheduledTaskManager.registerCron( + { + workflowId: workflow.id, + nodeId: nodeB, + timezone: workflow.timezone, + expression: everyMinute, + }, + onTick, + ); + + expect(scheduledTaskManager.cronsByWorkflow.get(workflow.id)?.size).toBe(2); + + scheduledTaskManager.deregisterCron(workflow.id, nodeA); + + const remaining = scheduledTaskManager.cronsByWorkflow.get(workflow.id); + expect(remaining?.size).toBe(1); + expect([...(remaining?.values() ?? [])][0].ctx.nodeId).toBe(nodeB); + }); + + it('should drop the workflow entry once its last node cron is deregistered', () => { + const nodeId = 'only-node'; + scheduledTaskManager.registerCron( + { workflowId: workflow.id, nodeId, timezone: workflow.timezone, expression: everyMinute }, + onTick, + ); + + scheduledTaskManager.deregisterCron(workflow.id, nodeId); + + expect(scheduledTaskManager.cronsByWorkflow.get(workflow.id)).toBeUndefined(); + expect(scheduledTaskManager.hasCrons(workflow.id)).toBe(false); + }); + + it('hasCrons reflects whether a workflow has registered crons', () => { + expect(scheduledTaskManager.hasCrons(workflow.id)).toBe(false); + + scheduledTaskManager.registerCron( + { + workflowId: workflow.id, + nodeId: 'n', + timezone: workflow.timezone, + expression: everyMinute, + }, + onTick, + ); + + expect(scheduledTaskManager.hasCrons(workflow.id)).toBe(true); + }); + it('should not set up log interval when activeInterval is 0', () => { const configWithZeroInterval = mock({ activeInterval: 0 }); const manager = new ScheduledTaskManager( diff --git a/packages/core/src/execution-engine/active-workflow-triggers.ts b/packages/core/src/execution-engine/active-workflow-triggers.ts index c908431278e..5a0edf3d651 100644 --- a/packages/core/src/execution-engine/active-workflow-triggers.ts +++ b/packages/core/src/execution-engine/active-workflow-triggers.ts @@ -67,13 +67,13 @@ export class ActiveWorkflowTriggers { } /** - * Makes a workflow active + * Makes a workflow active by registering all of its trigger and poll nodes. * * @param {string} workflowId The id of the workflow to activate * @param {Workflow} workflow The workflow to activate * @param {IWorkflowExecuteAdditionalData} additionalData The additional data which is needed to run workflows */ - async add( + async addAllTriggers( workflowId: string, workflow: Workflow, additionalData: IWorkflowExecuteAdditionalData, @@ -85,9 +85,46 @@ export class ActiveWorkflowTriggers { // Tear down any registration still lingering for this workflow before readding it. await this.remove(workflowId); - const triggerFunctionNodes = workflow.getTriggerNodes(); + const nodeIds = [...workflow.getTriggerNodes(), ...workflow.getPollNodes()].map( + (node) => node.id, + ); - const triggers = new WorkflowActiveTriggersState(); + await this.addTriggers( + workflowId, + workflow, + nodeIds, + additionalData, + mode, + activation, + getTriggerFunctions, + getPollFunctions, + ); + } + + /** + * Activates the given subset of a workflow's trigger and poll nodes, merging + * them into any triggers already active for the workflow. Used to apply a + * trigger-level diff during publication without disturbing unchanged triggers. + */ + async addTriggers( + workflowId: string, + workflow: Workflow, + nodeIds: string[], + additionalData: IWorkflowExecuteAdditionalData, + mode: WorkflowExecuteMode, + activation: WorkflowActivateMode, + getTriggerFunctions: IGetExecuteTriggerFunctions, + getPollFunctions: IGetExecutePollFunctions, + ) { + const nodeIdSet = new Set(nodeIds); + const existing = this.activeTriggersByWorkflowId.get(workflowId); + const triggers = existing ?? new WorkflowActiveTriggersState(); + const triggersAddedDuringThisCall = new WorkflowActiveTriggersState(); + const triggerNodeIdsAddedDuringThisCall: string[] = []; + + const triggerFunctionNodes = workflow + .getTriggerNodes() + .filter((node) => nodeIdSet.has(node.id)); for (const triggerNode of triggerFunctionNodes) { try { @@ -101,13 +138,22 @@ export class ActiveWorkflowTriggers { ); if (triggerResponse !== undefined) { triggers.add(triggerNode.id, triggerResponse); + triggersAddedDuringThisCall.add(triggerNode.id, triggerResponse); + triggerNodeIdsAddedDuringThisCall.push(triggerNode.id); } } catch (e) { const error = ensureError(e); // Tear down anything an earlier node already registered, so a failed // activation doesn't leave triggers or crons running. - await this.rollbackPartialActivation(workflowId, triggers); + await this.rollbackPartialActivation( + workflowId, + triggersAddedDuringThisCall, + existing ? nodeIdSet : undefined, + ); + for (const nodeId of triggerNodeIdsAddedDuringThisCall) { + triggers.delete(nodeId); + } throw new WorkflowActivationError( `There was a problem activating the workflow: "${error.message}"`, @@ -118,7 +164,7 @@ export class ActiveWorkflowTriggers { this.activeTriggersByWorkflowId.set(workflowId, triggers); - const pollTriggerNodes = workflow.getPollNodes(); + const pollTriggerNodes = workflow.getPollNodes().filter((node) => nodeIdSet.has(node.id)); if (pollTriggerNodes.length === 0) return; @@ -134,10 +180,17 @@ export class ActiveWorkflowTriggers { activation, ); } catch (e) { - // A failed activation must not leave the workflow half-active. Drop it - // from memory and tear down every trigger and cron registered so far. - this.activeTriggersByWorkflowId.delete(workflowId); - await this.rollbackPartialActivation(workflowId, triggers); + if (!existing) { + this.activeTriggersByWorkflowId.delete(workflowId); + } + await this.rollbackPartialActivation( + workflowId, + triggersAddedDuringThisCall, + existing ? nodeIdSet : undefined, + ); + for (const nodeId of triggerNodeIdsAddedDuringThisCall) { + triggers.delete(nodeId); + } const error = ensureError(e); @@ -149,6 +202,36 @@ export class ActiveWorkflowTriggers { } } + /** + * Deactivates the given subset of a workflow's trigger and poll nodes, + * leaving the rest active. Closes each node's trigger response and + * deregisters its poll crons. Drops the workflow from the active set only + * when no triggers or crons remain for it. + */ + async removeTriggers(workflowId: string, nodeIds: Set) { + const activeTriggers = this.activeTriggersByWorkflowId.get(workflowId); + if (!activeTriggers) { + for (const nodeId of nodeIds) { + this.scheduledTaskManager.deregisterCron(workflowId, nodeId); + } + return; + } + + for (const nodeId of nodeIds) { + this.scheduledTaskManager.deregisterCron(workflowId, nodeId); + + const response = activeTriggers.get(nodeId); + if (response) { + await this.closeTrigger(response, workflowId); + } + activeTriggers.delete(nodeId); + } + + if (activeTriggers.isEmpty && !this.scheduledTaskManager.hasCrons(workflowId)) { + this.activeTriggersByWorkflowId.delete(workflowId); + } + } + /** * Tears down everything an in-progress activation registered before it * failed — the trigger responses' close functions and any crons — so a @@ -159,10 +242,17 @@ export class ActiveWorkflowTriggers { private async rollbackPartialActivation( workflowId: string, triggers: WorkflowActiveTriggersState, + nodeIds?: Iterable, ) { // Stop the crons first: deregistration is synchronous and is what actually // prevents the failed activation from continuing to fire. - this.scheduledTaskManager.deregisterCrons(workflowId); + if (nodeIds) { + for (const nodeId of nodeIds) { + this.scheduledTaskManager.deregisterCron(workflowId, nodeId); + } + } else { + this.scheduledTaskManager.deregisterCrons(workflowId); + } for (const response of triggers.triggerResponses) { try { diff --git a/packages/core/src/execution-engine/scheduled-task-manager.ts b/packages/core/src/execution-engine/scheduled-task-manager.ts index cbe4d00be33..6fdeadf4d65 100644 --- a/packages/core/src/execution-engine/scheduled-task-manager.ts +++ b/packages/core/src/execution-engine/scheduled-task-manager.ts @@ -152,6 +152,39 @@ export class ScheduledTaskManager { return Array.from(this.cronsByWorkflow.keys()); } + /** Deregister the crons registered for a single node of a workflow. */ + deregisterCron(workflowId: string, nodeId: string) { + const workflowCrons = this.cronsByWorkflow.get(workflowId); + + if (!workflowCrons || workflowCrons.size === 0) return; + + const summaries: string[] = []; + + for (const [key, cron] of workflowCrons) { + if (cron.ctx.nodeId !== nodeId) continue; + summaries.push(cron.summary); + void cron.job.stop(); + workflowCrons.delete(key); + } + + if (workflowCrons.size === 0) this.cronsByWorkflow.delete(workflowId); + + if (summaries.length === 0) return; + + this.logger.info('Deregistered crons for node', { + workflowId, + nodeId, + crons: summaries, + instanceRole: this.instanceSettings.instanceRole, + }); + } + + /** Whether any crons are currently registered for the workflow. */ + hasCrons(workflowId: string) { + const workflowCrons = this.cronsByWorkflow.get(workflowId); + return workflowCrons !== undefined && workflowCrons.size > 0; + } + deregisterAllCrons() { for (const workflowId of this.cronsByWorkflow.keys()) { this.deregisterCrons(workflowId); diff --git a/packages/core/src/execution-engine/workflow-active-triggers-state.ts b/packages/core/src/execution-engine/workflow-active-triggers-state.ts index 6dc3f1e9f8a..278b8b7d7de 100644 --- a/packages/core/src/execution-engine/workflow-active-triggers-state.ts +++ b/packages/core/src/execution-engine/workflow-active-triggers-state.ts @@ -12,6 +12,21 @@ export class WorkflowActiveTriggersState { this.triggersByNodeId.set(nodeId, response); } + /** The trigger response recorded for a node, if any. */ + get(nodeId: string) { + return this.triggersByNodeId.get(nodeId); + } + + /** Whether a trigger response has been recorded for the given node. */ + has(nodeId: string) { + return this.triggersByNodeId.has(nodeId); + } + + /** Drops the trigger response recorded for a node. */ + delete(nodeId: string) { + this.triggersByNodeId.delete(nodeId); + } + /** Whether no trigger responses have been recorded yet. */ get isEmpty() { return this.triggersByNodeId.size === 0; diff --git a/packages/workflow/src/workflow-diff.ts b/packages/workflow/src/workflow-diff.ts index d5293e1cb24..0e027af82ba 100644 --- a/packages/workflow/src/workflow-diff.ts +++ b/packages/workflow/src/workflow-diff.ts @@ -35,7 +35,7 @@ export type NodeDiff = { node: T; }; -export type WorkflowDiff = Map>; +export type WorkflowDiff = Map>; export function compareNodes( base: T | undefined,