/** * RabbitMQ management API collector. * Enable: PEARDATA_RABBITMQ=1 * URL: PEARDATA_RABBITMQ_URL=http://guest:guest@127.0.0.1:15672 */ import { CollectorPlugin } from './plugin.js' import { registerChart } from '../../../shared/metrics.js' export function isRabbitmqEnabled() { const v = process.env.PEARDATA_RABBITMQ return v === '1' || v === 'on' || v === 'true' } function baseUrl() { return (process.env.PEARDATA_RABBITMQ_URL || 'http://guest:guest@127.0.0.1:15672').replace( /\/$/, '' ) } export class RabbitmqCollector extends CollectorPlugin { constructor() { super({ name: 'rabbitmq' }) this._prev = null this._prevTs = 0 } isEnabled() { return isRabbitmqEnabled() } async collect() { const ts = Date.now() const overview = await fetchJson(`${baseUrl()}/api/overview`) if (!overview) { registerChart({ id: 'rabbitmq.up', name: 'rabbitmq.up', context: 'rabbitmq.up', title: 'RabbitMQ up', units: 'boolean', family: 'rabbitmq', chartType: 'line', priority: 8600, plugin: 'rabbitmq', dimensions: [{ id: 'up', name: 'up', algorithm: 'absolute' }], }) return [{ chart: 'rabbitmq.up', context: 'rabbitmq.up', ts, values: { up: 0 } }] } const q = overview.queue_totals || {} const m = overview.message_stats || {} const dt = this._prevTs ? (ts - this._prevTs) / 1000 : 0 const prev = this._prev const rate = (key) => { const cur = m[key] || 0 if (!prev || dt <= 0) return 0 return Math.max(0, (cur - (prev[key] || 0)) / dt) } registerChart({ id: 'rabbitmq.up', name: 'rabbitmq.up', context: 'rabbitmq.up', title: 'RabbitMQ up', units: 'boolean', family: 'rabbitmq', chartType: 'line', priority: 8600, plugin: 'rabbitmq', dimensions: [{ id: 'up', name: 'up', algorithm: 'absolute' }], }) registerChart({ id: 'rabbitmq.queues', name: 'rabbitmq.queues', context: 'rabbitmq.queues', title: 'RabbitMQ queue messages', units: 'messages', family: 'rabbitmq', chartType: 'line', priority: 8610, plugin: 'rabbitmq', dimensions: [ { id: 'ready', name: 'ready', algorithm: 'absolute' }, { id: 'unacked', name: 'unacked', algorithm: 'absolute' }, { id: 'total', name: 'total', algorithm: 'absolute' }, ], }) registerChart({ id: 'rabbitmq.rates', name: 'rabbitmq.rates', context: 'rabbitmq.rates', title: 'RabbitMQ message rates', units: 'messages/s', family: 'rabbitmq', chartType: 'line', priority: 8620, plugin: 'rabbitmq', dimensions: [ { id: 'publish', name: 'publish', algorithm: 'incremental' }, { id: 'deliver', name: 'deliver', algorithm: 'incremental' }, { id: 'ack', name: 'ack', algorithm: 'incremental' }, ], }) this._prev = { publish: m.publish || 0, deliver: m.deliver_get || 0, ack: m.ack || 0 } this._prevTs = ts return [ { chart: 'rabbitmq.up', context: 'rabbitmq.up', ts, values: { up: 1 } }, { chart: 'rabbitmq.queues', context: 'rabbitmq.queues', ts, values: { ready: q.messages_ready || 0, unacked: q.messages_unacknowledged || 0, total: q.messages || 0, }, }, { chart: 'rabbitmq.rates', context: 'rabbitmq.rates', ts, values: { publish: rate('publish'), deliver: rate('deliver_get'), ack: rate('ack'), }, }, ] } } async function fetchJson(url) { try { const res = await fetch(url, { signal: AbortSignal.timeout(2000) }) if (!res.ok) return null return await res.json() } catch { return null } } let singleton = null export function getRabbitmqCollector() { if (!singleton) singleton = new RabbitmqCollector() return singleton }