From e8eab6e689a82a3256c0c8eee856de29cf3e81b3 Mon Sep 17 00:00:00 2001 From: yml2213 Date: Wed, 5 Aug 2026 20:05:43 +0800 Subject: [PATCH] =?UTF-8?q?fix(kuaishou-industry):=20=E4=BF=AE=E5=A4=8D?= =?UTF-8?q?=E7=94=B5=E5=AD=90=E5=87=AD=E8=AF=81=E5=8F=91=E8=B4=A7=E5=9B=9E?= =?UTF-8?q?=E8=B0=83=E5=B9=B6=E5=8F=91=E5=AF=BC=E8=87=B4=20807000=20?= =?UTF-8?q?=E8=8E=B7=E5=8F=96=E5=B9=B6=E5=8F=91=E9=94=81=E5=A4=B1=E8=B4=A5?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 同一 oid 的发货回调可能被多个来源并发调度(内存重试链 timer、worker_scan 兜底扫描、快手重复推送 send-code 后的再次调度),快手侧对同一 oid 有并发锁, 并发请求会返回 807000「获取并发锁失败」。 - runKuaishouIndustrySendCallbackForOid 增加进程内 in-flight 互斥, 正在执行时其它调度直接跳过,由当前执行链路决定成功或按 backoff 重试 (try/finally 保证所有 return 路径释放) - listKuaishouIndustryVouchersForSendCallbackRetry 增加 maxAttempts 参数 及 SQL 条件 COALESCE(send_callback_attempt_count,0) < ,缩小扫描集 --- .../kuaishou-industry-voucher-repo.ts | 8 +- .../kuaishou-industry/send-code-service.ts | 304 +++++++++++------- 2 files changed, 187 insertions(+), 125 deletions(-) diff --git a/apps/backend/src/repositories/kuaishou-industry-voucher-repo.ts b/apps/backend/src/repositories/kuaishou-industry-voucher-repo.ts index a8b8b837..b18aeb7c 100644 --- a/apps/backend/src/repositories/kuaishou-industry-voucher-repo.ts +++ b/apps/backend/src/repositories/kuaishou-industry-voucher-repo.ts @@ -203,16 +203,18 @@ export async function listKuaishouIndustryVouchersByOid( export async function listKuaishouIndustryVouchersForSendCallbackRetry( limit = 100, + maxAttempts = 8, ): Promise { const result = await query( ` SELECT * FROM kuaishou_industry_vouchers WHERE send_callback_status <> 'success' + AND COALESCE(send_callback_attempt_count, 0) < $2 ORDER BY COALESCE(send_callback_sent_at, created_at) ASC, id ASC LIMIT $1 `, - [normalizePositiveLimit(limit, 100)], + [normalizePositiveLimit(limit, 100), Math.max(1, Math.trunc(Number(maxAttempts) || 8))], ) return result.rows @@ -277,7 +279,9 @@ export async function listKuaishouIndustryVouchersForAdmin( where.push(`v.seller_id = $${params.length}`) } - const status = String(input.status || '').trim().toUpperCase() + const status = String(input.status || '') + .trim() + .toUpperCase() if (status) { params.push(status) where.push(`v.status = $${params.length}`) diff --git a/apps/backend/src/services/platforms/kuaishou-industry/send-code-service.ts b/apps/backend/src/services/platforms/kuaishou-industry/send-code-service.ts index 07b3e649..8e927e68 100644 --- a/apps/backend/src/services/platforms/kuaishou-industry/send-code-service.ts +++ b/apps/backend/src/services/platforms/kuaishou-industry/send-code-service.ts @@ -1,6 +1,6 @@ import { findLatestOrderByPlatformOrderId } from '../../../repositories/order-repo.js' import { listTasksByOrderId } from '../../../repositories/task-repo.js' -import { asJsonObject, type JsonObject } from '../../../types/json.js' +import { asJsonObject, type JsonObject } from '../../../types/json.js' import { findKuaishouIndustryVoucherByCode, listKuaishouIndustryVouchersByOid, @@ -13,10 +13,7 @@ import { retryOpen91Order } from '../ninetyone/order-service.js' import { syncWorkOrdersFromKuaishouSendCode } from '../../worker-platform/sync-work-orders-from-send-code.js' import { normalizeTimestampIso } from '../../../utils/time.js' import { logWarn, logIntegration } from '../../../utils/logger.js' -import { - getKuaishouIndustryConfig, - assertMatchingAppKey, -} from './config.js' +import { getKuaishouIndustryConfig, assertMatchingAppKey } from './config.js' import { normalizeSendCodePayload, assertSendCodePayload } from './payload.js' import { assertKuaishouIndustrySignature } from './crypto.js' import { @@ -51,6 +48,7 @@ const SEND_CALLBACK_MAX_ATTEMPTS = 8 const SEND_CALLBACK_BACKOFF_MS = [3_000, 10_000, 30_000, 60_000, 120_000, 300_000] const sendCallbackRetryTimers = new Map() +const sendCallbackInFlightOids = new Set() let sendCallbackWorkerTimer: NodeJS.Timeout | null = null let sendCallbackWorkerRunning = false @@ -231,73 +229,104 @@ async function runKuaishouIndustrySendCallbackForOid( } } - const vouchers = await listKuaishouIndustryVouchersByOid(normalizedOid) - const pendingVouchers = vouchers - .filter((voucher) => !isKuaishouIndustryVoucherSendCallbackSuccess(voucher)) - .filter((voucher) => - normalizePositiveInteger(voucher.send_callback_attempt_count) < SEND_CALLBACK_MAX_ATTEMPTS, + // 同一 oid 的回调可能被多个来源同时调度(内存重试链 timer、worker_scan 兜底扫描、 + // 快手重复推送 send-code 后的再次调度)。快手侧对同一 oid 的电子凭证发货回调有并发锁, + // 并发请求会拿到 807000「获取并发锁失败」。这里用 in-flight 集合做进程内互斥: + // 正在执行时其它调度直接跳过,由当前执行链路自己决定成功或按 backoff 重试。 + if (sendCallbackInFlightOids.has(normalizedOid)) { + logIntegration( + '[kuaishou-industry/send-code]', + '电子凭证发货回调正在执行中,跳过本次并发调度', + { + source, + oid: normalizedOid, + }, ) - - if (pendingVouchers.length === 0) { - return { - success: vouchers.length > 0, - skipped: true, - reason: vouchers.length > 0 ? 'send_callback_already_success_or_max_attempts' : 'voucher_missing', - } - } - - const rawPayload = parseJsonObject(pendingVouchers[0]?.raw_payload_json) - const body = parseJsonObject(rawPayload.body) - const firstVoucher = pendingVouchers[0] - if (!firstVoucher) { return { success: false, - error: '缺少待回调电子凭证', + skipped: true, + reason: 'callback_in_flight', } } + sendCallbackInFlightOids.add(normalizedOid) - const params = resolveSendCodeCallbackParams(firstVoucher, body) - const eticketType = String(params.eticketType || '').trim() - const etickets = pendingVouchers.map((voucher) => - buildKuaishouIndustryEticketFromVoucher(voucher, eticketType), - ) - const order = await findLatestOrderByPlatformOrderId({ - provider: OPEN_91_PROVIDER, - platform: OPEN_91_PLATFORM, - platformOrderId: normalizedOid, - }) - const result = await sendAndRecordCallback({ - oid: params.oid, - sendType: params.sendType, - etickets, - vouchers: pendingVouchers, - sellerId: params.sellerId, - token: params.token, - eticketType, - source, - preferredTotalGoodsValue: resolveSendCallbackPreferredTotalGoodsValue(order), - ...(params.ext ? { ext: params.ext } : {}), - }) + try { + const vouchers = await listKuaishouIndustryVouchersByOid(normalizedOid) + const pendingVouchers = vouchers + .filter((voucher) => !isKuaishouIndustryVoucherSendCallbackSuccess(voucher)) + .filter( + (voucher) => + normalizePositiveInteger(voucher.send_callback_attempt_count) < + SEND_CALLBACK_MAX_ATTEMPTS, + ) - if (result.success) { - await syncOpen91OrderAfterSendCallbackSuccess(normalizedOid, { source }) - return result - } + if (pendingVouchers.length === 0) { + return { + success: vouchers.length > 0, + skipped: true, + reason: + vouchers.length > 0 ? 'send_callback_already_success_or_max_attempts' : 'voucher_missing', + } + } - const nextAttemptCount = - Math.max(...pendingVouchers.map((voucher) => - normalizePositiveInteger(voucher.send_callback_attempt_count), - )) + 1 - if ( - nextAttemptCount < SEND_CALLBACK_MAX_ATTEMPTS && - isRetriableSendCallbackFailure(result) - ) { - scheduleKuaishouIndustrySendCallbackRetry(normalizedOid, resolveSendCallbackBackoffMs(nextAttemptCount), { - source: `${source}_retry`, + const rawPayload = parseJsonObject(pendingVouchers[0]?.raw_payload_json) + const body = parseJsonObject(rawPayload.body) + const firstVoucher = pendingVouchers[0] + if (!firstVoucher) { + return { + success: false, + error: '缺少待回调电子凭证', + } + } + + const params = resolveSendCodeCallbackParams(firstVoucher, body) + const eticketType = String(params.eticketType || '').trim() + const etickets = pendingVouchers.map((voucher) => + buildKuaishouIndustryEticketFromVoucher(voucher, eticketType), + ) + const order = await findLatestOrderByPlatformOrderId({ + provider: OPEN_91_PROVIDER, + platform: OPEN_91_PLATFORM, + platformOrderId: normalizedOid, + }) + const result = await sendAndRecordCallback({ + oid: params.oid, + sendType: params.sendType, + etickets, + vouchers: pendingVouchers, + sellerId: params.sellerId, + token: params.token, + eticketType, + source, + preferredTotalGoodsValue: resolveSendCallbackPreferredTotalGoodsValue(order), + ...(params.ext ? { ext: params.ext } : {}), }) - } - return result + if (result.success) { + await syncOpen91OrderAfterSendCallbackSuccess(normalizedOid, { source }) + return result + } + + const nextAttemptCount = + Math.max( + ...pendingVouchers.map((voucher) => + normalizePositiveInteger(voucher.send_callback_attempt_count), + ), + ) + 1 + if (nextAttemptCount < SEND_CALLBACK_MAX_ATTEMPTS && isRetriableSendCallbackFailure(result)) { + scheduleKuaishouIndustrySendCallbackRetry( + normalizedOid, + resolveSendCallbackBackoffMs(nextAttemptCount), + { + source: `${source}_retry`, + }, + ) + } + + return result + } finally { + sendCallbackInFlightOids.delete(normalizedOid) + } } function scheduleKuaishouIndustrySendCallbackRetry( @@ -314,11 +343,16 @@ function scheduleKuaishouIndustrySendCallbackRetry( const timer = setTimeout(() => { sendCallbackRetryTimers.delete(normalizedOid) void runKuaishouIndustrySendCallbackForOid(normalizedOid, { source }).catch((error) => { - logIntegration('[kuaishou-industry/send-code]', '异步电子凭证发货回调执行异常', { - source, - oid: normalizedOid, - error: error instanceof Error ? error.message : String(error), - }, { level: 'warn' }) + logIntegration( + '[kuaishou-industry/send-code]', + '异步电子凭证发货回调执行异常', + { + source, + oid: normalizedOid, + error: error instanceof Error ? error.message : String(error), + }, + { level: 'warn' }, + ) scheduleKuaishouIndustrySendCallbackRetry(normalizedOid, resolveSendCallbackBackoffMs(1), { source: `${source}_exception_retry`, }) @@ -339,12 +373,15 @@ function scheduleSendCallbackWorkerScan(delayMs: number) { return } - sendCallbackWorkerTimer = setTimeout(() => { - sendCallbackWorkerTimer = null - void runSendCallbackWorkerScan().finally(() => { - scheduleSendCallbackWorkerScan(SEND_CALLBACK_SCAN_INTERVAL_MS) - }) - }, Math.max(0, Number(delayMs) || 0)) + sendCallbackWorkerTimer = setTimeout( + () => { + sendCallbackWorkerTimer = null + void runSendCallbackWorkerScan().finally(() => { + scheduleSendCallbackWorkerScan(SEND_CALLBACK_SCAN_INTERVAL_MS) + }) + }, + Math.max(0, Number(delayMs) || 0), + ) } async function runSendCallbackWorkerScan() { @@ -354,7 +391,10 @@ async function runSendCallbackWorkerScan() { sendCallbackWorkerRunning = true try { - const vouchers = await listKuaishouIndustryVouchersForSendCallbackRetry(200) + const vouchers = await listKuaishouIndustryVouchersForSendCallbackRetry( + 200, + SEND_CALLBACK_MAX_ATTEMPTS, + ) const dueOids = new Set() for (const voucher of vouchers) { if (isVoucherDueForSendCallbackRetry(voucher)) { @@ -368,9 +408,14 @@ async function runSendCallbackWorkerScan() { }) } } catch (error) { - logIntegration('[kuaishou-industry/send-code]', '扫描待重试电子凭证发货回调失败', { - error: error instanceof Error ? error.message : String(error), - }, { level: 'warn' }) + logIntegration( + '[kuaishou-industry/send-code]', + '扫描待重试电子凭证发货回调失败', + { + error: error instanceof Error ? error.message : String(error), + }, + { level: 'warn' }, + ) } finally { sendCallbackWorkerRunning = false } @@ -389,18 +434,19 @@ function isVoucherDueForSendCallbackRetry(voucher: KuaishouIndustryVoucherRow) { const baseAt = voucher.send_callback_sent_at || voucher.created_at const baseTime = Date.parse(String(baseAt || '')) const elapsedMs = Date.now() - (Number.isFinite(baseTime) ? baseTime : 0) - const requiredDelayMs = attemptCount === 0 - ? SEND_CALLBACK_INITIAL_DELAY_MS - : resolveSendCallbackBackoffMs(attemptCount) + const requiredDelayMs = + attemptCount === 0 ? SEND_CALLBACK_INITIAL_DELAY_MS : resolveSendCallbackBackoffMs(attemptCount) return elapsedMs >= requiredDelayMs } function resolveSendCallbackBackoffMs(attemptCount: unknown) { const attempt = Math.max(1, normalizePositiveInteger(attemptCount)) - return SEND_CALLBACK_BACKOFF_MS[Math.min(attempt, SEND_CALLBACK_BACKOFF_MS.length - 1)] || + return ( + SEND_CALLBACK_BACKOFF_MS[Math.min(attempt, SEND_CALLBACK_BACKOFF_MS.length - 1)] || SEND_CALLBACK_BACKOFF_MS[SEND_CALLBACK_BACKOFF_MS.length - 1] || SEND_CALLBACK_INITIAL_DELAY_MS + ) } async function syncOpen91OrderAfterSendCallbackSuccess( @@ -439,12 +485,17 @@ async function syncOpen91OrderAfterSendCallbackSuccess( taskCount: result.tasks.length, }) } catch (error) { - logIntegration('[kuaishou-industry/send-code]', '电子凭证发货回调成功后重试 91 订单失败', { - source, - oid: normalizedOid, - orderId: order.id, - error: error instanceof Error ? error.message : String(error), - }, { level: 'warn' }) + logIntegration( + '[kuaishou-industry/send-code]', + '电子凭证发货回调成功后重试 91 订单失败', + { + source, + oid: normalizedOid, + orderId: order.id, + error: error instanceof Error ? error.message : String(error), + }, + { level: 'warn' }, + ) } } @@ -465,13 +516,11 @@ function isRetriableSendCallbackFailure(result: Awaited String(item || '').trim()).join(' ') + ] + .map((item) => String(item || '').trim()) + .join(' ') - return ( - marker.includes('9994') || - marker.includes('807000') || - marker.includes('并发冲突') - ) + return marker.includes('9994') || marker.includes('807000') || marker.includes('并发冲突') } async function sendAndRecordCallback(input: { @@ -538,26 +587,33 @@ async function sendAndRecordCallback(input: { } const sentAt = normalizeTimestampIso(new Date().toISOString()) - await Promise.all(input.vouchers.map((voucher) => - updateKuaishouIndustryVoucherByCode(voucher.voucher_code, { - sendCallbackStatus: result.success - ? KUAISHOU_INDUSTRY_SEND_CALLBACK_STATUS.SUCCESS - : KUAISHOU_INDUSTRY_SEND_CALLBACK_STATUS.FAILED, - sendCallbackAttemptCount: normalizePositiveInteger(voucher.send_callback_attempt_count) + 1, - sendCallbackLastError: result.success ? '' : resolveSendCallbackFailureMessage(result), - sendCallbackResponseJson: result.response || {}, - sendCallbackSentAt: sentAt, - updatedAt: sentAt, - }), - )) + await Promise.all( + input.vouchers.map((voucher) => + updateKuaishouIndustryVoucherByCode(voucher.voucher_code, { + sendCallbackStatus: result.success + ? KUAISHOU_INDUSTRY_SEND_CALLBACK_STATUS.SUCCESS + : KUAISHOU_INDUSTRY_SEND_CALLBACK_STATUS.FAILED, + sendCallbackAttemptCount: normalizePositiveInteger(voucher.send_callback_attempt_count) + 1, + sendCallbackLastError: result.success ? '' : resolveSendCallbackFailureMessage(result), + sendCallbackResponseJson: result.response || {}, + sendCallbackSentAt: sentAt, + updatedAt: sentAt, + }), + ), + ) if (!result.success) { - logIntegration('[kuaishou-industry/send-code]', '电子凭证发货回调执行失败,等待异步重试', { - source: input.source || 'send_code', - oid: input.oid, - error: result.error || '', - response: result.response || null, - }, { level: 'warn' }) + logIntegration( + '[kuaishou-industry/send-code]', + '电子凭证发货回调执行失败,等待异步重试', + { + source: input.source || 'send_code', + oid: input.oid, + error: result.error || '', + response: result.response || null, + }, + { level: 'warn' }, + ) return result } @@ -605,7 +661,7 @@ function resolveSendCallbackFailureMessage(result: Awaited, + etickets: Array<{ goodsValue?: unknown; [key: string]: unknown }>, ext = '', preferredTotalGoodsValue: unknown = null, ) { @@ -646,7 +702,9 @@ export function resolveSendCallbackPreferredTotalGoodsValue(order: OrderRow | nu const rawPayload = parseJsonObject(order.raw_payload_json) const body = parseJsonObject(rawPayload.body) - const maxAmountFen = normalizePositiveInteger(parseAmountToFen(body.maxAmount ?? rawPayload.maxAmount)) + const maxAmountFen = normalizePositiveInteger( + parseAmountToFen(body.maxAmount ?? rawPayload.maxAmount), + ) if (maxAmountFen > 0) { return maxAmountFen } @@ -677,11 +735,11 @@ function splitGoodsValue(totalGoodsValue: number, count: number): number[] { function resolveGoodsValueTotalFromExt(ext: unknown): number { const parsed = parseJsonObject(ext) return normalizePositiveInteger( - parsed.totalGoodsValue - ?? parsed.goodsValue - ?? parsed.payment - ?? parsed.payAmount - ?? parsed.amount, + parsed.totalGoodsValue ?? + parsed.goodsValue ?? + parsed.payment ?? + parsed.payAmount ?? + parsed.amount, ) } @@ -693,11 +751,11 @@ function normalizePositiveInteger(value: unknown): number { function resolveSendCodePaymentFen(ext: unknown): number { const parsed = parseJsonObject(ext) return normalizePositiveInteger( - parsed.payment - ?? parsed.totalGoodsValue - ?? parsed.goodsValue - ?? parsed.payAmount - ?? parsed.amount, + parsed.payment ?? + parsed.totalGoodsValue ?? + parsed.goodsValue ?? + parsed.payAmount ?? + parsed.amount, ) }