diff --git a/packages/@n8n/constants/src/index.ts b/packages/@n8n/constants/src/index.ts index cd8de86afb8..ac28ca2fa21 100644 --- a/packages/@n8n/constants/src/index.ts +++ b/packages/@n8n/constants/src/index.ts @@ -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', diff --git a/packages/@n8n/constants/src/uuid.ts b/packages/@n8n/constants/src/uuid.ts new file mode 100644 index 00000000000..e71e20e4634 --- /dev/null +++ b/packages/@n8n/constants/src/uuid.ts @@ -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; diff --git a/packages/@n8n/engine/package.json b/packages/@n8n/engine/package.json index 991c9afa5c4..c5a2103a2c8 100644 --- a/packages/@n8n/engine/package.json +++ b/packages/@n8n/engine/package.json @@ -25,6 +25,7 @@ ], "dependencies": { "@n8n/config": "workspace:*", + "@n8n/constants": "workspace:*", "@n8n/di": "workspace:*", "@n8n/typeorm": "workspace:*", "express": "catalog:", diff --git a/packages/@n8n/engine/src/database/__tests__/workflow-execution.integration.test.ts b/packages/@n8n/engine/src/database/__tests__/workflow-execution.integration.test.ts index 3952dd5ae7a..a4f47e25ad3 100644 --- a/packages/@n8n/engine/src/database/__tests__/workflow-execution.integration.test.ts +++ b/packages/@n8n/engine/src/database/__tests__/workflow-execution.integration.test.ts @@ -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', diff --git a/packages/@n8n/engine/src/database/__tests__/workflow-step-execution.integration.test.ts b/packages/@n8n/engine/src/database/__tests__/workflow-step-execution.integration.test.ts index 7c7a1cb844d..06fda3549d8 100644 --- a/packages/@n8n/engine/src/database/__tests__/workflow-step-execution.integration.test.ts +++ b/packages/@n8n/engine/src/database/__tests__/workflow-step-execution.integration.test.ts @@ -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 { const repo = dataSource.getRepository(WorkflowExecution); const execution = repo.create({ + id: generateId(), workflowId: 'wf-1', status: 'running', mode: 'production', diff --git a/packages/@n8n/engine/src/database/entities/workflow-execution.entity.ts b/packages/@n8n/engine/src/database/entities/workflow-execution.entity.ts index 9debbb82512..a1a3835c9e5 100644 --- a/packages/@n8n/engine/src/database/entities/workflow-execution.entity.ts +++ b/packages/@n8n/engine/src/database/entities/workflow-execution.entity.ts @@ -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(); - } } diff --git a/packages/@n8n/engine/src/database/typeorm-execution-store.ts b/packages/@n8n/engine/src/database/typeorm-execution-store.ts index b585dcf13b2..5f6f0861531 100644 --- a/packages/@n8n/engine/src/database/typeorm-execution-store.ts +++ b/packages/@n8n/engine/src/database/typeorm-execution-store.ts @@ -19,13 +19,12 @@ type InsertValues = Parameters['insert']>[0]; export class TypeOrmExecutionStore implements ExecutionStore { constructor(private readonly repo: Repository) {} - async createExecution(record: NewExecutionRecord): Promise<{ id: string }> { + async createExecution(record: NewExecutionRecord): Promise { 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 { diff --git a/packages/@n8n/engine/src/execution/__tests__/execution-start.integration.test.ts b/packages/@n8n/engine/src/execution/__tests__/execution-start.integration.test.ts index 9913c8b8886..e1c0af53019 100644 --- a/packages/@n8n/engine/src/execution/__tests__/execution-start.integration.test.ts +++ b/packages/@n8n/engine/src/execution/__tests__/execution-start.integration.test.ts @@ -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', diff --git a/packages/@n8n/engine/src/execution/__tests__/start-execution.service.test.ts b/packages/@n8n/engine/src/execution/__tests__/start-execution.service.test.ts index 48dd8fd3b4f..b6dae8f48dc 100644 --- a/packages/@n8n/engine/src/execution/__tests__/start-execution.service.test.ts +++ b/packages/@n8n/engine/src/execution/__tests__/start-execution.service.test.ts @@ -17,7 +17,7 @@ function makeQueue(): WorkQueue { function makeStore(overrides: Partial = {}): 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 { } 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(); diff --git a/packages/@n8n/engine/src/execution/__tests__/step-execution.integration.test.ts b/packages/@n8n/engine/src/execution/__tests__/step-execution.integration.test.ts index 708f3f61f35..fd0be4a3f28 100644 --- a/packages/@n8n/engine/src/execution/__tests__/step-execution.integration.test.ts +++ b/packages/@n8n/engine/src/execution/__tests__/step-execution.integration.test.ts @@ -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', diff --git a/packages/@n8n/engine/src/execution/execution-store.ts b/packages/@n8n/engine/src/execution/execution-store.ts index 5551eab8475..9703c1b6e12 100644 --- a/packages/@n8n/engine/src/execution/execution-store.ts +++ b/packages/@n8n/engine/src/execution/execution-store.ts @@ -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; /** Load a full execution by id. Throws `ExecutionNotFoundError` if absent. */ loadExecution(id: string): Promise; diff --git a/packages/@n8n/engine/src/execution/start-execution.service.ts b/packages/@n8n/engine/src/execution/start-execution.service.ts index fcf0c7b7b0a..8d4aca41ee0 100644 --- a/packages/@n8n/engine/src/execution/start-execution.service.ts +++ b/packages/@n8n/engine/src/execution/start-execution.service.ts @@ -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 }; } } diff --git a/packages/@n8n/engine/src/server/__tests__/workflow-executions.integration.test.ts b/packages/@n8n/engine/src/server/__tests__/workflow-executions.integration.test.ts index 7355be902d7..ba7adc4a152 100644 --- a/packages/@n8n/engine/src/server/__tests__/workflow-executions.integration.test.ts +++ b/packages/@n8n/engine/src/server/__tests__/workflow-executions.integration.test.ts @@ -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 = {}) => ({ + workflowId: 'wf-1', + graph: sampleGraph, + executionId: generateId(), + ...overrides, +}); + let container: StartedPostgreSqlContainer; let dataSource: DataSource; let workQueue: WorkQueue; @@ -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; diff --git a/packages/@n8n/engine/src/server/routes/workflow-executions.ts b/packages/@n8n/engine/src/server/routes/workflow-executions.ts index 8f644da8b55..aa8f562bafb 100644 --- a/packages/@n8n/engine/src/server/routes/workflow-executions.ts +++ b/packages/@n8n/engine/src/server/routes/workflow-executions.ts @@ -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 { diff --git a/packages/@n8n/node-engine-compatibility/package.json b/packages/@n8n/node-engine-compatibility/package.json index 56e3a62814c..18783acd7b8 100644 --- a/packages/@n8n/node-engine-compatibility/package.json +++ b/packages/@n8n/node-engine-compatibility/package.json @@ -35,6 +35,7 @@ "supertest": "^7.1.1", "testcontainers": "catalog:", "typescript": "catalog:typescript", + "uuid": "catalog:", "vitest": "catalog:" }, "license": "LicenseRef-n8n-sustainable-use" diff --git a/packages/@n8n/node-engine-compatibility/src/__tests__/acceptance-fixtures.ts b/packages/@n8n/node-engine-compatibility/src/__tests__/acceptance-fixtures.ts index dbec4a9c7ad..e12bd544161 100644 --- a/packages/@n8n/node-engine-compatibility/src/__tests__/acceptance-fixtures.ts +++ b/packages/@n8n/node-engine-compatibility/src/__tests__/acceptance-fixtures.ts @@ -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; diff --git a/packages/cli/eslint.config.mjs b/packages/cli/eslint.config.mjs index 7d002955bfa..8d0afb0eed5 100644 --- a/packages/cli/eslint.config.mjs +++ b/packages/cli/eslint.config.mjs @@ -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, + ], + }, ], }, }, diff --git a/packages/cli/src/executions/__tests__/execution-id.test.ts b/packages/cli/src/executions/__tests__/execution-id.test.ts index 7e1288a8598..82b496e498b 100644 --- a/packages/cli/src/executions/__tests__/execution-id.test.ts +++ b/packages/cli/src/executions/__tests__/execution-id.test.ts @@ -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()); + }); +}); diff --git a/packages/cli/src/executions/execution-id.ts b/packages/cli/src/executions/execution-id.ts index 5798f874689..1b2043f08b6 100644 --- a/packages/cli/src/executions/execution-id.ts +++ b/packages/cli/src/executions/execution-id.ts @@ -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; diff --git a/packages/cli/src/modules/engine-v2/__tests__/engine-data-plane-client.test.ts b/packages/cli/src/modules/engine-v2/__tests__/engine-data-plane-client.test.ts index 0d760ffb268..d1980009096 100644 --- a/packages/cli/src/modules/engine-v2/__tests__/engine-data-plane-client.test.ts +++ b/packages/cli/src/modules/engine-v2/__tests__/engine-data-plane-client.test.ts @@ -21,6 +21,7 @@ describe('EngineDataPlaneClient', () => { const request: StartExecutionRequest = { workflowId: 'wf-1', graph: { nodes: [], edges: [] }, + executionId: EXECUTION_ID, }; let http: HttpRequestClient; diff --git a/packages/cli/src/modules/engine-v2/__tests__/engine-v2.module.test.ts b/packages/cli/src/modules/engine-v2/__tests__/engine-v2.module.test.ts index 6dedd1d28a2..fbe2d97b4cf 100644 --- a/packages/cli/src/modules/engine-v2/__tests__/engine-v2.module.test.ts +++ b/packages/cli/src/modules/engine-v2/__tests__/engine-v2.module.test.ts @@ -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(); diff --git a/packages/cli/src/services/__tests__/engine-data-plane-proxy.service.test.ts b/packages/cli/src/services/__tests__/engine-data-plane-proxy.service.test.ts index 58bf59acd84..220f7fbc7fe 100644 --- a/packages/cli/src/services/__tests__/engine-data-plane-proxy.service.test.ts +++ b/packages/cli/src/services/__tests__/engine-data-plane-proxy.service.test.ts @@ -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(() => { diff --git a/packages/cli/src/services/__tests__/engine-v2-dispatcher.service.test.ts b/packages/cli/src/services/__tests__/engine-v2-dispatcher.service.test.ts index 30a39c79a6f..25a32ed2b07 100644 --- a/packages/cli/src/services/__tests__/engine-v2-dispatcher.service.test.ts +++ b/packages/cli/src/services/__tests__/engine-v2-dispatcher.service.test.ts @@ -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); }); }); }); diff --git a/packages/cli/src/services/engine-v2-dispatcher.service.ts b/packages/cli/src/services/engine-v2-dispatcher.service.ts index ca3ab8e8803..be875ab91ae 100644 --- a/packages/cli/src/services/engine-v2-dispatcher.service.ts +++ b/packages/cli/src/services/engine-v2-dispatcher.service.ts @@ -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 { 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, diff --git a/packages/cli/test/integration/workflows/workflows.controller.test.ts b/packages/cli/test/integration/workflows/workflows.controller.test.ts index a899b0fa9df..1e2b52bf088 100644 --- a/packages/cli/test/integration/workflows/workflows.controller.test.ts +++ b/packages/cli/test/integration/workflows/workflows.controller.test.ts @@ -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: {} }]], diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index a4e76b4bfe1..bb2eb4c1003 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -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))