n8n 是一个 Fair-code 工作流自动化平台,原生支持 AI Agent 能力。通过可视化画布结合自定义代码,连接 1500+ 集成,支持自托管或云端部署。本文从源码层面深入分析其架构设计。
1. 版本化节点系统(VersionedNodeType)
n8n 采用版本化节点架构,所有节点(包括 AI Agent)都通过 VersionedNodeType 基类管理多版本。这种设计允许节点在不破坏现有工作流的情况下持续演进。
源码位置: packages/@n8n/nodes-langchain/nodes/agents/Agent/Agent.node.ts
export class Agent extends VersionedNodeType {
constructor() {
const baseDescription: INodeTypeBaseDescription = {
displayName: 'AI Agent',
name: 'agent',
icon: 'node:ai-agent',
iconColor: 'black',
group: ['transform'],
description: 'Generates an action plan and executes it. Can use external tools.',
codex: {
alias: ['LangChain', 'Chat', 'Conversational', 'Plan and Execute', 'ReAct', 'Tools'],
categories: ['AI'],
subcategories: {
AI: ['Agents', 'Root Nodes'],
},
},
defaultVersion: 3.1,
};
const nodeVersions: IVersionedNodeType['nodeVersions'] = {
1: new AgentV1(baseDescription),
1.1: new AgentV1(baseDescription),
// ... 1.2-1.9 均为 AgentV1
2: new AgentV2(baseDescription),
2.1: new AgentV2(baseDescription),
// ... 2.2-2.3 均为 AgentV2
3: new AgentV3(baseDescription),
3.1: new AgentV3(baseDescription),
};
super(nodeVersions, baseDescription);
}
}
架构洞察: 版本映射表将语义版本号映射到具体实现类。V1 系列(1.0-1.9)共享同一实现,V2 和 V3 各有独立实现。defaultVersion: 3.1 表示新建工作流默认使用最新版本。这种设计使得旧工作流可以继续运行,同时新功能只在新版本中提供。
2. 核心执行引擎(WorkflowExecute)
执行引擎是 n8n 的心脏,负责工作流的调度和执行。它使用 PCancelable 实现可取消的异步执行,支持部分执行(Partial Execution)优化。
源码位置: packages/core/src/execution-engine/workflow-execute.ts
export class WorkflowExecute {
private status: ExecutionStatus = 'new';
private readonly abortController = new AbortController();
timedOut: boolean = false;
constructor(
private readonly additionalData: IWorkflowExecuteAdditionalData,
private readonly mode: WorkflowExecuteMode,
private runExecutionData: IRunExecutionData = createRunExecutionData(),
private readonly storedAt: ExecutionStorageLocation = 'db',
) {}
// IMPORTANT: Do not add "async" to this function, it will then convert the
// PCancelable to a regular Promise and does so not allow canceling
// active executions anymore
run({ workflow, startNode, destinationNode, pinData, triggerToStartFrom,
additionalRunFilterNodes }: RunWorkflowOptions): PCancelable {
this.status = 'running';
// Get the nodes to start workflow execution from
startNode = startNode || workflow.getStartNode(destinationNode?.nodeName);
// ...
return this.processRunExecutionData(workflow);
}
}
关键设计:
- PCancelable: 故意不使用
async关键字,以保持PCancelable的可取消特性 - AbortController: 支持执行超时和手动取消
- 状态机:
new → running → ...的状态转换 - 部分执行:
runPartialWorkflow2方法通过有向图分析,只执行必要的节点子图
3. 部分执行与有向图分析
n8n 的部分执行系统是最精妙的设计之一。当用户手动执行某个节点时,引擎会构建有向图,找到触发器到目标节点的最小子图,避免执行整个工作流。
源码位置: packages/core/src/execution-engine/workflow-execute.ts
runPartialWorkflow2(workflow, runData, pinData, dirtyNodeNames, destinationNode,
agentRequest?): PCancelable {
// 1. 构建有向图
let graph = DirectedGraph.fromWorkflow(workflow);
// 2. 工具节点特殊处理:重连图结构
if (NodeHelpers.isTool(destinationNodeType.description, destination.parameters)) {
graph = rewireGraph(destination, graph, agentRequest);
workflow = graph.toWorkflow({ ...workflow });
}
// 3. 查找触发器
let trigger = findTriggerForPartialExecution(workflow, destinationNode.nodeName, runData);
// 4. 查找子图
const filteredGraph = filterDisabledNodes(graph);
graph = findSubgraph({ graph: filteredGraph, destination, trigger });
// 5. 查找起始节点
let startNodes = findStartNodes({ graph, trigger, destination, runData, pinData });
// 6. 检测并处理循环
startNodes = handleCycles(graph, startNodes, trigger);
// 7. 重建执行栈
const { nodeExecutionStack, waitingExecution, waitingExecutionSource } =
recreateNodeExecutionStack(graph, startNodes, runData, pinData ?? {});
// 8. 执行
return this.processRunExecutionData(workflow);
}
算法流程: DirectedGraph.fromWorkflow → filterDisabledNodes → findSubgraph → findStartNodes → handleCycles → recreateNodeExecutionStack。这是一套完整的图论算法链,确保最小化执行范围。
4. LangChain Agent 集成架构
n8n 通过 LangChain 集成实现了 AI Agent 能力。Agent 节点作为根节点,通过子节点连接(Subnodes)获取模型、工具和记忆。
源码位置: packages/@n8n/nodes-langchain/nodes/agents/Agent/agents/ToolsAgent/common.ts
export async function getChatModel(
ctx: IExecuteFunctions | ISupplyDataFunctions | IWebhookFunctions,
index: number = 0,
): Promise {
const connectedModels = await ctx.getInputConnectionData(
NodeConnectionTypes.AiLanguageModel, 0
);
let model;
if (Array.isArray(connectedModels) && index !== undefined) {
if (connectedModels.length <= index) return undefined;
// We get the models in reversed order from the workflow
const reversedModels = [...connectedModels].reverse();
model = reversedModels[index] as BaseChatModel;
} else {
model = connectedModels as BaseChatModel;
}
if (!isChatInstance(model) || !model.bindTools) {
throw new NodeOperationError(ctx.getNode(),
'Tools Agent requires Chat Model which supports Tools calling');
}
return model;
}
连接模型: NodeConnectionTypes.AiLanguageModel 定义了 AI 模型连接类型。模型通过 getInputConnectionData 从子节点获取,支持主模型和备用模型(Fallback Model)。
5. Agent 序列构建与执行
Agent 的核心执行逻辑封装在 createAgentSequence 中,它将 LangChain 的 createToolCallingAgent 与输出解析器、备用模型组合成可执行的序列。
源码位置: packages/@n8n/nodes-langchain/nodes/agents/Agent/agents/ToolsAgent/V3/helpers/createAgentSequence.ts
export function createAgentSequence(
model: BaseChatModel,
tools: Array,
prompt: ChatPromptTemplate,
_options: { maxIterations?: number; returnIntermediateSteps?: boolean },
outputParser?: N8nOutputParser,
memory?: BaseChatMemory,
fallbackModel?: BaseChatModel | null,
forceToolCall = false,
) {
const allTools = getAllTools(model, tools);
const agent = createToolCallingAgent({
llm: forceToolCall && model.bindTools
? model.bindTools(allTools, { tool_choice: 'any' })
: model,
tools: allTools,
prompt,
streamRunnable: false,
});
let fallbackAgent: AgentRunnableSequence | undefined;
if (fallbackModel) {
fallbackAgent = createToolCallingAgent({
llm: forceToolCall && fallbackModel.bindTools
? fallbackModel.bindTools(fallbackTools, { tool_choice: 'any' })
: fallbackModel,
tools: fallbackTools,
prompt,
streamRunnable: false,
});
}
const runnableAgent = RunnableSequence.from([
fallbackAgent ? agent.withFallbacks([fallbackAgent]) : agent,
getAgentStepsParser(outputParser, memory),
fixEmptyContentMessage,
]) as AgentRunnableSequence;
runnableAgent.singleAction = true;
runnableAgent.streamRunnable = false;
return runnableAgent;
}
设计要点:
forceToolCall: 强制模型调用工具(tool_choice: 'any'),用于首次迭代withFallbacks: LangChain 原生的备用模型机制,主模型失败时自动切换singleAction = true: 单次动作模式,每次只调用一个工具- 输出链:
Agent → StepsParser → FixEmptyContent
6. 批处理与并发执行
V3 版本引入了批处理机制,支持配置批量大小和批次间延迟,通过 Promise.allSettled 实现并发执行。
源码位置: packages/@n8n/nodes-langchain/nodes/agents/Agent/agents/ToolsAgent/V3/helpers/executeBatch.ts
export async function executeBatch(
ctx: IExecuteFunctions | ISupplyDataFunctions,
batch: INodeExecutionData[],
startIndex: number,
model: BaseChatModel,
fallbackModel: BaseChatModel | null,
memory: BaseChatMemory | undefined,
response?: EngineResponse,
): Promise {
// 处理 HITL(Human-in-the-Loop)工具响应
const hitlResult = processHitlResponses(response, startIndex);
if (hitlResult.hasApprovedHitlTools && hitlResult.pendingGatedToolRequest) {
return { returnData: [], request: hitlResult.pendingGatedToolRequest, memoryHits };
}
// 检查最大迭代次数
const maxIterations = ctx.getNodeParameter('options.maxIterations', 0, 10);
// 并发处理批次中的所有项目
const batchPromises = batch.map(async (_item, batchItemIndex) => {
const itemIndex = startIndex + batchItemIndex;
checkMaxIterations(response, maxIterations, ctx.getNode());
const itemContext = await prepareItemContext(ctx, itemIndex, processedResponse, model);
const executor = createAgentSequence(model, tools, prompt, options,
outputParser, memory, fallbackModel, ...);
return await runAgent(ctx, executor, itemContext, model, memory,
processedResponse, memoryHits);
});
const batchResults = await Promise.allSettled(batchPromises);
// ... 处理结果
}
并发策略: 使用 Promise.allSettled 而非 Promise.all,确保单个项目的失败不会影响整个批次。continueOnFail() 机制允许在错误时返回错误信息而非抛出异常。
7. 流式执行与事件流处理
Agent 支持流式输出,通过 LangChain 的 streamEvents API 实现实时响应。
源码位置: packages/@n8n/nodes-langchain/nodes/agents/Agent/agents/ToolsAgent/V3/helpers/runAgent.ts
export async function runAgent(
ctx: IExecuteFunctions | ISupplyDataFunctions,
executor: AgentRunnableSequence,
itemContext: ItemContext,
model: BaseChatModel,
memory: BaseChatMemory | undefined,
response?: EngineResponse,
memoryHits?: { loads: number; saves: number },
): Promise {
const isStreamingAvailable = 'isStreaming' in ctx ? ctx.isStreaming?.() : undefined;
if ('isStreaming' in ctx && options.enableStreaming && isStreamingAvailable
&& ctx.getNode().typeVersion >= 2.1) {
// 流式执行路径
const chatHistory = await loadMemory(memory, model, options.maxTokensFromMemory);
const eventStream = executorWithTracing.streamEvents(
{ ...invokeParams, chat_history: chatHistory },
{ version: 'v2', ...executeOptions },
);
const result = await processEventStream(ctx, eventStream, itemIndex, finalizeOutput);
// 如果结果包含工具调用,构建请求对象
if (result.toolCalls && result.toolCalls.length > 0) {
const actions = createEngineRequests(result.toolCalls, itemIndex, tools);
return { actions, metadata: buildResponseMetadata(response, itemIndex) };
}
return result;
} else {
// 非流式执行路径
const modelResponse = await executorWithTracing.invoke(
{ ...invokeParams, chat_history: chatHistory }, executeOptions);
// ...
}
}
双路径设计: 流式和非流式执行共享相同的输入准备逻辑,但在执行阶段分叉。流式路径使用 streamEvents API,非流式使用 invoke。两种路径都支持工具调用的中断和恢复。
8. 子 Agent(Sub-Agent)与工具执行
n8n 支持多 Agent 协作,通过 AgentTool 节点实现 Agent 作为工具被其他 Agent 调用。子 Agent 的工具执行通过内联解析而非引擎调度。
源码位置: packages/@n8n/nodes-langchain/nodes/agents/Agent/agents/ToolsAgent/V3/helpers/resolveSubAgentRequest.ts
export async function resolveSubAgentRequest(
ctx: ISupplyDataFunctions,
request: EngineRequest,
deps: ResolveSubAgentRequestDeps,
): Promise {
let current: INodeExecutionData[][] | EngineRequest = request;
const node = ctx.getNode();
const tools = (await getTools(ctx)) as ConnectedTool[];
while (isEngineRequest(current)) {
assertNoHitlActions(node, current.actions);
ctx.getExecutionCancelSignal?.()?.throwIfAborted?.();
// 内联执行所有工具调用
const actionResponses = await Promise.all(
current.actions.map(async (action) =>
await executeEngineAction(node, action, tools)),
);
// 递归调用 Agent 执行
current = await deps.runAgentBatch({
actionResponses, metadata: current.metadata
});
}
return current;
}
function assertNoHitlActions(node: INode, actions: Action[]): void {
for (const action of actions) {
if (action.metadata?.hitl) {
throw new UserError(
`Human-in-the-Loop nodes cannot be used inside a sub-agent. Move "${action.nodeName}" to the top-level workflow.`);
}
}
}
关键限制: 子 Agent 不支持 HITL(Human-in-the-Loop)节点,因为审批流程需要顶层引擎的参与。这是一个合理的架构约束。
9. 路由节点与声明式 API 集成
n8n 的 RoutingNode 类实现了声明式 REST API 集成,支持分页、批处理和凭证管理。
源码位置: packages/core/src/execution-engine/routing-node.ts
export class RoutingNode {
constructor(
private readonly context: ExecuteContext,
private readonly nodeType: INodeType,
private readonly credentialsDecrypted?: ICredentialsDecrypted,
) {}
async runNode(): Promise {
const items = (inputData[NodeConnectionTypes.Main] ??
inputData[NodeConnectionTypes.AiTool])[0] as INodeExecutionData[];
const returnData: INodeExecutionData[] = [];
const { credentials, credentialDescription } = await this.prepareCredentials();
const { batching } = context.getNodeParameter('requestOptions', 0, {}) as {
batching: { batch: { batchSize: number; batchInterval: number } };
};
const batchSize = batching?.batch?.batchSize > 0 ? batching?.batch?.batchSize : 1;
const batchInterval = batching?.batch.batchInterval;
for (let itemIndex = 0; itemIndex 0 && batchSize >= 0 && batchInterval > 0) {
if (itemIndex % batchSize === 0) {
await sleep(batchInterval);
}
}
// ...
}
}
}
设计模式: 路由节点支持 Main 和 AiTool 两种输入连接类型,使得同一个节点既可以作为工作流节点使用,也可以作为 AI Agent 的工具使用。
10. 节点加载器与插件系统
n8n 的节点加载器支持从目录动态加载节点和凭证类型,支持自定义节点和第三方包。
源码位置: packages/core/src/nodes-loader/directory-loader.ts
export abstract class DirectoryLoader implements NodeLoader {
isLazyLoaded = false;
loadedNodes: INodeTypeNameVersion[] = [];
nodeTypes: INodeTypeData = {};
credentialTypes: ICredentialTypeData = {};
known: KnownNodesAndCredentials = { nodes: {}, credentials: {} };
types: Types = { nodes: [], credentials: [] };
nodesByCredential: Record = {};
constructor(
readonly directory: string,
protected excludeNodes: string[] = [],
protected includeNodes: string[] = [],
) {
// If `directory` is a symlink, we try to resolve it to its real path
try {
this.directory = realpathSync(directory);
} catch (error) {
if (error.code !== 'ENOENT') throw error;
}
}
}
插件架构: DirectoryLoader 是抽象基类,具体实现包括 PackageDirectoryLoader(npm 包)、LazyPackageDirectoryLoader(懒加载)和 CustomDirectoryLoader(自定义目录)。known 字段存储节点和凭证的位置信息,支持懒加载场景。
架构总结
n8n 的架构设计体现了以下核心理念:
- 版本化演进: 节点通过版本号管理,旧工作流永不被破坏
- 图论驱动执行: 有向图分析实现最小化部分执行
- LangChain 深度集成: 不是简单封装,而是深度利用 LangChain 的 Agent、Tool、Memory 原语
- 可取消异步: PCancelable + AbortController 实现优雅的执行控制
- 子 Agent 协作: 通过内联解析实现多 Agent 嵌套调用
- 声明式集成: RoutingNode 支持声明式 REST API 定义
- 插件化节点: DirectoryLoader 支持动态加载和懒加载
这种架构使得 n8n 既能作为通用工作流引擎运行,又能作为 AI Agent 平台提供智能自动化能力。
发表回复