orderPaymentWorker.js 4.7 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129
  1. const mq = require('./index')
  2. const { mq: mqName } = require('./mqPrefix')
  3. const { queryPaymentOrder, getPaymentConfig } = require('../../lib/PaymentClient')
  4. const { getGatewayOrderNos } = require('../../lib/OrderPaymentAttempt')
  5. const { closePendingOrder, completePaidOrder } = require('../../lib/OrderSettlement')
  6. const db = require('../DataBase/db')
  7. const ORDER_PAYMENT_QUEUE = mqName('order_payment_check')
  8. let orderPaymentWorkerStarted = false
  9. async function isOrderStillPending(orderId) {
  10. const rows = await db.query('SELECT state FROM orders WHERE orderId = ? LIMIT 1', [orderId])
  11. if (!rows?.length) return { exists: false, pending: false }
  12. return { exists: true, pending: Number(rows[0].state) === 0, state: Number(rows[0].state) }
  13. }
  14. function sleep(ms) {
  15. return new Promise(resolve => setTimeout(resolve, ms))
  16. }
  17. async function pollOrderPaymentStatus(orderId, logger) {
  18. const paymentConfig = await getPaymentConfig()
  19. if (!paymentConfig.pid || !paymentConfig.url || !paymentConfig.key) {
  20. logger.error('支付配置错误,无法轮询易支付状态')
  21. return
  22. }
  23. const maxRetries = 120
  24. const delayMs = 2500
  25. logger.info(`开始轮询订单支付状态,订单号:${orderId}`)
  26. for (let retry = 0; retry <= maxRetries; retry++) {
  27. const status = await isOrderStillPending(orderId)
  28. if (!status.exists) {
  29. logger.warn(`订单不存在,停止轮询:${orderId}`)
  30. return
  31. }
  32. if (!status.pending) {
  33. logger.info(`订单已处理(state=${status.state}),停止轮询:${orderId}`)
  34. return
  35. }
  36. if (retry >= maxRetries) {
  37. const closeResult = await closePendingOrder({ orderId, logger })
  38. if (closeResult.closed) logger.info(`订单超时未支付,自动取消:${orderId}`)
  39. return
  40. }
  41. try {
  42. const gatewayOrderNos = await getGatewayOrderNos(orderId)
  43. for (const gatewayOrderNo of gatewayOrderNos) {
  44. const queryData = await queryPaymentOrder(gatewayOrderNo, logger)
  45. logger.info(`轮询支付状态,订单号:${orderId},网关订单号:${gatewayOrderNo},次数:${retry + 1},结果:${JSON.stringify(queryData)}`)
  46. if (Number(queryData.code) === 1 && Number(queryData.status) === 1) {
  47. const result = await completePaidOrder({
  48. orderId,
  49. payType: queryData.type,
  50. payId: queryData.trade_no,
  51. payTime: Date.now(),
  52. paidAmount: queryData.money,
  53. logger
  54. })
  55. if (result.completed || result.reason === 'state_changed') return
  56. logger.error(`网关已支付但订单完成失败,订单号:${orderId},原因:${result.reason}`)
  57. return
  58. }
  59. }
  60. } catch (error) {
  61. logger.warn(`轮询支付状态失败,订单号:${orderId},原因:${error.message || error}`)
  62. }
  63. await sleep(delayMs)
  64. }
  65. }
  66. async function enqueueOrderPaymentCheck(orderId) {
  67. const ch = await mq.getChannel('order_payment')
  68. await ch.assertQueue(ORDER_PAYMENT_QUEUE, { durable: true })
  69. ch.sendToQueue(
  70. ORDER_PAYMENT_QUEUE,
  71. Buffer.from(JSON.stringify({ orderId, enqueueTime: Date.now() })),
  72. { persistent: true }
  73. )
  74. }
  75. async function startOrderPaymentWorker(logger) {
  76. if (orderPaymentWorkerStarted) return
  77. try {
  78. const ch = await mq.getChannel('order_payment')
  79. await ch.assertQueue(ORDER_PAYMENT_QUEUE, { durable: true })
  80. await ch.prefetch(1)
  81. logger.info(`订单支付结果轮询消费者已启动,队列:${ORDER_PAYMENT_QUEUE}`)
  82. orderPaymentWorkerStarted = true
  83. ch.consume(ORDER_PAYMENT_QUEUE, async (msg) => {
  84. if (!msg) return
  85. const content = JSON.parse(msg.content.toString() || '{}')
  86. const { orderId } = content
  87. if (!orderId) {
  88. logger.warn('收到无效的订单支付检查消息:缺少 orderId')
  89. ch.ack(msg)
  90. return
  91. }
  92. try {
  93. await pollOrderPaymentStatus(orderId, logger)
  94. ch.ack(msg)
  95. } catch (err) {
  96. logger.error(`订单支付轮询处理失败,订单号:${orderId},错误:${err.stack || err}`)
  97. ch.nack(msg, false, true)
  98. }
  99. })
  100. } catch (e) {
  101. logger.error(`启动订单支付 MQ 消费者失败:${e.stack || e}`)
  102. }
  103. }
  104. module.exports = {
  105. ORDER_PAYMENT_QUEUE,
  106. enqueueOrderPaymentCheck,
  107. startOrderPaymentWorker
  108. }