AI 任务提交成功,只代表服务端接受了请求,不代表最终结果已经可靠送达业务系统。本文用 Node.js 构建一套可运行的 Webhook 示例,覆盖签名校验、重复投递、乱序到达、超时重试和重放防护。重点是区分任务状态、事件版本与消费幂等,让异步交付具备可恢复性。

一、先区分任务完成与结果送达

同步接口通常在一次 HTTP 请求中返回结果,但 LLM 长文本生成、批量文档解析、知识库构建和智能体工作流往往需要持续运行数秒甚至更久。让客户端一直占用连接等待结果,容易受到网关超时、网络抖动和服务重启影响。

更常见的做法是拆成两步:

  1. 客户端提交任务,服务端立即返回 task_id
  2. 任务执行完成后,服务端向业务系统发送 Webhook 事件。

这两个步骤对应两个不同的可靠性问题:

阶段需要保证的事情常见误区
任务提交请求被服务端接受并持久化把 HTTP 200 当成任务已完成
任务执行状态能够持续更新并可查询只保存在进程内存中
事件投递结果最终能够送到接收方只发送一次,不处理超时
事件消费重复和乱序不会破坏业务数据每次收到都直接覆盖状态

因此,提交接口可以返回 202 Accepted,响应中明确表示任务正在处理;真正的结果则通过 Webhook 发送,同时保留任务查询接口作为补偿路径。Webhook 不是唯一数据源,而是推动业务系统及时更新的通知机制。

二、设计一个可恢复的事件协议

一次 Webhook 投递至少应包含事件 ID、事件类型、创建时间、任务 ID、状态和版本号。例如:

{
  "id": "evt_8f6d",
  "type": "ai.task.completed",
  "created_at": "2025-01-01T12:00:00.000Z",
  "data": {
    "task_id": "task_123",
    "status": "succeeded",
    "version": 2,
    "result": {
      "text": "处理完成"
    }
  }
}

id 用于消费端幂等。接收方处理成功后,应把事件 ID 写入数据库的唯一索引表,或者写入具有明确过期策略的幂等存储。不能只依赖内存集合,因为进程重启后会失去去重记录。

version 用于解决乱序到达。任务状态不是简单的“最后收到什么就写什么”,而是只接受版本更高的事件。例如版本 2 的 succeeded 先到达,随后版本 1 的 running 才到达,消费端应该忽略版本 1,而不是把已完成任务改回运行中。

发送方需要把事件记录和任务状态放在可恢复的存储中。生产环境通常会采用事务写入任务表和事件表,再由后台投递器读取未发送或待重试事件。本文示例使用内存结构,是为了便于直接运行,不能直接作为生产存储方案。

三、Node.js 实现签名与投递重试

下面的单文件示例只使用 Node.js 内置模块。它同时提供任务提交接口、Webhook 接收接口和一个简单的投递器。任务在一秒后完成,投递器会向本机的 Webhook 地址发送事件;接收端会验证签名、检查事件版本并执行幂等消费。

import http from "node:http";
import { createHmac, randomUUID, timingSafeEqual } from "node:crypto";

const PORT = 3000;
const SECRET = "replace-this-secret";
const tasks = new Map();
const consumedEvents = new Set();
const taskViews = new Map();

function sign(rawBody, timestamp) {
  return "sha256=" + createHmac("sha256", SECRET)
    .update(`${timestamp}.${rawBody}`)
    .digest("hex");
}

function sendJson(res, status, data) {
  const body = JSON.stringify(data);
  res.writeHead(status, {
    "content-type": "application/json; charset=utf-8",
    "content-length": Buffer.byteLength(body)
  });
  res.end(body);
}

async function readBody(req) {
  const chunks = [];
  for await (const chunk of req) chunks.push(chunk);
  return Buffer.concat(chunks).toString("utf8");
}

