Add prometheus exporter (#17392)

Create a metric endpoint and expose prometheus gauge for event stream
count.
This commit is contained in:
Thomas Trompette
2026-01-23 16:38:15 +01:00
committed by GitHub
parent 4c94e650a7
commit 2346efd71f
8 changed files with 124 additions and 9 deletions
@@ -1,4 +1,4 @@
import { Injectable } from '@nestjs/common';
import { Injectable, Logger, OnModuleInit } from '@nestjs/common';
import { type RecordGqlOperationSignature } from 'twenty-shared/types';
import { isDefined } from 'twenty-shared/utils';
@@ -9,6 +9,7 @@ import { WithLock } from 'src/engine/core-modules/cache-lock/with-lock.decorator
import { InjectCacheStorage } from 'src/engine/core-modules/cache-storage/decorators/cache-storage.decorator';
import { CacheStorageService } from 'src/engine/core-modules/cache-storage/services/cache-storage.service';
import { CacheStorageNamespace } from 'src/engine/core-modules/cache-storage/types/cache-storage-namespace.enum';
import { MetricsService } from 'src/engine/core-modules/metrics/metrics.service';
import { EVENT_STREAM_TTL_MS } from 'src/engine/subscriptions/constants/event-stream-ttl.constant';
import {
EventStreamException,
@@ -17,13 +18,38 @@ import {
import { type EventStreamData } from 'src/engine/subscriptions/types/event-stream-data.type';
@Injectable()
export class EventStreamService {
export class EventStreamService implements OnModuleInit {
private readonly logger = new Logger(EventStreamService.name);
constructor(
@InjectCacheStorage(CacheStorageNamespace.EngineSubscriptions)
private readonly cacheStorageService: CacheStorageService,
private readonly cacheLockService: CacheLockService,
private readonly metricsService: MetricsService,
) {}
onModuleInit() {
this.metricsService.createObservableGauge(
'twenty_event_streams_live_total',
{ description: 'Current number of live event streams' },
async (observableResult) => {
try {
const count = await this.getTotalActiveStreamCount();
observableResult.observe(count);
} catch (error) {
this.logger.error('Failed to collect event streams metrics', error);
}
},
);
}
async getTotalActiveStreamCount(): Promise<number> {
return this.cacheStorageService.scanAndCountSetMembers(
'workspace:*:activeStreams',
);
}
async createEventStream({
workspaceId,
eventStreamChannelId,
@@ -3,6 +3,7 @@ import { TypeOrmModule } from '@nestjs/typeorm';
import { CacheLockModule } from 'src/engine/core-modules/cache-lock/cache-lock.module';
import { CacheStorageModule } from 'src/engine/core-modules/cache-storage/cache-storage.module';
import { MetricsModule } from 'src/engine/core-modules/metrics/metrics.module';
import { RedisClientModule } from 'src/engine/core-modules/redis-client/redis-client.module';
import { WorkspaceEntity } from 'src/engine/core-modules/workspace/workspace.entity';
import { EventStreamService } from 'src/engine/subscriptions/event-stream.service';
@@ -13,6 +14,7 @@ import { SubscriptionService } from 'src/engine/subscriptions/subscription.servi
RedisClientModule,
CacheStorageModule,
CacheLockModule,
MetricsModule,
TypeOrmModule.forFeature([WorkspaceEntity]),
],
providers: [SubscriptionService, EventStreamService],