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
64 changes: 64 additions & 0 deletions packages/@n8n/backend-common/src/logging/__tests__/logger.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -231,6 +231,70 @@ describe('Logger', () => {
});
});

describe('optional metadata fields', () => {
afterEach(() => {
jest.resetAllMocks();
});

test('should include optional metadata fields in JSON output when defined', () => {
const stdoutSpy = jest.spyOn(process.stdout, 'write').mockReturnValue(true);
const globalConfig = mock<GlobalConfig>({
logging: {
format: 'json',
level: 'info',
outputs: ['console'],
scopes: [],
},
});
const logger = new Logger(globalConfig, mock<InstanceSettingsConfig>());

logger.info('Workflow execution started', {
workflowId: 'wf-1',
projectId: 'proj-1',
projectName: 'Test Project',
});

expect(stdoutSpy).toHaveBeenCalledTimes(1);
const output = stdoutSpy.mock.lastCall?.[0];
if (typeof output !== 'string') {
fail(`expected 'output' to be of type 'string', got ${typeof output}`);
}
Comment thread
krider2010 marked this conversation as resolved.

const parsedOutput = JSON.parse(output) as { metadata: Record<string, unknown> };
expect(parsedOutput.metadata).toMatchObject({
workflowId: 'wf-1',
projectId: 'proj-1',
projectName: 'Test Project',
});
});

test('should omit undefined metadata fields from JSON output', () => {
const stdoutSpy = jest.spyOn(process.stdout, 'write').mockReturnValue(true);
const globalConfig = mock<GlobalConfig>({
logging: {
format: 'json',
level: 'info',
outputs: ['console'],
scopes: [],
},
});
const logger = new Logger(globalConfig, mock<InstanceSettingsConfig>());

logger.info('Workflow execution started', { workflowId: 'wf-1' });

expect(stdoutSpy).toHaveBeenCalledTimes(1);
const output = stdoutSpy.mock.lastCall?.[0];
if (typeof output !== 'string') {
fail(`expected 'output' to be of type 'string', got ${typeof output}`);
}
Comment thread
krider2010 marked this conversation as resolved.

const parsedOutput = JSON.parse(output) as { metadata: Record<string, unknown> };
expect(parsedOutput.metadata.workflowId).toBe('wf-1');
expect(parsedOutput.metadata).not.toHaveProperty('projectId');
expect(parsedOutput.metadata).not.toHaveProperty('projectName');
});
});

