lepaoSchedulePublisher.js 3.5 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106
  1. const mq = require('./index')
  2. const db = require('../DataBase/db')
  3. const { publishRunforgeTask } = require('./runforgeTaskMq')
  4. const {
  5. popDueMessages,
  6. stripScheduleMeta,
  7. requeueAt,
  8. SCHEDULE_KEY
  9. } = require('./lepaoAutoScheduleRedis')
  10. const mqNames = require('./jkesMqNames')
  11. const { prepareAutoRunTask } = require('../jkes/prepareAutoRunTask')
  12. let intervalHandle = null
  13. const ACCOUNT_SQL = `
  14. SELECT student_num, auto_day, token, target_count,
  15. auto_run_distance_min_km, auto_run_distance_max_km, pace_min_sec_per_km, pace_max_sec_per_km
  16. FROM lepao_account
  17. WHERE student_num = ? AND auto_run = 1 AND state = 1
  18. LIMIT 1
  19. `
  20. async function refreshDelayedStartRunPayload(msg, logger) {
  21. const account = msg?.data?.account
  22. if (!account || msg?.type !== 'lepao.startRun') {
  23. return msg
  24. }
  25. const rows = await db.query(ACCOUNT_SQL, [account])
  26. if (!rows?.length) {
  27. logger.info?.(`[LepaoSchedule] 账号 ${account} 不可用,跳过延迟任务`)
  28. return null
  29. }
  30. const prepared = await prepareAutoRunTask(rows[0])
  31. if (!prepared.run) {
  32. logger.info?.(`[LepaoSchedule] ${account} 延迟任务跳过:${prepared.reason}`)
  33. return null
  34. }
  35. return {
  36. ...msg,
  37. data: {
  38. ...msg.data,
  39. targetKm: prepared.task.targetKm,
  40. autoDoubleSlot: prepared.task.autoDoubleSlot,
  41. paceRandomMinSecPerKm: prepared.task.paceRandomMinSecPerKm,
  42. paceRandomMaxSecPerKm: prepared.task.paceRandomMaxSecPerKm
  43. }
  44. }
  45. }
  46. /**
  47. * 定时将 Redis 中已到期的 JKES 乐跑任务写入主任务队列(jkes_runforge_task_queue,支持 mqPrefix)
  48. */
  49. function startLepaoSchedulePublisher(options = {}) {
  50. const logger = options.logger || console
  51. const intervalMs = options.intervalMs ?? 2000
  52. const batch = options.batch ?? 100
  53. if (intervalHandle) return
  54. intervalHandle = setInterval(async () => {
  55. try {
  56. const now = Date.now()
  57. const rawList = await popDueMessages(now, batch)
  58. if (!rawList.length) return
  59. const channel = await mq.getChannel(mqNames.channelScheduleTick)
  60. for (const raw of rawList) {
  61. let msg
  62. try {
  63. msg = JSON.parse(raw)
  64. } catch (e) {
  65. logger.error?.(
  66. `[LepaoSchedule] 调度 JSON 无效已丢弃: ${String(raw).slice(0, 120)}`
  67. )
  68. continue
  69. }
  70. try {
  71. const refreshed = await refreshDelayedStartRunPayload(msg, logger)
  72. if (!refreshed) continue
  73. publishRunforgeTask(channel, stripScheduleMeta(refreshed))
  74. } catch (e) {
  75. logger.error?.(`[LepaoSchedule] MQ 投递失败,5s 后重试: ${e.message || e}`)
  76. const retryAt = now + 5000
  77. try {
  78. await requeueAt(retryAt, msg)
  79. } catch (re2) {
  80. logger.error?.(`[LepaoSchedule] 写回 Redis 失败: ${re2.message || re2}`)
  81. }
  82. }
  83. }
  84. } catch (e) {
  85. logger.error?.(`[LepaoSchedule] tick 异常: ${e.message || e}`)
  86. }
  87. }, intervalMs)
  88. logger.info?.(
  89. `[LepaoSchedule] 已启动 Redis→MQ(JKES)调度(间隔 ${intervalMs}ms,每批最多 ${batch} 条,key=${SCHEDULE_KEY})`
  90. )
  91. }
  92. module.exports = { startLepaoSchedulePublisher }