diff --git a/packages/@n8n/api-types/src/index.ts b/packages/@n8n/api-types/src/index.ts index 8833f123780..94317af4a88 100644 --- a/packages/@n8n/api-types/src/index.ts +++ b/packages/@n8n/api-types/src/index.ts @@ -388,6 +388,13 @@ export { export type { AgentRunState, AgentNode } from './schemas/agent-run-reducer'; +export { + EVAL_PARALLEL_EXECUTION_FLAG, + startTestRunPayloadSchema, + StartTestRunRequestDto, + type StartTestRunPayload, +} from './schemas/evaluations.schema'; + export { ALLOWED_DOMAINS, isAllowedDomain } from './utils/allowed-domains'; export { diff --git a/packages/@n8n/api-types/src/schemas/evaluations.schema.ts b/packages/@n8n/api-types/src/schemas/evaluations.schema.ts new file mode 100644 index 00000000000..56d871bcb30 --- /dev/null +++ b/packages/@n8n/api-types/src/schemas/evaluations.schema.ts @@ -0,0 +1,26 @@ +import { z } from 'zod'; + +import { Z } from '../zod-class'; + +// Single source of truth for the parallel-execution rollout flag id, shared +// between the FE checkbox-rendering gate and the BE controller's safety net. +// Strings can drift if duplicated; importing from a shared module cannot. +export const EVAL_PARALLEL_EXECUTION_FLAG = '080_eval_parallel_execution'; + +// `concurrency` is the optional number of evaluation test cases to run in +// parallel for a single test run. Clamped 1–10. When omitted, the runner +// falls back to sequential execution (concurrency = 1). The PostHog +// rollout flag `080_eval_parallel_execution` gates whether the controller +// honours values > 1; flag-off requests are silently coerced to 1 so the +// flag id never leaks into HTTP responses. +const startTestRunPayloadShape = { + concurrency: z.number().int().min(1).max(10).optional(), +}; + +export const startTestRunPayloadSchema = z.object(startTestRunPayloadShape); +export type StartTestRunPayload = z.infer; + +// Controller-side DTO used by the @Body decorator's reflection-based +// validation. Shares the same shape as `startTestRunPayloadSchema` — +// single source of truth so the two validators cannot silently diverge. +export class StartTestRunRequestDto extends Z.class(startTestRunPayloadShape) {} diff --git a/packages/@n8n/config/src/configs/evaluation.config.ts b/packages/@n8n/config/src/configs/evaluation.config.ts new file mode 100644 index 00000000000..302c99dd114 --- /dev/null +++ b/packages/@n8n/config/src/configs/evaluation.config.ts @@ -0,0 +1,25 @@ +import { Config, Env } from '../decorators'; + +@Config +export class EvaluationConfig { + /** + * Force-enable the parallel-execution feature for evaluation test runs. + * + * Acts as an operator-level override of the `080_eval_parallel_execution` + * PostHog rollout flag. When set to `true`, the FE renders the + * concurrency UI for every user and the BE honours `concurrency` payloads + * regardless of PostHog cohort. When `false` (default), PostHog remains + * the source of truth — the rollout flag controls visibility per-user. + * + * Useful for: + * - Local development without PostHog wiring. + * - Operator escape hatch if PostHog is unreachable. + * - Self-hosted deployments that want the feature without PostHog + * dependency. + * + * Cannot force-disable: setting this to `false` falls back to PostHog, + * not a kill-switch. Use PostHog itself to disable a rolled-out flag. + */ + @Env('N8N_EVAL_PARALLEL_EXECUTION_ENABLED') + parallelExecutionEnabled: boolean = false; +} diff --git a/packages/@n8n/config/src/index.ts b/packages/@n8n/config/src/index.ts index f80eb7d33b9..0cf2aad304d 100644 --- a/packages/@n8n/config/src/index.ts +++ b/packages/@n8n/config/src/index.ts @@ -14,6 +14,7 @@ import { DeploymentConfig } from './configs/deployment.config'; import { DiagnosticsConfig } from './configs/diagnostics.config'; import { DynamicBannersConfig } from './configs/dynamic-banners.config'; import { EndpointsConfig } from './configs/endpoints.config'; +import { EvaluationConfig } from './configs/evaluation.config'; import { EventBusConfig } from './configs/event-bus.config'; import { ExecutionsConfig } from './configs/executions.config'; import { ExpressionEngineConfig } from './configs/expression-engine.config'; @@ -161,6 +162,9 @@ export class GlobalConfig { @Nested multiMainSetup: MultiMainSetupConfig; + @Nested + evaluation: EvaluationConfig; + @Nested generic: GenericConfig; diff --git a/packages/@n8n/config/test/config.test.ts b/packages/@n8n/config/test/config.test.ts index dc86b55f2af..db761c7a773 100644 --- a/packages/@n8n/config/test/config.test.ts +++ b/packages/@n8n/config/test/config.test.ts @@ -371,6 +371,9 @@ describe('GlobalConfig', () => { ttl: 10, interval: 3, }, + evaluation: { + parallelExecutionEnabled: false, + }, generic: { timezone: 'America/New_York', releaseChannel: 'dev', diff --git a/packages/cli/package.json b/packages/cli/package.json index 42d6651e3cb..b35d1081fbb 100644 --- a/packages/cli/package.json +++ b/packages/cli/package.json @@ -182,6 +182,7 @@ "otpauth": "9.1.1", "p-cancelable": "2.1.1", "p-lazy": "3.1.0", + "p-limit": "^3.1.0", "pg": "catalog:", "picocolors": "catalog:", "pkce-challenge": "5.0.0", diff --git a/packages/cli/src/evaluation.ee/__tests__/test-runs.controller.ee.test.ts b/packages/cli/src/evaluation.ee/__tests__/test-runs.controller.ee.test.ts index b86420dea12..1ba9a0c1fbe 100644 --- a/packages/cli/src/evaluation.ee/__tests__/test-runs.controller.ee.test.ts +++ b/packages/cli/src/evaluation.ee/__tests__/test-runs.controller.ee.test.ts @@ -1,9 +1,12 @@ +import type { Logger } from '@n8n/backend-common'; import type { TestCaseExecutionRepository, TestRun, TestRunRepository, User } from '@n8n/db'; +import type express from 'express'; import { NotFoundError } from '@/errors/response-errors/not-found.error'; import type { TestRunnerService } from '@/evaluation.ee/test-runner/test-runner.service.ee'; import { TestRunsController } from '@/evaluation.ee/test-runs.controller.ee'; import type { TestRunsRequest } from '@/evaluation.ee/test-runs.types.ee'; +import type { PostHogClient } from '@/posthog'; import type { Telemetry } from '@/telemetry'; import type { WorkflowFinderService } from '@/workflows/workflow-finder.service'; @@ -16,6 +19,8 @@ describe('TestRunsController', () => { let mockTestCaseExecutionRepository: jest.Mocked; let mockTestRunnerService: jest.Mocked; let mockTelemetry: jest.Mocked; + let mockPostHogClient: jest.Mocked; + let mockLogger: jest.Mocked; let mockUser: User; let mockWorkflowId: string; let mockTestRunId: string; @@ -47,15 +52,28 @@ describe('TestRunsController', () => { track: jest.fn(), } as unknown as jest.Mocked; + mockPostHogClient = { + getFeatureFlags: jest.fn().mockResolvedValue({}), + } as unknown as jest.Mocked; + + mockLogger = { + warn: jest.fn(), + debug: jest.fn(), + error: jest.fn(), + info: jest.fn(), + } as unknown as jest.Mocked; + testRunsController = new TestRunsController( mockTestRunRepository, mockWorkflowFinderService, mockTestCaseExecutionRepository, mockTestRunnerService, mockTelemetry, + mockPostHogClient, + mockLogger, ); - mockUser = { id: 'user123' } as User; + mockUser = { id: 'user123', createdAt: new Date('2024-01-01T00:00:00Z') } as User; mockWorkflowId = 'workflow123'; mockTestRunId = 'testrun123'; @@ -157,4 +175,111 @@ describe('TestRunsController', () => { }); }); }); + + describe('create', () => { + const buildCreateRequest = () => + ({ + params: { workflowId: mockWorkflowId }, + user: mockUser, + }) as unknown as TestRunsRequest.Create; + + const mockResponse = () => { + const res = { status: jest.fn(), json: jest.fn() } as unknown as express.Response; + (res.status as jest.Mock).mockReturnValue(res); + (res.json as jest.Mock).mockReturnValue(res); + return res; + }; + + it('flag-on user with concurrency=5 → service called with concurrency=5 and flagEnabledForUser=true', async () => { + mockPostHogClient.getFeatureFlags.mockResolvedValue({ '080_eval_parallel_execution': true }); + + await testRunsController.create( + buildCreateRequest(), + mockResponse() as any, + { concurrency: 5 } as any, + ); + + expect(mockPostHogClient.getFeatureFlags).toHaveBeenCalledWith(mockUser); + expect(mockTestRunnerService.runTest).toHaveBeenCalledWith(mockUser, mockWorkflowId, 5, true); + }); + + it('flag-off user with concurrency=5 → service called with concurrency=1 and flagEnabledForUser=false (cohort wall)', async () => { + mockPostHogClient.getFeatureFlags.mockResolvedValue({}); + + await testRunsController.create( + buildCreateRequest(), + mockResponse() as any, + { concurrency: 5 } as any, + ); + + expect(mockTestRunnerService.runTest).toHaveBeenCalledWith( + mockUser, + mockWorkflowId, + 1, + false, + ); + }); + + it('flag-on user with no concurrency body → service called with concurrency=1', async () => { + mockPostHogClient.getFeatureFlags.mockResolvedValue({ '080_eval_parallel_execution': true }); + + await testRunsController.create(buildCreateRequest(), mockResponse() as any, {} as any); + + expect(mockTestRunnerService.runTest).toHaveBeenCalledWith(mockUser, mockWorkflowId, 1, true); + }); + + it('flag-off user with no concurrency body → service called with concurrency=1', async () => { + mockPostHogClient.getFeatureFlags.mockResolvedValue({}); + + await testRunsController.create(buildCreateRequest(), mockResponse() as any, {} as any); + + expect(mockTestRunnerService.runTest).toHaveBeenCalledWith( + mockUser, + mockWorkflowId, + 1, + false, + ); + }); + + it('always returns 202 success regardless of flag state (no flag-id leak)', async () => { + mockPostHogClient.getFeatureFlags.mockResolvedValue({}); + + const res = mockResponse(); + await testRunsController.create(buildCreateRequest(), res as any, { concurrency: 7 } as any); + + expect(res.status).toHaveBeenCalledWith(202); + expect(res.json).toHaveBeenCalledWith({ success: true }); + }); + + it('resolves the feature flag exactly once per request', async () => { + mockPostHogClient.getFeatureFlags.mockResolvedValue({ '080_eval_parallel_execution': true }); + + await testRunsController.create( + buildCreateRequest(), + mockResponse() as any, + { concurrency: 3 } as any, + ); + + expect(mockPostHogClient.getFeatureFlags).toHaveBeenCalledTimes(1); + }); + + it('fails open to sequential when PostHog throws (rollout gate is non-critical)', async () => { + mockPostHogClient.getFeatureFlags.mockRejectedValue(new Error('posthog timeout')); + + const res = mockResponse(); + await testRunsController.create(buildCreateRequest(), res as any, { concurrency: 5 } as any); + + expect(mockTestRunnerService.runTest).toHaveBeenCalledWith( + mockUser, + mockWorkflowId, + 1, + false, + ); + expect(res.status).toHaveBeenCalledWith(202); + expect(mockLogger.warn).toHaveBeenCalledWith( + expect.stringContaining('Failed to resolve eval parallel-execution flag'), + expect.any(Object), + ); + }); + }); }); diff --git a/packages/cli/src/evaluation.ee/test-runner/__tests__/evaluation-metrics.ee.test.ts b/packages/cli/src/evaluation.ee/test-runner/__tests__/evaluation-metrics.ee.test.ts index 01205a8e052..946ddb83903 100644 --- a/packages/cli/src/evaluation.ee/test-runner/__tests__/evaluation-metrics.ee.test.ts +++ b/packages/cli/src/evaluation.ee/test-runner/__tests__/evaluation-metrics.ee.test.ts @@ -40,4 +40,69 @@ describe('EvaluationMetrics', () => { expect(info).toHaveProperty('addedMetrics'); expect(info!.addedMetrics).toEqual({ metric1: 1, metric2: 0 }); }); + + describe('build/merge split', () => { + test('buildContribution is pure and does not mutate aggregator state', () => { + const metrics = new EvaluationMetrics(); + + const contribution = EvaluationMetrics.buildContribution({ metric1: 1, metric2: 0.5 }); + + expect(contribution.addedMetrics).toEqual({ metric1: 1, metric2: 0.5 }); + expect(metrics.getAggregatedMetrics()).toEqual({}); + }); + + test('mergeContribution applies a built contribution to the aggregator', () => { + const metrics = new EvaluationMetrics(); + + const contribution = EvaluationMetrics.buildContribution({ metric1: 0.5, metric2: 1 }); + metrics.mergeContribution(contribution); + + expect(metrics.getAggregatedMetrics()).toEqual({ metric1: 0.5, metric2: 1 }); + }); + + test('parallel build + reverse-order merge yields the same averages as sequential addResults', () => { + // `metric2` values are deliberately chosen to produce IEEE-754 + // reduction drift between forward and reverse summation: + // forward: 0 + 0.2 + 0.4 + 0.6 = 1.2000000000000002 / 4 = 0.30000000000000004 + // reverse: 0.6 + 0.4 + 0.2 + 0 = 1.2 / 4 = 0.3 + // `toBeCloseTo(..., 10)` is the assertion that does the work here: + // the contract is "averages stable to floating-point precision", + // not "byte-identical". `metric1` values are exact powers of two + // and behave the same way regardless of order — the keyset check + // still validates the structural contract for them. + const inputs = [ + { metric1: 1, metric2: 0 }, + { metric1: 0.5, metric2: 0.2 }, + { metric1: 0.25, metric2: 0.4 }, + { metric1: 0.125, metric2: 0.6 }, + ]; + + // Sequential baseline. + const sequential = new EvaluationMetrics(); + for (const input of inputs) { + sequential.addResults(input); + } + + // Parallel-build (no shared state) + reverse-order merge. + // `[...contributions].reverse()` keeps the source array intact. + const parallel = new EvaluationMetrics(); + const contributions = inputs.map((input) => EvaluationMetrics.buildContribution(input)); + for (const contribution of [...contributions].reverse()) { + parallel.mergeContribution(contribution); + } + + const sequentialResult = sequential.getAggregatedMetrics(); + const parallelResult = parallel.getAggregatedMetrics(); + expect(Object.keys(parallelResult).sort()).toEqual(Object.keys(sequentialResult).sort()); + for (const key of Object.keys(sequentialResult)) { + expect(parallelResult[key]).toBeCloseTo(sequentialResult[key], 15); + } + }); + + test('buildContribution throws on non-numeric values', () => { + expect(() => EvaluationMetrics.buildContribution({ metric1: 'not a number' })).toThrow( + 'INVALID_METRICS', + ); + }); + }); }); diff --git a/packages/cli/src/evaluation.ee/test-runner/__tests__/test-runner.service.ee.test.ts b/packages/cli/src/evaluation.ee/test-runner/__tests__/test-runner.service.ee.test.ts index a2783bb5cfa..e4539ff69fa 100644 --- a/packages/cli/src/evaluation.ee/test-runner/__tests__/test-runner.service.ee.test.ts +++ b/packages/cli/src/evaluation.ee/test-runner/__tests__/test-runner.service.ee.test.ts @@ -22,6 +22,7 @@ import path from 'path'; import { TestRunnerService } from '../test-runner.service.ee'; import type { ActiveExecutions } from '@/active-executions'; +import type { ConcurrencyControlService } from '@/concurrency/concurrency-control.service'; import { TestRunError } from '@/evaluation.ee/test-runner/errors.ee'; import { LoadNodesAndCredentials } from '@/load-nodes-and-credentials'; import type { Publisher } from '@/scaling/pubsub/publisher.service'; @@ -45,6 +46,7 @@ describe('TestRunnerService', () => { const executionsConfig = mockInstance(ExecutionsConfig, { mode: 'regular' }); const publisher = mock(); const instanceSettings = mock({ hostId: 'test-host-id', isMultiMain: false }); + const concurrencyControlService = mock(); let testRunnerService: TestRunnerService; mockInstance(LoadNodesAndCredentials, { @@ -65,6 +67,7 @@ describe('TestRunnerService', () => { mock(), publisher, instanceSettings, + concurrencyControlService, ); testRunRepository.createTestRun.mockResolvedValue(mock({ id: 'test-run-id' })); @@ -507,6 +510,7 @@ describe('TestRunnerService', () => { mock(), publisher, instanceSettings, + concurrencyControlService, ); process.env.OFFLOAD_MANUAL_EXECUTIONS_TO_WORKERS = 'true'; @@ -824,6 +828,7 @@ describe('TestRunnerService', () => { mock(), publisher, instanceSettings, + concurrencyControlService, ); }); @@ -1940,4 +1945,468 @@ describe('TestRunnerService', () => { }).not.toThrow(); }); }); + + describe('runTest - parallel execution', () => { + const TRIGGER_NODE_NAME = 'Dataset Trigger'; + const METRICS_NODE_NAME = 'Set Metrics'; + const USER = mock<{ id: string }>({ id: 'user-1' }); + const WORKFLOW_ID = 'wf-1'; + + // Builds a minimal workflow that passes validateWorkflowConfiguration. + // Using a plain object cast (not mock) so per-node + // boolean fields like `disabled` read as undefined instead of being + // auto-mocked as truthy functions by jest-mock-extended's deep proxy. + const buildWorkflow = (): IWorkflowBase => + ({ + id: WORKFLOW_ID, + name: 'Eval Workflow', + active: false, + nodes: [ + { + id: 'trigger', + name: TRIGGER_NODE_NAME, + type: EVALUATION_TRIGGER_NODE_TYPE, + typeVersion: 4.7, + position: [0, 0] as [number, number], + parameters: { + source: 'dataTable', + dataTableId: 'dt-1', + }, + }, + { + id: 'metrics', + name: METRICS_NODE_NAME, + type: EVALUATION_NODE_TYPE, + typeVersion: 4.7, + position: [200, 0] as [number, number], + parameters: { + operation: 'setMetrics', + metric: 'customMetrics', + metrics: { + assignments: [{ id: '1', name: 'score', value: 1 }], + }, + }, + }, + ], + connections: {}, + settings: {}, + }) as unknown as IWorkflowBase; + + // Dataset-trigger execution result containing N test rows. + const buildDatasetExecution = (rowCount: number): IRun => + ({ + data: { + resultData: { + runData: { + [TRIGGER_NODE_NAME]: [ + { + data: { + [NodeConnectionTypes.Main]: [ + Array.from({ length: rowCount }, (_, i) => ({ + json: { caseId: i, prompt: `prompt ${i}` }, + })), + ], + }, + }, + ], + }, + }, + }, + }) as unknown as IRun; + + // Per-case execution result with a single user-defined metric. + const buildCaseExecution = (score: number): IRun => + ({ + data: { + resultData: { + runData: { + [METRICS_NODE_NAME]: [ + { + data: { + [NodeConnectionTypes.Main]: [[{ json: { score } }]], + }, + }, + ], + }, + }, + }, + }) as unknown as IRun; + + // Wires repository and runner mocks for an N-case happy-path runTest. + // Returns a counter object the caller can read after `runTest` resolves. + const setupHappyPathMocks = (caseCount: number) => { + const workflow = buildWorkflow(); + workflowRepository.findById.mockResolvedValue(workflow as never); + + // Default: throttle resolves immediately (capacity always available). + // Tests that exercise abort-during-throttle override this. + concurrencyControlService.throttle.mockResolvedValue(undefined as never); + + testRunRepository.markAsRunning.mockResolvedValue(undefined as never); + testRunRepository.markAsCompleted.mockResolvedValue(undefined as never); + testRunRepository.markAsCancelled.mockResolvedValue(undefined as never); + testRunRepository.clearInstanceTracking.mockResolvedValue(undefined as never); + testRunRepository.isCancellationRequested.mockResolvedValue(false); + testCaseExecutionRepository.createTestCaseExecution.mockResolvedValue(undefined as never); + testCaseExecutionRepository.markAllPendingAsCancelled.mockResolvedValue(undefined as never); + // `manager` is a TypeORM EntityManager not auto-deep-mocked by mock(). + // Provide a transaction stub that just invokes the callback so cancel + // paths run end-to-end. + Object.assign(testRunRepository, { + manager: { + transaction: jest + .fn() + .mockImplementation(async (cb: (trx: unknown) => Promise) => await cb({})), + }, + }); + + let runCallIndex = 0; + const inFlightTracker = { inFlight: 0, max: 0, perCaseStarted: 0 }; + + workflowRunner.run.mockImplementation(async () => { + const id = runCallIndex === 0 ? 'dataset-exec' : `case-exec-${runCallIndex}`; + runCallIndex++; + return id; + }); + + activeExecutions.getPostExecutePromise.mockImplementation(async (executionId) => { + if (executionId === 'dataset-exec') { + return buildDatasetExecution(caseCount); + } + inFlightTracker.inFlight++; + inFlightTracker.perCaseStarted++; + inFlightTracker.max = Math.max(inFlightTracker.max, inFlightTracker.inFlight); + // Yield to the event loop so other queued tasks observably overlap. + await new Promise((resolve) => setTimeout(resolve, 5)); + inFlightTracker.inFlight--; + const score = parseInt(executionId.replace('case-exec-', ''), 10) / 10; + return buildCaseExecution(score); + }); + + return { workflow, inFlightTracker }; + }; + + test('concurrency=1 runs cases sequentially (max in-flight = 1)', async () => { + const { inFlightTracker } = setupHappyPathMocks(4); + + await testRunnerService.runTest(USER as never, WORKFLOW_ID, 1); + + expect(inFlightTracker.perCaseStarted).toBe(4); + expect(inFlightTracker.max).toBe(1); + // 1 dataset trigger + 4 test cases. + expect(workflowRunner.run).toHaveBeenCalledTimes(5); + expect(testRunRepository.markAsCompleted).toHaveBeenCalledTimes(1); + expect(testRunRepository.markAsCancelled).not.toHaveBeenCalled(); + }); + + test('concurrency=3 fans out cases in parallel (max in-flight > 1)', async () => { + const { inFlightTracker } = setupHappyPathMocks(6); + + await testRunnerService.runTest(USER as never, WORKFLOW_ID, 3); + + expect(inFlightTracker.perCaseStarted).toBe(6); + expect(inFlightTracker.max).toBeGreaterThan(1); + expect(inFlightTracker.max).toBeLessThanOrEqual(3); + expect(workflowRunner.run).toHaveBeenCalledTimes(7); + expect(testRunRepository.markAsCompleted).toHaveBeenCalledTimes(1); + }); + + test('concurrency clamped 1-10 defensively (above-bound input does not exceed cap)', async () => { + const { inFlightTracker } = setupHappyPathMocks(12); + + await testRunnerService.runTest(USER as never, WORKFLOW_ID, 99); + + expect(inFlightTracker.perCaseStarted).toBe(12); + expect(inFlightTracker.max).toBeLessThanOrEqual(10); + }); + + test('aggregate metrics produce the same average regardless of concurrency', async () => { + setupHappyPathMocks(5); + await testRunnerService.runTest(USER as never, WORKFLOW_ID, 1); + const sequentialMetrics = testRunRepository.markAsCompleted.mock.calls[0][1]; + + // `clearAllMocks` resets call history but not implementations. The + // `createTestRun` stub is set in the outer `beforeEach`, so it needs + // re-stubbing here. `setupHappyPathMocks` re-wires everything else. + jest.clearAllMocks(); + testRunRepository.createTestRun.mockResolvedValue(mock({ id: 'test-run-id' })); + setupHappyPathMocks(5); + await testRunnerService.runTest(USER as never, WORKFLOW_ID, 4); + const parallelMetrics = testRunRepository.markAsCompleted.mock.calls[0][1]; + + expect(Object.keys(parallelMetrics ?? {}).sort()).toEqual( + Object.keys(sequentialMetrics ?? {}).sort(), + ); + for (const key of Object.keys(sequentialMetrics ?? {})) { + expect((parallelMetrics as Record)[key]).toBeCloseTo( + (sequentialMetrics as Record)[key], + 15, + ); + } + }); + + test('per-case error does not stop other cases from completing', async () => { + setupHappyPathMocks(4); + + // Override per-case mock so case 2 throws. + activeExecutions.getPostExecutePromise.mockImplementation(async (executionId) => { + if (executionId === 'dataset-exec') { + return buildDatasetExecution(4); + } + if (executionId === 'case-exec-2') { + throw new Error('synthetic failure'); + } + return buildCaseExecution(0.5); + }); + + await testRunnerService.runTest(USER as never, WORKFLOW_ID, 2); + + // 4 test-case executions attempted; 1 errored, 3 succeeded. + const createCalls = testCaseExecutionRepository.createTestCaseExecution.mock.calls; + const errorRows = createCalls.filter(([row]) => row.status === 'error'); + const successRows = createCalls.filter(([row]) => row.status === 'success'); + expect(errorRows).toHaveLength(1); + expect(successRows).toHaveLength(3); + expect(testRunRepository.markAsCompleted).toHaveBeenCalledTimes(1); + }); + + test('throttle is called once per case and release is called once per case', async () => { + setupHappyPathMocks(4); + + await testRunnerService.runTest(USER as never, WORKFLOW_ID, 2); + + expect(concurrencyControlService.throttle).toHaveBeenCalledTimes(4); + expect(concurrencyControlService.release).toHaveBeenCalledTimes(4); + expect(concurrencyControlService.throttle).toHaveBeenCalledWith({ + mode: 'evaluation', + executionId: expect.stringContaining('test-run-id-case-'), + }); + expect(concurrencyControlService.release).toHaveBeenCalledWith({ mode: 'evaluation' }); + }); + + test('release is called even when runTestCase throws', async () => { + // `setupHappyPathMocks` wires the dataset-trigger execution; the + // override below replaces only the per-case execution path so each + // case throws. The `dataset-exec` branch is preserved manually here + // to keep the dataset trigger intact. + setupHappyPathMocks(3); + + activeExecutions.getPostExecutePromise.mockImplementation(async (executionId) => { + if (executionId === 'dataset-exec') { + return { + data: { + resultData: { + runData: { + [TRIGGER_NODE_NAME]: [ + { + data: { + [NodeConnectionTypes.Main]: [ + [{ json: { id: 0 } }, { json: { id: 1 } }, { json: { id: 2 } }], + ], + }, + }, + ], + }, + }, + }, + } as unknown as IRun; + } + throw new Error('synthetic per-case failure'); + }); + + await testRunnerService.runTest(USER as never, WORKFLOW_ID, 2); + + expect(concurrencyControlService.throttle).toHaveBeenCalledTimes(3); + // release fires in the finally block — must run even though every + // case threw. + expect(concurrencyControlService.release).toHaveBeenCalledTimes(3); + }); + + test('telemetry payload includes concurrency, parallel_enabled, concurrency_limited_by_config, flag_enabled_for_user', async () => { + setupHappyPathMocks(2); + + await testRunnerService.runTest(USER as never, WORKFLOW_ID, 4, true); + + const trackCalls = telemetry.track.mock.calls.filter( + ([eventName]) => eventName === 'Test run finished', + ); + expect(trackCalls).toHaveLength(1); + const payload = trackCalls[0][1] as Record; + expect(payload).toEqual( + expect.objectContaining({ + concurrency: 4, + parallel_enabled: true, + concurrency_limited_by_config: false, + flag_enabled_for_user: true, + }), + ); + }); + + test('flag_enabled_for_user defaults to false when not passed', async () => { + setupHappyPathMocks(2); + + await testRunnerService.runTest(USER as never, WORKFLOW_ID, 1); + + const payload = telemetry.track.mock.calls.find( + ([eventName]) => eventName === 'Test run finished', + )?.[1] as Record; + expect(payload.flag_enabled_for_user).toBe(false); + }); + + test('telemetry parallel_enabled is false for sequential runs', async () => { + setupHappyPathMocks(2); + + await testRunnerService.runTest(USER as never, WORKFLOW_ID, 1); + + const trackCalls = telemetry.track.mock.calls.filter( + ([eventName]) => eventName === 'Test run finished', + ); + expect(trackCalls).toHaveLength(1); + const payload = trackCalls[0][1] as Record; + expect(payload).toEqual( + expect.objectContaining({ + concurrency: 1, + parallel_enabled: false, + concurrency_limited_by_config: false, + }), + ); + }); + + test('telemetry payload reports realised fan-out (cases_started, peak_in_flight)', async () => { + const { inFlightTracker } = setupHappyPathMocks(6); + + await testRunnerService.runTest(USER as never, WORKFLOW_ID, 3, true); + + const payload = telemetry.track.mock.calls.find( + ([eventName]) => eventName === 'Test run finished', + )?.[1] as Record; + expect(payload.cases_started).toBe(6); + // Should match what the test harness independently observed — + // proves the telemetry stat is the same number a watcher would see. + expect(payload.peak_in_flight).toBe(inFlightTracker.max); + expect(payload.peak_in_flight).toBeGreaterThan(1); + expect(payload.peak_in_flight).toBeLessThanOrEqual(3); + }); + + test('evaluationLimit clamps requested concurrency and flags concurrency_limited_by_config', async () => { + const cappedConfig = mockInstance(ExecutionsConfig, { + mode: 'regular', + concurrency: { productionLimit: -1, evaluationLimit: 2 } as never, + }); + const cappedService = new TestRunnerService( + logger, + telemetry, + workflowRepository, + workflowRunner, + activeExecutions, + testRunRepository, + testCaseExecutionRepository, + errorReporter, + cappedConfig, + mock(), + publisher, + instanceSettings, + concurrencyControlService, + ); + + const { inFlightTracker } = setupHappyPathMocks(6); + + await cappedService.runTest(USER as never, WORKFLOW_ID, 5); + + expect(inFlightTracker.max).toBeLessThanOrEqual(2); + const payload = telemetry.track.mock.calls.find( + ([eventName]) => eventName === 'Test run finished', + )?.[1] as Record; + expect(payload).toEqual( + expect.objectContaining({ + concurrency: 2, + parallel_enabled: true, + concurrency_limited_by_config: true, + }), + ); + }); + + test('abort during throttle wait evicts the queue entry and short-circuits without an UNKNOWN_ERROR row', async () => { + setupHappyPathMocks(3); + + // Hold throttle indefinitely so all per-case tasks are stuck in the + // queue when we fire the abort. + concurrencyControlService.throttle.mockImplementation( + async () => + await new Promise(() => { + /* never resolves */ + }), + ); + + // Kick off the run, then abort while cases are queued in throttle. + const runPromise = testRunnerService.runTest(USER as never, WORKFLOW_ID, 3); + await new Promise((resolve) => setTimeout(resolve, 5)); + + // Reach into the service's abort controllers map and trip it. + // `cancelTestRunLocally` is the path the public cancel API takes. + const cancelled = ( + testRunnerService as never as { + cancelTestRunLocally: (id: string) => boolean; + } + ).cancelTestRunLocally('test-run-id'); + expect(cancelled).toBe(true); + + await runPromise; + + // Each queued case should have been evicted via remove(). + expect(concurrencyControlService.remove).toHaveBeenCalledWith({ + mode: 'evaluation', + executionId: expect.stringContaining('test-run-id-case-'), + }); + + // And no test-case row should have been written for the evicted + // cases — they short-circuit before touching the DB. The legacy + // path would have produced UNKNOWN_ERROR rows here. + const errorRows = testCaseExecutionRepository.createTestCaseExecution.mock.calls.filter( + ([row]) => row.errorCode === 'UNKNOWN_ERROR', + ); + expect(errorRows).toHaveLength(0); + + // Run is marked cancelled, not completed. + expect(testRunRepository.markAsCancelled).toHaveBeenCalled(); + expect(testRunRepository.markAsCompleted).not.toHaveBeenCalled(); + }); + + test('multi-main DB cancel flag flipped mid-run aborts the controller', async () => { + const multiMainInstance = mock({ + hostId: 'main-a', + isMultiMain: true, + }); + const multiMainService = new TestRunnerService( + logger, + telemetry, + workflowRepository, + workflowRunner, + activeExecutions, + testRunRepository, + testCaseExecutionRepository, + errorReporter, + executionsConfig, + mock(), + publisher, + multiMainInstance, + concurrencyControlService, + ); + + setupHappyPathMocks(4); + + // Flip the cancellation flag after the first case sees it as false. + let pollCount = 0; + testRunRepository.isCancellationRequested.mockImplementation(async () => { + pollCount++; + return pollCount > 1; + }); + + await multiMainService.runTest(USER as never, WORKFLOW_ID, 1); + + expect(testRunRepository.isCancellationRequested).toHaveBeenCalled(); + expect(testRunRepository.markAsCancelled).toHaveBeenCalled(); + expect(testRunRepository.markAsCompleted).not.toHaveBeenCalled(); + }); + }); }); diff --git a/packages/cli/src/evaluation.ee/test-runner/evaluation-metrics.ee.ts b/packages/cli/src/evaluation.ee/test-runner/evaluation-metrics.ee.ts index 7375f9967be..453c9cc68c8 100644 --- a/packages/cli/src/evaluation.ee/test-runner/evaluation-metrics.ee.ts +++ b/packages/cli/src/evaluation.ee/test-runner/evaluation-metrics.ee.ts @@ -7,27 +7,37 @@ export interface EvaluationMetricsAddResultsInfo { incorrectTypeMetrics: Set; } +/** + * A single test case's contribution to the aggregated metrics. Built off the + * hot path so multiple cases can be processed in parallel; merged sequentially + * once all parallel work resolves. + * + * Note: `incorrectTypeMetrics` is deliberately not on this type. The legacy + * `addResults` shape includes it for backwards-compat reasons, but builds + * throw on the first non-numeric value, so the set could only ever hold a + * single name and is unreachable from the caller anyway. + */ +export interface MetricContribution { + addedMetrics: Record; +} + export class EvaluationMetrics { private readonly rawMetricsByName = new Map(); - addResults(result: IDataObject): EvaluationMetricsAddResultsInfo { - const addResultsInfo: EvaluationMetricsAddResultsInfo = { - addedMetrics: {}, - incorrectTypeMetrics: new Set(), - }; + /** + * Pure: builds a contribution from a single test case's result without + * mutating any shared state. Safe to call from parallel fan-out tasks. + * + * Throws `TestCaseExecutionError('INVALID_METRICS')` on non-numeric values, + * matching the historical `addResults` semantics. + */ + static buildContribution(result: IDataObject): MetricContribution { + const addedMetrics: Record = {}; for (const [metricName, metricValue] of Object.entries(result)) { if (typeof metricValue === 'number') { - addResultsInfo.addedMetrics[metricName] = metricValue; - - // Initialize the array if this is the first time we see this metric - if (!this.rawMetricsByName.has(metricName)) { - this.rawMetricsByName.set(metricName, []); - } - - this.rawMetricsByName.get(metricName)!.push(metricValue); + addedMetrics[metricName] = metricValue; } else { - addResultsInfo.incorrectTypeMetrics.add(metricName); throw new TestCaseExecutionError('INVALID_METRICS', { metricName, metricValue, @@ -35,7 +45,41 @@ export class EvaluationMetrics { } } - return addResultsInfo; + return { addedMetrics }; + } + + /** + * Single-threaded: applies a contribution to the aggregator. Call once per + * contribution after all `buildContribution` calls have resolved. + */ + mergeContribution(contribution: MetricContribution): void { + for (const [metricName, metricValue] of Object.entries(contribution.addedMetrics)) { + let bucket = this.rawMetricsByName.get(metricName); + if (!bucket) { + bucket = []; + this.rawMetricsByName.set(metricName, bucket); + } + bucket.push(metricValue); + } + } + + /** + * Backwards-compatible wrapper around the build/merge split. The runner + * still uses this in commit 2; the parallel fan-out (a later commit) calls + * `buildContribution` from inside the per-case task and `mergeContribution` + * once after `Promise.all`. + * + * The returned `incorrectTypeMetrics` is empty by construction (the build + * step throws on the first non-numeric value, so there is no caller-visible + * post-throw state). Kept for API compatibility with pre-split consumers. + */ + addResults(result: IDataObject): EvaluationMetricsAddResultsInfo { + const contribution = EvaluationMetrics.buildContribution(result); + this.mergeContribution(contribution); + return { + addedMetrics: contribution.addedMetrics, + incorrectTypeMetrics: new Set(), + }; } getAggregatedMetrics() { diff --git a/packages/cli/src/evaluation.ee/test-runner/test-runner.service.ee.ts b/packages/cli/src/evaluation.ee/test-runner/test-runner.service.ee.ts index 17802cb2b80..18f63b5901e 100644 --- a/packages/cli/src/evaluation.ee/test-runner/test-runner.service.ee.ts +++ b/packages/cli/src/evaluation.ee/test-runner/test-runner.service.ee.ts @@ -6,6 +6,7 @@ import { OnPubSubEvent } from '@n8n/decorators'; import { Service } from '@n8n/di'; import { ErrorReporter, InstanceSettings } from 'n8n-core'; +import { ConcurrencyControlService } from '@/concurrency/concurrency-control.service'; import { Publisher } from '@/scaling/pubsub/publisher.service'; import { EVALUATION_NODE_TYPE, @@ -28,6 +29,7 @@ import type { } from 'n8n-workflow'; import assert from 'node:assert'; import type { JsonObject } from 'openid-client'; +import pLimit from 'p-limit'; import { ActiveExecutions } from '@/active-executions'; import { EventService } from '@/events/event.service'; @@ -39,7 +41,7 @@ import { import { Telemetry } from '@/telemetry'; import { WorkflowRunner } from '@/workflow-runner'; -import { EvaluationMetrics } from './evaluation-metrics.ee'; +import { EvaluationMetrics, type MetricContribution } from './evaluation-metrics.ee'; export interface TestRunMetadata { testRunId: string; @@ -77,6 +79,7 @@ export class TestRunnerService { private readonly eventService: EventService, private readonly publisher: Publisher, private readonly instanceSettings: InstanceSettings, + private readonly concurrencyControlService: ConcurrencyControlService, ) {} /** @@ -500,10 +503,37 @@ export class TestRunnerService { } /** - * Creates a new test run for the given workflow + * Creates a new test run for the given workflow. + * + * `concurrency` is the requested number of test cases to run in parallel. + * The effective value is `min(user_request, 10, evaluationLimit)`: + * - Clamped 1–10 as a defensive UX guardrail (the controller already + * validates this via zod, but direct service callers must not exceed + * it either). + * - Further clamped to `evaluationLimit` (`N8N_CONCURRENCY_EVALUATION_LIMIT`) + * when an admin has set a positive cap. `concurrency_limited_by_config` + * is recorded in telemetry when this kicks in. + * + * `concurrency = 1` reproduces the legacy sequential behaviour exactly. */ - async runTest(user: User, workflowId: string): Promise { - this.logger.debug('Starting new test run', { workflowId }); + async runTest( + user: User, + workflowId: string, + concurrency: number = 1, + flagEnabledForUser: boolean = false, + ): Promise { + const requestedConcurrency = Math.max(1, Math.min(10, Math.floor(concurrency))); + const evaluationLimit = this.executionsConfig.concurrency.evaluationLimit; + const concurrencyLimitedByConfig = + evaluationLimit > 0 && requestedConcurrency > evaluationLimit; + const effectiveConcurrency = concurrencyLimitedByConfig + ? evaluationLimit + : requestedConcurrency; + + this.logger.debug( + `[Eval] runTest called: requestedConcurrency=${requestedConcurrency} effectiveConcurrency=${effectiveConcurrency} evaluationLimit=${evaluationLimit} flagEnabledForUser=${flagEnabledForUser}`, + { workflowId }, + ); const workflow = await this.workflowRepository.findById(workflowId); assert(workflow, 'Workflow not found'); @@ -524,6 +554,17 @@ export class TestRunnerService { metric_count: 0, error_message: '', duration: 0, + concurrency: effectiveConcurrency, + parallel_enabled: effectiveConcurrency > 1, + concurrency_limited_by_config: concurrencyLimitedByConfig, + flag_enabled_for_user: flagEnabledForUser, + // Realised parallelism observed at runtime — `cases_started` counts + // callbacks that actually began (post-throttle, pre-abort), and + // `peak_in_flight` is the high-water mark for in-flight cases. + // Updated in lockstep with fanOutMetrics inside the per-case callback; + // stays at 0 if the run aborts before the fan-out begins. + cases_started: 0, + peak_in_flight: 0, }; // 0.1 Initialize AbortController @@ -577,167 +618,306 @@ export class TestRunnerService { // 2. Run over all the test cases /// - for (const testCase of testCases) { - if (abortSignal.aborted) { - telemetryMeta.status = 'cancelled'; - this.logger.debug('Test run was cancelled', { - workflowId, - }); - break; - } + // pLimit(N) governs how many per-case tasks may be in flight at + // once. With concurrency=1 the per-case callback runs in serial, + // reproducing the legacy `for…of` loop exactly. Each callback + // returns the contributions it built; merging happens once on the + // main thread after Promise.all so EvaluationMetrics state is + // never touched concurrently. + // + // `telemetryMeta.*++` increments inside the callback are safe under + // JS's single-threaded event loop: `++` is synchronous, and there + // is no `await` between the read and the write of the counter. + const limit = pLimit(effectiveConcurrency); - // Check database cancellation flag (fallback for multi-main mode only) - // This ensures cancellation works even if the pub/sub message didn't reach this instance - // In single-main mode, the local abort controller check above is sufficient - if ( - this.instanceSettings.isMultiMain && - (await this.testRunRepository.isCancellationRequested(testRun.id)) - ) { - this.logger.debug('Test run cancellation requested via database flag', { - workflowId, - testRunId: testRun.id, - }); - abortController.abort(); - telemetryMeta.status = 'cancelled'; - break; - } + // Visibility for parallel fan-out. The `inFlight` counter is mutated + // from per-case callbacks but the increments are safe — JS's single- + // threaded event loop guarantees no interleaving between read and + // write of `++` within a sync block. + const fanOutMetrics = { inFlight: 0, peakInFlight: 0, casesStarted: 0 }; + this.logger.debug( + `[Eval] Fan-out begin: cases=${testCases.length} concurrency=${effectiveConcurrency}`, + { testRunId: testRun.id, workflowId }, + ); - this.logger.debug('Running test case'); - const runAt = new Date(); + const contributionResults = await Promise.all( + testCases.map( + async (testCase, caseIndex) => + await limit(async (): Promise => { + if (abortSignal.aborted) { + return []; + } - try { - const testCaseMetadata = { - ...testRunMetadata, - }; + // Multi-main DB cancellation poll, run per case as a defensive + // fallback for the rare case a foreign main flips the cancel + // flag but the pubsub broadcast doesn't reach this instance. + // Cheap (~1ms indexed PK lookup); kept as-is rather than + // optimised to once-per-run to preserve the existing safety net. + if ( + this.instanceSettings.isMultiMain && + (await this.testRunRepository.isCancellationRequested(testRun.id)) + ) { + this.logger.debug('Test run cancellation requested via database flag', { + workflowId, + testRunId: testRun.id, + }); + abortController.abort(); + return []; + } - // Run the test case and wait for it to finish - const testCaseResult = await this.runTestCase( - workflow, - testCaseMetadata, - testCase, - abortSignal, - ); - assert(testCaseResult); + // Layer onto the existing instance-wide concurrency control. The + // service is a no-op in queue mode (BullMQ governs there) and when + // `evaluationLimit` is unset (-1). pLimit and the eval queue cap + // the in-flight count at *the same number* by design — pLimit is + // per-run, the queue is shared across all test runs from all users + // on the instance, so they're complementary, not redundant. + // + // Abort-aware acquisition: if Stop is clicked while we're queued + // behind another evaluation's capacity, we evict ourselves from the + // queue so the slot returns to circulation and our task short- + // circuits promptly instead of waiting for an unrelated run to + // release. Without this, queued cases would block until they drained + // through the queue — and then `runTestCase` would return undefined + // (abort observed at its top), tripping the assert below and landing + // a misleading UNKNOWN_ERROR test-case row. + const caseTrackingId = `${testRun.id}-case-${caseIndex}`; + let abortHandler: (() => void) | undefined; + let throttleAcquired = false; + const abortRace = new Promise<'aborted'>((resolve) => { + abortHandler = () => resolve('aborted'); + abortSignal.addEventListener('abort', abortHandler, { once: true }); + }); + const acquireRace = this.concurrencyControlService + .throttle({ mode: 'evaluation', executionId: caseTrackingId }) + .then(() => { + throttleAcquired = true; + return 'acquired' as const; + }); + const acquired = await Promise.race([acquireRace, abortRace]); - const { executionId: testCaseExecutionId, executionData: testCaseExecution } = - testCaseResult; + if (acquired === 'aborted') { + // Two abort sub-cases handled defensively, distinguished by + // whether throttle's `.then` microtask managed to set + // `throttleAcquired` before abort won the race: + // + // 1. throttleAcquired = true — the eval queue had immediate + // capacity (no queue push, slot synchronously consumed), + // and the `.then` microtask fired before abort. Race + // *should* have picked 'acquired' in this ordering, but + // handle defensively against scheduler quirks: release + // the slot back to the queue. + // 2. throttleAcquired = false — we were either still queued + // (capacity wasn't available) or the immediate-acquire's + // microtask hadn't fired yet. Either way, remove() splices + // a queued entry (frees the slot via internal capacity++) + // and is a no-op for non-queued entries. The unawaited + // acquireRace becomes garbage. + if (throttleAcquired) { + this.concurrencyControlService.release({ mode: 'evaluation' }); + } else { + this.concurrencyControlService.remove({ + mode: 'evaluation', + executionId: caseTrackingId, + }); + } + return []; + } - assert(testCaseExecution); - assert(testCaseExecutionId); + if (abortHandler) { + abortSignal.removeEventListener('abort', abortHandler); + } - this.logger.debug('Test case execution finished'); + // Narrow window: abort could fire between `throttle` resolving and + // here. We have the capacity slot; release it and bail. + if (abortSignal.aborted) { + this.concurrencyControlService.release({ mode: 'evaluation' }); + return []; + } - // In case of a permission check issue, the test case execution will be undefined. - // If that happens, or if the test case execution produced an error, mark the test case as failed. - if (!testCaseExecution || testCaseExecution.data.resultData.error) { - // Save the failed test case execution in DB - await this.testCaseExecutionRepository.createTestCaseExecution({ - executionId: testCaseExecutionId, - testRun: { - id: testRun.id, - }, - status: 'error', - errorCode: 'FAILED_TO_EXECUTE_WORKFLOW', - metrics: {}, - }); - telemetryMeta.errored_test_case_count++; - continue; - } - const completedAt = new Date(); + // In-flight tracking — increment as we leave the throttle and + // decrement in the outer finally. The peak counter shows whether + // the runner actually fanned out concurrently. Mirror the two + // summary stats into telemetryMeta so they survive into the + // `Test run finished` event even if the run errors mid-fan-out. + fanOutMetrics.inFlight += 1; + fanOutMetrics.casesStarted += 1; + telemetryMeta.cases_started = fanOutMetrics.casesStarted; + if (fanOutMetrics.inFlight > fanOutMetrics.peakInFlight) { + fanOutMetrics.peakInFlight = fanOutMetrics.inFlight; + telemetryMeta.peak_in_flight = fanOutMetrics.peakInFlight; + } + this.logger.debug( + `[Eval] Case started: case=${caseIndex} inFlight=${fanOutMetrics.inFlight}/${effectiveConcurrency} peak=${fanOutMetrics.peakInFlight}`, + { testRunId: testRun.id }, + ); - // Collect common metrics - const { addedMetrics: addedPredefinedMetrics } = metrics.addResults( - this.extractPredefinedMetrics(testCaseExecution), - ); - this.logger.debug('Test case common metrics extracted', addedPredefinedMetrics); + const runAt = new Date(); - // Collect user-defined metrics - const { addedMetrics: addedUserDefinedMetrics } = metrics.addResults( - this.extractUserDefinedMetrics(testCaseExecution, workflow), - ); + try { + try { + const testCaseMetadata = { ...testRunMetadata }; - if (Object.keys(addedUserDefinedMetrics).length === 0) { - await this.testCaseExecutionRepository.createTestCaseExecution({ - executionId: testCaseExecutionId, - testRun: { - id: testRun.id, - }, - runAt, - completedAt, - status: 'error', - errorCode: 'NO_METRICS_COLLECTED', - }); - telemetryMeta.errored_test_case_count++; - } else { - const combinedMetrics = { - ...addedUserDefinedMetrics, - ...addedPredefinedMetrics, - }; + const testCaseResult = await this.runTestCase( + workflow, + testCaseMetadata, + testCase, + abortSignal, + ); - const inputs = this.getEvaluationData(testCaseExecution, workflow, 'setInputs'); - const outputs = this.getEvaluationData(testCaseExecution, workflow, 'setOutputs'); + // `runTestCase` returns undefined only when `abortSignal.aborted` + // is true at entry (see method body). Skip silently so the outer + // reconciliation can mark the run as cancelled — landing an + // UNKNOWN_ERROR test-case row here would be misleading. Asserting + // the abort invariant catches future regressions where the + // undefined return path widens. + if (!testCaseResult) { + assert( + abortSignal.aborted, + 'runTestCase returned undefined without abort being set', + ); + return []; + } - this.logger.debug( - 'Test case metrics extracted (user-defined)', - addedUserDefinedMetrics, - ); + const { executionId: testCaseExecutionId, executionData: testCaseExecution } = + testCaseResult; - // Create a new test case execution in DB - await this.testCaseExecutionRepository.createTestCaseExecution({ - executionId: testCaseExecutionId, - testRun: { - id: testRun.id, - }, - runAt, - completedAt, - status: 'success', - metrics: combinedMetrics, - inputs, - outputs, - }); - } - } catch (e) { - const completedAt = new Date(); - // FIXME: this is a temporary log - this.logger.error('Test case execution failed', { - workflowId, - testRunId: testRun.id, - error: e, - }); + assert(testCaseExecution); + assert(testCaseExecutionId); - telemetryMeta.errored_test_case_count++; + this.logger.debug('Test case execution finished'); - // In case of an unexpected error save it as failed test case execution and continue with the next test case - if (e instanceof TestCaseExecutionError) { - await this.testCaseExecutionRepository.createTestCaseExecution({ - testRun: { - id: testRun.id, - }, - runAt, - completedAt, - status: 'error', - errorCode: e.code, - errorDetails: e.extra as IDataObject, - }); - } else { - await this.testCaseExecutionRepository.createTestCaseExecution({ - testRun: { - id: testRun.id, - }, - runAt, - completedAt, - status: 'error', - errorCode: 'UNKNOWN_ERROR', - }); + if (!testCaseExecution || testCaseExecution.data.resultData.error) { + await this.testCaseExecutionRepository.createTestCaseExecution({ + executionId: testCaseExecutionId, + testRun: { id: testRun.id }, + status: 'error', + errorCode: 'FAILED_TO_EXECUTE_WORKFLOW', + metrics: {}, + }); + telemetryMeta.errored_test_case_count++; + return []; + } + const completedAt = new Date(); - // Report unexpected errors - this.errorReporter.error(e); - } + const predefinedContribution = EvaluationMetrics.buildContribution( + this.extractPredefinedMetrics(testCaseExecution), + ); + this.logger.debug( + 'Test case common metrics extracted', + predefinedContribution.addedMetrics, + ); + + const userDefinedContribution = EvaluationMetrics.buildContribution( + this.extractUserDefinedMetrics(testCaseExecution, workflow), + ); + + if (Object.keys(userDefinedContribution.addedMetrics).length === 0) { + await this.testCaseExecutionRepository.createTestCaseExecution({ + executionId: testCaseExecutionId, + testRun: { id: testRun.id }, + runAt, + completedAt, + status: 'error', + errorCode: 'NO_METRICS_COLLECTED', + }); + telemetryMeta.errored_test_case_count++; + // Predefined metrics still merge — the case ran, just had no user metrics. + return [predefinedContribution]; + } + + const combinedMetrics = { + ...userDefinedContribution.addedMetrics, + ...predefinedContribution.addedMetrics, + }; + + const inputs = this.getEvaluationData(testCaseExecution, workflow, 'setInputs'); + const outputs = this.getEvaluationData(testCaseExecution, workflow, 'setOutputs'); + + this.logger.debug( + 'Test case metrics extracted (user-defined)', + userDefinedContribution.addedMetrics, + ); + + await this.testCaseExecutionRepository.createTestCaseExecution({ + executionId: testCaseExecutionId, + testRun: { id: testRun.id }, + runAt, + completedAt, + status: 'success', + metrics: combinedMetrics, + inputs, + outputs, + }); + + return [predefinedContribution, userDefinedContribution]; + } catch (e) { + const completedAt = new Date(); + this.logger.error('[Eval] Test case execution failed', { + workflowId, + testRunId: testRun.id, + caseIndex, + errorName: e instanceof Error ? e.name : 'Unknown', + errorMessage: e instanceof Error ? e.message : String(e), + errorStack: e instanceof Error ? e.stack : undefined, + }); + + telemetryMeta.errored_test_case_count++; + + if (e instanceof TestCaseExecutionError) { + await this.testCaseExecutionRepository.createTestCaseExecution({ + testRun: { id: testRun.id }, + runAt, + completedAt, + status: 'error', + errorCode: e.code, + errorDetails: e.extra as IDataObject, + }); + } else { + await this.testCaseExecutionRepository.createTestCaseExecution({ + testRun: { id: testRun.id }, + runAt, + completedAt, + status: 'error', + errorCode: 'UNKNOWN_ERROR', + }); + this.errorReporter.error(e); + } + return []; + } + } finally { + // Always release capacity, even when runTestCase throws. + // The synthetic id is irrelevant — release dequeues by mode. + this.concurrencyControlService.release({ mode: 'evaluation' }); + fanOutMetrics.inFlight -= 1; + this.logger.debug( + `[Eval] Case finished: case=${caseIndex} inFlight=${fanOutMetrics.inFlight}/${effectiveConcurrency}`, + { testRunId: testRun.id }, + ); + } + }), + ), + ); + + this.logger.debug( + `[Eval] Fan-out complete: cases=${fanOutMetrics.casesStarted} peakInFlight=${fanOutMetrics.peakInFlight}/${effectiveConcurrency}`, + { testRunId: testRun.id }, + ); + + // Single-threaded merge step. Order is irrelevant for averages + // (within IEEE-754 precision; see evaluation-metrics tests). + for (const caseContributions of contributionResults) { + for (const contribution of caseContributions) { + metrics.mergeContribution(contribution); } } - // Mark the test run as completed or cancelled + // Mark the test run as completed or cancelled. The multi-main DB + // poll inside each per-case callback can flip `abortController.abort()` + // on its own; this branch is the only place telemetry status is set + // for cancellations, so both the user-initiated and poll-initiated + // paths converge here. if (abortSignal.aborted) { + this.logger.debug('Test run was cancelled', { workflowId }); await dbManager.transaction(async (trx) => { await this.testRunRepository.markAsCancelled(testRun.id, trx); await this.testCaseExecutionRepository.markAllPendingAsCancelled(testRun.id, trx); diff --git a/packages/cli/src/evaluation.ee/test-runs.controller.ee.ts b/packages/cli/src/evaluation.ee/test-runs.controller.ee.ts index 6f1eeaf54c3..e5082b63237 100644 --- a/packages/cli/src/evaluation.ee/test-runs.controller.ee.ts +++ b/packages/cli/src/evaluation.ee/test-runs.controller.ee.ts @@ -1,6 +1,8 @@ +import { EVAL_PARALLEL_EXECUTION_FLAG, StartTestRunRequestDto } from '@n8n/api-types'; +import { Logger } from '@n8n/backend-common'; import { TestCaseExecutionRepository, TestRunRepository } from '@n8n/db'; import type { User } from '@n8n/db'; -import { Delete, Get, Post, RestController } from '@n8n/decorators'; +import { Body, Delete, Get, Post, RestController } from '@n8n/decorators'; import express from 'express'; import { UnexpectedError } from 'n8n-workflow'; @@ -9,6 +11,7 @@ import { NotFoundError } from '@/errors/response-errors/not-found.error'; import { TestRunnerService } from '@/evaluation.ee/test-runner/test-runner.service.ee'; import { TestRunsRequest } from '@/evaluation.ee/test-runs.types.ee'; import { listQueryMiddleware } from '@/middlewares'; +import { PostHogClient } from '@/posthog'; import { Telemetry } from '@/telemetry'; import { WorkflowFinderService } from '@/workflows/workflow-finder.service'; @@ -20,8 +23,28 @@ export class TestRunsController { private readonly testCaseExecutionRepository: TestCaseExecutionRepository, private readonly testRunnerService: TestRunnerService, private readonly telemetry: Telemetry, + private readonly postHogClient: PostHogClient, + private readonly logger: Logger, ) {} + /** + * Resolves the parallel-execution rollout flag for a user, defaulting to + * `false` (sequential) on any PostHog failure. Fail-open semantics: a + * PostHog outage degrades the rollout cohort to the legacy sequential + * behaviour rather than 500ing the test-run start. + */ + private async isParallelExecutionFlagEnabled(user: User): Promise { + try { + const flags = await this.postHogClient.getFeatureFlags(user); + return flags?.[EVAL_PARALLEL_EXECUTION_FLAG] === true; + } catch (error) { + this.logger.warn('Failed to resolve eval parallel-execution flag', { + error: error instanceof Error ? error.message : String(error), + }); + return false; + } + } + private async assertUserHasAccessToWorkflow(workflowId: string, user: User) { const workflow = await this.workflowFinderService.findWorkflowForUser(workflowId, user, [ 'workflow:read', @@ -110,13 +133,27 @@ export class TestRunsController { } @Post('/:workflowId/test-runs/new') - async create(req: TestRunsRequest.Create, res: express.Response) { + async create( + req: TestRunsRequest.Create, + res: express.Response, + @Body payload: StartTestRunRequestDto, + ) { const { workflowId } = req.params; await this.assertUserHasAccessToWorkflow(workflowId, req.user); + // Resolve the rollout flag for this user. Cached 10 min by PostHogClient + // so the hot path is one outbound call per user per 10-min window at + // most. Flag-off users are silently coerced to sequential — no error, + // no flag-id in the response — so the cohort wall is invisible to + // direct API callers and stale tabs. Fail-open on PostHog errors. + const flagEnabledForUser = await this.isParallelExecutionFlagEnabled(req.user); + + const requestedConcurrency = payload.concurrency ?? 1; + const concurrency = flagEnabledForUser ? requestedConcurrency : 1; + // We do not await for the test run to complete - void this.testRunnerService.runTest(req.user, workflowId); + void this.testRunnerService.runTest(req.user, workflowId, concurrency, flagEnabledForUser); res.status(202).json({ success: true }); } diff --git a/packages/cli/src/posthog/index.ts b/packages/cli/src/posthog/index.ts index f8158a0b6a4..e00f3366717 100644 --- a/packages/cli/src/posthog/index.ts +++ b/packages/cli/src/posthog/index.ts @@ -1,3 +1,4 @@ +import { EVAL_PARALLEL_EXECUTION_FLAG } from '@n8n/api-types'; import { GlobalConfig } from '@n8n/config'; import type { PublicUser } from '@n8n/db'; import { Service } from '@n8n/di'; @@ -100,6 +101,22 @@ export class PostHogClient { } async getFeatureFlags(user: Pick): Promise { + // Catch PostHog errors here (rather than letting them propagate) so + // env-var overrides still apply when PostHog is unreachable. Without + // this, a transient PostHog outage would short-circuit the override + // path and leave operators without an escape hatch. + let flags: FeatureFlags = {}; + try { + flags = await this.fetchFlagsFromPostHog(user); + } catch { + // fall through to env overrides + } + return this.applyEnvOverrides(flags); + } + + private async fetchFlagsFromPostHog( + user: Pick, + ): Promise { if (!this.postHog) return {}; const { instanceId } = this.instanceSettings; @@ -125,4 +142,17 @@ export class PostHogClient { return flags ?? {}; } + + /** + * Applies env-var overrides on top of PostHog-resolved flags. The override + * is force-enable only — `false` defers to PostHog. Cached PostHog data is + * stored without overrides so changing the env var (across restarts) + * doesn't poison the cache. + */ + private applyEnvOverrides(flags: FeatureFlags): FeatureFlags { + if (this.globalConfig.evaluation.parallelExecutionEnabled) { + return { ...flags, [EVAL_PARALLEL_EXECUTION_FLAG]: true }; + } + return flags; + } } diff --git a/packages/cli/test/integration/evaluation/test-runs.api.test.ts b/packages/cli/test/integration/evaluation/test-runs.api.test.ts index 85e2fe79404..83500cdc4c5 100644 --- a/packages/cli/test/integration/evaluation/test-runs.api.test.ts +++ b/packages/cli/test/integration/evaluation/test-runs.api.test.ts @@ -282,6 +282,8 @@ describe('POST /workflows/:workflowId/test-runs/new', () => { expect(testRunner.runTest).toHaveBeenCalledWith( expect.objectContaining({ id: ownerShell.id }), workflowUnderTest.id, + 1, + false, ); }); diff --git a/packages/frontend/@n8n/design-system/src/css/index.scss b/packages/frontend/@n8n/design-system/src/css/index.scss index b9039b28538..8144b73d7c0 100644 --- a/packages/frontend/@n8n/design-system/src/css/index.scss +++ b/packages/frontend/@n8n/design-system/src/css/index.scss @@ -12,6 +12,7 @@ @use './radio.scss'; @use './checkbox.scss'; @use './switch.scss'; +@use './slider.scss'; @use './select.scss'; @use './skeleton.scss'; @use './table.scss'; diff --git a/packages/frontend/@n8n/design-system/src/css/slider.scss b/packages/frontend/@n8n/design-system/src/css/slider.scss new file mode 100644 index 00000000000..6aa88f5f1a3 --- /dev/null +++ b/packages/frontend/@n8n/design-system/src/css/slider.scss @@ -0,0 +1,7 @@ +// Pulls element-plus's default slider styles into the bundle. Without this +// the runway, button, and bar render at 0px height because the el-slider +// component CSS isn't auto-imported elsewhere in the design system. +// +// Component-level brand-token overrides (e.g. setting --el-slider-main-bg-color +// to --color--primary) live wherever is used. +@use 'element-plus/theme-chalk/src/slider'; diff --git a/packages/frontend/@n8n/i18n/src/locales/en.json b/packages/frontend/@n8n/i18n/src/locales/en.json index c60491fbc92..18b1c5eb93d 100644 --- a/packages/frontend/@n8n/i18n/src/locales/en.json +++ b/packages/frontend/@n8n/i18n/src/locales/en.json @@ -4871,6 +4871,9 @@ "evaluation.runDetail.notice.useSetInputs": "Tip: Show input columns from your dataset here by adding the evaluation node's 'set inputs' operation to your workflow", "evaluation.runTest": "Run Test", "evaluation.stopTest": "Stop Test", + "evaluation.runInParallel.label.sequential": "Sequential", + "evaluation.runInParallel.label.concurrent": "Concurrent · {count}", + "evaluation.runInParallel.tooltip": "Drag the slider to set how many evaluation test cases run at the same time. A value of 1 runs cases sequentially; higher values fan out for speed but may hit LLM rate limits.", "evaluation.cancelTestRun": "Cancel Test Run", "evaluation.notImplemented": "This feature is not implemented yet!", "evaluation.viewDetails": "View Details", diff --git a/packages/frontend/editor-ui/src/app/constants/localStorage.ts b/packages/frontend/editor-ui/src/app/constants/localStorage.ts index 3046658f2ff..ae2f0a40693 100644 --- a/packages/frontend/editor-ui/src/app/constants/localStorage.ts +++ b/packages/frontend/editor-ui/src/app/constants/localStorage.ts @@ -30,3 +30,4 @@ export const LOCAL_STORAGE_CHAT_HUB_HAD_CONVERSATION_BEFORE = (userId: string) = export const LOCAL_STORAGE_SIDEBAR_WIDTH = 'N8N_SIDEBAR_WIDTH'; export const LOCAL_STORAGE_BROWSER_NOTIFICATION_METADATA = 'N8N_BROWSER_NOTIFICATION_METADATA'; export const LOCAL_STORAGE_FLOATING_CHAT_WINDOW = 'N8N_FLOATING_CHAT_WINDOW'; +export const LOCAL_STORAGE_PARALLEL_EVAL_BY_WORKFLOW = 'N8N_PARALLEL_EVAL_BY_WORKFLOW'; diff --git a/packages/frontend/editor-ui/src/features/ai/evaluation.ee/components/ConcurrencySlider/ConcurrencySlider.vue b/packages/frontend/editor-ui/src/features/ai/evaluation.ee/components/ConcurrencySlider/ConcurrencySlider.vue new file mode 100644 index 00000000000..8eb6cb7a8b2 --- /dev/null +++ b/packages/frontend/editor-ui/src/features/ai/evaluation.ee/components/ConcurrencySlider/ConcurrencySlider.vue @@ -0,0 +1,98 @@ + + + + + diff --git a/packages/frontend/editor-ui/src/features/ai/evaluation.ee/components/ConcurrencySlider/index.ts b/packages/frontend/editor-ui/src/features/ai/evaluation.ee/components/ConcurrencySlider/index.ts new file mode 100644 index 00000000000..5da63befe42 --- /dev/null +++ b/packages/frontend/editor-ui/src/features/ai/evaluation.ee/components/ConcurrencySlider/index.ts @@ -0,0 +1,3 @@ +import ConcurrencySlider from './ConcurrencySlider.vue'; + +export default ConcurrencySlider; diff --git a/packages/frontend/editor-ui/src/features/ai/evaluation.ee/evaluation.api.ts b/packages/frontend/editor-ui/src/features/ai/evaluation.ee/evaluation.api.ts index 8b649e7dbba..41d15e1a9cd 100644 --- a/packages/frontend/editor-ui/src/features/ai/evaluation.ee/evaluation.api.ts +++ b/packages/frontend/editor-ui/src/features/ai/evaluation.ee/evaluation.api.ts @@ -1,3 +1,4 @@ +import type { StartTestRunPayload } from '@n8n/api-types'; import type { IRestApiContext } from '@n8n/rest-api-client'; import { makeRestApiRequest, request } from '@n8n/rest-api-client'; import type { JsonObject } from 'n8n-workflow'; @@ -58,13 +59,25 @@ export const getTestRun = async (context: IRestApiContext, params: GetTestRunPar ); }; +// FE alias of the shared payload contract from @n8n/api-types. Re-exporting +// instead of duplicating the shape avoids silent drift between FE and BE. +export type StartTestRunOptions = StartTestRunPayload; + // Start a new test run -export const startTestRun = async (context: IRestApiContext, workflowId: string) => { +export const startTestRun = async ( + context: IRestApiContext, + workflowId: string, + options?: StartTestRunOptions, +) => { const response = await request({ method: 'POST', baseURL: context.baseUrl, endpoint: `/workflows/${workflowId}/test-runs/new`, headers: { 'push-ref': context.pushRef }, + // `data: undefined` sends an empty POST body. Express's body-parser + // normalises that to `req.body = {}`, which the controller's zod parse + // resolves to `concurrency: undefined`, defaulting to sequential. + data: options?.concurrency !== undefined ? { concurrency: options.concurrency } : undefined, }); // CLI is returning the response without wrapping it in `data` key return response as { success: boolean }; diff --git a/packages/frontend/editor-ui/src/features/ai/evaluation.ee/evaluation.store.test.ts b/packages/frontend/editor-ui/src/features/ai/evaluation.ee/evaluation.store.test.ts index c2f13be030d..6c85436b9ae 100644 --- a/packages/frontend/editor-ui/src/features/ai/evaluation.ee/evaluation.store.test.ts +++ b/packages/frontend/editor-ui/src/features/ai/evaluation.ee/evaluation.store.test.ts @@ -84,7 +84,16 @@ describe('evaluation.store.ee', () => { test('Starting Test Run', async () => { const result = await store.startTestRun('1'); - expect(startTestRun).toHaveBeenCalledWith(rootStoreMock.restApiContext, '1'); + expect(startTestRun).toHaveBeenCalledWith(rootStoreMock.restApiContext, '1', undefined); + expect(result).toEqual({ success: true }); + }); + + test('Starting Test Run with concurrency', async () => { + const result = await store.startTestRun('1', { concurrency: 5 }); + + expect(startTestRun).toHaveBeenCalledWith(rootStoreMock.restApiContext, '1', { + concurrency: 5, + }); expect(result).toEqual({ success: true }); }); diff --git a/packages/frontend/editor-ui/src/features/ai/evaluation.ee/evaluation.store.ts b/packages/frontend/editor-ui/src/features/ai/evaluation.ee/evaluation.store.ts index b6323d3ab22..68a492b60e8 100644 --- a/packages/frontend/editor-ui/src/features/ai/evaluation.ee/evaluation.store.ts +++ b/packages/frontend/editor-ui/src/features/ai/evaluation.ee/evaluation.store.ts @@ -125,8 +125,15 @@ export const useEvaluationStore = defineStore( return run; }; - const startTestRun = async (workflowId: string) => { - const result = await evaluationsApi.startTestRun(rootStore.restApiContext, workflowId); + const startTestRun = async ( + workflowId: string, + options?: evaluationsApi.StartTestRunOptions, + ) => { + const result = await evaluationsApi.startTestRun( + rootStore.restApiContext, + workflowId, + options, + ); return result; }; diff --git a/packages/frontend/editor-ui/src/features/ai/evaluation.ee/parallelEval.store.test.ts b/packages/frontend/editor-ui/src/features/ai/evaluation.ee/parallelEval.store.test.ts new file mode 100644 index 00000000000..f10d04bf6bf --- /dev/null +++ b/packages/frontend/editor-ui/src/features/ai/evaluation.ee/parallelEval.store.test.ts @@ -0,0 +1,149 @@ +import { createPinia, setActivePinia } from 'pinia'; +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'; + +import { LOCAL_STORAGE_PARALLEL_EVAL_BY_WORKFLOW } from '@/app/constants/localStorage'; +import { usePostHog } from '@/app/stores/posthog.store'; +import { DEFAULT_PARALLEL_CONCURRENCY, useParallelEvalStore } from './parallelEval.store'; + +vi.mock('@/app/stores/posthog.store', () => ({ + usePostHog: vi.fn(() => ({ + isFeatureEnabled: vi.fn(() => false), + })), +})); + +describe('parallelEval.store', () => { + beforeEach(() => { + setActivePinia(createPinia()); + localStorage.removeItem(LOCAL_STORAGE_PARALLEL_EVAL_BY_WORKFLOW); + }); + + afterEach(() => { + vi.clearAllMocks(); + }); + + describe('isFeatureEnabled', () => { + it('reflects the PostHog flag (off by default)', () => { + vi.mocked(usePostHog).mockReturnValue({ + isFeatureEnabled: vi.fn(() => false), + } as never); + + const store = useParallelEvalStore(); + + expect(store.isFeatureEnabled).toBe(false); + }); + + it('returns true when the rollout flag resolves true', () => { + vi.mocked(usePostHog).mockReturnValue({ + isFeatureEnabled: vi.fn(() => true), + } as never); + + const store = useParallelEvalStore(); + + expect(store.isFeatureEnabled).toBe(true); + }); + }); + + describe('per-workflow defaults', () => { + it('isParallel defaults to true for an unsaved workflow', () => { + const store = useParallelEvalStore(); + expect(store.isParallel('wf-a')).toBe(true); + }); + + it('concurrencyValue defaults to 3 for an unsaved workflow', () => { + const store = useParallelEvalStore(); + expect(store.concurrencyValue('wf-a')).toBe(DEFAULT_PARALLEL_CONCURRENCY); + }); + + it('uses the "new" sentinel when workflowId is empty or undefined', () => { + const store = useParallelEvalStore(); + store.setParallel('', false); + store.setConcurrencyValue('', 7); + + expect(store.isParallel(undefined)).toBe(false); + expect(store.concurrencyValue('')).toBe(7); + }); + }); + + describe('persistence', () => { + it('per-workflow state is independent', () => { + const store = useParallelEvalStore(); + store.setParallel('wf-a', false); + store.setConcurrencyValue('wf-a', 5); + + expect(store.isParallel('wf-b')).toBe(true); + expect(store.concurrencyValue('wf-b')).toBe(DEFAULT_PARALLEL_CONCURRENCY); + }); + + it('reflects writes via the public getters', () => { + const store = useParallelEvalStore(); + store.setParallel('wf-a', false); + store.setConcurrencyValue('wf-a', 8); + + expect(store.isParallel('wf-a')).toBe(false); + expect(store.concurrencyValue('wf-a')).toBe(8); + }); + + it('"new" sentinel state survives a workflow ID assignment', () => { + const store = useParallelEvalStore(); + store.setParallel(undefined, false); + store.setConcurrencyValue(undefined, 4); + + // The user saves the workflow and gets a real id. The sentinel + // entry remains in localStorage (orphaned but harmless); the new + // id picks up the defaults — that's intentional, the choice was + // per-workflow scoped. + expect(store.isParallel('newly-assigned-id')).toBe(true); + expect(store.concurrencyValue('newly-assigned-id')).toBe(DEFAULT_PARALLEL_CONCURRENCY); + // And the sentinel still holds the user's pre-save selection. + expect(store.isParallel(undefined)).toBe(false); + expect(store.concurrencyValue(undefined)).toBe(4); + }); + }); + + describe('clamping', () => { + it('setConcurrencyValue clamps to 1-10', () => { + const store = useParallelEvalStore(); + store.setConcurrencyValue('wf-a', 0); + expect(store.concurrencyValue('wf-a')).toBe(1); + + store.setConcurrencyValue('wf-a', 99); + expect(store.concurrencyValue('wf-a')).toBe(10); + + store.setConcurrencyValue('wf-a', 4.7); + expect(store.concurrencyValue('wf-a')).toBe(4); + }); + + it('setConcurrencyValue falls back to the default for non-finite input', () => { + const store = useParallelEvalStore(); + + // NaN propagates unchanged through Math.floor/min/max, so without a + // guard a cleared N8nInputNumber would persist NaN and break the + // 1-10 contract. Verify it lands at the parallel default instead. + store.setConcurrencyValue('wf-a', Number.NaN); + expect(store.concurrencyValue('wf-a')).toBe(DEFAULT_PARALLEL_CONCURRENCY); + + // Positive Infinity already clamps to 10 via Math.min, but pin the + // behaviour explicitly so future refactors don't regress it. + store.setConcurrencyValue('wf-b', Number.POSITIVE_INFINITY); + expect(store.concurrencyValue('wf-b')).toBe(DEFAULT_PARALLEL_CONCURRENCY); + + store.setConcurrencyValue('wf-c', Number.NEGATIVE_INFINITY); + expect(store.concurrencyValue('wf-c')).toBe(DEFAULT_PARALLEL_CONCURRENCY); + }); + }); + + describe('effectiveConcurrency', () => { + it('returns the slider value when parallel is enabled', () => { + const store = useParallelEvalStore(); + store.setConcurrencyValue('wf-a', 6); + expect(store.effectiveConcurrency('wf-a')).toBe(6); + }); + + it('returns 1 (sequential) when parallel is disabled, regardless of slider', () => { + const store = useParallelEvalStore(); + store.setConcurrencyValue('wf-a', 6); + store.setParallel('wf-a', false); + expect(store.effectiveConcurrency('wf-a')).toBe(1); + }); + }); +}); diff --git a/packages/frontend/editor-ui/src/features/ai/evaluation.ee/parallelEval.store.ts b/packages/frontend/editor-ui/src/features/ai/evaluation.ee/parallelEval.store.ts new file mode 100644 index 00000000000..a81869b2440 --- /dev/null +++ b/packages/frontend/editor-ui/src/features/ai/evaluation.ee/parallelEval.store.ts @@ -0,0 +1,106 @@ +import { EVAL_PARALLEL_EXECUTION_FLAG } from '@n8n/api-types'; +import { useLocalStorage } from '@vueuse/core'; +import { defineStore } from 'pinia'; +import { computed } from 'vue'; + +import { LOCAL_STORAGE_PARALLEL_EVAL_BY_WORKFLOW } from '@/app/constants/localStorage'; +import { usePostHog } from '@/app/stores/posthog.store'; + +// Sentinel used for workflows that haven't been saved yet (no id assigned). +// Mirrors the per-workflow localStorage pattern used elsewhere in the editor. +const NEW_WORKFLOW_SENTINEL = 'new'; + +export const DEFAULT_PARALLEL_CONCURRENCY = 3; + +interface PerWorkflowState { + parallelEnabled: boolean; + concurrencyValue: number; +} + +type StoredState = Record; + +const buildDefaultState = (): PerWorkflowState => ({ + parallelEnabled: true, + concurrencyValue: DEFAULT_PARALLEL_CONCURRENCY, +}); + +/** + * Per-workflow UI state for the parallel-execution rollout. Visibility of the + * UI is gated on the `080_eval_parallel_execution` PostHog flag (FE primary gate). + * Backend has its own safety net that coerces flag-off requests to sequential. + * + * State shape: `{ [workflowId]: { parallelEnabled, concurrencyValue } }`. + * Workflow id `'new'` is a sentinel for unsaved workflows; the entry becomes + * orphaned in localStorage once the workflow gets a real id, but it's + * harmless and self-cleaning across sessions. + */ +export const useParallelEvalStore = defineStore('parallelEval', () => { + const postHog = usePostHog(); + const storage = useLocalStorage( + LOCAL_STORAGE_PARALLEL_EVAL_BY_WORKFLOW, + {}, + // `flush: 'sync'` so reads-after-writes don't race in tight loops. + { deep: true, flush: 'sync' }, + ); + + // Coerce `boolean | undefined` (PostHog's return shape) to a clean boolean. + const isFeatureEnabled = computed( + () => postHog.isFeatureEnabled(EVAL_PARALLEL_EXECUTION_FLAG) === true, + ); + + const resolveKey = (workflowId: string | undefined): string => + workflowId && workflowId.length > 0 ? workflowId : NEW_WORKFLOW_SENTINEL; + + // Always materialises a reactive entry in storage rather than returning a + // detached default object. Without this, a read-only consumer reading + // before any setter writes would observe a stale plain object that never + // invalidates when storage updates. + const ensureEntry = (key: string): PerWorkflowState => { + if (!storage.value[key]) { + storage.value[key] = buildDefaultState(); + } + return storage.value[key]; + }; + + const isParallel = (workflowId: string | undefined): boolean => + ensureEntry(resolveKey(workflowId)).parallelEnabled; + + const concurrencyValue = (workflowId: string | undefined): number => + ensureEntry(resolveKey(workflowId)).concurrencyValue; + + const setParallel = (workflowId: string | undefined, value: boolean): void => { + ensureEntry(resolveKey(workflowId)).parallelEnabled = value; + }; + + const setConcurrencyValue = (workflowId: string | undefined, value: number): void => { + // Guard against non-finite inputs (NaN from a cleared N8nInputNumber, + // Infinity from edge-case maths). NaN would propagate through the + // floor/min/max chain unchanged and persist a broken state, silently + // violating the 1-10 contract. Fall back to the parallel default so + // the checked-but-cleared UX feels natural rather than dropping to + // sequential behind the user's back. + const safe = Number.isFinite(value) ? value : DEFAULT_PARALLEL_CONCURRENCY; + const clamped = Math.max(1, Math.min(10, Math.floor(safe))); + ensureEntry(resolveKey(workflowId)).concurrencyValue = clamped; + }; + + /** + * The numeric concurrency the FE should send for a run. Returns `1` when + * the parallel checkbox is unchecked (sequential), the slider value when + * checked. Caller is responsible for skipping the field entirely when + * the feature flag is off. + */ + const effectiveConcurrency = (workflowId: string | undefined): number => { + const state = ensureEntry(resolveKey(workflowId)); + return state.parallelEnabled ? state.concurrencyValue : 1; + }; + + return { + isFeatureEnabled, + isParallel, + concurrencyValue, + setParallel, + setConcurrencyValue, + effectiveConcurrency, + }; +}); diff --git a/packages/frontend/editor-ui/src/features/ai/evaluation.ee/views/EvaluationsView.test.ts b/packages/frontend/editor-ui/src/features/ai/evaluation.ee/views/EvaluationsView.test.ts index 34e2820adae..1872f8a277f 100644 --- a/packages/frontend/editor-ui/src/features/ai/evaluation.ee/views/EvaluationsView.test.ts +++ b/packages/frontend/editor-ui/src/features/ai/evaluation.ee/views/EvaluationsView.test.ts @@ -5,6 +5,7 @@ import EvaluationsView from './EvaluationsView.vue'; import { mockedStore } from '@/__tests__/utils'; import { useEvaluationStore } from '../evaluation.store'; +import { useParallelEvalStore } from '../parallelEval.store'; import userEvent from '@testing-library/user-event'; import type { TestRunRecord } from '../evaluation.api'; import { waitFor } from '@testing-library/vue'; @@ -92,7 +93,15 @@ describe('EvaluationsView', () => { await userEvent.click(getByTestId('run-test-button')); - expect(evaluationStore.startTestRun).toHaveBeenCalledWith('workflow-id'); + // Assert only the workflow id (first arg). The second arg is the + // parallel-execution options payload — `undefined` when the + // rollout flag is off, an object when it's on — and isn't + // what this test is checking. + expect(evaluationStore.startTestRun).toHaveBeenCalled(); + const [firstCallWorkflowId] = ( + evaluationStore.startTestRun as unknown as ReturnType + ).mock.calls[0]; + expect(firstCallWorkflowId).toBe('workflow-id'); expect(evaluationStore.fetchTestRuns).toHaveBeenCalledWith('workflow-id'); }); @@ -116,6 +125,130 @@ describe('EvaluationsView', () => { await waitFor(() => expect(getByTestId('stop-test-button')).toBeInTheDocument()); }); + describe('parallel-execution UI', () => { + it('does not render the parallel controls when the rollout flag is off', async () => { + const evaluationStore = mockedStore(useEvaluationStore); + evaluationStore.testRunsById = {}; + const parallelEvalStore = mockedStore(useParallelEvalStore); + parallelEvalStore.isFeatureEnabled = false; + + const { queryByTestId, getByTestId } = renderComponent(); + await waitFor(() => expect(getByTestId('run-test-button')).toBeInTheDocument()); + + expect(queryByTestId('parallel-eval-controls')).toBeNull(); + expect(queryByTestId('run-in-parallel-concurrency')).toBeNull(); + expect(queryByTestId('run-in-parallel-mode-label')).toBeNull(); + }); + + it('renders the slider + tooltip + label when the flag is on', async () => { + const evaluationStore = mockedStore(useEvaluationStore); + evaluationStore.testRunsById = {}; + const parallelEvalStore = mockedStore(useParallelEvalStore); + parallelEvalStore.isFeatureEnabled = true; + parallelEvalStore.concurrencyValue.mockReturnValue(3); + + const { getByTestId } = renderComponent(); + + await waitFor(() => expect(getByTestId('parallel-eval-controls')).toBeInTheDocument()); + expect(getByTestId('run-in-parallel-mode-label')).toBeInTheDocument(); + expect(getByTestId('run-in-parallel-tooltip-icon')).toBeInTheDocument(); + expect(getByTestId('run-in-parallel-concurrency')).toBeInTheDocument(); + }); + + it('shows "Sequential" label when slider is at 1', async () => { + const evaluationStore = mockedStore(useEvaluationStore); + evaluationStore.testRunsById = {}; + const parallelEvalStore = mockedStore(useParallelEvalStore); + parallelEvalStore.isFeatureEnabled = true; + parallelEvalStore.concurrencyValue.mockReturnValue(1); + + const { getByTestId } = renderComponent(); + await waitFor(() => expect(getByTestId('run-in-parallel-mode-label')).toBeInTheDocument()); + + expect(getByTestId('run-in-parallel-mode-label').textContent?.trim()).toBe('Sequential'); + }); + + it('shows "Concurrent · N" label when slider is above 1', async () => { + const evaluationStore = mockedStore(useEvaluationStore); + evaluationStore.testRunsById = {}; + const parallelEvalStore = mockedStore(useParallelEvalStore); + parallelEvalStore.isFeatureEnabled = true; + parallelEvalStore.concurrencyValue.mockReturnValue(7); + + const { getByTestId } = renderComponent(); + await waitFor(() => expect(getByTestId('run-in-parallel-mode-label')).toBeInTheDocument()); + + expect(getByTestId('run-in-parallel-mode-label').textContent?.trim()).toBe( + 'Concurrent · 7', + ); + }); + + it('omits concurrency from the request when the flag is off', async () => { + const evaluationStore = mockedStore(useEvaluationStore); + evaluationStore.testRunsById = {}; + const parallelEvalStore = mockedStore(useParallelEvalStore); + parallelEvalStore.isFeatureEnabled = false; + + const { getByTestId } = renderComponent(); + await waitFor(() => expect(getByTestId('run-test-button')).toBeInTheDocument()); + + await userEvent.click(getByTestId('run-test-button')); + + expect(evaluationStore.startTestRun).toHaveBeenCalledWith('workflow-id', undefined); + }); + + it('sends the slider value as concurrency when the flag is on', async () => { + const evaluationStore = mockedStore(useEvaluationStore); + evaluationStore.testRunsById = {}; + const parallelEvalStore = mockedStore(useParallelEvalStore); + parallelEvalStore.isFeatureEnabled = true; + parallelEvalStore.concurrencyValue.mockReturnValue(3); + + const { getByTestId } = renderComponent(); + await waitFor(() => expect(getByTestId('run-test-button')).toBeInTheDocument()); + + await userEvent.click(getByTestId('run-test-button')); + + expect(evaluationStore.startTestRun).toHaveBeenCalledWith('workflow-id', { + concurrency: 3, + }); + }); + + it('sends concurrency=1 when the user has dragged the slider to the minimum', async () => { + const evaluationStore = mockedStore(useEvaluationStore); + evaluationStore.testRunsById = {}; + const parallelEvalStore = mockedStore(useParallelEvalStore); + parallelEvalStore.isFeatureEnabled = true; + parallelEvalStore.concurrencyValue.mockReturnValue(1); + + const { getByTestId } = renderComponent(); + await waitFor(() => expect(getByTestId('run-test-button')).toBeInTheDocument()); + + await userEvent.click(getByTestId('run-test-button')); + + expect(evaluationStore.startTestRun).toHaveBeenCalledWith('workflow-id', { + concurrency: 1, + }); + }); + + it('forwards the slider value when the user picked a non-default concurrency', async () => { + const evaluationStore = mockedStore(useEvaluationStore); + evaluationStore.testRunsById = {}; + const parallelEvalStore = mockedStore(useParallelEvalStore); + parallelEvalStore.isFeatureEnabled = true; + parallelEvalStore.concurrencyValue.mockReturnValue(7); + + const { getByTestId } = renderComponent(); + await waitFor(() => expect(getByTestId('run-test-button')).toBeInTheDocument()); + + await userEvent.click(getByTestId('run-test-button')); + + expect(evaluationStore.startTestRun).toHaveBeenCalledWith('workflow-id', { + concurrency: 7, + }); + }); + }); + it('should call cancelTestRun when stop button is clicked', async () => { const evaluationStore = mockedStore(useEvaluationStore); evaluationStore.cancelTestRun.mockResolvedValue({ success: true }); diff --git a/packages/frontend/editor-ui/src/features/ai/evaluation.ee/views/EvaluationsView.vue b/packages/frontend/editor-ui/src/features/ai/evaluation.ee/views/EvaluationsView.vue index 86b1dfa8434..cae584c6592 100644 --- a/packages/frontend/editor-ui/src/features/ai/evaluation.ee/views/EvaluationsView.vue +++ b/packages/frontend/editor-ui/src/features/ai/evaluation.ee/views/EvaluationsView.vue @@ -2,12 +2,15 @@ import { useI18n } from '@n8n/i18n'; import { computed, ref, watch } from 'vue'; +import ConcurrencySlider from '../components/ConcurrencySlider'; import RunsSection from '../components/ListRuns/RunsSection.vue'; import { useEvaluationStore } from '../evaluation.store'; +import { useParallelEvalStore } from '../parallelEval.store'; import orderBy from 'lodash/orderBy'; import { useToast } from '@/app/composables/useToast'; -import { N8nButton } from '@n8n/design-system'; +import { N8nButton, N8nIcon, N8nTooltip } from '@n8n/design-system'; + const props = defineProps<{ workflowId: string; }>(); @@ -16,15 +19,40 @@ const locale = useI18n(); const toast = useToast(); const evaluationStore = useEvaluationStore(); +const parallelEvalStore = useParallelEvalStore(); const selectedMetric = ref(''); const cancellingTestRun = ref(false); const runningTestRun = computed(() => runs.value.find((run) => run.status === 'running')); +const concurrencyModel = computed({ + get: () => parallelEvalStore.concurrencyValue(props.workflowId), + set: (value: number) => parallelEvalStore.setConcurrencyValue(props.workflowId, value), +}); + +// Slider value 1 = sequential, > 1 = concurrent. The slider is the single +// source of truth — no separate toggle. The current count is folded into +// the label ("Concurrent · 3") so the standalone numeric readout can be +// dropped from the layout. +const concurrencyLabel = computed(() => + concurrencyModel.value > 1 + ? locale.baseText('evaluation.runInParallel.label.concurrent', { + interpolate: { count: String(concurrencyModel.value) }, + }) + : locale.baseText('evaluation.runInParallel.label.sequential'), +); + async function runTest() { try { - await evaluationStore.startTestRun(props.workflowId); + // When the rollout flag is off, omit `concurrency` entirely so the BE + // safety net is a no-op and behaviour matches the legacy sequential + // path. Flag-on cohort sends the slider value (1 = sequential, >1 = + // concurrent fan-out). + const options = parallelEvalStore.isFeatureEnabled + ? { concurrency: concurrencyModel.value } + : undefined; + await evaluationStore.startTestRun(props.workflowId, options); } catch (error) { toast.showError(error, locale.baseText('evaluation.listRuns.error.cantStartTestRun')); } @@ -74,25 +102,56 @@ watch(runningTestRun, (run) => {