| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129 |
- 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
- }
|