新增 affiliate_dash 回调路由与幂等去重(阶段 5)

- POST /api/v1/open/affiliate-dash:验签 → 查任务(404 重试) → webhook_events 幂等闸门 → 状态同步
- migration 013:webhook_events(provider+event_id 唯一)
- task-repo findAffiliateDashTaskByOrder(仿 feifei 同构 SQL)
- notify-service:去重闸门在找到任务之后,任务未匹配不被 eventId 永久拦截
- typecheck + 212 测试通过;前置错误分支脚本验证通过
This commit is contained in:
yml2213
2026-08-05 16:07:33 +08:00
parent 00d490a2cd
commit 7fc3d60588
7 changed files with 294 additions and 0 deletions
+2
View File
@@ -3,6 +3,7 @@ import express from 'express'
import process from 'node:process'
import adminRouter from './routes/admin.js'
import affiliateDashRouter from './routes/affiliate-dash.js'
import collectRouter from './routes/collect.js'
import filesRouter from './routes/files.js'
import claimsRouter from './routes/claims.js'
@@ -93,6 +94,7 @@ export function createApp({ startupState, isShutdownStarted, config }: CreateApp
})
})
app.use('/api/v1/open/affiliate-dash', affiliateDashRouter)
app.use('/api/v1/open/91', open91Router)
app.use('/api/v1/open/kuaishou-industry', kuaishouIndustryRouter)
app.use('/api/v1/open/kuaishou-feifei', kuaishouFeifeiRouter)
@@ -0,0 +1,15 @@
CREATE TABLE IF NOT EXISTS webhook_events (
id BIGSERIAL PRIMARY KEY,
provider TEXT NOT NULL,
event_id TEXT NOT NULL,
event_type TEXT NOT NULL DEFAULT '',
task_id BIGINT,
payload JSONB NOT NULL,
processed_at TIMESTAMPTZ NOT NULL,
created_at TIMESTAMPTZ NOT NULL
);
CREATE UNIQUE INDEX IF NOT EXISTS idx_webhook_events_provider_event_id
ON webhook_events(provider, event_id);
COMMENT ON TABLE webhook_events IS '第三方回调事件去重表(provider + event_id 唯一,幂等处理)';
@@ -16,6 +16,7 @@ import type {
type TaskQueryExecutor = (text: string, params?: unknown[]) => Promise<{ rows: unknown[] }>
const KUAISHOU_FEIFEI_EXECUTOR_KEY = 'kuaishou_feifei'
const AFFILIATE_DASH_EXECUTOR_KEY = 'affiliate_dash'
type TaskRuntimeContextRow = {
runtime_session_id: string
@@ -355,6 +356,40 @@ export async function findKuaishouFeifeiTaskByOrder(input: {
return result.rows[0] || null
}
export async function findAffiliateDashTaskByOrder(input: {
orderNo?: unknown
clientOrderNo?: unknown
}): Promise<TaskRow | null> {
const orderNo = String(input.orderNo || '').trim()
const clientOrderNo = String(input.clientOrderNo || '').trim()
const params: unknown[] = [AFFILIATE_DASH_EXECUTOR_KEY]
const orderFilters: string[] = []
if (orderNo) {
params.push(orderNo)
orderFilters.push(`ft.context_json #>> '{affiliateDash,orderNo}' = $${params.length}`)
}
if (clientOrderNo) {
params.push(clientOrderNo)
orderFilters.push(`ft.context_json #>> '{affiliateDash,clientOrderNo}' = $${params.length}`)
}
if (orderFilters.length === 0) {
return null
}
const result = await query<TaskRow>(
`${buildTaskSelect()}
WHERE ft.executor_key = $1
AND (${orderFilters.join(' OR ')})
ORDER BY ft.id DESC
LIMIT 1`,
params,
)
return result.rows[0] || null
}
export async function listTasks({
page = 1,
pageSize = 20,
@@ -0,0 +1,43 @@
import { query } from '../db/client.js'
import type { JsonObject } from '../types/json.js'
import { nowIso } from '../utils/time.js'
/**
* 幂等写入一条 webhook 事件(provider + event_id 唯一)。
* 冲突时不做任何更新并返回 inserted=false —— 供回调去重闸门使用。
*/
export async function insertWebhookEventOnce(input: {
provider: string
eventId: string
eventType: string
taskId?: number | string | null
payload: JsonObject
processedAt?: string
}): Promise<{ inserted: boolean }> {
const processedAt = input.processedAt || nowIso()
const result = await query<{ id: number }>(
`
INSERT INTO webhook_events (
provider,
event_id,
event_type,
task_id,
payload,
processed_at,
created_at
) VALUES ($1, $2, $3, $4, $5::jsonb, $6, $6)
ON CONFLICT (provider, event_id) DO NOTHING
RETURNING id
`,
[
String(input.provider || '').trim(),
String(input.eventId || '').trim(),
String(input.eventType || '').trim(),
input.taskId == null ? null : Number(input.taskId),
JSON.stringify(input.payload || {}),
processedAt,
],
)
return { inserted: result.rows.length > 0 }
}
+56
View File
@@ -0,0 +1,56 @@
import { Router } from 'express'
import { createRateLimitMiddleware } from '../middleware/rate-limit.js'
import { handleAffiliateDashNotify } from '../services/platforms/affiliate-dash/notify-service.js'
import { sendRouteError } from '../utils/http.js'
import { createRequestId, logIntegration } from '../utils/logger.js'
const router = Router()
const notifyRateLimit = createRateLimitMiddleware({
scope: 'affiliateDashNotify',
windowMs: 60_000,
max: 300,
})
router.post('/', notifyRateLimit, async (req, res) => {
const requestId = createRequestId('adn')
const startedAt = Date.now()
logIntegration('[affiliate-dash/notify]', '收到 affiliate-dash 订单回调', {
requestId,
method: req.method,
originalUrl: req.originalUrl,
ip: req.ip,
headers: {
'content-type': req.headers['content-type'],
'x-event-id': req.headers['x-event-id'],
'x-timestamp': req.headers['x-timestamp'],
'x-sign': req.headers['x-sign'],
},
rawBody: req.rawBody,
})
try {
const result = await handleAffiliateDashNotify({
body: req.body,
headers: req.headers,
rawBody: req.rawBody,
})
logIntegration('[affiliate-dash/notify]', '订单回调处理完成', {
requestId,
durationMs: Date.now() - startedAt,
result,
})
res.status(200).json({ code: 0 })
} catch (error) {
logIntegration('[affiliate-dash/notify]', '订单回调处理失败', {
requestId,
durationMs: Date.now() - startedAt,
error,
}, { level: 'error' })
sendRouteError(res, error, 'affiliate-dash 回调处理失败', '[affiliate-dash/notify]')
}
})
export default router
@@ -0,0 +1,124 @@
import { createTaskEvent } from '../../../repositories/task-event-repo.js'
import { findAffiliateDashTaskByOrder } from '../../../repositories/task-repo.js'
import { insertWebhookEventOnce } from '../../../repositories/webhook-event-repo.js'
import type { JsonObject } from '../../../types/json.js'
import { syncAffiliateDashTaskStatus } from '../../fulfillment/affiliate-dash/index.js'
import { createHttpError } from '../../../utils/http.js'
import { nowIso } from '../../../utils/time.js'
import { verifyAffiliateDashCallback } from './verify-callback.js'
type NotifyHeaders = Record<string, string | string[] | undefined>
/**
* affiliate-dash 回调入口(验签 + 幂等去重 + 状态同步)。
*
* 处理顺序:
* 1. 验签(X-Timestamp 容差 → 重算签名)并解析事件;
* 2. 按 order_no / client_order_no 查本地任务(未匹配 → 404,触发对方重试);
* 3. 幂等闸门:webhook_events(provider + event_id 唯一),重复回调直接 200 跳过;
* 4. 记录任务事件 + 状态合并式同步(可重入)。
*/
export async function handleAffiliateDashNotify(input: {
body: unknown
headers: NotifyHeaders
rawBody?: string | undefined
}) {
const verified = verifyAffiliateDashCallback({
headers: input.headers,
rawBody: input.rawBody,
})
const data = isPlainObject(verified.data) ? verified.data : {}
const orderNo = String(data.order_no || '').trim()
const clientOrderNo = String(data.client_order_no || '').trim()
const orderStatus = String(data.order_status || '').trim()
if (!orderNo && !clientOrderNo) {
throw createHttpError('affiliate-dash 回调订单号缺失', {
statusCode: 400,
errorCode: 'affiliate_dash_callback_order_no_missing',
})
}
const task = await findAffiliateDashTaskByOrder({
orderNo,
clientOrderNo,
})
if (!task) {
throw createHttpError('affiliate-dash 回调未匹配到本地任务', {
statusCode: 404,
errorCode: 'affiliate_dash_callback_task_not_found',
context: {
orderNo,
clientOrderNo,
},
})
}
const rawBody = String(input.rawBody || '')
let payload: JsonObject = {}
try {
const parsed = JSON.parse(rawBody)
if (isPlainObject(parsed)) {
payload = parsed
}
} catch {
payload = { order_no: orderNo, client_order_no: clientOrderNo, order_status: orderStatus }
}
const dedupe = await insertWebhookEventOnce({
provider: 'affiliate_dash',
eventId: verified.eventId,
eventType: verified.event,
taskId: task.id,
payload,
})
if (!dedupe.inserted) {
return {
ok: true,
duplicate: true,
eventId: verified.eventId,
event: verified.event,
taskId: task.id,
orderNo,
orderStatus,
}
}
const now = nowIso()
await createTaskEvent(
task.id,
'affiliate_dash_notify_received',
{
event: verified.event,
eventId: verified.eventId,
occurredAt: verified.occurredAt,
orderNo,
clientOrderNo,
orderStatus,
raw: data,
},
now,
)
const updatedTask = await syncAffiliateDashTaskStatus(task)
return {
ok: true,
duplicate: false,
eventId: verified.eventId,
event: verified.event,
taskId: task.id,
orderNo,
clientOrderNo,
orderStatus,
taskStatus: updatedTask.task_status,
deliveryStatus: updatedTask.delivery_status,
}
}
function isPlainObject(value: unknown): value is JsonObject {
return Object.prototype.toString.call(value) === '[object Object]'
}