const mq = require('./index') const { mq: mqName } = require('./mqPrefix') const { queryPaymentOrder, getPaymentConfig } = require('../../lib/PaymentClient') const { getGatewayOrderNos } = require('../../lib/OrderPaymentAttempt') const { closePendingOrder, completePaidOrder } = require('../../lib/OrderSettlement') const db = require('../DataBase/db') const ORDER_PAYMENT_QUEUE = mqName('order_payment_check') let orderPaymentWorkerStarted = false async function isOrderStillPending(orderId) { const rows = await db.query('SELECT state FROM orders WHERE orderId = ? LIMIT 1', [orderId]) if (!rows?.length) return { exists: false, pending: false } return { exists: true, pending: Number(rows[0].state) === 0, state: Number(rows[0].state) } } function sleep(ms) { return new Promise(resolve => setTimeout(resolve, ms)) } async function pollOrderPaymentStatus(orderId, logger) { const paymentConfig = await getPaymentConfig() if (!paymentConfig.pid || !paymentConfig.url || !paymentConfig.key) { logger.error('支付配置错误,无法轮询易支付状态') return } const maxRetries = 120 const delayMs = 2500 logger.info(`开始轮询订单支付状态,订单号:${orderId}`) for (let retry = 0; retry <= maxRetries; retry++) { const status = await isOrderStillPending(orderId) if (!status.exists) { logger.warn(`订单不存在,停止轮询:${orderId}`) return } if (!status.pending) { logger.info(`订单已处理(state=${status.state}),停止轮询:${orderId}`) return } if (retry >= maxRetries) { const closeResult = await closePendingOrder({ orderId, logger }) if (closeResult.closed) logger.info(`订单超时未支付,自动取消:${orderId}`) return } try { const gatewayOrderNos = await getGatewayOrderNos(orderId) for (const gatewayOrderNo of gatewayOrderNos) { const queryData = await queryPaymentOrder(gatewayOrderNo, logger) logger.info(`轮询支付状态,订单号:${orderId},网关订单号:${gatewayOrderNo},次数:${retry + 1},结果:${JSON.stringify(queryData)}`) if (Number(queryData.code) === 1 && Number(queryData.status) === 1) { const result = await completePaidOrder({ orderId, payType: queryData.type, payId: queryData.trade_no, payTime: Date.now(), paidAmount: queryData.money, logger }) if (result.completed || result.reason === 'state_changed') return logger.error(`网关已支付但订单完成失败,订单号:${orderId},原因:${result.reason}`) return } } } catch (error) { logger.warn(`轮询支付状态失败,订单号:${orderId},原因:${error.message || error}`) } await sleep(delayMs) } } async function enqueueOrderPaymentCheck(orderId) { const ch = await mq.getChannel('order_payment') await ch.assertQueue(ORDER_PAYMENT_QUEUE, { durable: true }) ch.sendToQueue( ORDER_PAYMENT_QUEUE, Buffer.from(JSON.stringify({ orderId, enqueueTime: Date.now() })), { persistent: true } ) } async function startOrderPaymentWorker(logger) { if (orderPaymentWorkerStarted) return try { const ch = await mq.getChannel('order_payment') await ch.assertQueue(ORDER_PAYMENT_QUEUE, { durable: true }) await ch.prefetch(1) logger.info(`订单支付结果轮询消费者已启动,队列:${ORDER_PAYMENT_QUEUE}`) orderPaymentWorkerStarted = true ch.consume(ORDER_PAYMENT_QUEUE, async (msg) => { if (!msg) return const content = JSON.parse(msg.content.toString() || '{}') const { orderId } = content if (!orderId) { logger.warn('收到无效的订单支付检查消息:缺少 orderId') ch.ack(msg) return } try { await pollOrderPaymentStatus(orderId, logger) ch.ack(msg) } catch (err) { logger.error(`订单支付轮询处理失败,订单号:${orderId},错误:${err.stack || err}`) ch.nack(msg, false, true) } }) } catch (e) { logger.error(`启动订单支付 MQ 消费者失败:${e.stack || e}`) } } module.exports = { ORDER_PAYMENT_QUEUE, enqueueOrderPaymentCheck, startOrderPaymentWorker }