diff --git a/packages/dd-trace/src/opentelemetry/metrics/index.js b/packages/dd-trace/src/opentelemetry/metrics/index.js index e9e16910bc7..20c6d424f06 100644 --- a/packages/dd-trace/src/opentelemetry/metrics/index.js +++ b/packages/dd-trace/src/opentelemetry/metrics/index.js @@ -10,6 +10,8 @@ const MeterProvider = require('./meter_provider') const PeriodicMetricReader = require('./periodic_metric_reader') const OtlpHttpMetricExporter = require('./otlp_http_metric_exporter') +const RESERVED_TRACER_TAGS = new Set(['service', 'env', 'version', 'runtime_id', 'runtime-id']) + /** * @typedef {import('../../config')} Config */ @@ -78,7 +80,16 @@ function initializeOpenTelemetryMetrics (config) { metrics.setGlobalMeterProvider(meterProvider) } -function buildResourceAttributes (tags, { reportHostname, otelSemanticsEnabled, service, env, serviceVersion } = {}) { +/** + * @param {Record} tags + * @param {object} [options] + * @param {boolean} [options.reportHostname] + * @param {string} [options.service] + * @param {string} [options.env] + * @param {string} [options.serviceVersion] + * @returns {import('@opentelemetry/api').Attributes} + */ +function buildResourceAttributes (tags, { reportHostname, service, env, serviceVersion } = {}) { const attrs = { 'telemetry.sdk.name': 'datadog', 'telemetry.sdk.language': 'nodejs', @@ -89,14 +100,19 @@ function buildResourceAttributes (tags, { reportHostname, otelSemanticsEnabled, if (env) attrs['deployment.environment.name'] = env if (reportHostname) attrs['host.name'] = os.hostname() - if (!otelSemanticsEnabled) { - if (tags?.['runtime-id']) attrs['datadog.runtime_id'] = tags['runtime-id'] - const processTagsObject = processTags.tagsObject - if (processTagsObject) { - for (const key of Object.keys(processTagsObject)) { - attrs[`datadog.${key}`] = processTagsObject[key] - } - } + if (tags['runtime-id']) attrs['datadog.runtime_id'] = tags['runtime-id'] + const tracerTags = [] + for (const [key, value] of Object.entries(tags)) { + const valueType = typeof value + const supported = valueType === 'string' || valueType === 'boolean' || + (valueType === 'number' && Number.isFinite(value)) + if (!RESERVED_TRACER_TAGS.has(key) && supported) tracerTags.push(`${key}:${value}`) + } + if (tracerTags.length) attrs['datadog.tracer_tags'] = tracerTags + // Mirrors the legacy v0.6/stats ProcessTags shape (buildProcessTags().tagsArray); keep both in sync. + const processTagsArray = processTags.tagsArray + if (processTagsArray.length) { + attrs['datadog.process_tags'] = processTagsArray } return attrs } @@ -106,7 +122,6 @@ function createOtlpSpanStatsExporter (config) { const protocol = config.OTEL_EXPORTER_OTLP_METRICS_PROTOCOL || 'http/json' const resourceAttributes = buildResourceAttributes(config.tags, { reportHostname: config.reportHostname, - otelSemanticsEnabled: config.DD_TRACE_OTEL_SEMANTICS_ENABLED, service: config.service, env: config.env, serviceVersion: config.version, @@ -115,8 +130,6 @@ function createOtlpSpanStatsExporter (config) { config.OTEL_EXPORTER_OTLP_METRICS_ENDPOINT, protocol, resourceAttributes, - config.DD_TRACE_OTEL_SEMANTICS_ENABLED, - config.service, config.OTEL_EXPORTER_OTLP_METRICS_HEADERS, config.OTEL_EXPORTER_OTLP_METRICS_TIMEOUT ) diff --git a/packages/dd-trace/src/opentelemetry/metrics/otlp_span_stats_exporter.js b/packages/dd-trace/src/opentelemetry/metrics/otlp_span_stats_exporter.js index 018493dbf5c..b7809e00ffd 100644 --- a/packages/dd-trace/src/opentelemetry/metrics/otlp_span_stats_exporter.js +++ b/packages/dd-trace/src/opentelemetry/metrics/otlp_span_stats_exporter.js @@ -11,15 +11,12 @@ class OtlpStatsExporter extends OtlpHttpExporterBase { * @param {string} url * @param {string} protocol * @param {import('@opentelemetry/api').Attributes} resourceAttributes - * @param {boolean} [otelSemanticsEnabled] - * @param {string} [defaultService] * @param {Record} [headers] * @param {number} [timeout] */ - constructor (url, protocol, resourceAttributes, otelSemanticsEnabled = false, defaultService = '', - headers, timeout = 10_000) { + constructor (url, protocol, resourceAttributes, headers, timeout = 10_000) { super(url, headers, timeout, protocol, 'span-stats') - this.#transformer = new OtlpStatsTransformer(resourceAttributes, protocol, otelSemanticsEnabled, defaultService) + this.#transformer = new OtlpStatsTransformer(resourceAttributes, protocol) } /** diff --git a/packages/dd-trace/src/opentelemetry/metrics/otlp_span_stats_transformer.js b/packages/dd-trace/src/opentelemetry/metrics/otlp_span_stats_transformer.js index 8830919979b..eac7648cee1 100644 --- a/packages/dd-trace/src/opentelemetry/metrics/otlp_span_stats_transformer.js +++ b/packages/dd-trace/src/opentelemetry/metrics/otlp_span_stats_transformer.js @@ -1,10 +1,11 @@ 'use strict' -const { LogCollapsingLowestDenseDDSketch } = require('../../../../../vendor/dist/@datadog/sketches-js') const OtlpTransformerBase = require('../otlp/otlp_transformer_base') const { getProtobufTypes } = require('../otlp/protobuf_loader') const { GRPC_STATUS_NAMES } = require('../../constants') +const { stableStringify } = OtlpTransformerBase + const NS_PER_S = 1e9 // Must match libdatadog's EXPLICIT_BOUNDS_SECONDS and OTel spanmetrics connector defaults. @@ -12,6 +13,19 @@ const EXPLICIT_BOUNDS_SECONDS = [ 0.002, 0.004, 0.006, 0.008, 0.01, 0.05, 0.1, 0.2, 0.4, 0.8, 1, 1.4, 2, 5, 10, 15, ] +const SPAN_KIND_METRIC_MAP = { + internal: 'SPAN_KIND_INTERNAL', + SPAN_KIND_INTERNAL: 'SPAN_KIND_INTERNAL', + server: 'SPAN_KIND_SERVER', + SPAN_KIND_SERVER: 'SPAN_KIND_SERVER', + client: 'SPAN_KIND_CLIENT', + SPAN_KIND_CLIENT: 'SPAN_KIND_CLIENT', + producer: 'SPAN_KIND_PRODUCER', + SPAN_KIND_PRODUCER: 'SPAN_KIND_PRODUCER', + consumer: 'SPAN_KIND_CONSUMER', + SPAN_KIND_CONSUMER: 'SPAN_KIND_CONSUMER', +} + /** * @param {object} sketch * @returns {number[]} @@ -41,22 +55,16 @@ function getDeltaTemporality () { return _deltaTemporality } -const ERROR_STATUS_ATTR = { key: 'status.code', value: { intValue: 2 } } +const STATUS_CODE_OK = 'STATUS_CODE_OK' +const STATUS_CODE_ERROR = 'STATUS_CODE_ERROR' class OtlpStatsTransformer extends OtlpTransformerBase { - #otelSemanticsEnabled - #defaultService - /** * @param {import('@opentelemetry/api').Attributes} resourceAttributes * @param {string} protocol - * @param {boolean} [otelSemanticsEnabled] - * @param {string} [defaultService] */ - constructor (resourceAttributes, protocol, otelSemanticsEnabled = false, defaultService = '') { + constructor (resourceAttributes, protocol) { super(resourceAttributes, protocol, 'span-stats') - this.#otelSemanticsEnabled = otelSemanticsEnabled - this.#defaultService = defaultService } /** @@ -82,32 +90,34 @@ class OtlpStatsTransformer extends OtlpTransformerBase { const dataPoints = [] for (const { timeNs, bucket } of drained) { + const distributions = new Map() const endTimeNs = timeNs + bucketSizeNs const startNano = isJson ? String(timeNs) : timeNs const endNano = isJson ? String(endTimeNs) : endTimeNs for (const aggStats of bucket.values()) { - const baseAttrs = this.#buildAttributes(aggStats.aggKey) - - if (this.#otelSemanticsEnabled) { - const okDist = new LogCollapsingLowestDenseDDSketch() - okDist.merge(aggStats.topLevelOkDistribution) - okDist.merge(aggStats.nonTopLevelOkDistribution) - const errDist = new LogCollapsingLowestDenseDDSketch() - errDist.merge(aggStats.topLevelErrorDistribution) - errDist.merge(aggStats.nonTopLevelErrorDistribution) - this.#pushPoint(dataPoints, okDist, startNano, endNano, baseAttrs) - this.#pushPoint(dataPoints, errDist, startNano, endNano, [...baseAttrs, ERROR_STATUS_ATTR]) - } else { - const tlAttrs = [...baseAttrs, { key: 'datadog.span.top_level', value: { boolValue: true } }] - const ntlAttrs = [...baseAttrs, { key: 'datadog.span.top_level', value: { boolValue: false } }] - this.#pushPoint(dataPoints, aggStats.topLevelOkDistribution, startNano, endNano, tlAttrs) - this.#pushPoint(dataPoints, aggStats.topLevelErrorDistribution, startNano, endNano, - [...tlAttrs, ERROR_STATUS_ATTR]) - this.#pushPoint(dataPoints, aggStats.nonTopLevelOkDistribution, startNano, endNano, ntlAttrs) - this.#pushPoint(dataPoints, aggStats.nonTopLevelErrorDistribution, startNano, endNano, - [...ntlAttrs, ERROR_STATUS_ATTR]) - } + const baseAttributes = this.#buildAttributes(aggStats.aggKey) + + this.#addDistribution( + distributions, aggStats.topLevelOkDistribution, startNano, endNano, + baseAttributes, true, STATUS_CODE_OK + ) + this.#addDistribution( + distributions, aggStats.topLevelErrorDistribution, startNano, endNano, + baseAttributes, true, STATUS_CODE_ERROR + ) + this.#addDistribution( + distributions, aggStats.nonTopLevelOkDistribution, startNano, endNano, + baseAttributes, false, STATUS_CODE_OK + ) + this.#addDistribution( + distributions, aggStats.nonTopLevelErrorDistribution, startNano, endNano, + baseAttributes, false, STATUS_CODE_ERROR + ) + } + + for (const { sketch, startNano, endNano, attributes } of distributions.values()) { + this.#pushPoint(dataPoints, sketch, startNano, endNano, attributes) } } @@ -123,8 +133,44 @@ class OtlpStatsTransformer extends OtlpTransformerBase { }] } - #pushPoint (points, sketch, startNano, endNano, attributes) { + /** + * @param {Map} distributions + * @param {object} sketch + * @param {string | number} startNano + * @param {string | number} endNano + * @param {import('@opentelemetry/api').Attributes} baseAttributes + * @param {boolean} topLevel + * @param {string} statusCode + * @returns {void} + */ + #addDistribution (distributions, sketch, startNano, endNano, baseAttributes, topLevel, statusCode) { if (!sketch || sketch.count === 0) return + + const attributes = { + ...baseAttributes, + 'datadog.span.top_level': topLevel, + 'status.code': statusCode, + } + const key = stableStringify(attributes) + const existing = distributions.get(key) + if (existing) { + existing.sketch.merge(sketch) + } else { + distributions.set(key, { + sketch, + startNano, + endNano, + attributes: this.transformAttributes(attributes), + }) + } + } + + #pushPoint (points, sketch, startNano, endNano, attributes) { points.push({ attributes, startTimeUnixNano: startNano, @@ -140,15 +186,18 @@ class OtlpStatsTransformer extends OtlpTransformerBase { /** * @param {import('../../span_stats').SpanAggKey} aggKey + * @returns {import('@opentelemetry/api').Attributes} */ #buildAttributes (aggKey) { - const raw = { 'span.name': aggKey.resource } - - if (aggKey.service && aggKey.service !== this.#defaultService) { - raw['service.name'] = aggKey.service + const spanKind = Object.hasOwn(SPAN_KIND_METRIC_MAP, aggKey.spanKind) + ? SPAN_KIND_METRIC_MAP[aggKey.spanKind] + : 'SPAN_KIND_INTERNAL' + const raw = { + 'span.name': aggKey.resource, + 'service.name': aggKey.service, + 'span.kind': spanKind, } - if (aggKey.spanKind) raw['span.kind'] = aggKey.spanKind if (aggKey.statusCode) raw['http.response.status_code'] = Number(aggKey.statusCode) if (aggKey.method) raw['http.request.method'] = aggKey.method if (aggKey.endpoint) raw['http.route'] = aggKey.endpoint @@ -159,13 +208,15 @@ class OtlpStatsTransformer extends OtlpTransformerBase { : String(aggKey.rpcStatusCode).toUpperCase() } - if (!this.#otelSemanticsEnabled) { - raw['datadog.operation.name'] = aggKey.name - if (aggKey.type) raw['datadog.span.type'] = aggKey.type - if (aggKey.synthetics) raw['datadog.origin'] = 'synthetics' - } + // TODO: additional_metric_tags support is still evolving/TBD across most SDKs; not implemented here yet. + + raw['datadog.operation.name'] = aggKey.name + if (aggKey.type) raw['datadog.span.type'] = aggKey.type + if (aggKey.synthetics) raw['datadog.origin'] = 'synthetics' + if (aggKey.srvSrc) raw['datadog.svc_src'] = aggKey.srvSrc + if (typeof aggKey.isTraceRoot === 'boolean') raw['datadog.is_trace_root'] = aggKey.isTraceRoot - return this.transformAttributes(raw) + return raw } } diff --git a/packages/dd-trace/src/span_processor.js b/packages/dd-trace/src/span_processor.js index 727631a1085..34ba801fee3 100644 --- a/packages/dd-trace/src/span_processor.js +++ b/packages/dd-trace/src/span_processor.js @@ -6,10 +6,11 @@ const SpanSampler = require('./span_sampler') const GitMetadataTagger = require('./git_metadata_tagger') const processTags = require('./process-tags') const { applyHttpOtelSemantics } = require('./plugins/util/http-otel-semantics') -const { APM_TRACING_ENABLED_KEY } = require('./constants') +const { APM_TRACING_ENABLED_KEY, TOP_LEVEL_KEY } = require('./constants') const startedSpans = new WeakSet() const finishedSpans = new WeakSet() +const servicesByTrace = new WeakMap() class SpanProcessor { constructor (exporter, prioritySampler, config, otlpStatsExporter) { @@ -17,6 +18,7 @@ class SpanProcessor { this._prioritySampler = prioritySampler this._config = config this._killAll = false + this._trackServiceBoundaries = Boolean(otlpStatsExporter) if (config.stats?.DD_TRACE_STATS_COMPUTATION_ENABLED && !config.appsec?.standalone?.enabled) { const { SpanStatsProcessor } = require('./span_stats') @@ -56,8 +58,32 @@ class SpanProcessor { let isFirstSpanInChunk = true const stampApmDisabled = this._config.apmTracingEnabled === false + let serviceBySpanId + if (this._trackServiceBoundaries) { + serviceBySpanId = servicesByTrace.get(trace) + if (!serviceBySpanId) { + serviceBySpanId = new WeakMap() + servicesByTrace.set(trace, serviceBySpanId) + } + for (const span of started) { + const context = span.context() + if (context._spanId !== undefined) { + serviceBySpanId.set(context._spanId, context.getTag('service.name')) + } + } + } for (const span of started) { + if (serviceBySpanId) { + const context = span.context() + const parentId = context._parentId + const service = context.getTag('service.name') + const parentService = serviceBySpanId.get(parentId) + if (parentId && (!serviceBySpanId.has(parentId) || + (service !== undefined && parentService !== undefined && service !== parentService))) { + context.setTag(TOP_LEVEL_KEY, 1) + } + } if (span._duration === undefined) { active.push(span) } else { diff --git a/packages/dd-trace/src/span_stats.js b/packages/dd-trace/src/span_stats.js index 82d9d5e96ef..e5f66588eb6 100644 --- a/packages/dd-trace/src/span_stats.js +++ b/packages/dd-trace/src/span_stats.js @@ -14,8 +14,10 @@ const { GRPC_STATUS_CODE, } = require('../../../ext/tags') const { ORIGIN_KEY, TOP_LEVEL_KEY, SVC_SRC_KEY, GRPC_STATUS_NAMES } = require('./constants') +const id = require('./id') const GRPC_STATUS_CODE_MAP = Object.fromEntries(GRPC_STATUS_NAMES.map((name, i) => [name, String(i)])) +const ZERO_ID = id('0') const { version } = require('./pkg') const processTags = require('./process-tags') @@ -150,9 +152,24 @@ class SpanAggKey { } class SpanBuckets extends Map { + #includeTraceRoot + + /** + * @param {boolean} [includeTraceRoot] + */ + constructor (includeTraceRoot = false) { + super() + this.#includeTraceRoot = includeTraceRoot + } + forSpan (span) { const aggKey = new SpanAggKey(span) - const key = aggKey.toString() + const baseKey = aggKey.toString() + const parentId = span.parent_id + if (this.#includeTraceRoot && parentId !== undefined && parentId !== null) { + aggKey.isTraceRoot = parentId.equals(ZERO_ID) + } + const key = this.#includeTraceRoot ? `${baseKey},${aggKey.isTraceRoot}` : baseKey if (!this.has(key)) { this.set(key, new SpanAggStats(aggKey)) @@ -163,9 +180,19 @@ class SpanBuckets extends Map { } class TimeBuckets extends Map { + #includeTraceRoot + + /** + * @param {boolean} [includeTraceRoot] + */ + constructor (includeTraceRoot = false) { + super() + this.#includeTraceRoot = includeTraceRoot + } + forTime (time) { if (!this.has(time)) { - this.set(time, new SpanBuckets()) + this.set(time, new SpanBuckets(this.#includeTraceRoot)) } return this.get(time) @@ -192,7 +219,7 @@ class SpanStatsProcessor { const intervalMs = otlpExporter ? (flushIntervalMs ?? 10_000) : interval * 1e3 this.interval = intervalMs / 1e3 this.bucketSizeNs = intervalMs * 1e6 - this.buckets = new TimeBuckets() + this.buckets = new TimeBuckets(Boolean(otlpExporter)) this.hostname = os.hostname() this.enabled = enabled this.otlpExporter = otlpExporter || null diff --git a/packages/dd-trace/test/opentelemetry/metrics/otlp_span_stats_exporter.spec.js b/packages/dd-trace/test/opentelemetry/metrics/otlp_span_stats_exporter.spec.js index c1fae321295..ccceb847fd6 100644 --- a/packages/dd-trace/test/opentelemetry/metrics/otlp_span_stats_exporter.spec.js +++ b/packages/dd-trace/test/opentelemetry/metrics/otlp_span_stats_exporter.spec.js @@ -2,7 +2,7 @@ const assert = require('node:assert/strict') const http = require('node:http') -const { describe, it, beforeEach, afterEach } = require('mocha') +const { describe, it, before, beforeEach, afterEach } = require('mocha') const sinon = require('sinon') require('../../setup/core') @@ -11,10 +11,13 @@ const { OtlpStatsExporter } = require('../../../src/opentelemetry/metrics/otlp_s const { buildResourceAttributes, createOtlpSpanStatsExporter } = require('../../../src/opentelemetry/metrics') const { SpanBuckets } = require('../../../src/span_stats') const { HTTP_STATUS_CODE } = require('../../../../../ext/tags') +const processTags = require('../../../src/process-tags') const RESOURCE_ATTRS = { 'service.name': 'svc' } const BUCKET_SIZE_NS = 10 * 1e9 +before(() => processTags.initialize()) + function makeSpan (overrides = {}) { return { startTime: 12345 * 1e9, @@ -50,14 +53,49 @@ describe('buildResourceAttributes', () => { assert.strictEqual(attrs['service.version'], '1.0.0') }) - it('includes datadog.runtime_id from tags when otelSemanticsEnabled is false', () => { - const attrs = buildResourceAttributes({ 'runtime-id': 'abc-123' }, { otelSemanticsEnabled: false }) + it('includes datadog.runtime_id from tags', () => { + const attrs = buildResourceAttributes({ 'runtime-id': 'abc-123' }) assert.strictEqual(attrs['datadog.runtime_id'], 'abc-123') }) - it('omits dd.* attributes when otelSemanticsEnabled is true', () => { - const attrs = buildResourceAttributes({ 'runtime-id': 'abc-123' }, { otelSemanticsEnabled: true }) - assert.ok(!Object.keys(attrs).some(k => k.startsWith('datadog.'))) + it('includes supported non-reserved tracer tags as a single array attribute', () => { + const attrs = buildResourceAttributes({ + team: 'apm', + region: 'us-east-1', + retries: 3, + enabled: false, + service: 'ignored', + env: 'ignored', + version: 'ignored', + runtime_id: 'ignored', + 'runtime-id': 'abc-123', + }) + + assert.deepStrictEqual(attrs['datadog.tracer_tags'], [ + 'team:apm', + 'region:us-east-1', + 'retries:3', + 'enabled:false', + ]) + }) + + it('omits nullish, non-finite, and non-scalar tracer tags', () => { + const attrs = buildResourceAttributes({ + missing: undefined, + nullable: null, + infinite: Infinity, + notANumber: NaN, + nested: { team: 'apm' }, + list: ['apm'], + }) + + assert.ok(!('datadog.tracer_tags' in attrs)) + }) + + it('includes datadog.process_tags as a single array attribute', () => { + const attrs = buildResourceAttributes({}) + assert.ok(Array.isArray(attrs['datadog.process_tags'])) + assert.ok(!('datadog.entrypoint.type' in attrs)) }) }) @@ -76,6 +114,7 @@ describe('createOtlpSpanStatsExporter', () => { const exporter = createOtlpSpanStatsExporter({ OTEL_EXPORTER_OTLP_METRICS_ENDPOINT: 'http://localhost:4318/v1/metrics', service: 'svc', + tags: {}, }) assert.ok(exporter instanceof OtlpStatsExporter) }) diff --git a/packages/dd-trace/test/opentelemetry/metrics/otlp_span_stats_transformer.spec.js b/packages/dd-trace/test/opentelemetry/metrics/otlp_span_stats_transformer.spec.js index 7fcb44e3243..bc981ab8d91 100644 --- a/packages/dd-trace/test/opentelemetry/metrics/otlp_span_stats_transformer.spec.js +++ b/packages/dd-trace/test/opentelemetry/metrics/otlp_span_stats_transformer.spec.js @@ -10,7 +10,7 @@ const { EXPLICIT_BOUNDS_SECONDS } = OtlpStatsTransformer const { SpanBuckets } = require('../../../src/span_stats') const { getProtobufTypes } = require('../../../src/opentelemetry/otlp/protobuf_loader') const { HTTP_STATUS_CODE, HTTP_METHOD, HTTP_ROUTE, SPAN_KIND, GRPC_STATUS_CODE } = require('../../../../../ext/tags') -const { ORIGIN_KEY, TOP_LEVEL_KEY } = require('../../../src/constants') +const { ORIGIN_KEY, TOP_LEVEL_KEY, SVC_SRC_KEY } = require('../../../src/constants') const METRIC_NAME = 'traces.span.sdk.metrics.duration' const RESOURCE_ATTRS = { @@ -20,7 +20,6 @@ const RESOURCE_ATTRS = { 'service.version': '1.2.3', 'deployment.environment.name': 'test', } -const DEFAULT_SERVICE = 'svc' const BUCKET_SIZE_NS = 10 * 1e9 function makeSpan (overrides = {}) { @@ -42,16 +41,16 @@ function makeTopLevelSpan (overrides = {}) { return makeSpan({ metrics: { [TOP_LEVEL_KEY]: 1 }, ...overrides }) } -function makeBucket (spans) { - const bucket = new SpanBuckets() +function makeBucket (spans, includeTraceRoot) { + const bucket = new SpanBuckets(includeTraceRoot) for (const span of spans) { bucket.forSpan(span).record(span) } return bucket } -function makeDrained (timeNs, spans) { - return [{ timeNs, bucket: makeBucket(spans) }] +function makeDrained (timeNs, spans, includeTraceRoot) { + return [{ timeNs, bucket: makeBucket(spans, includeTraceRoot) }] } /** @@ -77,11 +76,11 @@ describe('OtlpStatsTransformer', () => { ({ protoMetricsService, protoAggregationTemporality } = getProtobufTypes()) }) - describe('JSON format (default mode)', () => { + describe('JSON format', () => { let transformer before(() => { - transformer = new OtlpStatsTransformer(RESOURCE_ATTRS, 'http/json', false, DEFAULT_SERVICE) + transformer = new OtlpStatsTransformer(RESOURCE_ATTRS, 'http/json') }) it('emits a single histogram metric with the correct name, unit and temporality', () => { @@ -97,6 +96,7 @@ describe('OtlpStatsTransformer', () => { it('maps span dimensions to OTel and dd.* data-point attributes', () => { const span = makeSpan({ + parent_id: { equals: () => true }, meta: { [HTTP_STATUS_CODE]: 404, [HTTP_METHOD]: 'POST', @@ -104,22 +104,104 @@ describe('OtlpStatsTransformer', () => { [SPAN_KIND]: 'server', [GRPC_STATUS_CODE]: 'OK', [ORIGIN_KEY]: 'synthetics', + [SVC_SRC_KEY]: 'integration', }, }) - const payload = JSON.parse(transformer.transform(makeDrained(12340000000000, [span]), BUCKET_SIZE_NS)) + const payload = JSON.parse(transformer.transform( + makeDrained(12340000000000, [span], true), + BUCKET_SIZE_NS + )) + const dataPoint = dataPointsOf(payload)[0] - assert.deepStrictEqual(attrMapOf(dataPointsOf(payload)[0]), { + assert.deepStrictEqual(attrMapOf(dataPoint), { 'span.name': 'GET /foo', - 'span.kind': 'server', + 'service.name': 'svc', + 'span.kind': 'SPAN_KIND_SERVER', 'http.response.status_code': 404, 'http.request.method': 'POST', 'http.route': '/users/:id', 'rpc.response.status_code': 'OK', + 'status.code': 'STATUS_CODE_OK', 'datadog.operation.name': 'test.op', 'datadog.span.type': 'web', 'datadog.origin': 'synthetics', + 'datadog.svc_src': 'integration', 'datadog.span.top_level': false, + 'datadog.is_trace_root': true, }) + assert.deepStrictEqual( + dataPoint.attributes.filter(({ key }) => key === 'datadog.span.top_level' || key === 'datadog.is_trace_root'), + [ + { key: 'datadog.is_trace_root', value: { boolValue: true } }, + { key: 'datadog.span.top_level', value: { boolValue: false } }, + ] + ) + }) + + it('coalesces span.kind aliases that map to the same exported attribute', () => { + const spans = [ + makeSpan({ meta: { [SPAN_KIND]: 'server' } }), + makeSpan({ meta: { [SPAN_KIND]: 'SPAN_KIND_SERVER' } }), + ] + const drained = makeDrained(12340000000000, spans) + assert.strictEqual(drained[0].bucket.size, 2) + + const payload = JSON.parse(transformer.transform(drained, BUCKET_SIZE_NS)) + const points = dataPointsOf(payload) + + assert.strictEqual(points.length, 1) + assert.strictEqual(points[0].count, 2) + assert.strictEqual(attrMapOf(points[0])['span.kind'], 'SPAN_KIND_SERVER') + }) + + it('defaults missing and unknown span.kind values to SPAN_KIND_INTERNAL', () => { + const spans = [ + makeSpan(), + makeSpan({ meta: { [HTTP_STATUS_CODE]: 200, [SPAN_KIND]: 'unknown' } }), + makeSpan({ meta: { [HTTP_STATUS_CODE]: 200, [SPAN_KIND]: 'SPAN_KIND_UNSPECIFIED' } }), + makeSpan({ meta: { [HTTP_STATUS_CODE]: 200, [SPAN_KIND]: 'toString' } }), + makeSpan({ meta: { [HTTP_STATUS_CODE]: 200, [SPAN_KIND]: 'constructor' } }), + ] + const drained = makeDrained(12340000000000, spans) + assert.strictEqual(drained[0].bucket.size, 5) + + const payload = JSON.parse(transformer.transform(drained, BUCKET_SIZE_NS)) + const points = dataPointsOf(payload) + + assert.strictEqual(points.length, 1) + assert.strictEqual(points[0].count, 5) + assert.strictEqual(attrMapOf(points[0])['span.kind'], 'SPAN_KIND_INTERNAL') + }) + + it('keeps root and non-root distributions separate when datadog.is_trace_root is exported', () => { + const spans = [ + makeSpan({ parent_id: { equals: () => true } }), + makeSpan({ parent_id: { equals: () => false } }), + ] + const payload = JSON.parse(transformer.transform( + makeDrained(12340000000000, spans, true), + BUCKET_SIZE_NS + )) + const points = dataPointsOf(payload) + + assert.strictEqual(points.length, 2) + assert.deepStrictEqual(points.map(point => + point.attributes.find(({ key }) => key === 'datadog.is_trace_root').value + ), [{ boolValue: true }, { boolValue: false }]) + }) + + it('omits datadog.is_trace_root when its value is unknown', () => { + const drained = makeDrained(12340000000000, [makeSpan()], true) + + const payload = JSON.parse(transformer.transform(drained, BUCKET_SIZE_NS)) + + assert.ok(!dataPointsOf(payload)[0].attributes.some(({ key }) => key === 'datadog.is_trace_root')) + }) + + it('omits datadog.svc_src when service source is empty', () => { + const payload = JSON.parse(transformer.transform(makeDrained(12340000000000, [makeSpan()]), BUCKET_SIZE_NS)) + + assert.ok(!dataPointsOf(payload)[0].attributes.some(({ key }) => key === 'datadog.svc_src')) }) it('emits the raw grpc.status.code name upper-cased as rpc.response.status_code', () => { @@ -136,15 +218,22 @@ describe('OtlpStatsTransformer', () => { assert.strictEqual(attrMapOf(dataPointsOf(payload)[0])['rpc.response.status_code'], 'UNAVAILABLE') }) - it('omits optional attributes when not present on the span', () => { + it('omits optional dimensions when not present on the span', () => { const payload = JSON.parse( - transformer.transform(makeDrained(12340000000000, [makeSpan({ meta: {} })]), BUCKET_SIZE_NS) + transformer.transform(makeDrained(12340000000000, [makeSpan({ type: '', meta: {} })]), BUCKET_SIZE_NS) ) - const keys = dataPointsOf(payload)[0].attributes.map(a => a.key) + const attrs = attrMapOf(dataPointsOf(payload)[0]) - for (const key of ['http.response.status_code', 'http.request.method', 'http.route', 'span.kind']) { - assert.ok(!keys.includes(key), `${key} should be omitted`) + for (const key of [ + 'http.response.status_code', + 'http.request.method', + 'http.route', + 'datadog.span.type', + 'datadog.svc_src', + ]) { + assert.ok(!(key in attrs), `${key} should be omitted`) } + assert.strictEqual(attrs['span.kind'], 'SPAN_KIND_INTERNAL') }) it('converts duration to seconds with fixed bounds and a sketch-derived distribution', () => { @@ -162,14 +251,14 @@ describe('OtlpStatsTransformer', () => { assert.strictEqual(dp.bucketCounts.filter(c => c > 0).length, 2) }) - it('marks error data points with status.code=ERROR and ok data points without it', () => { + it('marks error data points with status.code=STATUS_CODE_ERROR and ok data points with STATUS_CODE_OK', () => { const spans = [makeTopLevelSpan(), makeTopLevelSpan({ error: 1 })] const payload = JSON.parse(transformer.transform(makeDrained(12340000000000, spans), BUCKET_SIZE_NS)) const points = dataPointsOf(payload) - const ok = points.find(dp => attrMapOf(dp)['datadog.span.top_level'] === true && !attrMapOf(dp)['status.code']) - const err = points.find(dp => attrMapOf(dp)['status.code'] === 2) - assert.ok(ok, 'ok data point should carry no status.code') + const ok = points.find(dp => attrMapOf(dp)['status.code'] === 'STATUS_CODE_OK') + const err = points.find(dp => attrMapOf(dp)['status.code'] === 'STATUS_CODE_ERROR') + assert.ok(ok, 'ok data point should carry status.code=STATUS_CODE_OK') assert.strictEqual(attrMapOf(err)['datadog.span.top_level'], true) }) @@ -179,8 +268,8 @@ describe('OtlpStatsTransformer', () => { const points = dataPointsOf(payload) assert.strictEqual(points.length, 2) - const ok = points.find(dp => !attrMapOf(dp)['status.code']) - const err = points.find(dp => attrMapOf(dp)['status.code'] === 2) + const ok = points.find(dp => attrMapOf(dp)['status.code'] === 'STATUS_CODE_OK') + const err = points.find(dp => attrMapOf(dp)['status.code'] === 'STATUS_CODE_ERROR') assert.strictEqual(ok.count, 2) assert.strictEqual(err.count, 1) assert.strictEqual(attrMapOf(ok)['datadog.span.top_level'], true) @@ -221,7 +310,7 @@ describe('OtlpStatsTransformer', () => { assert.strictEqual(resourceAttrs['deployment.environment.name'], 'test') }) - it('emits a single scopeMetrics and tags data points whose service differs from the default', () => { + it('emits a single scopeMetrics and tags every data point with service.name, including the default service', () => { const drained = makeDrained(12340000000000, [ makeSpan({ service: 'svc', resource: 'GET /foo' }), makeSpan({ service: 'svc-other', resource: 'GET /bar' }), @@ -234,7 +323,7 @@ describe('OtlpStatsTransformer', () => { const serviceByResource = Object.fromEntries( dataPointsOf(payload).map(dp => [attrMapOf(dp)['span.name'], attrMapOf(dp)['service.name']]) ) - assert.strictEqual(serviceByResource['GET /foo'], undefined) + assert.strictEqual(serviceByResource['GET /foo'], 'svc') assert.strictEqual(serviceByResource['GET /bar'], 'svc-other') }) @@ -258,37 +347,11 @@ describe('OtlpStatsTransformer', () => { }) }) - describe('JSON format (OTel-semantics mode)', () => { - let transformer - - before(() => { - transformer = new OtlpStatsTransformer(RESOURCE_ATTRS, 'http/json', true, DEFAULT_SERVICE) - }) - - it('emits only OTel attributes (no dd.*) while keeping status.code on errors', () => { - const span = makeTopLevelSpan({ - error: 1, - meta: { [HTTP_STATUS_CODE]: 500, [HTTP_METHOD]: 'GET' }, - }) - const payload = JSON.parse(transformer.transform(makeDrained(12340000000000, [span]), BUCKET_SIZE_NS)) - const attrs = attrMapOf(dataPointsOf(payload)[0]) - - assert.ok( - !Object.keys(attrs).some(k => k.startsWith('datadog.')), - 'no datadog.* attributes in OTel-semantics mode' - ) - assert.deepStrictEqual( - { name: attrs['span.name'], method: attrs['http.request.method'], status: attrs['status.code'] }, - { name: 'GET /foo', method: 'GET', status: 2 } - ) - }) - }) - describe('protobuf format', () => { let transformer before(() => { - transformer = new OtlpStatsTransformer(RESOURCE_ATTRS, 'http/protobuf', false, DEFAULT_SERVICE) + transformer = new OtlpStatsTransformer(RESOURCE_ATTRS, 'http/protobuf') }) it('emits a valid ExportMetricsServiceRequest with a single duration metric', () => { @@ -302,7 +365,10 @@ describe('OtlpStatsTransformer', () => { it('uses delta temporality and native typed attribute values', () => { const delta = protoAggregationTemporality.values.AGGREGATION_TEMPORALITY_DELTA - const spans = [makeSpan({ resource: 'GET /a' }), makeTopLevelSpan({ error: 1, resource: 'GET /b' })] + const spans = [ + makeSpan({ resource: 'GET /a' }), + makeTopLevelSpan({ error: 1, resource: 'GET /b', meta: { [SVC_SRC_KEY]: 'integration' } }), + ] const buf = transformer.transform(makeDrained(12340000000000, spans), BUCKET_SIZE_NS) const decoded = protoMetricsService.decode(buf) const metric = decoded.resourceMetrics[0].scopeMetrics[0].metrics[0] @@ -310,14 +376,17 @@ describe('OtlpStatsTransformer', () => { assert.strictEqual(metric.histogram.aggregationTemporality, delta) const okNotTopLevel = metric.histogram.dataPoints.find(dp => dp.attributes.some(a => a.key === 'datadog.span.top_level' && a.value.boolValue === false) && - !dp.attributes.some(a => a.key === 'status.code') + dp.attributes.some(a => a.key === 'status.code' && a.value.stringValue === 'STATUS_CODE_OK') ) const errTopLevel = metric.histogram.dataPoints.find(dp => - dp.attributes.some(a => a.key === 'status.code' && Number(a.value.intValue) === 2) && + dp.attributes.some(a => a.key === 'status.code' && a.value.stringValue === 'STATUS_CODE_ERROR') && dp.attributes.some(a => a.key === 'datadog.span.top_level' && a.value.boolValue === true) ) assert.ok(okNotTopLevel, 'should have ok not-top-level data point') assert.ok(errTopLevel, 'should have error top-level data point') + assert.ok(errTopLevel.attributes.some(a => + a.key === 'datadog.svc_src' && a.value.stringValue === 'integration' + ), 'service source should be an OTLP string') }) }) }) diff --git a/packages/dd-trace/test/span_processor.spec.js b/packages/dd-trace/test/span_processor.spec.js index 06433980cd6..d7a9c8620df 100644 --- a/packages/dd-trace/test/span_processor.spec.js +++ b/packages/dd-trace/test/span_processor.spec.js @@ -9,7 +9,7 @@ const proxyquire = require('proxyquire') require('./setup/core') -const { APM_TRACING_ENABLED_KEY } = require('../src/constants') +const { APM_TRACING_ENABLED_KEY, TOP_LEVEL_KEY } = require('../src/constants') describe('SpanProcessor', () => { let prioritySampler @@ -25,6 +25,24 @@ describe('SpanProcessor', () => { let SpanSampler let sample + function makeSpan (spanId, parentSpanId, service, duration) { + const tags = { 'service.name': service } + const context = { + _trace: trace, + _spanId: spanId, + _parentId: parentSpanId, + _sampling: {}, + getTag: (key) => tags[key], + getTags: () => tags, + setTag: (key, value) => { tags[key] = value }, + } + return { + _duration: duration, + tracer: sinon.stub().returns(tracer), + context: sinon.stub().returns(context), + } + } + before(() => { require('../src/process-tags').initialize() }) @@ -296,6 +314,89 @@ describe('SpanProcessor', () => { assert.ok(!Object.hasOwn(formattedSpan.metrics, APM_TRACING_ENABLED_KEY)) }) + it('should preserve service-entry top-level state across a partial flush', () => { + const rootId = {} + const sameServiceChildId = {} + const differentServiceChildId = {} + const remoteParentId = {} + const root = makeSpan(rootId, remoteParentId, 'web', 100) + const sameServiceChild = makeSpan(sameServiceChildId, rootId, 'web') + const differentServiceChild = makeSpan(differentServiceChildId, rootId, 'postgres') + spanFormat.callsFake(span => { + const context = span.context() + const topLevel = context.getTag(TOP_LEVEL_KEY) + return { + service: context.getTag('service.name'), + parent_id: context._parentId, + metrics: topLevel === undefined ? {} : { [TOP_LEVEL_KEY]: topLevel }, + } + }) + processor = new SpanProcessor(exporter, prioritySampler, config, {}) + processor._stats = { onSpanFinished: sinon.stub() } + config.flushMinSpans = 1 + trace.started = [root, sameServiceChild, differentServiceChild] + trace.finished = [root] + + processor.process(root) + + assert.deepStrictEqual(trace.started, [sameServiceChild, differentServiceChild]) + assert.strictEqual(sameServiceChild.context().getTag(TOP_LEVEL_KEY), undefined) + assert.strictEqual(differentServiceChild.context().getTag(TOP_LEVEL_KEY), 1) + + sameServiceChild._duration = 100 + differentServiceChild._duration = 100 + trace.finished = [sameServiceChild, differentServiceChild] + processor.process(differentServiceChild) + + const sameServiceChildFormatted = spanFormat.getCall(1).returnValue + const differentServiceChildFormatted = spanFormat.getCall(2).returnValue + assert.ok(!Object.hasOwn(sameServiceChildFormatted.metrics, TOP_LEVEL_KEY)) + assert.strictEqual(differentServiceChildFormatted.metrics[TOP_LEVEL_KEY], 1) + }) + + it('should mark a different-service child created after its parent is partially flushed', () => { + const rootId = {} + const activeId = {} + const childId = {} + const remoteParentId = {} + const root = makeSpan(rootId, remoteParentId, 'web', 100) + const active = makeSpan(activeId, rootId, 'web') + spanFormat.callsFake(span => { + const context = span.context() + const topLevel = context.getTag(TOP_LEVEL_KEY) + return { metrics: topLevel === undefined ? {} : { [TOP_LEVEL_KEY]: topLevel } } + }) + processor = new SpanProcessor(exporter, prioritySampler, config, {}) + processor._stats = { onSpanFinished: sinon.stub() } + config.flushMinSpans = 1 + trace.started = [root, active] + trace.finished = [root] + + processor.process(root) + + const child = makeSpan(childId, rootId, 'postgres', 100) + trace.started.push(child) + trace.finished = [child] + processor.process(child) + + assert.strictEqual(spanFormat.secondCall.returnValue.metrics[TOP_LEVEL_KEY], 1) + }) + + it('should mark a span with a remote parent as top level', () => { + const span = makeSpan({}, {}, 'web', 100) + spanFormat.callsFake(span => { + const context = span.context() + return { metrics: { [TOP_LEVEL_KEY]: context.getTag(TOP_LEVEL_KEY) } } + }) + processor = new SpanProcessor(exporter, prioritySampler, config, {}) + trace.started = [span] + trace.finished = [span] + + processor.process(span) + + assert.strictEqual(spanFormat.firstCall.returnValue.metrics[TOP_LEVEL_KEY], 1) + }) + describe('with DD_TRACE_OTEL_SEMANTICS_ENABLED', () => { function formattedHttpSpan () { return { diff --git a/packages/dd-trace/test/span_stats.spec.js b/packages/dd-trace/test/span_stats.spec.js index 515d5e564fc..3652c8fe168 100644 --- a/packages/dd-trace/test/span_stats.spec.js +++ b/packages/dd-trace/test/span_stats.spec.js @@ -123,6 +123,7 @@ describe('SpanAggKey', () => { it('should use sensible defaults', () => { const key = new SpanAggKey({ meta: {}, metrics: {} }) assert.strictEqual(key.toString(), `${DEFAULT_SPAN_NAME},${DEFAULT_SERVICE_NAME},,,0,false,,,,,`) + assert.strictEqual(key.isTraceRoot, undefined) }) it('should include HTTP method and route in aggregation key', () => { @@ -203,6 +204,16 @@ describe('SpanAggKey', () => { key.toString(), 'basic-span,service-name,resource-name,span-type,0,false,,,,,14') }) + it('should defer trace-root detection until bucketing requests it', () => { + const equals = sinon.stub().returns(false) + const span = { ...basicSpan, parent_id: { equals } } + const key = new SpanAggKey(span) + sinon.assert.notCalled(equals) + assert.strictEqual(key.isTraceRoot, undefined) + assert.strictEqual( + key.toString(), 'basic-span,service-name,resource-name,span-type,200,false,,,integration,,') + }) + it('should use rpc.grpc.status_code OTel alias when grpc.status.code is absent', () => { const span = { ...basicSpan, meta: { ...basicSpan.meta, 'rpc.grpc.status_code': '2' }, metrics: {} } const key = new SpanAggKey(span) @@ -333,10 +344,60 @@ describe('SpanBuckets', () => { assert.strictEqual(buckets.size, 1) }) + it('should keep top-level and non-top-level distributions separate within one aggregation bucket', () => { + const localBuckets = new SpanBuckets() + const topLevelBasicSpan = { + ...basicSpan, + metrics: { ...basicSpan.metrics, [TOP_LEVEL_KEY]: 1 }, + } + + localBuckets.forSpan(basicSpan).record(basicSpan) + localBuckets.forSpan(topLevelBasicSpan).record(topLevelBasicSpan) + + assert.strictEqual(localBuckets.size, 1) + const stats = localBuckets.values().next().value + assert.strictEqual(stats.nonTopLevelOkDistribution.count, 1) + assert.strictEqual(stats.topLevelOkDistribution.count, 1) + }) + it('should add a new entry when new span does not match existing agg keys', () => { buckets.forSpan(errorSpan) assert.strictEqual(buckets.size, 2) }) + + it('should split trace roots only when requested by the OTLP exporter', () => { + const rootIdEquals = sinon.stub().returns(true) + const childIdEquals = sinon.stub().returns(false) + const rootSpan = { ...basicSpan, parent_id: { equals: rootIdEquals } } + const childSpan = { ...basicSpan, parent_id: { equals: childIdEquals } } + const legacyBuckets = new SpanBuckets() + const otlpBuckets = new SpanBuckets(true) + + legacyBuckets.forSpan(rootSpan) + legacyBuckets.forSpan(childSpan) + sinon.assert.notCalled(rootIdEquals) + sinon.assert.notCalled(childIdEquals) + + otlpBuckets.forSpan(rootSpan) + otlpBuckets.forSpan(childSpan) + + assert.strictEqual(legacyBuckets.size, 1) + assert.strictEqual(legacyBuckets.values().next().value.aggKey.isTraceRoot, undefined) + assert.strictEqual(otlpBuckets.size, 2) + assert.deepStrictEqual([...otlpBuckets.values()].map(({ aggKey }) => aggKey.isTraceRoot), [true, false]) + sinon.assert.calledOnce(rootIdEquals) + sinon.assert.calledOnce(childIdEquals) + }) + + it('should leave trace-root unknown when parent_id is missing or null', () => { + for (const parentId of [undefined, null]) { + const otlpBuckets = new SpanBuckets(true) + + otlpBuckets.forSpan({ ...basicSpan, parent_id: parentId }) + + assert.strictEqual(otlpBuckets.values().next().value.aggKey.isTraceRoot, undefined) + } + }) }) describe('TimeBuckets', () => { @@ -550,6 +611,17 @@ describe('SpanStatsProcessor', () => { assert.strictEqual(bucketSizeNs, p.bucketSizeNs) }) + it('should split OTLP trace roots when their attribute is exported', () => { + const childSpan = { ...topLevelSpan, parent_id: { equals: () => false } } + const processor = new SpanStatsProcessor(config, otlpExporter) + clearTimeout(processor.timer) + + processor.onSpanFinished(topLevelSpan) + processor.onSpanFinished(childSpan) + + assert.strictEqual(processor.buckets.values().next().value.size, 2) + }) + it('should not call OTLP exporter on interval when drained is empty', () => { otlpExporter.export.resetHistory() const p = new SpanStatsProcessor(config, otlpExporter)