feat(core): Add tags support to MCP server (#31446)

This commit is contained in:
Ricardo Espinoza
2026-06-11 13:25:52 +00:00
committed by GitHub
parent 2f3ebb620c
commit 1640a69ea6
17 changed files with 1145 additions and 31 deletions
@@ -0,0 +1,159 @@
import { mockInstance } from '@n8n/backend-test-utils';
import type { ITagWithCountDb } from '@n8n/db';
import { User } from '@n8n/db';
import { TagService } from '@/services/tag.service';
import { Telemetry } from '@/telemetry';
import { createListTagsTool, listTags } from '../tools/list-tags.tool';
const buildTag = (overrides: Partial<ITagWithCountDb> = {}): ITagWithCountDb =>
({
id: 'tag-1',
name: 'production',
usageCount: 3,
createdAt: new Date('2024-01-01T00:00:00.000Z'),
updatedAt: new Date('2024-01-02T00:00:00.000Z'),
...overrides,
}) as ITagWithCountDb;
const userWithScopes = (scopeSlugs: string[]) =>
Object.assign(new User(), {
id: 'user-1',
role: { slug: 'global:test', scopes: scopeSlugs.map((slug) => ({ slug })) },
});
describe('list-tags MCP tool', () => {
const user = userWithScopes(['tag:list']);
const createMocks = (
tagsOrError: ITagWithCountDb[] | Error = [],
opts: { totalCount?: number } = {},
) => {
const listWithUsageCount =
tagsOrError instanceof Error
? jest.fn().mockRejectedValue(tagsOrError)
: jest.fn().mockResolvedValue({
data: tagsOrError,
totalCount: opts.totalCount ?? tagsOrError.length,
});
const tagService = mockInstance(TagService, { listWithUsageCount });
const telemetry = mockInstance(Telemetry, { track: jest.fn() });
return { tagService, telemetry };
};
describe('smoke tests', () => {
test('creates the tool correctly', () => {
const { tagService, telemetry } = createMocks();
const tool = createListTagsTool(user, tagService, telemetry);
expect(tool.name).toBe('list_tags');
expect(tool.config.description).toEqual(expect.any(String));
expect(tool.config.inputSchema).toBeDefined();
expect(tool.config.outputSchema).toBeDefined();
expect(tool.config.annotations).toMatchObject({
readOnlyHint: true,
destructiveHint: false,
idempotentHint: true,
openWorldHint: false,
});
expect(typeof tool.handler).toBe('function');
});
});
describe('handler', () => {
test('requests usage counts and returns the formatted payload', async () => {
const tags = [buildTag(), buildTag({ id: 'tag-2', name: 'staging', usageCount: 0 })];
const { tagService } = createMocks(tags);
const result = await listTags(tagService);
expect(tagService.listWithUsageCount).toHaveBeenCalledWith({ limit: 500 });
expect(result.count).toBe(2);
expect(result.totalCount).toBe(2);
expect(result.data).toEqual([
{
id: 'tag-1',
name: 'production',
usageCount: 3,
createdAt: new Date('2024-01-01T00:00:00.000Z').toISOString(),
updatedAt: new Date('2024-01-02T00:00:00.000Z').toISOString(),
},
{
id: 'tag-2',
name: 'staging',
usageCount: 0,
createdAt: new Date('2024-01-01T00:00:00.000Z').toISOString(),
updatedAt: new Date('2024-01-02T00:00:00.000Z').toISOString(),
},
]);
});
test('defaults usageCount to 0 when missing', async () => {
const tags = [buildTag({ usageCount: undefined as unknown as number })];
const { tagService } = createMocks(tags);
const result = await listTags(tagService);
expect(result.data[0].usageCount).toBe(0);
});
test('pushes the requested limit into the query and reports totalCount', async () => {
const tags = [buildTag({ id: 'tag-1' }), buildTag({ id: 'tag-2' })];
const { tagService } = createMocks(tags, { totalCount: 3 });
const result = await listTags(tagService, { limit: 2 });
expect(tagService.listWithUsageCount).toHaveBeenCalledWith({ limit: 2 });
expect(result.data.map((t) => t.id)).toEqual(['tag-1', 'tag-2']);
expect(result.count).toBe(2);
expect(result.totalCount).toBe(3);
});
test('emits telemetry on success', async () => {
const { tagService, telemetry } = createMocks([buildTag()]);
const tool = createListTagsTool(user, tagService, telemetry);
await tool.handler({ limit: undefined as unknown as number }, {} as never);
expect(telemetry.track).toHaveBeenCalledWith(
expect.any(String),
expect.objectContaining({
tool_name: 'list_tags',
user_id: 'user-1',
results: { success: true, data: { count: 1 } },
}),
);
});
test('emits telemetry and rethrows on failure', async () => {
const { tagService, telemetry } = createMocks(new Error('boom'));
const tool = createListTagsTool(user, tagService, telemetry);
await expect(
tool.handler({ limit: undefined as unknown as number }, {} as never),
).rejects.toThrow('boom');
expect(telemetry.track).toHaveBeenCalledWith(
expect.any(String),
expect.objectContaining({
tool_name: 'list_tags',
results: { success: false, error: 'boom' },
}),
);
});
test('rejects when the user does not have tag:list scope', async () => {
const noScopeUser = userWithScopes([]);
const { tagService, telemetry } = createMocks([buildTag()]);
const tool = createListTagsTool(noScopeUser, tagService, telemetry);
await expect(
tool.handler({ limit: undefined as unknown as number }, {} as never),
).rejects.toThrow('permission to list tags');
expect(tagService.listWithUsageCount).not.toHaveBeenCalled();
});
});
});
@@ -42,6 +42,7 @@ import { NodeTypes } from '@/node-types';
import { PostHogClient } from '@/posthog';
import { ProjectService } from '@/services/project.service.ee';
import { RoleService } from '@/services/role.service';
import { TagService } from '@/services/tag.service';
import { UrlService } from '@/services/url.service';
import { Telemetry } from '@/telemetry';
import { WorkflowRunner } from '@/workflow-runner';
@@ -93,6 +94,7 @@ describe('McpService', () => {
mockInstance(ExecutionService),
mockInstance(DataTableProxyService),
mockInstance(CollaborationService),
mockInstance(TagService),
mockInstance(LicenseState),
mockInstance(PostHogClient),
);
@@ -135,6 +137,7 @@ describe('McpService', () => {
mockInstance(ExecutionService),
mockInstance(DataTableProxyService),
mockInstance(CollaborationService),
mockInstance(TagService),
mockInstance(LicenseState),
mockInstance(PostHogClient),
);
@@ -330,6 +333,7 @@ describe('McpService', () => {
mockInstance(ExecutionService),
mockInstance(DataTableProxyService),
mockInstance(CollaborationService),
mockInstance(TagService),
mockInstance(LicenseState),
opts.postHogClient,
);
@@ -433,6 +437,7 @@ describe('McpService', () => {
mockInstance(ExecutionService),
mockInstance(DataTableProxyService),
mockInstance(CollaborationService),
mockInstance(TagService),
mockInstance(LicenseState),
mockInstance(PostHogClient),
);
@@ -477,6 +482,7 @@ describe('McpService', () => {
mockInstance(ExecutionService),
mockInstance(DataTableProxyService),
mockInstance(CollaborationService),
mockInstance(TagService),
mockInstance(LicenseState),
mockInstance(PostHogClient),
);
@@ -545,6 +551,7 @@ describe('McpService', () => {
mockInstance(ExecutionService),
mockInstance(DataTableProxyService),
mockInstance(CollaborationService),
mockInstance(TagService),
mockInstance(LicenseState),
postHogClient,
);
@@ -1,5 +1,6 @@
import { mockInstance } from '@n8n/backend-test-utils';
import { User } from '@n8n/db';
import type { WorkflowEntity } from '@n8n/db';
import type { INode } from 'n8n-workflow';
import {
@@ -109,6 +110,7 @@ describe('search-workflows MCP tool', () => {
scopes: ['workflow:read', 'workflow:execute'],
canExecute: true,
availableInMCP: true,
tags: [],
},
{
id: 'b',
@@ -121,10 +123,66 @@ describe('search-workflows MCP tool', () => {
scopes: ['workflow:read'],
canExecute: false,
availableInMCP: true,
tags: [],
},
]);
});
test('forwards tags filter and surfaces workflow tags in output', async () => {
const tags = [
{ id: 'tag-1', name: 'production' },
{ id: 'tag-2', name: 'critical' },
] as unknown as WorkflowEntity['tags'];
const workflows = [
createWorkflow({
id: 'tagged',
activeVersionId: uuid(),
tags,
}),
];
const workflowService = mockInstance(WorkflowService, {
getMany: jest.fn().mockResolvedValue({ workflows, count: 1 }),
});
const result = await searchWorkflows(user, workflowService as unknown as WorkflowService, {
tags: ['production', 'critical'],
});
const [, optionsArg] = (workflowService.getMany as jest.Mock).mock.calls[0];
expect(optionsArg.filter).toMatchObject({ tags: ['production', 'critical'] });
expect(optionsArg.select).toMatchObject({ tags: true });
expect(result.data[0].tags).toEqual([
{ id: 'tag-1', name: 'production' },
{ id: 'tag-2', name: 'critical' },
]);
});
test('drops empty tag entries and omits filter when no tags remain', async () => {
const workflowService = mockInstance(WorkflowService, {
getMany: jest.fn().mockResolvedValue({ workflows: [], count: 0 }),
});
await searchWorkflows(user, workflowService as unknown as WorkflowService, {
tags: ['', ''],
});
const [, optionsArg] = (workflowService.getMany as jest.Mock).mock.calls[0];
expect(optionsArg.filter.tags).toBeUndefined();
});
test('deduplicates repeated tag names before forwarding the filter', async () => {
const workflowService = mockInstance(WorkflowService, {
getMany: jest.fn().mockResolvedValue({ workflows: [], count: 0 }),
});
await searchWorkflows(user, workflowService as unknown as WorkflowService, {
tags: ['production', 'production', 'critical', 'production'],
});
const [, optionsArg] = (workflowService.getMany as jest.Mock).mock.calls[0];
expect(optionsArg.filter.tags).toEqual(['production', 'critical']);
});
test('applies provided filters and clamps high limit', async () => {
const workflows = [createWorkflow({ id: 'x', activeVersionId: uuid() })];
const workflowService = mockInstance(WorkflowService, {
@@ -1,4 +1,5 @@
import { mockInstance } from '@n8n/backend-test-utils';
import { GlobalConfig } from '@n8n/config';
import { SharedWorkflowRepository, User, WorkflowEntity } from '@n8n/db';
import { NodeConnectionTypes, type IConnections, type INode } from 'n8n-workflow';
@@ -8,6 +9,7 @@ import { CollaborationService } from '@/collaboration/collaboration.service';
import { CredentialsService } from '@/credentials/credentials.service';
import { NotFoundError } from '@/errors/response-errors/not-found.error';
import { NodeTypes } from '@/node-types';
import { TagService } from '@/services/tag.service';
import { UrlService } from '@/services/url.service';
import { Telemetry } from '@/telemetry';
import { WorkflowFinderService } from '@/workflows/workflow-finder.service';
@@ -48,8 +50,14 @@ type DataTableOpsMock = {
getManyAndCount: jest.Mock;
};
const userWithScopes = (scopeSlugs: string[]) =>
Object.assign(new User(), {
id: 'user-1',
role: { slug: 'global:test', scopes: scopeSlugs.map((slug) => ({ slug })) },
});
describe('update-workflow MCP tool', () => {
const user = Object.assign(new User(), { id: 'user-1' });
const user = userWithScopes(['tag:create']);
let workflowFinderService: WorkflowFinderService;
let findWorkflowMock: jest.Mock;
let workflowService: WorkflowService;
@@ -61,6 +69,10 @@ describe('update-workflow MCP tool', () => {
let nodeTypes: ReturnType<typeof mockInstance<NodeTypes>>;
let collaborationService: CollaborationService;
let dataTableOps: DataTableOpsMock;
let tagService: TagService;
let findOrCreateByNamesMock: jest.Mock;
let findByNamesMock: jest.Mock;
let globalConfig: GlobalConfig;
const buildExistingWorkflow = () =>
Object.assign(new WorkflowEntity(), {
@@ -122,6 +134,14 @@ describe('update-workflow MCP tool', () => {
dataTableOps = {
getManyAndCount: jest.fn().mockResolvedValue({ data: [], count: 0 }),
};
findOrCreateByNamesMock = jest.fn();
findByNamesMock = jest.fn();
tagService = mockInstance(TagService, {
findOrCreateByNames: findOrCreateByNamesMock,
findByNames: findByNamesMock,
});
globalConfig = mockInstance(GlobalConfig, { tags: { disabled: false } });
});
const createTool = () =>
@@ -136,6 +156,8 @@ describe('update-workflow MCP tool', () => {
sharedWorkflowRepository,
collaborationService,
dataTableOps as never,
tagService,
globalConfig,
);
const callHandler = async (
@@ -945,5 +967,201 @@ describe('update-workflow MCP tool', () => {
expect(workflowService.update).toHaveBeenCalled();
});
});
describe('tag operations', () => {
const workflowWithTags = (tagNames: string[]) =>
Object.assign(buildExistingWorkflow(), {
tags: tagNames.map((name, i) => ({ id: `tag-${i}`, name })),
});
test('resolves added tag names and passes tagIds to workflow update', async () => {
findWorkflowMock.mockResolvedValue(workflowWithTags(['production']));
findOrCreateByNamesMock.mockResolvedValue([
{ id: 'tag-0', name: 'production' },
{ id: 'tag-new', name: 'critical' },
]);
const result = await callHandler({
workflowId: 'wf-1',
operations: [{ type: 'addTags', names: ['critical'] }],
});
expect(result.isError).toBeUndefined();
expect(findWorkflowMock).toHaveBeenCalledWith(
'wf-1',
user,
['workflow:update'],
expect.objectContaining({ includeTags: true }),
);
expect(findOrCreateByNamesMock).toHaveBeenCalledTimes(1);
const passedNames = findOrCreateByNamesMock.mock.calls[0][0] as string[];
expect(passedNames.sort()).toEqual(['critical', 'production']);
const [, , , updateOptions] = updateMock.mock.calls[0];
expect(updateOptions.tagIds.sort()).toEqual(['tag-0', 'tag-new']);
});
test('removeTags drops names from the resolved set', async () => {
findWorkflowMock.mockResolvedValue(workflowWithTags(['production', 'critical']));
findOrCreateByNamesMock.mockResolvedValue([{ id: 'tag-0', name: 'production' }]);
await callHandler({
workflowId: 'wf-1',
operations: [{ type: 'removeTags', names: ['critical'] }],
});
expect(findOrCreateByNamesMock).toHaveBeenCalledWith(['production']);
const [, , , updateOptions] = updateMock.mock.calls[0];
expect(updateOptions.tagIds).toEqual(['tag-0']);
});
test('removing the last tag passes an empty tagIds array', async () => {
findWorkflowMock.mockResolvedValue(workflowWithTags(['production']));
findOrCreateByNamesMock.mockResolvedValue([]);
await callHandler({
workflowId: 'wf-1',
operations: [{ type: 'removeTags', names: ['production'] }],
});
expect(findOrCreateByNamesMock).toHaveBeenCalledWith([]);
const [, , , updateOptions] = updateMock.mock.calls[0];
expect(updateOptions.tagIds).toEqual([]);
});
test('does not call tagService or pass tagIds when no tag ops are present', async () => {
await callHandler({
workflowId: 'wf-1',
operations: [{ type: 'setWorkflowMetadata', name: 'renamed' }],
});
expect(findOrCreateByNamesMock).not.toHaveBeenCalled();
const [, , , updateOptions] = updateMock.mock.calls[0];
expect(updateOptions.tagIds).toBeUndefined();
// Tags should not be loaded when there are no tag ops
expect(findWorkflowMock).toHaveBeenCalledWith(
'wf-1',
user,
['workflow:update'],
expect.objectContaining({ includeTags: false }),
);
});
test('rejects tag operations when tags are disabled instance-wide', async () => {
globalConfig = mockInstance(GlobalConfig, { tags: { disabled: true } });
const result = await callHandler({
workflowId: 'wf-1',
operations: [{ type: 'addTags', names: ['anything'] }],
});
expect(result.isError).toBe(true);
expect(findOrCreateByNamesMock).not.toHaveBeenCalled();
expect(workflowService.update).not.toHaveBeenCalled();
expect(findWorkflowMock).not.toHaveBeenCalled();
});
test('without tag:create scope, attaches only existing tags', async () => {
const memberUser = userWithScopes([]);
findWorkflowMock.mockResolvedValue(workflowWithTags([]));
findByNamesMock.mockResolvedValue([{ id: 'tag-existing', name: 'production' }]);
const tool = createUpdateWorkflowTool(
memberUser,
workflowFinderService,
workflowService,
urlService,
telemetry,
nodeTypes,
credentialsService,
sharedWorkflowRepository,
collaborationService,
dataTableOps as never,
tagService,
globalConfig,
);
await callHandler(
{
workflowId: 'wf-1',
operations: [{ type: 'addTags', names: ['production'] }],
},
tool,
);
expect(findByNamesMock).toHaveBeenCalledWith(['production']);
expect(findOrCreateByNamesMock).not.toHaveBeenCalled();
const [, , , updateOptions] = updateMock.mock.calls[0];
expect(updateOptions.tagIds).toEqual(['tag-existing']);
});
test('without tag:create scope, fails when a tag name does not exist', async () => {
const memberUser = userWithScopes([]);
findWorkflowMock.mockResolvedValue(workflowWithTags([]));
findByNamesMock.mockResolvedValue([{ id: 'tag-existing', name: 'production' }]);
const tool = createUpdateWorkflowTool(
memberUser,
workflowFinderService,
workflowService,
urlService,
telemetry,
nodeTypes,
credentialsService,
sharedWorkflowRepository,
collaborationService,
dataTableOps as never,
tagService,
globalConfig,
);
const result = await callHandler(
{
workflowId: 'wf-1',
operations: [{ type: 'addTags', names: ['production', 'novel-tag'] }],
},
tool,
);
expect(result.isError).toBe(true);
expect(findOrCreateByNamesMock).not.toHaveBeenCalled();
expect(workflowService.update).not.toHaveBeenCalled();
});
test('does not flip aiBuilderAssisted when the batch contains only tag operations', async () => {
findWorkflowMock.mockResolvedValue(workflowWithTags(['existing']));
findOrCreateByNamesMock.mockResolvedValue([{ id: 'tag-0', name: 'existing' }]);
await callHandler({
workflowId: 'wf-1',
operations: [{ type: 'addTags', names: ['existing'] }],
});
const [, workflowArg, , updateOptions] = updateMock.mock.calls[0];
expect(updateOptions.aiBuilderAssisted).toBe(false);
expect(workflowArg.meta).not.toEqual(
expect.objectContaining({ aiBuilderAssisted: true, builderVariant: 'mcp' }),
);
});
test('keeps aiBuilderAssisted=true when tag ops are mixed with node ops', async () => {
findWorkflowMock.mockResolvedValue(workflowWithTags([]));
findOrCreateByNamesMock.mockResolvedValue([{ id: 'tag-0', name: 'foo' }]);
await callHandler({
workflowId: 'wf-1',
operations: [
{ type: 'setWorkflowMetadata', name: 'Renamed' },
{ type: 'addTags', names: ['foo'] },
],
});
const [, workflowArg, , updateOptions] = updateMock.mock.calls[0];
expect(updateOptions.aiBuilderAssisted).toBe(true);
expect(workflowArg.meta).toEqual(
expect.objectContaining({ aiBuilderAssisted: true, builderVariant: 'mcp' }),
);
});
});
});
});
@@ -794,4 +794,80 @@ describe('applyOperations', () => {
expect(result.success).toBe(false);
});
});
describe('addTags / removeTags', () => {
test('addTags adds names to the existing set without duplicating', () => {
const wf = { ...baseWorkflow(), tagNames: ['production'] };
const ops: PartialUpdateOperation[] = [
{ type: 'addTags', names: ['critical', 'production'] },
];
const result = applyOperations(wf, ops);
if (!result.success) throw new Error('expected success');
expect(result.tagNames?.sort()).toEqual(['critical', 'production']);
expect(result.workflow.tagNames?.sort()).toEqual(['critical', 'production']);
});
test('removeTags removes only the matching names', () => {
const wf = { ...baseWorkflow(), tagNames: ['production', 'critical', 'wip'] };
const ops: PartialUpdateOperation[] = [{ type: 'removeTags', names: ['wip', 'missing'] }];
const result = applyOperations(wf, ops);
if (!result.success) throw new Error('expected success');
expect(result.tagNames?.sort()).toEqual(['critical', 'production']);
});
test('sequential ops apply in order (add then remove yields empty)', () => {
const wf = { ...baseWorkflow(), tagNames: [] };
const ops: PartialUpdateOperation[] = [
{ type: 'addTags', names: ['a', 'b'] },
{ type: 'removeTags', names: ['a'] },
{ type: 'addTags', names: ['c'] },
];
const result = applyOperations(wf, ops);
if (!result.success) throw new Error('expected success');
expect(result.tagNames?.sort()).toEqual(['b', 'c']);
});
test('returns undefined tagNames when batch has no tag operations', () => {
const wf = { ...baseWorkflow(), tagNames: ['production'] };
const ops: PartialUpdateOperation[] = [{ type: 'setWorkflowMetadata', name: 'renamed' }];
const result = applyOperations(wf, ops);
if (!result.success) throw new Error('expected success');
expect(result.tagNames).toBeUndefined();
});
test('fails when tag operations are present but existing tags were not loaded', () => {
const wf = baseWorkflow();
const ops: PartialUpdateOperation[] = [{ type: 'addTags', names: ['production'] }];
const result = applyOperations(wf, ops);
expect(result.success).toBe(false);
});
test('does not mutate input tagNames on success', () => {
const wf = { ...baseWorkflow(), tagNames: ['production'] };
const before = [...wf.tagNames];
applyOperations(wf, [{ type: 'addTags', names: ['critical'] }]);
expect(wf.tagNames).toEqual(before);
});
test('schema rejects empty names array', () => {
const parsed = partialUpdateOperationSchema.safeParse({ type: 'addTags', names: [] });
expect(parsed.success).toBe(false);
});
test('schema rejects empty string names', () => {
const parsed = partialUpdateOperationSchema.safeParse({ type: 'removeTags', names: [''] });
expect(parsed.success).toBe(false);
});
test('schema trims whitespace from names', () => {
const parsed = partialUpdateOperationSchema.safeParse({
type: 'addTags',
names: [' spaced '],
});
expect(parsed.success).toBe(true);
if (parsed.success && parsed.data.type === 'addTags') {
expect(parsed.data.names).toEqual(['spaced']);
}
});
});
});
@@ -38,6 +38,7 @@ import { createGetExecutionTool } from './tools/get-execution.tool';
import { createSearchExecutionsTool } from './tools/search-executions.tool';
import { createWorkflowDetailsTool } from './tools/get-workflow-details.tool';
import { createListCredentialsTool } from './tools/list-credentials.tool';
import { createListTagsTool } from './tools/list-tags.tool';
import { createPublishWorkflowTool } from './tools/publish-workflow.tool';
import { createSearchFoldersTool } from './tools/search-folders.tool';
import { createSearchProjectsTool } from './tools/search-projects.tool';
@@ -65,6 +66,7 @@ import { NodeTypes } from '@/node-types';
import { PostHogClient } from '@/posthog';
import { ProjectService } from '@/services/project.service.ee';
import { RoleService } from '@/services/role.service';
import { TagService } from '@/services/tag.service';
import { UrlService } from '@/services/url.service';
import { Telemetry } from '@/telemetry';
import { WorkflowRunner } from '@/workflow-runner';
@@ -128,6 +130,7 @@ export class McpService {
private readonly executionService: ExecutionService,
private readonly dataTableProxyService: DataTableProxyService,
private readonly collaborationService: CollaborationService,
private readonly tagService: TagService,
private readonly licenseState: LicenseState,
private readonly postHogClient: PostHogClient,
) {}
@@ -329,6 +332,11 @@ export class McpService {
listCredentialsTool.handler,
);
if (!this.globalConfig.tags.disabled) {
const listTagsTool = createListTagsTool(user, this.tagService, this.telemetry);
server.registerTool(listTagsTool.name, listTagsTool.config, listTagsTool.handler);
}
// Data table tools
const dataTableOps = this.dataTableProxyService.makeDataTableOperationsForUser(user);
@@ -519,6 +527,8 @@ export class McpService {
this.sharedWorkflowRepository,
this.collaborationService,
dataTableOps,
this.tagService,
this.globalConfig,
);
server.registerTool(updateTool.name, updateTool.config, updateTool.handler);
@@ -40,6 +40,7 @@ export type SearchWorkflowsParams = {
limit?: number;
query?: string;
projectId?: string;
tags?: string[];
sortBy?: SearchWorkflowsSortBy;
};
@@ -54,6 +55,7 @@ export type SearchWorkflowsItem = {
scopes: string[];
canExecute: boolean;
availableInMCP: boolean;
tags: Array<{ id: string; name: string }>;
};
export type SearchWorkflowsResult = {
@@ -7,7 +7,7 @@ import type {
WorkflowDetailsResult,
UserCalledMCPToolEventPayload,
} from '../mcp.types';
import { workflowDetailsOutputSchema } from './schemas';
import { toTagSummary, workflowDetailsOutputSchema } from './schemas';
import { getTriggerDetails, type WebhookEndpoints } from './webhook-utils';
import { getMcpWorkflow } from './workflow-validation.utils';
@@ -164,7 +164,7 @@ export async function getWorkflowDetails(
connections,
nodes: nodes.map(({ credentials: _credentials, ...node }) => node),
activeVersion,
tags: (workflow.tags ?? []).map((tag) => ({ id: tag.id, name: tag.name })),
tags: toTagSummary(workflow.tags),
meta: workflow.meta ?? null,
parentFolderId: workflow.parentFolder?.id ?? null,
description: workflow.description ?? undefined,
@@ -0,0 +1,129 @@
import type { User } from '@n8n/db';
import { hasGlobalScope } from '@n8n/permissions';
import z from 'zod';
import type { TagService } from '@/services/tag.service';
import type { Telemetry } from '@/telemetry';
import { USER_CALLED_MCP_TOOL_EVENT } from '../mcp.constants';
import type { ToolDefinition, UserCalledMCPToolEventPayload } from '../mcp.types';
import { createLimitSchema } from './schemas';
const MAX_RESULTS = 500;
const inputSchema = {
limit: createLimitSchema(MAX_RESULTS),
} satisfies z.ZodRawShape;
const outputSchema = {
data: z
.array(
z.object({
id: z.string().describe('The unique identifier of the tag'),
name: z.string().describe('The display name of the tag'),
usageCount: z
.number()
.int()
.min(0)
.describe('Number of non-archived workflows using this tag'),
createdAt: z.string().describe('The ISO timestamp when the tag was created'),
updatedAt: z.string().describe('The ISO timestamp when the tag was last updated'),
}),
)
.describe('Workflow tags available in the instance'),
count: z.number().int().min(0).describe('Number of tags returned'),
totalCount: z.number().int().min(0).describe('Total number of tags before applying the limit'),
} satisfies z.ZodRawShape;
export type ListTagsParams = {
limit?: number;
};
export type ListTagsItem = {
id: string;
name: string;
usageCount: number;
createdAt: string;
updatedAt: string;
};
export type ListTagsResult = {
data: ListTagsItem[];
count: number;
totalCount: number;
};
/**
* Creates mcp tool definition for listing all workflow tags in the instance.
* Tags are global (not project-scoped) and can be used with search_workflows to filter results.
*/
export const createListTagsTool = (
user: User,
tagService: TagService,
telemetry: Telemetry,
): ToolDefinition<typeof inputSchema> => ({
name: 'list_tags',
config: {
description: 'List all workflow tags in the instance.',
inputSchema,
outputSchema,
annotations: {
title: 'List Tags',
readOnlyHint: true,
destructiveHint: false,
idempotentHint: true,
openWorldHint: false,
},
},
handler: async ({ limit = MAX_RESULTS }: ListTagsParams) => {
const telemetryPayload: UserCalledMCPToolEventPayload = {
user_id: user.id,
tool_name: 'list_tags',
parameters: { limit },
};
try {
if (!hasGlobalScope(user, 'tag:list')) {
throw new Error('User does not have permission to list tags');
}
const payload = await listTags(tagService, { limit });
telemetryPayload.results = {
success: true,
data: { count: payload.count },
};
telemetry.track(USER_CALLED_MCP_TOOL_EVENT, telemetryPayload);
return {
content: [{ type: 'text', text: JSON.stringify(payload) }],
structuredContent: payload,
};
} catch (error) {
telemetryPayload.results = {
success: false,
error: error instanceof Error ? error.message : String(error),
};
telemetry.track(USER_CALLED_MCP_TOOL_EVENT, telemetryPayload);
throw error;
}
},
});
export async function listTags(
tagService: TagService,
{ limit = MAX_RESULTS }: ListTagsParams = {},
): Promise<ListTagsResult> {
const safeLimit = Math.min(Math.max(1, limit), MAX_RESULTS);
const { data: tags, totalCount } = await tagService.listWithUsageCount({ limit: safeLimit });
const data: ListTagsItem[] = tags.map((tag) => ({
id: tag.id,
name: tag.name,
usageCount: tag.usageCount ?? 0,
createdAt: tag.createdAt.toISOString(),
updatedAt: tag.updatedAt.toISOString(),
}));
return { data, count: data.length, totalCount };
}
@@ -10,6 +10,9 @@ export const nodeSchema = z
export const tagSchema = z.object({ id: z.string(), name: z.string() }).passthrough();
export const toTagSummary = (tags: Array<{ id: string; name: string }> | undefined | null) =>
(tags ?? []).map((tag) => ({ id: tag.id, name: tag.name }));
export const workflowSettingsSchema = z
.custom<IWorkflowSettings>((_value): _value is IWorkflowSettings => true)
.nullable();
@@ -15,7 +15,7 @@ import type {
import type { ListQuery } from '@/requests';
import type { Telemetry } from '@/telemetry';
import type { WorkflowService } from '@/workflows/workflow.service';
import { createLimitSchema } from './schemas';
import { createLimitSchema, tagSchema, toTagSummary } from './schemas';
const MAX_RESULTS = 200;
@@ -25,6 +25,10 @@ const inputSchema = {
limit: createLimitSchema(MAX_RESULTS),
query: z.string().optional().describe('Filter by name or description'),
projectId: z.string().optional(),
tags: z
.array(z.string())
.optional()
.describe('Filter by tag names (AND semantics — workflow must have all).'),
sortBy: z
.enum(SEARCH_WORKFLOWS_SORT_BY_VALUES)
.optional()
@@ -60,6 +64,7 @@ const outputSchema = {
.boolean()
.describe('Whether the user has permission to execute this workflow'),
availableInMCP: z.boolean().describe('Whether the workflow is visible to MCP tools'),
tags: z.array(tagSchema).describe('Tags assigned to the workflow'),
}),
)
.describe('List of workflows matching the query'),
@@ -67,8 +72,8 @@ const outputSchema = {
} satisfies z.ZodRawShape;
/**
* Creates mcp tool definition for searching workflows with optional filters. Workflows can be filtered by name, active status, and project ID.
* Returns a preview of each workflow including id, name, active status, creation and update timestamps, and trigger count.
* Creates mcp tool definition for searching workflows with optional filters. Workflows can be filtered by name, project ID, and tags.
* Returns a preview of each workflow including id, name, active status, creation and update timestamps, trigger count, and tags.
*/
export const createSearchWorkflowsTool = (
user: User,
@@ -94,14 +99,10 @@ export const createSearchWorkflowsTool = (
limit = MAX_RESULTS,
query,
projectId,
tags,
sortBy,
}: {
limit?: number;
query?: string;
projectId?: string;
sortBy?: SearchWorkflowsSortBy;
}) => {
const parameters = { limit, query, projectId, sortBy };
}: SearchWorkflowsParams) => {
const parameters = { limit, query, projectId, tags, sortBy };
const telemetryPayload: UserCalledMCPToolEventPayload = {
user_id: user.id,
tool_name: 'search_workflows',
@@ -113,6 +114,7 @@ export const createSearchWorkflowsTool = (
limit,
query,
projectId,
tags,
sortBy,
});
@@ -151,9 +153,10 @@ export const createSearchWorkflowsTool = (
export async function searchWorkflows(
user: User,
workflowService: WorkflowService,
{ limit = MAX_RESULTS, query, projectId, sortBy = DEFAULT_SORT_BY }: SearchWorkflowsParams,
{ limit = MAX_RESULTS, query, projectId, tags, sortBy = DEFAULT_SORT_BY }: SearchWorkflowsParams,
): Promise<SearchWorkflowsResult> {
const safeLimit = Math.min(Math.max(1, limit), MAX_RESULTS);
const filterTags = tags && Array.from(new Set(tags.filter((tag) => tag.length > 0)));
const options: ListQuery.Options = {
take: safeLimit,
@@ -162,6 +165,7 @@ export async function searchWorkflows(
isArchived: false,
...(query ? { query } : {}),
...(projectId ? { projectId } : {}),
...(filterTags && filterTags.length > 0 ? { tags: filterTags } : {}),
},
select: {
id: true,
@@ -173,6 +177,7 @@ export async function searchWorkflows(
triggerCount: true,
ownedBy: true, // Required for loading 'shared' relation used in scope computation
settings: true,
tags: true,
},
};
@@ -185,8 +190,17 @@ export async function searchWorkflows(
);
const formattedWorkflows: SearchWorkflowsItem[] = workflows.map((workflow) => {
const { id, name, description, activeVersionId, createdAt, updatedAt, triggerCount, settings } =
workflow as WorkflowEntity;
const {
id,
name,
description,
activeVersionId,
createdAt,
updatedAt,
triggerCount,
settings,
tags: workflowTags,
} = workflow as WorkflowEntity;
const scopes = ('scopes' in workflow ? (workflow.scopes as string[]) : undefined) ?? [];
return {
@@ -200,6 +214,7 @@ export async function searchWorkflows(
scopes,
canExecute: scopes.includes('workflow:execute'),
availableInMCP: settings?.availableInMCP ?? false,
tags: toTagSummary(workflowTags),
};
});
@@ -1,4 +1,6 @@
import type { GlobalConfig } from '@n8n/config';
import { type User, type SharedWorkflowRepository, WorkflowEntity } from '@n8n/db';
import { hasGlobalScope } from '@n8n/permissions';
import type { WorkflowJSON } from '@n8n/workflow-sdk';
import z from 'zod';
@@ -21,6 +23,7 @@ import type { CollaborationService } from '@/collaboration/collaboration.service
import type { CredentialsService } from '@/credentials/credentials.service';
import type { DataTableUserOperations } from '@/modules/data-table/data-table-proxy.service';
import type { NodeTypes } from '@/node-types';
import type { TagService } from '@/services/tag.service';
import type { UrlService } from '@/services/url.service';
import type { Telemetry } from '@/telemetry';
import { resolveNodeWebhookIds } from '@/workflow-helpers';
@@ -44,6 +47,8 @@ const operationTypeSchema = z.enum([
'setNodeDisabled',
'setNodeSettings',
'setWorkflowMetadata',
'addTags',
'removeTags',
]);
const positionInputSchema = z.array(z.number()).length(2).describe('Canvas [x, y].');
@@ -111,6 +116,7 @@ const operationInputSchema = z
settings: nodeSettingsInputSchema.optional().describe('For setNodeSettings.'),
name: z.string().max(128).optional().describe('Only used for setWorkflowMetadata.'),
description: z.string().max(255).optional().describe('Only used for setWorkflowMetadata.'),
names: z.array(z.string()).optional().describe('For addTags / removeTags.'),
})
.describe('Workflow update operation. Provide fields matching type.');
@@ -228,6 +234,8 @@ export const createUpdateWorkflowTool = (
sharedWorkflowRepository: SharedWorkflowRepository,
collaborationService: CollaborationService,
dataTableOps: DataTableUserOperations,
tagService: TagService,
globalConfig: GlobalConfig,
): ToolDefinition<typeof inputSchema> => ({
name: MCP_UPDATE_WORKFLOW_TOOL.toolName,
config: {
@@ -266,16 +274,31 @@ export const createUpdateWorkflowTool = (
try {
const strictOperations = parseStrictOperations(operations);
const hasTagOperations = strictOperations.some(
(op) => op.type === 'addTags' || op.type === 'removeTags',
);
if (hasTagOperations && globalConfig.tags.disabled) {
throw new Error(
'Tag operations are not supported on this instance because tags are disabled.',
);
}
const existingWorkflow = await getMcpWorkflow(
workflowId,
user,
['workflow:update'],
workflowFinderService,
{ includeTags: hasTagOperations },
);
await collaborationService.ensureWorkflowEditable(existingWorkflow.id);
const result = applyOperations(toWorkflowSlice(existingWorkflow), strictOperations);
const result = applyOperations(
toWorkflowSlice(existingWorkflow, { includeTags: hasTagOperations }),
strictOperations,
);
if (!result.success) {
throw new Error(result.error);
@@ -316,6 +339,10 @@ export const createUpdateWorkflowTool = (
throw new Error(dataTableCheck.error);
}
const hasNonTagOperations = strictOperations.some(
(op) => op.type !== 'addTags' && op.type !== 'removeTags',
);
const workflowUpdateData = new WorkflowEntity();
Object.assign(workflowUpdateData, {
name: result.workflow.name,
@@ -324,11 +351,13 @@ export const createUpdateWorkflowTool = (
: {}),
nodes: result.workflow.nodes,
connections: result.workflow.connections,
meta: {
...(existingWorkflow.meta ?? {}),
aiBuilderAssisted: true,
builderVariant: 'mcp',
},
meta: hasNonTagOperations
? {
...(existingWorkflow.meta ?? {}),
aiBuilderAssisted: true,
builderVariant: 'mcp',
}
: (existingWorkflow.meta ?? {}),
});
resolveNodeWebhookIds(workflowUpdateData, nodeTypes);
@@ -363,9 +392,30 @@ export const createUpdateWorkflowTool = (
connections: workflowUpdateData.connections,
} as unknown as WorkflowJSON);
let tagIds: string[] | undefined;
if (result.tagNames !== undefined) {
if (hasGlobalScope(user, 'tag:create')) {
const resolvedTags = await tagService.findOrCreateByNames(result.tagNames);
tagIds = resolvedTags.map((t) => t.id);
} else {
const resolvedTags = await tagService.findByNames(result.tagNames);
const resolvedNames = new Set(resolvedTags.map((t) => t.name));
const missing = result.tagNames
.map((n) => n.trim())
.filter((name) => name.length > 0 && !resolvedNames.has(name));
if (missing.length > 0) {
throw new Error(
`Cannot apply the following tags because they don't exist and your account does not have permission to create them: ${missing.join(', ')}`,
);
}
tagIds = resolvedTags.map((t) => t.id);
}
}
const updatedWorkflow = await workflowService.update(user, workflowUpdateData, workflowId, {
aiBuilderAssisted: true,
aiBuilderAssisted: hasNonTagOperations,
source: 'n8n-mcp',
...(tagIds !== undefined ? { tagIds } : {}),
});
void collaborationService.broadcastWorkflowUpdate(workflowId, user.id).catch(() => {});
@@ -168,6 +168,22 @@ export const partialUpdateOperationSchema = z.discriminatedUnion('type', [
name: z.string().max(128).optional(),
description: z.string().max(255).optional(),
}),
z.object({
type: z.literal('addTags'),
names: z
.array(z.string().trim().min(1).max(24))
.min(1)
.max(50)
.describe('Tag names to attach. Unknown names are auto-created. Idempotent.'),
}),
z.object({
type: z.literal('removeTags'),
names: z
.array(z.string().trim().min(1).max(24))
.min(1)
.max(50)
.describe('Tag names to detach from the workflow. Unknown names are ignored.'),
}),
]);
export type PartialUpdateOperation = z.infer<typeof partialUpdateOperationSchema>;
@@ -177,12 +193,16 @@ interface WorkflowSlice {
description?: string;
nodes: INode[];
connections: IConnections;
/** Existing tag names on the workflow. Undefined when not loaded; tag ops require this. */
tagNames?: string[];
}
export interface ApplyOperationsSuccess {
success: true;
workflow: WorkflowSlice;
addedNodeNames: string[];
/** Final tag set after applying tag ops. Undefined means "leave unchanged". */
tagNames?: string[];
}
export interface ApplyOperationsFailure {
@@ -198,6 +218,7 @@ const cloneWorkflow = (workflow: WorkflowSlice): WorkflowSlice => ({
description: workflow.description,
nodes: workflow.nodes.map((node) => structuredClone(node)),
connections: structuredClone(workflow.connections),
tagNames: workflow.tagNames ? [...workflow.tagNames] : undefined,
});
const isPlainObject = (value: unknown): value is Record<string, unknown> =>
@@ -380,6 +401,9 @@ export function applyOperations(
const workflow = cloneWorkflow(input);
const nodeByName = new Map(workflow.nodes.map((n) => [n.name, n]));
const addedNodeNames = new Set<string>();
// Tag set is null until the first tag op runs; that keeps "no tag ops"
// distinguishable from "tag ops applied to an empty set" at return time.
let tagSet: Set<string> | null = null;
for (let i = 0; i < operations.length; i++) {
const op = operations[i];
@@ -566,6 +590,20 @@ export function applyOperations(
break;
}
case 'addTags':
case 'removeTags': {
if (workflow.tagNames === undefined) {
return fail(i, 'tag operations require existing tags to be loaded');
}
if (tagSet === null) tagSet = new Set(workflow.tagNames);
if (op.type === 'addTags') {
for (const name of op.names) tagSet.add(name);
} else {
for (const name of op.names) tagSet.delete(name);
}
break;
}
default: {
op satisfies never;
return fail(i, 'unknown operation type');
@@ -573,18 +611,39 @@ export function applyOperations(
}
}
return { success: true, workflow, addedNodeNames: [...addedNodeNames] };
if (tagSet !== null) {
workflow.tagNames = [...tagSet];
}
return {
success: true,
workflow,
addedNodeNames: [...addedNodeNames],
tagNames: tagSet !== null ? [...tagSet] : undefined,
};
}
/**
* Pick only the fields the partial-update path needs from a workflow entity.
* Keeps the surface explicit and avoids mutating the loaded entity.
*/
export function toWorkflowSlice(workflow: IWorkflowBase): WorkflowSlice {
export function toWorkflowSlice(
workflow: IWorkflowBase,
options: { includeTags?: boolean } = {},
): WorkflowSlice {
let tagNames: string[] | undefined;
if (options.includeTags) {
const tags = (workflow as { tags?: Array<{ name: string }> }).tags;
if (tags === undefined) {
throw new Error('toWorkflowSlice: includeTags=true requires the tags relation to be loaded.');
}
tagNames = tags.map((t) => t.name);
}
return {
name: workflow.name ?? '',
description: (workflow as { description?: string }).description,
nodes: workflow.nodes,
connections: workflow.connections,
tagNames,
};
}
@@ -39,6 +39,7 @@ export type FoundWorkflow = NonNullable<
export type GetMcpWorkflowOptions = {
includeActiveVersion?: boolean;
includeTags?: boolean;
};
/**
@@ -56,6 +57,7 @@ export async function getMcpWorkflow(
): Promise<FoundWorkflow> {
const workflow = await workflowFinderService.findWorkflowForUser(workflowId, user, scopes, {
includeActiveVersion: options?.includeActiveVersion,
includeTags: options?.includeTags,
});
if (!workflow) {
@@ -0,0 +1,223 @@
import type { TagEntity, TagRepository } from '@n8n/db';
// eslint-disable-next-line n8n-local-rules/misplaced-n8n-typeorm-import
import { FindOperator, QueryFailedError } from '@n8n/typeorm';
import { mock } from 'jest-mock-extended';
import type { ExternalHooks } from '@/external-hooks';
import { TagService } from '@/services/tag.service';
const makeTag = (overrides: Partial<TagEntity> = {}): TagEntity =>
({
id: 'tag-id',
name: 'tag',
createdAt: new Date('2026-01-01'),
updatedAt: new Date('2026-01-01'),
...overrides,
}) as TagEntity;
describe('TagService', () => {
const tagRepository = mock<TagRepository>();
const externalHooks = mock<ExternalHooks>();
const tagService = new TagService(externalHooks, tagRepository);
beforeEach(() => {
jest.resetAllMocks();
});
describe('listWithUsageCount', () => {
test('builds a limited ordered query and returns data + totalCount in parallel', async () => {
const limitFn = jest.fn().mockReturnThis();
const orderByFn = jest.fn().mockReturnThis();
const orderByCallOrder: jest.Mock = orderByFn;
const limitCallOrder: jest.Mock = limitFn;
const getMany = jest.fn().mockResolvedValue([makeTag()]);
const builder = {
select: jest.fn().mockReturnThis(),
loadRelationCountAndMap: jest.fn().mockReturnThis(),
orderBy: orderByFn,
limit: limitFn,
getMany,
};
tagRepository.createQueryBuilder.mockReturnValue(builder as never);
tagRepository.count.mockResolvedValue(42);
const result = await tagService.listWithUsageCount({ limit: 10 });
expect(orderByFn).toHaveBeenCalledWith('tag.name', 'ASC');
expect(limitFn).toHaveBeenCalledWith(10);
// orderBy must run before limit, or generated SQL is invalid
expect(orderByCallOrder.mock.invocationCallOrder[0]).toBeLessThan(
limitCallOrder.mock.invocationCallOrder[0],
);
expect(tagRepository.count).toHaveBeenCalledTimes(1);
expect(result.totalCount).toBe(42);
expect(result.data).toHaveLength(1);
});
test('does not order when called via getAll without orderByName', async () => {
const orderByFn = jest.fn().mockReturnThis();
const builder = {
select: jest.fn().mockReturnThis(),
loadRelationCountAndMap: jest.fn().mockReturnThis(),
orderBy: orderByFn,
limit: jest.fn().mockReturnThis(),
getMany: jest.fn().mockResolvedValue([]),
};
tagRepository.createQueryBuilder.mockReturnValue(builder as never);
await tagService.getAll({ withUsageCount: true });
expect(orderByFn).not.toHaveBeenCalled();
});
});
describe('findByNames', () => {
test('returns existing tags by name, preserving input order', async () => {
const tags = [
makeTag({ id: 'tag-1', name: 'production' }),
makeTag({ id: 'tag-2', name: 'critical' }),
];
tagRepository.find.mockResolvedValue(tags);
const result = await tagService.findByNames(['critical', 'production', 'missing']);
expect(result.map((t) => t.name)).toEqual(['critical', 'production']);
expect(tagRepository.save).not.toHaveBeenCalled();
});
test('returns empty for whitespace-only inputs without querying', async () => {
const result = await tagService.findByNames(['', ' ']);
expect(result).toEqual([]);
expect(tagRepository.find).not.toHaveBeenCalled();
});
});
describe('findOrCreateByNames', () => {
test('returns empty array for empty input', async () => {
const result = await tagService.findOrCreateByNames([]);
expect(result).toEqual([]);
expect(tagRepository.find).not.toHaveBeenCalled();
});
test('returns empty array when all inputs are whitespace', async () => {
const result = await tagService.findOrCreateByNames(['', ' ']);
expect(result).toEqual([]);
expect(tagRepository.find).not.toHaveBeenCalled();
});
test('returns existing tags without creating', async () => {
const existing = [
makeTag({ id: 'tag-1', name: 'production' }),
makeTag({ id: 'tag-2', name: 'critical' }),
];
tagRepository.find.mockResolvedValue(existing);
const result = await tagService.findOrCreateByNames(['production', 'critical']);
expect(result.map((t) => t.id)).toEqual(['tag-1', 'tag-2']);
expect(tagRepository.save).not.toHaveBeenCalled();
});
test('deduplicates the input case-insensitively before lookup', async () => {
tagRepository.find.mockResolvedValue([]);
const createdTag = makeTag({ id: 'tag-new', name: 'Prod' });
tagRepository.create.mockReturnValue(createdTag);
tagRepository.save.mockResolvedValue(createdTag);
const result = await tagService.findOrCreateByNames(['Prod', 'prod', 'PROD']);
expect(tagRepository.find).toHaveBeenCalledTimes(1);
const findArg = tagRepository.find.mock.calls[0][0] as unknown as {
where: { name: FindOperator<string[]> };
};
expect(findArg.where.name).toBeInstanceOf(FindOperator);
expect(findArg.where.name.value).toEqual(['Prod']);
// One create using the first-seen original case.
expect(tagRepository.save).toHaveBeenCalledTimes(1);
expect(tagRepository.create).toHaveBeenCalledWith({ name: 'Prod' });
expect(result.map((t) => t.name)).toEqual(['Prod']);
});
test('matches the DB case-sensitively (parity with REST tags API)', async () => {
tagRepository.find.mockResolvedValue([makeTag({ id: 'tag-1', name: 'Production' })]);
const createdTag = makeTag({ id: 'tag-new', name: 'production' });
tagRepository.create.mockReturnValue(createdTag);
tagRepository.save.mockResolvedValue(createdTag);
// 'Production' exists in DB; user asks for 'production' — the existing
// REST contract lets these coexist, so MCP creates a new tag too.
const result = await tagService.findOrCreateByNames(['production']);
expect(tagRepository.save).toHaveBeenCalledTimes(1);
expect(tagRepository.create).toHaveBeenCalledWith({ name: 'production' });
expect(result.map((t) => t.id)).toEqual(['tag-new']);
});
const uniqueViolationError = (code: string | number) => {
const driver = Object.assign(new Error('duplicate key'), { code });
return new QueryFailedError('insert', undefined, driver);
};
test('returns the now-existing row when a concurrent caller wins the create race (postgres)', async () => {
tagRepository.find.mockResolvedValue([]);
const racedTag = makeTag({ id: 'tag-raced', name: 'critical' });
tagRepository.create.mockReturnValue(racedTag);
tagRepository.save.mockRejectedValueOnce(uniqueViolationError('23505'));
tagRepository.findOneBy.mockResolvedValue(racedTag);
const result = await tagService.findOrCreateByNames(['critical']);
expect(result).toEqual([racedTag]);
expect(tagRepository.findOneBy).toHaveBeenCalledWith({ name: 'critical' });
});
test('recognises sqlite unique-constraint code', async () => {
tagRepository.find.mockResolvedValue([]);
const racedTag = makeTag({ id: 'tag-raced', name: 'critical' });
tagRepository.create.mockReturnValue(racedTag);
tagRepository.save.mockRejectedValueOnce(uniqueViolationError('SQLITE_CONSTRAINT_UNIQUE'));
tagRepository.findOneBy.mockResolvedValue(racedTag);
const result = await tagService.findOrCreateByNames(['critical']);
expect(result).toEqual([racedTag]);
});
test('rethrows unrelated QueryFailedError instead of masking it as a race', async () => {
tagRepository.find.mockResolvedValue([]);
tagRepository.create.mockReturnValue(makeTag({ name: 'critical' }));
const unrelated = new QueryFailedError('insert', undefined, new Error('connection lost'));
tagRepository.save.mockRejectedValueOnce(unrelated);
await expect(tagService.findOrCreateByNames(['critical'])).rejects.toBe(unrelated);
expect(tagRepository.findOneBy).not.toHaveBeenCalled();
});
test('rethrows when the loser of the race cannot find the row afterwards', async () => {
tagRepository.find.mockResolvedValue([]);
tagRepository.create.mockReturnValue(makeTag({ name: 'critical' }));
const err = uniqueViolationError('23505');
tagRepository.save.mockRejectedValueOnce(err);
tagRepository.findOneBy.mockResolvedValue(null);
await expect(tagService.findOrCreateByNames(['critical'])).rejects.toBe(err);
});
test('creates missing tags and merges with existing', async () => {
tagRepository.find.mockResolvedValue([makeTag({ id: 'tag-1', name: 'production' })]);
const createdTag = makeTag({ id: 'tag-new', name: 'critical' });
tagRepository.create.mockReturnValue(createdTag);
tagRepository.save.mockResolvedValue(createdTag);
const result = await tagService.findOrCreateByNames(['production', 'critical']);
expect(result.map((t) => t.name)).toEqual(['production', 'critical']);
expect(tagRepository.save).toHaveBeenCalledTimes(1);
});
});
});
+108 -6
View File
@@ -1,6 +1,8 @@
import type { TagEntity, ITagWithCountDb } from '@n8n/db';
import { TagRepository } from '@n8n/db';
import { Service } from '@n8n/di';
// eslint-disable-next-line n8n-local-rules/misplaced-n8n-typeorm-import
import { In, QueryFailedError } from '@n8n/typeorm';
import { ExternalHooks } from '@/external-hooks';
import { validateEntity } from '@/generic-helpers';
@@ -9,6 +11,34 @@ type GetAllResult<T> = T extends { withUsageCount: true } ? ITagWithCountDb[] :
type Action = 'Create' | 'Update';
// Trim and dedupe input names case-insensitively, keeping the first-seen case.
// Inputs are matched against the DB exactly, so this only collapses obvious
// duplicates like ['Prod','prod','PROD'] within a single batch — it does NOT
// change the case-sensitive contract that the REST tag API exposes.
function dedupeNamesPreservingCase(names: string[]): string[] {
const seen = new Set<string>();
const result: string[] = [];
for (const raw of names) {
const trimmed = raw.trim();
if (trimmed.length === 0) continue;
const key = trimmed.toLowerCase();
if (seen.has(key)) continue;
seen.add(key);
result.push(trimmed);
}
return result;
}
// n8n supports postgres (SQLSTATE 23505) and sqlite (SQLITE_CONSTRAINT_UNIQUE,
// or the older SQLITE_CONSTRAINT with a "UNIQUE constraint" message).
function isUniqueConstraintViolation(error: unknown): error is QueryFailedError {
if (!(error instanceof QueryFailedError)) return false;
const driver = (error as { driverError?: { code?: unknown } }).driverError;
const code = driver && typeof driver.code !== 'undefined' ? String(driver.code) : undefined;
if (code === '23505' || code === 'SQLITE_CONSTRAINT_UNIQUE') return true;
return code === 'SQLITE_CONSTRAINT' && /UNIQUE constraint/i.test(error.message);
}
@Service()
export class TagService {
constructor(
@@ -46,23 +76,29 @@ export class TagService {
return await deleteResult;
}
async getAll<T extends { withUsageCount: boolean }>(options?: T): Promise<GetAllResult<T>> {
async getAll<T extends { withUsageCount: boolean; limit?: number; orderByName?: boolean }>(
options?: T,
): Promise<GetAllResult<T>> {
if (options?.withUsageCount) {
const tags = await this.tagRepository
const qb = this.tagRepository
.createQueryBuilder('tag')
.select(['tag.id', 'tag.name', 'tag.createdAt', 'tag.updatedAt'])
.loadRelationCountAndMap('tag.usageCount', 'tag.workflowMappings', 'wm', (qb) =>
qb.leftJoin('wm.workflows', 'workflow').where('workflow.isArchived = :isArchived', {
.loadRelationCountAndMap('tag.usageCount', 'tag.workflowMappings', 'wm', (qb2) =>
qb2.leftJoin('wm.workflows', 'workflow').where('workflow.isArchived = :isArchived', {
isArchived: false,
}),
)
.getMany();
);
if (options.orderByName) qb.orderBy('tag.name', 'ASC');
if (options.limit !== undefined) qb.limit(options.limit);
const tags = await qb.getMany();
return tags as GetAllResult<T>;
}
return await (this.tagRepository.find({
select: ['id', 'name', 'createdAt', 'updatedAt'],
...(options?.orderByName ? { order: { name: 'ASC' as const } } : {}),
...(options?.limit !== undefined ? { take: options.limit } : {}),
}) as Promise<GetAllResult<T>>);
}
@@ -72,6 +108,21 @@ export class TagService {
});
}
/**
* Paginated tags with non-archived usage counts plus the total count, both
* via DB-level queries. Runs the data query and `count` in parallel.
*/
async listWithUsageCount({ limit }: { limit: number }): Promise<{
data: ITagWithCountDb[];
totalCount: number;
}> {
const [data, totalCount] = await Promise.all([
this.getAll({ withUsageCount: true, limit, orderByName: true }),
this.tagRepository.count(),
]);
return { data, totalCount };
}
/**
* Sort tags based on the order of the tag IDs in the request.
*/
@@ -83,4 +134,55 @@ export class TagService {
return requestOrder.map((tagId) => tagMap[tagId]);
}
/**
* Look up tags by name; never creates. Input is deduped case-insensitively
* but matched against the DB exactly (REST tag API contract).
*/
async findByNames(names: string[]): Promise<TagEntity[]> {
const uniqueNames = dedupeNamesPreservingCase(names);
if (uniqueNames.length === 0) return [];
const existing = await this.tagRepository.find({ where: { name: In(uniqueNames) } });
const existingByName = new Map(existing.map((t) => [t.name, t]));
const result: TagEntity[] = [];
for (const name of uniqueNames) {
const hit = existingByName.get(name);
if (hit) result.push(hit);
}
return result;
}
/**
* Resolve names to tag entities, creating any missing. Input is deduped
* case-insensitively (first-seen case wins) but matched against the DB
* exactly. Race-safe against concurrent same-name creates.
*/
async findOrCreateByNames(names: string[]): Promise<TagEntity[]> {
const uniqueNames = dedupeNamesPreservingCase(names);
if (uniqueNames.length === 0) return [];
const existing = await this.tagRepository.find({ where: { name: In(uniqueNames) } });
const existingByName = new Map(existing.map((t) => [t.name, t]));
const result: TagEntity[] = [];
for (const name of uniqueNames) {
const hit = existingByName.get(name);
if (hit) {
result.push(hit);
continue;
}
try {
const created = await this.save(this.toEntity({ name }), 'create');
result.push(created);
} catch (error) {
if (!isUniqueConstraintViolation(error)) throw error;
const raced = await this.tagRepository.findOneBy({ name });
if (!raced) throw error;
result.push(raced);
}
}
return result;
}
}
@@ -91,6 +91,7 @@ export interface SearchWorkflowsResult {
scopes: string[];
canExecute: boolean;
availableInMCP: boolean;
tags: Array<{ id: string; name: string }>;
}>;
count: number;
}