From c47b185f04de8c72240733fc9dfe80ff9be3818e Mon Sep 17 00:00:00 2001 From: mfsiega <93014743+mfsiega@users.noreply.github.com> Date: Thu, 30 Oct 2025 14:22:47 +0100 Subject: [PATCH] feat(core): Implement function to index a single workflow (no-changelog) (#21373) --- .../entities/workflow-dependency-entity.ts | 6 +- .../workflow-dependency.repository.ts | 7 +- .../__tests__/workflow-index.service.test.ts | 240 ++++++++++++++++++ .../workflow-index/workflow-index.service.ts | 139 ++++++++++ .../workflow-dependency.repository.test.ts | 34 +-- packages/workflow/src/interfaces.ts | 1 + 6 files changed, 404 insertions(+), 23 deletions(-) create mode 100644 packages/cli/src/modules/workflow-index/__tests__/workflow-index.service.test.ts create mode 100644 packages/cli/src/modules/workflow-index/workflow-index.service.ts diff --git a/packages/@n8n/db/src/entities/workflow-dependency-entity.ts b/packages/@n8n/db/src/entities/workflow-dependency-entity.ts index 06007fbffaf..02ff1458896 100644 --- a/packages/@n8n/db/src/entities/workflow-dependency-entity.ts +++ b/packages/@n8n/db/src/entities/workflow-dependency-entity.ts @@ -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() diff --git a/packages/@n8n/db/src/repositories/workflow-dependency.repository.ts b/packages/@n8n/db/src/repositories/workflow-dependency.repository.ts index f832967c763..3b4588c7bbf 100644 --- a/packages/@n8n/db/src/repositories/workflow-dependency.repository.ts +++ b/packages/@n8n/db/src/repositories/workflow-dependency.repository.ts @@ -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 | 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); } diff --git a/packages/cli/src/modules/workflow-index/__tests__/workflow-index.service.test.ts b/packages/cli/src/modules/workflow-index/__tests__/workflow-index.service.test.ts new file mode 100644 index 00000000000..d9d22b56a4d --- /dev/null +++ b/packages/cli/src/modules/workflow-index/__tests__/workflow-index.service.test.ts @@ -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; + let mockLogger: jest.Mocked; + let mockErrorReporter: jest.Mocked; + + beforeEach(() => { + mockRepository = { + updateDependenciesForWorkflow: jest.fn(), + } as unknown as jest.Mocked; + + mockLogger = { + debug: jest.fn(), + error: jest.fn(), + } as unknown as jest.Mocked; + + mockErrorReporter = { + error: jest.fn(), + warn: jest.fn(), + } as unknown as jest.Mocked; + + service = new WorkflowIndexService(mockRepository, mockLogger, mockErrorReporter); + }); + + const createNode = (overrides: Partial & { 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 }, + }), + ]), + }), + ); + }); +}); diff --git a/packages/cli/src/modules/workflow-index/workflow-index.service.ts b/packages/cli/src/modules/workflow-index/workflow-index.service.ts new file mode 100644 index 00000000000..0f7740da719 --- /dev/null +++ b/packages/cli/src/modules/workflow-index/workflow-index.service.ts @@ -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; + } +} diff --git a/packages/cli/test/integration/database/repositories/workflow-dependency.repository.test.ts b/packages/cli/test/integration/database/repositories/workflow-dependency.repository.test.ts index 0a691ab556a..6f2663f0694 100644 --- a/packages/cli/test/integration/database/repositories/workflow-dependency.repository.test.ts +++ b/packages/cli/test/integration/database/repositories/workflow-dependency.repository.test.ts @@ -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, }); diff --git a/packages/workflow/src/interfaces.ts b/packages/workflow/src/interfaces.ts index c9e896a2106..d850465f3ff 100644 --- a/packages/workflow/src/interfaces.ts +++ b/packages/workflow/src/interfaces.ts @@ -2576,6 +2576,7 @@ export interface IWorkflowBase { staticData?: IDataObject; pinData?: IPinData; versionId?: string; + versionCounter?: number; meta?: WorkflowFEMeta; }