describe('transports', () => {
afterEach(() => {
jest.restoreAllMocks();
Expand Down
120 changes: 120 additions & 0 deletions packages/cli/src/scaling/__tests__/job-processor.service.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -919,4 +919,124 @@ describe('JobProcessor', () => {
expect(lastResponse.response).toBe('supply data tool result');
});
});

describe('project info in log metadata', () => {
Comment thread
krider2010 marked this conversation as resolved.
beforeEach(() => {
jest.clearAllMocks();
});
it('should include project info in log metadata when present in job data', async () => {
const executionRepository = mock<ExecutionRepository>();
executionRepository.findSingleExecution.mockResolvedValueOnce(
mock<IExecutionResponse>({
mode: 'manual',
workflowData: { id: 'wf-1', name: 'Test Workflow', nodes: [] },
data: mock<IRunExecutionData>({
executionData: undefined,
}),
}),
);
// Second call for checking errors
executionRepository.findSingleExecution.mockResolvedValueOnce(
mock<IExecutionResponse>({
status: 'success',
data: mock<IRunExecutionData>({ resultData: { runData: {} } }),
}),
);

const manualExecutionService = mock<ManualExecutionService>();
const jobProcessor = new JobProcessor(
logger,
executionRepository,
mock(),
mock(),
mock(),
manualExecutionService,
executionsConfig,
mock(),
);

const job = mock<Job>();
job.data = {
workflowId: 'wf-1',
executionId: 'exec-1',
loadStaticData: false,
projectId: 'proj-123',
projectName: 'My Project',
};

await jobProcessor.processJob(job);

// "Worker started" log should include project info
expect(logger.info).toHaveBeenCalledWith(
expect.stringContaining('Worker started execution'),
expect.objectContaining({
workflowId: 'wf-1',
workflowName: 'Test Workflow',
projectId: 'proj-123',
projectName: 'My Project',
}),
);

// "Worker finished" log should include project info
expect(logger.info).toHaveBeenCalledWith(
expect.stringContaining('Worker finished execution'),
expect.objectContaining({
workflowId: 'wf-1',
workflowName: 'Test Workflow',
projectId: 'proj-123',
projectName: 'My Project',
}),
);
});

it('should not include project info in log metadata when absent from job data', async () => {
const executionRepository = mock<ExecutionRepository>();
executionRepository.findSingleExecution.mockResolvedValueOnce(
mock<IExecutionResponse>({
mode: 'manual',
workflowData: { id: 'wf-1', name: 'Test Workflow', nodes: [] },
data: mock<IRunExecutionData>({
executionData: undefined,
}),
}),
);
executionRepository.findSingleExecution.mockResolvedValueOnce(
mock<IExecutionResponse>({
status: 'success',
data: mock<IRunExecutionData>({ resultData: { runData: {} } }),
}),
);

const manualExecutionService = mock<ManualExecutionService>();
const jobProcessor = new JobProcessor(
logger,
executionRepository,
mock(),
mock(),
mock(),
manualExecutionService,
executionsConfig,
mock(),
);

const job = mock<Job>();
job.data = {
workflowId: 'wf-1',
executionId: 'exec-1',
loadStaticData: false,
};

await jobProcessor.processJob(job);

// "Worker started" log should not include project fields
const startedCall = (logger.info as jest.Mock).mock.calls.find(
(call: unknown[]) =>
typeof call[0] === 'string' && call[0].includes('Worker started execution'),
) as [string, Record<string, unknown>] | undefined;
expect(startedCall).toBeDefined();
expect(startedCall![1].workflowId).toBe('wf-1');
expect(startedCall![1]).not.toHaveProperty('projectId');
expect(startedCall![1]).not.toHaveProperty('projectName');
});
});
});
9 changes: 9 additions & 0 deletions packages/cli/src/scaling/job-processor.ts
Original file line number Diff line number Diff line change
Expand Up @@ -95,7 +95,10 @@ export class JobProcessor {
this.logger.info(`Worker started execution ${executionId} (job ${job.id})`, {
executionId,
workflowId,
workflowName: execution.workflowData.name,
jobId: job.id,
...(job.data.projectId !== undefined && { projectId: job.data.projectId }),
...(job.data.projectName !== undefined && { projectName: job.data.projectName }),
});

const startedAt = await this.executionRepository.setRunning(executionId);
Expand Down Expand Up @@ -214,7 +217,10 @@ export class JobProcessor {
{
executionId,
workflowId,
workflowName: execution.workflowData.name,
jobId: job.id,
...(job.data.projectId && { projectId: job.data.projectId }),
...(job.data.projectName && { projectName: job.data.projectName }),
},
);
};
Expand Down Expand Up @@ -294,8 +300,11 @@ export class JobProcessor {
this.logger.info(`Worker finished execution ${executionId} (job ${job.id})`, {
executionId,
workflowId,
workflowName: execution.workflowData.name,
jobId: job.id,
success: props.success,
...(job.data.projectId && { projectId: job.data.projectId }),
...(job.data.projectName && { projectName: job.data.projectName }),
});

const msg: JobFinishedMessage = {
Expand Down
2 changes: 2 additions & 0 deletions packages/cli/src/scaling/scaling.types.ts
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,8 @@ export type JobData = {
pushRef?: string;
streamingEnabled?: boolean;
restartExecutionId?: string;
projectId?: string;
projectName?: string;

// MCP-specific fields for queue mode support
/** Whether this execution was triggered by an MCP tool call. */
Expand Down
1 change: 1 addition & 0 deletions packages/cli/src/webhooks/webhook-helpers.ts
Original file line number Diff line number Diff line change
Expand Up @@ -641,6 +641,7 @@ export async function executeWebhook(
workflowData,
pinData,
projectId: project?.id,
projectName: project?.name,
};

// When resuming from a wait node, copy over the pushRef from the execution-data
Expand Down
6 changes: 6 additions & 0 deletions packages/cli/src/workflow-runner.ts
Original file line number Diff line number Diff line change
Expand Up @@ -147,6 +147,7 @@ export class WorkflowRunner {
const executionId = await this.activeExecutions.add(data, restartExecutionId);

const { id: workflowId, nodes } = data.workflowData;

try {
await this.credentialsPermissionChecker.check(workflowId, nodes);
} catch (error) {
Expand Down Expand Up @@ -205,6 +206,9 @@ export class WorkflowRunner {
error,
executionId,
workflowId,
workflowName: data.workflowData.name,
...(data.projectId && { projectId: data.projectId }),
...(data.projectName && { projectName: data.projectName }),
});
});
}
Expand Down Expand Up @@ -390,6 +394,8 @@ export class WorkflowRunner {
pushRef: data.pushRef,
streamingEnabled: data.streamingEnabled,
restartExecutionId,
projectId: data.projectId,
projectName: data.projectName,
// MCP-specific fields for queue mode support
isMcpExecution: data.isMcpExecution,
mcpType: data.mcpType,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@ import {

import type { IWorkflowErrorData } from '@/interfaces';
import type { NodeTypes } from '@/node-types';
import type { OwnershipService } from '@/services/ownership.service';
import type { TestWebhooks } from '@/webhooks/test-webhooks';
import * as WorkflowExecuteAdditionalData from '@/workflow-execute-additional-data';
import type { WorkflowRunner } from '@/workflow-runner';
Expand Down Expand Up @@ -75,6 +76,14 @@ const secondHackerNewsNode: INode = {
position: [0, 0],
};

const mockOwnershipService = () => {
const ownershipService = mock<OwnershipService>();
ownershipService.getWorkflowProjectCached.mockResolvedValue(
mock<Project>({ id: 'test-project-id', name: 'Test Project' }),
);
return ownershipService;
};

describe('WorkflowExecutionService', () => {
const nodeTypes = mock<NodeTypes>();
const workflowRunner = mock<WorkflowRunner>();
Expand All @@ -90,6 +99,7 @@ describe('WorkflowExecutionService', () => {
mock(),
mock(),
mock(),
mockOwnershipService(),
);

const additionalData = mock<IWorkflowExecuteAdditionalData>({});
Expand Down Expand Up @@ -154,6 +164,8 @@ describe('WorkflowExecutionService', () => {
workflowData,
userId,
dirtyNodeNames: runPayload.dirtyNodeNames,
projectId: 'test-project-id',
projectName: 'Test Project',
});
expect(result).toEqual({ executionId });
});
Expand Down Expand Up @@ -187,6 +199,8 @@ describe('WorkflowExecutionService', () => {
pushRef: undefined,
workflowData,
userId,
projectId: 'test-project-id',
projectName: 'Test Project',
});
expect(result).toEqual({ executionId });
});
Expand Down Expand Up @@ -247,6 +261,8 @@ describe('WorkflowExecutionService', () => {
workflowData,
userId,
triggerToStartFrom: { name: pinnedTrigger.name },
projectId: 'test-project-id',
projectName: 'Test Project',
});
expect(result).toEqual({ executionId });
});
Expand Down Expand Up @@ -303,6 +319,8 @@ describe('WorkflowExecutionService', () => {
userId,
// pass unexecuted trigger to start from
triggerToStartFrom: runPayload.triggerToStartFrom,
projectId: 'test-project-id',
projectName: 'Test Project',
});
expect(result).toEqual({ executionId });
});
Expand Down Expand Up @@ -409,6 +427,7 @@ describe('WorkflowExecutionService', () => {
mock(),
mock(),
mock(),
mockOwnershipService(),
);

const runPayload: WorkflowRequest.FullManualExecutionFromKnownTriggerPayload = {
Expand Down Expand Up @@ -589,6 +608,7 @@ describe('WorkflowExecutionService', () => {
mock(),
mock(),
mock(),
mockOwnershipService(),
);
});

Expand Down Expand Up @@ -740,12 +760,13 @@ describe('WorkflowExecutionService', () => {
mock(),
mock(),
mock(),
mockOwnershipService(),
);

await service.executeErrorWorkflow(
'error-workflow-id',
workflowErrorData,
mock<Project>({ id: 'project-id' }),
mock<Project>({ id: 'project-id', name: 'Error Project' }),
);

expect(workflowRunnerMock.run).toHaveBeenCalledTimes(1);
Expand Down Expand Up @@ -786,6 +807,7 @@ describe('WorkflowExecutionService', () => {
}),
workflowData: errorWorkflow,
projectId: 'project-id',
projectName: 'Error Project',
});
});

Expand Down Expand Up @@ -868,6 +890,7 @@ describe('WorkflowExecutionService', () => {
mock(),
mock(),
mock(),
mock(),
);

await service.executeErrorWorkflow(
Expand Down
Loading
Loading