lepaoAutoScheduleRedis.js 3.4 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125
  1. const Redis = require('../DataBase/Redis')
  2. const { TASK_QUEUE } = require('./runforgeTaskMq')
  3. /** JKES 延迟调度 ZSET(唯一) */
  4. const SCHEDULE_KEY = 'jkes_lepao:mq:scheduled'
  5. const POP_DUE_LUA = `
  6. local key = KEYS[1]
  7. local now = tonumber(ARGV[1])
  8. local limit = tonumber(ARGV[2])
  9. local items = redis.call('ZRANGEBYSCORE', key, '-inf', now, 'LIMIT', 0, limit)
  10. for i = 1, #items do
  11. redis.call('ZREM', key, items[i])
  12. end
  13. return items
  14. `
  15. function stripScheduleMeta(msg) {
  16. if (!msg || typeof msg !== 'object') return msg
  17. const copy = { ...msg }
  18. delete copy._scheduleMeta
  19. return copy
  20. }
  21. async function scheduleDelayedRunforgeTask(fireAt, messageObject, meta = null) {
  22. const toStore =
  23. meta != null
  24. ? {
  25. ...messageObject,
  26. _scheduleMeta: {
  27. ...meta,
  28. fireAt
  29. }
  30. }
  31. : messageObject
  32. const member = JSON.stringify(toStore)
  33. await Redis.sendCommand(['ZADD', SCHEDULE_KEY, String(fireAt), member])
  34. }
  35. async function requeueAt(fireAt, messageObject) {
  36. const member = JSON.stringify(messageObject)
  37. await Redis.sendCommand(['ZADD', SCHEDULE_KEY, String(fireAt), member])
  38. }
  39. async function popDueMessages(now = Date.now(), limit = 100) {
  40. const res = await Redis.sendCommand([
  41. 'EVAL',
  42. POP_DUE_LUA,
  43. '1',
  44. SCHEDULE_KEY,
  45. String(now),
  46. String(limit)
  47. ])
  48. if (!Array.isArray(res)) return []
  49. return res.map((x) => (Buffer.isBuffer(x) ? x.toString('utf8') : String(x)))
  50. }
  51. async function pruneStaleScheduled(beforeScore, now = Date.now()) {
  52. await Redis.sendCommand(['ZREMRANGEBYSCORE', SCHEDULE_KEY, '-inf', String(beforeScore)])
  53. }
  54. async function listPendingScheduledForAdmin(now = Date.now(), limitTotal = 800) {
  55. await pruneStaleScheduled(now - 48 * 3600 * 1000, now)
  56. const raw = await Redis.sendCommand([
  57. 'ZRANGEBYSCORE',
  58. SCHEDULE_KEY,
  59. `(${String(now)}`,
  60. '+inf',
  61. 'WITHSCORES',
  62. 'LIMIT',
  63. '0',
  64. String(limitTotal)
  65. ])
  66. const items = []
  67. for (let i = 0; i < raw.length; i += 2) {
  68. const score = Number(raw[i + 1])
  69. const value = Buffer.isBuffer(raw[i]) ? raw[i].toString('utf8') : raw[i]
  70. let parsed
  71. try {
  72. parsed = JSON.parse(value)
  73. } catch {
  74. parsed = { raw: value }
  75. }
  76. const meta = parsed._scheduleMeta || {}
  77. const fireAt = meta.fireAt != null ? meta.fireAt : score
  78. items.push({
  79. taskId: parsed.id,
  80. type: parsed.type,
  81. account: parsed.data?.account,
  82. name: meta.name,
  83. fireAt,
  84. score,
  85. delayMs: meta.delayMs,
  86. remainMs: Math.max(0, fireAt - now),
  87. payloadPreview: stripScheduleMeta(parsed)
  88. })
  89. }
  90. return {
  91. items,
  92. note: `Redis 调度:未到 fireAt 前不会进入 ${TASK_QUEUE};到期由定时任务写入 MQ。`
  93. }
  94. }
  95. async function countPendingScheduled(now = Date.now()) {
  96. const n = await Redis.sendCommand([
  97. 'ZCOUNT',
  98. SCHEDULE_KEY,
  99. `(${String(now)}`,
  100. '+inf'
  101. ])
  102. return Number(n) || 0
  103. }
  104. module.exports = {
  105. SCHEDULE_KEY,
  106. stripScheduleMeta,
  107. scheduleDelayedRunforgeTask,
  108. requeueAt,
  109. popDueMessages,
  110. listPendingScheduledForAdmin,
  111. countPendingScheduled
  112. }