function requestWebhook(event) {
  return new Promise((resolve, reject) => {
    const body = JSON.stringify(event);
    const timestamp = Math.floor(Date.now() / 1000).toString();
    const req = http.request("http://127.0.0.1:3000/webhook/ai", {
      method: "POST",
      headers: {
        "content-type": "application/json",
        "content-length": Buffer.byteLength(body),
        "x-webhook-id": event.id,
        "x-webhook-timestamp": timestamp,
        "x-webhook-signature": sign(body, timestamp)
      }
    }, (res) => {
      res.resume();
      res.on("end", () => resolve(res.statusCode));
    });
    req.on("error", reject);
    req.setTimeout(3000, () => req.destroy(new Error("webhook timeout")));
    req.end(body);
  });
}

async function deliver(event) {
  for (let attempt = 1; attempt <= 5; attempt += 1) {
    try {
      const status = await requestWebhook(event);
      if (status >= 200 && status < 300) return;
      throw new Error(`receiver returned ${status}`);
    } catch (error) {
      if (attempt === 5) {
        console.error("delivery failed", event.id, error.message);
        return;
      }
      const delay = Math.min(30000, 500 * 2 ** (attempt - 1));
      await new Promise((resolve) => setTimeout(resolve, delay));
    }
  }
}

function startTask() {
  const taskId = `task_${randomUUID()}`;
  tasks.set(taskId, { status: "pending", version: 1 });
  setTimeout(() => {
    const task = tasks.get(taskId);
    task.status = "succeeded";
    task.version = 2;
    const event = {
      id: `evt_${randomUUID()}`,
      type: "ai.task.completed",
      created_at: new Date().toISOString(),
      data: {
        task_id: taskId,
        status: task.status,
        version: task.version,
        result: { text: "处理完成" }
      }
    };
    deliver(event);
  }, 1000);
  return taskId;
}

function verifySignature(req, rawBody) {
  const timestamp = Number(req.headers["x-webhook-timestamp"]);
  const received = String(req.headers["x-webhook-signature"] || "");
  if (!Number.isInteger(timestamp)) return false;
  if (Math.abs(Date.now() / 1000 - timestamp) > 300) return false;
  const expected = Buffer.from(sign(rawBody, timestamp.toString()));
  const actual = Buffer.from(received);
  return expected.length === actual.length && timingSafeEqual(expected, actual);
}

const server = http.createServer(async (req, res) => {
  if (req.method === "POST" && req.url === "/tasks") {
    const taskId = startTask();
    return sendJson(res, 202, { task_id: taskId, status: "pending" });
  }

  if (req.method === "GET" && req.url.startsWith("/tasks/")) {
    const taskId = req.url.slice("/tasks/".length);
    const task = tasks.get(taskId);
    return task ? sendJson(res, 200, { task_id: taskId, ...task })
      : sendJson(res, 404, { error: "task not found" });
  }

  if (req.method === "POST" && req.url === "/webhook/ai") {
    const rawBody = await readBody(req);
    if (!verifySignature(req, rawBody)) {
      return sendJson(res, 401, { error: "invalid webhook signature" });
    }

    const event = JSON.parse(rawBody);
    if (consumedEvents.has(event.id)) {
      return sendJson(res, 200, { received: true, duplicate: true });
    }

    const current = taskViews.get(event.data.task_id);
    if (!current || event.data.version > current.version) {
      taskViews.set(event.data.task_id, {
        status: event.data.status,
        version: event.data.version,
        result: event.data.result
      });
    }
    consumedEvents.add(event.id);
    return sendJson(res, 200, { received: true });
  }

  sendJson(res, 404, { error: "not found" });
});

server.listen(PORT, () => {
  console.log(`server listening on http://127.0.0.1:${PORT}`);
});

保存为 webhook-demo.mjs 后运行:

node webhook-demo.mjs
curl -i -X POST http://127.0.0.1:3000/tasks

示例中的重试策略使用指数退避,最多尝试五次。真实系统还应加入随机抖动,避免接收方恢复时大量请求同时涌入。对于永久性错误,例如签名配置错误或请求格式错误,通常应记录到死信表,不应无限重试。

四、签名校验、重放与幂等消费

