feat(core): Add trigger diffing to workflow publication (no-changelog) (#31971)

This commit is contained in:
Tomi Turtiainen
2026-06-10 14:57:59 +00:00
committed by GitHub
parent 2e8e81ed7f
commit ef2f21fed3
14 changed files with 1181 additions and 194 deletions
@@ -67,6 +67,49 @@ describe('ActiveWorkflowManager', () => {
);
});
describe('getEnabledTriggerNodes', () => {
function node(id: string, type: string, overrides: Partial<INode> = {}): 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<typeof ActiveWorkflowTriggers>[4],
);
await realActiveWorkflowTriggers.add(
await realActiveWorkflowTriggers.addAllTriggers(
'wf-1',
workflow,
additionalData,
+201 -7
View File
@@ -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<string>,
) {
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<string>,
) {
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<INode['id']>,
) {
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<INode['id']>,
) {
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<IWorkflowBase>;
nodeIds?: Set<string>;
},
) {
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,
@@ -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));
@@ -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> = {}): 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() });
});
});
@@ -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<ErrorReporter>();
const outboxRepository = mock<WorkflowPublicationOutboxRepository>();
const workflowRepository = mock<WorkflowRepository>();
const workflowHistoryRepository = mock<WorkflowHistoryRepository>();
const workflowPublishedVersionRepository = mock<WorkflowPublishedVersionRepository>();
const activeWorkflowManager = mock<ActiveWorkflowManager>();
const activationErrorsService = mock<ActivationErrorsService>();
const entityManager = mock<EntityManager>();
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> = {}): 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<EntityManager>({ 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);
});
});
});
@@ -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<INode['id']>;
toRemove: Set<INode['id']>;
}
/**
* 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<INode['id']> = new Set();
const toRemove: Set<INode['id']> = 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 };
}
@@ -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<never> {
throw new UnexpectedError(
'Workflow publication trigger activation failure handling is not implemented yet',
{
cause: error,
extra: { workflowId, outboxId: recordId },
},
);
}
}
@@ -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();
});
});
@@ -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<INode>({ id: 'a' });
const triggerNodeB = mock<INode>({ 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<INode>({ id: 'a' });
const triggerNodeB = mock<INode>({ id: 'b' });
const pollNodeP = mock<INode>({ id: 'p' });
const responseA = mock<ITriggerResponse>();
const responseB = mock<ITriggerResponse>();
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<typeof mock<Logger>>;
@@ -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<INode>(), mock<INode>()]);
@@ -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,
@@ -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(
@@ -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<INode['id']>) {
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<string>,
) {
// 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 {
@@ -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);
@@ -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;
+1 -1
View File
@@ -35,7 +35,7 @@ export type NodeDiff<T> = {
node: T;
};
export type WorkflowDiff<T> = Map<string, NodeDiff<T>>;
export type WorkflowDiff<T> = Map<INode['id'], NodeDiff<T>>;
export function compareNodes<T extends DiffableNode>(
base: T | undefined,