GetQueueTasks.js 10 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297
  1. const API = require('../../../lib/API')
  2. const axios = require('axios')
  3. const mq = require('../../../plugin/mq')
  4. const config = require('../../../config.json')
  5. const AccessControl = require('../../../lib/AccessControl')
  6. const { BaseStdResponse } = require('../../../BaseStdResponse')
  7. const {
  8. SCHEDULE_KEY,
  9. listPendingScheduledForAdmin,
  10. countPendingScheduled
  11. } = require('../../../plugin/mq/lepaoAutoScheduleRedis')
  12. const { mq: mqPrefixName, PREFIX } = require('../../../plugin/mq/mqPrefix')
  13. const { QUEUE_BASE_NAMES } = require('../../../plugin/mq/jkesMqNames')
  14. const { TASK_QUEUE } = require('../../../plugin/mq/runforgeTaskMq')
  15. /** 允许通过管理接口查看的队列(防任意队列名探测) */
  16. const ALLOWED_QUEUES = QUEUE_BASE_NAMES.map(mqPrefixName)
  17. function canonicalQueueName(q) {
  18. const raw = q || TASK_QUEUE
  19. if (ALLOWED_QUEUES.includes(raw)) return raw
  20. if (PREFIX && QUEUE_BASE_NAMES.includes(raw)) return mqPrefixName(raw)
  21. return raw
  22. }
  23. function parseAmqpHttpBase(amqpUrl) {
  24. const u = new URL(String(amqpUrl).replace(/^amqp:/, 'http:'))
  25. return {
  26. user: decodeURIComponent(u.username || ''),
  27. password: decodeURIComponent(u.password || ''),
  28. hostname: u.hostname
  29. }
  30. }
  31. function managementEndpoint(queueName) {
  32. const rm = config.rabbitmq || {}
  33. const creds = parseAmqpHttpBase(rm.url || '')
  34. const base = (rm.managementBaseUrl || `http://${creds.hostname}:15672`).replace(/\/$/, '')
  35. const vhost = encodeURIComponent(rm.vhost != null ? rm.vhost : '/')
  36. const q = encodeURIComponent(queueName)
  37. return {
  38. url: `${base}/api/queues/${vhost}/${q}/get`,
  39. auth: {
  40. username: rm.managementUser || creds.user,
  41. password: rm.managementPassword || creds.password
  42. }
  43. }
  44. }
  45. function decodePayload(msg) {
  46. const enc = msg.payload_encoding || 'string'
  47. let raw = msg.payload
  48. if (enc === 'base64' && typeof raw === 'string') {
  49. try {
  50. raw = Buffer.from(raw, 'base64').toString('utf8')
  51. } catch {
  52. return { raw: msg.payload, encoding: enc, parseError: true }
  53. }
  54. }
  55. if (typeof raw === 'string') {
  56. try {
  57. return { body: JSON.parse(raw), encoding: enc }
  58. } catch {
  59. return { body: raw, encoding: enc }
  60. }
  61. }
  62. return { body: raw, encoding: enc }
  63. }
  64. /**
  65. * 通过 Management API 窥视队列消息(reject_requeue_true:看完后重新入队,不消费)
  66. */
  67. async function peekQueueMessages(queueName, limit) {
  68. const { url, auth } = managementEndpoint(queueName)
  69. const { data, status } = await axios.post(
  70. url,
  71. {
  72. count: limit,
  73. ackmode: 'reject_requeue_true',
  74. encoding: 'auto'
  75. },
  76. {
  77. auth,
  78. timeout: 12000,
  79. validateStatus: () => true
  80. }
  81. )
  82. if (status >= 400) {
  83. const reason =
  84. typeof data === 'object' && data !== null
  85. ? data.reason || data.error || JSON.stringify(data)
  86. : String(data)
  87. const err = new Error(reason || `Management HTTP ${status}`)
  88. err.status = status
  89. throw err
  90. }
  91. const list = Array.isArray(data) ? data : []
  92. return list.map((m) => {
  93. const decoded = decodePayload(m)
  94. return {
  95. redelivered: m.redelivered,
  96. routing_key: m.routing_key,
  97. exchange: m.exchange,
  98. properties: m.properties
  99. ? {
  100. messageId: m.properties.message_id,
  101. timestamp: m.properties.timestamp,
  102. contentType: m.properties.content_type,
  103. headers: m.properties.headers
  104. }
  105. : undefined,
  106. payload: decoded.body,
  107. payload_encoding: decoded.encoding
  108. }
  109. })
  110. }
  111. class GetQueueTasks extends API {
  112. constructor() {
  113. super()
  114. this.setPath('/Admin/MQ/GetQueueTasks')
  115. this.setMethod('get')
  116. }
  117. async onRequest(req, res) {
  118. const {
  119. uuid,
  120. session,
  121. queue,
  122. limit: limitStr,
  123. summary,
  124. includeScheduled,
  125. scheduledLimit: scheduledLimitStr
  126. } = req.query
  127. if ([uuid, session].some((v) => v === '' || v == null))
  128. return res.json({
  129. ...BaseStdResponse.MISSING_PARAMETER
  130. })
  131. if (!(await AccessControl.checkSession(uuid, session)))
  132. return res.status(401).json({
  133. ...BaseStdResponse.ACCESS_DENIED
  134. })
  135. const permission = await AccessControl.getPermission(uuid)
  136. if (!permission.includes('admin') && !permission.includes('service'))
  137. return res.json({
  138. ...BaseStdResponse.PERMISSION_DENIED
  139. })
  140. const wantSummary = summary === '1' || summary === 'true'
  141. try {
  142. const ch = await mq.getChannel('admin_queue_inspect')
  143. if (wantSummary) {
  144. const queues = {}
  145. for (const name of ALLOWED_QUEUES) {
  146. try {
  147. const info = await ch.checkQueue(name)
  148. queues[name] = {
  149. messageCount: info.messageCount,
  150. consumerCount: info.consumerCount
  151. }
  152. } catch (e) {
  153. queues[name] = {
  154. messageCount: null,
  155. consumerCount: null,
  156. error: e.message || String(e)
  157. }
  158. }
  159. }
  160. const slimit = Math.min(
  161. 2000,
  162. Math.max(1, parseInt(scheduledLimitStr, 10) || 800)
  163. )
  164. const pendingCount = await countPendingScheduled(Date.now())
  165. const scheduledMirror = await listPendingScheduledForAdmin(Date.now(), slimit)
  166. return res.json({
  167. ...BaseStdResponse.OK,
  168. data: {
  169. summary: true,
  170. queues,
  171. redisScheduler: {
  172. keys: [SCHEDULE_KEY],
  173. pendingCount,
  174. note: `到期任务由本服务定时写入 ${TASK_QUEUE}。`
  175. },
  176. autoRunScheduledMirror: {
  177. pendingCount,
  178. note: scheduledMirror.note,
  179. sample: scheduledMirror.items.slice(0, 20)
  180. },
  181. fetchedAt: Date.now()
  182. }
  183. })
  184. }
  185. const queueName = canonicalQueueName(queue)
  186. if (!ALLOWED_QUEUES.includes(queueName))
  187. return res.json({
  188. ...BaseStdResponse.ERR,
  189. msg: '不支持的队列名称'
  190. })
  191. let limit = parseInt(limitStr, 10)
  192. if (Number.isNaN(limit) || limit < 1) limit = 30
  193. if (limit > 100) limit = 100
  194. let messageCount = null
  195. let consumerCount = null
  196. try {
  197. const info = await ch.checkQueue(queueName)
  198. messageCount = info.messageCount
  199. consumerCount = info.consumerCount
  200. } catch (e) {
  201. return res.json({
  202. ...BaseStdResponse.ERR,
  203. msg: `无法访问队列:${e.message || e}`
  204. })
  205. }
  206. let tasks = []
  207. let managementError = null
  208. try {
  209. tasks = await peekQueueMessages(queueName, limit)
  210. } catch (e) {
  211. managementError =
  212. e.status === 401 || e.status === 403
  213. ? 'Management 鉴权失败,请在 config.json 的 rabbitmq 中配置 managementUser / managementPassword,或检查管理插件用户权限'
  214. : e.code === 'ECONNREFUSED' || e.code === 'ETIMEDOUT'
  215. ? '无法连接 RabbitMQ Management(默认 15672)。请开启管理插件并开放端口,或配置 rabbitmq.managementBaseUrl'
  216. : (e.message || String(e))
  217. this.logger.warn(`[GetQueueTasks] Management 窥视失败: ${managementError}`)
  218. }
  219. const wantScheduled =
  220. includeScheduled !== '0' &&
  221. includeScheduled !== 'false' &&
  222. queueName === TASK_QUEUE
  223. let autoRunScheduledMirror = null
  224. let pendingScheduledCount = null
  225. if (queueName === TASK_QUEUE) {
  226. pendingScheduledCount = await countPendingScheduled(Date.now())
  227. if (wantScheduled) {
  228. const slimit = Math.min(
  229. 500,
  230. Math.max(1, parseInt(scheduledLimitStr, 10) || 200)
  231. )
  232. autoRunScheduledMirror = await listPendingScheduledForAdmin(Date.now(), slimit)
  233. }
  234. }
  235. const detail = {
  236. queue: queueName,
  237. messageCount,
  238. consumerCount,
  239. peekLimit: limit,
  240. tasks,
  241. managementError,
  242. redisScheduler:
  243. queueName === TASK_QUEUE
  244. ? {
  245. keys: [SCHEDULE_KEY],
  246. pendingCount: pendingScheduledCount
  247. }
  248. : undefined,
  249. autoRunScheduledMirror,
  250. fetchedAt: Date.now()
  251. }
  252. if (queueName === TASK_QUEUE) {
  253. detail.peekNote =
  254. `tasks:已在主队列中的消息;autoRunScheduledMirror:尚未到 fireAt、仍只在 Redis 中的调度(到期后由服务写入 ${TASK_QUEUE})。`
  255. }
  256. return res.json({
  257. ...BaseStdResponse.OK,
  258. data: detail
  259. })
  260. } catch (error) {
  261. this.logger.error(`GetQueueTasks: ${error.stack || error}`)
  262. return res.json({
  263. ...BaseStdResponse.ERR,
  264. msg: error.message || '查询失败'
  265. })
  266. }
  267. }
  268. }
  269. module.exports.GetQueueTasks = GetQueueTasks