diff --git a/dev-packages/deno-integration-tests/suites/orchestrion-kafkajs/test.ts b/dev-packages/deno-integration-tests/suites/orchestrion-kafkajs/test.ts new file mode 100644 index 000000000000..819de70e2c40 --- /dev/null +++ b/dev-packages/deno-integration-tests/suites/orchestrion-kafkajs/test.ts @@ -0,0 +1,95 @@ +// + +import { tracingChannel } from 'node:diagnostics_channel'; +import type { TransactionEvent } from '@sentry/core'; +import type { DenoClient } from '@sentry/deno'; +import { getCurrentScope, getGlobalScope, getIsolationScope, init, startSpan } from '@sentry/deno'; +import { assert } from 'https://deno.land/std@0.212.0/assert/assert.ts'; +import { assertEquals } from 'https://deno.land/std@0.212.0/assert/assert_equals.ts'; +import { assertExists } from 'https://deno.land/std@0.212.0/assert/assert_exists.ts'; + +function resetGlobals(): void { + getCurrentScope().clear(); + getCurrentScope().setClient(undefined); + getIsolationScope().clear(); + getGlobalScope().clear(); +} + +/** See deno-redis.test.ts — same sink shape, deduped for clarity. */ +function transactionSink(): { + beforeSendTransaction: (event: TransactionEvent) => null; + waitFor: (predicate: (event: TransactionEvent) => boolean) => Promise; +} { + const transactions: TransactionEvent[] = []; + const waiters: { predicate: (e: TransactionEvent) => boolean; resolve: (e: TransactionEvent) => void }[] = []; + return { + beforeSendTransaction(event) { + transactions.push(event); + for (let i = waiters.length - 1; i >= 0; i--) { + const w = waiters[i]!; + if (w.predicate(event)) { + waiters.splice(i, 1); + w.resolve(event); + } + } + return null; + }, + waitFor(predicate) { + const already = transactions.find(predicate); + if (already) return Promise.resolve(already); + return new Promise(resolve => { + waiters.push({ predicate, resolve }); + }); + }, + }; +} + +function withTimeout(p: Promise, ms: number, what: string): Promise { + let timer: ReturnType | undefined; + const timeout = new Promise((_, reject) => { + timer = setTimeout(() => reject(new Error(`Timed out waiting for ${what} after ${ms}ms`)), ms); + }); + return Promise.race([p, timeout]).finally(() => { + if (timer !== undefined) clearTimeout(timer); + }); +} + +Deno.test('kafkajs instrumentation: included in default integrations (Deno 2.8.0+)', () => { + resetGlobals(); + const client = init({ dsn: 'https://username@domain/123' }) as DenoClient; + const names = client.getOptions().integrations.map(i => i.name); + assert(names.includes('Kafka'), `Kafka should be in defaults, got ${names.join(', ')}`); +}); + +Deno.test('kafkajs instrumentation: orchestrion:kafkajs:send_batch channel produces a nested producer span', async () => { + resetGlobals(); + const sink = transactionSink(); + init({ + dsn: 'https://username@domain/123', + tracesSampleRate: 1, + beforeSendTransaction: sink.beforeSendTransaction, + }); + + const channel = tracingChannel('orchestrion:kafkajs:send_batch'); + + // `arguments[0]` is the `{ topicMessages }` batch; a producer span is opened per message. + const ctx = { arguments: [{ topicMessages: [{ topic: 'my-topic', messages: [{ value: 'hi' }] }] }] }; + + startSpan({ name: 'parent', op: 'test' }, () => { + channel.start.publish(ctx); + channel.asyncEnd.publish(ctx); + }); + + const parent = await withTimeout( + sink.waitFor(t => t.transaction === 'parent'), + 5000, + "'parent' transaction", + ); + + const kafkaSpan = parent.spans?.find(s => s.op === 'message'); + assertExists(kafkaSpan, `expected a message child span, got ops: ${parent.spans?.map(s => s.op).join(', ')}`); + assertEquals(kafkaSpan!.description, 'send my-topic'); + assertEquals(kafkaSpan!.data?.['messaging.system'], 'kafka'); + assertEquals(kafkaSpan!.data?.['messaging.destination.name'], 'my-topic'); + assertEquals(kafkaSpan!.data?.['sentry.origin'], 'auto.kafkajs.orchestrion.producer'); +}); diff --git a/packages/deno/src/index.ts b/packages/deno/src/index.ts index dd73b7bd97d4..3bd9cca9fbbc 100644 --- a/packages/deno/src/index.ts +++ b/packages/deno/src/index.ts @@ -121,6 +121,7 @@ export { expressChannelIntegration, genericPoolChannelIntegration, hapiChannelIntegration, + kafkajsChannelIntegration, knexChannelIntegration, koaChannelIntegration, lruMemoizerChannelIntegration, diff --git a/packages/deno/src/sdk.ts b/packages/deno/src/sdk.ts index 97f990dead6c..8672aaee0dc1 100644 --- a/packages/deno/src/sdk.ts +++ b/packages/deno/src/sdk.ts @@ -16,6 +16,7 @@ import { expressChannelIntegration, genericPoolChannelIntegration, hapiChannelIntegration, + kafkajsChannelIntegration, koaChannelIntegration, lruMemoizerChannelIntegration, mongodbChannelIntegration, @@ -82,6 +83,7 @@ export function getDefaultIntegrations(_options: Options): Integration[] { expressChannelIntegration(), genericPoolChannelIntegration(), hapiChannelIntegration(), + kafkajsChannelIntegration(), koaChannelIntegration(), lruMemoizerChannelIntegration(), mongodbChannelIntegration(), diff --git a/packages/deno/test/__snapshots__/mod.test.ts.snap b/packages/deno/test/__snapshots__/mod.test.ts.snap index a64cec594012..8a7e56e1fb33 100644 --- a/packages/deno/test/__snapshots__/mod.test.ts.snap +++ b/packages/deno/test/__snapshots__/mod.test.ts.snap @@ -119,6 +119,7 @@ snapshot[`captureException 1`] = ` "Express", "GenericPool", "Hapi", + "Kafka", "Koa", "LruMemoizer", "Mongo", @@ -207,6 +208,7 @@ snapshot[`captureMessage 1`] = ` "Express", "GenericPool", "Hapi", + "Kafka", "Koa", "LruMemoizer", "Mongo", @@ -302,6 +304,7 @@ snapshot[`captureMessage twice 1`] = ` "Express", "GenericPool", "Hapi", + "Kafka", "Koa", "LruMemoizer", "Mongo", @@ -404,6 +407,7 @@ snapshot[`captureMessage twice 2`] = ` "Express", "GenericPool", "Hapi", + "Kafka", "Koa", "LruMemoizer", "Mongo",