Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -406,4 +406,48 @@ describe('ExecutionRepository', () => {
expect(result.data.resultData.error?.message).toBe('The execution was cancelled manually');
});
});

describe('setRunning', () => {
beforeEach(() => {
entityManager.transaction.mockImplementation(async (fn: unknown) => {
return await (fn as (em: typeof entityManager) => Promise<unknown>)(entityManager);
});
});

test('should set startedAt when not already set', async () => {
const executionId = '123';

entityManager.findOneBy.mockResolvedValueOnce({ startedAt: null });

const result = await executionRepository.setRunning(executionId);

expect(entityManager.transaction).toHaveBeenCalled();
expect(entityManager.findOneBy).toHaveBeenCalledWith(ExecutionEntity, {
id: executionId,
});
expect(entityManager.update).toHaveBeenCalledWith(
ExecutionEntity,
{ id: executionId },
{ status: 'running', startedAt: expect.any(Date) },
);
expect(result).toBeInstanceOf(Date);
});

test('should preserve existing startedAt for resumed executions', async () => {
const executionId = '456';
const existingStartedAt = new Date('2025-12-02T09:04:47.150Z');

entityManager.findOneBy.mockResolvedValueOnce({ startedAt: existingStartedAt });

const result = await executionRepository.setRunning(executionId);

expect(entityManager.transaction).toHaveBeenCalled();
expect(entityManager.update).toHaveBeenCalledWith(
ExecutionEntity,
{ id: executionId },
{ status: 'running', startedAt: existingStartedAt },
);
expect(result).toBe(existingStartedAt);
});
});
});
15 changes: 13 additions & 2 deletions packages/@n8n/db/src/repositories/execution.repository.ts
Original file line number Diff line number Diff line change
Expand Up @@ -369,9 +369,20 @@ export class ExecutionRepository extends Repository<ExecutionEntity> {
async setRunning(executionId: string) {
const startedAt = new Date();

await this.update({ id: executionId }, { status: 'running', startedAt });
return await this.manager.transaction(async (manager) => {
const existing = await manager.findOneBy(ExecutionEntity, { id: executionId });

return startedAt;
// Preserve original startedAt for resumed executions
const effectiveStartedAt = existing?.startedAt ?? startedAt;

await manager.update(
ExecutionEntity,
{ id: executionId },
{ status: 'running', startedAt: effectiveStartedAt },
);

return effectiveStartedAt;
});
}

/**
Expand Down
23 changes: 23 additions & 0 deletions packages/cli/src/__tests__/wait-tracker.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -146,6 +146,29 @@ describe('WaitTracker', () => {
);
});

it('should preserve original startedAt timestamp when resuming execution', async () => {
const originalStartedAt = new Date('2025-12-02T09:04:47.150Z');
const executionWithStartedAt = {
...execution,
startedAt: originalStartedAt,
};

executionRepository.findSingleExecution
.calledWith(execution.id)
.mockResolvedValue(executionWithStartedAt);

await waitTracker.startExecution(execution.id);

expect(workflowRunner.run).toHaveBeenCalledWith(
expect.objectContaining({
startedAt: originalStartedAt,
}),
false,
false,
execution.id,
);
});

describe('parent execution with waiting sub-workflow', () => {
const setupParentExecutionTest = (shouldResume: boolean | undefined) => {
const parentExecution = mock<IExecutionResponse>({
Expand Down
9 changes: 8 additions & 1 deletion packages/cli/src/__tests__/workflow-runner.test.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,12 @@
import { testDb, createWorkflow, mockInstance } from '@n8n/backend-test-utils';
import { GlobalConfig } from '@n8n/config';
import { type User, type ExecutionEntity, GLOBAL_OWNER_ROLE, Project } from '@n8n/db';
import {
type User,
type ExecutionEntity,
GLOBAL_OWNER_ROLE,
Project,
ExecutionRepository,
} from '@n8n/db';
import { Container, Service } from '@n8n/di';
import type { Response } from 'express';
import { mock } from 'jest-mock-extended';
Expand Down Expand Up @@ -60,6 +66,7 @@ afterAll(() => {
beforeEach(async () => {
await testDb.truncate(['WorkflowEntity', 'SharedWorkflow']);
jest.clearAllMocks();
jest.spyOn(Container.get(ExecutionRepository), 'setRunning').mockResolvedValue(new Date());
});

describe('processError', () => {
Expand Down
32 changes: 32 additions & 0 deletions packages/cli/test/integration/execution.repository.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -67,6 +67,38 @@ describe('UserRepository', () => {
});
});

describe('setRunning', () => {
test('should set startedAt when execution has no startedAt', async () => {
const workflow = await createWorkflow();
const execution = await createExecution({ status: 'new', startedAt: null }, workflow);

const result = await executionRepository.setRunning(execution.id);

expect(result).toBeInstanceOf(Date);

const row = await executionRepository.findOneBy({ id: execution.id });
expect(row?.status).toBe('running');
expect(row?.startedAt).toEqual(result);
});

test('should preserve original startedAt for resumed executions', async () => {
const originalStartedAt = new Date('2025-12-02T09:04:47.150Z');
const workflow = await createWorkflow();
const execution = await createExecution(
{ status: 'waiting', startedAt: originalStartedAt },
workflow,
);

const result = await executionRepository.setRunning(execution.id);

expect(result.getTime()).toBe(originalStartedAt.getTime());

const row = await executionRepository.findOneBy({ id: execution.id });
expect(row?.status).toBe('running');
expect(row?.startedAt?.getTime()).toBe(originalStartedAt.getTime());
});
});

describe('updateExistingExecution with conditions', () => {
test.each([
{
Expand Down
Original file line number Diff line number Diff line change
@@ -1,5 +1,8 @@
import flatted from 'flatted';

import { test, expect } from '../../../../fixtures/base';
import executionOutOfMemoryResponse from '../../../../fixtures/execution-out-of-memory-server-response.json';
import { retryUntil } from '../../../../utils/retry-utils';

test.describe(
'Executions Filter',
Expand Down Expand Up @@ -260,6 +263,54 @@ test.describe('Workflow Executions', () => {
});
});

test.describe('execution timing', () => {
test('should preserve execution start time for standard workflow', async ({ api }) => {
const { webhookPath, workflowId, createdWorkflow } =
await api.workflows.importWorkflowFromFile('simple-webhook-test.json');

await api.workflows.activate(workflowId, createdWorkflow.versionId!);

const webhookResponse = await api.request.post(`/webhook/${webhookPath}`, { data: {} });
expect(webhookResponse.ok()).toBe(true);

const execution = await api.workflows.waitForExecution(workflowId, 10000);
const originalStartedAt = execution.startedAt;

const finalExecution = await api.workflows.getExecution(execution.id);
expect(finalExecution.startedAt).toBe(originalStartedAt);
});

test('should preserve execution start time after resuming from wait node', async ({ api }) => {
const { webhookPath, workflowId, createdWorkflow } =
await api.workflows.importWorkflowFromFile('cat-1854-wait-execution-history.json');

await api.workflows.activate(workflowId, createdWorkflow.versionId!);

const webhookResponse = await api.request.get(`/webhook/${webhookPath}`);
expect(webhookResponse.ok()).toBe(true);

const execution = await api.workflows.waitForWorkflowStatus(workflowId, 'waiting', 10000);
const originalStartedAt = execution.startedAt;

const fullExecution = await api.workflows.getExecution(execution.id);
const executionData = flatted.parse(fullExecution.data);
const resumeUrl = new URL(
executionData.resultData.runData['Capture Resume URL'][0].data.main[0][0].json
.resumeUrl as string,
);

const resumeResponse = await api.webhooks.trigger(`${resumeUrl.pathname}${resumeUrl.search}`);
expect(resumeResponse.ok()).toBe(true);

await api.workflows.waitForExecution(workflowId, 15000);

await retryUntil(async () => {
const completedExecution = await api.workflows.getExecution(execution.id);
expect(completedExecution.startedAt).toBe(originalStartedAt);
});
});
});

test.describe('when new workflow is not saved', () => {
test.beforeEach(async ({ n8n }) => {
await n8n.start.fromBlankCanvas();
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,110 @@
{
"name": "CAT-1854 Wait Execution History Test",
"nodes": [
{
"parameters": {
"httpMethod": "GET",
"path": "cat-1854-test",
"options": {}
},
"type": "n8n-nodes-base.webhook",
"typeVersion": 2,
"position": [0, 0],
"id": "webhook-trigger-id",
"name": "Webhook",
"webhookId": "cat-1854-webhook-id"
},
{
"parameters": {
"assignments": {
"assignments": [
{
"id": "capture-resume-url",
"name": "resumeUrl",
"value": "={{ $execution.resumeUrl }}",
"type": "string"
}
]
},
"options": {}
},
"type": "n8n-nodes-base.set",
"typeVersion": 3.4,
"position": [220, 0],
"id": "capture-url-id",
"name": "Capture Resume URL"
},
{
"parameters": {
"resume": "webhook",
"options": {}
},
"type": "n8n-nodes-base.wait",
"typeVersion": 1.1,
"position": [440, 0],
"id": "wait-node-id",
"name": "Wait",
"webhookId": "cat-1854-wait-webhook-id"
},
{
"parameters": {
"assignments": {
"assignments": [
{
"id": "output-field-id",
"name": "result",
"value": "completed after wait",
"type": "string"
}
]
},
"options": {}
},
"type": "n8n-nodes-base.set",
"typeVersion": 3.4,
"position": [660, 0],
"id": "set-node-id",
"name": "Edit Fields"
}
],
"connections": {
"Webhook": {
"main": [
[
{
"node": "Capture Resume URL",
"type": "main",
"index": 0
}
]
]
},
"Capture Resume URL": {
"main": [
[
{
"node": "Wait",
"type": "main",
"index": 0
}
]
]
},
"Wait": {
"main": [
[
{
"node": "Edit Fields",
"type": "main",
"index": 0
}
]
]
}
},
"pinData": {},
"active": true,
"meta": {
"instanceId": "test-instance"
}
}
Loading