From e4109a63f0d4c382fee089dbf274fc7d72188e2c Mon Sep 17 00:00:00 2001 From: Jackson Harper Date: Wed, 24 Jan 2024 12:15:38 +0800 Subject: [PATCH 1/6] Expose custom metrics for the queue processor --- packages/api/src/queue-processor.ts | 30 +++++++++++++++++++++++------ 1 file changed, 24 insertions(+), 6 deletions(-) diff --git a/packages/api/src/queue-processor.ts b/packages/api/src/queue-processor.ts index 7ccd5e272..b35cf7778 100644 --- a/packages/api/src/queue-processor.ts +++ b/packages/api/src/queue-processor.ts @@ -2,7 +2,7 @@ /* eslint-disable @typescript-eslint/restrict-template-expressions */ /* eslint-disable @typescript-eslint/require-await */ /* eslint-disable @typescript-eslint/no-misused-promises */ -import { Job, QueueEvents, Worker, Queue } from 'bullmq' +import { Job, QueueEvents, Worker, Queue, JobType } from 'bullmq' import express, { Express } from 'express' import { SnakeNamingStrategy } from 'typeorm-naming-strategies' import { appDataSource } from './data_source' @@ -63,6 +63,24 @@ const main = async () => { // respond healthy to auto-scaler. app.get('/_ah/health', (req, res) => res.sendStatus(200)) + app.get('/metrics', async (_, res) => { + const queue = await getBackendQueue() + if (!queue) { + res.sendStatus(400) + return + } + let output = '' + const metrics: JobType[] = ['active', 'failed', 'completed', 'prioritized'] + const counts = await queue.getJobCounts(...metrics) + console.log('counts: ', counts) + + metrics.forEach((metric, idx) => { + output += `omnivore_backend_queue_${metric}{} ${counts[metric]}\n` + }) + + res.status(200).setHeader('Content-Type', 'text/plain').send(output) + }) + const server = app.listen(port, () => { console.log(`[queue-processor]: started`) }) @@ -78,17 +96,17 @@ const main = async () => { throw '[queue-processor] error redis is not initialized' } - const queue = new Queue(QUEUE_NAME, { - connection: workerRedisClient, - }) + // Init the backend queue now + await getBackendQueue() const worker = new Worker( QUEUE_NAME, async (job: Job) => { switch (job.name) { case 'refresh-all-feeds': { - const counts = await queue.getJobCounts('wait') - if (counts.wait > 1000) { + const queue = await getBackendQueue() + const counts = await queue?.getJobCounts('prioritized') + if (counts && counts.wait > 1000) { return } return await refreshAllFeeds(appDataSource) From c29f31d618c01a2aa482a6c049f16d09d60af37c Mon Sep 17 00:00:00 2001 From: Jackson Harper Date: Wed, 24 Jan 2024 12:39:04 +0800 Subject: [PATCH 2/6] Add _count to prom metrics names --- packages/api/src/queue-processor.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/packages/api/src/queue-processor.ts b/packages/api/src/queue-processor.ts index b35cf7778..9f82e26e1 100644 --- a/packages/api/src/queue-processor.ts +++ b/packages/api/src/queue-processor.ts @@ -75,7 +75,7 @@ const main = async () => { console.log('counts: ', counts) metrics.forEach((metric, idx) => { - output += `omnivore_backend_queue_${metric}{} ${counts[metric]}\n` + output += `omnivore_backend_queue_${metric}_count{} ${counts[metric]}\n` }) res.status(200).setHeader('Content-Type', 'text/plain').send(output) From d24db92e031c16ebfffa4b3b0b2f416f989416b9 Mon Sep 17 00:00:00 2001 From: Jackson Harper Date: Wed, 24 Jan 2024 12:40:00 +0800 Subject: [PATCH 3/6] Dont force init this queue --- packages/api/src/queue-processor.ts | 3 --- 1 file changed, 3 deletions(-) diff --git a/packages/api/src/queue-processor.ts b/packages/api/src/queue-processor.ts index 9f82e26e1..69b495ea8 100644 --- a/packages/api/src/queue-processor.ts +++ b/packages/api/src/queue-processor.ts @@ -96,9 +96,6 @@ const main = async () => { throw '[queue-processor] error redis is not initialized' } - // Init the backend queue now - await getBackendQueue() - const worker = new Worker( QUEUE_NAME, async (job: Job) => { From c581372a93bf1b57313d6e2e7654dc21c84a36d2 Mon Sep 17 00:00:00 2001 From: Jackson Harper Date: Wed, 24 Jan 2024 21:13:13 +0800 Subject: [PATCH 4/6] Improve exported metrics naming convention --- packages/api/src/queue-processor.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/packages/api/src/queue-processor.ts b/packages/api/src/queue-processor.ts index 69b495ea8..4cfe45d85 100644 --- a/packages/api/src/queue-processor.ts +++ b/packages/api/src/queue-processor.ts @@ -75,7 +75,7 @@ const main = async () => { console.log('counts: ', counts) metrics.forEach((metric, idx) => { - output += `omnivore_backend_queue_${metric}_count{} ${counts[metric]}\n` + output += `omnivore_queue_messages_${metric}\{queue="${QUEUE_NAME}"\} ${counts[metric]}\n` }) res.status(200).setHeader('Content-Type', 'text/plain').send(output) From 5f9be385b7be047bd75c6016c325123d6983b06c Mon Sep 17 00:00:00 2001 From: Jackson Harper Date: Wed, 24 Jan 2024 21:34:51 +0800 Subject: [PATCH 5/6] Remove unneeded escapes --- packages/api/src/queue-processor.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/packages/api/src/queue-processor.ts b/packages/api/src/queue-processor.ts index 4cfe45d85..5aac7012d 100644 --- a/packages/api/src/queue-processor.ts +++ b/packages/api/src/queue-processor.ts @@ -75,7 +75,7 @@ const main = async () => { console.log('counts: ', counts) metrics.forEach((metric, idx) => { - output += `omnivore_queue_messages_${metric}\{queue="${QUEUE_NAME}"\} ${counts[metric]}\n` + output += `omnivore_queue_messages_${metric}{queue="${QUEUE_NAME}"} ${counts[metric]}\n` }) res.status(200).setHeader('Content-Type', 'text/plain').send(output) From a639a2fa102642214f39b476c23a52dc1ed84881 Mon Sep 17 00:00:00 2001 From: Jackson Harper Date: Thu, 25 Jan 2024 12:37:02 +0800 Subject: [PATCH 6/6] Add types to exported metrics --- packages/api/src/queue-processor.ts | 2 ++ 1 file changed, 2 insertions(+) diff --git a/packages/api/src/queue-processor.ts b/packages/api/src/queue-processor.ts index 5aac7012d..233bf36e7 100644 --- a/packages/api/src/queue-processor.ts +++ b/packages/api/src/queue-processor.ts @@ -69,12 +69,14 @@ const main = async () => { res.sendStatus(400) return } + let output = '' const metrics: JobType[] = ['active', 'failed', 'completed', 'prioritized'] const counts = await queue.getJobCounts(...metrics) console.log('counts: ', counts) metrics.forEach((metric, idx) => { + output += `# TYPE omnivore_queue_messages_${metric} gauge\n` output += `omnivore_queue_messages_${metric}{queue="${QUEUE_NAME}"} ${counts[metric]}\n` })