feat(core): Implement function to index a single workflow (no-changelog) (#21373)

This commit is contained in:
mfsiega
2025-10-30 14:22:47 +01:00
committed by GitHub
parent 3b53649165
commit c47b185f04
6 changed files with 404 additions and 23 deletions
@@ -11,7 +11,7 @@ import {
import { WithCreatedAt } from './abstract-entity';
import type { WorkflowEntity } from './workflow-entity';
export type DependencyType = 'credential' | 'nodeType' | 'webhookPath' | 'workflowCall';
export type DependencyType = 'credentialId' | 'nodeType' | 'webhookPath' | 'workflowCall';
@Entity({ name: 'workflow_dependency' })
export class WorkflowDependency extends WithCreatedAt {
@@ -34,7 +34,7 @@ export class WorkflowDependency extends WithCreatedAt {
/**
* The type of the dependency.
* credential | nodeType | webhookPath | workflowCall
* credentialId | nodeType | webhookPath | workflowCall
*/
@Column({ length: 32 })
@Index()
@@ -42,7 +42,7 @@ export class WorkflowDependency extends WithCreatedAt {
/**
* The ID of the dependency, interpreted based on the dependency type.
* E.g., for 'credential' it would be the credential ID, for 'nodeType' the node type name, etc.
* E.g., for 'credentialId' it would be the credential ID, for 'nodeType' the node type name, etc.
*/
@Column({ length: 255 })
@Index()
@@ -4,6 +4,8 @@ import { DataSource, EntityManager, LessThan, Repository } from '@n8n/typeorm';
import { WorkflowDependency } from '../entities';
const INDEX_VERSION_ID = 1;
/**
* Helper class to collect workflow dependencies before writing them to the database.
*/
@@ -13,20 +15,19 @@ export class WorkflowDependencies {
constructor(
readonly workflowId: string,
readonly workflowVersionId: number | undefined,
readonly indexVersionId: number,
) {}
add(dependency: {
dependencyType: string;
dependencyKey: string | null;
dependencyInfo: unknown;
dependencyInfo: Record<string, unknown> | null;
}) {
const dep = new WorkflowDependency();
Object.assign(dep, dependency);
Object.assign(dep, {
workflowId: this.workflowId,
workflowVersionId: this.workflowVersionId,
indexVersionId: this.indexVersionId,
indexVersionId: INDEX_VERSION_ID,
});
this.dependencies.push(dep);
}
@@ -0,0 +1,240 @@
/* eslint-disable @typescript-eslint/unbound-method */
/* eslint-disable @typescript-eslint/no-unsafe-assignment */
import type { Logger } from '@n8n/backend-common';
import type { WorkflowDependencyRepository } from '@n8n/db';
import type { ErrorReporter } from 'n8n-core';
import type { INode, IWorkflowBase } from 'n8n-workflow';
import { WorkflowIndexService } from '../workflow-index.service';
describe('WorkflowIndexService', () => {
let service: WorkflowIndexService;
let mockRepository: jest.Mocked<WorkflowDependencyRepository>;
let mockLogger: jest.Mocked<Logger>;
let mockErrorReporter: jest.Mocked<ErrorReporter>;
beforeEach(() => {
mockRepository = {
updateDependenciesForWorkflow: jest.fn(),
} as unknown as jest.Mocked<WorkflowDependencyRepository>;
mockLogger = {
debug: jest.fn(),
error: jest.fn(),
} as unknown as jest.Mocked<Logger>;
mockErrorReporter = {
error: jest.fn(),
warn: jest.fn(),
} as unknown as jest.Mocked<ErrorReporter>;
service = new WorkflowIndexService(mockRepository, mockLogger, mockErrorReporter);
});
const createNode = (overrides: Partial<INode> & { id: string; type: string }): INode => {
const { parameters, ...rest } = overrides;
return {
name: overrides.id,
typeVersion: 1,
position: [0, 0],
parameters: parameters ?? {},
...rest,
};
};
const createWorkflow = (nodes: INode[]): IWorkflowBase => ({
id: 'workflow-123',
name: 'Test Workflow',
active: true,
isArchived: false,
createdAt: new Date(),
updatedAt: new Date(),
versionCounter: 1,
nodes,
connections: {},
});
it('should extract all dependency types correctly', async () => {
mockRepository.updateDependenciesForWorkflow.mockResolvedValue(true);
const workflow = createWorkflow([
createNode({
id: 'node-1',
type: 'n8n-nodes-base.webhook',
parameters: { path: 'webhook-1' },
}),
createNode({
id: 'node-2',
type: 'n8n-nodes-base.httpRequest',
credentials: {
httpAuth: { id: 'cred-1', name: 'Auth 1' },
apiKey: { id: 'cred-2', name: 'Auth 2' },
},
}),
createNode({
id: 'node-3',
type: 'n8n-nodes-base.executeWorkflow',
parameters: { workflowId: 'sub-workflow-1' },
}),
createNode({
id: 'node-4',
type: 'n8n-nodes-base.executeWorkflow',
parameters: { workflowId: { value: 'sub-workflow-2' } },
}),
]);
await service.updateIndexFor(workflow);
expect(mockRepository.updateDependenciesForWorkflow).toHaveBeenCalledWith(
'workflow-123',
expect.objectContaining({
dependencies: expect.arrayContaining([
// nodeType dependencies
expect.objectContaining({
dependencyType: 'nodeType',
dependencyKey: 'n8n-nodes-base.webhook',
dependencyInfo: { nodeId: 'node-1', nodeVersion: 1 },
}),
expect.objectContaining({
dependencyType: 'nodeType',
dependencyKey: 'n8n-nodes-base.httpRequest',
dependencyInfo: { nodeId: 'node-2', nodeVersion: 1 },
}),
expect.objectContaining({
dependencyType: 'nodeType',
dependencyKey: 'n8n-nodes-base.executeWorkflow',
dependencyInfo: { nodeId: 'node-3', nodeVersion: 1 },
}),
expect.objectContaining({
dependencyType: 'nodeType',
dependencyKey: 'n8n-nodes-base.executeWorkflow',
dependencyInfo: { nodeId: 'node-4', nodeVersion: 1 },
}),
// webhookPath dependencies
expect.objectContaining({
dependencyType: 'webhookPath',
dependencyKey: 'webhook-1',
dependencyInfo: { nodeId: 'node-1', nodeVersion: 1 },
}),
// credentialId dependencies
expect.objectContaining({
dependencyType: 'credentialId',
dependencyKey: 'cred-1',
dependencyInfo: { nodeId: 'node-2', nodeVersion: 1 },
}),
expect.objectContaining({
dependencyType: 'credentialId',
dependencyKey: 'cred-2',
dependencyInfo: { nodeId: 'node-2', nodeVersion: 1 },
}),
// workflowCall dependencies (both string and object format)
expect.objectContaining({
dependencyType: 'workflowCall',
dependencyKey: 'sub-workflow-1',
dependencyInfo: { nodeId: 'node-3', nodeVersion: 1 },
}),
expect.objectContaining({
dependencyType: 'workflowCall',
dependencyKey: 'sub-workflow-2',
dependencyInfo: { nodeId: 'node-4', nodeVersion: 1 },
}),
]),
}),
);
});
it('should handle repository errors gracefully', async () => {
const error = new Error('Database error');
mockRepository.updateDependenciesForWorkflow.mockRejectedValue(error);
const workflow = createWorkflow([
createNode({
id: 'node-1',
type: 'n8n-nodes-base.start',
}),
]);
await service.updateIndexFor(workflow);
expect(mockLogger.error).toHaveBeenCalledWith(
'Failed to update workflow dependency index for workflow workflow-123: Database error',
);
expect(mockErrorReporter.error).toHaveBeenCalledWith(error);
});
it('should not create workflowCall dependencies for parameter, localFile, and url sources', async () => {
mockRepository.updateDependenciesForWorkflow.mockResolvedValue(true);
const workflow = createWorkflow([
createNode({
id: 'node-1',
type: 'n8n-nodes-base.executeWorkflow',
parameters: { source: 'parameter' },
}),
createNode({
id: 'node-2',
type: 'n8n-nodes-base.executeWorkflow',
parameters: { source: 'localFile' },
}),
createNode({
id: 'node-3',
type: 'n8n-nodes-base.executeWorkflow',
parameters: { source: 'url' },
}),
]);
await service.updateIndexFor(workflow);
expect(mockRepository.updateDependenciesForWorkflow).toHaveBeenCalledWith(
'workflow-123',
expect.objectContaining({
dependencies: expect.not.arrayContaining([
expect.objectContaining({
dependencyType: 'workflowCall',
}),
]),
}),
);
});
it('should handle multiple credentials on a single node', async () => {
mockRepository.updateDependenciesForWorkflow.mockResolvedValue(true);
const workflow = createWorkflow([
createNode({
id: 'node-1',
type: 'n8n-nodes-base.httpRequest',
credentials: {
httpAuth: { id: 'cred-1', name: 'Basic Auth' },
apiKey: { id: 'cred-2', name: 'API Key' },
oAuth2: { id: 'cred-3', name: 'OAuth2' },
},
}),
]);
await service.updateIndexFor(workflow);
expect(mockRepository.updateDependenciesForWorkflow).toHaveBeenCalledWith(
'workflow-123',
expect.objectContaining({
dependencies: expect.arrayContaining([
expect.objectContaining({
dependencyType: 'credentialId',
dependencyKey: 'cred-1',
dependencyInfo: { nodeId: 'node-1', nodeVersion: 1 },
}),
expect.objectContaining({
dependencyType: 'credentialId',
dependencyKey: 'cred-2',
dependencyInfo: { nodeId: 'node-1', nodeVersion: 1 },
}),
expect.objectContaining({
dependencyType: 'credentialId',
dependencyKey: 'cred-3',
dependencyInfo: { nodeId: 'node-1', nodeVersion: 1 },
}),
]),
}),
);
});
});
@@ -0,0 +1,139 @@
import { Logger } from '@n8n/backend-common';
import { WorkflowDependencies, WorkflowDependencyRepository } from '@n8n/db';
import { Service } from '@n8n/di';
import { ErrorReporter } from 'n8n-core';
import { ensureError, INode, IWorkflowBase } from 'n8n-workflow';
/**
* Service for managing the workflow dependency index. The index tracks dependencies such as node types,
* credentials, workflow calls, and webhook paths used by each workflow. The service builds the index on server start
* and updates it in response to workflow-related events.
*
* TODO(CAT-1595): Build the index on startup.
* TODO(CAT-1597): Update the index in realtime.
*/
@Service()
export class WorkflowIndexService {
constructor(
private readonly dependencyRepository: WorkflowDependencyRepository,
private readonly logger: Logger,
private readonly errorReporter: ErrorReporter,
) {}
/**
* Update the dependency index for a given workflow.
*
* NOTE: this should generally be handled via events, rather than called directly.
* The exception is during workflow imports where it's simpler to call directly.
*
*/
async updateIndexFor(workflow: IWorkflowBase) {
// TODO: input validation.
// Generate the dependency updates for the given workflow.
const dependencyUpdates = new WorkflowDependencies(workflow.id, workflow.versionCounter);
workflow.nodes.forEach((node) => {
this.addNodeTypeDependencies(node, dependencyUpdates);
this.addCredentialDependencies(node, dependencyUpdates);
this.addWorkflowCallDependencies(node, dependencyUpdates);
this.addWebhookPathDependencies(node, dependencyUpdates);
});
let updated: boolean;
try {
updated = await this.dependencyRepository.updateDependenciesForWorkflow(
workflow.id,
dependencyUpdates,
);
} catch (e) {
const error = ensureError(e);
this.logger.error(
`Failed to update workflow dependency index for workflow ${workflow.id}: ${error.message}`,
);
this.errorReporter.error(error);
return;
}
this.logger.debug(
`Workflow dependency index ${updated ? 'updated' : 'skipped'} for workflow ${workflow.id}`,
);
}
private addNodeTypeDependencies(node: INode, dependencyUpdates: WorkflowDependencies): void {
dependencyUpdates.add({
dependencyType: 'nodeType',
dependencyKey: node.type,
dependencyInfo: { nodeId: node.id, nodeVersion: node.typeVersion },
});
}
private addCredentialDependencies(node: INode, dependencyUpdates: WorkflowDependencies): void {
if (!node.credentials) {
return;
}
for (const credentialDetails of Object.values(node.credentials)) {
const { id } = credentialDetails;
dependencyUpdates.add({
dependencyType: 'credentialId',
dependencyKey: id,
dependencyInfo: { nodeId: node.id, nodeVersion: node.typeVersion },
});
}
}
private addWorkflowCallDependencies(node: INode, dependencyUpdates: WorkflowDependencies): void {
if (node.type !== 'n8n-nodes-base.executeWorkflow') {
return;
}
const calledWorkflowId: string | undefined = this.getCalledWorkflowIdFrom(node);
if (!calledWorkflowId) {
return;
}
dependencyUpdates.add({
dependencyType: 'workflowCall',
dependencyKey: calledWorkflowId,
dependencyInfo: { nodeId: node.id, nodeVersion: node.typeVersion },
});
}
private addWebhookPathDependencies(node: INode, dependencyUpdates: WorkflowDependencies): void {
if (node.type !== 'n8n-nodes-base.webhook') {
return;
}
const webhookPath = node.parameters.path as string;
dependencyUpdates.add({
dependencyType: 'webhookPath',
dependencyKey: webhookPath,
dependencyInfo: { nodeId: node.id, nodeVersion: node.typeVersion },
});
}
private getCalledWorkflowIdFrom(node: INode): string | undefined {
if (node.parameters?.['source'] === 'parameter') {
return undefined; // The sub-workflow is provided directly in the parameters, so no dependency to track.
}
if (node.parameters?.['source'] === 'localFile') {
return undefined; // The sub-workflow is provided via a local file, so no dependency to track.
}
if (node.parameters?.['source'] === 'url') {
return undefined; // The sub-workflow is provided via a URL, so no dependency to track.
}
// If it's none of those sources, it must be 'workflowId'. This might be either directly as a string, or an object.
if (typeof node.parameters?.['workflowId'] === 'string') {
return node.parameters?.['workflowId'];
}
if (
node.parameters &&
typeof node.parameters['workflowId'] === 'object' &&
node.parameters['workflowId'] !== null &&
'value' in node.parameters['workflowId'] &&
typeof node.parameters['workflowId']['value'] === 'string'
) {
return node.parameters['workflowId']['value'];
}
this.errorReporter.warn(
`While indexing, could not determine called workflow ID from executeWorkflow node ${node.id}`,
{ extra: node.parameters },
);
return undefined;
}
}
@@ -51,9 +51,9 @@ if (globalConfig.database.isLegacySqlite) {
// ARRANGE
//
const workflow = await createWorkflow({ versionId: 'v1' });
const dependencies = new WorkflowDependencies(workflow.id, 1, 1);
const dependencies = new WorkflowDependencies(workflow.id, 1);
dependencies.add({
dependencyType: 'credential',
dependencyType: 'credentialId',
dependencyKey: 'cred-123',
dependencyInfo: { name: 'Test Credential' },
});
@@ -83,7 +83,7 @@ if (globalConfig.database.isLegacySqlite) {
expect(savedDependencies[0]).toMatchObject({
workflowId: workflow.id,
workflowVersionId: 1,
dependencyType: 'credential',
dependencyType: 'credentialId',
dependencyKey: 'cred-123',
dependencyInfo: { name: 'Test Credential' },
indexVersionId: 1,
@@ -105,18 +105,18 @@ if (globalConfig.database.isLegacySqlite) {
const workflow = await createWorkflow({ versionId: 'v1' });
// Insert initial dependencies with version 1
const initialDeps = new WorkflowDependencies(workflow.id, 1, 1);
const initialDeps = new WorkflowDependencies(workflow.id, 1);
initialDeps.add({
dependencyType: 'credential',
dependencyType: 'credentialId',
dependencyKey: 'cred-old',
dependencyInfo: null,
});
await workflowDependencyRepository.updateDependenciesForWorkflow(workflow.id, initialDeps);
// Create new dependencies with version 2
const updatedDeps = new WorkflowDependencies(workflow.id, 2, 1);
const updatedDeps = new WorkflowDependencies(workflow.id, 2);
updatedDeps.add({
dependencyType: 'credential',
dependencyType: 'credentialId',
dependencyKey: 'cred-new',
dependencyInfo: { updated: true },
});
@@ -156,18 +156,18 @@ if (globalConfig.database.isLegacySqlite) {
const workflow = await createWorkflow({ versionId: 'v2' });
// Insert dependencies with version 2
const newerDeps = new WorkflowDependencies(workflow.id, 2, 1);
const newerDeps = new WorkflowDependencies(workflow.id, 2);
newerDeps.add({
dependencyType: 'credential',
dependencyType: 'credentialId',
dependencyKey: 'cred-new',
dependencyInfo: null,
});
await workflowDependencyRepository.updateDependenciesForWorkflow(workflow.id, newerDeps);
// Try to update with older version 1
const olderDeps = new WorkflowDependencies(workflow.id, 1, 1);
const olderDeps = new WorkflowDependencies(workflow.id, 1);
olderDeps.add({
dependencyType: 'credential',
dependencyType: 'credentialId',
dependencyKey: 'cred-old',
dependencyInfo: null,
});
@@ -198,16 +198,16 @@ if (globalConfig.database.isLegacySqlite) {
//
const workflow = await createWorkflow({ versionId: '2' });
const depsVersion1 = new WorkflowDependencies(workflow.id, 1, 1);
const depsVersion1 = new WorkflowDependencies(workflow.id, 1);
depsVersion1.add({
dependencyType: 'credential',
dependencyType: 'credentialId',
dependencyKey: 'cred-1',
dependencyInfo: null,
});
const depsVersion2 = new WorkflowDependencies(workflow.id, 2, 1);
const depsVersion2 = new WorkflowDependencies(workflow.id, 2);
depsVersion2.add({
dependencyType: 'credential',
dependencyType: 'credentialId',
dependencyKey: 'cred-2',
dependencyInfo: null,
});
@@ -242,9 +242,9 @@ if (globalConfig.database.isLegacySqlite) {
// ARRANGE
//
const workflow = await createWorkflow({ versionId: 'v1' });
const dependencies = new WorkflowDependencies(workflow.id, 1, 1);
const dependencies = new WorkflowDependencies(workflow.id, 1);
dependencies.add({
dependencyType: 'credential',
dependencyType: 'credentialId',
dependencyKey: 'cred-1',
dependencyInfo: null,
});
+1
View File
@@ -2576,6 +2576,7 @@ export interface IWorkflowBase {
staticData?: IDataObject;
pinData?: IPinData;
versionId?: string;
versionCounter?: number;
meta?: WorkflowFEMeta;
}