orderPaymentWorker.js 6.8 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180
  1. const db = require('../DataBase/db')
  2. const config = require('../../config.json')
  3. const mq = require('./index')
  4. const { ORDER_PAYMENT_QUEUE } = require('./jkesMqNames')
  5. const { insertLedgerRecord } = require('../../lib/Lepao/CountLedger')
  6. const { releaseUsageForOrder } = require('../../lib/CouponService')
  7. const { queryPaymentOrder } = require('../../lib/PaymentClient')
  8. let orderPaymentWorkerStarted = false
  9. async function writePurchaseLedger(orderId, userUuid, addCount, logger) {
  10. const delta = Number(addCount || 0)
  11. if (!orderId || !userUuid || delta === 0) return
  12. try {
  13. const userRows = await db.query(
  14. 'SELECT lepao_count FROM users WHERE uuid = ?',
  15. [userUuid]
  16. )
  17. if (!userRows || userRows.length !== 1) return
  18. const afterCount = Number(userRows[0].lepao_count || 0)
  19. const beforeCount = afterCount - delta
  20. await insertLedgerRecord({
  21. userUuid,
  22. delta,
  23. balanceBefore: beforeCount,
  24. balanceAfter: afterCount,
  25. bizType: 'purchase',
  26. bizId: orderId,
  27. remark: `订单号:${orderId}`
  28. })
  29. } catch (error) {
  30. logger?.error?.(`写入购买里程流水失败 ${orderId}: ${error.stack || error}`)
  31. }
  32. }
  33. async function pollOrderPaymentStatus(orderId, logger) {
  34. const paymentConfig = config.pay || {}
  35. if (!paymentConfig.pid || !paymentConfig.url || !paymentConfig.key) {
  36. logger.error('支付配置错误,无法轮询易支付状态')
  37. return
  38. }
  39. const MAX_RETRIES = 120
  40. const DELAY = 2500
  41. const pollOrderStatus = async (retry = 0) => {
  42. if (retry >= MAX_RETRIES) {
  43. const closeRes = await db.query(
  44. 'UPDATE orders SET state = 3 WHERE orderId = ? AND state = 0',
  45. [orderId]
  46. )
  47. if (closeRes?.affectedRows > 0) {
  48. await releaseUsageForOrder(orderId)
  49. logger.info(`订单超时未支付,自动取消,订单号:${orderId}`)
  50. }
  51. return
  52. }
  53. try {
  54. const existing = await db.query('SELECT state FROM orders WHERE orderId = ?', [orderId])
  55. if (!existing?.length) {
  56. logger.warn(`订单不存在,停止轮询:${orderId}`)
  57. return
  58. }
  59. if (Number(existing[0].state) !== 0) {
  60. logger.info(`订单已处理(state=${existing[0].state}),停止轮询:${orderId}`)
  61. return
  62. }
  63. const queryData = await queryPaymentOrder(orderId, logger)
  64. logger.info(`轮询订单支付状态,订单号:${orderId},尝试次数:${retry + 1},查询结果:${JSON.stringify(queryData)}`)
  65. if (queryData.code == 1 && queryData.status == 1) {
  66. const { trade_no, out_trade_no, type } = queryData
  67. const time = Date.now()
  68. let sql = 'UPDATE orders SET state = 1, pay_type = ?, pay_id = ?, pay_time = ? WHERE orderId = ? AND state = 0'
  69. const result = await db.query(sql, [type, trade_no, time, out_trade_no])
  70. if (result.affectedRows > 0) {
  71. sql = `
  72. SELECT g.lepao_count, a.create_user
  73. FROM orders a
  74. LEFT JOIN goods g ON a.goods_id = g.id
  75. WHERE a.orderId = ?
  76. `
  77. const rows = await db.query(sql, [out_trade_no])
  78. if (!rows || rows.length !== 1) {
  79. logger.error(`订单商品信息异常,订单号:${out_trade_no}`)
  80. await db.query('UPDATE orders SET state = 4 WHERE orderId = ?', [out_trade_no])
  81. return
  82. }
  83. const { lepao_count, create_user } = rows[0]
  84. sql = 'UPDATE users SET lepao_count = lepao_count + ? WHERE uuid = ?'
  85. const updateUser = await db.query(sql, [lepao_count, create_user])
  86. if (!updateUser || updateUser.affectedRows !== 1) {
  87. logger.error(`更新用户失败,UUID: ${create_user}`)
  88. await db.query('UPDATE orders SET state = 4 WHERE orderId = ?', [out_trade_no])
  89. }
  90. sql = 'UPDATE orders SET state = 2 WHERE orderId = ?'
  91. await db.query(sql, [out_trade_no])
  92. await writePurchaseLedger(out_trade_no, create_user, lepao_count, logger)
  93. logger.info(`订单处理成功:${out_trade_no}`)
  94. return
  95. }
  96. logger.info(`支付网关已确认付款,订单已由其他进程/回调处理,停止轮询:${out_trade_no}`)
  97. return
  98. }
  99. setTimeout(() => pollOrderStatus(retry + 1), DELAY)
  100. } catch (error) {
  101. logger.warn(`轮询支付状态失败,订单号:${orderId},原因:${error.message || error}`)
  102. setTimeout(() => pollOrderStatus(retry + 1), DELAY)
  103. }
  104. }
  105. logger.info(`开始轮询订单支付状态,订单号:${orderId}`)
  106. pollOrderStatus()
  107. }
  108. async function enqueueOrderPaymentCheck(orderId) {
  109. const ch = await mq.getChannel('order_payment')
  110. await ch.assertQueue(ORDER_PAYMENT_QUEUE, { durable: true })
  111. ch.sendToQueue(
  112. ORDER_PAYMENT_QUEUE,
  113. Buffer.from(JSON.stringify({ orderId, enqueueTime: Date.now() })),
  114. { persistent: true }
  115. )
  116. }
  117. async function startOrderPaymentWorker(logger) {
  118. if (orderPaymentWorkerStarted) {
  119. return
  120. }
  121. try {
  122. const ch = await mq.getChannel('order_payment')
  123. await ch.assertQueue(ORDER_PAYMENT_QUEUE, { durable: true })
  124. await ch.prefetch(1)
  125. logger.info(`订单支付结果轮询消费者已启动,队列:${ORDER_PAYMENT_QUEUE}`)
  126. orderPaymentWorkerStarted = true
  127. ch.consume(ORDER_PAYMENT_QUEUE, async (msg) => {
  128. if (!msg) return
  129. const content = JSON.parse(msg.content.toString() || '{}')
  130. const { orderId } = content
  131. if (!orderId) {
  132. logger.warn('收到无效的订单支付检查消息(缺少 orderId)')
  133. ch.ack(msg)
  134. return
  135. }
  136. try {
  137. await pollOrderPaymentStatus(orderId, logger)
  138. ch.ack(msg)
  139. } catch (err) {
  140. logger.error(`订单支付轮询处理失败,订单号:${orderId},错误:${err.stack || err}`)
  141. ch.nack(msg, false, true)
  142. }
  143. })
  144. } catch (e) {
  145. logger.error(`启动订单支付 MQ 消费者失败:${e.stack || e}`)
  146. }
  147. }
  148. module.exports = {
  149. ORDER_PAYMENT_QUEUE,
  150. enqueueOrderPaymentCheck,
  151. startOrderPaymentWorker
  152. }