☰
05.03 · n8n 源码剖析:Webhook 接入层 Webhook Ingress
2026/10/11 4:36:36 网站建设 项目流程

本文是专栏「n8n 工作流引擎剖析」第 05 章(组件深度剖析)的第 03/11 篇,承接上一篇《触发器注册中心 Trigger Registry》。讲的是从一个原始 HTTP 请求到调用WorkflowRunner.run之间发生的一切。

你在这里:读完本文,你能从 Express 一路追踪请求直到第一个数据项,能解释各种响应模式,并能说清为什么 Webhook 节点会在执行记录存在之前就先运行。

缩写:HTTP(Hypertext Transfer Protocol,超文本传输协议)、CORS(Cross-Origin Resource Sharing,跨域资源共享)、ID(Identifier,标识符)、DB(Database,数据库)、JSON(JavaScript Object Notation,JavaScript 对象表示法)、V8(Google’s JavaScript engine,谷歌的 JavaScript 引擎)。


角色回顾

负责:HTTP 请求 →(找到对应 webhook、加载已发布工作流、生成首个数据项、给出 HTTP 响应)。
掌握:路由 → 数据行的映射(通过WebhookService)、响应模式、请求/响应对象。
不做:不遍历图;不决定跑在主进程还是工作进程;不持久化执行记录。
出现于:S3、S4、S5(以及 S6 末尾那个被延迟发出的回复)。

内部设计

存在 workflowData

没有 workflowData

Express:app.all('/webhook/*path')

createWebhookHandlerFor(liveWebhooks,'webhook')

WebhookRequestHandler.handleRequest
方法检查 · CORS
OPTIONS→204

LiveWebhooks.executeWebhook

findWebhook(path, method)
缓存 → DB 静态 → DB 动态

加载已发布版本
→ new Workflow(...)

sanitizeWebhookRequest
(除非节点在认证白名单中)

WebhookHelpers.executeWebhook

evaluateResponseOptions, parseRequestBody

invokeWebhook → WebhookService.runWebhook → node.webhook(ctx)

handleImmediateWebhookResponse

prepareExecutionData → WorkflowRunner.run(...)

直接响应,不创建执行

延迟的 'onReceived' 回复

图注:接入层的处理流水线——注意节点自己的webhook()是先执行的,它甚至可能决定根本不启动任何执行。

  1. 路由。AbstractServer把app.all('/<endpointWebhook>/*path', createWebhookHandlerFor(liveWebhooks, 'webhook'))挂载在请求体解析器之前,这样 Webhook 节点才能自己流式处理二进制请求体(cli/src/abstract-server.ts第 240–252 行)。表单(Form)以及等待中的 webhook/表单恢复接口有对应的兄弟路由;测试 webhook 由TestWebhooks在webhook-test前缀下提供服务。

  2. 处理器。WebhookRequestHandler.handleRequest拒绝不支持的方法,只在存在origin请求头时才应用 CORS,用204响应OPTIONS,然后把请求交给webhookManager.executeWebhook。错误会被转换成 HTTP 错误响应(未知的 webhook →WebhookNotFoundError)。

  3. 查找。WebhookService.findWebhook→ 静态缓存(webhook:${method}-${path})→ 静态数据库记录 → 动态路径探测(路径中带:param段的情况;路径参数会被复制进request.params)。

  4. 加载已发布的版本。loadWebhookExecutionData使用workflow.activeVersion(或者在某个开关后面的新发布服务)——nodes/connections来自已发布的版本,绝不是草稿。构建出一个Workflow;getBase(...)创建additionalData(凭据辅助对象、钩子占位、设置、生产调用下的userId= 发布者)。

  5. 表达式隔离实例——仅在需要时才创建。webhookPhaseNeedsIsolate会跳过创建 V8 隔离实例,前提是满足一个很常见的场景:Webhook 节点 v2 及以上版本、参数是静态的、描述字段能原生解析。只要有任何一点无法证明是静态的,就会去获取一个隔离实例。

  6. WebhookHelpers.executeWebhook(webhook-helpers.ts第 804 行):

    • 通过evaluateResponseOptions解析出responseMode等;不支持的模式 → HTTP 500。支持的模式有:onReceived、lastNode、responseNode、formPage、streaming、hostedChat;
    • parseRequestBody;
    • invokeWebhook→WebhookService.runWebhook(webhook.service.ts第 646 行)构建一个WebhookContext,调用节点的webhook()。这一步出错会被上报,并以一个通用错误作答;不会创建任何执行记录。
  7. 由节点来决定。Webhook 节点的webhook()(Webhook.node.ts第 225 行)会检查 IP 白名单、机器人过滤、认证方式(基础认证/请求头/JWT/n8n OAuth)、一个可选的"仅在满足条件时运行"表达式,然后构建出:

    {json:{headers:req.headers,params:req.params,query:req.query,body:req.body}}

    并返回{ webhookResponse, workflowData: [[item]] }。如果认证失败,它会自己写出403/401,并返回{ noWebhookResponse: true }(没有workflowData→ 不会创建执行)。

  8. 响应模式(决定调用方在等什么):

    模式调用方会收到……
    onReceived(默认)一旦执行记录存在,就立刻收到一个 JSON 回复:如果配置了responseData就用它,否则是{ "message": "Workflow was started" }
    lastNode运行结束时最后一个执行节点的输出
    responseNodeRespond to Webhook节点发送的任何内容
    streaming运行过程中的分块流(sendChunk钩子)
    formPage/hostedChatForm / Chat 触发器对应的界面页面流程
  9. 启动运行。prepareExecutionData构建出初始的IRunExecutionData——关键是executionData.nodeExecutionStack = [{ node: <起始节点>, data: { main: <workflowData> }, source: null }]——然后:

    executionId=awaitContainer.get(WorkflowRunner).run(runData,/*loadStaticData*/true,/*realtime*/!didSendResponse&&!shouldDeferOnReceivedResponse,existingExecution/* 只有在恢复一次 Wait 时才有值 */,responsePromise);

    对于onReceived模式,回复会延迟到run返回之后,这样responseData表达式里就能用上$execution.id(webhook-helpers.ts约第 1262–1290 行)。

