Uh oh!
There was an error while loading. Please reload this page.
Uh oh!
There was an error while loading. Please reload this page.
- Notifications
You must be signed in to change notification settings - Fork 1.8k
feat(server-utils): Migrate @opentelemetry/instrumentation-kafkajs to orchestrion#21923
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Uh oh!
There was an error while loading. Please reload this page.
Changes from all commits
02bced20e70a4558919da50522a72a8712cFile filter
Filter by extension
Conversations
Uh oh!
There was an error while loading. Please reload this page.
Jump to
Uh oh!
There was an error while loading. Please reload this page.
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,100 @@ | ||
| /* | ||
| * Copyright The OpenTelemetry Authors, Aspecto | ||
| * SPDX-License-Identifier: Apache-2.0 | ||
| * | ||
| * NOTICE from the Sentry authors: | ||
| * - Vendored from: https://github.com/open-telemetry/opentelemetry-js-contrib/tree/15ef7506553f631ea4181391e0c5725a56f0d082/packages/instrumentation-kafkajs | ||
| * - Upstream version: @opentelemetry/instrumentation-kafkajs@0.27.0 | ||
| * - Span-creating wrappers for the consumer callbacks (`_getConsumerEachMessagePatch`/ | ||
| * `_getConsumerEachBatchPatch`), migrated to the `@sentry/core` span API. The `run` channel's `start` | ||
| * subscriber swaps the user's `eachMessage`/`eachBatch` for these before the original runs. | ||
| */ | ||
| import { MESSAGING_BATCH_MESSAGE_COUNT } from '@sentry/conventions/attributes'; | ||
| import type { Span } from '@sentry/core'; | ||
| import { continueTrace, startNewTrace, withActiveSpan } from '@sentry/core'; | ||
| import { | ||
| ATTR_MESSAGING_DESTINATION_PARTITION_ID, | ||
| MESSAGING_OPERATION_TYPE_VALUE_PROCESS, | ||
| MESSAGING_OPERATION_TYPE_VALUE_RECEIVE, | ||
| } from './semconv'; | ||
| import { endSpansOnPromise, getHeaderAsString, getLinksFromHeaders, startConsumerSpan } from './spans'; | ||
| import type { EachBatchHandler, EachMessageHandler, KafkaMessage } from './types'; | ||
| // Marks a callback we've already wrapped. A user can reuse one config object across multiple | ||
| // `consumer.run(config)` calls (e.g. a second consumer); without this guard the second `start` would | ||
| // wrap the wrapper and emit duplicate spans per message. | ||
| const consumerCallbackWrapped: unique symbol = Symbol('sentry-kafkajs-consumer-callback-wrapped'); | ||
| type MaybeWrapped = { [consumerCallbackWrapped]?: true }; | ||
| /** Whether `fn` is a callback this module already wrapped, so callers skip re-wrapping it. */ | ||
| export function isWrappedConsumerCallback(fn: unknown): boolean { | ||
| return typeof fn === 'function' && (fn as MaybeWrapped)[consumerCallbackWrapped] === true; | ||
| } | ||
| /** Wraps `eachMessage` so each processed message becomes a consumer span parented to the message's producer. */ | ||
| export function wrapEachMessage(original: EachMessageHandler): EachMessageHandler { | ||
| const wrapped: EachMessageHandler & MaybeWrapped = function eachMessage(this: unknown, payload) { | ||
| const sentryTrace = getHeaderAsString(payload.message.headers, 'sentry-trace'); | ||
| const baggage = getHeaderAsString(payload.message.headers, 'baggage'); | ||
| // Continue the producer's trace so the consumer span is parented to the message's producer. | ||
| return continueTrace({ sentryTrace, baggage }, () => { | ||
| const span = startConsumerSpan({ | ||
| topic: payload.topic, | ||
| message: payload.message, | ||
| operationType: MESSAGING_OPERATION_TYPE_VALUE_PROCESS, | ||
| attributes: { | ||
| [ATTR_MESSAGING_DESTINATION_PARTITION_ID]: String(payload.partition), | ||
| }, | ||
| }); | ||
| const promise = withActiveSpan(span, () => original.call(this, payload)); | ||
| return endSpansOnPromise([span], promise); | ||
| }); | ||
| }; | ||
| wrapped[consumerCallbackWrapped] = true; | ||
| return wrapped; | ||
| } | ||
| /** Wraps `eachBatch` so the batch pull becomes a fresh-root receiving span with a process span per message. */ | ||
| export function wrapEachBatch(original: EachBatchHandler): EachBatchHandler { | ||
| const wrapped: EachBatchHandler & MaybeWrapped = function eachBatch(this: unknown, payload) { | ||
| // A batch pull aggregates messages from many producers, so the receiving span is a fresh root | ||
| // trace and each processed message links back to its own producer span. Mirrors the OTel messaging | ||
| // semantic conventions for a topic with multiple consumers. | ||
| const receivingSpan = startNewTrace(() => | ||
| startConsumerSpan({ | ||
| topic: payload.batch.topic, | ||
| message: undefined, | ||
| operationType: MESSAGING_OPERATION_TYPE_VALUE_RECEIVE, | ||
| attributes: { | ||
| [MESSAGING_BATCH_MESSAGE_COUNT]: payload.batch.messages.length, | ||
| [ATTR_MESSAGING_DESTINATION_PARTITION_ID]: String(payload.batch.partition), | ||
| }, | ||
| }), | ||
| ); | ||
| return withActiveSpan(receivingSpan, () => { | ||
| const spans: Span[] = [receivingSpan]; | ||
| payload.batch.messages.forEach((message: KafkaMessage) => { | ||
| spans.push( | ||
| startConsumerSpan({ | ||
| topic: payload.batch.topic, | ||
| message, | ||
| operationType: MESSAGING_OPERATION_TYPE_VALUE_PROCESS, | ||
| links: getLinksFromHeaders(message.headers), | ||
| attributes: { | ||
| [ATTR_MESSAGING_DESTINATION_PARTITION_ID]: String(payload.batch.partition), | ||
| }, | ||
| }), | ||
| ); | ||
| }); | ||
| const promise = original.call(this, payload); | ||
| return endSpansOnPromise(spans, promise); | ||
| }); | ||
| }; | ||
| wrapped[consumerCallbackWrapped] = true; | ||
| return wrapped; | ||
| } |
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,113 @@ | ||
| import * as diagnosticsChannel from 'node:diagnostics_channel'; | ||
| import type { TracingChannelSubscribers } from 'node:diagnostics_channel'; | ||
| import type { IntegrationFn, Span } from '@sentry/core'; | ||
| import { debug, defineIntegration } from '@sentry/core'; | ||
| import { DEBUG_BUILD } from '../../../debug-build'; | ||
| import { CHANNELS } from '../../../orchestrion/channels'; | ||
| import { isWrappedConsumerCallback, wrapEachBatch, wrapEachMessage } from './consumer'; | ||
| import { applyErrorToSpans, startProducerSpan } from './spans'; | ||
| import type { ConsumerRunConfig, ProducerBatch } from './types'; | ||
| // NOTE: this uses the same name as the OTel `Kafka` integration by design. When enabled, the OTel | ||
| // integration is omitted from the default set (see `experimentalUseDiagnosticsChannelInjection`). | ||
| const INTEGRATION_NAME = 'Kafka' as const; | ||
| /** The tracing-channel context the transform attaches around `messageProducer.js`'s `sendBatch`. */ | ||
| interface SendBatchChannelContext { | ||
| // `arguments[0]` is the `{ topicMessages }` batch (kafkajs normalizes `send` into `sendBatch`). | ||
| arguments: [ProducerBatch?, ...unknown[]]; | ||
| error?: unknown; | ||
| // The producer spans opened at `start`, ended on `asyncEnd` (and marked errored on `error`). | ||
| _sentrySpans?: Span[]; | ||
| } | ||
| /** The tracing-channel context the transform attaches around `consumer/index.js`'s `run`. */ | ||
| interface ConsumerRunChannelContext { | ||
| // `arguments[0]` is the `run(config)` config whose `eachMessage`/`eachBatch` we swap in place. | ||
| arguments: [ConsumerRunConfig?, ...unknown[]]; | ||
| } | ||
| function subscribeToProducer(): void { | ||
| const channel = diagnosticsChannel.tracingChannel<SendBatchChannelContext>(CHANNELS.KAFKAJS_SEND_BATCH); | ||
| // Node types `subscribe` as requiring every lifecycle handler; runtime accepts a partial set, so we | ||
| // pass only the ones we use (matching `bindTracingChannelToSpan`'s handling in `tracing-channel.ts`). | ||
| const subscribers: Partial<TracingChannelSubscribers<SendBatchChannelContext>> = { | ||
| start(ctx) { | ||
| const spans: Span[] = []; | ||
| // `startProducerSpan` mutates each message's headers; doing it at `start` means the mutation | ||
| // reaches the real call, propagating the producer's trace to consumers. | ||
| (ctx.arguments[0]?.topicMessages ?? []).forEach(topicMessage => { | ||
| topicMessage.messages.forEach(message => { | ||
| spans.push(startProducerSpan(topicMessage.topic, message)); | ||
| }); | ||
| }); | ||
| ctx._sentrySpans = spans; | ||
| }, | ||
| error(ctx) { | ||
| if (ctx._sentrySpans) { | ||
| applyErrorToSpans(ctx._sentrySpans, ctx.error); | ||
| } | ||
| }, | ||
| asyncEnd(ctx) { | ||
| // `asyncEnd` fires on both success and failure; `error` (above) has already set the status. | ||
| ctx._sentrySpans?.forEach(span => span.end()); | ||
| }, | ||
| }; | ||
| channel.subscribe(subscribers as TracingChannelSubscribers<SendBatchChannelContext>); | ||
| } | ||
| function subscribeToConsumer(): void { | ||
| const channel = diagnosticsChannel.tracingChannel<ConsumerRunChannelContext>(CHANNELS.KAFKAJS_CONSUMER_RUN); | ||
| const subscribers: Partial<TracingChannelSubscribers<ConsumerRunChannelContext>> = { | ||
| start(ctx) { | ||
| const config = ctx.arguments[0]; | ||
| if (!config || typeof config !== 'object') { | ||
| return; | ||
| } | ||
| // Swap the user callbacks for span-creating wrappers before `run` destructures its config. The | ||
| // `isWrappedConsumerCallback` guard keeps this idempotent: a config object reused across another | ||
| // `run` (or a second consumer) must not have its wrapper wrapped again, which would double spans. | ||
| if (typeof config.eachMessage === 'function' && !isWrappedConsumerCallback(config.eachMessage)) { | ||
| config.eachMessage = wrapEachMessage(config.eachMessage); | ||
| } | ||
| if (typeof config.eachBatch === 'function' && !isWrappedConsumerCallback(config.eachBatch)) { | ||
| config.eachBatch = wrapEachBatch(config.eachBatch); | ||
| } | ||
| }, | ||
| }; | ||
| channel.subscribe(subscribers as TracingChannelSubscribers<ConsumerRunChannelContext>); | ||
| } | ||
| const _kafkajsChannelIntegration = (() => { | ||
| return { | ||
| name: INTEGRATION_NAME, | ||
| setupOnce() { | ||
| if (!diagnosticsChannel.tracingChannel) { | ||
| return; | ||
| } | ||
| DEBUG_BUILD && | ||
| debug.log( | ||
| `[orchestrion:kafkajs] subscribing to channels "${CHANNELS.KAFKAJS_SEND_BATCH}", "${CHANNELS.KAFKAJS_CONSUMER_RUN}"`, | ||
| ); | ||
| subscribeToProducer(); | ||
| subscribeToConsumer(); | ||
| }, | ||
| }; | ||
| }) satisfies IntegrationFn; | ||
| /** | ||
| * EXPERIMENTAL — orchestrion-driven kafkajs integration. | ||
| * | ||
| * Subscribes to the `orchestrion:kafkajs:*` diagnostics_channels that the orchestrion code transform | ||
| * injects into `kafkajs`'s `producer/messageProducer.js` (`sendBatch`) and `consumer/index.js` (`run`). | ||
| * Requires the orchestrion runtime hook or bundler plugin to be active — wire that up via | ||
| * `experimentalUseDiagnosticsChannelInjection`. | ||
| * | ||
| * Known limitation vs. the OTel integration it replaces: the wrapping producer-`transaction` span is | ||
| * not emitted (the transformer can't replace `transaction()`'s return value to patch commit/abort). | ||
| * Transactional `send`/`sendBatch` calls still produce producer spans, since they route through the | ||
| * same instrumented `sendBatch`. | ||
| */ | ||
| export const kafkajsChannelIntegration = defineIntegration(_kafkajsChannelIntegration); | ||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,17 @@ | ||
| /* | ||
| * Unstable OTel semantic-convention keys/values without a `@sentry/conventions` constant, inlined to | ||
| * keep this integration free of `@opentelemetry/*` deps. Values are byte-identical to the vendored OTel | ||
| * instrumentation (`@sentry/node`'s `integrations/tracing/kafka/vendored/semconv.ts`) for span parity. | ||
| */ | ||
| export const ATTR_MESSAGING_DESTINATION_PARTITION_ID = 'messaging.destination.partition.id' as const; | ||
| export const ATTR_MESSAGING_KAFKA_MESSAGE_KEY = 'messaging.kafka.message.key' as const; | ||
| export const ATTR_MESSAGING_KAFKA_MESSAGE_TOMBSTONE = 'messaging.kafka.message.tombstone' as const; | ||
| export const ATTR_MESSAGING_KAFKA_OFFSET = 'messaging.kafka.offset' as const; | ||
| export const MESSAGING_OPERATION_TYPE_VALUE_PROCESS = 'process' as const; | ||
| export const MESSAGING_OPERATION_TYPE_VALUE_RECEIVE = 'receive' as const; | ||
| export const MESSAGING_OPERATION_TYPE_VALUE_SEND = 'send' as const; | ||
| export const MESSAGING_SYSTEM_VALUE_KAFKA = 'kafka' as const; | ||
andreiborza marked this conversation as resolved.
Uh oh!There was an error while loading. Please reload this page. | ||
| // `_OTHER` is OTel's fallback bucket when no more specific `error.type` is known. | ||
| export const ERROR_TYPE_VALUE_OTHER = '_OTHER' as const; | ||
Uh oh!
There was an error while loading. Please reload this page.
Uh oh!
There was an error while loading. Please reload this page.
Uh oh!
There was an error while loading. Please reload this page.