Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
@@ -0,0 +1,7 @@
services:
redis:
image: redis:8
restart: always
container_name: e2e-tests-node-bullmq-redis
ports:
- '6379:6379'
Original file line number Diff line number Diff line change
@@ -0,0 +1,13 @@
import { execSync } from 'child_process';
import { dirname } from 'path';
import { fileURLToPath } from 'url';

const __dirname = dirname(fileURLToPath(import.meta.url));

export default async function globalSetup() {
// Start Redis via Docker Compose
execSync('docker compose up -d --wait', {
cwd: __dirname,
stdio: 'inherit',
});
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,13 @@
import { execSync } from 'child_process';
import { dirname } from 'path';
import { fileURLToPath } from 'url';

const __dirname = dirname(fileURLToPath(import.meta.url));

export default async function globalTeardown() {
// Stop Redis and remove containers
execSync('docker compose down --volumes', {
cwd: __dirname,
stdio: 'inherit',
});
}
25 changes: 25 additions & 0 deletions dev-packages/e2e-tests/test-applications/node-bullmq/package.json
Original file line number Diff line number Diff line change
@@ -0,0 +1,25 @@
{
"name": "node-bullmq",
"version": "0.0.1",
"private": true,
"type": "module",
"scripts": {
"start": "node src/app.mjs",
"clean": "npx rimraf node_modules pnpm-lock.yaml",
"test": "playwright test",
"test:build": "pnpm install",
"test:assert": "pnpm test"
},
"dependencies": {
"@sentry/node": "file:../../packed/sentry-node-packed.tgz",
"bullmq": "^5.0.0",
"express": "^4.21.0"
},
"devDependencies": {
"@playwright/test": "~1.56.0",
"@sentry-internal/test-utils": "link:../../../test-utils"
},
"volta": {
"extends": "../../package.json"
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,11 @@
import { getPlaywrightConfig } from '@sentry-internal/test-utils';

const config = getPlaywrightConfig({
startCommand: `pnpm start`,
});

export default {
...config,
globalSetup: './global-setup.mjs',
globalTeardown: './global-teardown.mjs',
};
64 changes: 64 additions & 0 deletions dev-packages/e2e-tests/test-applications/node-bullmq/src/app.mjs
Original file line number Diff line number Diff line change
@@ -0,0 +1,64 @@
import './instrument.mjs';

import * as Sentry from '@sentry/node';
import { Queue, Worker } from 'bullmq';
import express from 'express';

const app = express();
const port = 3030;

const connection = { host: '127.0.0.1', port: 6379 };
const telemetry = new Sentry.BullMQTelemetry({ enableMetrics: true });

const testQueue = new Queue('test-queue', { connection, telemetry });

const worker = new Worker(
'test-queue',
async job => {
if (job.name === 'fail-job') {
throw new Error('Test error from BullMQ processor');
}

if (job.name === 'breadcrumb-job') {
Sentry.addBreadcrumb({ message: 'breadcrumb-from-bullmq-processor' });
}

return { success: true };
},
{ connection, telemetry },
);

worker.on('error', err => {
console.error('Worker error:', err);
});

app.get('/enqueue/success', async (req, res) => {
await testQueue.add('success-job', { data: 'test' });
res.send('Job enqueued');
});

app.get('/enqueue/fail', async (req, res) => {
await testQueue.add('fail-job', { data: 'test' });
res.send('Job enqueued');
});

app.get('/enqueue/breadcrumb-test', async (req, res) => {
await testQueue.add('breadcrumb-job', { data: 'test' });
res.send('Job enqueued');
});

app.get('/enqueue/link-test', async (req, res) => {
await testQueue.add('link-job', { data: 'test' });
res.send('Job enqueued');
});

app.get('/check-isolation', async (req, res) => {
Sentry.captureException(new Error('Isolation check'));
res.send('Isolation check');
});

Sentry.setupExpressErrorHandler(app);

app.listen(port, () => {
console.log(`App listening on port ${port}`);
});
Original file line number Diff line number Diff line change
@@ -0,0 +1,7 @@
import * as Sentry from '@sentry/node';

Sentry.init({
dsn: process.env.E2E_TEST_DSN,
tunnel: 'http://localhost:3031/',
tracesSampleRate: 1.0,
});
Original file line number Diff line number Diff line change
@@ -0,0 +1,6 @@
import { startEventProxyServer } from '@sentry-internal/test-utils';

startEventProxyServer({
port: 3031,
proxyServerName: 'node-bullmq',
});
Original file line number Diff line number Diff line change
@@ -0,0 +1,134 @@
import { expect, test } from '@playwright/test';
import { getSpanOp, waitForError, waitForMetric, waitForStreamedSpan } from '@sentry-internal/test-utils';

test('Creates a queue.publish span when adding a job', async ({ baseURL }) => {
const publishSpanPromise = waitForStreamedSpan('node-bullmq', span => {
return (
getSpanOp(span) === 'queue.publish' && span.attributes['sentry.segment.name']?.value === 'GET /enqueue/success'
);
});

await fetch(`${baseURL}/enqueue/success`);

const publishSpan = await publishSpanPromise;

expect(publishSpan.is_segment).toBe(false);
expect(publishSpan.attributes['sentry.origin']?.value).toBe('auto.queue.bullmq.producer');
expect(publishSpan.attributes['messaging.system']?.value).toBe('bullmq');
});

test('Creates a segment span for queue.process when processing a job', async ({ baseURL }) => {
const processSpanPromise = waitForStreamedSpan('node-bullmq', span => {
return span.is_segment && getSpanOp(span) === 'queue.process';
});

await fetch(`${baseURL}/enqueue/success`);

const processSpan = await processSpanPromise;

expect(processSpan.attributes['sentry.origin']?.value).toBe('auto.queue.bullmq.consumer');
expect(processSpan.attributes['messaging.system']?.value).toBe('bullmq');
});

test('Sends exception to Sentry on error in job processor', async ({ baseURL }) => {
const errorEventPromise = waitForError('node-bullmq', event => {
return (
!event.type &&
event.exception?.values?.[0]?.value === 'Test error from BullMQ processor' &&
event.exception?.values?.[0]?.mechanism?.type === 'auto.queue.bullmq'
);
});

await fetch(`${baseURL}/enqueue/fail`);

const errorEvent = await errorEventPromise;

expect(errorEvent.exception?.values).toHaveLength(1);
expect(errorEvent.exception?.values?.[0]?.mechanism).toEqual({
handled: false,
type: 'auto.queue.bullmq',
});
});

test('BullMQ processor breadcrumbs do not leak into subsequent HTTP requests', async ({ baseURL }) => {
const processSpanPromise = waitForStreamedSpan('node-bullmq', span => {
return (
span.is_segment &&
getSpanOp(span) === 'queue.process' &&
span.attributes['bullmq.job.name']?.value === 'breadcrumb-job'
);
});

await fetch(`${baseURL}/enqueue/breadcrumb-test`);

await processSpanPromise;

const errorEventPromise = waitForError('node-bullmq', event => {
return event.exception?.values?.[0]?.value === 'Isolation check';
});

await fetch(`${baseURL}/check-isolation`);

const errorEvent = await errorEventPromise;

const leakedBreadcrumb = (errorEvent.breadcrumbs || []).find(
(b: { message?: string }) => b.message === 'breadcrumb-from-bullmq-processor',
);
expect(leakedBreadcrumb).toBeUndefined();
});

test('Links the queue.process segment span to its producer span', async ({ baseURL }) => {
const producerSpanPromise = waitForStreamedSpan('node-bullmq', span => {
return getSpanOp(span) === 'queue.publish' && span.attributes['bullmq.job.name']?.value === 'link-job';
});

const consumerSpanPromise = waitForStreamedSpan('node-bullmq', span => {
return (
span.is_segment && getSpanOp(span) === 'queue.process' && span.attributes['bullmq.job.name']?.value === 'link-job'
);
});

await fetch(`${baseURL}/enqueue/link-test`);

const producerSpan = await producerSpanPromise;
const consumerSpan = await consumerSpanPromise;

expect(producerSpan.attributes['sentry.segment.name']?.value).toBe('GET /enqueue/link-test');
expect(consumerSpan.attributes['sentry.previous_trace']).toBeUndefined();
expect(consumerSpan.links).toEqual([
{
trace_id: producerSpan.trace_id,
span_id: producerSpan.span_id,
sampled: true,
attributes: { 'sentry.link.type': { type: 'string', value: 'previous_trace' } },
},
]);
});

test('Emits bullmq.jobs.completed counter metric on successful job', async ({ baseURL }) => {
const metricPromise = waitForMetric('node-bullmq', metric => {
return metric.name === 'bullmq.jobs.completed' && metric.type === 'counter';
});

await fetch(`${baseURL}/enqueue/success`);

const metric = await metricPromise;

expect(metric.name).toBe('bullmq.jobs.completed');
expect(metric.type).toBe('counter');
expect(metric.value).toEqual(expect.any(Number));
});

test('Emits bullmq.job.duration histogram metric on job completion', async ({ baseURL }) => {
const metricPromise = waitForMetric('node-bullmq', metric => {
return metric.name === 'bullmq.job.duration' && metric.type === 'distribution';
});

await fetch(`${baseURL}/enqueue/success`);

const metric = await metricPromise;

expect(metric.name).toBe('bullmq.job.duration');
expect(metric.type).toBe('distribution');
expect(metric.value).toEqual(expect.any(Number));
});
1 change: 1 addition & 0 deletions dev-packages/node-integration-tests/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -62,6 +62,7 @@
"@types/pg": "^8.6.5",
"ai": "^4.3.16",
"amqplib": "^0.10.9",
"bullmq": "^5.79.1",
"body-parser": "^2.3.0",
"consola": "^3.2.3",
"cors": "^2.8.5",
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,12 @@
services:
db:
image: redis:latest
restart: always
ports:
- '6384:6379'
healthcheck:
test: ['CMD-SHELL', 'redis-cli ping | grep -q PONG']
interval: 2s
timeout: 3s
retries: 30
start_period: 5s
Original file line number Diff line number Diff line change
@@ -0,0 +1,9 @@
import * as Sentry from '@sentry/node';
import { loggingTransport } from '@sentry-internal/node-integration-tests';

Sentry.init({
dsn: 'https://public@dsn.ingest.sentry.io/1337',
release: '1.0',
tracesSampleRate: 1.0,
transport: loggingTransport,
});
Original file line number Diff line number Diff line change
@@ -0,0 +1,68 @@
import * as Sentry from '@sentry/node';
import { FlowProducer, Queue, Worker } from 'bullmq';

const telemetry = new Sentry.BullMQTelemetry({ enableMetrics: true });
const connection = { host: '127.0.0.1', port: 6384 };

async function run() {
const opsQueue = new Queue('attributes-ops-queue', { connection, telemetry });
const opsWorker = new Worker('attributes-ops-queue', null, { connection, telemetry, name: 'attributes-ops-worker' });
const flowProducer = new FlowProducer({ connection, telemetry });
const processQueue = new Queue('attributes-process-queue', { connection, telemetry });
const processWorker = new Worker('attributes-process-queue', async () => {}, {
connection,
telemetry,
name: 'attributes-worker',
});

const jobProcessed = new Promise(resolve => {
processWorker.once('completed', () => resolve());
});

await Sentry.startSpan({ name: 'bullmq attributes' }, async () => {
await processQueue.add('process-job', { data: 'test-data' });

await opsQueue.addBulk([
{ name: 'bulk-job-1', data: {} },
{ name: 'bulk-job-2', data: {} },
]);
await flowProducer.add({
name: 'flow-root',
queueName: 'attributes-ops-queue',
children: [{ name: 'flow-child', queueName: 'attributes-ops-queue' }],
});
await opsQueue.upsertJobScheduler('attributes-scheduler', { every: 60_000 }, { name: 'scheduled-job' });

const progressJob = await opsQueue.add('progress-job', {});
await opsQueue.updateJobProgress(progressJob.id, { percentage: 50 });
await opsQueue.remove(progressJob.id, { removeChildren: true });
await opsQueue.removeDeduplicationKey('attributes-deduplication-id');
await opsQueue.removeDebounceKey('attributes-debounce-id');
await opsQueue.retryJobs({ count: 10 });
await opsQueue.promoteJobs({ count: 10 });
await opsQueue.trimEvents(100);
await opsQueue.clean(0, 100, 'completed');

await opsWorker.getNextJob('attributes-token', { block: false });
// BullMQ runs these two on timers. Calling them directly runs them once, inside this span.
await opsWorker['moveStalledJobsToWait']();
await opsWorker['lockManager'].extendLocks([]);
await opsWorker.rateLimit(1);
await opsWorker.pause(true);
await opsWorker.close(true);

await opsQueue.rateLimit(1);
await opsQueue.drain(true);
});

await jobProcessed;
await processQueue.recordJobCountsMetric('waiting', 'completed');
await processWorker.close();

await opsQueue.obliterate({ force: true });
await processQueue.obliterate({ force: true });
await Promise.all([flowProducer.close(), opsQueue.close(), processQueue.close()]);
await Sentry.flush();
}

run();
Loading
Loading