签名不能只覆盖请求体,还应覆盖时间戳。示例签名原文是 timestamp.rawBody,这样攻击者即使拿到一份合法请求,也不能轻易把它改成另一个时间再次提交。

接收端验证时有三个关键步骤:

  1. 读取原始请求体。签名校验必须使用收到的原始字节,不能先解析 JSON 再重新序列化,因为空格、字段顺序和转义形式都可能变化。
  2. 检查时间窗口。示例允许五分钟偏差,生产环境应根据网络延迟和重试策略设置,并确保服务器时间同步。
  3. 使用常量时间比较。timingSafeEqual 可以减少逐字节比较造成的时序信息泄露,同时要先检查两个 Buffer 长度是否一致,否则 Node.js 会抛出异常。

时间戳只能降低重放风险,不能替代幂等。一个合法请求在五分钟内仍可能被重复发送,因此必须持久化 event.id。推荐建立如下唯一约束:

数据约束或索引作用
event_id唯一索引阻止同一事件重复执行业务逻辑
task_idversion唯一索引或版本条件更新防止旧状态覆盖新状态
delivery_id普通索引追踪一次投递的多次尝试

消费过程还要考虑并发。如果两个相同事件同时到达,不能采用“先查询、后插入”的非原子流程,否则两个请求都可能认为事件尚未处理。更可靠的方式是先插入幂等记录,利用唯一键冲突判断重复;随后在同一事务中更新业务表。

如果业务处理包含外部副作用,例如发送邮件、扣减额度或创建下游订单,则仅记录事件 ID 还不够。应为副作用建立独立的操作 ID,或者使用事务消息、Outbox 和下游幂等键,确保重试不会重复执行不可逆操作。

五、状态机与故障处理边界

任务状态最好定义成有限状态机,而不是任意字符串。一个简单模型可以是 pending -> running -> succeeded,失败路径为 pending -> running -> failed,取消路径为 pending/running -> canceled。每次状态变更都递增版本号,并把版本号随事件发送。

状态事件和结果事件也可以分开。例如 ai.task.progress 只用于展示进度,ai.task.completed 才代表最终结果。消费端不应仅凭事件到达顺序判断最终状态,而应根据版本号和状态机规则校验转换是否合法。

Webhook 接收端应尽快返回成功,不要在 HTTP 请求中执行耗时的文档入库、向量化或二次模型调用。常见流程是:验签、校验基本字段、原子写入事件表,然后返回 2xx;真正的业务处理交给内部队列。这样可以缩短发送方的等待时间,也能把业务失败与网络投递失败分开处理。

返回状态码可以按下面的方式划分:

状态含义发送方行为
2xx已接收并持久化结束本次投递
400JSON 或字段非法记录错误,通常不重试
401/403签名或权限失败告警并人工处理配置
408/429超时或限流延迟后重试
500-599接收方暂时不可用指数退避重试

发送方还应记录每次投递的事件 ID、目标地址、尝试次数、响应状态、耗时和最后错误。接收方则需要提供事件查询、失败重放和任务状态查询能力。这样即使事件永久失败,也可以通过后台任务修复,而不是要求用户重新提交一次 AI 任务。

总结

Webhook 可靠交付的核心不是“发出一个 POST 请求”,而是建立一套可重试、可验证、可去重和可恢复的事件协议。

要点可以归纳为:

  • 202 Accepted 表示任务已接受,不把请求成功等同于结果完成。
  • 每个事件分配稳定的事件 ID,并在消费端使用持久化唯一约束实现幂等。
  • 使用任务版本号处理乱序到达,避免旧状态覆盖新状态。
  • 使用时间戳加 HMAC 签名校验来源,并用常量时间比较验证签名。
  • 对连接超时、限流和服务端错误进行指数退避,对永久错误进入死信或人工处理流程。
  • 接收端先快速验签和落库,再异步执行耗时业务;任务查询接口作为 Webhook 失败后的补偿路径。

对于长时间运行的 LLM、文档处理和智能体任务,这些机制比单纯增加请求超时时间更重要。它们把一次不可靠的网络通知,变成能够审计、重试和恢复的交付流程。