Add gauge for awaiting jobs (#17747)
<img width="774" height="333" alt="Capture d’écran 2026-02-05 à 16 18 06" src="https://github.com/user-attachments/assets/ec938a39-801d-4d71-a2de-7ced1795a0bd" />
This commit is contained in:
+29
-2
@@ -1,4 +1,8 @@
|
||||
import { Logger, type OnModuleDestroy } from '@nestjs/common';
|
||||
import {
|
||||
Logger,
|
||||
type OnModuleDestroy,
|
||||
type OnModuleInit,
|
||||
} from '@nestjs/common';
|
||||
|
||||
import {
|
||||
type JobsOptions,
|
||||
@@ -29,7 +33,9 @@ export type BullMQDriverOptions = QueueOptions;
|
||||
|
||||
const V4_LENGTH = 36;
|
||||
|
||||
export class BullMQDriver implements MessageQueueDriver, OnModuleDestroy {
|
||||
export class BullMQDriver
|
||||
implements MessageQueueDriver, OnModuleDestroy, OnModuleInit
|
||||
{
|
||||
private logger = new Logger(BullMQDriver.name);
|
||||
private queueMap: Record<MessageQueue, Queue> = {} as Record<
|
||||
MessageQueue,
|
||||
@@ -45,6 +51,27 @@ export class BullMQDriver implements MessageQueueDriver, OnModuleDestroy {
|
||||
private metricsService: MetricsService,
|
||||
) {}
|
||||
|
||||
onModuleInit() {
|
||||
this.metricsService.createObservableGauge(
|
||||
'twenty_queue_jobs_waiting_total',
|
||||
{ description: 'Current number of jobs waiting in queue' },
|
||||
async (observableResult) => {
|
||||
for (const [queueName, queue] of Object.entries(this.queueMap)) {
|
||||
try {
|
||||
const waitingCount = await queue.count();
|
||||
|
||||
observableResult.observe(waitingCount, { queue: queueName });
|
||||
} catch (error) {
|
||||
this.logger.error(
|
||||
`Failed to collect waiting jobs metrics for queue ${queueName}`,
|
||||
error,
|
||||
);
|
||||
}
|
||||
}
|
||||
},
|
||||
);
|
||||
}
|
||||
|
||||
register(queueName: MessageQueue): void {
|
||||
this.queueMap[queueName] = new Queue(queueName, this.options);
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user