交互关系

对象契约
触发器注册中心读取它写入的数据行;启动时会预先填充静态缓存(populateCache)
节点调用webhook(ctx);期望得到IWebhookResponseData(workflowData、webhookResponse、noWebhookResponse)
工作流运行器(下一篇)run(IWorkflowExecutionDataProcess, …)→ 返回一个执行 ID
队列模式对于responseNode/lastNode,工作进程会把响应转发回这个进程(伸缩队列是本系列后续文章)

⚓ 回到示例 —— S3 → S5

POST /webhook/orders,内容为{"customer":"ACME","amount":100}(Content-Type: application/json):

  1. createWebhookHandlerFor拼出params.path=orders。调试日志:Received webhook "POST" for path "orders"。

  2. findWebhook('POST','orders')→ 缓存未命中 → 数据库命中(orders, POST)→ 写入缓存。

  3. 加载已发布版本;webhookPhaseNeedsIsolate返回假(Webhook v2、参数是静态的)→ 不构建隔离实例。

  4. runWebhook调用Webhook.webhook(ctx);没有配置认证,validateAuth通过;结果:

    {"workflowData":[[{"json":{"headers":{"content-type":"application/json","…":"…"},"params":{},"query":{},"body":{"customer":"ACME","amount":100}}}]]}
  5. prepareExecutionData把这个数据项放进初始栈里;WorkflowRunner.run被调用(S6 开始)。它返回例如"1042";直到这时接入层才发出:

    HTTP/1.1 200 OK {"message":"Workflow was started"}

失败行为

失败情形结果
找不到(path, method)对应的记录抛出WebhookNotFoundError→ 404(错误信息里会列出这个路径已注册的其他方法)
webhook()中认证失败节点自己写出 401/403,返回noWebhookResponse;不创建执行
webhook()抛出异常上报给错误报告器;返回一个通用错误响应;不创建执行
不支持的响应模式500,提示The response mode '…' is not valid!
调用方断开连接对onReceived模式没有影响——执行记录已经创建好了
WorkflowRunner.run抛出异常(例如执行前被阻断)异常会一路传到 HTTP 响应;PreExecuteBlockedError会在运行器里被解包

下一篇:《工作流运行器 Workflow Runner》,讲清楚一次执行是怎么被创建、又被决定跑在哪里的。


📚 返回专栏目录

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询