feat(engine): Let the control plane mint the engine 2.0 execution id (#37224)

This commit is contained in:
Tomi Turtiainen
2026-08-31 14:04:01 +00:00
committed by GitHub
parent 1f19b11c40
commit b04acabea4
26 changed files with 225 additions and 93 deletions
+1
View File
@@ -7,6 +7,7 @@ export * from './execution';
export * from './logstreaming';
export * from './nodes';
export * from './scheduler';
export * from './uuid';
export const LICENSE_FEATURES = {
SHARING: 'feat:sharing',
+6
View File
@@ -0,0 +1,6 @@
/**
* Matches a uuidv7: version `7` in the third group, RFC 4122 variant in the
* fourth. Time-ordered ids are validated against this on the wire.
*/
export const UUID_V7_PATTERN =
/^[0-9a-f]{8}-[0-9a-f]{4}-7[0-9a-f]{3}-[89ab][0-9a-f]{3}-[0-9a-f]{12}$/i;
+1
View File
@@ -25,6 +25,7 @@
],
"dependencies": {
"@n8n/config": "workspace:*",
"@n8n/constants": "workspace:*",
"@n8n/di": "workspace:*",
"@n8n/typeorm": "workspace:*",
"express": "catalog:",
@@ -7,6 +7,7 @@ import { ExecutionNotFoundError } from '../../execution/execution-store';
import { createDataSource } from '../data-source';
import { WorkflowExecution } from '../entities/workflow-execution.entity';
import { WorkflowStepExecution } from '../entities/workflow-step-execution.entity';
import { generateId } from '../generate-id';
import { TypeOrmExecutionViewStore } from '../typeorm-execution-view-store';
describe('workflow_execution table (integration)', () => {
@@ -29,6 +30,7 @@ describe('workflow_execution table (integration)', () => {
const repo = dataSource.getRepository(WorkflowExecution);
const created = repo.create({
id: generateId(),
workflowId: 'wf-1',
status: 'running',
mode: 'production',
@@ -57,6 +59,7 @@ describe('workflow_execution table (integration)', () => {
const repo = dataSource.getRepository(WorkflowExecution);
const finishedAt = new Date();
const created = repo.create({
id: generateId(),
workflowId: 'wf-3',
status: 'completed',
mode: 'manual',
@@ -100,6 +103,7 @@ describe('workflow_execution table (integration)', () => {
await repo.save(
repo.create({
id: generateId(),
workflowId: 'wf-2',
status: 'running',
mode: 'production',
@@ -110,6 +114,7 @@ describe('workflow_execution table (integration)', () => {
);
await repo.save(
repo.create({
id: generateId(),
workflowId: 'wf-2',
status: 'completed',
mode: 'production',
@@ -8,6 +8,7 @@ import { StepNotFoundError, type NewStepRecord } from '../../execution/step-stor
import { createDataSource } from '../data-source';
import { WorkflowExecution } from '../entities/workflow-execution.entity';
import { WorkflowStepExecution } from '../entities/workflow-step-execution.entity';
import { generateId } from '../generate-id';
import { TypeOrmExecutionViewStore } from '../typeorm-execution-view-store';
import { TypeOrmStepStore } from '../typeorm-step-store';
@@ -39,6 +40,7 @@ describe('workflow_step_execution table (integration)', () => {
async function createExecution(): Promise<string> {
const repo = dataSource.getRepository(WorkflowExecution);
const execution = repo.create({
id: generateId(),
workflowId: 'wf-1',
status: 'running',
mode: 'production',
@@ -1,5 +1,4 @@
import {
BeforeInsert,
Column,
CreateDateColumn,
Entity,
@@ -14,7 +13,6 @@ import type {
TriggerOutputs,
} from '../../execution/execution.types';
import type { WorkflowGraph } from '../../graph';
import { generateId } from '../generate-id';
@Entity('workflow_execution')
@Index('idx_workflow_execution_workflow_id', ['workflowId'])
@@ -46,9 +44,4 @@ export class WorkflowExecution {
@Column({ name: 'finished_at', type: 'timestamptz', nullable: true, precision: 3 })
finishedAt!: Date | null;
@BeforeInsert()
setId(): void {
if (!this.id) this.id = generateId();
}
}
@@ -19,13 +19,12 @@ type InsertValues = Parameters<Repository<WorkflowExecution>['insert']>[0];
export class TypeOrmExecutionStore implements ExecutionStore {
constructor(private readonly repo: Repository<WorkflowExecution>) {}
async createExecution(record: NewExecutionRecord): Promise<{ id: string }> {
async createExecution(record: NewExecutionRecord): Promise<void> {
const execution = this.repo.create({ ...record, finishedAt: null });
// The cast is needed because the insert payload type recurses into the
// opaque `graph` jsonb and rejects `StepConfig`'s deliberate `unknown`.
// NOTE: prefer insert to save for performance reasons.
await this.repo.insert(execution as InsertValues);
return { id: execution.id };
}
async loadExecution(id: string): Promise<ExecutionRecord> {
@@ -11,6 +11,7 @@ import {
WorkflowExecution,
WorkflowStepExecution,
} from '../../database';
import { generateId } from '../../database/generate-id';
import type { WorkflowGraph } from '../../graph';
import { noopLifecycleEventPublisher } from '../../lifecycle-events';
import {
@@ -97,6 +98,7 @@ describe('execution start (integration)', () => {
workflowId: 'wf-1',
graph,
triggerOutputs: [[{ json: { hello: 'world' } }]],
executionId: generateId(),
});
await ready;
@@ -134,7 +136,9 @@ describe('execution start (integration)', () => {
noopLifecycleEventPublisher,
);
const { id: executionId } = await executionStore.createExecution({
const executionId = generateId();
await executionStore.createExecution({
id: executionId,
workflowId: 'wf-2',
status: 'queued',
mode: 'production',
@@ -17,7 +17,7 @@ function makeQueue(): WorkQueue<OrchestrationMessage> {
function makeStore(overrides: Partial<ExecutionStore> = {}): ExecutionStore {
return {
createExecution: vi.fn().mockResolvedValue({ id: 'exec-id-1' }),
createExecution: vi.fn(),
loadExecution: vi.fn(),
transitionStatus: vi.fn().mockResolvedValue(true),
finishExecution: vi.fn().mockResolvedValue(true),
@@ -26,7 +26,7 @@ function makeStore(overrides: Partial<ExecutionStore> = {}): ExecutionStore {
}
describe('StartExecutionService', () => {
it('admits, persists a queued execution, publishes execution:enqueued, returns id', async () => {
it('admits, persists a queued execution under the caller-minted id, publishes execution:enqueued', async () => {
const admittance: AdmittanceService = {
evaluate: vi.fn().mockResolvedValue({ accept: true }),
};
@@ -38,11 +38,13 @@ describe('StartExecutionService', () => {
workflowId: 'wf-1',
graph: sampleGraph,
triggerOutputs: [[{ json: { hello: 'world' } }]],
executionId: 'exec-id-1',
});
expect(result.executionId).toBe('exec-id-1');
expect(admittance.evaluate).toHaveBeenCalledWith({ workflowId: 'wf-1' });
expect(store.createExecution).toHaveBeenCalledWith({
id: 'exec-id-1',
workflowId: 'wf-1',
status: 'queued',
mode: 'production',
@@ -63,7 +65,7 @@ describe('StartExecutionService', () => {
const queue = makeQueue();
const service = new StartExecutionService(admittance, store, queue);
await service.start({ workflowId: 'wf-1', graph: sampleGraph });
await service.start({ workflowId: 'wf-1', graph: sampleGraph, executionId: 'exec-id-1' });
expect(store.createExecution).toHaveBeenCalledWith(
expect.objectContaining({ mode: 'production', triggerOutputs: null }),
@@ -77,7 +79,7 @@ describe('StartExecutionService', () => {
const validateGraph = vi.fn();
const service = new StartExecutionService(admittance, makeStore(), makeQueue(), validateGraph);
await service.start({ workflowId: 'wf-1', graph: sampleGraph });
await service.start({ workflowId: 'wf-1', graph: sampleGraph, executionId: 'exec-id-1' });
expect(validateGraph).toHaveBeenCalledExactlyOnceWith(sampleGraph);
});
@@ -94,7 +96,9 @@ describe('StartExecutionService', () => {
});
const service = new StartExecutionService(admittance, store, queue, validateGraph);
await expect(service.start({ workflowId: 'wf-1', graph: sampleGraph })).rejects.toBe(rejection);
await expect(
service.start({ workflowId: 'wf-1', graph: sampleGraph, executionId: 'exec-id-1' }),
).rejects.toBe(rejection);
expect(store.createExecution).not.toHaveBeenCalled();
expect(queue.publish).not.toHaveBeenCalled();
@@ -108,9 +112,9 @@ describe('StartExecutionService', () => {
const queue = makeQueue();
const service = new StartExecutionService(admittance, store, queue);
await expect(service.start({ workflowId: 'wf-1', graph: sampleGraph })).rejects.toBeInstanceOf(
AdmittanceRejectedError,
);
await expect(
service.start({ workflowId: 'wf-1', graph: sampleGraph, executionId: 'exec-id-1' }),
).rejects.toBeInstanceOf(AdmittanceRejectedError);
expect(store.createExecution).not.toHaveBeenCalled();
expect(queue.publish).not.toHaveBeenCalled();
@@ -12,6 +12,7 @@ import {
WorkflowExecution,
WorkflowStepExecution,
} from '../../database';
import { generateId } from '../../database/generate-id';
import type { IStepExecutor, StepExecutionRequest } from '../../dependencies';
import type { WorkflowGraph } from '../../graph';
import { noopLifecycleEventPublisher } from '../../lifecycle-events';
@@ -101,7 +102,7 @@ describe('step execution (integration)', () => {
const response = await request(runtime.app)
.post('/api/workflow-executions')
.set(authHeader())
.send({ workflowId, graph: workflowGraph, triggerOutputs })
.send({ workflowId, graph: workflowGraph, triggerOutputs, executionId: generateId() })
.expect(201);
const { executionId } = response.body as StartExecutionResult;
await finished;
@@ -381,7 +382,9 @@ describe('step execution (integration)', () => {
noopLifecycleEventPublisher,
);
const { id: executionId } = await executionStore.createExecution({
const executionId = generateId();
await executionStore.createExecution({
id: executionId,
workflowId: 'wf-2',
status: 'running',
mode: 'production',
@@ -1,8 +1,10 @@
import type { WorkflowGraph } from '../graph';
import type { ExecutionMode, ExecutionStatus, TriggerOutputs } from './execution.types';
/** A new execution to persist. `id` and timestamps are assigned by the store. */
/** A new execution to persist. Timestamps are assigned by the store. */
export interface NewExecutionRecord {
/** Caller-minted id. The store never mints one. */
id: string;
workflowId: string;
status: ExecutionStatus;
mode: ExecutionMode;
@@ -34,8 +36,8 @@ export class ExecutionNotFoundError extends Error {
/** Persistence interface for executions. */
export interface ExecutionStore {
/** Persist a new execution record; returns its generated id. */
createExecution(record: NewExecutionRecord): Promise<{ id: string }>;
/** Persist a new execution record under the caller-minted `record.id`. */
createExecution(record: NewExecutionRecord): Promise<void>;
/** Load a full execution by id. Throws `ExecutionNotFoundError` if absent. */
loadExecution(id: string): Promise<ExecutionRecord>;
@@ -10,6 +10,11 @@ export interface StartExecutionRequest {
/** Trigger step's output slots, one entry per output. */
triggerOutputs?: TriggerOutputs | null;
mode?: ExecutionMode;
/**
* Caller-minted, so the caller can record state against the run before it
* starts. The engine never mints one.
*/
executionId: string;
}
export interface StartExecutionResult {
@@ -34,7 +39,12 @@ export class StartExecutionService {
throw new AdmittanceRejectedError(decision.reason);
}
const { id } = await this.executionStore.createExecution({
// The caller's id is authoritative: it already has a session registered
// against it, so the store never gets to rename the run.
const { executionId } = request;
await this.executionStore.createExecution({
id: executionId,
workflowId: request.workflowId,
// admitted; a worker flips this to 'running' when it starts
status: 'queued',
@@ -48,9 +58,9 @@ export class StartExecutionService {
// reconciliation sweep (not yet built) re-dispatches it.
await this.workQueue.publish({
type: 'execution:enqueued',
executionId: id,
executionId,
});
return { executionId: id };
return { executionId };
}
}
@@ -7,6 +7,7 @@ import { afterAll, afterEach, beforeAll, beforeEach, describe, expect, it, vi }
import { AllowAllAdmittance } from '../../admittance';
import { mintIdentityToken, SharedSecretIdentityVerifier } from '../../auth';
import { createDataSource, createStores, WorkflowExecution } from '../../database';
import { generateId } from '../../database/generate-id';
import { ExecutionQueryService, StartExecutionService } from '../../execution';
import type { WorkflowGraph } from '../../graph';
import type { OrchestrationMessage, WorkQueue } from '../../queue';
@@ -24,6 +25,14 @@ const authHeader = () => ({
authorization: `Bearer ${mintIdentityToken(secret, { cpId: 'cp-1', tenantId: 'tenant-1' })}`,
});
/** The caller always mints the id, so every valid body carries one. */
const startBody = (overrides: Record<string, unknown> = {}) => ({
workflowId: 'wf-1',
graph: sampleGraph,
executionId: generateId(),
...overrides,
});
let container: StartedPostgreSqlContainer;
let dataSource: DataSource;
let workQueue: WorkQueue<OrchestrationMessage>;
@@ -59,19 +68,17 @@ afterAll(async () => {
});
describe('POST /api/workflow-executions (integration)', () => {
it('creates an execution row, publishes execution:enqueued, returns 201', async () => {
it('creates the row under the caller-minted id, publishes execution:enqueued, returns 201', async () => {
const body = startBody({ triggerOutputs: [[{ json: { hello: 'world' } }]] });
const response = await request(url)
.post('/api/workflow-executions')
.set(authHeader())
.send({
workflowId: 'wf-1',
graph: sampleGraph,
triggerOutputs: [[{ json: { hello: 'world' } }]],
});
.send(body);
expect(response.status).toBe(201);
const { executionId } = response.body as { executionId: string };
expect(executionId).toBeTruthy();
expect(executionId).toBe(body.executionId);
const repo = dataSource.getRepository(WorkflowExecution);
// `findOne({ where })`, not `findOneByOrFail`: the latter's overload exceeds
@@ -86,6 +93,21 @@ describe('POST /api/workflow-executions (integration)', () => {
expect(workQueue.publish).toHaveBeenCalledWith({ type: 'execution:enqueued', executionId });
});
// v4 included: the id has to be time-ordered. `undefined` covers the omitted
// case — the engine never mints a replacement.
it.each(['not-a-uuid', '9f1b7d0e-2c4a-4f8b-9d3e-6a5c1b2d3e4f', undefined])(
'rejects the execution id %p with 400',
async (executionId) => {
const response = await request(url)
.post('/api/workflow-executions')
.set(authHeader())
.send({ workflowId: 'wf-1', graph: sampleGraph, executionId });
expect(response.status).toBe(400);
expect((response.body as { error: string }).error).toBe('invalid_request');
},
);
it('rejects an invalid body with 400', async () => {
const response = await request(url)
.post('/api/workflow-executions')
@@ -100,11 +122,7 @@ describe('POST /api/workflow-executions (integration)', () => {
const response = await request(url)
.post('/api/workflow-executions')
.set(authHeader())
.send({
workflowId: 'wf-1',
graph: sampleGraph,
triggerOutputs: { hello: 'world' },
});
.send(startBody({ triggerOutputs: { hello: 'world' } }));
expect(response.status).toBe(400);
expect((response.body as { error: string }).error).toBe('invalid_request');
@@ -114,18 +132,17 @@ describe('POST /api/workflow-executions (integration)', () => {
const response = await request(url)
.post('/api/workflow-executions')
.set(authHeader())
.send({ workflowId: 'wf-1', graph: sampleGraph, triggerOutputs });
.send(startBody({ triggerOutputs }));
expect(response.status).toBe(400);
expect((response.body as { error: string }).error).toBe('invalid_request');
});
it('rejects an empty-array triggerOutputs with 400 (send null or omit for "no payload")', async () => {
const response = await request(url).post('/api/workflow-executions').set(authHeader()).send({
workflowId: 'wf-1',
graph: sampleGraph,
triggerOutputs: [],
});
const response = await request(url)
.post('/api/workflow-executions')
.set(authHeader())
.send(startBody({ triggerOutputs: [] }));
expect(response.status).toBe(400);
expect((response.body as { error: string }).error).toBe('invalid_request');
@@ -135,11 +152,7 @@ describe('POST /api/workflow-executions (integration)', () => {
const response = await request(url)
.post('/api/workflow-executions')
.set(authHeader())
.send({
workflowId: 'wf-1',
graph: sampleGraph,
triggerOutputs: Array.from({ length: 102 }, () => []),
});
.send(startBody({ triggerOutputs: Array.from({ length: 102 }, () => []) }));
expect(response.status).toBe(400);
expect((response.body as { error: string }).error).toBe('invalid_request');
@@ -149,10 +162,7 @@ describe('POST /api/workflow-executions (integration)', () => {
const response = await request(url)
.post('/api/workflow-executions')
.set(authHeader())
.send({
workflowId: 'wf-1',
graph: { nodes: [{ id: 'a', name: 'A', type: 'v1-node' }], edges: [] },
});
.send(startBody({ graph: { nodes: [{ id: 'a', name: 'A', type: 'v1-node' }], edges: [] } }));
expect(response.status).toBe(400);
expect((response.body as { error: string }).error).toBe('invalid_graph');
@@ -165,7 +175,7 @@ describe('POST /api/workflow-executions (integration)', () => {
.post('/api/workflow-executions')
.set(authHeader())
.send({
workflowId: 'wf-1',
...startBody(),
graph: {
nodes: [
{ id: 'trigger', name: 'T', type: 'trigger' },
@@ -189,7 +199,7 @@ describe('POST /api/workflow-executions (integration)', () => {
.post('/api/workflow-executions')
.set(authHeader())
.send({
workflowId: 'wf-1',
...startBody(),
graph: {
nodes: [
{ id: 'trigger', name: 'T', type: 'trigger' },
@@ -222,11 +232,7 @@ describe('GET /api/workflow-executions/:id (integration)', () => {
const response = await request(url)
.post('/api/workflow-executions')
.set(authHeader())
.send({
workflowId: 'wf-1',
graph: sampleGraph,
triggerOutputs: [[{ json: { hello: 'world' } }]],
});
.send(startBody({ triggerOutputs: [[{ json: { hello: 'world' } }]] }));
return (response.body as { executionId: string }).executionId;
}
@@ -293,11 +299,7 @@ describe('GET /api/workflow-executions/:id (integration)', () => {
const postResponse = await request(runtime.app)
.post('/api/workflow-executions')
.set(authHeader())
.send({
workflowId: 'wf-1',
graph: sampleGraph,
triggerOutputs: [[{ json: { hello: 'world' } }]],
});
.send(startBody({ triggerOutputs: [[{ json: { hello: 'world' } }]] }));
const { executionId } = postResponse.body as { executionId: string };
await finished;
@@ -1,3 +1,4 @@
import { UUID_V7_PATTERN } from '@n8n/constants';
import { Router, type Router as RouterType } from 'express';
import { z } from 'zod';
@@ -44,6 +45,8 @@ const StartExecutionBody = z.object({
/** Trigger output slots. Empty means "no payload" — send `null` or omit instead. */
triggerOutputs: z.array(jsonValueSchema).min(1).max(MAX_TRIGGER_SLOTS).nullable().optional(),
mode: z.enum(['production', 'manual']).optional(),
/** The caller mints the id. v7 only, so ids stay time-ordered. */
executionId: z.string().regex(UUID_V7_PATTERN),
});
export function createWorkflowExecutionsRouter(deps: EngineServerDeps): RouterType {
@@ -35,6 +35,7 @@
"supertest": "^7.1.1",
"testcontainers": "catalog:",
"typescript": "catalog:typescript",
"uuid": "catalog:",
"vitest": "catalog:"
},
"license": "LicenseRef-n8n-sustainable-use"
@@ -23,6 +23,7 @@ import { SplitOut } from 'n8n-nodes-base/nodes/Transform/SplitOut/SplitOut.node'
import type { IDataObject, INodeType, INodeTypes, IVersionedNodeType } from 'n8n-workflow';
import { NodeHelpers } from 'n8n-workflow';
import request from 'supertest';
import { v7 as uuidv7 } from 'uuid';
import { vi } from 'vitest';
import { createEngineStepDataLoader } from '../engine-step-data-loader';
@@ -345,7 +346,8 @@ export function makeRunWorkflow(getDataSource: () => EngineDataSource) {
const response = await request(runtime.app)
.post('/api/workflow-executions')
.set('Authorization', `Bearer ${mintIdentityToken(authSecret, caller)}`)
.send({ workflowId: 'wf-m1', graph, triggerOutputs, mode })
// The caller mints the execution id; the engine never mints one.
.send({ workflowId: 'wf-m1', graph, triggerOutputs, mode, executionId: uuidv7() })
.expect(201);
const { executionId } = response.body as StartExecutionResult;
+28 -1
View File
@@ -31,6 +31,13 @@ const instanceAiLazyRuntimeImports = [
message: INSTANCE_AI_LAZY_IMPORT_MESSAGE,
}));
const engineV2ModuleOnlyImport = {
name: '@n8n/engine',
allowTypeImports: true,
message:
'Only src/modules/engine-v2/** may import @n8n/engine at runtime. Use a type import, or reach the engine through EngineDataPlaneProxyService.',
};
export default defineConfig(
globalIgnores(['scripts/**/*.mjs', 'vitest.*.ts', 'coverage/**']),
nodeConfig,
@@ -177,13 +184,33 @@ export default defineConfig(
'n8n-local-rules/require-public-api-controller': 'off',
},
},
{
files: ['./src/**/*.ts'],
ignores: ['./src/modules/engine-v2/**/*.ts'],
rules: {
// Repeats the policy restriction: a later block replaces the rule's options
// wholesale rather than merging them.
'@typescript-eslint/no-restricted-imports': [
'error',
{ paths: [POLICY_INTERNAL_RESTRICTION, engineV2ModuleOnlyImport] },
],
},
},
{
files: ['./src/modules/instance-ai/**/*.ts'],
ignores: ['./src/modules/instance-ai/**/__tests__/**/*.ts'],
rules: {
// Repeats the engine restriction: a later block replaces the rule's options
// wholesale rather than merging them.
'@typescript-eslint/no-restricted-imports': [
'error',
{ paths: [POLICY_INTERNAL_RESTRICTION, ...instanceAiLazyRuntimeImports] },
{
paths: [
POLICY_INTERNAL_RESTRICTION,
...instanceAiLazyRuntimeImports,
engineV2ModuleOnlyImport,
],
},
],
},
},
@@ -1,6 +1,7 @@
import { UUID_V7_PATTERN } from '@n8n/constants';
import { v7 as uuidv7 } from 'uuid';
import { isExecutionIdV2 } from '../execution-id';
import { createExecutionIdV2, isExecutionIdV2 } from '../execution-id';
describe('isExecutionIdV2', () => {
it('should accept an id the engine mints', () => {
@@ -19,3 +20,14 @@ describe('isExecutionIdV2', () => {
},
);
});
describe('createExecutionIdV2', () => {
// The engine rejects any other shape on the wire, so the format is the contract.
it('should mint an id the engine accepts', () => {
expect(createExecutionIdV2()).toMatch(UUID_V7_PATTERN);
});
it('should mint a distinct id each call', () => {
expect(createExecutionIdV2()).not.toBe(createExecutionIdV2());
});
});
@@ -1,3 +1,5 @@
import { v7 as uuidv7 } from 'uuid';
/** Tagged so a v1 id cannot reach a function that only handles v2 executions. */
export type ExecutionIdV2 = string & { readonly __brand: 'ExecutionIdV2' };
@@ -6,3 +8,11 @@ const UUID = /^[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}$/i;
/** A v1 id is numeric and a v2 id is a UUID, so the shape alone picks the backend. */
export const isExecutionIdV2 = (id: string): id is ExecutionIdV2 => UUID.test(id);
/**
* Mints an id for a run the control plane hands to the engine.
*
* The engine validates this shape on the wire and rejects anything else, so the
* format is a contract, not a reason to import the engine's generator here.
*/
export const createExecutionIdV2 = (): ExecutionIdV2 => uuidv7() as ExecutionIdV2;
@@ -21,6 +21,7 @@ describe('EngineDataPlaneClient', () => {
const request: StartExecutionRequest = {
workflowId: 'wf-1',
graph: { nodes: [], edges: [] },
executionId: EXECUTION_ID,
};
let http: HttpRequestClient;
@@ -46,7 +46,11 @@ describe('EngineV2Module', () => {
it('registers the client as the data plane provider', async () => {
const proxy = Container.get(EngineDataPlaneProxyService);
const request = { workflowId: 'wf-1', graph: { nodes: [], edges: [] } };
const request = {
workflowId: 'wf-1',
graph: { nodes: [], edges: [] },
executionId: '01a038ae-c4a8-7799-8a3e-e3c2ca055cfa',
};
await expect(proxy.startExecution(request)).rejects.toThrow('N8N_ENABLED_MODULES');
await module.init();
@@ -8,13 +8,14 @@ import type { EngineDataPlaneProvider } from '../engine-data-plane-proxy.service
import { EngineDataPlaneProxyService } from '../engine-data-plane-proxy.service';
describe('EngineDataPlaneProxyService', () => {
const executionId = '01a038ae-c4a8-7799-8a3e-e3c2ca055cfa' as ExecutionIdV2;
const request: StartExecutionRequest = {
workflowId: 'wf-1',
graph: { nodes: [], edges: [] },
executionId,
};
const executionId = '01a038ae-c4a8-7799-8a3e-e3c2ca055cfa' as ExecutionIdV2;
let proxy: EngineDataPlaneProxyService;
beforeEach(() => {
@@ -1,3 +1,4 @@
import { UUID_V7_PATTERN } from '@n8n/constants';
import type {
INode,
IPinData,
@@ -112,11 +113,12 @@ describe('EngineV2Dispatcher', () => {
});
describe('start', () => {
it('returns the data plane execution id', async () => {
await expect(dispatcher.start(runData())).resolves.toBe('dp-uuid');
it('mints the execution id and sends it to the data plane', async () => {
const executionId = await dispatcher.start(runData());
expect(executionId).toMatch(UUID_V7_PATTERN);
expect(proxy.startExecution).toHaveBeenCalledWith(
expect.objectContaining({ workflowId: 'wf-1', mode: 'manual' }),
expect.objectContaining({ executionId, workflowId: 'wf-1', mode: 'manual' }),
);
});
@@ -348,16 +350,30 @@ describe('EngineV2Dispatcher', () => {
});
describe('the push session', () => {
it('records the run against the data plane execution id', async () => {
await dispatcher.start(runData({ pushRef: 'push-1' }));
it('records the run against the minted execution id', async () => {
const executionId = await dispatcher.start(runData({ pushRef: 'push-1' }));
expect(pushRegistry.register).toHaveBeenCalledExactlyOnceWith('dp-uuid', {
expect(pushRegistry.register).toHaveBeenCalledExactlyOnceWith(executionId, {
pushRef: 'push-1',
workflowId: 'wf-1',
trigger: { nodeName: MANUAL_TRIGGER.name, outputs: [[{ json: {} }]] },
});
});
it('records the run before it dispatches, so no event can arrive first', async () => {
let registeredBeforeDispatch = false;
proxy.startExecution.mockImplementationOnce(async ({ executionId }) => {
registeredBeforeDispatch = pushRegistry.register.mock.calls.some(
([id]) => id === executionId,
);
return { executionId };
});
await dispatcher.start(runData({ pushRef: 'push-1' }));
expect(registeredBeforeDispatch).toBe(true);
});
it('records the trigger payload the engine was given', async () => {
const data = runData({
pushRef: 'push-1',
@@ -381,12 +397,13 @@ describe('EngineV2Dispatcher', () => {
expect(pushRegistry.register).not.toHaveBeenCalled();
});
it('records nothing when the data plane refused the run', async () => {
it('releases the session when the data plane refused the run', async () => {
proxy.startExecution.mockRejectedValueOnce(new Error('down'));
await expect(dispatcher.start(runData({ pushRef: 'push-1' }))).rejects.toThrow('down');
expect(pushRegistry.register).not.toHaveBeenCalled();
const [executionId] = pushRegistry.register.mock.calls[0];
expect(pushRegistry.release).toHaveBeenCalledExactlyOnceWith(executionId);
});
});
});
@@ -3,6 +3,7 @@ import type { StepSlots, TriggerOutputs } from '@n8n/engine';
import type { INodeExecutionData, IWorkflowExecutionDataProcess } from 'n8n-workflow';
import { isTriggerNodeType, MANUAL_TRIGGER_NODE_TYPE, UserError } from 'n8n-workflow';
import { createExecutionIdV2 } from '@/executions/execution-id';
import { CredentialsPermissionChecker } from '@/executions/pre-execution-checks';
import type { ResumableExecution } from '@/interfaces';
import { EngineDataPlaneProxyService } from '@/services/engine-data-plane-proxy.service';
@@ -23,7 +24,8 @@ const DEFAULT_MAIN_OUTPUT: INodeExecutionData[][] = [[{ json: {} }]];
* not pick.
*
* No control-plane execution row is created: the data plane is the source of
* truth, and the returned execution id is its UUID.
* truth for the run. Only the execution id is minted here, so the push session
* can be recorded before dispatch.
*/
@Service()
export class EngineV2Dispatcher {
@@ -49,7 +51,7 @@ export class EngineV2Dispatcher {
);
}
/** Returns the data plane's execution id. */
/** Returns the execution id this dispatch minted. */
async start(data: IWorkflowExecutionDataProcess): Promise<string> {
this.assertSupported(data);
@@ -64,24 +66,30 @@ export class EngineV2Dispatcher {
const graph = new V1WorkflowConverter().convert(workflowData, data.triggerToStartFrom?.name);
const triggerMain = this.triggerMainOutputs(data);
const { executionId } = await this.proxy.startExecution({
workflowId: workflowData.id,
graph,
triggerOutputs: this.toTriggerOutputs(triggerMain, toStepOutputs),
mode: 'manual',
});
// TODO(CAT-4255): the engine can publish lifecycle events before this line
// runs, and the relay drops them because no session exists yet. Let the
// control plane mint the execution id so this can register before dispatch.
const executionId = createExecutionIdV2();
// At the session cap this can evict another run's session, uncaught below. Rare; not worth fixing.
this.registerPushSession(executionId, data, triggerMain);
try {
await this.proxy.startExecution({
executionId,
workflowId: workflowData.id,
graph,
triggerOutputs: this.toTriggerOutputs(triggerMain, toStepOutputs),
mode: 'manual',
});
} catch (error) {
// Assumes rejection: a dropped success response also releases a still-live session.
this.pushRegistry.release(executionId);
throw error;
}
return executionId;
}
/**
* Lifecycle events carry no session id, so the push ref is recorded here,
* keyed by execution id, before any events can arrive.
* keyed by execution id.
*/
private registerPushSession(
executionId: string,
@@ -12,6 +12,7 @@ import {
testDb,
mockInstance,
} from '@n8n/backend-test-utils';
import { UUID_V7_PATTERN } from '@n8n/constants';
import type {
User,
ListQueryDb,
@@ -4983,6 +4984,8 @@ describe('POST /workflows/:workflowId/run', () => {
});
beforeEach(() => {
// Deliberately not the id the control plane minted, so a response echoing
// the data plane back would fail the assertion below.
startExecution.mockResolvedValue({ executionId: 'a3c1e0f2-0000-4000-8000-000000000001' });
});
@@ -5023,9 +5026,14 @@ describe('POST /workflows/:workflowId/run', () => {
.send({ triggerToStartFrom: { name: TRIGGER_NAME } });
expect(response.statusCode).toBe(200);
expect(response.body.data.executionId).toBe('a3c1e0f2-0000-4000-8000-000000000001');
// The control plane mints the id, dispatches under it, and reports that
// same id — the data plane's response never renames the run.
const { executionId } = response.body.data;
expect(executionId).toMatch(UUID_V7_PATTERN);
expect(startExecution).toHaveBeenCalledWith(
objectContaining({
executionId,
workflowId: dbWorkflow.id,
mode: 'manual',
triggerOutputs: [[{ json: {} }]],
+6
View File
@@ -2290,6 +2290,9 @@ importers:
'@n8n/config':
specifier: workspace:*
version: link:../config
'@n8n/constants':
specifier: workspace:*
version: link:../constants
'@n8n/di':
specifier: workspace:*
version: link:../di
@@ -3190,6 +3193,9 @@ importers:
typescript:
specifier: catalog:typescript
version: 7.0.2
uuid:
specifier: 'catalog:'
version: 11.1.1
vitest:
specifier: 'catalog:'
version: 4.1.9(@opentelemetry/api@1.9.1)(@types/node@20.19.41)(@vitest/browser-playwright@4.1.9)(@vitest/coverage-v8@4.1.9)(jsdom@23.0.1(bufferutil@4.0.9)(supports-color@8.1.1)(utf-8-validate@5.0.10))(vite@8.0.2(@emnapi/core@1.11.1)(@emnapi/runtime@1.11.1)(@types/node@20.19.41)(esbuild@0.28.1)(jiti@2.6.1)(sass-embedded@1.98.0)(sass@1.98.0)(terser@5.16.1)(tsx@4.19.3)(yaml@2.8.3))