
你大概率已经在生产环境里埋了这个 bug,只是自己不知道。
你的应用处理一个订单:往数据库插一条记录,给合作伙伴的 API 发一个 webhook,两个操作都成功了。然后出问题了——可能是约束冲突、网络抖动、代码抛异常——数据库事务回滚了。
订单从来没存在过。但 webhook 已经发出去了。
你的合作伙伴收到了一条幽灵订单的通知,给一个从来没付款的用户发了货。你的系统里没有任何记录表明这一切发生过。
这就是幽灵 webhook 问题。它静默发生,在生产环境里,大多数 webhook 实现都有这个毛病。
标准的 webhook 实现大概长这样:
1. 开始事务
2. INSERT order INTO orders
3. 提交事务
4. POST /webhook 给合作伙伴 ← 在提交之后发生 看起来很安全,对吧?提交成功之后才发 webhook。
但有三种没人想到的失败场景:
场景一:崩溃窗口
1. INSERT order ──────────────────→ 成功
2. COMMIT ───────────────────────→ 成功
3. [进程崩溃 / Pod 重启]
4. POST /webhook ─────────────────→ 根本没发生 订单在数据库里存在。Webhook 永远没发出去。合作伙伴完全不知道这件事。
场景二:超时竞态
1. INSERT order ──────────────────→ 成功
2. COMMIT ───────────────────────→ 成功
3. POST /webhook ─────────────────→ 超时(合作伙伴接口慢)
4. 重试 POST /webhook ───────────→ 成功(重复发送了!) 你发了两次 webhook。合作伙伴扣了用户两次钱。
场景三:静默回滚
1. INSERT order ──────────────────→ 成功
2. INSERT order_items ────────────→ 约束冲突 ← 第二次插入失败
3. ROLLBACK ──────────────────────→ 订单从数据库消失
4. [但你在第 1.5 步已经调用了 POST /webhook] 这种情况发生在 webhook 调用被埋在一个大事务里,后来事务回滚了。Webhook 发给了不存在的数据。
解决方案是:别把 webhook 当成独立的副作用,而是让它和业务数据属于同一个原子操作。
❌ 你现在的做法: ┌─────────────────────────────┐ │ 事务 │ │ ├── INSERT order │ │ └── COMMIT │ └─────────────────────────────┘ │ └── POST /webhook ← 在事务外面 可能独立失败 可能永远不会发生 可能发生两次 ✅ 事务性发件箱: ┌─────────────────────────────────────────┐ │ 事务 │ │ ├── INSERT order │ │ ├── INSERT webhook_event ← 同一事务 │ │ └── COMMIT(或者 ROLLBACK,原子操作) │ └─────────────────────────────────────────┘ │ └── 后台 Worker(常驻运行) │ ├── 读取 webhook_events 表 ├── POST /webhook ├── 失败时重试(指数退避) └── N 次失败后进入死信队列
要么订单和 webhook 事件同时存在,要么都不存在。后台 Worker 负责实际的 HTTP 调用,有重试逻辑和完整的审计日志。
这个模式有名字——事务性发件箱(transactional outbox)——在 Java 和 Go 圈子里已经是标准实践了。Rust 生态之前没有对应的库,直到现在。
我做了 webhooksmith 来解决 Rust + PostgreSQL 应用的问题。添加到项目里:
[dependencies]
webhooksmith = "0.1"
tokio = { version = "1", features = ["full"] }
serde_json = "1" 初始化大概 10 行代码:
use webhooksmith::WebhookEngine;
use serde_json::json; let engine = WebhookEngine::builder() .database_url("postgres://user:pass@localhost/mydb") .build() .await?; engine.migrate().await?; // 创建 3 张表,启动时调用一次即可 let endpoint = engine .register("https://partner.example.com/webhooks", "your-secret-key") .await?; 用 webhooksmith 实现发件箱模式:
let mut tx = engine.pool().begin().await?; // 你的业务逻辑
sqlx::query!("INSERT INTO orders (id, total) VALUES ($1, $2)", order_id, total) .execute(&mut *tx) .await?; // Webhook 在同一个事务里
engine .send_in_tx("order.created", json!({"id": order_id, "total": total}), endpoint.id, &mut tx) .await?; tx.commit().await?;
// 提交成功:订单存在 AND webhook 已入队
// 提交失败:订单不存在 AND webhook 没入队
// 没有幽灵 webhook,没有静默丢失 在 main() 末尾启动后台 Worker:
// 投递事件,失败重试,耗尽后移入 DLQ
engine.run().await; // 或者优雅关闭前先清空队列:
engine.run_graceful(async { tokio::signal::ctrl_c().await.ok();
}).await; Worker 用指数退避重试。如果合作伙伴接口挂了 6 小时,事件在数据库里等着,恢复之后继续投递——不需要你写任何代码。
第 1 次尝试 ──→ 合作伙伴返回 503 ──→ 约 2 秒后重试
第 2 次尝试 ──→ 合作伙伴返回 503 ──→ 约 4 秒后重试
第 3 次尝试 ──→ 合作伙伴返回 503 ──→ 约 8 秒后重试
...
第 10 次尝试 ──→ 合作伙伴返回 503 ──→ 移入死信队列 你可以查看 DLQ 并手动重新入队:
// 查看失败的事件
let dead = engine.dead_events(endpoint.id).await?;
for event in &dead { println!("{}: {} 次尝试", event.event_type, event.attempts);
} // 重新入队所有失败事件
let requeued = engine.retry_all_dead(endpoint.id).await?;
println!("已重新入队 {} 个事件", requeued); 每个事件都有投递日志——你可以看到具体发了什么、什么时候发的、HTTP 响应状态码、每次尝试花了多长时间。
如果有多个端点(多个合作伙伴、多个环境),broadcast 会原子化地投递到所有端点:
// 一次调用 = 每个端点一个事件,全在一条 SQL 里完成
let events = engine.broadcast_in_tx("order.created", json!({"id": order_id}), &mut tx).await?;
println!("已在 {} 个端点入队 {} 个 webhook", events.len(), events.len()); 如果你的代码在网络出错时会重试,可能对同一个事件调用两次 send。幂等键(idempotency key)防止重复:
// 多次调用安全——只会创建一个 webhook
engine .send_idempotent("order.created", payload, endpoint.id, "order-1001-created") .await?; 第二次调用返回第一次调用创建的事件。合作伙伴只会收到一个 webhook。
Worker 发送 HTTP POST,带有 HMAC-SHA256 签名,合作伙伴可以验证:
POST https://partner.example.com/webhooks
Content-Type: application/json
x-webhooksmith-signature: v1,a3f4b2...
x-webhooksmith-timestamp: 1735689600
x-webhooksmith-event-id: 550e8400-e29b-41d4-a716-446655440000
x-webhooksmith-event-type: order.created {"id": 1001, "total": 49.99} 如果你也用 axum 写接收端,webhooksmith-axum 可以自动验证签名:
use axum::{Router, routing::post, http::StatusCode};
use webhooksmith_axum::{WebhookSecretLayer, VerifiedWebhook}; async fn receive(VerifiedWebhook(payload): VerifiedWebhook) -> StatusCode { println!("{}: {:?}", payload.event_type, payload.body); StatusCode::OK
} let app: Router = Router::new() .route("/webhooks", post(receive)) .layer(WebhookSecretLayer::new("your-secret-key")); Extractor 会自动拒绝重放请求(5 分钟时间戳窗口)、签名错误请求、以及超过 1MB 的请求体。
let stats = engine.queue_stats().await?;
println!("pending={} delivering={} failed={} dead={}", stats.pending, stats.delivering, stats.failed, stats.dead); 如果 dead > 0 就报警——说明事件已经耗尽所有重试次数,需要人工介入。
cargo add webhooksmith webhooksmith 在 crates.io
Source + examples on GitHub
webhooksmith-axum 用于接收端
如果你在做有合作伙伴 webhook 的 SaaS,或者任何"发完就不管"不可接受的系统,这个库让你用大约 10 行代码实现发件箱模式。
幽灵订单问题是真实存在的——我在生产环境里见过。发件箱模式是正确的解法。这是我第一次踩坑时希望已经存在的那套实现。