Skip to content

Commit 41d90a4

Browse files
feat(backend): add Bull Board metrics history (#1604)
* feat(backend): add Bull Board metrics history * docs: add Bull Board metrics changelog entry * feedback
1 parent 251c8e2 commit 41d90a4

6 files changed

Lines changed: 136 additions & 105 deletions

File tree

packages/backend/package.json

Lines changed: 4 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -22,9 +22,10 @@
2222
"vitest": "^4.1.4"
2323
},
2424
"dependencies": {
25-
"@bull-board/api": "6.11.2",
26-
"@bull-board/express": "6.11.2",
27-
"@bull-board/ui": "6.11.2",
25+
"@bull-board/api": "9.0.0",
26+
"@bull-board/express": "9.0.0",
27+
"@bull-board/metrics": "9.0.0",
28+
"@bull-board/ui": "9.0.0",
2829
"@coderabbitai/bitbucket": "^1.1.3",
2930
"@gitbeaker/rest": "^40.5.1",
3031
"@octokit/app": "^16.1.1",

packages/backend/src/api.ts

Lines changed: 31 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,12 +1,17 @@
11
import { createBullBoard } from '@bull-board/api';
2-
import { BullMQAdapter } from '@bull-board/api/bullMQAdapter.js';
2+
import { BullMQAdapter } from '@bull-board/api/bullMQAdapter';
33
import { ExpressAdapter } from '@bull-board/express';
4+
import {
5+
MetricsRecorder,
6+
RedisMetricsHistoryProvider,
7+
} from '@bull-board/metrics';
48
import { Octokit } from '@octokit/rest';
59
import * as Sentry from "@sentry/node";
610
import { PrismaClient, RepoIndexingJobType } from '@sourcebot/db';
711
import { createLogger, env, JOB_PRIORITIES } from '@sourcebot/shared';
812
import express, { NextFunction, Request, Response } from 'express';
913
import 'express-async-errors';
14+
import type { Redis } from 'ioredis';
1015
import * as http from "http";
1116
import z from 'zod';
1217
import { SINGLE_TENANT_ORG_ID } from './constants.js';
@@ -22,21 +27,41 @@ const PORT = Number(workerApiUrl.port) || (workerApiUrl.protocol === "https:" ?
2227

2328
export class Api {
2429
private server: http.Server;
30+
private metricsRecorder: MetricsRecorder;
31+
private metricsHistoryProvider: RedisMetricsHistoryProvider;
2532

2633
constructor(
2734
promClient: PromClient,
2835
private prisma: PrismaClient,
2936
private jobManager: JobManager,
37+
redis: Redis,
3038
) {
3139
const app = express();
3240
app.use(express.json());
3341
app.use(express.urlencoded({ extended: true }));
3442

3543
const bullBoardAdapter = new ExpressAdapter();
3644
bullBoardAdapter.setBasePath('/admin/queues');
45+
const queueAdapters = jobManager
46+
.getQueues()
47+
.map(queue => new BullMQAdapter(queue));
48+
this.metricsHistoryProvider = new RedisMetricsHistoryProvider({
49+
connection: redis,
50+
});
51+
this.metricsRecorder = new MetricsRecorder({
52+
queues: queueAdapters,
53+
connection: redis,
54+
onLatencyError: (error, queueName) => {
55+
logger.error(
56+
`Failed to record BullMQ latency metrics for queue "${queueName}"`,
57+
error,
58+
);
59+
},
60+
});
3761
createBullBoard({
38-
queues: jobManager.getQueues().map(queue => new BullMQAdapter(queue, { readOnlyMode: true })),
62+
queues: queueAdapters,
3963
serverAdapter: bullBoardAdapter,
64+
options: { historyProvider: this.metricsHistoryProvider },
4065
});
4166
app.use('/admin/queues', bullBoardAdapter.getRouter());
4267

@@ -58,6 +83,7 @@ export class Api {
5883
logger.debug(`API server is running on port ${PORT}`);
5984
logger.debug(`Bull Board is available at ${workerApiUrl.origin}/admin/queues`);
6085
});
86+
this.metricsRecorder.start();
6187
}
6288

6389
private async experimental_addGithubRepo(req: Request, res: Response) {
@@ -125,6 +151,9 @@ export class Api {
125151
}
126152

127153
public async dispose() {
154+
this.metricsRecorder.stop();
155+
this.metricsHistoryProvider.disconnect();
156+
128157
return new Promise<void>((resolve, reject) => {
129158
this.server.close((err) => {
130159
if (err) reject(err);

packages/backend/src/index.ts

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -84,7 +84,7 @@ jobManager.register(repoPermissionSyncWorkload);
8484
jobManager.register(attachmentPruneWorkload);
8585
jobManager.register(auditLogPruneWorkload);
8686

87-
const api = new Api(promClient, prisma, jobManager);
87+
const api = new Api(promClient, prisma, jobManager, redis);
8888

8989
await cleanupOrphanedRepoResources(prisma);
9090

@@ -115,11 +115,11 @@ const listenToShutdownSignals = () => {
115115
logger.info(`Received ${signal}, cleaning up...`);
116116

117117
await configManager.dispose()
118+
await api.dispose();
118119
await jobManager.stop();
119120

120121
await prisma.$disconnect();
121122
await redis.quit();
122-
await api.dispose();
123123
await shutdownPosthog();
124124

125125
logger.info('All workers shut down gracefully');

packages/backend/src/jobManager.test.ts

Lines changed: 15 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -33,6 +33,7 @@ const mocks = vi.hoisted(() => {
3333
workers: [] as Array<{
3434
processor: (job: unknown) => Promise<unknown>;
3535
handlers: Map<string, (...args: unknown[]) => void>;
36+
options: unknown;
3637
}>,
3738
};
3839
});
@@ -81,8 +82,9 @@ vi.mock("bullmq", () => ({
8182
constructor(
8283
_name: string,
8384
processor: (job: unknown) => Promise<unknown>,
85+
options: unknown,
8486
) {
85-
this.record = { processor, handlers: new Map() };
87+
this.record = { processor, handlers: new Map(), options };
8688
mocks.workers.push(this.record);
8789
}
8890

@@ -92,6 +94,7 @@ vi.mock("bullmq", () => ({
9294

9395
close = mocks.workerClose;
9496
},
97+
MetricsTime: { ONE_WEEK: 10_080 },
9598
}));
9699

97100
import { BullMQJobManager } from "./jobManager.js";
@@ -234,6 +237,17 @@ describe("BullMQJobManager lifecycle", () => {
234237
expect(mocks.producerClose).toHaveBeenCalledOnce();
235238
});
236239

240+
test("retains one week of native metrics for history recording", async () => {
241+
const manager = new BullMQJobManager({} as Redis);
242+
manager.register(createWorkload());
243+
244+
await manager.start();
245+
246+
expect(mocks.workers[0].options).toMatchObject({
247+
metrics: { maxDataPoints: 10_080 },
248+
});
249+
});
250+
237251
test("calls onStarted before processing and onCompleted after completion", async () => {
238252
const calls: string[] = [];
239253
const workload = createWorkload({

packages/backend/src/jobManager.ts

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -10,7 +10,7 @@ import {
1010
scheduleToMs,
1111
runWithJobLogContext,
1212
} from "@sourcebot/shared";
13-
import { Job, Queue, Worker } from "bullmq";
13+
import { Job, MetricsTime, Queue, Worker } from "bullmq";
1414
import { Redis } from "ioredis";
1515
import { WORKER_STOP_GRACEFUL_TIMEOUT_MS } from "./constants.js";
1616
import { createExecutionLockRunner } from "./executionLock.js";
@@ -220,6 +220,7 @@ export class BullMQJobManager implements JobManager {
220220
connection: this.connection,
221221
concurrency,
222222
maxStalledCount: 1,
223+
metrics: { maxDataPoints: MetricsTime.ONE_WEEK },
223224
...(rateLimit
224225
? {
225226
limiter: {

0 commit comments

Comments
 (0)