至少要写三个接口
| 接口 | 干什么 | 别做什么 |
|---|---|---|
POST /api/checkout/create | 校验商品和金额,建本地订单,再调 Clink 创建 Session | 别让前端决定价格 |
GET /api/orders/:id | 给回跳页面查订单状态 | 别把 Clink 的原始响应透出去 |
POST /api/webhooks/clink | 接收事件,验签,更新订单,触发发货 | 别在没验签的情况下信任请求内容 |
商户订单要存什么
至少这些字段,不然出问题时排查不了: 支付尝试、退款记录和人工对账证据要分开保存。-- 商户订单:一次下单一行
CREATE TABLE merchant_orders (
id TEXT PRIMARY KEY, -- 商户订单号
customer_id TEXT NOT NULL, -- 商户系统里的用户
product_snapshot JSONB NOT NULL, -- 商品名、单价、数量(下单那一刻的快照)
original_amount NUMERIC(18,4) NOT NULL,
original_currency TEXT NOT NULL,
refundable_paid_clink_order_id TEXT, -- 建立退款基准的成功 Order
refundable_paid_amount NUMERIC(18,4), -- 权威的可退款现金金额
refundable_paid_currency TEXT, -- 这笔现金支付/退款使用的币种
clink_session_id TEXT, -- Clink 返回的 sessionId
clink_session_status TEXT NOT NULL DEFAULT 'open'
CHECK (clink_session_status IN ('open', 'completed', 'expired')),
clink_session_last_event_created BIGINT,
current_clink_order_id TEXT, -- 当前支付尝试,可能是推导值或 Session 确认值
current_attempt_confirmed BOOLEAN NOT NULL DEFAULT FALSE, -- 只有 Session 查询后才为 true
payment_status TEXT NOT NULL, -- 由下面的支付尝试聚合而来
fulfillment_status TEXT NOT NULL, -- 发货状态,和支付状态分开存
refunded_amount NUMERIC(18,4) NOT NULL DEFAULT 0,
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
CONSTRAINT ck_refundable_payment_basis CHECK (
(refundable_paid_clink_order_id IS NULL AND refundable_paid_amount IS NULL AND refundable_paid_currency IS NULL)
OR
(refundable_paid_clink_order_id IS NOT NULL AND refundable_paid_amount IS NOT NULL AND refundable_paid_currency IS NOT NULL)
),
CONSTRAINT ck_clink_session_event_created CHECK (
clink_session_last_event_created IS NULL OR clink_session_last_event_created >= 0
)
);
-- 支付尝试:一个 Clink Order 一行
CREATE TABLE payment_attempts (
clink_order_id TEXT PRIMARY KEY, -- obj.orderId
merchant_order_id TEXT NOT NULL REFERENCES merchant_orders(id),
clink_session_id TEXT NOT NULL,
status TEXT NOT NULL, -- pending / action_required / succeeded / failed
attempt_created_at BIGINT, -- 只来自 order.created 的 event.created
last_event_created BIGINT NOT NULL, -- 已处理状态事件的版本,只用于同一 Order 内比较
failure_code TEXT,
failure_message TEXT,
received_at TIMESTAMPTZ NOT NULL DEFAULT NOW(), -- 本地收到时间,仅审计用
updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW()
);
-- 排序只认 attempt_created_at;为空的排在后面,代表顺序未知
CREATE INDEX idx_attempts_order
ON payment_attempts (merchant_order_id, attempt_created_at DESC NULLS LAST);
-- 复合归属键:商户订单指针和退款不能引用另一商户订单的 Clink Order
ALTER TABLE payment_attempts
ADD CONSTRAINT uq_payment_attempt_owner UNIQUE (merchant_order_id, clink_order_id);
ALTER TABLE merchant_orders
ADD CONSTRAINT fk_current_attempt_owner
FOREIGN KEY (id, current_clink_order_id)
REFERENCES payment_attempts (merchant_order_id, clink_order_id),
ADD CONSTRAINT fk_refundable_attempt_owner
FOREIGN KEY (id, refundable_paid_clink_order_id)
REFERENCES payment_attempts (merchant_order_id, clink_order_id);
-- 成功退款:业务字段不可变,每个 Clink refundId 一行
CREATE TABLE refunds (
refund_id TEXT PRIMARY KEY,
merchant_order_id TEXT NOT NULL REFERENCES merchant_orders(id),
clink_order_id TEXT NOT NULL REFERENCES payment_attempts(clink_order_id),
amount NUMERIC(18,4) NOT NULL CHECK (amount > 0),
currency TEXT NOT NULL CHECK (currency ~ '^[A-Z]{3}$'),
status TEXT NOT NULL CHECK (status = 'success'),
first_event_id TEXT NOT NULL, -- 仅供审计;精确重放时不改写
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
CONSTRAINT fk_refund_attempt_owner
FOREIGN KEY (merchant_order_id, clink_order_id)
REFERENCES payment_attempts (merchant_order_id, clink_order_id)
);
CREATE INDEX idx_refunds_success
ON refunds (merchant_order_id, status);
-- 自动流程不可重试或覆盖的数据冲突,必须把证据持久化
CREATE TABLE reconciliation_cases (
dedupe_key TEXT PRIMARY KEY, -- 稳定业务键,不是随机任务 ID
event_id TEXT NOT NULL,
merchant_order_id TEXT NOT NULL REFERENCES merchant_orders(id),
clink_order_id TEXT,
refund_id TEXT,
reason TEXT NOT NULL,
evidence JSONB NOT NULL,
status TEXT NOT NULL DEFAULT 'pending',
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW()
);
merchant_orders 表,应先增加 Session 事件版本列,再部署处理器:
ALTER TABLE merchant_orders
ADD COLUMN IF NOT EXISTS clink_session_last_event_created BIGINT;
ALTER TABLE merchant_orders
DROP CONSTRAINT IF EXISTS ck_clink_session_event_created;
ALTER TABLE merchant_orders
ADD CONSTRAINT ck_clink_session_event_created CHECK (
clink_session_last_event_created IS NULL OR clink_session_last_event_created >= 0
);
paymentAttempts.insertIfAbsent() 不是「先查再写」。它必须让 clink_order_id 主键裁决并发首次投递:
INSERT INTO payment_attempts
(clink_order_id, merchant_order_id, clink_session_id, status,
attempt_created_at, last_event_created, failure_code, failure_message)
VALUES
($1, $2, $3, $4, $5, $6, $7, $8)
ON CONFLICT (clink_order_id) DO NOTHING
RETURNING *;
merchant_order_id 和 clink_session_id。归属不同就返回 manual_review;绝不能去锁另一商户订单,也不能覆盖它的 attempt。
refunds 行不可变。下文的 insertIfAbsent() 指执行以下 SQL;没有返回行时,再读取数据库中已有的记录:
INSERT INTO refunds
(refund_id, merchant_order_id, clink_order_id, amount, currency, status, first_event_id)
VALUES
($1, $2, $3, $4, $5, $6, $7)
ON CONFLICT (refund_id) DO NOTHING
RETURNING *;
refund_id 才算精确重放。绝不能使用 DO UPDATE:不可变字段不同是人工对账证据,不是覆盖旧行的理由。reconciliation_cases.insertIfAbsent() 同样采用 INSERT ... ON CONFLICT (dedupe_key) DO NOTHING。
INSERT INTO reconciliation_cases
(dedupe_key, event_id, merchant_order_id, clink_order_id, refund_id, reason, evidence)
VALUES
($1, $2, $3, $4, $5, $6, $7::jsonb)
ON CONFLICT (dedupe_key) DO NOTHING;
下面用 项目已有的 Decimal 库也可以。金额在数据库里用
decimal.js 做金额运算:npm install decimal.js
NUMERIC,应用层用 Decimal,绝不能用二进制浮点数。**退款基准的适用范围:**下面的自动退款逻辑只适用于 Hosted Checkout 的单一现金资金来源支付,并且成功 Order 的
amountTotal 就是权威可退款现金额。公开 Order 契约有 amountTotal 和 paymentCurrency,但余额/积分加现金的拆分支付没有单独暴露 cashAmount。如果启用了混合资金来源,不能拿 originalAmount 或总 amountTotal 代替。必须先从后端契约取得权威的可退款现金字段;在此之前,这类订单应转人工对账,不运行下面的自动退款状态计算。确定性数据冲突返回
manual_review,不能作为可重试异常抛出。返回前,同一事务必须写入不可变的 reconciliation_cases 行和去重后的 manual_reconciliation Outbox 任务。只有这两项持久化后,Webhook 才能标成 processed 并确认接收。Worker 再用稳定 case key 幂等创建或更新运营工单和告警。不要用本地插入时间判断哪一笔是当前尝试。
received_at 是商户库自己的写入时刻,受推送顺序、重试和处理延迟影响,和 Clink 那边 Order 的创建先后没有关系。它只用来排查问题。能代表 Order 创建顺序的只有一个来源:order.created 事件的外层 event.created。其他 order.* 事件的 event.created 是各自状态发生的时间,不是 Order 创建时间;data.object 里没有 Order 的创建时间,GET /order/{orderId} 目前也不返回。所以 order.created 必须订阅。四个事件一个都不能少:order.created、order.next_action、order.succeeded、order.failed。order.* 事件都写进商户订单那一行,第一笔的失败事件迟到就会把第二笔的成功盖掉。
所以按 obj.orderId 一笔支付尝试一行,商户订单的支付状态由这些尝试聚合出来。
下面的 SQL 用
snake_case 列名,JavaScript 示例统一用 camelCase 属性名(clinkOrderId、clinkSessionStatus、paymentStatus、attemptCreatedAt)。两者指的是同一个字段,具体转换交给你的 ORM 或查询层。因此本页所有 tx.* / db.* repository 方法都假定返回 camelCase 对象。如果绕过 repository 直接执行 node-postgres SQL,必须像后面的 claimTask() 一样为每个会读取的列显式加 alias。不得返回“业务 ID 是 camelCase、claim_token 却是 snake_case”的混合对象。session.complete 只把 clinkSessionStatus 推进到 completed,session.expired 只推进到 expired;两者都不能决定支付、退款、履约或 Attempt 状态。clinkSessionLastEventCreated 保存 Session 终态事件的 Unix 毫秒版本,防止迟到的旧终态覆盖较新的状态。面向用户的提示应综合 Session 状态、聚合支付状态,以及是否仍有 pending/action-required Attempt。
1. 创建订单和 Session
const CLINK_API = 'https://uat-api.clinkbill.com/api';
function assertClinkPath(path) {
let decodedPath = '';
try {
decodedPath = decodeURIComponent(String(path).split(/[?#]/, 1)[0]);
} catch {
throw new Error('Clink API path must be a controlled relative path');
}
if (
typeof path !== 'string' ||
!path.startsWith('/') ||
path.startsWith('//') ||
path.includes('\\') ||
decodedPath.startsWith('//') ||
decodedPath.split('/').includes('..')
) {
throw new Error('Clink API path must be a controlled relative path');
}
}
function throwSanitizedClinkError(method, path, detail) {
// Only method, controlled path, HTTP status, and numeric envelope code are
// allowed into the error. Never copy the key, response body, or json.msg.
throw new Error(`Clink ${method} ${path} failed${detail ? `: ${detail}` : ''}`);
}
async function clinkApiRequest(path, { method, body, signal } = {}) {
assertClinkPath(path);
if (method !== 'GET' && method !== 'POST') {
throw new Error('Clink API method must be GET or POST');
}
const apiKey = process.env.CLINK_SECRET_KEY;
if (!apiKey) throw new Error('CLINK_SECRET_KEY is not configured');
const headers = {
'X-API-Key': apiKey,
'X-Timestamp': String(Date.now()),
Accept: 'application/json',
};
const options = { method, headers, signal };
if (method === 'POST') {
headers['Content-Type'] = 'application/json';
options.body = JSON.stringify(body);
}
let res;
try {
res = await fetch(`${CLINK_API}${path}`, options);
} catch (err) {
if (err?.name === 'AbortError') throw err;
throwSanitizedClinkError(method, path, 'request error');
}
let json;
try {
json = await res.json();
} catch (err) {
if (err?.name === 'AbortError') throw err;
throwSanitizedClinkError(method, path, 'invalid JSON response');
}
const dataIsObject =
json?.data !== null && typeof json?.data === 'object' && !Array.isArray(json.data);
if (!res.ok || json?.code !== 200 || !dataIsObject) {
const details = [];
if (Number.isInteger(res.status)) details.push(`HTTP ${res.status}`);
if (Number.isInteger(json?.code)) details.push(`code ${json.code}`);
throwSanitizedClinkError(method, path, details.join(', ') || 'invalid response envelope');
}
return json.data;
}
async function clinkRequest(path, body, { signal } = {}) {
return clinkApiRequest(path, { method: 'POST', body, signal });
}
async function clinkGet(path, { signal } = {}) {
return clinkApiRequest(path, { method: 'GET', signal });
}
app.post('/api/checkout/create', async (req, res) => {
const { productId, quantity } = req.body;
const user = req.user;
// 价格从本地数据库取,不要信前端传来的金额
const product = await db.products.findById(productId);
// 价格按整数分存,算完再转回元
const unitAmount = product.unitAmountMinor / 100; // 1999 → 19.99
const amount = (product.unitAmountMinor * quantity) / 100;
const order = await db.merchantOrders.create({
customerId: user.id,
productSnapshot: { name: product.name, unitAmount, quantity },
originalAmount: amount,
originalCurrency: 'USD',
paymentStatus: 'created',
fulfillmentStatus: 'none',
});
const session = await clinkRequest('/checkout/session', {
customerEmail: user.email,
referenceCustomerId: user.id,
originalAmount: amount,
originalCurrency: 'USD',
merchantReferenceId: order.id,
uiMode: 'hostedPage',
priceDataList: [
{ name: product.name, quantity, unitAmount, currency: 'USD' },
],
successUrl: `https://your-site.com/pay/result?orderId=${order.id}`,
cancelUrl: `https://your-site.com/pay/cancel?orderId=${order.id}`,
});
await db.merchantOrders.update(order.id, {
clinkSessionId: session.sessionId,
clinkSessionStatus: 'open',
paymentStatus: 'pending',
});
res.json({ merchantOrderId: order.id, checkoutUrl: session.url });
});
import { ClinkPayClient } from '@clink-ai/clink-typescript-sdk';
const client = new ClinkPayClient({
apiKey: process.env.CLINK_SECRET_KEY,
env: 'sandbox',
});
app.post('/api/checkout/create', async (req, res) => {
const { productId, quantity } = req.body;
const user = req.user;
const product = await db.products.findById(productId);
const unitAmount = product.unitAmountMinor / 100;
const amount = (product.unitAmountMinor * quantity) / 100;
const order = await db.merchantOrders.create({
customerId: user.id,
productSnapshot: { name: product.name, unitAmount, quantity },
originalAmount: amount,
originalCurrency: 'USD',
paymentStatus: 'created',
fulfillmentStatus: 'none',
});
const session = await client.createCheckoutSession({
customerEmail: user.email,
referenceCustomerId: user.id,
originalAmount: amount,
originalCurrency: 'USD',
merchantReferenceId: order.id,
uiMode: 'hostedPage',
priceDataList: [
{ name: product.name, quantity, unitAmount, currency: 'USD' },
],
successUrl: `https://your-site.com/pay/result?orderId=${order.id}`,
cancelUrl: `https://your-site.com/pay/cancel?orderId=${order.id}`,
});
await db.merchantOrders.update(order.id, {
clinkSessionId: session.sessionId,
clinkSessionStatus: 'open',
paymentStatus: 'pending',
});
res.json({ merchantOrderId: order.id, checkoutUrl: session.url });
});
服务端也可以用
@clink-ai/clink-typescript-sdk,它会处理认证请求头并提供类型定义。上面直接调 API 的写法是为了让每个字段的位置更清楚。- 金额从数据库取。前端传过来的价格一律不信,否则客户改个请求就能一块钱买走会员。
- 金额是主货币单位,不是分。19.99 美元传
19.99。传1999会被当成 1999 美元,扣款差 100 倍。日元、韩元、印尼盾这类零小数位币种只能传整数。 - 别用浮点直接乘金额。JavaScript 里
0.1 * 3等于0.30000000000000004,而 Clink 会精确比对商品明细之和与总额,这点误差就会被拒。价格按整数分存,乘完数量再转回主单位;金额结构复杂时用 decimal 类库。 merchantReferenceId填商户自己的订单号。Clink 不拿它做幂等——同一个值调两次会得到两个 Session。防重复下单由商户负责。referenceCustomerId填商户系统里的用户 ID。以后再给同一个客户创建 Session,Clink 能自动关联到同一个 Clink 客户。
2. 回跳页面
客户付完款跳回successUrl。这个页面要做的事只有一件:查本地后端的订单状态,把结果显示出来。
// 回跳页面
const orderId = new URLSearchParams(location.search).get('orderId');
const order = await fetch(`/api/orders/${orderId}`).then((r) => r.json());
if (order.paymentStatus === 'paid') {
showSuccess();
} else if (order.paymentStatus === 'pending') {
showPending(); // 「支付确认中」,几秒后再查一次
} else {
showFailed();
}
pending。这时候显示「支付确认中」,隔几秒再查一次,别直接显示失败。
如果一直是 pending,后端可以主动查一次 GET /checkout/session/{id} 补上状态。
3. 接收 Webhook
这是整个接入里最要紧的一段。付款结果以这里为准。验签
Clink 用 HMAC SHA-256 签名,签的是时间戳 + "." + 原始请求体。
完整的验签要做三件事,缺一件都不算验过:校验 signType、校验时间戳新鲜度、恒定时间比对签名。下面这段是完整实现,直接抄:
import crypto from 'node:crypto';
const REPLAY_WINDOW_MS = 5 * 60 * 1000; // 5 分钟
function verifyWebhook(rawBody, headers) {
// 1. 算法必须是 SHA256,不是就别往下算
if (headers['x-clink-signtype'] !== 'SHA256') return false;
// 2. 时间戳是 Unix 毫秒。超出窗口或不是数字,直接拒
const ts = Number(headers['x-clink-timestamp']);
if (!Number.isFinite(ts)) return false;
if (Math.abs(Date.now() - ts) > REPLAY_WINDOW_MS) return false;
// 3. 用原始请求体算 HMAC,此时还没有 JSON.parse
const expected = crypto
.createHmac('sha256', process.env.CLINK_WEBHOOK_SIGNING_KEY)
.update(`${headers['x-clink-timestamp']}.${rawBody}`)
.digest('hex');
// 4. 恒定时间比较。长度不同要先挡掉,timingSafeEqual 长度不等会抛
const a = Buffer.from(expected, 'utf8');
const b = Buffer.from(headers['x-clink-signature'] ?? '', 'utf8');
return a.length === b.length && crypto.timingSafeEqual(a, b);
}
X-Clink-Timestamp 是 Unix 毫秒时间戳,所以直接和 Date.now() 比。5 分钟是这里选的窗口,平台没有强制值——窗口越大,被截获的请求可重放的时间越长;越小则对时钟漂移越敏感。服务器请用 NTP 校时。
时间校验和 HMAC 都必须在
JSON.parse 之前完成,而且用原始请求体。 解析后再序列化回去,字段顺序和空格都可能变,签名一定对不上。Express 里要用 express.raw(),不能用 express.json()。时间窗口只防重放,不防重复。 Clink 的重试也在窗口内到达,同一个事件仍然会收到多遍。重复必须另外按
event.id 原子去重,见下一节。用 TypeScript SDK 的话
用 TypeScript SDK 的话
@clink-ai/clink-typescript-sdk@1.0.1 的 ClinkWebhook.verifyAndGet() 只做 HMAC 那一步,上面三件事里它只覆盖了一件:- 不校验
X-Clink-SignType - 不校验时间戳新鲜度
- 签名用
===比较,不是恒定时间比较
import { ClinkWebhook, ClinkWebhookSignatureError } from '@clink-ai/clink-typescript-sdk';
const webhook = new ClinkWebhook({
signatureKey: process.env.CLINK_WEBHOOK_SIGNING_KEY,
});
function verifyWithSdk(rawBody, headers) {
if (headers['x-clink-signtype'] !== 'SHA256') return null;
const ts = Number(headers['x-clink-timestamp']);
if (!Number.isFinite(ts) || Math.abs(Date.now() - ts) > REPLAY_WINDOW_MS) return null;
try {
return webhook.verifyAndGet({
timestamp: headers['x-clink-timestamp'],
body: rawBody,
headerSignature: headers['x-clink-signature'],
});
} catch (e) {
if (e instanceof ClinkWebhookSignatureError) return null;
throw e;
}
}
完整的处理器
webhook_events 不只是一张去重表,它同时是待处理队列。表上至少要有这几列:
CREATE TABLE webhook_events (
id TEXT PRIMARY KEY, -- event.id;约束名为 webhook_events_pkey
type TEXT NOT NULL,
payload TEXT NOT NULL, -- 精确的原始事件体,重处理时要用
status TEXT NOT NULL DEFAULT 'pending'
CHECK (status IN ('pending', 'processing', 'processed', 'failed')),
retry_count INT NOT NULL DEFAULT 0,
next_retry_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
last_error TEXT,
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW()
);
// 注意:express.raw(),不是 express.json()
app.post(
'/api/webhooks/clink',
express.raw({ type: 'application/json' }),
async (req, res) => {
const rawBody = req.body.toString('utf8');
if (!verifyWebhook(rawBody, req.headers)) {
return res.status(401).send('invalid signature');
}
const event = JSON.parse(rawBody);
try {
await db.transaction(async (tx) => {
// 1. 落库 + 去重。event.id 上有唯一索引,
// 并发投递里只有一个事务能插进去
await tx.webhookEvents.insert({
id: event.id,
type: event.type,
payload: rawBody,
status: 'pending',
});
// 2. 正常完成和已可靠隔离的冲突都是终态;缺少依赖才留 pending
const outcome = await handleEvent(event, tx);
if (outcome === 'done' || outcome === 'manual_review') {
await tx.webhookEvents.update(event.id, { status: 'processed' });
} else {
// outcome === 'deferred':依赖对象还不存在,登记一个重处理任务,
// 事件本身保持 pending,等 Worker 回来再试
await tx.outbox.insert({
eventId: event.id,
task: 'reprocess_webhook',
dedupeKey: `reprocess_webhook:${event.id}`,
});
}
});
} catch (err) {
// 重复投递:只认 webhook_events.id 这一个约束
if (isDuplicateEvent(err)) return res.status(200).send('ok');
// 这里只处理意外或暂时性失败。确定性冲突已经和事件一起提交人工
// 对账证据,不会进入这里;本分支事务已回滚,应让 Clink 重试
logger.error({ err, eventId: event.id }, 'webhook failed');
return res.status(500).send('retry later');
}
// 3. 事件与重处理任务都已可靠落库,才能确认
res.status(200).send('ok');
}
);
// 按约束名精确识别,不要用「是不是唯一键冲突」这种粗判断
function isDuplicateEvent(err) {
return err.code === '23505' && err.constraint === 'webhook_events_pkey';
}
依赖对象不存在时,绝对不能把事件标成
processed。refund.succeeded 完全可能早于 order.succeeded 到达。如果这时候直接 return 并返回 200,事件就被记成处理过了——Clink 不会再推,这笔退款永远不会入账。依赖缺失时应保持 pending 并登记重处理任务,等依赖具备后由 Worker 再试。只有业务已完成,或确定性冲突及其人工对账任务已在同一事务可靠隔离后,才能置 processed。handleEvent 返回 'done'、'deferred' 或 'manual_review'。manual_review 之所以是终态,是因为冲突证据和运营处理任务已经在同一事务内持久化。
事务里的步骤顺序是固定的,退款和订单事件都走同一套,锁顺序一致才不会死锁:
1
插入并去重 webhook_events
2
定位商户订单,SELECT ... FOR UPDATE 锁住这一行
3
校验必要标识和商户订单的 Session
4
插入或重读 payment_attempts,再校验商户与 Session 归属
5
锁内重新读取该订单的全部 attempts
6
聚合并更新 paymentStatus
7
同一事务写入 fulfillment / reconciliation Outbox
8
标记 processed,或依赖缺失时保持 pending
import Decimal from 'decimal.js';
const ATTEMPT_STATUS = {
'order.created': 'pending',
'order.next_action': 'action_required',
'order.succeeded': 'succeeded',
'order.failed': 'failed',
};
const TERMINAL = new Set(['succeeded', 'failed']);
// Single-attempt status -> merchant order status
const ATTEMPT_TO_ORDER = {
succeeded: 'paid',
failed: 'payment_failed',
action_required: 'action_required',
pending: 'pending',
};
function tryNormalizeCurrency(value) {
const currency = String(value ?? '').trim().toUpperCase();
return /^[A-Z]{3}$/.test(currency) ? currency : null;
}
function tryPositiveDecimal(value) {
try {
const amount = new Decimal(String(value));
return amount.isFinite() && amount.isPositive() ? amount : null;
} catch {
return null;
}
}
// Manual-review evidence is deliberately allowlisted. Never copy a whole
// Webhook/API object: it may contain customerEmail, metadata, or payment data.
function attemptEvidence(attempt) {
if (!attempt) return null;
return {
clinkOrderId: attempt.clinkOrderId,
merchantOrderId: attempt.merchantOrderId,
clinkSessionId: attempt.clinkSessionId,
status: attempt.status,
attemptCreatedAt: attempt.attemptCreatedAt ?? null,
lastEventCreated: attempt.lastEventCreated,
};
}
function orderEvidence(order) {
if (!order) return null;
return {
orderId: order.orderId ?? null,
sessionId: order.sessionId ?? null,
merchantReferenceId: order.merchantReferenceId ?? null,
status: order.status ?? null,
amountTotal: order.amountTotal == null ? null : String(order.amountTotal),
paymentCurrency: order.paymentCurrency ?? null,
failureCode: order.failureCode ?? null,
};
}
function refundEvidence(refund) {
if (!refund) return null;
return {
refundId: refund.refundId ?? null,
merchantOrderId: refund.merchantOrderId ?? null,
clinkOrderId: refund.clinkOrderId ?? refund.orderId ?? null,
amount: refund.amount == null ? (refund.refundAmount == null ? null : String(refund.refundAmount)) : String(refund.amount),
currency: refund.currency ?? refund.refundCurrency ?? null,
status: refund.status ?? null,
};
}
function sessionEvidence(session) {
if (!session) return null;
return {
sessionId: session.sessionId ?? null,
orderId: session.orderId ?? null,
merchantReferenceId: session.merchantReferenceId ?? null,
status: session.status ?? null,
};
}
// Persist both the case and the notification task before returning a terminal
// manual_review result. Neither insert updates the first evidence on a replay.
async function recordManualReview(tx, review) {
await tx.reconciliationCases.insertIfAbsent({
dedupeKey: review.dedupeKey,
eventId: review.eventId,
merchantOrderId: review.merchantOrderId,
clinkOrderId: review.clinkOrderId ?? null,
refundId: review.refundId ?? null,
reason: review.reason,
evidence: review.evidence,
status: 'pending',
});
await tx.outbox.insert({
eventId: review.eventId,
task: 'manual_reconciliation',
orderId: review.merchantOrderId,
clinkOrderId: review.clinkOrderId ?? null,
payload: { caseKey: review.dedupeKey },
dedupeKey: `manual_reconciliation:${review.dedupeKey}`,
});
return 'manual_review';
}
function recordOrderOwnershipReview(tx, merchantOrder, event, incoming, storedAttempt, reason) {
const orderKey = incoming.clinkOrderId ?? `event:${event.id}`;
return recordManualReview(tx, {
dedupeKey: `order_ownership_conflict:${merchantOrder.id}:${orderKey}`,
eventId: event.id,
merchantOrderId: merchantOrder.id,
clinkOrderId: incoming.clinkOrderId ?? null,
reason,
evidence: {
reason,
eventId: event.id,
eventType: event.type,
eventCreated: Number.isFinite(event.created) ? event.created : null,
incoming: {
merchantOrderId: merchantOrder.id,
clinkSessionId: incoming.clinkSessionId ?? null,
clinkOrderId: incoming.clinkOrderId ?? null,
},
storedMerchantOrder: {
merchantOrderId: merchantOrder.id,
clinkSessionId: merchantOrder.clinkSessionId,
currentClinkOrderId: merchantOrder.currentClinkOrderId ?? null,
},
storedAttempt: attemptEvidence(storedAttempt),
},
});
}
function recordSessionOwnershipReview(tx, merchantOrder, event, session, reason) {
const incomingSessionId = session?.sessionId ?? null;
return recordManualReview(tx, {
dedupeKey:
`session_review:${merchantOrder.id}:${reason}:${event.type}:${incomingSessionId ?? 'missing'}`,
eventId: event.id,
merchantOrderId: merchantOrder.id,
clinkOrderId: session?.orderId ?? null,
reason,
evidence: {
reason,
eventId: event.id,
eventType: event.type,
eventCreated: Number.isFinite(event.created) ? event.created : null,
incoming: {
merchantOrderId: session?.merchantReferenceId ?? null,
clinkSessionId: incomingSessionId,
clinkOrderId: session?.orderId ?? null,
status: session?.status ?? null,
},
storedMerchantOrder: {
merchantOrderId: merchantOrder.id,
clinkSessionId: merchantOrder.clinkSessionId,
clinkSessionStatus: merchantOrder.clinkSessionStatus,
clinkSessionLastEventCreated:
merchantOrder.clinkSessionLastEventCreated ?? null,
},
},
});
}
const SESSION_STATUS_BY_EVENT = {
'session.complete': 'completed',
'session.expired': 'expired',
};
async function handleSessionTerminal(event, session, tx) {
if (!session || typeof session !== 'object' || Array.isArray(session)) {
return 'deferred';
}
// Without a merchant reference there is no safe owner row for a durable case.
if (!session.merchantReferenceId) return 'deferred';
const merchantOrder = await tx.merchantOrders.findByIdForUpdate(
session.merchantReferenceId,
);
if (!merchantOrder) return 'deferred';
if (!session.sessionId) {
return recordSessionOwnershipReview(
tx,
merchantOrder,
event,
session,
'session_event_missing_required_identifier',
);
}
if (merchantOrder.clinkSessionId !== session.sessionId) {
return recordSessionOwnershipReview(
tx,
merchantOrder,
event,
session,
'session_event_ownership_conflict',
);
}
const nextStatus = SESSION_STATUS_BY_EVENT[event.type];
if (!nextStatus) return 'deferred';
if (!Number.isSafeInteger(event.created) || event.created < 0) {
return recordSessionOwnershipReview(
tx,
merchantOrder,
event,
session,
'session_event_invalid_created',
);
}
if (session.status != null && session.status !== nextStatus) {
return recordSessionOwnershipReview(
tx,
merchantOrder,
event,
session,
'session_event_status_mismatch',
);
}
const previousCreated = merchantOrder.clinkSessionLastEventCreated ?? null;
if (previousCreated != null && event.created < previousCreated) {
return 'done';
}
if (previousCreated === event.created) {
if (merchantOrder.clinkSessionStatus === nextStatus) return 'done';
return recordSessionOwnershipReview(
tx,
merchantOrder,
event,
session,
'session_terminal_same_version_conflict',
);
}
await tx.merchantOrders.update(merchantOrder.id, {
clinkSessionStatus: nextStatus,
clinkSessionLastEventCreated: event.created,
});
merchantOrder.clinkSessionStatus = nextStatus;
merchantOrder.clinkSessionLastEventCreated = event.created;
// Session lifecycle never changes the payment aggregate, refund status,
// fulfillment status, or an Attempt terminal state.
return 'done';
}
// Call only while holding the merchant-order row lock. The first successful
// Order establishes the refund basis; exact replays are harmless, but a
// different Order, amount, or currency is a duplicate-payment/data conflict.
async function persistRefundablePaymentBasis(tx, merchantOrder, payment, context) {
const existing = {
clinkOrderId: merchantOrder.refundablePaidClinkOrderId,
amount: merchantOrder.refundablePaidAmount,
currency: merchantOrder.refundablePaidCurrency,
};
const incoming = {
clinkOrderId: payment.clinkOrderId ?? null,
amount: payment.amountTotal == null ? null : String(payment.amountTotal),
currency: payment.paymentCurrency ?? null,
};
const dedupeKey = context.dedupeKey ??
`refund_basis_conflict:${merchantOrder.id}:${payment.clinkOrderId ?? 'missing'}`;
const amount = tryPositiveDecimal(payment.amountTotal);
const currency = tryNormalizeCurrency(payment.paymentCurrency);
const successfulAttempt = payment.clinkOrderId
? await tx.paymentAttempts.findByClinkOrderId(payment.clinkOrderId)
: null;
if (
!successfulAttempt ||
successfulAttempt.merchantOrderId !== merchantOrder.id ||
successfulAttempt.clinkSessionId !== merchantOrder.clinkSessionId ||
successfulAttempt.status !== 'succeeded'
) {
return recordManualReview(tx, {
dedupeKey,
eventId: context.eventId,
merchantOrderId: merchantOrder.id,
clinkOrderId: payment.clinkOrderId ?? null,
refundId: context.refundId ?? null,
reason: 'refund_basis_attempt_ownership_conflict',
evidence: {
storedBasis: existing,
incomingPayment: incoming,
storedAttempt: attemptEvidence(successfulAttempt),
},
});
}
if (!payment.clinkOrderId || !amount || !currency) {
return recordManualReview(tx, {
dedupeKey,
eventId: context.eventId,
merchantOrderId: merchantOrder.id,
clinkOrderId: payment.clinkOrderId ?? null,
refundId: context.refundId ?? null,
reason: 'invalid_refund_payment_basis',
evidence: { stored: existing, incoming },
});
}
const isEmpty = existing.clinkOrderId == null && existing.amount == null && existing.currency == null;
if (isEmpty) {
const basis = {
refundablePaidClinkOrderId: payment.clinkOrderId,
refundablePaidAmount: amount.toString(),
refundablePaidCurrency: currency,
};
await tx.merchantOrders.update(merchantOrder.id, basis);
Object.assign(merchantOrder, basis);
return 'stored';
}
const existingAmount = tryPositiveDecimal(existing.amount);
const existingCurrency = tryNormalizeCurrency(existing.currency);
const isExactReplay =
existing.clinkOrderId === payment.clinkOrderId &&
existingAmount?.equals(amount) === true &&
existingCurrency === currency;
if (!isExactReplay) {
return recordManualReview(tx, {
dedupeKey,
eventId: context.eventId,
merchantOrderId: merchantOrder.id,
clinkOrderId: payment.clinkOrderId,
refundId: context.refundId ?? null,
reason: 'refund_payment_basis_conflict',
evidence: { stored: existing, incoming: { ...incoming, amount: amount.toString(), currency } },
});
}
return 'replayed';
}
async function handleEvent(event, tx) {
const obj = event.data?.object;
if (SESSION_STATUS_BY_EVENT[event.type]) {
return handleSessionTerminal(event, obj, tx);
}
if (event.type.startsWith('refund.')) {
return obj ? handleRefund(event, obj, tx) : 'deferred';
}
if (!event.type.startsWith('order.')) return 'done';
const nextStatus = ATTEMPT_STATUS[event.type];
if (!nextStatus) return 'done';
// Without the merchant reference there is no merchant row under which a
// durable case can be filed. Treat it as unresolved input, not a DB error.
if (!obj?.merchantReferenceId) return 'deferred';
// Lock the merchant order. This serializes every event for that order —
// the webhook_events unique index only stops the same event.id, not two
// different events racing on one order
const merchantOrder = await tx.merchantOrders.findByIdForUpdate(obj.merchantReferenceId);
if (!merchantOrder) return 'deferred';
const incomingOwnership = {
clinkOrderId: obj.orderId ?? null,
clinkSessionId: obj.sessionId ?? null,
};
if (!obj.orderId || !obj.sessionId || !Number.isSafeInteger(event.created)) {
return recordOrderOwnershipReview(
tx,
merchantOrder,
event,
incomingOwnership,
null,
'order_event_missing_required_identifier',
);
}
if (merchantOrder.clinkSessionId !== obj.sessionId) {
const storedAttempt = await tx.paymentAttempts.findByClinkOrderId(obj.orderId);
return recordOrderOwnershipReview(
tx,
merchantOrder,
event,
incomingOwnership,
storedAttempt,
'order_event_session_mismatch',
);
}
const attemptResult = await upsertAttempt(tx, {
eventId: event.id, // needed by reconcile_attempt
clinkOrderId: obj.orderId,
merchantOrderId: merchantOrder.id,
clinkSessionId: obj.sessionId,
status: nextStatus,
eventType: event.type,
eventCreated: event.created,
failureCode: obj.failureCode,
failureMessage: obj.failureMessage,
}, merchantOrder, event);
if (attemptResult.outcome === 'manual_review') return attemptResult.outcome;
if (
attemptResult.inserted &&
attemptResult.attempt.attemptCreatedAt == null &&
merchantOrder.currentAttemptConfirmed &&
merchantOrder.currentClinkOrderId !== obj.orderId
) {
// A new unordered Order was not covered by the old Session confirmation.
await tx.merchantOrders.update(merchantOrder.id, {
currentAttemptConfirmed: false,
});
merchantOrder.currentAttemptConfirmed = false;
}
if (event.type === 'order.succeeded') {
// Persist even when order.created has not arrived yet. Do not trust a
// conflicting succeeded event unless the stored attempt actually converged
// to succeeded.
const storedAttempt = attemptResult.attempt;
if (storedAttempt?.status === 'succeeded') {
const basisOutcome = await persistRefundablePaymentBasis(
tx,
merchantOrder,
{
clinkOrderId: obj.orderId,
amountTotal: obj.amountTotal,
paymentCurrency: obj.paymentCurrency,
},
{ eventId: event.id },
);
// The attempt write above and the manual-review evidence now commit
// together. Do not aggregate or throw after quarantining the conflict.
if (basisOutcome === 'manual_review') return basisOutcome;
}
}
// order.created brought a creation time — see if the current attempt moves
if (event.type === 'order.created') {
const pointerOutcome = await maybeAdvanceCurrentAttempt(
tx,
merchantOrder,
obj.orderId,
event.created,
event,
attemptResult.inserted,
);
if (pointerOutcome === 'manual_review') return pointerOutcome;
}
const aggregateOutcome = await recomputeMerchantOrder(tx, merchantOrder, event.id);
if (aggregateOutcome !== 'done') return aggregateOutcome;
// A same-timestamp status ambiguity still needs its Order query even when
// the merchant aggregate can be kept safe in the meantime.
return attemptResult.deferReason === 'status_timestamp_tie'
? 'deferred'
: 'done';
}
// order.created can infer the current attempt only while creation times are
// unambiguous. A tie invalidates that inference and makes the Session decide.
async function maybeAdvanceCurrentAttempt(
tx,
merchantOrder,
clinkOrderId,
attemptCreatedAt,
event,
isNewAttempt,
) {
const currentId = merchantOrder.currentClinkOrderId;
if (!currentId) {
await tx.merchantOrders.update(merchantOrder.id, {
currentClinkOrderId: clinkOrderId,
currentAttemptConfirmed: false,
});
merchantOrder.currentClinkOrderId = clinkOrderId;
merchantOrder.currentAttemptConfirmed = false;
return;
}
const cur = await tx.paymentAttempts.findByClinkOrderId(currentId);
if (
!cur ||
cur.merchantOrderId !== merchantOrder.id ||
cur.clinkSessionId !== merchantOrder.clinkSessionId
) {
return recordOrderOwnershipReview(
tx,
merchantOrder,
event,
{ clinkOrderId: currentId, clinkSessionId: merchantOrder.clinkSessionId },
cur,
'current_attempt_pointer_ownership_conflict',
);
}
if (
merchantOrder.currentAttemptConfirmed &&
clinkOrderId !== currentId &&
!isNewAttempt
) {
// A delayed order.created for an Attempt already covered by the last
// Session snapshot only fills its creation time; it cannot override that
// newer authoritative Session answer.
return;
}
// The incumbent has no known creation time, or the newcomer is strictly later
if (cur.attemptCreatedAt == null || attemptCreatedAt > cur.attemptCreatedAt) {
await tx.merchantOrders.update(merchantOrder.id, {
currentClinkOrderId: clinkOrderId,
currentAttemptConfirmed: false,
});
merchantOrder.currentClinkOrderId = clinkOrderId;
merchantOrder.currentAttemptConfirmed = false;
return;
}
if (
clinkOrderId !== currentId &&
attemptCreatedAt === cur.attemptCreatedAt
) {
// A different Order has the same creation time. The stored pointer may
// merely reflect arrival order, and even an older Session answer is stale
// once this newly observed Order exists. Force a fresh Session query.
await tx.merchantOrders.update(merchantOrder.id, {
currentAttemptConfirmed: false,
});
merchantOrder.currentAttemptConfirmed = false;
}
}
// Returns { outcome: applied | ignored | deferred | manual_review, attempt }.
// The primary key, not a read-then-insert race, chooses the first owner.
async function upsertAttempt(tx, a, merchantOrder, event) {
let cur = await tx.paymentAttempts.findByClinkOrderId(a.clinkOrderId);
const isCreated = a.eventType === 'order.created';
let inserted = false;
if (!cur) {
const candidate = {
clinkOrderId: a.clinkOrderId,
merchantOrderId: a.merchantOrderId,
clinkSessionId: a.clinkSessionId,
status: a.status,
// Only order.created's event.created represents Order creation time
attemptCreatedAt: isCreated ? a.eventCreated : null,
lastEventCreated: a.eventCreated,
failureCode: a.failureCode ?? null,
failureMessage: a.failureMessage ?? null,
};
const winner = await tx.paymentAttempts.insertIfAbsent(candidate);
cur = winner ?? await tx.paymentAttempts.findByClinkOrderId(a.clinkOrderId);
if (!cur) throw new Error(`attempt ${a.clinkOrderId} disappeared after insert`);
inserted = Boolean(winner);
}
if (
cur.merchantOrderId !== merchantOrder.id ||
cur.clinkSessionId !== a.clinkSessionId
) {
const outcome = await recordOrderOwnershipReview(
tx,
merchantOrder,
event,
{ clinkOrderId: a.clinkOrderId, clinkSessionId: a.clinkSessionId },
cur,
'payment_attempt_ownership_conflict',
);
return { outcome, attempt: cur, inserted };
}
if (inserted) {
// Success is globally decisive even without order.created. Other status
// events need Session reconciliation before they can replace an older result.
const canAggregate = isCreated || a.status === 'succeeded';
return {
outcome: canAggregate ? 'applied' : 'deferred',
attempt: cur,
inserted: true,
deferReason: canAggregate ? null : 'attempt_creation_unknown',
};
}
if (isCreated) {
// A late order.created only fills in the creation time; it never
// drags the status back to pending
if (cur.attemptCreatedAt == null) {
await tx.paymentAttempts.update(a.clinkOrderId, { attemptCreatedAt: a.eventCreated });
return {
outcome: 'applied',
attempt: { ...cur, attemptCreatedAt: a.eventCreated },
inserted: false,
}; // ordering is now known, so re-aggregate
}
return { outcome: 'ignored', attempt: cur, inserted: false };
}
// Within one Order, event.created orders the status events
if (a.eventCreated < cur.lastEventCreated) {
return { outcome: 'ignored', attempt: cur, inserted: false };
}
if (a.eventCreated === cur.lastEventCreated) {
if (a.status === cur.status) {
// A replay is not complete while order.created is still missing.
// Once ordering is known, re-run aggregation even though no row changes.
const creationUnknown = cur.attemptCreatedAt == null && cur.status !== 'succeeded';
return {
outcome: creationUnknown ? 'deferred' : 'ignored',
attempt: cur,
inserted: false,
deferReason: creationUnknown ? 'attempt_creation_unknown' : null,
};
}
// Equal timestamps, different statuses: do not guess from arrival order
await tx.outbox.insert({
eventId: a.eventId,
task: 'reconcile_attempt',
clinkOrderId: a.clinkOrderId,
// One key per ambiguity, not one per Order
dedupeKey: `reconcile_attempt:${a.clinkOrderId}:${a.eventId}`,
});
return {
outcome: 'deferred',
attempt: cur,
inserted: false,
deferReason: 'status_timestamp_tie',
};
}
if (TERMINAL.has(cur.status)) {
return { outcome: 'ignored', attempt: cur, inserted: false }; // terminal never rolls back
}
const update = {
status: a.status,
lastEventCreated: a.eventCreated,
failureCode: a.failureCode ?? null,
failureMessage: a.failureMessage ?? null,
};
await tx.paymentAttempts.update(a.clinkOrderId, update);
// Creation time still missing, so this attempt cannot be treated as latest
const creationUnknown = cur.attemptCreatedAt == null && a.status !== 'succeeded';
return {
outcome: creationUnknown ? 'deferred' : 'applied',
attempt: { ...cur, ...update },
inserted: false,
deferReason: creationUnknown ? 'attempt_creation_unknown' : null,
};
}
async function recomputeMerchantOrder(tx, merchantOrder, eventId) {
// Re-read inside the lock to work from a current snapshot
const attempts = await tx.paymentAttempts.listByMerchantOrder(merchantOrder.id);
const succeededAttempts = attempts.filter((x) => x.status === 'succeeded');
const succeeded = succeededAttempts[0] ?? null;
const unorderedAttempts = attempts.filter((x) => x.attemptCreatedAt == null);
// A persisted pointer inferred from order.created cannot break a timestamp
// tie. Only a pointer confirmed by a Session query can do that.
const current = merchantOrder.currentClinkOrderId
? attempts.find((x) => x.clinkOrderId === merchantOrder.currentClinkOrderId)
: null;
if (merchantOrder.currentClinkOrderId && !current) {
const storedAttempt = await tx.paymentAttempts.findByClinkOrderId(
merchantOrder.currentClinkOrderId,
);
return recordManualReview(tx, {
dedupeKey:
`current_attempt_ownership_conflict:${merchantOrder.id}:${merchantOrder.currentClinkOrderId}`,
eventId,
merchantOrderId: merchantOrder.id,
clinkOrderId: merchantOrder.currentClinkOrderId,
reason: 'current_attempt_pointer_ownership_conflict',
evidence: {
merchantOrderId: merchantOrder.id,
currentClinkOrderId: merchantOrder.currentClinkOrderId,
storedAttempt: attemptEvidence(storedAttempt),
},
});
}
// Sort on attemptCreatedAt only — never receivedAt, an autoincrement id, or orderId
const ordered = attempts
.filter((x) => x.attemptCreatedAt != null)
.sort((a, b) => b.attemptCreatedAt - a.attemptCreatedAt);
const latestIsTied =
ordered.length > 1 &&
ordered[0].attemptCreatedAt === ordered[1].attemptCreatedAt;
const confirmedCurrent = merchantOrder.currentAttemptConfirmed ? current : null;
let paymentStatus;
let needsSessionReconciliation = false;
if (succeededAttempts.length > 0) {
paymentStatus = 'paid'; // one success is enough
} else if (confirmedCurrent) {
// A fresh Session answer may select an Attempt whose order.created has not
// arrived. The confirmed pointer is authoritative for the current Attempt.
paymentStatus = ATTEMPT_TO_ORDER[confirmedCurrent.status] ?? 'pending';
} else if (unorderedAttempts.length > 0) {
paymentStatus = 'pending';
needsSessionReconciliation = true;
} else if (latestIsTied) {
paymentStatus = 'pending';
needsSessionReconciliation = true;
} else if (ordered.length > 0) {
paymentStatus = ATTEMPT_TO_ORDER[ordered[0].status] ?? 'pending';
} else {
paymentStatus = 'pending';
}
// Second line of defence: terminal states are only rewritten by the refund flow
const merchantStatusIsTerminal =
['paid', 'refunded', 'partial_refunded'].includes(merchantOrder.paymentStatus);
if (!merchantStatusIsTerminal && merchantOrder.paymentStatus !== paymentStatus) {
await tx.merchantOrders.updateIfStatus(
merchantOrder.id,
merchantOrder.paymentStatus, // conditional UPDATE, paired with the row lock
{ paymentStatus },
);
merchantOrder.paymentStatus = paymentStatus;
}
if (
paymentStatus === 'paid' &&
!['refunded', 'partial_refunded'].includes(merchantOrder.paymentStatus)
) {
await tx.outbox.insert({
eventId,
task: 'fulfill',
orderId: merchantOrder.id,
clinkOrderId: succeeded.clinkOrderId,
dedupeKey: `fulfill:${merchantOrder.id}`, // unique; duplicate fulfillment stops here
});
}
if (needsSessionReconciliation) {
// One key per concrete ambiguity. Reprocessing the same Webhook is silent,
// while a later new Order gets its own Session query.
await tx.outbox.insert({
eventId,
task: 'reconcile_session',
orderId: merchantOrder.id,
dedupeKey: `reconcile_session:${merchantOrder.id}:${eventId}`,
});
return 'deferred';
}
return 'done';
}
function recordRefundReview(tx, context, reason, evidence) {
const refundKey = context.refundId ?? `event:${context.eventId}`;
return recordManualReview(tx, {
dedupeKey: `refund_conflict:${refundKey}`,
eventId: context.eventId,
merchantOrderId: context.merchantOrderId,
clinkOrderId: context.clinkOrderId,
refundId: context.refundId ?? null,
reason,
evidence,
});
}
function normalizeRefundCandidate(incoming, eventId) {
const amount = tryPositiveDecimal(incoming.amount);
const currency = tryNormalizeCurrency(incoming.currency);
if (
!incoming.refundId ||
!incoming.merchantOrderId ||
!incoming.clinkOrderId ||
!amount ||
!currency
) return null;
return {
refundId: incoming.refundId,
merchantOrderId: incoming.merchantOrderId,
clinkOrderId: incoming.clinkOrderId,
amount: amount.toString(),
currency,
status: 'success',
firstEventId: eventId,
};
}
function isExactRefundReplay(stored, candidate) {
const storedAmount = tryPositiveDecimal(stored?.amount);
return Boolean(stored) &&
stored.merchantOrderId === candidate.merchantOrderId &&
stored.clinkOrderId === candidate.clinkOrderId &&
storedAmount?.equals(new Decimal(candidate.amount)) === true &&
tryNormalizeCurrency(stored.currency) === candidate.currency &&
stored.status === candidate.status;
}
// INSERT ... ON CONFLICT DO NOTHING, then compare the immutable business
// tuple. A conflict never executes UPDATE and therefore never destroys proof.
async function persistRefundOnce(tx, incoming, context) {
const candidate = normalizeRefundCandidate(incoming, context.eventId);
if (!candidate) {
return {
outcome: await recordRefundReview(tx, context, 'invalid_refund_record', {
incoming: refundEvidence(incoming),
}),
refund: null,
};
}
const inserted = await tx.refunds.insertIfAbsent(candidate);
const stored = inserted ?? await tx.refunds.findByRefundId(incoming.refundId);
if (!stored) throw new Error(`refund ${incoming.refundId} disappeared after insert`);
if (!isExactRefundReplay(stored, candidate)) {
return {
outcome: await recordRefundReview(tx, context, 'refund_id_immutable_fields_conflict', {
stored: refundEvidence(stored),
incoming: refundEvidence(candidate),
}),
refund: stored,
};
}
return { outcome: inserted ? 'stored' : 'replayed', refund: stored };
}
async function handleRefund(event, obj, tx) {
if (event.type !== 'refund.succeeded') return 'done';
// Check an existing immutable refund before deferring on a missing attempt.
// Reusing its refundId with an unknown/different Order is deterministic.
const preexistingRefund = obj.refundId
? await tx.refunds.findByRefundId(obj.refundId)
: null;
// A refund event carries only orderId, so find the payment attempt first.
const attempt = await tx.paymentAttempts.findByClinkOrderId(obj.orderId);
if (!attempt) {
if (!preexistingRefund) return 'deferred';
const merchantOrder = await tx.merchantOrders.findByIdForUpdate(
preexistingRefund.merchantOrderId,
);
if (!merchantOrder) throw new Error(`merchant order ${preexistingRefund.merchantOrderId} missing`);
return recordRefundReview(tx, {
eventId: event.id,
merchantOrderId: merchantOrder.id,
clinkOrderId: obj.orderId ?? null,
refundId: obj.refundId,
}, 'refund_id_reused_with_missing_attempt', {
stored: refundEvidence(preexistingRefund),
incoming: refundEvidence({
refundId: obj.refundId,
orderId: obj.orderId ?? null,
refundAmount: obj.refundAmount ?? null,
refundCurrency: obj.refundCurrency ?? null,
status: 'success',
}),
});
}
// Same lock, same lock order as the order events — otherwise the two
// paths deadlock against each other
const merchantOrder = await tx.merchantOrders.findByIdForUpdate(attempt.merchantOrderId);
if (!merchantOrder) return 'deferred';
const reviewContext = {
eventId: event.id,
merchantOrderId: merchantOrder.id,
clinkOrderId: attempt.clinkOrderId,
refundId: obj.refundId,
};
const incomingRefund = {
refundId: obj.refundId,
merchantOrderId: merchantOrder.id,
clinkOrderId: attempt.clinkOrderId,
amount: obj.refundAmount,
currency: obj.refundCurrency,
status: 'success',
};
const candidate = normalizeRefundCandidate(incomingRefund, event.id);
if (!candidate) {
return recordRefundReview(tx, reviewContext, 'invalid_refund_record', {
incoming: refundEvidence(incomingRefund),
});
}
const storedRefund = await tx.refunds.findByRefundId(obj.refundId);
if (storedRefund && !isExactRefundReplay(storedRefund, candidate)) {
return recordRefundReview(tx, reviewContext, 'refund_id_immutable_fields_conflict', {
stored: refundEvidence(storedRefund),
incoming: refundEvidence(candidate),
});
}
const basis = [
merchantOrder.refundablePaidClinkOrderId,
merchantOrder.refundablePaidAmount,
merchantOrder.refundablePaidCurrency,
];
if (basis.every((x) => x == null)) return 'deferred';
if (basis.some((x) => x == null)) {
return recordRefundReview(tx, reviewContext, 'stored_refund_basis_incomplete', {
storedBasis: basis,
incomingRefund: refundEvidence(incomingRefund),
});
}
const paidCurrency = tryNormalizeCurrency(merchantOrder.refundablePaidCurrency);
const refundCurrency = tryNormalizeCurrency(obj.refundCurrency);
if (
attempt.clinkOrderId !== merchantOrder.refundablePaidClinkOrderId ||
!paidCurrency ||
refundCurrency !== paidCurrency
) {
return recordRefundReview(tx, reviewContext, 'refund_order_or_currency_conflict', {
storedBasis: {
clinkOrderId: merchantOrder.refundablePaidClinkOrderId,
amount: merchantOrder.refundablePaidAmount,
currency: merchantOrder.refundablePaidCurrency,
},
storedRefund: refundEvidence(storedRefund),
incomingRefund: refundEvidence(incomingRefund),
});
}
const persisted = await persistRefundOnce(tx, incomingRefund, reviewContext);
if (persisted.outcome === 'manual_review') return persisted.outcome;
const refunds = await tx.refunds.listSuccessByMerchantOrder(merchantOrder.id);
const invalid = refunds.find((x) => !tryPositiveDecimal(x.amount));
const wrongOrder = refunds.find(
(x) => x.clinkOrderId !== merchantOrder.refundablePaidClinkOrderId,
);
const wrongCurrency = refunds.find(
(x) => tryNormalizeCurrency(x.currency) !== paidCurrency,
);
const conflicting = invalid ?? wrongOrder ?? wrongCurrency;
if (conflicting) {
return recordRefundReview(
tx,
{ ...reviewContext, refundId: conflicting.refundId },
'stored_successful_refund_conflict',
{ storedBasis: basis, conflictingRefund: refundEvidence(conflicting) },
);
}
const refundedAmount = refunds.reduce(
(sum, x) => sum.plus(tryPositiveDecimal(x.amount)),
new Decimal(0),
);
const refundablePaidAmount = tryPositiveDecimal(merchantOrder.refundablePaidAmount);
if (!refundablePaidAmount || refundedAmount.greaterThan(refundablePaidAmount)) {
return recordRefundReview(tx, reviewContext, 'refund_total_exceeds_payment_basis', {
refundablePaidAmount: merchantOrder.refundablePaidAmount,
refundedAmount: refundedAmount.toString(),
refundIds: refunds.map((x) => x.refundId),
});
}
await tx.merchantOrders.update(merchantOrder.id, { refundedAmount: refundedAmount.toString() });
// Do not write refunded here; the reconcile task settles the final status
// from the Order query
await tx.outbox.insert({
eventId: event.id,
task: 'refund_reconcile',
orderId: merchantOrder.id,
clinkOrderId: attempt.clinkOrderId,
payload: { refundId: obj.refundId },
dedupeKey: `refund_reconcile:${obj.refundId}`,
});
return 'done';
}
webhook_events 的唯一索引挡不住并发改同一个订单。 它只保证同一个 event.id 处理一次;order.succeeded 和另一笔 Order 的 order.failed 是两个不同事件,可以同时进来,各自读到旧快照再各自写回。所以要在动 attempts 之前先 SELECT ... FOR UPDATE 锁住商户订单那一行,把同一个订单的所有状态处理串行化。带状态条件的 UPDATE 可以作为第二层防御,但不能拿它替代行锁——它挡得住覆盖,挡不住基于旧快照算出的错误聚合结果。退款分支用同一把锁、同一个加锁顺序,否则两条路径交叉加锁会死锁。顺序定不下来时不要猜。 succeeded Attempt 对聚合结果具有最高优先级;除此以外,Session 已确认的当前 Attempt 即使
attemptCreatedAt 为空也可以决定状态。如果这两个安全锚点都不存在:- 任一 Attempt 尚未收到
order.created时,聚合状态回到pending;除了事件自己的reprocess_webhook,还要登记reconcile_session,让GET /checkout/session/{sessionId}主动确认当前 Order,不无限等待创建时间 - 两笔 Attempt 的已知、非空
attemptCreatedAt相同,也登记reconcile_session
order.created 只补该 Attempt 的创建时间,绝不回退已有状态。同一笔 Order 的两个状态事件 event.created 相同时同理,登记 reconcile_attempt,用 GET /order/{orderId} 的 status 收敛。Session 查询持续无结果时保持 pending,并在重试策略耗尽后升级处理。received_at、数据库自增 ID、Webhook 到达顺序、Clink Order ID 字典序都不能用来打破平局——它们和真实创建顺序无关。| 层 | 靠什么保证 | 管什么 |
|---|---|---|
| 事件幂等 | webhook_events.id 唯一索引 + 单事务 | 同一个 event.id 只产生一次副作用 |
| 单笔尝试收敛 | last_event_created 版本比较 + 终态不回退 | 一个 clink_order_id 的状态只能单调前进 |
| 商户订单聚合 | 遍历该订单的全部尝试 | 多笔 Order 之间谁说了算 |
order.next_action 只代表当前有效尝试。 它不能由一笔已经被替代的旧 Order 写回商户订单——上面 recomputeMerchantOrder 里,action_required 只能来自 Session 已确认的 current,或创建时间唯一最新时的 ordered[0];未解决的平局保持 pending。同理,某一笔尝试失败不代表订单失败。只要还有一笔成功,聚合结果就是成功。事件顺序或字段不足以判定时,以查询结果为准。
event.created 只能用来比较同一笔 Order 的事件先后。如果两个事件时间戳相同、或者拿不到可靠的版本字段,就按 obj.orderId 调 GET /order/{orderId} 查一次,用返回的 status 收敛这笔尝试——服务端查询是消除乱序歧义最稳妥的依据,前端事件和推送顺序都不是。这次查询可以放进 Outbox 任务里做,别在 Webhook 请求线程里同步调。dedupeKey 上建唯一索引,重复登记同一件事时靠数据库挡掉,不会堆出一堆同样的任务。
refund.succeeded 不等于全额退款。 部分退款成功也发这个事件,一律置成 refunded 会让只退了一部分的订单被当成全退。也不要在收到事件的当下立刻回查 Order。 退款服务是先更新退款记录并推事件,之后才异步去更新 Order 状态的。这一刻查 GET /order/{orderId},很可能读到的还是退款前的旧状态。另外,这个事件说的是退款在 Clink 侧处理成功,不代表钱已经回到客户账上。原路退回银行卡通常还要几个工作日,客服话术别写成”已到账”。Outbox Worker
Webhook 只负责在一个事务里记账并登记待办,真正干活的是一个独立的 Worker。 领取和执行必须分开:领取用一个短事务,外部调用在事务提交之后做。外部 API 慢的时候,长时间持着数据库行锁会拖垮整个连接池。 完整的表结构如下。默认值都给全了,INSERT 只写业务字段就能跑:
CREATE TABLE outbox (
id BIGSERIAL PRIMARY KEY,
task TEXT NOT NULL
CHECK (task IN (
'fulfill', 'refund_reconcile', 'reprocess_webhook',
'reconcile_attempt', 'reconcile_session',
'apply_refund_policy', 'manual_reconciliation',
'page_outbox_failure'
)),
event_id TEXT NOT NULL,
order_id TEXT, -- 商户订单号,Worker 要读
clink_order_id TEXT, -- Clink Order ID,退款对账要读
payload JSONB, -- 任务附带参数
dedupe_key TEXT NOT NULL, -- 防止重复登记同一件事
status TEXT NOT NULL DEFAULT 'pending'
CHECK (status IN ('pending', 'processing', 'succeeded', 'failed')),
attempt INT NOT NULL DEFAULT 0,
next_retry_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
claim_token TEXT, -- 每次领取生成的新随机值
worker_id TEXT,
lease_until TIMESTAMPTZ,
last_error TEXT,
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW()
);
-- 领取查询走这个索引
CREATE INDEX idx_outbox_claim ON outbox (status, next_retry_at);
-- 同一件事只登记一次
CREATE UNIQUE INDEX uq_outbox_dedupe ON outbox (dedupe_key);
status、attempt、next_retry_at 由默认值补齐:
INSERT INTO outbox
(task, event_id, order_id, clink_order_id, payload, dedupe_key)
VALUES
($1, $2, $3, $4, $5::jsonb, $6)
ON CONFLICT (dedupe_key) DO NOTHING;
ON CONFLICT DO NOTHING 让重复登记变成静默跳过,而不是抛异常。放在 Webhook 事务里时这点很重要——重复登记不该让整个事务回滚。两个 reconcile 任务不需要额外字段:
reconcile_attempt用clink_order_id就能调GET /order/{clinkOrderId}reconcile_session用order_id反查商户订单,再从订单行上取clink_session_idrefund_reconcile从payload读取refundId;累计refundedAmount要在商户订单锁内重新读取最新值apply_refund_policy从payload读取refundId、orderStatus、refundedAmount、paymentCurrencymanual_reconciliation从payload读取稳定的caseKey,加载不可变证据行,再以该 key 幂等创建或更新运营工单和告警page_outbox_failure只携带失败任务 ID/类型以及白名单化的错误名称/代码;它持久化重试告警,且绝不再创建另一条告警任务
payload 这个 JSONB 列,不要为每种任务单独加列。import crypto from 'node:crypto';
import Decimal from 'decimal.js';
import {
db,
fulfill,
applyRefundPolicy,
reconciliationDesk,
alerting,
logger,
} from './merchant-adapters.js';
const LEASE_MS = 5 * 60 * 1000; // Lease length; leave room for the slowest external call
const HEARTBEAT_MS = Math.floor(LEASE_MS / 3);
const MAX_ATTEMPT = 12;
const MAX_BACKOFF_MS = 60 * 60 * 1000;
function backoff(attempt, { now = Date.now(), random = Math.random } = {}) {
const numericAttempt = Number(attempt);
const normalizedAttempt = Number.isFinite(numericAttempt)
? Math.max(1, Math.floor(numericAttempt))
: 1;
const exponential = Math.min(
MAX_BACKOFF_MS,
1000 * (2 ** Math.min(normalizedAttempt - 1, 30)),
);
const randomValue = Math.min(1, Math.max(0, Number(random()) || 0));
const jittered = Math.min(
MAX_BACKOFF_MS,
Math.round(exponential * (0.75 + randomValue * 0.5)),
);
return new Date(Number(now) + jittered);
}
// Step 1: claim exactly one task atomically. A worker never lets a batch wait
// behind the first slow task while every lease counts down from the same instant.
async function claimTask(workerId) {
// A new token on every claim — never reuse the old one
const claimToken = crypto.randomUUID();
const { rows } = await db.query(
`UPDATE outbox SET
status = 'processing',
claim_token = $1,
worker_id = $2,
lease_until = NOW() + ($3 || ' milliseconds')::interval
WHERE id IN (
SELECT id FROM outbox
WHERE (status = 'pending' AND next_retry_at <= NOW())
OR (status = 'processing' AND lease_until < NOW()) -- expired lease, reclaim
ORDER BY id
LIMIT 1
FOR UPDATE SKIP LOCKED
)
RETURNING
id,
task,
event_id AS "eventId",
order_id AS "orderId",
clink_order_id AS "clinkOrderId",
payload,
dedupe_key AS "dedupeKey",
status,
attempt,
next_retry_at AS "nextRetryAt",
claim_token AS "claimToken",
worker_id AS "workerId",
lease_until AS "leaseUntil",
last_error AS "lastError",
created_at AS "createdAt",
updated_at AS "updatedAt"`,
[claimToken, workerId, LEASE_MS]
);
return rows[0] ?? null;
}
function startLeaseHeartbeat(t) {
let leaseLost = false;
let stopped = false;
let renewal = Promise.resolve(true);
const controller = new AbortController();
function markLeaseLost(reason) {
if (leaseLost) return;
leaseLost = true;
controller.abort(reason);
logger.warn({ taskId: t.id }, 'lease lost, stop claimed task');
}
function ensureOwned() {
renewal = renewal.then(async () => {
if (stopped || leaseLost) return false;
try {
const owned = await renewLease(t);
if (!owned) markLeaseLost(new Error('outbox lease lost'));
return owned;
} catch (err) {
// A renewal error means ownership is unknown. Stop rather than risk a
// side effect under an expired token.
markLeaseLost(err);
return false;
}
});
return renewal;
}
const timer = setInterval(() => { void ensureOwned(); }, HEARTBEAT_MS);
timer.unref?.();
return {
signal: controller.signal,
isLost: () => leaseLost,
ensureOwned,
async stop() {
stopped = true;
clearInterval(timer);
await renewal;
},
};
}
function throwIfAborted(signal) {
if (signal?.aborted) throw signal.reason ?? new Error('operation aborted');
}
async function runOutboxWorker(workerId) {
// The claim statement commits here; the row lock is released with it.
const t = await claimTask(workerId);
if (!t) return;
// A task may have waited between claim and execution. Confirm ownership
// before starting any API call or external side effect.
if (!await renewLease(t)) {
logger.warn({ taskId: t.id }, 'lease lost before task start');
return;
}
const lease = startLeaseHeartbeat(t);
try {
await executeClaimedTask(t, lease);
if (lease.isLost() || !await lease.ensureOwned()) return;
await finishTask(t, 'succeeded');
} catch (err) {
if (lease.isLost() || !await lease.ensureOwned()) return;
await failTask(t, err);
} finally {
await lease.stop();
}
}
// Step 2: external calls happen outside row-lock transactions. Every adapter
// receives its stable business idempotency key and AbortSignal where supported.
async function executeClaimedTask(t, lease, adapters = {
fulfill,
reconcileRefund,
applyRefundPolicy,
reprocessWebhook,
reconcileAttempt,
reconcileSession,
reconciliationDesk,
alerting,
findReconciliationCase: (caseKey) =>
db.reconciliationCases.findByDedupeKey(caseKey),
}) {
if (t.task === 'fulfill') {
await adapters.fulfill(t.orderId, {
eventId: t.eventId,
idempotencyKey: t.orderId,
signal: lease.signal,
});
} else if (t.task === 'refund_reconcile') {
await adapters.reconcileRefund(t.orderId, t.clinkOrderId, t.eventId, t.payload, {
signal: lease.signal,
});
} else if (t.task === 'apply_refund_policy') {
// Deduplicate on refundId, and only advance the policy's stored cumulative
// amount so an older task cannot roll entitlement backward.
await adapters.applyRefundPolicy({
merchantOrderId: t.orderId,
refundId: t.payload.refundId,
orderStatus: t.payload.orderStatus,
refundedAmount: t.payload.refundedAmount,
paymentCurrency: t.payload.paymentCurrency,
idempotencyKey: t.payload.refundId,
signal: lease.signal,
});
} else if (t.task === 'manual_reconciliation') {
const review = await adapters.findReconciliationCase(t.payload.caseKey);
if (!review) throw new Error(`manual-review case ${t.payload.caseKey} missing`);
const reviewContext = {
caseKey: review.dedupeKey,
eventId: review.eventId,
merchantOrderId: review.merchantOrderId,
clinkOrderId: review.clinkOrderId ?? null,
refundId: review.refundId ?? null,
reason: review.reason,
evidence: review.evidence,
};
await adapters.reconciliationDesk.openOrUpdate({
...reviewContext,
externalId: t.payload.caseKey,
signal: lease.signal,
});
// Desk and paging are separate external steps. Re-check the token so a
// lost worker never starts the second step.
if (!await lease.ensureOwned()) return;
await adapters.alerting.openOrUpdate({
dedupeKey: t.payload.caseKey,
title: 'payment data requires manual reconciliation',
context: reviewContext,
signal: lease.signal,
});
} else if (t.task === 'reprocess_webhook') {
await adapters.reprocessWebhook(t.eventId, { signal: lease.signal });
} else if (t.task === 'reconcile_attempt') {
await adapters.reconcileAttempt(t.clinkOrderId, t.eventId, { signal: lease.signal });
} else if (t.task === 'reconcile_session') {
await adapters.reconcileSession(t.orderId, t.eventId, { signal: lease.signal });
} else if (t.task === 'page_outbox_failure') {
await adapters.alerting.openOrUpdate({
dedupeKey: t.dedupeKey,
title: 'outbox task exhausted',
context: {
failedTaskId: t.payload.failedTaskId,
failedTask: t.payload.failedTask,
errorName: t.payload.errorName,
errorCode: t.payload.errorCode,
},
signal: lease.signal,
});
} else {
// Never mark an unimplemented task succeeded.
throw new Error(`unknown outbox task: ${t.task}`);
}
}
// Raw SQL is snake_case, but claimTask aliases every returned task property to
// camelCase. Worker code must use that one shape exclusively.
// Step 3: every status update must carry claim_token and status='processing'
async function finishTask(t, status) {
const { rowCount } = await db.query(
`UPDATE outbox SET
status = $1, claim_token = NULL, worker_id = NULL, lease_until = NULL
WHERE id = $2 AND claim_token = $3 AND status = 'processing'`,
[status, t.id, t.claimToken]
);
if (rowCount === 0) {
// The lease expired and another worker owns this task now. Let go
logger.warn({ taskId: t.id }, 'lease lost, result discarded');
}
return rowCount;
}
async function failTask(t, err) {
const attempt = t.attempt + 1;
// Paging failures retry forever (with a capped delay) and never recursively
// create another page_outbox_failure task.
const exhausted = t.task !== 'page_outbox_failure' && attempt >= MAX_ATTEMPT;
const errorName = typeof err?.name === 'string' ? err.name.slice(0, 100) : 'Error';
const errorCode = ['string', 'number'].includes(typeof err?.code)
? String(err.code).slice(0, 100)
: null;
const safeLastError = errorCode ? `${errorName}:${errorCode}` : errorName;
if (!exhausted) {
const { rowCount } = await db.query(
`UPDATE outbox SET
status = 'pending', claim_token = NULL, worker_id = NULL, lease_until = NULL,
attempt = $1, next_retry_at = $2, last_error = $3
WHERE id = $4 AND claim_token = $5 AND status = 'processing'`,
[attempt, backoff(attempt), safeLastError, t.id, t.claimToken]
);
if (rowCount === 0) logger.warn({ taskId: t.id }, 'lease lost, result discarded');
return rowCount;
}
// Marking the original task failed and creating its paging task is one
// transaction. A crash after commit still leaves durable work to claim.
const rowCount = await db.transaction(async (tx) => {
const result = await tx.query(
`UPDATE outbox SET
status = 'failed', claim_token = NULL, worker_id = NULL, lease_until = NULL,
attempt = $1, next_retry_at = $2, last_error = $3
WHERE id = $4 AND claim_token = $5 AND status = 'processing'`,
[attempt, backoff(attempt), safeLastError, t.id, t.claimToken]
);
if (result.rowCount === 0) return 0;
await tx.outbox.insert({
eventId: t.eventId,
task: 'page_outbox_failure',
payload: {
failedTaskId: t.id,
failedTask: t.task,
errorName,
errorCode,
},
dedupeKey: `outbox_failure_alert:${t.id}`,
});
return result.rowCount;
});
if (rowCount === 0) logger.warn({ taskId: t.id }, 'lease lost, result discarded');
return rowCount;
}
// Renewal carries the token too, or it would extend someone else's lease
async function renewLease(t) {
const { rowCount } = await db.query(
`UPDATE outbox SET lease_until = NOW() + ($1 || ' milliseconds')::interval
WHERE id = $2 AND claim_token = $3 AND status = 'processing'`,
[LEASE_MS, t.id, t.claimToken]
);
return rowCount === 1; // false means the lease is gone — stop work immediately
}
// Retry an event that was deferred until its dependency existed
async function reprocessWebhook(eventId, { signal } = {}) {
throwIfAborted(signal);
const outcome = await db.transaction(async (tx) => {
const row = await tx.webhookEvents.findById(eventId);
if (!row || row.status === 'processed') return 'done';
const outcome = await handleEvent(JSON.parse(row.payload), tx);
if (outcome === 'done' || outcome === 'manual_review') {
await tx.webhookEvents.update(eventId, { status: 'processed' });
}
throwIfAborted(signal); // abort rolls the transaction back before commit
return outcome;
});
// Throw only after the transaction commits. Any reconciliation task queued
// by handleEvent stays durable while this task goes through backoff.
if (outcome === 'deferred') throw new Error('dependency still missing');
}
// Two status events on one Order shared a timestamp — converge from the API
async function reconcileAttempt(clinkOrderId, eventId, { signal } = {}) {
throwIfAborted(signal);
// Capture the local version before the external query.
const snapshot = await db.paymentAttempts.findByClinkOrderId(clinkOrderId);
if (!snapshot) throw new Error(`attempt ${clinkOrderId} missing`);
// Query first, then open the transaction. External calls never run under a lock.
const order = await clinkGet(`/order/${clinkOrderId}`, { signal });
throwIfAborted(signal);
const outcome = await db.transaction(async (tx) => {
// Same lock, same order as the webhook handler. Re-read the attempt only
// after taking this lock, because every local attempt write takes it first.
const merchantOrder = await tx.merchantOrders.findByIdForUpdate(snapshot.merchantOrderId);
if (!merchantOrder) throw new Error(`merchant order ${snapshot.merchantOrderId} missing`);
const attempt = await tx.paymentAttempts.findByClinkOrderId(clinkOrderId);
if (
!attempt ||
merchantOrder.clinkSessionId !== snapshot.clinkSessionId ||
attempt.merchantOrderId !== snapshot.merchantOrderId ||
attempt.clinkSessionId !== snapshot.clinkSessionId ||
attempt.status !== snapshot.status ||
attempt.lastEventCreated !== snapshot.lastEventCreated ||
attempt.attemptCreatedAt !== snapshot.attemptCreatedAt
) {
// Local state advanced while Clink was being queried. Roll back and let
// the worker retry from a fresh snapshot; never apply this stale result.
throw new Error(`attempt ${clinkOrderId} changed during reconciliation`);
}
if (
order.orderId !== clinkOrderId ||
order.sessionId !== attempt.clinkSessionId ||
order.merchantReferenceId !== merchantOrder.id
) {
const manualOutcome = await recordManualReview(tx, {
dedupeKey: `order_query_ownership_conflict:${merchantOrder.id}:${clinkOrderId}`,
eventId,
merchantOrderId: merchantOrder.id,
clinkOrderId,
reason: 'order_query_ownership_conflict',
evidence: {
requestedOrderId: clinkOrderId,
returnedOrder: orderEvidence(order),
localAttempt: attemptEvidence(attempt),
storedBasis: {
clinkOrderId: merchantOrder.refundablePaidClinkOrderId,
amount: merchantOrder.refundablePaidAmount,
currency: merchantOrder.refundablePaidCurrency,
},
},
});
await tx.webhookEvents.update(eventId, { status: 'processed' });
return manualOutcome;
}
// An already-refunded order takes the refund path, never the
// "paid -> queue fulfillment" one
let aggregateOutcome = 'done';
if (REFUND_STATUS.has(order.status)) {
await tx.paymentAttempts.update(clinkOrderId, {
status: 'succeeded',
failureCode: null,
failureMessage: null,
});
} else {
const mapped = ORDER_STATUS_TO_ATTEMPT[order.status];
if (!mapped) {
// Still pending, or an unknown value — not settled yet
throw new Error(`order ${clinkOrderId} not settled yet: ${order.status}`);
}
await tx.paymentAttempts.update(clinkOrderId, {
status: mapped,
failureCode: mapped === 'failed' ? order.failureCode : null,
failureMessage: mapped === 'failed' ? order.failureMessage : null,
});
}
if (order.status === 'success' || REFUND_STATUS.has(order.status)) {
// GET /order exposes amountTotal and paymentCurrency for Hosted Checkout,
// so this path can establish the same basis as order.succeeded. The
// attempt update above remains durable if the basis is quarantined.
const basisOutcome = await persistRefundablePaymentBasis(
tx,
merchantOrder,
{
clinkOrderId,
amountTotal: order.amountTotal,
paymentCurrency: order.paymentCurrency,
},
{ eventId },
);
if (basisOutcome === 'manual_review') {
await tx.webhookEvents.update(eventId, { status: 'processed' });
return basisOutcome;
}
}
if (REFUND_STATUS.has(order.status)) {
await advanceRefundStatus(tx, merchantOrder, order.status);
} else {
aggregateOutcome = await recomputeMerchantOrder(tx, merchantOrder, eventId);
}
// Close the event only after both the attempt and merchant-order aggregate
// have converged. A Session tie remains pending.
if (aggregateOutcome === 'done' || aggregateOutcome === 'manual_review') {
await tx.webhookEvents.update(eventId, { status: 'processed' });
}
return aggregateOutcome;
});
if (outcome === 'deferred') throw new Error('merchant order still needs Session reconciliation');
}
const ORDER_STATUS_TO_ATTEMPT = {
success: 'succeeded',
failed: 'failed',
requires_action: 'action_required',
};
// Refund states are handled separately and never feed the payment aggregation
const REFUND_STATUS = new Set(['partial_refunded', 'refunded']);
// Refund state is monotonic: refunded is never overwritten by partial_refunded
const REFUND_RANK = { partial_refunded: 1, refunded: 2 };
async function advanceRefundStatus(tx, merchantOrder, orderStatus) {
const curRank = REFUND_RANK[merchantOrder.paymentStatus] ?? 0;
if (REFUND_RANK[orderStatus] <= curRank) return; // forward only
await tx.merchantOrders.updateIfStatus(
merchantOrder.id,
merchantOrder.paymentStatus,
{ paymentStatus: orderStatus },
);
}
async function queueRefundPolicy(tx, merchantOrder, refund) {
// Queue once per successful refund, even when the Order remains
// partial_refunded. State rank and policy evaluation are separate concerns.
await tx.outbox.insert({
eventId: refund.eventId,
task: 'apply_refund_policy',
orderId: merchantOrder.id,
payload: {
refundId: refund.refundId,
orderStatus: refund.orderStatus,
refundedAmount: refund.refundedAmount,
paymentCurrency: refund.paymentCurrency,
},
dedupeKey: `apply_refund_policy:${refund.refundId}`,
});
}
// Stable comparison only. Lexical Order ID sorting here does not choose the
// current attempt; it merely makes identical local snapshots hash the same way.
function fingerprintAttempts(attempts) {
return attempts
.map((x) => JSON.stringify([
x.clinkOrderId,
x.merchantOrderId,
x.clinkSessionId,
x.status,
x.attemptCreatedAt ?? null,
x.lastEventCreated,
]))
.sort()
.join('|');
}
// Two attempts share an attemptCreatedAt — ask the Session which one is current
async function reconcileSession(merchantOrderId, eventId, { signal } = {}) {
throwIfAborted(signal);
const local = await db.merchantOrders.findById(merchantOrderId);
if (!local) throw new Error(`merchant order ${merchantOrderId} missing`);
const localAttempts = await db.paymentAttempts.listByMerchantOrder(merchantOrderId);
const snapshot = {
clinkSessionId: local.clinkSessionId,
currentClinkOrderId: local.currentClinkOrderId,
currentAttemptConfirmed: local.currentAttemptConfirmed,
attemptsFingerprint: fingerprintAttempts(localAttempts),
};
const session = await clinkGet(`/checkout/session/${snapshot.clinkSessionId}`, { signal });
throwIfAborted(signal);
await db.transaction(async (tx) => {
const merchantOrder = await tx.merchantOrders.findByIdForUpdate(merchantOrderId);
if (!merchantOrder) throw new Error(`merchant order ${merchantOrderId} missing`);
const lockedAttempts = await tx.paymentAttempts.listByMerchantOrder(merchantOrderId);
if (
merchantOrder.clinkSessionId !== snapshot.clinkSessionId ||
merchantOrder.currentClinkOrderId !== snapshot.currentClinkOrderId ||
merchantOrder.currentAttemptConfirmed !== snapshot.currentAttemptConfirmed ||
fingerprintAttempts(lockedAttempts) !== snapshot.attemptsFingerprint
) {
// Another Order event or reconciliation changed local truth after the
// query began. Discard the stale Session answer and retry.
throw new Error(`attempt set for ${merchantOrderId} changed during Session reconciliation`);
}
if (
session.sessionId !== snapshot.clinkSessionId ||
session.merchantReferenceId !== merchantOrder.id
) {
const manualOutcome = await recordManualReview(tx, {
dedupeKey: `session_response_conflict:${merchantOrderId}:${snapshot.clinkSessionId}`,
eventId,
merchantOrderId,
clinkOrderId: session.orderId ?? null,
reason: 'session_response_ownership_conflict',
evidence: {
querySnapshot: snapshot,
returnedSession: sessionEvidence(session),
},
});
await tx.webhookEvents.update(eventId, { status: 'processed' });
return manualOutcome;
}
if (!session.orderId) {
// No Order yet, or the state is unsettled — back off and retry.
// Never treat this as success.
throw new Error(`session ${snapshot.clinkSessionId} has no current orderId yet`);
}
// Read by global Order ID inside the lock. Missing is retryable; an Order
// already owned by another merchant/session is a durable manual review.
const attempt = await tx.paymentAttempts.findByClinkOrderId(session.orderId);
if (!attempt) {
throw new Error(`session order ${session.orderId} not persisted for ${merchantOrderId}`);
}
if (
attempt.merchantOrderId !== merchantOrderId ||
attempt.clinkSessionId !== snapshot.clinkSessionId
) {
const manualOutcome = await recordManualReview(tx, {
dedupeKey: `session_order_conflict:${merchantOrderId}:${session.orderId}`,
eventId,
merchantOrderId,
clinkOrderId: session.orderId,
reason: 'session_order_ownership_conflict',
evidence: {
querySnapshot: snapshot,
returnedSession: sessionEvidence(session),
storedAttempt: attemptEvidence(attempt),
storedBasis: {
clinkOrderId: merchantOrder.refundablePaidClinkOrderId,
amount: merchantOrder.refundablePaidAmount,
currency: merchantOrder.refundablePaidCurrency,
},
},
});
await tx.webhookEvents.update(eventId, { status: 'processed' });
return manualOutcome;
}
// Persist both the Session answer and its provenance. Only this confirmed
// pointer may break an attemptCreatedAt tie.
await tx.merchantOrders.update(merchantOrderId, {
currentClinkOrderId: session.orderId,
currentAttemptConfirmed: true,
});
merchantOrder.currentClinkOrderId = session.orderId;
merchantOrder.currentAttemptConfirmed = true;
const aggregateOutcome = await recomputeMerchantOrder(tx, merchantOrder, eventId);
if (aggregateOutcome === 'deferred') {
throw new Error(`session ${snapshot.clinkSessionId} did not resolve the current attempt`);
}
await tx.webhookEvents.update(eventId, { status: 'processed' });
return aggregateOutcome;
});
}
async function reconcileRefund(
merchantOrderId,
clinkOrderId,
eventId,
refund,
{ signal } = {},
) {
throwIfAborted(signal);
if (!refund?.refundId) {
throw new Error('refund reconciliation payload is incomplete');
}
// The query runs outside the transaction
const order = await clinkGet(`/order/${clinkOrderId}`, { signal });
throwIfAborted(signal);
// Only once a result is in hand: open the transaction and take the same
// lock, in the same order, as every other path
await db.transaction(async (tx) => {
const merchantOrder = await tx.merchantOrders.findByIdForUpdate(merchantOrderId);
if (!merchantOrder) throw new Error(`merchant order ${merchantOrderId} missing`);
const reviewContext = { eventId, merchantOrderId, clinkOrderId, refundId: refund.refundId };
const storedAttempt = await tx.paymentAttempts.findByClinkOrderId(clinkOrderId);
if (
!storedAttempt ||
storedAttempt.merchantOrderId !== merchantOrder.id ||
storedAttempt.clinkSessionId !== merchantOrder.clinkSessionId ||
order.orderId !== clinkOrderId ||
order.sessionId !== merchantOrder.clinkSessionId ||
order.merchantReferenceId !== merchantOrder.id
) {
const storedRefund = await tx.refunds.findByRefundId(refund.refundId);
return recordRefundReview(tx, reviewContext, 'refund_order_query_ownership_conflict', {
requestedOrderId: clinkOrderId,
returnedOrder: orderEvidence(order),
storedAttempt: attemptEvidence(storedAttempt),
storedRefund: refundEvidence(storedRefund),
storedBasis: {
clinkOrderId: merchantOrder.refundablePaidClinkOrderId,
amount: merchantOrder.refundablePaidAmount,
currency: merchantOrder.refundablePaidCurrency,
},
});
}
// Ownership is valid but the refund state has not settled yet. Back off
// instead of reading a pre-refund status as a completed reconciliation.
if (!REFUND_STATUS.has(order.status)) {
throw new Error('order status not settled yet');
}
const basisOutcome = await persistRefundablePaymentBasis(
tx,
merchantOrder,
{
clinkOrderId,
amountTotal: order.amountTotal,
paymentCurrency: order.paymentCurrency,
},
{
eventId,
refundId: refund.refundId,
dedupeKey: `refund_conflict:${refund.refundId}`,
},
);
if (basisOutcome === 'manual_review') return basisOutcome;
if (
merchantOrder.refundablePaidClinkOrderId == null ||
merchantOrder.refundablePaidAmount == null ||
merchantOrder.refundablePaidCurrency == null
) {
return recordRefundReview(tx, reviewContext, 'refund_payment_basis_missing', {
storedBasis: {
clinkOrderId: merchantOrder.refundablePaidClinkOrderId,
amount: merchantOrder.refundablePaidAmount,
currency: merchantOrder.refundablePaidCurrency,
},
});
}
if (merchantOrder.refundablePaidClinkOrderId !== clinkOrderId) {
return recordRefundReview(tx, reviewContext, 'refund_reconciliation_order_conflict', {
storedOrderId: merchantOrder.refundablePaidClinkOrderId,
incomingOrderId: clinkOrderId,
});
}
// Compute the minimum state the Order API must have reached from the
// latest locked totals. Decimal avoids rounding a full refund down into a
// partial one. A lagging API result must retry; there is no later
// order.refunded event that can repair a premature success.
const refundedAmount = tryPositiveDecimal(merchantOrder.refundedAmount);
const refundablePaidAmount = tryPositiveDecimal(merchantOrder.refundablePaidAmount);
const paymentCurrency = tryNormalizeCurrency(merchantOrder.refundablePaidCurrency);
if (!refundedAmount || !refundablePaidAmount || !paymentCurrency) {
return recordRefundReview(tx, reviewContext, 'invalid_locked_refund_totals', {
refundedAmount: merchantOrder.refundedAmount,
refundablePaidAmount: merchantOrder.refundablePaidAmount,
paymentCurrency: merchantOrder.refundablePaidCurrency,
});
}
if (refundedAmount.greaterThan(refundablePaidAmount)) {
return recordRefundReview(tx, reviewContext, 'refund_total_exceeds_payment_basis', {
refundedAmount: refundedAmount.toString(),
refundablePaidAmount: refundablePaidAmount.toString(),
});
}
const expectedStatus = refundedAmount.greaterThanOrEqualTo(refundablePaidAmount)
? 'refunded'
: 'partial_refunded';
if (REFUND_RANK[order.status] < REFUND_RANK[expectedStatus]) {
throw new Error('order refund status has not caught up yet');
}
// Status advances only when its rank increases, but every refundId gets a
// policy task so repeated partial refunds are never collapsed together.
await advanceRefundStatus(tx, merchantOrder, order.status);
await queueRefundPolicy(tx, merchantOrder, {
eventId,
refundId: refund.refundId,
orderStatus: order.status,
// Use the latest value read under the lock, never the task's old snapshot.
refundedAmount: refundedAmount.toString(),
paymentCurrency,
});
});
}
pending ──领取(新 token)──> processing ──成功──> succeeded
↑ │
├──失败(未达上限)─────────────┤
│ └──达到上限──> failed + pending 告警任务
└──租约过期,被别的 Worker 重新领取
重试耗尽和告警任务必须原子提交。 带 claim token 条件的原任务
failed 更新,与去重的 page_outbox_failure 插入放在同一个数据库事务。告警 payload 只能包含 failedTaskId、failedTask、errorName、errorCode;不得复制异常消息、Webhook/API 对象、HTTP body、客户数据或支付详情。page_outbox_failure 必须等待告警适配器完成,投递成功后才能标为 succeeded。投递失败时,同一任务按带上限的退避回到 pending,不得递归创建另一条告警任务。告警适配器以 outbox_failure_alert:{failedTaskId} 作为外部幂等键。基础设施仍要监控:任何 status = 'failed' 的任务、长时间停在 pending/processing 的任务,以及长期未成功的 page_outbox_failure。持久化重试能防止进程崩溃后静默丢告警,但不能替代队列健康监控。SKIP LOCKED 本身不是幂等保证。 它只做一件事:让并发的 SELECT 跳过别人已经锁住的行,避免互相阻塞。真正保证”一条任务只被一个 Worker 领走”的,是在同一条语句里把状态改成 processing 并写入新的 claim_token。上面用的是 UPDATE ... RETURNING,锁定和标记一次完成,中间没有别的 Worker 能插进来。写成先 SELECT ... FOR UPDATE SKIP LOCKED、事务提交后再 UPDATE 的两段式,就会留出窗口。真要分两步,也必须把两步放进同一个短事务里。每次更新任务状态都要带
claim_token。 这一条防的是租约过期后的状态覆盖:- Worker A 领到任务,拿到
tokenA - A 卡住了,租约到期
- Worker B 用
tokenB重新领走同一条任务 - A 恢复过来,上报成功——但它的
WHERE claim_token = tokenA匹配不到任何行,影响行数为 0 - B 正常完成,状态不会被 A 覆盖
Worker 崩了怎么办。 进程如果在领取之后、更新状态之前挂掉,任务会一直卡在
processing。lease_until 就是为这个准备的:领取条件里带上 status = 'processing' AND lease_until < NOW(),租约到期的任务会被下一个 Worker 重新领走。本示例每次只领取一条任务。启动任何外部副作用前先续租一次,执行期间每隔 LEASE_MS / 3 心跳续租。续租失败会把本地标记为失租,取消支持 AbortSignal 的 HTTP 调用,跳过所有后续步骤,也不再调用 finishTask / failTask。人工对账路径在建工单和发告警之间还会再次确认租约。业务幂等仍然必须保留,但它只是最后一道保险,绝不能成为明知租约已失仍启动任务的理由。worker_id 和 lease_until 也是排查用的——能看出哪台机器卡住了、卡了多久。两个 reconcile 任务都必须真的发出查询,不能空跑。
reconcile_attempt 调 GET /order/{clinkOrderId},把返回的 status 映射成尝试状态;reconcile_session 调 GET /checkout/session/{sessionId},用返回的 orderId 认定当前尝试。两者的查询都在事务外执行。查询前先保存本地版本:reconcile_attempt 保存 attempt owner、Session ID、status、lastEventCreated、attemptCreatedAt;reconcile_session 保存 Session ID、当前指针/确认标记和全部 attempts 的归属及状态指纹。拿到结果后才进事务锁商户订单——和 Webhook 处理器同一把锁、同一加锁顺序——然后重读并比较快照。任一字段不一致都说明 API 结果可能已过期,必须回滚重试,不能写回。结果还不稳定时要抛错走退避重试,绝不能因为”这次没查出来”就把任务标成 succeeded:GET /order 返回的状态不在映射表里、或 GET /checkout/session 还没有 orderId,都属于这种情况。当前尝试落在 merchant_orders.current_clink_order_id,current_attempt_confirmed 则记录它是否来自 Session 查询。由 order.created 推导出的指针不能打破时间平局。不要去改 attempt_created_at——那是 Clink 给的真实创建时间,改了就污染了原始数据,而且「调成组内最大值」也未必真能解除相等。写入前要在锁内按全局 orderId 查询 Session 返回的 Order。查不到 attempt,说明那笔 Order 还没落到本地,应抛错重试;如果 attempt 已存在却属于另一商户或 Session,这是确定性冲突,应先持久化人工对账证据、把事件置 processed,再停止自动重试。租约机制保证的是”不会永久卡死”,不是”只执行一次”。 崩溃或租约过期都会导致任务被重新领取并再跑一遍,所以
fulfill()、reconcileRefund()、applyRefundPolicy()、reprocessWebhook() 自身必须用业务幂等键保护——发货按订单号、退款记账和退款策略都按 refundId、事件处理按 event.id。退款策略还必须比较累计 refundedAmount(或单调退款版本),迟到的旧任务不能把权益状态回退。外部接口若支持幂等键,调用时一并带上。clinkGet 是前面共享服务端 helper 的无 body GET 包装。它只接受受控相对路径,每次生成新的毫秒 X-Timestamp,透传 AbortSignal,同时要求 HTTP 成功和 { code: 200, data: object } 信封;失败时只抛脱敏错误,不复制 Secret Key 或响应 body。
不想做退款对账的话,还有个更简单的口径:自己按 refundId 汇总成功退款金额,用 Decimal 和 refundable_paid_amount 比,退满了就是全退。Clink 没有”该订单累计已退金额”这个字段,也没有 order.refunded 事件,所以这个汇总本来就要自行维护。
别在 Webhook 请求线程里
await 这个 Worker。一旦在返回 200 之前执行外部动作,发货失败就会导致返回非 2xx,Clink 重投——而重投会被 event.id 去重直接挡回 200,这单货就永远不会再发了。事务提交即返回 200,外部动作全部交给 Worker 重试,才不会出现这个死角。Outbox 表至少要有
status、attempt、next_retry_at、last_error 四个字段。Worker 幂等键用 eventId 或订单号,因为它一定会被重跑。发货、发邮件、调第三方接口这些动作出了数据库事务就回滚不了,所以不能写在事务里;反过来,enqueueXxx() 这种在事务里调外部队列的写法同样不行——事务回滚了任务却已经发出去,就成了没人认领的孤儿任务。待办必须和业务状态写进同一个事务。权益要不要收回,不该由支付状态直接决定——退一半钱是停服、按比例缩短还是保持不变,取决于所售商品的性质。把这个判断收在
applyRefundPolicy 这类地方,别散在 Webhook 处理器里。Order 状态只在 rank 上升时推进,但每个成功的 refundId 都要在对账事务内登记自己的 apply_refund_policy 任务。提交后再由 Worker 执行;持有商户订单行锁时绝不能直接调用。匹配订单、终态保护、登记履约这三件事要对所有
order.* 事件生效,不能只写在成功分支里。常见的写法是成功事件做了完整校验,失败事件却直接按订单号改成失败——一个迟到的 order.failed 就能把已经收到的钱标记成失败,货也退了。refund.* 的对象里没有 merchantReferenceId,只有 orderId 和 refundMerchantOrderId。所以要用付款时存下的 clinkOrderId 反查本地订单。
这五件事一件都不能少
1
验签
用原始请求体算 HMAC SHA-256,比对
X-Clink-Signature。签名对不上就返回 401。2
按事件 ID 去重
投递失败 Clink 会重试,同一个事件会收到多遍。在
event.id 上建唯一索引,把这条记录和业务写入放进同一个事务,靠数据库挡住并发。先查再写挡不住同时到达的两份。3
双重匹配订单
同时校验
merchantReferenceId 和 sessionId。只对上一个就当异常处理,别更新状态。4
容忍乱序
Clink 不保证事件按发生顺序到达,要分两种情况处理。旧事件后到:已经是
paid、refunded 或 partial_refunded 的订单,不能被后到的旧事件改回 pending。同一笔 Clink Order 内部,用 event.created 和该尝试的 last_event_created 比较,旧的直接丢弃。依赖事件还没到:比如 refund.succeeded 早于 order.succeeded 到达,此时本地订单还不存在。这种事件要保持 pending 并登记重处理任务,绝不能标成 processed——一旦标了,Clink 不会再推,这笔退款就永远丢了。5
写库成功后再返回 2xx
返回 200 表示”我处理完了”。还没落库就返回 200,这个事件就永远丢了。
这套机制要验的场景
上线前把这几种情况跑一遍,光看正常顺序是测不出问题的。 支付尝试顺序| 场景 | 期望结果 |
|---|---|
非成功 Attempt 尚未收到 order.created,且没有已确认 current | 保存其状态,把旧的非终态聚合安全收敛到 pending,各登记一条 reconcile_session 和 reprocess_webhook,事件保持 pending |
order.succeeded 先于 order.created 到达 | 成功具有全局优先级:写退款基准、置 paid、完成事件并只登记一条履约;后到 created 只补时间 |
order.created 迟到,且该 attempt 已是 succeeded | 只补 attemptCreatedAt,状态不回退到 pending |
已知 Order A 失败后,未排序 Order B 收到 next_action | B 状态落库,旧 payment_failed 聚合回到 pending,主动查询 Session 确认当前 Order |
| Session 对账返回未排序的 B | 把 B 持久化为已确认 current;按 B 收敛到 action_required 或 payment_failed,并完成原 Webhook |
| 已确认 current 之后又出现新的未排序 Order | 让 currentAttemptConfirmed 失效,回到 pending 并重新查询 Session |
Order A 的 order.created 先于 B 到达,且两笔 attemptCreatedAt 相同 | 将推导指针标为未确认,登记 reconcile_session,由 Session 决定,不按接收顺序猜 |
同一 Order 两个状态事件 event.created 相同 | 登记 reconcile_attempt,查 Order 收敛,不按到达顺序定 |
reconcile_attempt 先收敛为 failed,后续查询又收敛到非失败状态 | failed 时写入失败码和原因,后续状态清空旧失败信息,不残留过期原因 |
| Order A 建立退款基准后,Order B 成功但 Order ID、金额或币种不同 | 不覆盖旧基准;B 的 attempt 与冲突证据一并提交,Webhook 正常确认,语义重放只产生一张人工对账 case 和一条任务 |
Order 事件的 merchantReferenceId 正确,但 Session ID 属于别的 Session | 不把 attempt 关联到订单;用一张持久化人工对账 case/任务保存两个 Session ID,并正常确认 Webhook |
| 场景 | 期望结果 |
|---|---|
| 商户订单 A / Session A 第一次收到一个新 Clink Order | 只插入一条归属 A / Session A 的 attempt,然后继续正常聚合 |
| 同一 Order、同一商户订单、同一 Session 重放 | 不可变归属不变,重放保持幂等 |
已有 attempt 分别收到 order.created、order.next_action、order.succeeded、order.failed | 任何字段更新前都同时校验商户订单与 Session 归属 |
order_X 已属于 A / Session A,B / Session B 又发送它的 order.created | A 的状态和创建时间不变;B 不写当前指针;只产生一张 case 和一条任务 |
上述跨归属冲突改为 order.next_action、order.failed 或 order.succeeded | A 的状态、失败信息、创建时间不变;B 不写退款基准、不履约、不聚合 |
| 商户订单相同但 Session 不同 | 事件持久化为 manual_review,attempt 不修改 |
| Session 相同但商户订单不同 | 事件持久化为 manual_review,attempt 不修改 |
同一归属冲突以相同或不同 event.id 重放 | 稳定键 order_ownership_conflict:{incomingMerchantOrderId}:{clinkOrderId} 只保留一张业务 case 和一条 manual_reconciliation Outbox |
A 与 B 并发首次插入同一 clinkOrderId | 主键只选出一个 owner;失败方重读胜出行、记录 manual_review 并正常确认,不返回通用 500 |
同一 owner 的两个事件并发首次插入同一 clinkOrderId | 不产生错误人工工单;最终仍遵守 event.created 顺序和终态规则 |
| 并发归属冲突后,胜出 owner 继续发送正常事件 | 仍能正常处理,不锁定也不修改失败方商户订单 |
归属冲突的 order.succeeded 指向别处已有 attempt | 不写入或覆盖 refundablePaidClinkOrderId、金额和币种 |
| 任一归属冲突被隔离 | 不会从该冲突事件生成 fulfill、refund_reconcile 或 reconcile_session 任务 |
| 写 case、Outbox 或 processed 标记的人工对账事务中途失败 | case、任务、attempt 证据和 webhook_events.processed 变更一起回滚 |
租约导致 manual_reconciliation Worker 重放 | 工单和告警适配器都用 caseKey 做外部幂等键,不重复建单或告警 |
| 场景 | 期望结果 |
|---|---|
旧 Order A 失败后,新 Order B 的 next_action 到达但 B.created 延迟 | 聚合先回到 pending;Session 返回 B 后收敛为 action_required |
order.succeeded 与另一 Order 的 order.failed 并发 | 最终为 paid,且只产生一条 fulfillment Outbox |
paid 之后迟到的 next_action / failed | 状态不回退 |
| Session 查询返回 Order B,Worker 拿锁前 Order C 已提交 | attempts 指纹不一致,丢弃 B 的过期结果并重试,C 不会被覆盖 |
| Session 查询返回的 Order 已属于另一商户或 Session | 归属证据只隔离一次并关闭事件,确定性冲突不消耗自动重试 |
| Order 查询返回状态后,Worker 拿锁前该 attempt 已变化 | 本地版本不一致,丢弃过期状态并重试 |
| 场景 | 期望结果 |
|---|---|
同一 Order 先 order.failed,再收到更旧的 order.next_action | 该尝试保持 failed,商户订单保持 payment_failed,不回退 |
已成功的尝试收到迟到的 order.failed | 该尝试保持 succeeded,商户订单保持 paid |
| 同 Session:Order A 失败、Order B 成功 | 商户订单最终为 paid |
| Order B 成为当前尝试后,Order A 的迟到事件到达 | Order B 和聚合订单都不受影响 |
同一 event.id 重复投递 | 业务状态与副作用只发生一次 |
| 场景 | 期望结果 |
|---|---|
匹配的 session.complete | 只设置 clinkSessionStatus = completed,保存其 event.created 版本,并把 Webhook 标记 processed |
匹配的 session.expired | 只设置 clinkSessionStatus = expired,保存其 event.created 版本,并把 Webhook 标记 processed |
任一 Session 终态发生在 paid、partial_refunded 或 refunded 之后 | 支付、退款、履约和所有 Attempt 状态保持不变 |
缺少 merchantReferenceId 或商户订单尚不存在 | Webhook 保持 pending,并登记一条 reprocess_webhook |
| Session ID 缺失/冲突,或 payload status 与事件类型不一致 | 生命周期不变;原子持久化一张 manual_review case 和一条任务 |
| 较旧的 Session 终态晚于较新终态到达 | 用 event.created 忽略旧更新;不得使用到达时间 |
| 同一毫秒出现不同 Session 终态 | 不猜测;保留原生命周期并持久化 manual_review |
| 同一 Session 终态事件重放 | 生命周期版本、状态和任务保持幂等 |
| 场景 | 期望结果 |
|---|---|
refund.succeeded 先到,order.succeeded 后到 | 第一次处理时事件停在 pending 并登记 reprocess_webhook;订单建好后 Worker 重新处理,退款记账成功,事件转 processed |
正常 order.succeeded 后收到 refund.succeeded | 成功 Order 写入 Order ID、amountTotal、paymentCurrency,退款币种校验和对账可以完成 |
reconcile_attempt 查询收敛为 success 后再退款 | Order 查询会补齐同一套支付基准,不依赖成功 Webhook 路径也能完成退款币种校验和对账 |
| 当前退款或任一已存成功退款的币种与支付币种不同 | 汇总前停止,保留原支付基准并转人工对账 |
同一个 refund.succeeded 重复投递三次 | webhook_events 唯一索引挡掉后两次,退款记录只有一条,refundedAmount 不翻倍 |
同一 refundId 以完全相同的不可变业务字段重放 | 精确复用第一次退款行,不更新任何字段,也不重复累计金额 |
同一 refundId 携带不同金额、币种、Clink Order 或商户订单 | 第一条退款行保持不变;一张 refund_conflict:{refundId} case 和一条任务同时保存旧值与新值 |
已有 refundId 换成一个 attempt 尚不存在的未知 Order ID | 用旧退款行定位并锁住商户订单,直接返回 manual_review,不能误判为依赖缺失 |
| 退款金额为零、负数、非有限值或格式非法 | 不写退款行和自动对账任务;用一张持久化人工对账 case 保存非法输入 |
| 成功退款累计金额恰好等于退款基准 | 写入锁内累计值,退款对账继续收敛到 refunded |
| 成功退款累计金额超过退款基准 | 不把 refundedAmount 写到基准之上,也不登记 refund_reconcile;冲突持久化后转人工 |
连续两次部分退款成功,Order 状态一直是 partial_refunded | 每个 refundId 各有一条策略任务,并携带各自对账时锁内读取的最新累计 refundedAmount;任务乱序也不能让已应用金额回退 |
先部分退款,最后一笔补足全额;第一次 Order 查询仍返回 partial_refunded | 锁内 refundedAmount 已达到 refundablePaidAmount,任务必须重试到 Order 查询返回 refunded;本地状态不会永久停在部分退款 |
| 依赖订单始终不存在 | 重处理任务退避到上限;一个事务把它置为 failed 并创建 pending 的 page_outbox_failure,事件不会被记成 processed |
| 正常顺序(order 先、refund 后) | 与改动前行为一致,一次处理完成 |
| 步骤 | 期望结果 |
|---|---|
Worker A 领取任务,得到 tokenA | 任务 status = processing,claim_token = tokenA |
| Worker 准备启动,但执行前续租影响 0 行 | 不启动任何外部调用或副作用,只记录失租 |
| 领取查询被要求返回一批任务 | 本实现仍最多返回一行,后排任务不会在等待首条任务时消耗租约 |
| 长任务仍属于当前 Worker | 每隔 LEASE_MS / 3 心跳续租,持续后移 lease_until,其他 Worker 不能回收 |
A 卡住,租约到期;Worker B 用 tokenB 重新领取 | claim_token 变成 tokenB,lease_until 顺延 |
| 心跳续租影响 0 行 | 标记本地失租,取消支持的 HTTP 调用,不启动后续步骤,也不再 finish/fail |
| 人工工单建好后才失租 | 发告警前再次确认;不能启动告警步骤 |
| A 恢复后上报成功或失败 | 条件更新影响行数为 0,只记日志,不改状态 |
| B 正常完成 | 任务置 succeeded,没有被 A 覆盖 |
A 在执行中途调 renewLease() | 返回 false,A 应立即停止执行 |
refundId,人工操作按 caseKey,Webhook 重放按 event.id。
Outbox 耗尽与告警
| 场景 | 期望结果 |
|---|---|
任务在 MAX_ATTEMPT 前失败 | 回到 pending,不创建告警任务 |
任务达到 MAX_ATTEMPT | 原子标成 failed,并插入一条 outbox_failure_alert:{taskId} 告警任务 |
| 耗尽事务回滚或 claim token 已失效 | failed 状态和告警任务都不提交 |
| 耗尽事务提交后 Worker 进程立即退出 | 独立告警任务仍持久化为 pending |
| 告警成功 | 等待适配器完成,再把 page_outbox_failure 标为 succeeded |
| 告警失败 | 同一告警任务按带上限的退避回到 pending,不得递归 |
订阅哪些事件
一次性支付至少订阅这几个:| 事件 | 对应处理 |
|---|---|
order.created | 必须订阅。它的 event.created 是判断支付尝试先后顺序的唯一依据 |
order.succeeded | 确认收款,触发发货 |
order.failed | 记录失败原因,允许客户重试 |
order.next_action | 客户需要额外验证,订单进入等待状态 |
session.complete | 把独立 Session 生命周期设为 completed;不得推断支付成功 |
session.expired | 把独立的 clinkSessionStatus 设为 expired;绝不能把 Session 过期转换成支付失败 |
refund.succeeded | 退款处理成功。按 refundId 记账,Order 状态另行对账 |
GET /webhook/events 查当前支持的事件名。
订阅时要写完整的事件名。
subscription.* 这类通配符不是能提交给 API 的值。本地怎么调
Webhook 地址必须是公网可达的 HTTPS,localhost、回环地址和内网 IP 都会被拒绝。 已经有可用的公网地址(预览环境、Vercel/Netlify 之类的部署地址、自有域名)的话,直接用那个。纯本地开发才需要隧道:cloudflared tunnel --url http://127.0.0.1:3000 --no-autoupdate
--protocol http2 重试。
隧道地址每次重启都会变,换了地址记得回后台更新 Webhook 端点。
不刷卡也能调验签
隧道加真实付款能验证端到端,但反复调验签逻辑时,每次都刷一遍卡太慢。下面这些事件是商户自己构造的本地 fixture,不是 Clink 服务端发出的真实事件;对应商户订单和clinkSessionId 必须已经存在于本地数据库。
// send-test-event.mjs
import crypto from 'node:crypto';
const WEBHOOK_URL = 'http://127.0.0.1:3000/api/webhooks/clink';
const signingKey = process.env.CLINK_WEBHOOK_SIGNING_KEY;
if (!signingKey) throw new Error('CLINK_WEBHOOK_SIGNING_KEY is required');
async function sendSignedEvent(event, { signatureOverride } = {}) {
// Sign each exact serialized body with a fresh delivery timestamp.
const rawBody = JSON.stringify(event);
const timestamp = String(Date.now());
const signature = signatureOverride ?? crypto
.createHmac('sha256', signingKey)
.update(`${timestamp}.${rawBody}`)
.digest('hex');
const res = await fetch(WEBHOOK_URL, {
method: 'POST',
headers: {
'Content-Type': 'application/json',
'X-Clink-Timestamp': timestamp,
'X-Clink-Signature': signature,
'X-Clink-SignType': 'SHA256',
},
body: rawBody,
});
console.log(event.id, res.status, await res.text());
return res.status;
}
const merchantReferenceId = 'order_10001'; // must already exist locally
const sessionId = 'sess_test_001'; // must match local clinkSessionId
const orderId = 'order_test_001';
const createdEvent = {
id: 'event_order_created_test',
object: 'event',
type: 'order.created',
created: 1750000000000,
data: {
object: {
orderId,
merchantReferenceId,
sessionId,
status: 'created',
},
},
};
const succeededEvent = {
id: 'event_order_succeeded_test',
object: 'event',
type: 'order.succeeded',
created: 1750000001000,
data: {
object: {
orderId,
merchantReferenceId,
sessionId,
status: 'success',
amountTotal: 19.99,
paymentCurrency: 'USD',
},
},
};
const olderFailedEvent = {
id: 'event_order_failed_old_test',
object: 'event',
type: 'order.failed',
created: 1750000000500, // older than the succeeded status version
data: {
object: {
orderId,
merchantReferenceId,
sessionId,
status: 'failed',
failureCode: 'fixture_declined',
failureMessage: 'Local fixture only',
},
},
};
// Executable happy path, exact duplicate, then an out-of-order old failure.
await sendSignedEvent(createdEvent);
await sendSignedEvent(succeededEvent);
await sendSignedEvent(succeededEvent); // same Event ID and identical business object
await sendSignedEvent(olderFailedEvent);
// Negative fixtures: run individually when checking these branches.
const mismatchedSessionEvent = {
...succeededEvent,
id: 'event_order_succeeded_wrong_session_test',
data: {
object: { ...succeededEvent.data.object, sessionId: 'sess_other' },
},
};
// await sendSignedEvent(mismatchedSessionEvent); // durable manual_review, no payment update
// await sendSignedEvent(createdEvent, { signatureOverride: '0'.repeat(64) }); // expect 401
| 场景 | 期望结果 |
|---|---|
| 先 created、再 succeeded | paymentStatus = paid;退款基准为 order_test_001 / 19.99 / USD;只有一条履约任务 |
| 重复投递 | 完全相同的 succeededEvent 发送两次;第二次返回 200,不重复履约 |
| 乱序 | 成功后再投递状态时间更旧的 olderFailedEvent;订单仍为 paid |
| 错误签名 | 单独运行注释中的 signatureOverride 调用;预期 401 |
| Session 不匹配 | 单独运行注释中的 mismatch fixture;只产生一张持久化 manual_review case/任务,不应用支付状态 |
谁负责什么
| 模块 | 负责 | 不能做 |
|---|---|---|
| 后端 | 建业务订单、调 Clink API、存各种 ID、提供订单状态查询 | 不能把 Secret Key 返回给浏览器;结果未知时不能重复扣款 |
| 前端 | 调本地后端、跳转或挂载收银台、展示支付状态 | 不能直接调 Clink API;不能根据回跳或 SDK 事件发货 |
| Webhook 处理器 | 验签、去重、匹配订单、更新状态、触发发货 | 不能依赖事件顺序;不能让旧状态覆盖新状态 |
常见错误
在回跳页面发货。successUrl 是个普通地址,客户手动敲一遍也能打开。发货只能由验过签的 Webhook 触发。
Webhook 收到就发货,不去重。 Clink 会重试,同一个事件会收到好几遍。不去重的话客户下单一次收到三件货。
pending 当成失败,让客户重新付。 pending 是”还不知道”。这时候重新扣款,很可能扣两次。等 Webhook,或者查 GET /order/{id}。
用 express.json() 解析 Webhook 请求。 解析过的 body 再序列化回去,签名对不上,所有事件都会验签失败。
Secret Key 出现在前端代码里。 拿到它的人可以用该账户发起收款和退款。前端只能拿 Publishable Key。
接下来
密钥与 Webhook 配置
密钥轮换、IP 限制、投递重试规则。
上线检查
切生产前要验的场景清单。