diff --git a/packages/twenty-server/src/engine/core-modules/message-queue/drivers/bullmq.driver.ts b/packages/twenty-server/src/engine/core-modules/message-queue/drivers/bullmq.driver.ts index 13b6e4782f..5daf3e0134 100644 --- a/packages/twenty-server/src/engine/core-modules/message-queue/drivers/bullmq.driver.ts +++ b/packages/twenty-server/src/engine/core-modules/message-queue/drivers/bullmq.driver.ts @@ -54,17 +54,23 @@ export class BullMQDriver ) {} onModuleInit() { - this.metricsService.createObservableGauge({ + this.metricsService.createMultiObservableGauge({ metricName: 'twenty_queue_jobs_waiting_total', options: { description: 'Current number of jobs waiting in queue' }, callback: async () => { - let totalWaiting = 0; + const observations: Array<{ + value: number; + attributes: { queue: string }; + }> = []; for (const [queueName, queue] of Object.entries(this.queueMap)) { try { const waitingCount = await queue.count(); - totalWaiting += waitingCount; + observations.push({ + value: waitingCount, + attributes: { queue: queueName }, + }); } catch (error) { this.logger.error( `Failed to collect waiting jobs metrics for queue ${queueName}`, @@ -73,7 +79,7 @@ export class BullMQDriver } } - return totalWaiting; + return observations; }, }); } diff --git a/packages/twenty-server/src/engine/core-modules/metrics/metrics.service.ts b/packages/twenty-server/src/engine/core-modules/metrics/metrics.service.ts index 668ad58402..4d6a538ded 100644 --- a/packages/twenty-server/src/engine/core-modules/metrics/metrics.service.ts +++ b/packages/twenty-server/src/engine/core-modules/metrics/metrics.service.ts @@ -6,6 +6,7 @@ import { type Meter, type MetricOptions, type ObservableGauge, + type ObservableResult, } from '@opentelemetry/api'; import { isDefined } from 'twenty-shared/utils'; @@ -42,16 +43,63 @@ export class MetricsService { options: MetricOptions; callback: () => number | Promise; cacheValue?: boolean; + }): ObservableGauge { + return this.createObservableGaugeInternal({ + metricName, + options, + callback, + cacheValue, + observeResult: (observableResult, result) => { + observableResult.observe(result); + }, + }); + } + + createMultiObservableGauge({ + metricName, + options, + callback, + cacheValue = false, + }: { + metricName: string; + options: MetricOptions; + callback: () => Promise>; + cacheValue?: boolean; + }): ObservableGauge { + return this.createObservableGaugeInternal({ + metricName, + options, + callback, + cacheValue, + observeResult: (observableResult, observations) => { + for (const observation of observations) { + observableResult.observe(observation.value, observation.attributes); + } + }, + }); + } + + private createObservableGaugeInternal({ + metricName, + options, + callback, + cacheValue, + observeResult, + }: { + metricName: string; + options: MetricOptions; + callback: () => T | Promise; + cacheValue: boolean; + observeResult: (observableResult: ObservableResult, result: T) => void; }): ObservableGauge { const gauge = this.getMeter().createObservableGauge(metricName, options); gauge.addCallback(async (observableResult) => { if (cacheValue) { - const cachedResult = - await this.healthCacheStorage.get(metricName); + const cachedResult = await this.healthCacheStorage.get(metricName); if (isDefined(cachedResult)) { - observableResult.observe(cachedResult); + observeResult(observableResult, cachedResult); return; } @@ -60,7 +108,7 @@ export class MetricsService { try { const result = await callback(); - observableResult.observe(result); + observeResult(observableResult, result); if (cacheValue) { await this.healthCacheStorage.set(