AI 任务提交成功,只代表服务端接受了请求,不代表最终结果已经可靠送达业务系统。本文用 Node.js 构建一套可运行的 Webhook 示例,覆盖签名校验、重复投递、乱序到达、超时重试和重放防护。重点是区分任务状态、事件版本与消费幂等,让异步交付具备可恢复性。
一、先区分任务完成与结果送达
同步接口通常在一次 HTTP 请求中返回结果,但 LLM 长文本生成、批量文档解析、知识库构建和智能体工作流往往需要持续运行数秒甚至更久。让客户端一直占用连接等待结果,容易受到网关超时、网络抖动和服务重启影响。
更常见的做法是拆成两步:
- 客户端提交任务,服务端立即返回
task_id。 - 任务执行完成后,服务端向业务系统发送 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,这样攻击者即使拿到一份合法请求,也不能轻易把它改成另一个时间再次提交。
接收端验证时有三个关键步骤:
- 读取原始请求体。签名校验必须使用收到的原始字节,不能先解析 JSON 再重新序列化,因为空格、字段顺序和转义形式都可能变化。
- 检查时间窗口。示例允许五分钟偏差,生产环境应根据网络延迟和重试策略设置,并确保服务器时间同步。
- 使用常量时间比较。
timingSafeEqual可以减少逐字节比较造成的时序信息泄露,同时要先检查两个 Buffer 长度是否一致,否则 Node.js 会抛出异常。
时间戳只能降低重放风险,不能替代幂等。一个合法请求在五分钟内仍可能被重复发送,因此必须持久化 event.id。推荐建立如下唯一约束:
| 数据 | 约束或索引 | 作用 |
|---|---|---|
event_id | 唯一索引 | 阻止同一事件重复执行业务逻辑 |
task_id、version | 唯一索引或版本条件更新 | 防止旧状态覆盖新状态 |
delivery_id | 普通索引 | 追踪一次投递的多次尝试 |
消费过程还要考虑并发。如果两个相同事件同时到达,不能采用“先查询、后插入”的非原子流程,否则两个请求都可能认为事件尚未处理。更可靠的方式是先插入幂等记录,利用唯一键冲突判断重复;随后在同一事务中更新业务表。
如果业务处理包含外部副作用,例如发送邮件、扣减额度或创建下游订单,则仅记录事件 ID 还不够。应为副作用建立独立的操作 ID,或者使用事务消息、Outbox 和下游幂等键,确保重试不会重复执行不可逆操作。
五、状态机与故障处理边界
任务状态最好定义成有限状态机,而不是任意字符串。一个简单模型可以是 pending -> running -> succeeded,失败路径为 pending -> running -> failed,取消路径为 pending/running -> canceled。每次状态变更都递增版本号,并把版本号随事件发送。
状态事件和结果事件也可以分开。例如 ai.task.progress 只用于展示进度,ai.task.completed 才代表最终结果。消费端不应仅凭事件到达顺序判断最终状态,而应根据版本号和状态机规则校验转换是否合法。
Webhook 接收端应尽快返回成功,不要在 HTTP 请求中执行耗时的文档入库、向量化或二次模型调用。常见流程是:验签、校验基本字段、原子写入事件表,然后返回 2xx;真正的业务处理交给内部队列。这样可以缩短发送方的等待时间,也能把业务失败与网络投递失败分开处理。
返回状态码可以按下面的方式划分:
| 状态 | 含义 | 发送方行为 |
|---|---|---|
2xx | 已接收并持久化 | 结束本次投递 |
400 | JSON 或字段非法 | 记录错误,通常不重试 |
401/403 | 签名或权限失败 | 告警并人工处理配置 |
408/429 | 超时或限流 | 延迟后重试 |
500-599 | 接收方暂时不可用 | 指数退避重试 |
发送方还应记录每次投递的事件 ID、目标地址、尝试次数、响应状态、耗时和最后错误。接收方则需要提供事件查询、失败重放和任务状态查询能力。这样即使事件永久失败,也可以通过后台任务修复,而不是要求用户重新提交一次 AI 任务。
总结
Webhook 可靠交付的核心不是“发出一个 POST 请求”,而是建立一套可重试、可验证、可去重和可恢复的事件协议。
要点可以归纳为:
- 用
202 Accepted表示任务已接受,不把请求成功等同于结果完成。 - 每个事件分配稳定的事件 ID,并在消费端使用持久化唯一约束实现幂等。
- 使用任务版本号处理乱序到达,避免旧状态覆盖新状态。
- 使用时间戳加 HMAC 签名校验来源,并用常量时间比较验证签名。
- 对连接超时、限流和服务端错误进行指数退避,对永久错误进入死信或人工处理流程。
- 接收端先快速验签和落库,再异步执行耗时业务;任务查询接口作为 Webhook 失败后的补偿路径。
对于长时间运行的 LLM、文档处理和智能体任务,这些机制比单纯增加请求超时时间更重要。它们把一次不可靠的网络通知,变成能够审计、重试和恢复的交付流程。
评论