All files / packages/telemetry/src traces.ts

90% Statements 54/60
91.66% Branches 11/12
100% Functions 14/14
90.56% Lines 48/53

Press n or j to go to the next uncovered block, b, p or k for the previous block.

1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182                                                                            34x 34x     34x 34x 34x       34x   34x         375x 238x       613x 34x           375x                     1035x             375x   1035x 1035x           2x 2x             375x   375x   375x 375x           238x 238x   238x 11x 11x 11x           238x   238x   238x               1035x 1035x                     1035x 1035x           1035x 1035x   1035x           1035x 1035x       1037x 1037x 1035x 2x 2x 2x   1035x 1035x      
import { channel as plainChannel, tracingChannel } from 'node:diagnostics_channel'
 
import {
	type DiagLogger,
	type Span,
	SpanStatusCode,
	type Tracer,
	context,
	trace,
} from '@opentelemetry/api'
 
import type { DriverIdentity } from '@ydbjs/core'
 
import {
	EVENT_CHANNELS,
	type EventChannelEntry,
	type TracingChannelEntry,
	buildTracingChannels,
} from './channels.js'
import { getActiveSubscriberSpan, spanStorage } from './context.js'
import {
	BASE_ATTRIBUTES,
	coerceError,
	identityAttrs,
	recordErrorAttributes,
} from './semconv/index.js'
 
export type YdbTracesPipelineOptions = {
	captureQueryText: boolean
	emitAcquireSessionSpan: boolean
}
 
type SpanState = { span: Span }
 
export class YdbTracesPipeline {
	#tracer: Tracer
	#diag: DiagLogger
	#opts: YdbTracesPipelineOptions
	#subs: Disposable[] = []
	#stateMap: WeakMap<object, SpanState> = new WeakMap()
 
	constructor(tracer: Tracer, diag: DiagLogger, opts: YdbTracesPipelineOptions) {
		this.#tracer = tracer
		this.#diag = diag
		this.#opts = opts
	}
 
	enable(): void {
		Iif (this.#subs.length > 0) return
 
		let table = buildTracingChannels({
			captureQueryText: this.#opts.captureQueryText,
			emitAcquireSessionSpan: this.#opts.emitAcquireSessionSpan,
		})
 
		for (let entry of table) this.#subs.push(this.#subscribeTracing(entry))
		for (let entry of EVENT_CHANNELS) this.#subs.push(this.#subscribeEvent(entry))
	}
 
	disable(): void {
		for (let s of this.#subs) s[Symbol.dispose]()
		this.#subs.length = 0
	}
 
	#subscribeTracing(entry: TracingChannelEntry): Disposable {
		// StoreType=Span | undefined matches spanStorage; ContextType=object is
		// the payload shape (each row's lambda narrows it).
		let ch = tracingChannel<Span | undefined, object>(entry.channel)
 
		// Skip `end` (it fires on callback return — too early for async fns,
		// which would close the span before the awaited work runs) and
		// `asyncStart` (no trace-relevant signal). `asyncEnd` runs after
		// `error` as well, but `#finish` is a no-op the second time because
		// the WeakMap entry is already gone.
		//
		// When `#safeCreateSpan` swallows a span-creation error, bindStore
		// stores `undefined`; children fall back to `context.active()` via
		// `getActiveSubscriberSpan`. No cast needed — the ALS type allows it.
		ch.start.bindStore(spanStorage, (ctx) => this.#safeCreateSpan(entry, ctx))
 
		// @types/node declares the start/end/asyncStart fields as required even
		// though `subscribe` accepts a partial set at runtime — we only need
		// `asyncEnd` and `error`, so we type the literal as Partial and open
		// it to the full shape at the boundary.
		type Subs = Parameters<typeof ch.subscribe>[0]
		let handlers: Partial<Subs> = {
			asyncEnd: (ctx) => {
				try {
					this.#finish(ctx, undefined)
				} catch (err) {
					this.#diag.error(`telemetry asyncEnd in ${entry.channel}`, err as Error)
				}
			},
			error: (ctx) => {
				try {
					this.#finish(ctx, ctx.error)
				} catch (err) {
					this.#diag.error(`telemetry error in ${entry.channel}`, err as Error)
				}
			},
		}
 
		ch.subscribe(handlers as Subs)
 
		return {
			[Symbol.dispose]() {
				ch.unsubscribe(handlers as Subs)
				ch.start.unbindStore(spanStorage)
			},
		}
	}
 
	#subscribeEvent(entry: EventChannelEntry): Disposable {
		let ch = plainChannel(entry.channel)
		let diag = this.#diag
 
		let fn = (msg: unknown) => {
			try {
				let span = getActiveSubscriberSpan()
				if (span) entry.apply(msg as never, span)
			} catch (err) {
				diag.error(`telemetry event in ${entry.channel}`, err as Error)
			}
		}
 
		ch.subscribe(fn)
 
		return {
			[Symbol.dispose]() {
				ch.unsubscribe(fn)
			},
		}
	}
 
	// Swallows errors from #createSpan so a thrown attribute hook can't poison
	// the parent ALS scope — children just see the prior active span.
	#safeCreateSpan(entry: TracingChannelEntry, ctx: object): Span | undefined {
		try {
			return this.#createSpan(entry, ctx)
		} catch (err) {
			this.#diag.error(`telemetry start in ${entry.channel}`, err as Error)
			return undefined
		}
	}
 
	#createSpan(entry: TracingChannelEntry, ctx: object): Span {
		// Identity rides in the payload, not in ALS — that's what makes
		// multi-driver attribution work without producers having to wrap their
		// own entry points in any subscriber context.
		let identity = (ctx as { driver?: DriverIdentity }).driver
		let attrs = {
			...BASE_ATTRIBUTES,
			...identityAttrs(identity),
			...(entry.attrs ? entry.attrs(ctx as never) : {}),
		}
 
		let parentSpan = getActiveSubscriberSpan()
		let parentCtx = parentSpan ? trace.setSpan(context.active(), parentSpan) : context.active()
 
		let span = this.#tracer.startSpan(
			entry.span,
			{ kind: entry.kind, attributes: attrs as never },
			parentCtx
		)
 
		this.#stateMap.set(ctx, { span })
		return span
	}
 
	#finish(ctx: object, error: unknown): void {
		let state = this.#stateMap.get(ctx)
		if (!state) return
		if (error !== undefined) {
			state.span.setAttributes(recordErrorAttributes(error))
			state.span.recordException(coerceError(error))
			state.span.setStatus({ code: SpanStatusCode.ERROR })
		}
		state.span.end()
		this.#stateMap.delete(ctx)
	}
}