Add metrics for completed and failed jobs (#16969)

Add a counter for completed and failed jobs
This commit is contained in:
Thomas Trompette
2026-01-06 18:51:54 +01:00
committed by GitHub
parent 2c5a7570fc
commit 6043edd53f
6 changed files with 52 additions and 9 deletions
@@ -19,9 +19,11 @@ import { type MessageQueueJob } from 'src/engine/core-modules/message-queue/inte
import { type MessageQueueWorkerOptions } from 'src/engine/core-modules/message-queue/interfaces/message-queue-worker-options.interface';
import { QUEUE_RETENTION } from 'src/engine/core-modules/message-queue/constants/queue-retention.constants';
import { MESSAGE_QUEUE_PRIORITY } from 'src/engine/core-modules/message-queue/message-queue-priority.constant';
import { type MessageQueue } from 'src/engine/core-modules/message-queue/message-queue.constants';
import { getJobKey } from 'src/engine/core-modules/message-queue/utils/get-job-key.util';
import { MESSAGE_QUEUE_PRIORITY } from 'src/engine/core-modules/message-queue/message-queue-priority.constant';
import { type MetricsService } from 'src/engine/core-modules/metrics/metrics.service';
import { MetricsKeys } from 'src/engine/core-modules/metrics/types/metrics-keys.type';
export type BullMQDriverOptions = QueueOptions;
@@ -38,7 +40,10 @@ export class BullMQDriver implements MessageQueueDriver, OnModuleDestroy {
Worker
>;
constructor(private options: BullMQDriverOptions) {}
constructor(
private options: BullMQDriverOptions,
private metricsService: MetricsService,
) {}
register(queueName: MessageQueue): void {
this.queueMap[queueName] = new Queue(queueName, this.options);
@@ -89,6 +94,30 @@ export class BullMQDriver implements MessageQueueDriver, OnModuleDestroy {
},
workerOptions,
);
this.workerMap[queueName].on('completed', (job) => {
this.metricsService.incrementCounter({
key: MetricsKeys.JobCompleted,
attributes: { queue: queueName, job_name: job?.name ?? '' },
shouldStoreInCache: false,
});
});
this.workerMap[queueName].on('failed', (job, error) => {
if (!isDefined(job) || !isDefined(error)) {
return;
}
this.metricsService.incrementCounter({
key: MetricsKeys.JobFailed,
attributes: {
queue: queueName,
job_name: job.name,
error_type: error.name,
},
shouldStoreInCache: false,
});
});
}
async addCron<T>({
@@ -1,4 +1,5 @@
import { type BullMQDriverOptions } from 'src/engine/core-modules/message-queue/drivers/bullmq.driver';
import { type MetricsService } from 'src/engine/core-modules/metrics/metrics.service';
export enum MessageQueueDriverType {
BullMQ = 'bull-mq',
@@ -8,6 +9,7 @@ export enum MessageQueueDriverType {
export interface BullMQDriverFactoryOptions {
type: MessageQueueDriverType.BullMQ;
options: BullMQDriverOptions;
metricsService: MetricsService;
}
export interface SyncDriverFactoryOptions {
@@ -22,6 +22,7 @@ import {
} from 'src/engine/core-modules/message-queue/message-queue.module-definition';
import { MessageQueueService } from 'src/engine/core-modules/message-queue/services/message-queue.service';
import { getQueueToken } from 'src/engine/core-modules/message-queue/utils/get-queue-token.util';
import { MetricsModule } from 'src/engine/core-modules/metrics/metrics.module';
@Global()
@Module({})
@@ -77,6 +78,7 @@ export class MessageQueueCoreModule extends ConfigurableModuleClass {
return {
...dynamicModule,
imports: [...(dynamicModule.imports ?? []), MetricsModule],
providers: [
...(dynamicModule.providers ?? []),
driverProvider,
@@ -91,17 +93,17 @@ export class MessageQueueCoreModule extends ConfigurableModuleClass {
};
}
static async createDriver({ type, options }: typeof OPTIONS_TYPE) {
switch (type) {
static async createDriver(config: typeof OPTIONS_TYPE) {
switch (config.type) {
case MessageQueueDriverType.BullMQ: {
return new BullMQDriver(options);
return new BullMQDriver(config.options, config.metricsService);
}
case MessageQueueDriverType.Sync: {
return new SyncDriver();
}
default: {
this.logger.warn(
`Unsupported message queue driver type: ${type}. Using SyncDriver by default.`,
`Unsupported message queue driver type: ${(config as { type: string })?.type}. Using SyncDriver by default.`,
);
return new SyncDriver();
@@ -3,6 +3,7 @@ import {
MessageQueueDriverType,
type MessageQueueModuleOptions,
} from 'src/engine/core-modules/message-queue/interfaces';
import { type MetricsService } from 'src/engine/core-modules/metrics/metrics.service';
import { type RedisClientService } from 'src/engine/core-modules/redis-client/redis-client.service';
import { type TwentyConfigService } from 'src/engine/core-modules/twenty-config/twenty-config.service';
@@ -10,10 +11,13 @@ import { type TwentyConfigService } from 'src/engine/core-modules/twenty-config/
* MessageQueue Module factory
* @returns MessageQueueModuleOptions
* @param twentyConfigService
* @param redisClientService
* @param metricsService
*/
export const messageQueueModuleFactory = async (
_twentyConfigService: TwentyConfigService,
redisClientService: RedisClientService,
metricsService: MetricsService,
): Promise<MessageQueueModuleOptions> => {
const driverType = MessageQueueDriverType.BullMQ;
@@ -24,6 +28,7 @@ export const messageQueueModuleFactory = async (
options: {
connection: redisClientService.getQueueClient(),
},
metricsService,
} satisfies BullMQDriverFactoryOptions;
}
default: