mirror of
https://github.com/n8n-io/n8n.git
synced 2026-09-24 23:22:38 +08:00
feat: Include active version id in workflowActivated broadcast event (no-changelog) (#22805)
Co-authored-by: Svetoslav Dekov <svetoslav.dekov@n8n.io>
This commit is contained in:
@@ -2,6 +2,7 @@ export type WorkflowActivated = {
|
||||
type: 'workflowActivated';
|
||||
data: {
|
||||
workflowId: string;
|
||||
activeVersionId: string;
|
||||
};
|
||||
};
|
||||
|
||||
|
||||
@@ -565,10 +565,24 @@ export class ActiveWorkflowManager {
|
||||
) {
|
||||
const added = { webhooks: false, triggersAndPollers: false };
|
||||
|
||||
const dbWorkflow = existingWorkflow ?? (await this.workflowRepository.findById(workflowId));
|
||||
|
||||
if (!dbWorkflow) {
|
||||
throw new WorkflowActivationError(`Failed to find workflow with ID "${workflowId}"`, {
|
||||
level: 'warning',
|
||||
});
|
||||
}
|
||||
|
||||
if (this.instanceSettings.isMultiMain && shouldPublish) {
|
||||
if (!dbWorkflow?.activeVersionId) {
|
||||
throw new UnexpectedError('Active version ID not found for workflow', {
|
||||
extra: { workflowId },
|
||||
});
|
||||
}
|
||||
|
||||
void this.publisher.publishCommand({
|
||||
command: 'add-webhooks-triggers-and-pollers',
|
||||
payload: { workflowId },
|
||||
payload: { workflowId, activeVersionId: dbWorkflow.activeVersionId },
|
||||
});
|
||||
|
||||
return added;
|
||||
@@ -580,14 +594,6 @@ export class ActiveWorkflowManager {
|
||||
const shouldAddTriggersAndPollers = this.shouldAddTriggersAndPollers();
|
||||
|
||||
try {
|
||||
const dbWorkflow = existingWorkflow ?? (await this.workflowRepository.findById(workflowId));
|
||||
|
||||
if (!dbWorkflow) {
|
||||
throw new WorkflowActivationError(`Failed to find workflow with ID "${workflowId}"`, {
|
||||
level: 'warning',
|
||||
});
|
||||
}
|
||||
|
||||
if (['init', 'leadershipChange'].includes(activationMode) && !dbWorkflow.activeVersion) {
|
||||
this.logger.debug(
|
||||
`Skipping workflow ${formatWorkflow(dbWorkflow)} as it is no longer active`,
|
||||
@@ -671,8 +677,11 @@ export class ActiveWorkflowManager {
|
||||
}
|
||||
|
||||
@OnPubSubEvent('display-workflow-activation', { instanceType: 'main' })
|
||||
handleDisplayWorkflowActivation({ workflowId }: PubSubCommandMap['display-workflow-activation']) {
|
||||
this.push.broadcast({ type: 'workflowActivated', data: { workflowId } });
|
||||
handleDisplayWorkflowActivation({
|
||||
workflowId,
|
||||
activeVersionId,
|
||||
}: PubSubCommandMap['display-workflow-activation']) {
|
||||
this.push.broadcast({ type: 'workflowActivated', data: { workflowId, activeVersionId } });
|
||||
}
|
||||
|
||||
@OnPubSubEvent('display-workflow-deactivation', { instanceType: 'main' })
|
||||
@@ -700,17 +709,18 @@ export class ActiveWorkflowManager {
|
||||
})
|
||||
async handleAddWebhooksTriggersAndPollers({
|
||||
workflowId,
|
||||
activeVersionId,
|
||||
}: PubSubCommandMap['add-webhooks-triggers-and-pollers']) {
|
||||
try {
|
||||
await this.add(workflowId, 'activate', undefined, {
|
||||
shouldPublish: false, // prevent leader from re-publishing message
|
||||
});
|
||||
|
||||
this.push.broadcast({ type: 'workflowActivated', data: { workflowId } });
|
||||
this.push.broadcast({ type: 'workflowActivated', data: { workflowId, activeVersionId } });
|
||||
|
||||
await this.publisher.publishCommand({
|
||||
command: 'display-workflow-activation',
|
||||
payload: { workflowId },
|
||||
payload: { workflowId, activeVersionId },
|
||||
}); // instruct followers to show activation in UI
|
||||
} catch (e) {
|
||||
const error = ensureError(e);
|
||||
|
||||
@@ -12,6 +12,7 @@ describe('PubSubRegistry', () => {
|
||||
let pubsubEventBus: PubSubEventBus;
|
||||
let logger: ReturnType<typeof mockLogger>;
|
||||
const workflowId = 'test-workflow-id';
|
||||
const activeVersionId = 'test-version-id';
|
||||
|
||||
const createTestServiceClass = () => {
|
||||
@Service()
|
||||
@@ -130,9 +131,9 @@ describe('PubSubRegistry', () => {
|
||||
);
|
||||
pubSubRegistry.init();
|
||||
|
||||
pubsubEventBus.emit('add-webhooks-triggers-and-pollers', { workflowId });
|
||||
pubsubEventBus.emit('add-webhooks-triggers-and-pollers', { workflowId, activeVersionId });
|
||||
expect(onLeaderInstanceSpy).toHaveBeenCalledTimes(1);
|
||||
expect(onLeaderInstanceSpy).toHaveBeenCalledWith({ workflowId });
|
||||
expect(onLeaderInstanceSpy).toHaveBeenCalledWith({ workflowId, activeVersionId });
|
||||
|
||||
pubsubEventBus.emit('restart-event-bus');
|
||||
expect(onFollowerInstanceSpy).not.toHaveBeenCalled();
|
||||
@@ -152,7 +153,7 @@ describe('PubSubRegistry', () => {
|
||||
);
|
||||
followerPubSubRegistry.init();
|
||||
|
||||
pubsubEventBus.emit('add-webhooks-triggers-and-pollers', { workflowId });
|
||||
pubsubEventBus.emit('add-webhooks-triggers-and-pollers', { workflowId, activeVersionId });
|
||||
expect(onLeaderInstanceSpy).not.toHaveBeenCalled();
|
||||
|
||||
pubsubEventBus.emit('restart-event-bus');
|
||||
@@ -176,9 +177,9 @@ describe('PubSubRegistry', () => {
|
||||
);
|
||||
pubSubRegistry.init();
|
||||
|
||||
pubsubEventBus.emit('add-webhooks-triggers-and-pollers', { workflowId });
|
||||
pubsubEventBus.emit('add-webhooks-triggers-and-pollers', { workflowId, activeVersionId });
|
||||
expect(onLeaderInstanceSpy).toHaveBeenCalledTimes(1);
|
||||
expect(onLeaderInstanceSpy).toHaveBeenCalledWith({ workflowId });
|
||||
expect(onLeaderInstanceSpy).toHaveBeenCalledWith({ workflowId, activeVersionId });
|
||||
});
|
||||
|
||||
it('should handle dynamic role changes at runtime', () => {
|
||||
@@ -196,19 +197,19 @@ describe('PubSubRegistry', () => {
|
||||
pubSubRegistry.init();
|
||||
|
||||
// Initially as follower, event should be ignored
|
||||
pubsubEventBus.emit('add-webhooks-triggers-and-pollers', { workflowId });
|
||||
pubsubEventBus.emit('add-webhooks-triggers-and-pollers', { workflowId, activeVersionId });
|
||||
expect(onLeaderInstanceSpy).not.toHaveBeenCalled();
|
||||
|
||||
// Change role to leader
|
||||
instanceSettings.instanceRole = 'leader';
|
||||
pubsubEventBus.emit('add-webhooks-triggers-and-pollers', { workflowId });
|
||||
pubsubEventBus.emit('add-webhooks-triggers-and-pollers', { workflowId, activeVersionId });
|
||||
expect(onLeaderInstanceSpy).toHaveBeenCalledTimes(1);
|
||||
expect(onLeaderInstanceSpy).toHaveBeenCalledWith({ workflowId });
|
||||
expect(onLeaderInstanceSpy).toHaveBeenCalledWith({ workflowId, activeVersionId });
|
||||
|
||||
// Change back to follower
|
||||
onLeaderInstanceSpy.mockClear();
|
||||
instanceSettings.instanceRole = 'follower';
|
||||
pubsubEventBus.emit('add-webhooks-triggers-and-pollers', { workflowId });
|
||||
pubsubEventBus.emit('add-webhooks-triggers-and-pollers', { workflowId, activeVersionId });
|
||||
expect(onLeaderInstanceSpy).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
|
||||
@@ -56,6 +56,7 @@ export type PubSubCommandMap = {
|
||||
|
||||
'add-webhooks-triggers-and-pollers': {
|
||||
workflowId: string;
|
||||
activeVersionId: string;
|
||||
};
|
||||
|
||||
'remove-triggers-and-pollers': {
|
||||
@@ -64,6 +65,7 @@ export type PubSubCommandMap = {
|
||||
|
||||
'display-workflow-activation': {
|
||||
workflowId: string;
|
||||
activeVersionId: string;
|
||||
};
|
||||
|
||||
'display-workflow-deactivation': {
|
||||
|
||||
@@ -378,7 +378,7 @@ export class WorkflowService {
|
||||
await this.workflowTagMappingRepository.overwriteTaggings(workflowId, tagIds);
|
||||
}
|
||||
|
||||
const relations = tagsDisabled ? [] : ['tags'];
|
||||
const relations = tagsDisabled ? ['activeVersion'] : ['tags', 'activeVersion'];
|
||||
|
||||
// We sadly get nothing back from "update". Neither if it updated a record
|
||||
// nor the new value. So query now the hopefully updated entry.
|
||||
|
||||
@@ -1442,6 +1442,38 @@ describe('PATCH /workflows/:workflowId', () => {
|
||||
expect(historyVersion!.nodes).toEqual(payload.nodes);
|
||||
});
|
||||
});
|
||||
|
||||
test('should include activeVersion relation in response for active workflows', async () => {
|
||||
const teamProject = await createTeamProject();
|
||||
const workflow = await createActiveWorkflow({}, teamProject);
|
||||
|
||||
const response = await authOwnerAgent.patch(`/workflows/${workflow.id}`).send({
|
||||
name: 'Updated Name',
|
||||
versionId: workflow.versionId,
|
||||
});
|
||||
|
||||
expect(response.statusCode).toBe(200);
|
||||
expect(response.body.data).toHaveProperty('activeVersion');
|
||||
expect(response.body.data.activeVersion).toBeDefined();
|
||||
expect(response.body.data.activeVersion).toMatchObject({
|
||||
versionId: expect.any(String),
|
||||
workflowId: workflow.id,
|
||||
});
|
||||
});
|
||||
|
||||
test('should include activeVersion as null for inactive workflows', async () => {
|
||||
const teamProject = await createTeamProject();
|
||||
const workflow = await createWorkflow({}, teamProject);
|
||||
|
||||
const response = await authOwnerAgent.patch(`/workflows/${workflow.id}`).send({
|
||||
name: 'Updated Name',
|
||||
versionId: workflow.versionId,
|
||||
});
|
||||
|
||||
expect(response.statusCode).toBe(200);
|
||||
expect(response.body.data).toHaveProperty('activeVersion');
|
||||
expect(response.body.data.activeVersion).toBeNull();
|
||||
});
|
||||
});
|
||||
|
||||
describe('PUT /:workflowId/transfer', () => {
|
||||
|
||||
+24
-2
@@ -1,15 +1,37 @@
|
||||
import type { WorkflowActivated } from '@n8n/api-types/push/workflow';
|
||||
import { useWorkflowsStore } from '@/app/stores/workflows.store';
|
||||
import { useBannersStore } from '@/features/shared/banners/banners.store';
|
||||
import { useUIStore } from '@/app/stores/ui.store';
|
||||
import { getWorkflowVersion } from '@n8n/rest-api-client';
|
||||
import { useRootStore } from '@n8n/stores/useRootStore';
|
||||
|
||||
export async function workflowActivated({ data }: WorkflowActivated) {
|
||||
const workflowsStore = useWorkflowsStore();
|
||||
const bannersStore = useBannersStore();
|
||||
const rootStore = useRootStore();
|
||||
const uiStore = useUIStore();
|
||||
|
||||
workflowsStore.setWorkflowActive(data.workflowId);
|
||||
const { workflowId, activeVersionId } = data;
|
||||
|
||||
const workflowIsBeingViewed = workflowsStore.workflowId === workflowId;
|
||||
const activeVersionIsSet = workflowsStore.workflow.activeVersionId !== activeVersionId;
|
||||
if (workflowIsBeingViewed && activeVersionIsSet) {
|
||||
const activeVersion = await getWorkflowVersion(
|
||||
rootStore.restApiContext,
|
||||
workflowId,
|
||||
activeVersionId,
|
||||
);
|
||||
|
||||
workflowsStore.setWorkflowActive(workflowId, activeVersion, false);
|
||||
|
||||
// Only update checksum if there are no unsaved changes
|
||||
if (!uiStore.stateIsDirty) {
|
||||
await workflowsStore.updateWorkflowChecksum();
|
||||
}
|
||||
}
|
||||
|
||||
// Remove auto-deactivated banner if viewing this workflow
|
||||
if (workflowsStore.workflowId === data.workflowId) {
|
||||
if (workflowIsBeingViewed) {
|
||||
bannersStore.removeBannerFromStack('WORKFLOW_AUTO_DEACTIVATED');
|
||||
}
|
||||
}
|
||||
|
||||
+7
@@ -1,14 +1,21 @@
|
||||
import type { WorkflowAutoDeactivated } from '@n8n/api-types/push/workflow';
|
||||
import { useWorkflowsStore } from '@/app/stores/workflows.store';
|
||||
import { useBannersStore } from '@/features/shared/banners/banners.store';
|
||||
import { useUIStore } from '@/app/stores/ui.store';
|
||||
|
||||
export async function workflowAutoDeactivated({ data }: WorkflowAutoDeactivated) {
|
||||
const workflowsStore = useWorkflowsStore();
|
||||
const bannersStore = useBannersStore();
|
||||
const uiStore = useUIStore();
|
||||
|
||||
workflowsStore.setWorkflowInactive(data.workflowId);
|
||||
|
||||
if (workflowsStore.workflowId === data.workflowId) {
|
||||
// Only update checksum if there are no unsaved changes
|
||||
if (!uiStore.stateIsDirty) {
|
||||
await workflowsStore.updateWorkflowChecksum();
|
||||
}
|
||||
|
||||
bannersStore.pushBannerToStack('WORKFLOW_AUTO_DEACTIVATED');
|
||||
}
|
||||
}
|
||||
|
||||
+9
@@ -1,8 +1,17 @@
|
||||
import type { WorkflowDeactivated } from '@n8n/api-types/push/workflow';
|
||||
import { useWorkflowsStore } from '@/app/stores/workflows.store';
|
||||
import { useUIStore } from '@/app/stores/ui.store';
|
||||
|
||||
export async function workflowDeactivated({ data }: WorkflowDeactivated) {
|
||||
const workflowsStore = useWorkflowsStore();
|
||||
const uiStore = useUIStore();
|
||||
|
||||
workflowsStore.setWorkflowInactive(data.workflowId);
|
||||
|
||||
if (workflowsStore.workflowId === data.workflowId) {
|
||||
// Only update checksum if there are no unsaved changes
|
||||
if (!uiStore.stateIsDirty) {
|
||||
await workflowsStore.updateWorkflowChecksum();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -110,8 +110,8 @@ export function useWorkflowActivate() {
|
||||
}
|
||||
|
||||
// Update local state
|
||||
if (workflow.activeVersionId !== null) {
|
||||
workflowsStore.setWorkflowActive(currWorkflowId);
|
||||
if (workflow.activeVersion) {
|
||||
workflowsStore.setWorkflowActive(currWorkflowId, workflow.activeVersion);
|
||||
} else {
|
||||
workflowsStore.setWorkflowInactive(currWorkflowId);
|
||||
}
|
||||
@@ -214,6 +214,7 @@ export function useWorkflowActivate() {
|
||||
interpolate: { newStateName: 'published' },
|
||||
}) + ':',
|
||||
);
|
||||
workflowsStore.setWorkflowInactive(workflowId);
|
||||
return false;
|
||||
} finally {
|
||||
updatingWorkflowActivation.value = false;
|
||||
|
||||
@@ -856,8 +856,8 @@ export function useWorkflowHelpers() {
|
||||
uiStore.stateIsDirty = false;
|
||||
}
|
||||
|
||||
if (workflow.activeVersionId !== null) {
|
||||
workflowsStore.setWorkflowActive(workflowId);
|
||||
if (workflow.activeVersion) {
|
||||
workflowsStore.setWorkflowActive(workflowId, workflow.activeVersion, isCurrentWorkflow);
|
||||
} else {
|
||||
workflowsStore.setWorkflowInactive(workflowId);
|
||||
}
|
||||
|
||||
@@ -45,6 +45,7 @@ import {
|
||||
import { waitFor } from '@testing-library/vue';
|
||||
import { useWorkflowState } from '@/app/composables/useWorkflowState';
|
||||
import { useSourceControlStore } from '@/features/integrations/sourceControl.ee/sourceControl.store';
|
||||
import type { WorkflowHistory } from '@n8n/rest-api-client';
|
||||
|
||||
vi.mock('@/features/ndv/shared/ndv.store', () => ({
|
||||
useNDVStore: vi.fn(() => ({
|
||||
@@ -892,12 +893,22 @@ describe('useWorkflowsStore', () => {
|
||||
});
|
||||
|
||||
describe('setWorkflowActive()', () => {
|
||||
it('should set workflow as active when it is not already active', () => {
|
||||
it('should set workflow as active when it is not already active', async () => {
|
||||
uiStore.stateIsDirty = true;
|
||||
workflowsStore.workflowsById = { '1': { active: false } as IWorkflowDb };
|
||||
workflowsStore.workflow.id = '1';
|
||||
|
||||
workflowsStore.setWorkflowActive('1');
|
||||
const mockActiveVersion: WorkflowHistory = {
|
||||
versionId: 'test-version-id',
|
||||
name: 'Test Version',
|
||||
authors: 'Test Author',
|
||||
description: 'A test workflow version',
|
||||
createdAt: new Date().toISOString(),
|
||||
updatedAt: new Date().toISOString(),
|
||||
workflowPublishHistory: [],
|
||||
};
|
||||
|
||||
workflowsStore.setWorkflowActive('1', mockActiveVersion);
|
||||
|
||||
expect(workflowsStore.activeWorkflows).toContain('1');
|
||||
expect(workflowsStore.workflowsById['1'].active).toBe(true);
|
||||
@@ -910,7 +921,17 @@ describe('useWorkflowsStore', () => {
|
||||
workflowsStore.workflowsById = { '1': { active: true } as IWorkflowDb };
|
||||
workflowsStore.workflow.id = '1';
|
||||
|
||||
workflowsStore.setWorkflowActive('1');
|
||||
const mockActiveVersion: WorkflowHistory = {
|
||||
versionId: 'test-version-id',
|
||||
name: 'Test Version',
|
||||
authors: 'Test Author',
|
||||
description: 'A test workflow version',
|
||||
createdAt: new Date().toISOString(),
|
||||
updatedAt: new Date().toISOString(),
|
||||
workflowPublishHistory: [],
|
||||
};
|
||||
|
||||
workflowsStore.setWorkflowActive('1', mockActiveVersion);
|
||||
|
||||
expect(workflowsStore.activeWorkflows).toEqual(['1']);
|
||||
expect(workflowsStore.workflowsById['1'].active).toBe(true);
|
||||
@@ -921,7 +942,18 @@ describe('useWorkflowsStore', () => {
|
||||
uiStore.stateIsDirty = true;
|
||||
workflowsStore.workflow.id = '1';
|
||||
workflowsStore.workflowsById = { '1': { active: false } as IWorkflowDb };
|
||||
workflowsStore.setWorkflowActive('2');
|
||||
|
||||
const mockActiveVersion: WorkflowHistory = {
|
||||
versionId: 'test-version-id',
|
||||
name: 'Test Version',
|
||||
authors: 'Test Author',
|
||||
description: 'A test workflow version',
|
||||
createdAt: new Date().toISOString(),
|
||||
updatedAt: new Date().toISOString(),
|
||||
workflowPublishHistory: [],
|
||||
};
|
||||
|
||||
workflowsStore.setWorkflowActive('2', mockActiveVersion);
|
||||
expect(workflowsStore.workflowsById['1'].active).toBe(false);
|
||||
expect(uiStore.stateIsDirty).toBe(true);
|
||||
});
|
||||
|
||||
@@ -749,6 +749,11 @@ export const useWorkflowsStore = defineStore(STORES.WORKFLOWS, () => {
|
||||
workflowChecksum.value = checksum;
|
||||
}
|
||||
|
||||
async function updateWorkflowChecksum() {
|
||||
const checksum = await calculateWorkflowChecksum(workflow.value);
|
||||
setWorkflowChecksum(checksum);
|
||||
}
|
||||
|
||||
function setWorkflowActiveVersion(version: WorkflowHistory) {
|
||||
workflow.value.activeVersion = deepCopy(version);
|
||||
}
|
||||
@@ -899,21 +904,28 @@ export const useWorkflowsStore = defineStore(STORES.WORKFLOWS, () => {
|
||||
};
|
||||
}
|
||||
|
||||
function setWorkflowActive(targetWorkflowId: string, activeVersion?: WorkflowHistory) {
|
||||
const index = activeWorkflows.value.indexOf(targetWorkflowId);
|
||||
if (index === -1) {
|
||||
function setWorkflowActive(
|
||||
targetWorkflowId: string,
|
||||
activeVersion: WorkflowHistory,
|
||||
clearDirtyState: boolean = true,
|
||||
) {
|
||||
if (activeWorkflows.value.indexOf(targetWorkflowId) === -1) {
|
||||
activeWorkflows.value.push(targetWorkflowId);
|
||||
}
|
||||
const targetWorkflow = workflowsById.value[targetWorkflowId];
|
||||
if (targetWorkflow) {
|
||||
targetWorkflow.active = true;
|
||||
targetWorkflow.activeVersionId = activeVersion?.versionId ?? targetWorkflow.versionId;
|
||||
targetWorkflow.activeVersion = activeVersion;
|
||||
|
||||
const cachedWorkflow = workflowsById.value[targetWorkflowId];
|
||||
if (cachedWorkflow) {
|
||||
cachedWorkflow.active = true;
|
||||
cachedWorkflow.activeVersionId = activeVersion.versionId;
|
||||
cachedWorkflow.activeVersion = activeVersion;
|
||||
}
|
||||
|
||||
if (targetWorkflowId === workflow.value.id) {
|
||||
uiStore.stateIsDirty = false;
|
||||
if (clearDirtyState) {
|
||||
uiStore.stateIsDirty = false;
|
||||
}
|
||||
workflow.value.active = true;
|
||||
workflow.value.activeVersionId = activeVersion?.versionId ?? workflow.value.versionId;
|
||||
workflow.value.activeVersionId = activeVersion.versionId;
|
||||
workflow.value.activeVersion = activeVersion;
|
||||
}
|
||||
}
|
||||
@@ -2018,6 +2030,7 @@ export const useWorkflowsStore = defineStore(STORES.WORKFLOWS, () => {
|
||||
setUsedCredentials,
|
||||
setWorkflowVersionId,
|
||||
setWorkflowChecksum,
|
||||
updateWorkflowChecksum,
|
||||
setWorkflowActiveVersion,
|
||||
replaceInvalidWorkflowCredentials,
|
||||
assignCredentialToMatchingNodes,
|
||||
|
||||
Reference in New Issue
Block a user