feat(core): Run evaluation test cases in parallel behind PostHog rollout flag (#29412)

This commit is contained in:
Arvin A
2026-05-04 13:18:01 +00:00
committed by GitHub
parent e35042999f
commit 4c76aa1467
28 changed files with 1849 additions and 192 deletions
+7
View File
@@ -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 {
@@ -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<typeof startTestRunPayloadSchema>;
// 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) {}
@@ -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;
}
+4
View File
@@ -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;
+3
View File
@@ -371,6 +371,9 @@ describe('GlobalConfig', () => {
ttl: 10,
interval: 3,
},
evaluation: {
parallelExecutionEnabled: false,
},
generic: {
timezone: 'America/New_York',
releaseChannel: 'dev',
+1
View File
@@ -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",
@@ -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<TestCaseExecutionRepository>;
let mockTestRunnerService: jest.Mocked<TestRunnerService>;
let mockTelemetry: jest.Mocked<Telemetry>;
let mockPostHogClient: jest.Mocked<PostHogClient>;
let mockLogger: jest.Mocked<Logger>;
let mockUser: User;
let mockWorkflowId: string;
let mockTestRunId: string;
@@ -47,15 +52,28 @@ describe('TestRunsController', () => {
track: jest.fn(),
} as unknown as jest.Mocked<Telemetry>;
mockPostHogClient = {
getFeatureFlags: jest.fn().mockResolvedValue({}),
} as unknown as jest.Mocked<PostHogClient>;
mockLogger = {
warn: jest.fn(),
debug: jest.fn(),
error: jest.fn(),
info: jest.fn(),
} as unknown as jest.Mocked<Logger>;
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),
);
});
});
});
@@ -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',
);
});
});
});
@@ -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<Publisher>();
const instanceSettings = mock<InstanceSettings>({ hostId: 'test-host-id', isMultiMain: false });
const concurrencyControlService = mock<ConcurrencyControlService>();
let testRunnerService: TestRunnerService;
mockInstance(LoadNodesAndCredentials, {
@@ -65,6 +67,7 @@ describe('TestRunnerService', () => {
mock(),
publisher,
instanceSettings,
concurrencyControlService,
);
testRunRepository.createTestRun.mockResolvedValue(mock<TestRun>({ 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<IWorkflowBase>) 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<T>().
// 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<unknown>) => 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<TestRun>({ 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<string, number>)[key]).toBeCloseTo(
(sequentialMetrics as Record<string, number>)[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<string, unknown>;
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<string, unknown>;
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<string, unknown>;
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<string, unknown>;
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<string, unknown>;
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<void>(() => {
/* 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<InstanceSettings>({
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();
});
});
});
@@ -7,27 +7,37 @@ export interface EvaluationMetricsAddResultsInfo {
incorrectTypeMetrics: Set<string>;
}
/**
* 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<string, number>;
}
export class EvaluationMetrics {
private readonly rawMetricsByName = new Map<string, number[]>();
addResults(result: IDataObject): EvaluationMetricsAddResultsInfo {
const addResultsInfo: EvaluationMetricsAddResultsInfo = {
addedMetrics: {},
incorrectTypeMetrics: new Set<string>(),
};
/**
* 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<string, number> = {};
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() {
@@ -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<void> {
this.logger.debug('Starting new test run', { workflowId });
async runTest(
user: User,
workflowId: string,
concurrency: number = 1,
flagEnabledForUser: boolean = false,
): Promise<void> {
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<MetricContribution[]> => {
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);
@@ -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<boolean> {
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 });
}
+30
View File
@@ -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<PublicUser, 'id' | 'createdAt'>): Promise<FeatureFlags> {
// 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<PublicUser, 'id' | 'createdAt'>,
): Promise<FeatureFlags> {
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;
}
}
@@ -282,6 +282,8 @@ describe('POST /workflows/:workflowId/test-runs/new', () => {
expect(testRunner.runTest).toHaveBeenCalledWith(
expect.objectContaining({ id: ownerShell.id }),
workflowUnderTest.id,
1,
false,
);
});
@@ -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';
@@ -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 <ElSlider> is used.
@use 'element-plus/theme-chalk/src/slider';
@@ -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",
@@ -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';
@@ -0,0 +1,98 @@
<script setup lang="ts">
import { ElSlider } from 'element-plus';
// Feature-local slider for the parallel-evaluation concurrency control.
// Lives here (not in `@n8n/design-system`) because it bakes in a single
// hard-coded brand-token theme — that bypasses the variant story DS
// components are expected to support and would need design sign-off
// before being promoted. Promote when both conditions hold.
interface ConcurrencySliderProps {
min?: number;
max?: number;
step?: number;
disabled?: boolean;
showStops?: boolean;
showTooltip?: boolean;
}
withDefaults(defineProps<ConcurrencySliderProps>(), {
min: 0,
max: 100,
step: 1,
disabled: false,
showStops: false,
showTooltip: true,
});
// Element-plus's `.el-slider` rule sets these CSS vars at class-level
// specificity. Inline `:style` wins by specificity, guaranteeing the
// override regardless of stylesheet load order.
//
// Discord-like pill aesthetic: tall rounded runway, white thumb sitting
// inside the track without a visible coloured border ring, and subtle
// stops as dots rather than dividers.
//
// `--background--brand` is the semantic alias for the n8n orange. The
// runway/disabled tokens fall back to legacy `--color--foreground*`
// because the design system has no semantic alias for the
// "neutral surface behind an interactive control" role yet.
const brandTokens = {
'--el-slider-main-bg-color': 'var(--background--brand)',
'--el-slider-runway-bg-color': 'var(--color--foreground)',
'--el-slider-stop-bg-color': 'rgba(0, 0, 0, 0.28)',
'--el-slider-disabled-color': 'var(--color--foreground--shade-1)',
'--el-color-white': '#fff',
'--el-slider-button-size': '16px',
'--el-slider-button-wrapper-size': '28px',
'--el-slider-height': '20px',
'--el-slider-border-radius': '20px',
// Vertically centre the thumb wrapper inside the runway. Element-plus's
// default `-15px` was tuned for a 6px runway and 36px wrapper. With our
// 20px runway and 28px wrapper, the offset is `(runway - wrapper) / 2`
// = `(20 - 28) / 2` = -4px.
'--el-slider-button-wrapper-offset': '-4px',
} as const;
</script>
<template>
<ElSlider
:min="min"
:max="max"
:step="step"
:disabled="disabled"
:show-stops="showStops"
:show-tooltip="showTooltip"
:style="brandTokens"
:class="$style.slider"
v-bind="$attrs"
/>
</template>
<style lang="scss" module>
// Element-plus's slider thumb has `border: 2px solid var(--el-slider-main-bg-color)`,
// which puts a primary-coloured ring around the thumb. The Discord-like look
// the design wants is a borderless white thumb sitting inside the track. Force
// the border to match the thumb fill (effectively removing the ring) and add
// a subtle shadow for depth. `:global` is needed because the .el-slider__button
// class is rendered by ElSlider (outside this component's scoped CSS).
.slider {
:global(.el-slider__button) {
border-color: var(--el-color-white);
box-shadow:
0 2px 4px rgba(0, 0, 0, 0.12),
0 0 0 1px rgba(0, 0, 0, 0.04);
}
// Element-plus's default stops are 4px-wide vertical rectangles spanning
// the full runway height, which look like dividers on a tall pill track.
// Rounded dots match the Discord-style visual.
:global(.el-slider__stop) {
width: 6px;
height: 6px;
border-radius: 50%;
top: 50%;
transform: translate(-50%, -50%);
}
}
</style>
@@ -0,0 +1,3 @@
import ConcurrencySlider from './ConcurrencySlider.vue';
export default ConcurrencySlider;
@@ -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 };
@@ -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 });
});
@@ -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;
};
@@ -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);
});
});
});
@@ -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<string, PerWorkflowState>;
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<StoredState>(
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,
};
});
@@ -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<typeof vi.fn>
).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 });
@@ -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<string>('');
const cancellingTestRun = ref<boolean>(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) => {
<template>
<div :class="$style.evaluationsView">
<div :class="$style.header">
<N8nButton
variant="subtle"
v-if="runningTestRun"
:disabled="cancellingTestRun"
:class="$style.runOrStopTestButton"
size="small"
data-test-id="stop-test-button"
:label="locale.baseText('evaluation.stopTest')"
@click="stopTest"
/>
<N8nButton
variant="solid"
v-else
:class="$style.runOrStopTestButton"
size="small"
data-test-id="run-test-button"
:label="locale.baseText('evaluation.runTest')"
@click="runTest"
/>
<div :class="$style.headerInner">
<div
v-if="parallelEvalStore.isFeatureEnabled && !runningTestRun"
:class="$style.parallelControls"
data-test-id="parallel-eval-controls"
>
<span :class="$style.concurrencyLabel" data-test-id="run-in-parallel-mode-label">
{{ concurrencyLabel }}
</span>
<N8nTooltip
placement="top"
:content="locale.baseText('evaluation.runInParallel.tooltip')"
>
<N8nIcon
icon="circle-help"
size="small"
:class="$style.tooltipIcon"
data-test-id="run-in-parallel-tooltip-icon"
/>
</N8nTooltip>
<ConcurrencySlider
v-model="concurrencyModel"
:min="1"
:max="10"
:step="1"
show-stops
:class="$style.concurrencySlider"
data-test-id="run-in-parallel-concurrency"
/>
</div>
<N8nButton
variant="subtle"
v-if="runningTestRun"
:disabled="cancellingTestRun"
:class="$style.runOrStopTestButton"
size="small"
data-test-id="stop-test-button"
:label="locale.baseText('evaluation.stopTest')"
@click="stopTest"
/>
<N8nButton
variant="solid"
v-else
:class="$style.runOrStopTestButton"
size="small"
data-test-id="run-test-button"
:label="locale.baseText('evaluation.runTest')"
@click="runTest"
/>
</div>
</div>
<div :class="$style.wrapper">
<div :class="$style.content">
@@ -121,10 +180,12 @@ watch(runningTestRun, (run) => {
.header {
display: flex;
justify-content: end;
align-items: center;
justify-content: center;
padding: var(--spacing--md) var(--spacing--lg);
padding-left: 27px;
// Match `.wrapper` so the inner 1024px box anchors at the same left
// origin as the runs section below (the 58px reserves space for the
// editor's collapsed sidebar in the wrapper layout).
padding-left: 58px;
padding-bottom: 8px;
position: sticky;
top: 0;
@@ -133,6 +194,17 @@ watch(runningTestRun, (run) => {
z-index: 2;
}
// Inner container: width-capped at the same 1024px as `.runs` and
// horizontally centred via the parent's `justify-content: center`. This
// makes the controls + Run button hug the same horizontal bounds as the
// chart and table below, instead of floating against the viewport edges.
.headerInner {
display: flex;
align-items: center;
width: 100%;
max-width: 1024px;
}
.wrapper {
padding: 0 var(--spacing--lg);
padding-left: 58px;
@@ -140,6 +212,40 @@ watch(runningTestRun, (run) => {
.runOrStopTestButton {
white-space: nowrap;
// Anchor to the right edge of `.headerInner` regardless of which
// siblings render. Works for the flag-off case (only the button is
// rendered) and the flag-on case (parallel controls take the left).
margin-left: auto;
}
.parallelControls {
display: flex;
align-items: center;
gap: var(--spacing--xs);
}
.tooltipIcon {
color: var(--color--text--tint-1);
cursor: help;
}
.concurrencyLabel {
color: var(--color--text);
font-size: var(--font-size--sm);
// Tabular numerals so the width doesn't jitter as the user drags the
// slider through different digits.
font-variant-numeric: tabular-nums;
// Reserve enough room for the longest variant ("Concurrent · 10") so
// the slider doesn't shift left/right as the label content changes.
min-width: 100px;
}
.concurrencySlider {
// `flex: 0 0 160px` (basis 160px, no grow/shrink) wins over element-plus's
// default `.el-slider { width: 100% }` because flex-basis is what the flex
// algorithm uses for sizing. Plain `width: 160px` would lose the cascade.
flex: 0 0 160px;
margin: 0 var(--spacing--2xs);
}
.runs {
+3
View File
@@ -2729,6 +2729,9 @@ importers:
p-lazy:
specifier: 3.1.0
version: 3.1.0
p-limit:
specifier: ^3.1.0
version: 3.1.0
pg:
specifier: 'catalog:'
version: 8.17.0