chore(ai-builder): Add workflow naming, compaction, and session cleanup to multi-agent (no-changelog) (#22646)

Signed-off-by: Oleg Ivaniv <me@olegivaniv.com>
This commit is contained in:
oleg
2025-12-03 17:50:22 +01:00
committed by GitHub
parent f5d144cfaf
commit f2930e2db9
29 changed files with 1021 additions and 98 deletions
@@ -88,7 +88,6 @@ export function createAgent(
llmSimpleTask: llm,
llmComplexTask: llm,
checkpointer: new MemorySaver(),
enableMultiAgent: true,
tracer,
});
}
@@ -77,6 +77,8 @@ export interface ResponderContext {
discoveryContext?: DiscoveryContext | null;
/** Current workflow state */
workflowJSON: SimpleWorkflow;
/** Summary of previous conversation (from compaction) */
previousSummary?: string;
}
/**
@@ -98,7 +100,20 @@ export class ResponderAgent {
private buildContextMessage(context: ResponderContext): HumanMessage | null {
const contextParts: string[] = [];
// Check for errors first - if there's an error, surface it prominently
// Previous conversation summary (from compaction)
if (context.previousSummary) {
contextParts.push(`**Previous Conversation Summary:**\n${context.previousSummary}`);
}
// Check for state management actions (compact/clear)
const stateManagementEntry = context.coordinationLog.find(
(e) => e.phase === 'state_management',
);
if (stateManagementEntry) {
contextParts.push(`**State Management:** ${stateManagementEntry.summary}`);
}
// Check for errors - if there's an error, surface it prominently
const errorEntry = getErrorEntry(context.coordinationLog);
if (errorEntry) {
contextParts.push(
@@ -94,6 +94,8 @@ export interface SupervisorContext {
workflowJSON: SimpleWorkflow;
/** Coordination log tracking subgraph completion */
coordinationLog: CoordinationLogEntry[];
/** Summary of previous conversation (from compaction) */
previousSummary?: string;
}
/**
@@ -115,14 +117,21 @@ export class SupervisorAgent {
private buildContextMessage(context: SupervisorContext): HumanMessage | null {
const contextParts: string[] = [];
// 1. Workflow summary (node count and names only)
// 1. Previous conversation summary (from compaction)
if (context.previousSummary) {
contextParts.push('<previous_conversation_summary>');
contextParts.push(context.previousSummary);
contextParts.push('</previous_conversation_summary>');
}
// 2. Workflow summary (node count and names only)
if (context.workflowJSON.nodes.length > 0) {
contextParts.push('<workflow_summary>');
contextParts.push(buildWorkflowSummary(context.workflowJSON));
contextParts.push('</workflow_summary>');
}
// 2. Coordination log summary (what phases completed)
// 3. Coordination log summary (what phases completed)
if (context.coordinationLog.length > 0) {
contextParts.push('<completed_phases>');
contextParts.push(summarizeCoordinationLog(context.coordinationLog));
@@ -12,7 +12,11 @@ import type { IUser, INodeTypeDescription, ITelemetryTrackProperties } from 'n8n
import { LLMServiceError } from '@/errors';
import { anthropicClaudeSonnet45 } from '@/llm-config';
import { SessionManagerService } from '@/session-manager.service';
import { WorkflowBuilderAgent, type ChatPayload } from '@/workflow-builder-agent';
import {
BuilderFeatureFlags,
WorkflowBuilderAgent,
type ChatPayload,
} from '@/workflow-builder-agent';
type OnCreditsUpdated = (userId: string, creditsQuota: number, creditsClaimed: number) => void;
@@ -29,6 +33,7 @@ export class AiWorkflowBuilderService {
private readonly logger?: Logger,
private readonly instanceId?: string,
private readonly instanceUrl?: string,
private readonly n8nVersion?: string,
private readonly onCreditsUpdated?: OnCreditsUpdated,
private readonly onTelemetryEvent?: OnTelemetryEvent,
) {
@@ -153,7 +158,7 @@ export class AiWorkflowBuilderService {
});
}
private async getAgent(user: IUser) {
private async getAgent(user: IUser, featureFlags?: BuilderFeatureFlags) {
const { anthropicClaude, tracingClient, authHeaders } = await this.setupModels(user);
const agent = new WorkflowBuilderAgent({
@@ -162,7 +167,6 @@ export class AiWorkflowBuilderService {
llmSimpleTask: anthropicClaude,
llmComplexTask: anthropicClaude,
logger: this.logger,
enableMultiAgent: process.env.N8N_ENABLE_MULTI_AGENT === 'true',
checkpointer: this.sessionManager.getCheckpointer(),
tracer: tracingClient
? new LangChainTracer({ client: tracingClient, projectName: 'n8n-workflow-builder' })
@@ -171,6 +175,10 @@ export class AiWorkflowBuilderService {
onGenerationSuccess: async () => {
await this.onGenerationSuccess(user, authHeaders);
},
runMetadata: {
n8nVersion: this.n8nVersion,
featureFlags: featureFlags ?? {},
},
});
return agent;
@@ -200,7 +208,7 @@ export class AiWorkflowBuilderService {
}
async *chat(payload: ChatPayload, user: IUser, abortSignal?: AbortSignal) {
const agent = await this.getAgent(user);
const agent = await this.getAgent(user, payload.featureFlags);
const userId = user?.id?.toString();
const workflowId = payload.workflowContext?.currentWorkflow?.id;
@@ -0,0 +1,124 @@
export const SWITCH_NODE_EXAMPLES = `
### Switch Node Examples
#### Example 1: Route by Amount Tiers (Purchase Approval)
Current Parameters: { "mode": "rules" }
Requested Changes: Route under $100 to auto-approve, $100-$1000 to manager, over $1000 to finance
Expected Output:
{
"mode": "rules",
"rules": {
"values": [
{
"conditions": {
"options": { "caseSensitive": true, "leftValue": "", "typeValidation": "strict" },
"conditions": [
{
"leftValue": "={{ $json.amount }}",
"rightValue": 100,
"operator": { "type": "number", "operation": "lt" }
}
],
"combinator": "and"
},
"renameOutput": true,
"outputKey": "Auto-Approve"
},
{
"conditions": {
"options": { "caseSensitive": true, "leftValue": "", "typeValidation": "strict" },
"conditions": [
{
"leftValue": "={{ $json.amount }}",
"rightValue": 100,
"operator": { "type": "number", "operation": "gte" }
},
{
"leftValue": "={{ $json.amount }}",
"rightValue": 1000,
"operator": { "type": "number", "operation": "lte" }
}
],
"combinator": "and"
},
"renameOutput": true,
"outputKey": "Manager Review"
},
{
"conditions": {
"options": { "caseSensitive": true, "leftValue": "", "typeValidation": "strict" },
"conditions": [
{
"leftValue": "={{ $json.amount }}",
"rightValue": 1000,
"operator": { "type": "number", "operation": "gt" }
}
],
"combinator": "and"
},
"renameOutput": true,
"outputKey": "Finance Review"
}
]
}
}
#### Example 2: Route by Status String
Current Parameters: { "mode": "rules" }
Requested Changes: Route by order status - pending, processing, completed
Expected Output:
{
"mode": "rules",
"rules": {
"values": [
{
"conditions": {
"options": { "caseSensitive": false, "leftValue": "", "typeValidation": "loose" },
"conditions": [
{
"leftValue": "={{ $json.status }}",
"rightValue": "pending",
"operator": { "type": "string", "operation": "equals" }
}
],
"combinator": "and"
},
"renameOutput": true,
"outputKey": "Pending"
},
{
"conditions": {
"options": { "caseSensitive": false, "leftValue": "", "typeValidation": "loose" },
"conditions": [
{
"leftValue": "={{ $json.status }}",
"rightValue": "processing",
"operator": { "type": "string", "operation": "equals" }
}
],
"combinator": "and"
},
"renameOutput": true,
"outputKey": "Processing"
},
{
"conditions": {
"options": { "caseSensitive": false, "leftValue": "", "typeValidation": "loose" },
"conditions": [
{
"leftValue": "={{ $json.status }}",
"rightValue": "completed",
"operator": { "type": "string", "operation": "equals" }
}
],
"combinator": "and"
},
"renameOutput": true,
"outputKey": "Completed"
}
]
}
}
`;
@@ -0,0 +1,68 @@
export const SWITCH_NODE_GUIDE = `
### Switch Node Configuration Guide
The Switch node routes items to different outputs based on conditions. Uses the same filter structure as IF node but for multi-way branching.
#### Switch Node Structure (mode: 'rules')
\`\`\`json
{
"mode": "rules",
"rules": {
"values": [
{
"conditions": {
"options": {
"caseSensitive": true,
"leftValue": "",
"typeValidation": "strict"
},
"conditions": [
{
"leftValue": "={{ $json.amount }}",
"rightValue": 100,
"operator": {
"type": "number",
"operation": "lt"
}
}
],
"combinator": "and"
},
"renameOutput": true,
"outputKey": "Under $100"
}
]
}
}
\`\`\`
#### Key Points:
1. Each entry in rules.values[] creates ONE output
2. Conditions use the same filter structure as IF node
3. Multiple conditions per rule are combined with combinator ("and" or "or")
4. Use renameOutput: true + outputKey to label outputs descriptively
#### Numeric Operators
- lt: Less than
- gt: Greater than
- lte: Less than or equal
- gte: Greater than or equal
- equals: Equal to
#### String Operators
- equals: Exact match
- contains: Contains substring
- startsWith: Starts with
- endsWith: Ends with
#### Common Patterns:
**Numeric Range Routing** (for ranges like $100-$1000):
Use two conditions with combinator: "and":
- First condition: gte (greater than or equal to lower bound)
- Second condition: lte (less than or equal to upper bound)
**String-Based Routing** (status/type values):
- Use type: "string" with operation: "equals"
- Set caseSensitive: false in options for case-insensitive matching
`;
@@ -9,9 +9,11 @@ import { TOOL_NODE_EXAMPLES } from './examples/advanced/tool-node-examples';
import { IF_NODE_EXAMPLES } from './examples/basic/if-node-examples';
import { SET_NODE_EXAMPLES } from './examples/basic/set-node-examples';
import { SIMPLE_UPDATE_EXAMPLES } from './examples/basic/simple-updates';
import { SWITCH_NODE_EXAMPLES } from './examples/basic/switch-node-examples';
import { HTTP_REQUEST_GUIDE } from './node-types/http-request';
import { IF_NODE_GUIDE } from './node-types/if-node';
import { SET_NODE_GUIDE } from './node-types/set-node';
import { SWITCH_NODE_GUIDE } from './node-types/switch-node';
import { TOOL_NODES_GUIDE } from './node-types/tool-nodes';
import { RESOURCE_LOCATOR_GUIDE } from './parameter-types/resource-locator';
import { SYSTEM_MESSAGE_GUIDE } from './parameter-types/system-message';
@@ -42,6 +44,8 @@ export class ParameterUpdatePromptBuilder {
sections.push(SET_NODE_GUIDE);
} else if (this.isIfNode(context.nodeType)) {
sections.push(IF_NODE_GUIDE);
} else if (this.isSwitchNode(context.nodeType)) {
sections.push(SWITCH_NODE_GUIDE);
} else if (this.isHttpRequestNode(context.nodeType)) {
sections.push(HTTP_REQUEST_GUIDE);
}
@@ -130,6 +134,14 @@ export class ParameterUpdatePromptBuilder {
return category === 'if';
}
/**
* Checks if node is a Switch node
*/
private static isSwitchNode(nodeType: string): boolean {
const category = getNodeTypeCategory(nodeType);
return category === 'switch';
}
/**
* Checks if node is an HTTP Request node
*/
@@ -179,6 +191,8 @@ export class ParameterUpdatePromptBuilder {
examples.push(SET_NODE_EXAMPLES);
} else if (this.isIfNode(context.nodeType)) {
examples.push(IF_NODE_EXAMPLES);
} else if (this.isSwitchNode(context.nodeType)) {
examples.push(SWITCH_NODE_EXAMPLES);
}
// Add resource locator examples if needed
if (context.hasResourceLocatorParams) {
@@ -4,6 +4,7 @@ export const DEFAULT_PROMPT_CONFIG: NodePromptConfig = {
nodeTypePatterns: {
set: ['n8n-nodes-base.set', 'set'],
if: ['n8n-nodes-base.if', 'if', 'filter'],
switch: ['n8n-nodes-base.switch', 'switch'],
httpRequest: ['n8n-nodes-base.httpRequest', 'httprequest', 'webhook', 'n8n-nodes-base.webhook'],
tool: ['Tool', '.tool'],
},
@@ -7,6 +7,7 @@ import type { INodeTypeDescription } from 'n8n-workflow';
import { ResponderAgent } from './agents/responder.agent';
import { SupervisorAgent } from './agents/supervisor.agent';
import {
DEFAULT_AUTO_COMPACT_THRESHOLD_TOKENS,
MAX_BUILDER_ITERATIONS,
MAX_CONFIGURATOR_ITERATIONS,
MAX_DISCOVERY_ITERATIONS,
@@ -20,6 +21,13 @@ import type { SubgraphPhase } from './types/coordination';
import { createErrorMetadata } from './types/coordination';
import { getNextPhaseFromLog } from './utils/coordination-log';
import { processOperations } from './utils/operations-processor';
import {
determineStateAction,
handleCleanupDangling,
handleCompactMessages,
handleCreateWorkflowName,
handleDeleteMessages,
} from './utils/state-modifier';
import type { BuilderFeatureFlags } from './workflow-builder-agent';
/**
@@ -43,6 +51,8 @@ export interface MultiAgentSubgraphConfig {
logger?: Logger;
instanceUrl?: string;
checkpointer?: MemorySaver;
/** Token threshold for auto-compaction. Defaults to DEFAULT_AUTO_COMPACT_THRESHOLD_TOKENS */
autoCompactThresholdTokens?: number;
featureFlags?: BuilderFeatureFlags;
}
@@ -107,8 +117,15 @@ function createSubgraphNodeHandler<
* Parent graph orchestrates between subgraphs with minimal shared state.
*/
export function createMultiAgentWorkflowWithSubgraphs(config: MultiAgentSubgraphConfig) {
const { parsedNodeTypes, llmComplexTask, logger, instanceUrl, checkpointer, featureFlags } =
config;
const {
parsedNodeTypes,
llmComplexTask,
logger,
instanceUrl,
checkpointer,
autoCompactThresholdTokens = DEFAULT_AUTO_COMPACT_THRESHOLD_TOKENS,
featureFlags,
} = config;
const supervisorAgent = new SupervisorAgent({ llm: llmComplexTask });
const responderAgent = new ResponderAgent({ llm: llmComplexTask });
@@ -142,6 +159,7 @@ export function createMultiAgentWorkflowWithSubgraphs(config: MultiAgentSubgraph
messages: state.messages,
workflowJSON: state.workflowJSON,
coordinationLog: state.coordinationLog,
previousSummary: state.previousSummary,
});
return {
@@ -155,6 +173,7 @@ export function createMultiAgentWorkflowWithSubgraphs(config: MultiAgentSubgraph
coordinationLog: state.coordinationLog,
discoveryContext: state.discoveryContext,
workflowJSON: state.workflowJSON,
previousSummary: state.previousSummary,
});
return {
@@ -171,6 +190,31 @@ export function createMultiAgentWorkflowWithSubgraphs(config: MultiAgentSubgraph
workflowOperations: [], // Clear operations after processing
};
})
// State modification nodes (preprocessing)
.addNode('check_state', (state) => ({
nextPhase: determineStateAction(state, autoCompactThresholdTokens),
}))
.addNode('cleanup_dangling', (state) => handleCleanupDangling(state.messages, logger))
.addNode('compact_messages', async (state) => {
const isAutoCompact = state.messages[state.messages.length - 1]?.content !== '/compact';
return await handleCompactMessages(
state.messages,
state.previousSummary ?? '',
llmComplexTask,
isAutoCompact,
);
})
.addNode('delete_messages', (state) => handleDeleteMessages(state.messages))
.addNode(
'create_workflow_name',
async (state) =>
await handleCreateWorkflowName(
state.messages,
state.workflowJSON,
llmComplexTask,
logger,
),
)
// Add Subgraph Nodes (using helper to reduce duplication)
.addNode(
'discovery_subgraph',
@@ -206,8 +250,31 @@ export function createMultiAgentWorkflowWithSubgraphs(config: MultiAgentSubgraph
.addEdge('discovery_subgraph', 'process_operations')
.addEdge('builder_subgraph', 'process_operations')
.addEdge('configurator_subgraph', 'process_operations')
// Start flows to supervisor (initial routing only)
.addEdge(START, 'supervisor')
// Start flows to check_state (preprocessing)
.addEdge(START, 'check_state')
// Conditional routing from check_state
.addConditionalEdges('check_state', (state) => {
const routes: Record<string, string> = {
cleanup_dangling: 'cleanup_dangling',
compact_messages: 'compact_messages',
delete_messages: 'delete_messages',
create_workflow_name: 'create_workflow_name',
auto_compact_messages: 'compact_messages', // Reuse same node
continue: 'supervisor',
};
return routes[state.nextPhase] ?? 'supervisor';
})
// Route after state modification nodes
.addEdge('cleanup_dangling', 'check_state') // Re-check after cleanup
.addEdge('delete_messages', 'responder') // Clear → responder for acknowledgment
.addEdge('create_workflow_name', 'supervisor') // Continue after naming
// Compact has conditional routing: auto → continue, manual → responder
.addConditionalEdges('compact_messages', (state) => {
// Auto-compact preserves the last user message, manual /compact clears all
// If messages remain after compaction, it's auto-compact → continue processing
const hasMessages = state.messages.length > 0;
return hasMessages ? 'check_state' : 'responder';
})
// Conditional Edge for Supervisor (initial routing via LLM)
.addConditionalEdges('supervisor', (state) => routeToNode(state.nextPhase))
// Deterministic routing after subgraphs complete (based on coordination log)
@@ -1,5 +1,5 @@
import type { BaseMessage } from '@langchain/core/messages';
import { Annotation } from '@langchain/langgraph';
import { Annotation, messagesStateReducer } from '@langchain/langgraph';
import type { CoordinationLogEntry } from './types/coordination';
import type { DiscoveryContext } from './types/discovery-types';
@@ -17,7 +17,7 @@ import type { ChatPayload } from './workflow-builder-agent';
export const ParentGraphState = Annotation.Root({
// Shared: User's conversation history (for responder)
messages: Annotation<BaseMessage[]>({
reducer: (x, y) => x.concat(y),
reducer: messagesStateReducer,
default: () => [],
}),
@@ -56,6 +56,12 @@ export const ParentGraphState = Annotation.Root({
default: () => [],
}),
// For conversation compaction - stores summarized history
previousSummary: Annotation<string>({
reducer: (x, y) => y ?? x,
default: () => '',
}),
// Template IDs fetched from workflow examples for telemetry
templateIds: Annotation<number[]>({
reducer: appendArrayReducer,
@@ -55,6 +55,7 @@ STEP 3: VALIDATE (REQUIRED)
- After ALL nodes and connections are created, call validate_structure
- This step is MANDATORY - you cannot finish without it
- If validation finds issues (missing trigger, invalid connections), fix them and validate again
- MAXIMUM 3 VALIDATION ATTEMPTS: After 3 calls to validate_structure, proceed to respond regardless of remaining issues
STEP 4: RESPOND TO USER
- Only after validation passes, provide your brief summary
@@ -83,6 +84,12 @@ Placement rules:
For AI-generated structured data, prefer Structured Output Parser nodes over Code nodes.
For binary file data, use Extract From File node to extract content from files before processing.
Use Code nodes only for custom business logic beyond parsing.
STRUCTURED OUTPUT PARSER RULE:
When Discovery results include Structured Output Parser:
1. Create the Structured Output Parser node
2. Set AI Agent's hasOutputParser: true in connectionParameters
3. Connect: Structured Output Parser → AI Agent (ai_outputParser connection)
</data_parsing_strategy>
<proactive_design>
@@ -108,9 +115,19 @@ ALWAYS check node details and set connectionParameters explicitly.
CONNECTION PARAMETERS EXAMPLES:
- Static nodes (HTTP Request, Set, Code): reasoning="Static inputs/outputs", parameters={{}}
- AI Agent with parser: reasoning="hasOutputParser creates additional input", parameters={{ hasOutputParser: true }}
- AI Agent with structured output: reasoning="hasOutputParser enables ai_outputParser input for Structured Output Parser", parameters={{ hasOutputParser: true }}
- Vector Store insert: reasoning="Insert mode requires document input", parameters={{ mode: "insert" }}
- Document Loader custom: reasoning="Custom mode enables text splitter input", parameters={{ textSplittingMode: "custom" }}
- Switch with routing rules: reasoning="Switch needs N outputs, creating N rules.values entries with outputKeys", parameters={{ mode: "rules", rules: {{ values: [...] }} }} - see <switch_node_pattern> for full structure
<structured_output_parser_guidance>
WHEN TO SET hasOutputParser: true on AI Agent:
- Discovery found Structured Output Parser node → MUST set hasOutputParser: true
- AI output will be used in conditions (IF/Switch nodes checking $json.field)
- AI output will be formatted/displayed (HTML emails, reports with specific sections)
- AI output will be stored in database/data tables with specific fields
- AI is classifying, scoring, or extracting specific data fields
</structured_output_parser_guidance>
<node_connections_understanding>
n8n connections flow from SOURCE (output) to TARGET (input).
@@ -158,6 +175,51 @@ Common mistake to avoid:
- Document Loader is an AI sub-node that gives Vector Store document processing capability
</rag_workflow_pattern>
<switch_node_pattern>
For Switch nodes with multiple routing paths:
- The number of outputs is determined by the number of entries in rules.values[]
- You MUST create the rules.values[] array with placeholder entries for each output branch
- Each entry needs: conditions structure (with empty leftValue/rightValue) + renameOutput: true + descriptive outputKey
- Configurator will fill in the actual condition values later
- Use descriptive node names like "Route by Amount" or "Route by Status"
Example connectionParameters for 3-way routing:
{{
"mode": "rules",
"rules": {{
"values": [
{{
"conditions": {{
"options": {{ "caseSensitive": true, "leftValue": "", "typeValidation": "strict" }},
"conditions": [{{ "leftValue": "", "rightValue": "", "operator": {{ "type": "string", "operation": "equals" }} }}],
"combinator": "and"
}},
"renameOutput": true,
"outputKey": "Output 1 Name"
}},
{{
"conditions": {{
"options": {{ "caseSensitive": true, "leftValue": "", "typeValidation": "strict" }},
"conditions": [{{ "leftValue": "", "rightValue": "", "operator": {{ "type": "string", "operation": "equals" }} }}],
"combinator": "and"
}},
"renameOutput": true,
"outputKey": "Output 2 Name"
}},
{{
"conditions": {{
"options": {{ "caseSensitive": true, "leftValue": "", "typeValidation": "strict" }},
"conditions": [{{ "leftValue": "", "rightValue": "", "operator": {{ "type": "string", "operation": "equals" }} }}],
"combinator": "and"
}},
"renameOutput": true,
"outputKey": "Output 3 Name"
}}
]
}}
}}
</switch_node_pattern>
<connection_type_examples>
**Main Connections** (regular data flow):
- Trigger → HTTP Request → Set → Email
@@ -50,6 +50,7 @@ STEP 2: VALIDATE (REQUIRED)
- After ALL configurations complete, call validate_configuration
- This step is MANDATORY - you cannot finish without it
- If validation finds issues, fix them and validate again
- MAXIMUM 3 VALIDATION ATTEMPTS: After 3 calls to validate_configuration, proceed to respond regardless of remaining issues
STEP 3: RESPOND TO USER
- Only after validation passes, provide your response
@@ -85,6 +86,7 @@ CRITICAL PARAMETERS TO ALWAYS SET:
- Set node: Fields to set with values
- Code node: The actual code to execute
- IF node: Conditions to check
- Switch node: Configure rules.values[] with conditions for each output branch (uses same filter structure as IF node)
- Document Loader: dataType parameter ('binary' for files like PDF, 'json' for JSON data)
- AI nodes: Prompts, models, configurations
- Tool nodes: Use $fromAI for dynamic recipient/subject/message fields
@@ -95,6 +97,34 @@ Defaults are traps that cause runtime failures. Examples:
- HTTP Request defaults to GET but APIs often need POST
- Vector Store mode affects available connections - set explicitly (retrieve-as-tool when using with AI Agent)
<switch_node_configuration>
Switch nodes require configuring rules.values[] array - each entry creates one output:
Structure per rule:
{{
"conditions": {{
"options": {{ "caseSensitive": true, "leftValue": "", "typeValidation": "strict" }},
"conditions": [
{{
"leftValue": "={{{{ $json.fieldName }}}}",
"rightValue": <value>,
"operator": {{ "type": "number|string", "operation": "lt|gt|equals|etc" }}
}}
],
"combinator": "and"
}},
"renameOutput": true,
"outputKey": "Descriptive Label"
}}
For numeric ranges (e.g., $100-$1000):
- Use TWO conditions with combinator: "and"
- First: gte (greater than or equal)
- Second: lte (less than or equal)
Always set renameOutput: true and provide descriptive outputKey labels.
</switch_node_configuration>
<response_format>
After validation passes, provide a concise summary:
- List any placeholders requiring user configuration (e.g., "URL placeholder needs actual endpoint")
@@ -325,11 +325,38 @@ A parameter is connection-changing ONLY IF it appears in <input> or <output> exp
- Merge: numberInputs (appears in <input> expression)
- Webhook: responseMode (appears in <output> expression)
<dynamic_output_nodes>
Some nodes have DYNAMIC outputs that depend on parameter values:
**Switch Node** (n8n-nodes-base.switch):
- When mode is "rules", the number of outputs equals the number of routing rules
- Connection parameter: mode: "rules" - CRITICAL for enabling rule-based routing
- Each rule in rules.values[] creates one output
- The rules parameter uses the same filter structure as IF node conditions
- ALWAYS flag mode as connection-changing with possibleValues: ["rules", "expression"]
**Merge Node** (n8n-nodes-base.merge):
- numberInputs parameter controls how many inputs the node accepts
When you find these nodes, ALWAYS flag mode/numberInputs as connection-changing parameters with possibleValues.
</dynamic_output_nodes>
SUB-NODES SEARCHES:
When searching for AI nodes, ALSO search for their required sub-nodes:
- "AI Agent" → also search for "Chat Model", "Memory", "Output Parser"
- "Basic LLM Chain" → also search for "Chat Model", "Output Parser"
- "Vector Store" → also search for "Embeddings", "Document Loader"
STRUCTURED OUTPUT PARSER - WHEN TO INCLUDE:
Search for "Structured Output Parser" (@n8n/n8n-nodes-langchain.outputParserStructured) when:
- AI output will be used programmatically (conditions, formatting, database storage, API calls)
- AI needs to extract specific fields (e.g., score, category, priority, action items)
- AI needs to classify/categorize data into defined categories
- Downstream nodes need to access specific fields from AI response (e.g., $json.score, $json.category)
- Output will be displayed in a formatted way (e.g., HTML email with specific sections)
- Data needs validation against a schema before processing
- Always use search_nodes to find the exact node names and versions - NEVER guess versions
CRITICAL RULES:
@@ -209,6 +209,7 @@ describe('AiWorkflowBuilderService', () => {
mockLogger,
'test-instance-id',
'https://n8n.example.com',
'1.0.0',
mockOnCreditsUpdated,
);
});
@@ -221,6 +222,7 @@ describe('AiWorkflowBuilderService', () => {
mockLogger,
'test-instance-id',
'https://test.com',
'1.0.0',
mockOnCreditsUpdated,
);
@@ -247,6 +249,7 @@ describe('AiWorkflowBuilderService', () => {
mockLogger,
'test-instance-id',
'https://test.com',
'1.0.0',
mockOnCreditsUpdated,
);
@@ -271,6 +274,7 @@ describe('AiWorkflowBuilderService', () => {
mockLogger,
'test-instance-id',
'https://test.com',
'1.0.0',
mockOnCreditsUpdated,
);
@@ -29,6 +29,7 @@ export interface NodePromptConfig {
nodeTypePatterns: {
set: string[];
if: string[];
switch: string[];
httpRequest: string[];
tool: string[];
};
@@ -5,7 +5,7 @@
* and enable deterministic routing without polluting the messages array.
*/
export type SubgraphPhase = 'discovery' | 'builder' | 'configurator';
export type SubgraphPhase = 'discovery' | 'builder' | 'configurator' | 'state_management';
/**
* Entry in the coordination log tracking subgraph completion.
@@ -34,6 +34,7 @@ export type CoordinationMetadata =
| DiscoveryMetadata
| BuilderMetadata
| ConfiguratorMetadata
| StateManagementMetadata
| ErrorMetadata;
export interface DiscoveryMetadata {
@@ -72,6 +73,14 @@ export interface ErrorMetadata {
errorMessage: string;
}
export interface StateManagementMetadata {
phase: 'state_management';
/** Type of state management action */
action: 'compact' | 'clear';
/** Number of messages removed during compaction */
messagesRemoved?: number;
}
/**
* Helper functions to create typed metadata objects.
* These eliminate the need for type assertions when creating coordination log entries.
@@ -93,3 +102,9 @@ export function createConfiguratorMetadata(
export function createErrorMetadata(data: Omit<ErrorMetadata, 'phase'>): ErrorMetadata {
return { phase: 'error', ...data };
}
export function createStateManagementMetadata(
data: Omit<StateManagementMetadata, 'phase'>,
): StateManagementMetadata {
return { phase: 'state_management', ...data };
}
@@ -13,6 +13,7 @@ import type {
DiscoveryMetadata,
BuilderMetadata,
ConfiguratorMetadata,
StateManagementMetadata,
} from '../types/coordination';
export type RoutingDecision = 'discovery' | 'builder' | 'configurator' | 'responder';
@@ -80,10 +81,14 @@ export function getPhaseMetadata(
log: CoordinationLogEntry[],
phase: 'configurator',
): ConfiguratorMetadata | null;
export function getPhaseMetadata(
log: CoordinationLogEntry[],
phase: 'state_management',
): StateManagementMetadata | null;
export function getPhaseMetadata(
log: CoordinationLogEntry[],
phase: SubgraphPhase,
): DiscoveryMetadata | BuilderMetadata | ConfiguratorMetadata | null {
): DiscoveryMetadata | BuilderMetadata | ConfiguratorMetadata | StateManagementMetadata | null {
const entry = getPhaseEntry(log, phase);
if (!entry) return null;
@@ -0,0 +1,195 @@
import type { BaseChatModel } from '@langchain/core/language_models/chat_models';
import type { BaseMessage } from '@langchain/core/messages';
import { HumanMessage, RemoveMessage } from '@langchain/core/messages';
import type { Logger } from '@n8n/backend-common';
import { cleanupDanglingToolCallMessages } from './cleanup-dangling-tool-call-messages';
import { estimateTokenCountFromMessages } from './token-usage';
import { conversationCompactChain } from '../chains/conversation-compact';
import { workflowNameChain } from '../chains/workflow-name';
import type { CoordinationLogEntry } from '../types/coordination';
import { createStateManagementMetadata } from '../types/coordination';
import type { SimpleWorkflow } from '../types/workflow';
export type StateModificationAction =
| 'compact_messages'
| 'delete_messages'
| 'create_workflow_name'
| 'auto_compact_messages'
| 'cleanup_dangling'
| 'continue';
export interface StateModifierInput {
messages: BaseMessage[];
workflowJSON: SimpleWorkflow;
previousSummary?: string;
}
/**
* Determines if state modifications are needed before agent processing.
* Pure function - no side effects, easily testable.
*/
export function determineStateAction(
input: StateModifierInput,
autoCompactThresholdTokens: number,
): StateModificationAction {
const { messages, workflowJSON } = input;
// First check for dangling tool calls (from interrupted sessions)
const danglingMessages = cleanupDanglingToolCallMessages(messages);
if (danglingMessages.length > 0) {
return 'cleanup_dangling';
}
const lastHumanMessage = messages.findLast((m) => m instanceof HumanMessage);
if (!lastHumanMessage) return 'continue';
// Manual /compact command
if (lastHumanMessage.content === '/compact') {
return 'compact_messages';
}
// Manual /clear command
if (lastHumanMessage.content === '/clear') {
return 'delete_messages';
}
// Auto-generate workflow name on first message with empty workflow
const workflowName = workflowJSON?.name;
const nodesLength = workflowJSON?.nodes?.length ?? 0;
const isDefaultName = !workflowName || /^My workflow( \d+)?$/.test(workflowName);
if (isDefaultName && nodesLength === 0 && messages.length === 1) {
return 'create_workflow_name';
}
// Auto-compact when token threshold exceeded
const estimatedTokens = estimateTokenCountFromMessages(messages);
if (estimatedTokens > autoCompactThresholdTokens) {
return 'auto_compact_messages';
}
return 'continue';
}
/**
* Cleans up dangling tool call messages from interrupted sessions.
* Returns state update with RemoveMessage instances.
*/
export function handleCleanupDangling(
messages: BaseMessage[],
logger?: Logger,
): { messages: RemoveMessage[] } {
const messagesToRemove = cleanupDanglingToolCallMessages(messages);
if (messagesToRemove.length > 0) {
logger?.warn('Cleaning up dangling tool call messages', {
count: messagesToRemove.length,
});
}
return { messages: messagesToRemove };
}
/**
* Compacts conversation history by summarizing it.
* Used for both manual /compact and auto-compaction.
*
* For manual /compact: Removes all messages, routes to responder for acknowledgment.
* For auto-compact: Removes old messages, preserves last user message to continue processing.
*/
export async function handleCompactMessages(
messages: BaseMessage[],
previousSummary: string,
llm: BaseChatModel,
isAutoCompact: boolean,
): Promise<{
previousSummary: string;
messages: BaseMessage[];
coordinationLog: CoordinationLogEntry[];
}> {
const lastHumanMessage = messages.findLast((m) => m instanceof HumanMessage);
if (!lastHumanMessage) {
throw new Error('Cannot compact messages: no HumanMessage found');
}
const compactedMessages = await conversationCompactChain(llm, messages, previousSummary);
// For manual /compact: just remove messages, responder will generate acknowledgment
// For auto-compact: remove messages but preserve the last user message to continue processing
const newMessages: BaseMessage[] = [
...messages.map((m) => new RemoveMessage({ id: m.id! })),
...(isAutoCompact ? [new HumanMessage({ content: lastHumanMessage.content })] : []),
];
return {
previousSummary: compactedMessages.summaryPlain,
messages: newMessages,
coordinationLog: [
{
phase: 'state_management',
status: 'completed',
timestamp: Date.now(),
summary: isAutoCompact
? 'Auto-compacted conversation due to token limit'
: 'Manually compacted conversation history',
metadata: createStateManagementMetadata({
action: 'compact',
messagesRemoved: messages.length,
}),
},
],
};
}
/**
* Clears the session by removing all messages and resetting workflow.
*/
export function handleDeleteMessages(messages: BaseMessage[]): {
messages: RemoveMessage[];
workflowJSON: SimpleWorkflow;
previousSummary: string;
discoveryContext: null;
coordinationLog: CoordinationLogEntry[];
workflowOperations: [];
} {
return {
messages: messages.map((m) => new RemoveMessage({ id: m.id! })),
workflowJSON: { nodes: [], connections: {}, name: '' },
previousSummary: '',
discoveryContext: null,
coordinationLog: [
{
phase: 'state_management',
status: 'completed',
timestamp: Date.now(),
summary: 'Cleared session and reset workflow',
metadata: createStateManagementMetadata({ action: 'clear' }),
},
],
workflowOperations: [],
};
}
/**
* Generates a workflow name from the initial user message.
*/
export async function handleCreateWorkflowName(
messages: BaseMessage[],
workflowJSON: SimpleWorkflow,
llm: BaseChatModel,
logger?: Logger,
): Promise<{ workflowJSON: SimpleWorkflow }> {
if (messages.length === 1 && messages[0] instanceof HumanMessage) {
const initialMessage = messages[0];
if (typeof initialMessage.content !== 'string') {
logger?.debug('Initial message content is not a string, skipping workflow name generation');
return { workflowJSON };
}
logger?.debug('Generating workflow name');
const { name } = await workflowNameChain(llm, initialMessage.content);
return {
workflowJSON: { ...workflowJSON, name },
};
}
return { workflowJSON };
}
@@ -159,35 +159,6 @@ export function cleanContextTags(text: string): string {
// CHUNK PROCESSORS
// ============================================================================
/** Handle delete_messages node update */
function processDeleteMessages(update: unknown): StreamOutput | null {
const typed = update as { messages?: MessageContent[] } | undefined;
if (!typed?.messages?.length) return null;
const messageChunk: AgentMessageChunk = {
role: 'assistant',
type: 'message',
text: 'Deleted, refresh?',
};
return { messages: [messageChunk] };
}
/** Handle compact_messages node update */
function processCompactMessages(update: unknown): StreamOutput | null {
const typed = update as { messages?: MessageContent[] } | undefined;
if (!typed?.messages?.length) return null;
const content = extractMessageContent(typed.messages);
if (!content) return null;
const messageChunk: AgentMessageChunk = {
role: 'assistant',
type: 'message',
text: content,
};
return { messages: [messageChunk] };
}
/** Handle process_operations node update */
function processOperationsUpdate(update: unknown): StreamOutput | null {
const typed = update as { workflowJSON?: unknown; workflowOperations?: unknown } | undefined;
@@ -234,16 +205,13 @@ function processToolChunk(chunk: unknown): StreamOutput | null {
/** Process a single chunk from updates stream mode */
function processUpdatesChunk(nodeUpdate: Record<string, unknown>): StreamOutput | null {
// Guard against null/undefined chunks
if (!nodeUpdate || typeof nodeUpdate !== 'object') return null;
// Special nodes first (backward compatibility)
if (nodeUpdate.delete_messages) {
return processDeleteMessages(nodeUpdate.delete_messages);
}
if (nodeUpdate.compact_messages) {
return processCompactMessages(nodeUpdate.compact_messages);
if (nodeUpdate.delete_messages || nodeUpdate.compact_messages) {
return null;
}
// Process operations emits workflow updates
if (nodeUpdate.process_operations) {
return processOperationsUpdate(nodeUpdate.process_operations);
}
@@ -0,0 +1,297 @@
import { AIMessage, HumanMessage, RemoveMessage } from '@langchain/core/messages';
import { cleanupDanglingToolCallMessages } from '../cleanup-dangling-tool-call-messages';
import {
determineStateAction,
handleCleanupDangling,
handleDeleteMessages,
} from '../state-modifier';
import { estimateTokenCountFromMessages } from '../token-usage';
jest.mock('../cleanup-dangling-tool-call-messages');
jest.mock('../token-usage');
const mockCleanupDanglingToolCallMessages = cleanupDanglingToolCallMessages as jest.MockedFunction<
typeof cleanupDanglingToolCallMessages
>;
const mockEstimateTokenCountFromMessages = estimateTokenCountFromMessages as jest.MockedFunction<
typeof estimateTokenCountFromMessages
>;
describe('state-modifier', () => {
beforeEach(() => {
jest.clearAllMocks();
mockCleanupDanglingToolCallMessages.mockReturnValue([]);
mockEstimateTokenCountFromMessages.mockReturnValue(100);
});
describe('determineStateAction', () => {
const emptyWorkflow = { nodes: [], connections: {}, name: '' };
const defaultNameWorkflow = { nodes: [], connections: {}, name: 'My workflow' };
const defaultNameNumberedWorkflow = { nodes: [], connections: {}, name: 'My workflow 5' };
const customNameWorkflow = { nodes: [], connections: {}, name: 'Email automation' };
const workflowWithNodes = {
nodes: [
{
id: '1',
name: 'Start',
type: 'n8n-nodes-base.start',
position: [0, 0] as [number, number],
typeVersion: 1,
parameters: {},
},
],
connections: {},
name: 'My workflow',
};
it('should return cleanup_dangling when dangling tool calls exist', () => {
mockCleanupDanglingToolCallMessages.mockReturnValue([
new RemoveMessage({ id: 'dangling-1' }),
]);
const result = determineStateAction(
{
messages: [new HumanMessage({ id: 'h1', content: 'Hello' })],
workflowJSON: emptyWorkflow,
},
40000,
);
expect(result).toBe('cleanup_dangling');
});
it('should return compact_messages for /compact command', () => {
const result = determineStateAction(
{
messages: [new HumanMessage({ id: 'h1', content: '/compact' })],
workflowJSON: emptyWorkflow,
},
40000,
);
expect(result).toBe('compact_messages');
});
it('should return delete_messages for /clear command', () => {
const result = determineStateAction(
{
messages: [new HumanMessage({ id: 'h1', content: '/clear' })],
workflowJSON: emptyWorkflow,
},
40000,
);
expect(result).toBe('delete_messages');
});
it('should return create_workflow_name for first message with default workflow name', () => {
const result = determineStateAction(
{
messages: [new HumanMessage({ id: 'h1', content: 'Create an email workflow' })],
workflowJSON: defaultNameWorkflow,
},
40000,
);
expect(result).toBe('create_workflow_name');
});
it('should return create_workflow_name for first message with numbered default name', () => {
const result = determineStateAction(
{
messages: [new HumanMessage({ id: 'h1', content: 'Create an email workflow' })],
workflowJSON: defaultNameNumberedWorkflow,
},
40000,
);
expect(result).toBe('create_workflow_name');
});
it('should return create_workflow_name for first message with empty workflow name', () => {
const result = determineStateAction(
{
messages: [new HumanMessage({ id: 'h1', content: 'Create an email workflow' })],
workflowJSON: emptyWorkflow,
},
40000,
);
expect(result).toBe('create_workflow_name');
});
it('should NOT return create_workflow_name for custom workflow name', () => {
const result = determineStateAction(
{
messages: [new HumanMessage({ id: 'h1', content: 'Create an email workflow' })],
workflowJSON: customNameWorkflow,
},
40000,
);
expect(result).toBe('continue');
});
it('should NOT return create_workflow_name when workflow has nodes', () => {
const result = determineStateAction(
{
messages: [new HumanMessage({ id: 'h1', content: 'Add another node' })],
workflowJSON: workflowWithNodes,
},
40000,
);
expect(result).toBe('continue');
});
it('should NOT return create_workflow_name when there are multiple messages', () => {
const result = determineStateAction(
{
messages: [
new HumanMessage({ id: 'h1', content: 'First message' }),
new AIMessage({ id: 'a1', content: 'Response' }),
new HumanMessage({ id: 'h2', content: 'Second message' }),
],
workflowJSON: defaultNameWorkflow,
},
40000,
);
expect(result).toBe('continue');
});
it('should return auto_compact_messages when tokens exceed threshold', () => {
mockEstimateTokenCountFromMessages.mockReturnValue(50000);
const result = determineStateAction(
{
messages: [new HumanMessage({ id: 'h1', content: 'A very long conversation' })],
workflowJSON: customNameWorkflow,
},
40000,
);
expect(result).toBe('auto_compact_messages');
});
it('should return continue as default when no conditions match', () => {
const result = determineStateAction(
{
messages: [
new HumanMessage({ id: 'h1', content: 'Hello' }),
new AIMessage({ id: 'a1', content: 'Hi!' }),
],
workflowJSON: customNameWorkflow,
},
40000,
);
expect(result).toBe('continue');
});
it('should return continue when there are no human messages', () => {
const result = determineStateAction(
{
messages: [new AIMessage({ id: 'a1', content: 'AI response' })],
workflowJSON: emptyWorkflow,
},
40000,
);
expect(result).toBe('continue');
});
});
describe('handleCleanupDangling', () => {
it('should return RemoveMessage array for dangling messages', () => {
const danglingRemoveMessages = [
new RemoveMessage({ id: 'ai-1' }),
new RemoveMessage({ id: 'ai-2' }),
];
mockCleanupDanglingToolCallMessages.mockReturnValue(danglingRemoveMessages);
const messages = [
new AIMessage({
id: 'ai-1',
content: 'Call',
tool_calls: [{ id: 'tc1', name: 'tool', args: {} }],
}),
new AIMessage({
id: 'ai-2',
content: 'Call',
tool_calls: [{ id: 'tc2', name: 'tool', args: {} }],
}),
];
const result = handleCleanupDangling(messages);
expect(result.messages).toHaveLength(2);
expect(result.messages[0]).toBeInstanceOf(RemoveMessage);
expect(result.messages[1]).toBeInstanceOf(RemoveMessage);
});
it('should return empty array when no dangling messages', () => {
mockCleanupDanglingToolCallMessages.mockReturnValue([]);
const result = handleCleanupDangling([new HumanMessage({ id: 'h1', content: 'Hello' })]);
expect(result.messages).toHaveLength(0);
});
});
describe('handleDeleteMessages', () => {
it('should return RemoveMessage for each input message', () => {
const messages = [
new HumanMessage({ id: 'h1', content: 'Hello' }),
new AIMessage({ id: 'a1', content: 'Hi' }),
new HumanMessage({ id: 'h2', content: 'Bye' }),
];
const result = handleDeleteMessages(messages);
expect(result.messages).toHaveLength(3);
expect(result.messages[0]).toBeInstanceOf(RemoveMessage);
expect(result.messages[0].id).toBe('h1');
expect(result.messages[1].id).toBe('a1');
expect(result.messages[2].id).toBe('h2');
});
it('should reset workflowJSON to empty state', () => {
const result = handleDeleteMessages([new HumanMessage({ id: 'h1', content: 'Hello' })]);
expect(result.workflowJSON).toEqual({
nodes: [],
connections: {},
name: '',
});
});
it('should clear previousSummary', () => {
const result = handleDeleteMessages([]);
expect(result.previousSummary).toBe('');
});
it('should set discoveryContext to null', () => {
const result = handleDeleteMessages([]);
expect(result.discoveryContext).toBeNull();
});
it('should add coordination log entry for clear action', () => {
const result = handleDeleteMessages([]);
expect(result.coordinationLog).toHaveLength(1);
expect(result.coordinationLog[0].phase).toBe('state_management');
expect(result.coordinationLog[0].status).toBe('completed');
expect(result.coordinationLog[0].summary).toBe('Cleared session and reset workflow');
});
it('should return empty workflowOperations array', () => {
const result = handleDeleteMessages([]);
expect(result.workflowOperations).toEqual([]);
});
});
});
@@ -57,7 +57,7 @@ describe('stream-processor', () => {
expect(message.text).toBe('Part 1\nPart 2');
});
it('should handle delete_messages with refresh message', () => {
it('should skip delete_messages (responder handles user message)', () => {
const chunk = {
delete_messages: {
messages: [{ content: 'Some deleted message' }],
@@ -66,13 +66,10 @@ describe('stream-processor', () => {
const result = processStreamChunk('updates', chunk);
expect(result).toBeDefined();
expect(result?.messages).toHaveLength(1);
const message = result?.messages[0] as AgentMessageChunk;
expect(message.text).toBe('Deleted, refresh?');
expect(result).toBeNull();
});
it('should handle compact_messages returning last message', () => {
it('should skip compact_messages (responder handles user message)', () => {
const chunk = {
compact_messages: {
messages: [
@@ -85,10 +82,7 @@ describe('stream-processor', () => {
const result = processStreamChunk('updates', chunk);
expect(result).toBeDefined();
expect(result?.messages).toHaveLength(1);
const message = result?.messages[0] as AgentMessageChunk;
expect(message.text).toBe('Last message to display');
expect(result).toBeNull();
});
it('should handle compact_messages with empty content', () => {
@@ -294,10 +288,9 @@ describe('stream-processor', () => {
results.push(output);
}
expect(results).toHaveLength(3);
expect(results).toHaveLength(2);
expect((results[0].messages[0] as AgentMessageChunk).text).toBe('Message 1');
expect((results[1].messages[0] as ToolProgressChunk).toolName).toBe('test_tool');
expect((results[2].messages[0] as AgentMessageChunk).text).toBe('Deleted, refresh?');
});
it('should handle empty stream', async () => {
@@ -86,32 +86,28 @@ function checkMergeNodeConnections(
const issues: SingleEvaluatorResult['violations'] = [];
if (/\.merge$/.test(nodeInfo.node.type)) {
const providedInputTypes = getProvidedInputTypes(nodeConnections);
// Merge node's number of inputs is controlled by the numberInputs parameter (default 2)
// The node type definition has static inputs, so we must read from parameters directly
const numberInputsParam = nodeInfo.node.parameters?.numberInputs;
const expectedInputs = typeof numberInputsParam === 'number' ? numberInputsParam : 2;
const totalInputConnections = providedInputTypes.get('main') ?? 0;
const mainConnections = nodeConnections?.main ?? [];
if (totalInputConnections < 2) {
// Count actual input slots that have connections (not total connections)
const connectedSlots = mainConnections.filter(
(slot) => Array.isArray(slot) && slot.length > 0,
).length;
if (connectedSlots < 2) {
issues.push({
name: 'node-merge-single-input',
type: 'major',
description: `Merge node ${nodeInfo.node.name} has only ${totalInputConnections} input connection(s). Merge nodes require at least 2 inputs to function properly.`,
description: `Merge node ${nodeInfo.node.name} has only ${connectedSlots} input connection(s). Merge nodes require at least 2 inputs to function properly.`,
pointsDeducted: 20,
});
}
const expectedInputs =
nodeInfo.resolvedInputs?.filter((input) => input.type === 'main').length ?? 1;
if (totalInputConnections !== expectedInputs) {
issues.push({
name: 'node-merge-incorrect-num-inputs',
type: 'minor',
description: `Merge node ${nodeInfo.node.name} has ${totalInputConnections} input connections but is configured to accept ${expectedInputs}.`,
pointsDeducted: 10,
});
}
const mainConnections = nodeConnections?.main ?? [];
// Check if all expected input slots have connections
const missingIndexes: number[] = [];
for (let inputIndex = 0; inputIndex < expectedInputs; inputIndex++) {
@@ -141,12 +141,8 @@ export interface WorkflowBuilderAgentConfig {
autoCompactThresholdTokens?: number;
instanceUrl?: string;
onGenerationSuccess?: () => Promise<void>;
/**
* Enable multi-agent supervisor architecture (experimental)
* When true, uses specialized agents (Discovery, Builder, Configurator) with a Supervisor
* When false, uses the legacy single-agent architecture
*/
enableMultiAgent?: boolean;
/** Metadata to include in LangSmith traces */
runMetadata?: Record<string, unknown>;
}
export interface ExpressionValue {
@@ -157,6 +153,7 @@ export interface ExpressionValue {
export interface BuilderFeatureFlags {
templateExamples?: boolean;
multiAgent?: boolean;
}
export interface ChatPayload {
@@ -180,7 +177,7 @@ export class WorkflowBuilderAgent {
private autoCompactThresholdTokens: number;
private instanceUrl?: string;
private onGenerationSuccess?: () => Promise<void>;
private enableMultiAgent: boolean;
private runMetadata?: Record<string, unknown>;
constructor(config: WorkflowBuilderAgentConfig) {
this.parsedNodeTypes = config.parsedNodeTypes;
@@ -193,7 +190,7 @@ export class WorkflowBuilderAgent {
config.autoCompactThresholdTokens ?? DEFAULT_AUTO_COMPACT_THRESHOLD_TOKENS;
this.instanceUrl = config.instanceUrl;
this.onGenerationSuccess = config.onGenerationSuccess;
this.enableMultiAgent = config.enableMultiAgent ?? false;
this.runMetadata = config.runMetadata;
}
private getBuilderTools(featureFlags?: BuilderFeatureFlags): BuilderTool[] {
@@ -440,9 +437,12 @@ export class WorkflowBuilderAgent {
/**
* Create the workflow graph based on configuration
* Controlled by feature flag only
*/
private createWorkflow(featureFlags?: BuilderFeatureFlags) {
if (this.enableMultiAgent) {
const useMultiAgent = featureFlags?.multiAgent ?? false;
if (useMultiAgent) {
this.logger?.debug('Using multi-agent supervisor architecture');
return this.createMultiAgentGraph(featureFlags);
}
@@ -515,8 +515,9 @@ export class WorkflowBuilderAgent {
recursionLimit: 50,
signal: abortSignal,
callbacks: this.tracer ? [this.tracer] : undefined,
metadata: this.runMetadata,
// Enable subgraph streaming when using multi-agent architecture
subgraphs: this.enableMultiAgent,
subgraphs: payload.featureFlags?.multiAgent ?? false,
};
return { agent, threadConfig, streamConfig };
@@ -60,6 +60,7 @@ export class AiBuilderChatRequestDto extends Z.class({
featureFlags: z
.object({
templateExamples: z.boolean().optional(),
multiAgent: z.boolean().optional(),
})
.optional(),
}),
@@ -121,6 +121,7 @@ describe('WorkflowBuilderService', () => {
mockLogger,
'test-instance-id', // instanceId
'https://instance.test.com', // instanceUrl
expect.any(String), // n8nVersion
expect.any(Function), // onCreditsUpdated callback
expect.any(Function), // onTelemetryEvent callback
);
@@ -160,6 +161,7 @@ describe('WorkflowBuilderService', () => {
mockLogger,
'test-instance-id', // instanceId
'https://instance.test.com', // instanceUrl
expect.any(String), // n8nVersion
expect.any(Function), // onCreditsUpdated callback
expect.any(Function), // onTelemetryEvent callback
);
@@ -279,7 +281,7 @@ describe('WorkflowBuilderService', () => {
MockedAiWorkflowBuilderService.mockImplementation(((...args: any[]) => {
// eslint-disable-next-line @typescript-eslint/no-unsafe-assignment
const callback = args[5]; // onCreditsUpdated is the 6th parameter
const callback = args[6]; // onCreditsUpdated is the 7th parameter (after n8nVersion)
// eslint-disable-next-line @typescript-eslint/no-unsafe-assignment
capturedCallback = callback;
return mockAiService;
@@ -327,7 +329,7 @@ describe('WorkflowBuilderService', () => {
MockedAiWorkflowBuilderService.mockImplementation(((...args: any[]) => {
// eslint-disable-next-line @typescript-eslint/no-unsafe-assignment
const callback = args[5]; // onCreditsUpdated is the 6th parameter
const callback = args[6]; // onCreditsUpdated is the 7th parameter (after n8nVersion)
// eslint-disable-next-line @typescript-eslint/no-unsafe-assignment
capturedCallback = callback;
return mockAiService;
@@ -387,7 +389,7 @@ describe('WorkflowBuilderService', () => {
MockedAiWorkflowBuilderService.mockImplementation(((...args: any[]) => {
// eslint-disable-next-line @typescript-eslint/no-unsafe-assignment
const telemetryCallback = args[6]; // onTelemetryEvent is the 7th parameter
const telemetryCallback = args[7]; // onTelemetryEvent is the 8th parameter (after n8nVersion)
// eslint-disable-next-line @typescript-eslint/no-unsafe-assignment
capturedTelemetryCallback = telemetryCallback;
return mockAiService;
@@ -433,7 +435,7 @@ describe('WorkflowBuilderService', () => {
MockedAiWorkflowBuilderService.mockImplementation(((...args: any[]) => {
// eslint-disable-next-line @typescript-eslint/no-unsafe-assignment
const telemetryCallback = args[6]; // onTelemetryEvent is the 7th parameter
const telemetryCallback = args[7]; // onTelemetryEvent is the 8th parameter (after n8nVersion)
// eslint-disable-next-line @typescript-eslint/no-unsafe-assignment
capturedTelemetryCallback = telemetryCallback;
return mockAiService;
@@ -477,7 +479,7 @@ describe('WorkflowBuilderService', () => {
MockedAiWorkflowBuilderService.mockImplementation(((...args: any[]) => {
// eslint-disable-next-line @typescript-eslint/no-unsafe-assignment
const telemetryCallback = args[6]; // onTelemetryEvent is the 7th parameter
const telemetryCallback = args[7]; // onTelemetryEvent is the 8th parameter (after n8nVersion)
// eslint-disable-next-line @typescript-eslint/no-unsafe-assignment
capturedTelemetryCallback = telemetryCallback;
return mockAiService;
@@ -79,6 +79,7 @@ export class WorkflowBuilderService {
this.logger,
this.instanceSettings.instanceId,
this.urlService.getInstanceBaseUrl(),
N8N_VERSION,
onCreditsUpdated,
onTelemetryEvent,
);
@@ -99,6 +99,12 @@ export const AI_BUILDER_TEMPLATE_EXAMPLES_EXPERIMENT = {
variant: 'variant',
};
export const AI_BUILDER_MULTI_AGENT_EXPERIMENT = {
name: '057_ai_builder_multi_agent',
control: 'control',
variant: 'variant',
};
export const EXPERIMENTS_TO_TRACK = [
EXTRA_TEMPLATE_LINKS_EXPERIMENT.name,
TEMPLATE_ONBOARDING_EXPERIMENT.name,
@@ -109,6 +115,7 @@ export const EXPERIMENTS_TO_TRACK = [
TEMPLATES_DATA_QUALITY_EXPERIMENT.name,
READY_TO_RUN_V2_PART2_EXPERIMENT.name,
AI_BUILDER_TEMPLATE_EXAMPLES_EXPERIMENT.name,
AI_BUILDER_MULTI_AGENT_EXPERIMENT.name,
TIME_SAVED_NODE_EXPERIMENT.name,
TEMPLATE_SETUP_EXPERIENCE.name,
];
@@ -95,6 +95,7 @@ export namespace ChatRequest {
export interface BuilderFeatureFlags {
templateExamples?: boolean;
multiAgent?: boolean;
}
export interface UserChatMessage {
@@ -1,7 +1,10 @@
import type { ChatRequest } from '@/features/ai/assistant/assistant.types';
import { useAIAssistantHelpers } from '@/features/ai/assistant/composables/useAIAssistantHelpers';
import { usePostHog } from '@/app/stores/posthog.store';
import { AI_BUILDER_TEMPLATE_EXAMPLES_EXPERIMENT } from '@/app/constants/experiments';
import {
AI_BUILDER_MULTI_AGENT_EXPERIMENT,
AI_BUILDER_TEMPLATE_EXAMPLES_EXPERIMENT,
} from '@/app/constants/experiments';
import type { IRunExecutionData } from 'n8n-workflow';
import type { IWorkflowDb } from '@/Interface';
@@ -56,6 +59,9 @@ export function createBuilderPayload(
templateExamples:
posthogStore.getVariant(AI_BUILDER_TEMPLATE_EXAMPLES_EXPERIMENT.name) ===
AI_BUILDER_TEMPLATE_EXAMPLES_EXPERIMENT.variant,
multiAgent:
posthogStore.getVariant(AI_BUILDER_MULTI_AGENT_EXPERIMENT.name) ===
AI_BUILDER_MULTI_AGENT_EXPERIMENT.variant,
};
return {