diff --git a/.github/CODEOWNERS b/.github/CODEOWNERS index 76f97ada1d4..e9f14e13e91 100644 --- a/.github/CODEOWNERS +++ b/.github/CODEOWNERS @@ -62,6 +62,8 @@ /packages/dd-trace/src/lambda/ @DataDog/serverless-aws @DataDog/apm-serverless /packages/dd-trace/src/azure_metadata.js @DataDog/apm-serverless /packages/dd-trace/src/serverless.js @DataDog/apm-serverless +/packages/dd-trace/src/serverless/vercel.js @DataDog/apm-serverless +/packages/dd-trace/src/flush.js @DataDog/apm-serverless @DataDog/apm-sdk-capabilities-js /packages/dd-trace/test/lambda/ @DataDog/serverless-aws @DataDog/apm-serverless /packages/dd-trace/test/azure_metadata.spec.js @DataDog/apm-serverless /packages/dd-trace/test/serverless.spec.js @DataDog/apm-serverless diff --git a/packages/datadog-plugin-http2/test/server.spec.js b/packages/datadog-plugin-http2/test/server.spec.js index ee4c52be7da..3358c024735 100644 --- a/packages/datadog-plugin-http2/test/server.spec.js +++ b/packages/datadog-plugin-http2/test/server.spec.js @@ -9,6 +9,7 @@ const { setImmediate } = require('node:timers/promises') const { afterEach, beforeEach, describe, it } = require('mocha') const sinon = require('sinon') +const { channel } = require('dc-polyfill') const agent = require('../../dd-trace/test/plugins/agent') const web = require('../../dd-trace/src/plugins/util/web') @@ -309,6 +310,19 @@ describe('Plugin', () => { rawExpectedSchema.server ) + it('publishes a close response event', async () => { + const emit = sinon.spy() + const emitChannel = channel('apm:http2:server:response:emit') + emitChannel.subscribe(emit) + + try { + await request(http2, `http://localhost:${port}/user`) + sinon.assert.calledWithMatch(emit, { eventName: 'close' }) + } finally { + emitChannel.unsubscribe(emit) + } + }) + it('should do automatic instrumentation', done => { agent .assertFirstTraceSpan({ diff --git a/packages/dd-trace/src/dogstatsd.js b/packages/dd-trace/src/dogstatsd.js index 19a580ad2db..c5c31014dd6 100644 --- a/packages/dd-trace/src/dogstatsd.js +++ b/packages/dd-trace/src/dogstatsd.js @@ -8,6 +8,7 @@ const request = require('./exporters/common/request') const log = require('./log') const Histogram = require('./histogram') const { entityId } = require('./exporters/common/docker') +const { registerTelemetryFlusher } = require('./flush') const legacyStorage = storage('legacy') @@ -41,6 +42,7 @@ class DogStatsDClient { this._tags = options.tags this.#tagsPrefix = this._tags?.length ? `|#${this._tags.join(',')}` : '' this._queue = [] + this._activeFlushes = new Set() this._buffer = '' this._offset = 0 this._udp4 = this._socket('udp4') @@ -67,23 +69,44 @@ class DogStatsDClient { this._add(stat, value, TYPE_HISTOGRAM, tags) } - flush () { + flush (done) { const queue = this._enqueue() + const activeFlushes = [...this._activeFlushes] - if (queue.length === 0) return + if (queue.length === 0) return this._joinFlushes(activeFlushes, done) log.debug('Flushing %s metrics via %s', queue.length, this._httpOptions ? 'HTTP' : 'UDP') this._queue = [] + const flush = { callbacks: [] } + this._activeFlushes.add(flush) + activeFlushes.push(flush) + this._joinFlushes(activeFlushes, done) + if (this._httpOptions) { - this._sendHttp(queue) + this._sendHttp(queue, () => this._completeFlush(flush)) } else { - this._sendUdp(queue) + this._sendUdp(queue, () => this._completeFlush(flush)) } } - _sendHttp (queue) { + _joinFlushes (flushes, done) { + if (!done) return + let pending = flushes.length + if (pending === 0) return done() + const complete = () => { + if (--pending === 0) done() + } + for (const flush of flushes) flush.callbacks.push(complete) + } + + _completeFlush (flush) { + this._activeFlushes.delete(flush) + for (const done of flush.callbacks) done() + } + + _sendHttp (queue, done) { const buffer = Buffer.concat(queue) request(buffer, this._httpOptions, (err) => { if (err) { @@ -95,32 +118,46 @@ class DogStatsDClient { // options. Either way, we can give UDP a try. this._httpOptions = undefined } - this._sendUdp(queue) + this._sendUdp(queue, done) + } else { + done?.() } }) } - _sendUdp (queue) { + _sendUdp (queue, done) { // dgram resolves the local address via the instrumented dns.lookup when it // binds on first send; the noop store keeps that self-traffic off the trace. legacyStorage.run({ noop: true }, () => { if (this._family === 0) { this.#lookup(this._host, (error, address, family) => { - if (error) return log.error('DogStatsDClient: Host not found', error) - this._sendUdpFromQueue(queue, address, family) + if (error) { + log.error('DogStatsDClient: Host not found', error) + return done?.() + } + this._sendUdpFromQueue(queue, address, family, done) }) } else { - this._sendUdpFromQueue(queue, this._host, this._family) + this._sendUdpFromQueue(queue, this._host, this._family, done) } }) } - _sendUdpFromQueue (queue, address, family) { + _sendUdpFromQueue (queue, address, family, done) { const socket = family === 6 ? this._udp6 : this._udp4 + let pending = queue.length + const complete = () => { + if (--pending === 0) done?.() + } for (const buffer of queue) { log.debug('Sending to DogStatsD: %s', buffer) - socket.send(buffer, 0, buffer.length, this._port, address) + try { + socket.send(buffer, 0, buffer.length, this._port, address, complete) + } catch (error) { + log.error('DogStatsDClient: UDP error sending metrics', error) + complete() + } } } @@ -212,12 +249,12 @@ class MetricsAggregationClient { this.reset() } - flush () { + flush (done) { this._captureCounters() this._captureGauges() this._captureHistograms() - this._client.flush() + this._client.flush(done) } reset () { @@ -370,6 +407,7 @@ class CustomMetrics { setInterval(flush, 10 * 1000).unref?.() globalThis[Symbol.for('dd-trace')].beforeExitHandlers.add(flush) + registerTelemetryFlusher(done => this.flush(done)) } increment (stat, value = 1, tags) { @@ -392,8 +430,8 @@ class CustomMetrics { this.#client.histogram(stat, value, CustomMetrics.tagTranslator(tags)) } - flush () { - return this.#client.flush() + flush (done) { + return this.#client.flush(done) } /** diff --git a/packages/dd-trace/src/exporters/agent/index.js b/packages/dd-trace/src/exporters/agent/index.js index e951048ab5b..695cd846ae4 100644 --- a/packages/dd-trace/src/exporters/agent/index.js +++ b/packages/dd-trace/src/exporters/agent/index.js @@ -6,6 +6,7 @@ const Writer = require('./writer') class AgentExporter { #timer + #activeFlushes = new Set() constructor (config, prioritySampler) { this._config = config @@ -23,6 +24,7 @@ class AgentExporter { lookup, protocolVersion, headers, + onFlush: this.#trackWriterFlush.bind(this), }) globalThis[Symbol.for('dd-trace')].beforeExitHandlers.add(this.flush.bind(this)) @@ -44,10 +46,10 @@ class AgentExporter { const { flushInterval } = this._config if (flushInterval === 0) { - this._writer.flush() + this.#flush() } else if (this.#timer === undefined) { this.#timer = setTimeout(() => { - this._writer.flush() + this.#flush() this.#timer = undefined }, flushInterval) this.#timer.unref?.() @@ -57,7 +59,54 @@ class AgentExporter { flush (done = () => {}) { clearTimeout(this.#timer) this.#timer = undefined - this._writer.flush(done) + + // Snapshot before the boundary flush so a failed encoding cannot cause a + // Vercel lifecycle flush to abandon exports that were already in flight. + let activeFlushes = [...this.#activeFlushes] + try { + this.#flush() + } catch (error) { + log.error('Failed to flush traces: %s', error.message) + } + activeFlushes = [...new Set([...activeFlushes, ...this.#activeFlushes])] + if (activeFlushes.length === 0) return done() + + let pending = activeFlushes.length + const complete = () => { + if (--pending === 0) done() + } + for (const flush of activeFlushes) flush.callbacks.push(complete) + } + + #flush (done) { + const flush = { callbacks: done ? [done] : [] } + this.#activeFlushes.add(flush) + const complete = () => { + this.#activeFlushes.delete(flush) + for (const callback of flush.callbacks) callback() + } + try { + const flush = this._writer.flushDirect ?? this._writer.flush + flush.call(this._writer, complete) + } catch (error) { + complete() + throw error + } + } + + #trackWriterFlush (flush, done) { + const activeFlush = { callbacks: done ? [done] : [] } + this.#activeFlushes.add(activeFlush) + const complete = () => { + this.#activeFlushes.delete(activeFlush) + for (const callback of activeFlush.callbacks) callback() + } + try { + flush(complete) + } catch (error) { + complete() + throw error + } } } diff --git a/packages/dd-trace/src/exporters/agent/writer.js b/packages/dd-trace/src/exporters/agent/writer.js index d81f197395a..3439b0b0e6e 100644 --- a/packages/dd-trace/src/exporters/agent/writer.js +++ b/packages/dd-trace/src/exporters/agent/writer.js @@ -17,19 +17,21 @@ const firstFlushChannel = channel('dd-trace:exporter:first-flush') class AgentWriter extends BaseWriter { #request = commonRequest #requestTracker + #onFlush constructor (...args) { super({ ...args[0], beforeFirstFlush: () => firstFlushChannel.publish(), }) - const { prioritySampler, lookup, protocolVersion, headers, isTestOptimization } = args[0] + const { prioritySampler, lookup, protocolVersion, headers, isTestOptimization, onFlush } = args[0] const AgentEncoder = getEncoder(protocolVersion) this._prioritySampler = prioritySampler this._lookup = lookup this._protocolVersion = protocolVersion this._headers = headers + this.#onFlush = onFlush this._encoder = new AgentEncoder(this) if (isTestOptimization) { this.#request = require('../../ci-visibility/exporters/request') @@ -46,6 +48,12 @@ class AgentWriter extends BaseWriter { * @returns {void} */ flush (done, options) { + const flush = callback => this.flushDirect(callback, options) + if (this.#onFlush) return this.#onFlush(flush, done) + flush(done) + } + + flushDirect (done, options) { if (this.#requestTracker) { this.#requestTracker.flush(done, options) return diff --git a/packages/dd-trace/src/exporters/span-stats/index.js b/packages/dd-trace/src/exporters/span-stats/index.js index 9fa10de4f8a..99c5b73e14f 100644 --- a/packages/dd-trace/src/exporters/span-stats/index.js +++ b/packages/dd-trace/src/exporters/span-stats/index.js @@ -1,16 +1,76 @@ 'use strict' +const log = require('../../log') const { Writer } = require('./writer') class SpanStatsExporter { + #activeFlushes = new Set() + constructor (config) { this._url = config.url - this._writer = new Writer({ url: this._url }) + this._writer = new Writer({ url: this._url, onFlush: this.#trackWriterFlush.bind(this) }) } - export (payload) { + export (payload, done) { + if (done) { + const activeFlushes = [...this.#activeFlushes] + let pending = activeFlushes.length + 1 + const complete = () => { + if (--pending === 0) done() + } + for (const flush of activeFlushes) flush.callbacks.push(complete) + this._writer.append(payload) + try { + this.#flush(complete) + } catch (error) { + // `#flush` has notified the boundary request; keep waiting for prior exports. + log.error('Failed to flush span stats: %s', error.message) + } + return + } this._writer.append(payload) - this._writer.flush() + this.#flush() + } + + flush (done) { + const activeFlushes = [...this.#activeFlushes] + let pending = activeFlushes.length + 1 + const complete = () => { + if (--pending === 0) done?.() + } + for (const flush of activeFlushes) flush.callbacks.push(complete) + this.#flush(complete) + } + + #flush (done) { + const flush = { callbacks: done ? [done] : [] } + this.#activeFlushes.add(flush) + const complete = () => { + this.#activeFlushes.delete(flush) + for (const callback of flush.callbacks) callback() + } + try { + const flushWriter = this._writer.flushDirect ?? this._writer.flush + flushWriter.call(this._writer, complete) + } catch (error) { + complete() + throw error + } + } + + #trackWriterFlush (flush, done) { + const activeFlush = { callbacks: done ? [done] : [] } + this.#activeFlushes.add(activeFlush) + const complete = () => { + this.#activeFlushes.delete(activeFlush) + for (const callback of activeFlush.callbacks) callback() + } + try { + flush(complete) + } catch (error) { + complete() + throw error + } } } diff --git a/packages/dd-trace/src/exporters/span-stats/writer.js b/packages/dd-trace/src/exporters/span-stats/writer.js index a6a6cecb3a4..863ef36f747 100644 --- a/packages/dd-trace/src/exporters/span-stats/writer.js +++ b/packages/dd-trace/src/exporters/span-stats/writer.js @@ -9,12 +9,25 @@ const request = require('../common/request') const log = require('../../log') class Writer extends BaseWriter { - constructor ({ url }) { + #onFlush + + constructor ({ url, onFlush }) { super(...arguments) this._url = url + this.#onFlush = onFlush this._encoder = new SpanStatsEncoder(this) } + flush (done, options) { + const flush = callback => this.flushDirect(callback, options) + if (this.#onFlush) return this.#onFlush(flush, done) + flush(done) + } + + flushDirect (done, options) { + super.flush(done, options) + } + _sendPayload (data, _, done) { makeRequest(data, this._url, (err, res) => { if (err) { diff --git a/packages/dd-trace/src/flush.js b/packages/dd-trace/src/flush.js new file mode 100644 index 00000000000..2b62e5dc64b --- /dev/null +++ b/packages/dd-trace/src/flush.js @@ -0,0 +1,95 @@ +'use strict' + +const log = require('./log') + +/** + * @typedef {(done: () => void) => void | Promise} TelemetryFlusher + */ + +/** @type {Set} */ +const telemetryFlushers = new Set() +const postTraceTelemetryFlushers = new Set() + +/** + * Registers a configured telemetry pipeline so serverless lifecycle retention + * waits for its final export alongside trace delivery. + * @param {TelemetryFlusher} flusher + * @param {{ afterTrace?: boolean }} [options] + * @returns {() => void} Removes this pipeline when its provider is replaced. + */ +function registerTelemetryFlusher (flusher, options) { + const flushers = options?.afterTrace ? postTraceTelemetryFlushers : telemetryFlushers + flushers.add(flusher) + // Avoid retaining a replaced provider or flushing it alongside the new one. + return () => flushers.delete(flusher) +} + +/** + * Flushes the trace exporter and every registered telemetry pipeline. + * @param {{ + * _exporter?: { flush?: TelemetryFlusher }, + * _processor?: { _stats?: { forceFlush?: TelemetryFlusher } } + * }|undefined} tracer + * @param {() => void} [done] + * @param {{ timeout?: number }} [options] + */ +function flushAll (tracer, done, options) { + const traceExporter = tracer?._exporter + const traceFlusher = traceExporter?.flush + const spanStatsFlusher = tracer?._processor?._stats?.forceFlush + // TODO: Include DSM after DataStreamsProcessor exposes a completion-aware flush API. + let pending = telemetryFlushers.size + postTraceTelemetryFlushers.size + + (typeof traceFlusher === 'function' ? 1 : 0) + + (typeof spanStatsFlusher === 'function' ? 1 : 0) + let completed = false + let timeout + + const finish = () => { + if (completed) return + completed = true + clearTimeout(timeout) + done?.() + } + const complete = () => { + if (--pending === 0) finish() + } + + if (pending === 0) return finish() + if (options?.timeout) { + timeout = setTimeout(() => { + log.warn('Timed out waiting for telemetry flush after %dms', options.timeout) + finish() + }, options.timeout) + } + + const flush = (flusher, afterFlushed) => { + let flushed = false + const onFlushed = error => { + if (flushed) return + flushed = true + if (error) log.error('Error flushing telemetry pipeline:', error) + afterFlushed?.() + complete() + } + try { + const result = flusher(onFlushed) + result?.then(onFlushed, error => onFlushed(error)) + } catch (error) { + onFlushed(error) + } + } + + if (typeof traceFlusher === 'function') { + flush(done => traceFlusher.call(traceExporter, done), () => { + for (const flusher of postTraceTelemetryFlushers) flush(flusher) + }) + } else { + for (const flusher of postTraceTelemetryFlushers) flush(flusher) + } + if (typeof spanStatsFlusher === 'function') { + flush(done => spanStatsFlusher.call(tracer._processor._stats, done)) + } + for (const flusher of telemetryFlushers) flush(flusher) +} + +module.exports = { flushAll, registerTelemetryFlusher } diff --git a/packages/dd-trace/src/opentelemetry/logs/batch_log_processor.js b/packages/dd-trace/src/opentelemetry/logs/batch_log_processor.js index 46e8ba6c16a..67fe227988f 100644 --- a/packages/dd-trace/src/opentelemetry/logs/batch_log_processor.js +++ b/packages/dd-trace/src/opentelemetry/logs/batch_log_processor.js @@ -54,10 +54,34 @@ class BatchLogRecordProcessor { /** * Forces an immediate flush of all pending log records. - * @returns {undefined} Promise that resolves when flush is complete + * @param {Function} [done] Called after all pending log exports complete */ - forceFlush () { - this.#export() + forceFlush (done) { + this.#clearTimer() + // Flush only records present at this boundary. New records belong to the + // later request that produced them and must not extend this lifecycle flush. + const logRecords = this.#logRecords + this.#logRecords = [] + let pending = 2 + const complete = () => { + if (--pending === 0) done?.() + } + + // Join exports already active at this boundary before draining this snapshot. + if (typeof this.exporter.flush === 'function') this.exporter.flush(complete) + else complete() + + const flushNext = () => { + if (logRecords.length === 0) { + complete() + return + } + + // Drain the boundary snapshot one batch at a time. + const batch = logRecords.splice(0, this.#maxExportBatchSize) + this.exporter.export(batch, flushNext) + } + flushNext() } /** @@ -79,6 +103,7 @@ class BatchLogRecordProcessor { * @private */ #export () { + if (this.#logRecords.length === 0) return const logRecords = this.#logRecords.slice(0, this.#maxExportBatchSize) this.#logRecords = this.#logRecords.slice(this.#maxExportBatchSize) diff --git a/packages/dd-trace/src/opentelemetry/logs/index.js b/packages/dd-trace/src/opentelemetry/logs/index.js index d6a40ad0122..03606c94af9 100644 --- a/packages/dd-trace/src/opentelemetry/logs/index.js +++ b/packages/dd-trace/src/opentelemetry/logs/index.js @@ -27,10 +27,13 @@ const os = require('os') * @package */ +const { registerTelemetryFlusher } = require('../../flush') const LoggerProvider = require('./logger_provider') const BatchLogRecordProcessor = require('./batch_log_processor') const OtlpHttpLogExporter = require('./otlp_http_log_exporter') +let unregisterTelemetryFlusher + /** * Initializes OpenTelemetry Logs support * @param {import('../../config/config-base')} config - Tracer configuration instance @@ -77,8 +80,12 @@ function initializeOpenTelemetryLogs (config) { // Create logger provider with processor for Datadog Agent export const loggerProvider = new LoggerProvider({ processor }) - // Register the logger provider globally with OpenTelemetry API + // Expose this provider to application calls through the OpenTelemetry Logs API. loggerProvider.register() + // Remove the old provider callback so lifecycle retention flushes only this global provider. + unregisterTelemetryFlusher?.() + // Include final log batches in lifecycle retention with trace delivery. + unregisterTelemetryFlusher = registerTelemetryFlusher(done => loggerProvider.forceFlush(done)) } module.exports = { diff --git a/packages/dd-trace/src/opentelemetry/logs/logger_provider.js b/packages/dd-trace/src/opentelemetry/logs/logger_provider.js index 820d3da574f..81a2cf340a4 100644 --- a/packages/dd-trace/src/opentelemetry/logs/logger_provider.js +++ b/packages/dd-trace/src/opentelemetry/logs/logger_provider.js @@ -84,12 +84,15 @@ class LoggerProvider { /** * Forces a flush of all pending log records. - * @returns {undefined} Promise that resolves when flush is n ssue cncomplete + * @param {Function} [done] Called after all pending log exports complete */ - forceFlush () { - if (!this.isShutdown) { - return this.processor?.forceFlush() + forceFlush (done) { + if (this.isShutdown || !this.processor) { + done?.() + return } + + this.processor.forceFlush(done) } /** diff --git a/packages/dd-trace/src/opentelemetry/metrics/index.js b/packages/dd-trace/src/opentelemetry/metrics/index.js index 20c6d424f06..d90d43050c1 100644 --- a/packages/dd-trace/src/opentelemetry/metrics/index.js +++ b/packages/dd-trace/src/opentelemetry/metrics/index.js @@ -6,11 +6,13 @@ const { metrics } = require('@opentelemetry/api') const { VERSION } = require('../../../../../version') const processTags = require('../../process-tags') +const { registerTelemetryFlusher } = require('../../flush') 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']) +let unregisterTelemetryFlusher /** * @typedef {import('../../config')} Config @@ -78,6 +80,10 @@ function initializeOpenTelemetryMetrics (config) { const meterProvider = new MeterProvider({ reader }) metrics.setGlobalMeterProvider(meterProvider) + // Remove the old provider callback so lifecycle retention flushes only this global provider. + unregisterTelemetryFlusher?.() + // Include the final metric collection and export in lifecycle retention. + unregisterTelemetryFlusher = registerTelemetryFlusher(done => meterProvider.forceFlush(done)) } /** diff --git a/packages/dd-trace/src/opentelemetry/metrics/meter_provider.js b/packages/dd-trace/src/opentelemetry/metrics/meter_provider.js index ebc9eeb1910..53cfcb20a57 100644 --- a/packages/dd-trace/src/opentelemetry/metrics/meter_provider.js +++ b/packages/dd-trace/src/opentelemetry/metrics/meter_provider.js @@ -49,6 +49,14 @@ class MeterProvider { } return meter } + + /** + * @param {Function} [done] Called after the metric export completes + */ + forceFlush (done) { + if (this.reader) this.reader.forceFlush(done) + else done?.() + } } module.exports = MeterProvider diff --git a/packages/dd-trace/src/opentelemetry/metrics/otlp_http_metric_exporter.js b/packages/dd-trace/src/opentelemetry/metrics/otlp_http_metric_exporter.js index 8af42b70854..3bd2b305d03 100644 --- a/packages/dd-trace/src/opentelemetry/metrics/otlp_http_metric_exporter.js +++ b/packages/dd-trace/src/opentelemetry/metrics/otlp_http_metric_exporter.js @@ -34,10 +34,11 @@ class OtlpHttpMetricExporter extends OtlpHttpExporterBase { * * @param {Map} metrics - Map of metric data to export * - * @returns {void} + * @param {Function} [done] Called after the HTTP export completes */ - export (metrics) { + export (metrics, done) { if (metrics.size === 0) { + done?.({ code: 0 }) return } @@ -56,6 +57,7 @@ class OtlpHttpMetricExporter extends OtlpHttpExporterBase { if (result.code === 0) { this.recordTelemetry('otel.metrics_export_successes', 1, additionalTags) } + done?.(result) }) } } 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 b7809e00ffd..890f1ba1953 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 @@ -22,14 +22,16 @@ class OtlpStatsExporter extends OtlpHttpExporterBase { /** * @param {Array<{timeNs: number, bucket: import('../../span_stats').SpanBuckets}>} drained * @param {number} bucketSizeNs + * @param {Function} [done] Called after the HTTP export completes */ - export (drained, bucketSizeNs) { - if (drained.length === 0) return + export (drained, bucketSizeNs, done) { + if (drained.length === 0) return done?.() const payload = this.#transformer.transform(drained, bucketSizeNs) this.sendPayload(payload, (result) => { if (result.code !== 0) { log.error('Failed to export span stats: %s', result.error?.message) } + done?.() }) } } diff --git a/packages/dd-trace/src/opentelemetry/metrics/periodic_metric_reader.js b/packages/dd-trace/src/opentelemetry/metrics/periodic_metric_reader.js index a97b5cb6c99..32de8d9ae99 100644 --- a/packages/dd-trace/src/opentelemetry/metrics/periodic_metric_reader.js +++ b/packages/dd-trace/src/opentelemetry/metrics/periodic_metric_reader.js @@ -197,14 +197,23 @@ class PeriodicMetricReader { /** * Forces an immediate collection and export of all metrics. - * @returns {void} + * @param {Function} [done] Called after the metric export completes */ - forceFlush () { + forceFlush (done) { if (this.#isShutdown) { log.warn('PeriodicMetricReader is shutdown. %d measurement(s) were dropped', this.#droppedCount) + done?.() return } - this.#collectAndExport() + let pending = 2 + const complete = () => { + if (--pending === 0) done?.() + } + + // Snapshot requests already active before starting this flush's export. + if (typeof this.exporter.flush === 'function') this.exporter.flush(complete) + else complete() + this.#collectAndExport(complete) } /** @@ -250,7 +259,8 @@ class PeriodicMetricReader { * * @param {Function} [callback] - Called after export completes */ - #collectAndExport (callback = () => {}) { + #collectAndExport (callback) { + // Observable instruments must be collected even without synchronous measurements. // Atomically drain measurements for export. New measurements can be recorded // during export without interfering with this batch. const allMeasurements = this.#measurements @@ -292,7 +302,7 @@ class PeriodicMetricReader { } if (allMeasurements.length === 0) { - callback() + callback?.() return } diff --git a/packages/dd-trace/src/opentelemetry/otlp/otlp_http_exporter_base.js b/packages/dd-trace/src/opentelemetry/otlp/otlp_http_exporter_base.js index fe27cb643dd..d9bba361ad8 100644 --- a/packages/dd-trace/src/opentelemetry/otlp/otlp_http_exporter_base.js +++ b/packages/dd-trace/src/opentelemetry/otlp/otlp_http_exporter_base.js @@ -20,6 +20,7 @@ const legacyStorage = storage('legacy') */ class OtlpHttpExporterBase { #transport = https + #activeRequests = new Set() /** * Creates a new OtlpHttpExporterBase instance. @@ -88,39 +89,77 @@ class OtlpHttpExporterBase { }, } - legacyStorage.run({ noop: true }, () => { - const req = this.#transport.request(options, (res) => { - let data = '' + const activeRequest = { callbacks: [] } + this.#activeRequests.add(activeRequest) + let completed = false + const complete = result => { + if (completed) return + completed = true + this.#activeRequests.delete(activeRequest) + resultCallback(result) + for (const callback of activeRequest.callbacks) callback() + } - res.on('data', (chunk) => { - data += chunk + try { + legacyStorage.run({ noop: true }, () => { + const req = this.#transport.request(options, (res) => { + let data = '' + + res.on('data', (chunk) => { + data += chunk + }) + + res.once('error', (error) => { + complete({ code: 1, error }) + }) + + res.once('end', () => { + // @ts-expect-error - res.statusCode can be undefined + if (res.statusCode >= 200 && res.statusCode < 300) { + complete({ code: 0 }) + } else { + const error = new Error(`HTTP ${res.statusCode}: ${data}`) + complete({ code: 1, error }) + } + }) }) - res.once('end', () => { - // @ts-expect-error - res.statusCode can be undefined - if (res.statusCode >= 200 && res.statusCode < 300) { - resultCallback({ code: 0 }) - } else { - const error = new Error(`HTTP ${res.statusCode}: ${data}`) - resultCallback({ code: 1, error }) - } + req.on('error', (error) => { + log.error('Error sending OTLP %s:', this.signalType, error) + complete({ code: 1, error }) }) - }) - req.on('error', (error) => { - log.error('Error sending OTLP %s:', this.signalType, error) - resultCallback({ code: 1, error }) - }) + req.once('timeout', () => { + req.destroy() + const error = new Error('Request timeout') + complete({ code: 1, error }) + }) - req.once('timeout', () => { - req.destroy() - const error = new Error('Request timeout') - resultCallback({ code: 1, error }) + req.write(payload) + req.end() }) + } catch (error) { + log.error('Error sending OTLP %s:', this.signalType, error) + complete({ code: 1, error }) + } + } - req.write(payload) - req.end() - }) + /** + * Calls back once OTLP requests active at the flush boundary have completed. + * @param {Function} [done] + */ + flush (done) { + if (!done) return + const activeRequests = [...this.#activeRequests] + if (activeRequests.length === 0) { + done() + return + } + let pending = activeRequests.length + const complete = () => { + if (--pending === 0) done() + } + for (const request of activeRequests) request.callbacks.push(complete) } /** diff --git a/packages/dd-trace/src/proxy.js b/packages/dd-trace/src/proxy.js index 3330061d831..19ac7b63bd5 100644 --- a/packages/dd-trace/src/proxy.js +++ b/packages/dd-trace/src/proxy.js @@ -12,7 +12,8 @@ const telemetry = require('./telemetry') const nomenclature = require('./service-naming') const PluginManager = require('./plugin_manager') const NoopDogStatsDClient = require('./noop/dogstatsd') -const { IS_SERVERLESS } = require('./serverless') +const { IS_SERVERLESS, initializeServerlessTelemetry } = require('./serverless') +const { flushAll, registerTelemetryFlusher } = require('./flush') const processTags = require('./process-tags') const { isTrue } = require('./util') const { @@ -41,6 +42,8 @@ const OPENFEATURE_STATE_NOOP = 0 const OPENFEATURE_STATE_LAZY = 1 const OPENFEATURE_STATE_ACTIVE = 2 +let unregisterRuntimeMetricsFlusher + class LazyModule { constructor (provider) { this.provider = provider @@ -100,6 +103,11 @@ class Tracer extends NoopProxy { this._pluginManager = new PluginManager(this) this.dogstatsd = new NoopDogStatsDClient() this._tracingInitialized = false + // Keep a stable lifecycle owner even when tracing is disabled. In that + // configuration logs and metrics can still have registered flushers. + this._serverlessTelemetry = { + flushAll: (done, options) => flushAll(this._tracer, done, options), + } this._flare = new LazyModule(() => require('./flare')) this.setBaggageItem = setBaggageItem this.getBaggageItem = getBaggageItem @@ -253,8 +261,14 @@ class Tracer extends NoopProxy { initializeOpenTelemetryMetrics(config) } + unregisterRuntimeMetricsFlusher?.() + unregisterRuntimeMetricsFlusher = undefined if (config.runtimeMetrics.enabled) { runtimeMetrics.start(config) + // Agent trace response metrics are recorded asynchronously, so drain + // runtime metrics after the trace export has completed. + unregisterRuntimeMetricsFlusher = registerTelemetryFlusher( + done => runtimeMetrics.flush(done), { afterTrace: true }) } this.#updateTracing(config) @@ -393,6 +407,8 @@ class Tracer extends NoopProxy { setStartupLogPluginManager(this._pluginManager) startupLog() } + + initializeServerlessTelemetry(this._serverlessTelemetry) } /** diff --git a/packages/dd-trace/src/runtime_metrics/index.js b/packages/dd-trace/src/runtime_metrics/index.js index f9451beb359..600132de8e7 100644 --- a/packages/dd-trace/src/runtime_metrics/index.js +++ b/packages/dd-trace/src/runtime_metrics/index.js @@ -13,6 +13,7 @@ const noop = runtimeMetrics = { gauge () {}, increment () {}, decrement () {}, + flush (done) { done?.() }, } module.exports = { @@ -42,6 +43,10 @@ module.exports = { runtimeMetrics = noop Object.setPrototypeOf(module.exports, noop) }, + + flush (done) { + runtimeMetrics.flush(done) + }, } Object.setPrototypeOf(module.exports, noop) diff --git a/packages/dd-trace/src/runtime_metrics/otlp_runtime_metrics.js b/packages/dd-trace/src/runtime_metrics/otlp_runtime_metrics.js index fbb6e1d7675..3c9ec42c45b 100644 --- a/packages/dd-trace/src/runtime_metrics/otlp_runtime_metrics.js +++ b/packages/dd-trace/src/runtime_metrics/otlp_runtime_metrics.js @@ -245,6 +245,11 @@ module.exports = { decrement (name, tag) { this.count(name, -1, tag) }, + + flush (done) { + if (client) return client.flush(done) + done?.() + }, } /** diff --git a/packages/dd-trace/src/runtime_metrics/runtime_metrics.js b/packages/dd-trace/src/runtime_metrics/runtime_metrics.js index 74e522b047d..1f7d0a431c0 100644 --- a/packages/dd-trace/src/runtime_metrics/runtime_metrics.js +++ b/packages/dd-trace/src/runtime_metrics/runtime_metrics.js @@ -27,6 +27,7 @@ let client = null let lastTime = 0 let lastCpuUsage = null let eventLoopDelayObserver = null +let capture = null // !!!!!!!!!!! // IMPORTANT @@ -76,11 +77,10 @@ module.exports = { lastTime = performance.now() if (nativeMetrics) { - interval = setInterval(() => { + capture = () => { captureNativeMetrics(trackEventLoop, trackGc) captureCommonMetrics(trackEventLoop) - client.flush() - }, flushIntervalMs) + } } else { lastCpuUsage = process.cpuUsage() @@ -92,17 +92,21 @@ module.exports = { eventLoopDelayObserver.enable() } - interval = setInterval(() => { + capture = () => { captureCpuUsage() captureCommonMetrics(trackEventLoop) captureHeapSpace() if (trackEventLoop) { captureEventLoopDelay() } - client.flush() - }, flushIntervalMs) + } } + interval = setInterval(() => { + capture() + client.flush() + }, flushIntervalMs) + interval.unref?.() }, @@ -114,6 +118,7 @@ module.exports = { interval = null client = null + capture = null lastCpuUsage = null gcObserver?.disconnect() @@ -158,6 +163,12 @@ module.exports = { decrement (name, tag) { this.count(name, -1, tag) }, + + flush (done) { + if (!client) return done?.() + capture?.() + client.flush(done) + }, } function captureCpuUsage () { diff --git a/packages/dd-trace/src/serverless.js b/packages/dd-trace/src/serverless.js index 23e424b4baa..4fc0e94221a 100644 --- a/packages/dd-trace/src/serverless.js +++ b/packages/dd-trace/src/serverless.js @@ -45,44 +45,40 @@ function isInServerlessEnvironment () { /** * Gets tags describing the serverless platform where the tracer is running. * + * @param {{ isVercel: boolean }} [platform] Detected serverless platform. * @returns {string[]|undefined} */ -function getServerlessPlatformTags () { - if (getEnvironmentVariable('VERCEL') === '1') { - return getVercelPlatformTags() +function getServerlessPlatformTags (platform = getServerlessPlatform()) { + if (platform.isVercel) { + return require('./serverless/vercel').getVercelPlatformTags() } } /** - * @returns {string[]|undefined} + * Detects the serverless platform once while configuration is built. + * @returns {{ isVercel: boolean }} */ -function getVercelPlatformTags () { - let tags - const projectId = getEnvironmentVariable('VERCEL_PROJECT_ID') - if (projectId) { - tags = ['vercel.project_id', projectId] - } - - const environment = getEnvironmentVariable('VERCEL_ENV') - if (environment) { - tags ??= [] - tags.push('vercel.environment', environment) - } +function getServerlessPlatform () { + return { isVercel: getEnvironmentVariable('VERCEL') === '1' } +} - const region = getEnvironmentVariable('VERCEL_REGION') - if (region) { - tags ??= [] - tags.push('vercel.region', region) +/** + * Registers the lifecycle adapter selected by the detected serverless platform. + * @param {{ flushAll?: (done: () => void) => void }} tracer + */ +function initializeServerlessTelemetry (tracer) { + if (getServerlessPlatform().isVercel) { + return require('./serverless/vercel').registerVercelTelemetryRetention(tracer) } - - return tags } module.exports = { getServerlessPlatformTags, + getServerlessPlatform, getIsGCPFunction, getIsAzureFunction, enableGCPPubSubPushSubscription, getIsFlexConsumptionAzureFunction, + initializeServerlessTelemetry, IS_SERVERLESS: isInServerlessEnvironment(), } diff --git a/packages/dd-trace/src/serverless/vercel.js b/packages/dd-trace/src/serverless/vercel.js new file mode 100644 index 00000000000..eb65fcb0bf3 --- /dev/null +++ b/packages/dd-trace/src/serverless/vercel.js @@ -0,0 +1,111 @@ +'use strict' + +const { channel } = require('dc-polyfill') + +const { getEnvironmentVariable } = require('../config/helper') + +const httpRequestFinishChannel = channel('apm:http:server:request:finish') +const http2ResponseEmitChannel = channel('apm:http2:server:response:emit') +const VERCEL_REQUEST_CONTEXT = Symbol.for('@vercel/request-context') +const VERCEL_FLUSH_TIMEOUT = 2000 +const vercelRetentionHandlers = new WeakMap() + +/** + * @typedef {{ flushAll?: (done: () => void, options?: { timeout?: number }) => void }} TelemetryFlusher + */ + +/** + * @param {TelemetryFlusher} tracer + * @param {() => void} done + * @returns {void} + */ +function flushVercelTelemetry (tracer, done) { + setImmediate(() => { + try { + tracer.flushAll(done, { timeout: VERCEL_FLUSH_TIMEOUT }) + } catch { + done() + } + }) +} + +function registerVercelRequestFlush (tracer) { + const requestContext = getVercelRequestContext() + if (!requestContext) return + + const { waitUntil } = requestContext + if (typeof waitUntil !== 'function') return + + // Retain the invocation synchronously, then flush after the response completes. + let done + const pending = new Promise(resolve => { done = resolve }) + try { + waitUntil(pending) + flushVercelTelemetry(tracer, done) + } catch { + done() + } +} + +function getVercelRequestContext () { + return globalThis[VERCEL_REQUEST_CONTEXT]?.get?.() +} + +/** + * Retains a Vercel Node Function until configured telemetry exporters complete. + * + * @param {TelemetryFlusher} tracer + * @returns {(() => void)|undefined} + */ +function registerVercelTelemetryRetention (tracer) { + const existing = vercelRetentionHandlers.get(tracer) + if (existing) return existing + + if (typeof tracer?.flushAll !== 'function') return + const flushRequest = () => registerVercelRequestFlush(tracer) + const flushHttp2Response = ({ eventName }) => { + if (eventName === 'close') flushRequest() + } + httpRequestFinishChannel.subscribe(flushRequest) + http2ResponseEmitChannel.subscribe(flushHttp2Response) + + const unregister = () => { + httpRequestFinishChannel.unsubscribe(flushRequest) + http2ResponseEmitChannel.unsubscribe(flushHttp2Response) + vercelRetentionHandlers.delete(tracer) + } + vercelRetentionHandlers.set(tracer, unregister) + return unregister +} + +/** + * Gets Vercel deployment tags to attach to spans. + * + * @returns {string[]|undefined} + */ +function getVercelPlatformTags () { + let tags + const projectId = getEnvironmentVariable('VERCEL_PROJECT_ID') + if (projectId) { + tags = ['vercel.project_id', projectId] + } + + const environment = getEnvironmentVariable('VERCEL_ENV') + if (environment) { + tags ??= [] + tags.push('vercel.environment', environment) + } + + const region = getEnvironmentVariable('VERCEL_REGION') + if (region) { + tags ??= [] + tags.push('vercel.region', region) + } + + return tags +} + +module.exports = { + getVercelPlatformTags, + registerVercelTelemetryRetention, +} diff --git a/packages/dd-trace/src/span_stats.js b/packages/dd-trace/src/span_stats.js index e5f66588eb6..6e10d751df5 100644 --- a/packages/dd-trace/src/span_stats.js +++ b/packages/dd-trace/src/span_stats.js @@ -235,6 +235,18 @@ class SpanStatsProcessor { } onInterval () { + this.#flush() + } + + /** + * Drains pending span statistics and waits for their export. + * @param {Function} [done] + */ + forceFlush (done) { + this.#flush(done) + } + + #flush (done) { const drained = this.#drainBuckets() if (this.enabled && !this.otlpExporter) { @@ -248,10 +260,24 @@ class SpanStatsProcessor { RuntimeID: this.tags['runtime-id'], Sequence: ++this.sequence, ProcessTags: processTags.serialized, - }) + }, done) } else if (this.otlpExporter && drained.length > 0) { - this.otlpExporter.export(drained, this.bucketSizeNs) - } + if (typeof this.otlpExporter.flush === 'function' && done) { + // Snapshot requests already in flight before starting this boundary + // export, so a later invocation cannot extend this lifecycle barrier. + let pending = 2 + const complete = () => { + if (--pending === 0) done() + } + this.otlpExporter.flush(complete) + this.otlpExporter.export(drained, this.bucketSizeNs, complete) + } else { + this.otlpExporter.export(drained, this.bucketSizeNs, done) + } + } else if (this.otlpExporter) { + if (typeof this.otlpExporter.flush === 'function') this.otlpExporter.flush(done) + else done?.() + } else done?.() } onSpanFinished (span) { diff --git a/packages/dd-trace/src/tracer.js b/packages/dd-trace/src/tracer.js index 28e7df78d5f..5f5b473bd2c 100644 --- a/packages/dd-trace/src/tracer.js +++ b/packages/dd-trace/src/tracer.js @@ -13,6 +13,7 @@ const { isError } = require('./util') const { setStartupLogConfig } = require('./startup-log') const { DataStreamsCheckpointer, DataStreamsManager, DataStreamsProcessor } = require('./datastreams') const { IS_SERVERLESS } = require('./serverless') +const { flushAll } = require('./flush') const log = require('./log') // Always-on writer (console.warn), not the channel-gated `log`: these surface regardless of // DD_TRACE_DEBUG. @@ -150,6 +151,15 @@ class DatadogTracer extends Tracer { this._dataStreamsProcessor.setUrl(url) } + /** + * Flushes every configured telemetry pipeline. + * @param {Function} [done] Called after every configured export completes + * @param {{ timeout?: number }} [options] Bounds this flush operation. + */ + flushAll (done, options) { + flushAll(this, done, options) + } + scope () { return this._scope } diff --git a/packages/dd-trace/test/dogstatsd.spec.js b/packages/dd-trace/test/dogstatsd.spec.js index 2544531201b..2a25eca09c1 100644 --- a/packages/dd-trace/test/dogstatsd.spec.js +++ b/packages/dd-trace/test/dogstatsd.spec.js @@ -32,6 +32,7 @@ describe('dogstatsd', () => { let assertData let docker let log + let registerTelemetryFlusher beforeEach((done) => { udp6 = { @@ -74,10 +75,12 @@ describe('dogstatsd', () => { docker = {} log = { debug: sinon.stub(), error: sinon.stub() } + registerTelemetryFlusher = sinon.stub() const dogstatsd = proxyquire.noPreserveCache().noCallThru()('../src/dogstatsd', { dgram, '../../datadog-core': datadogCore, + './flush': { registerTelemetryFlusher }, './exporters/common/docker': docker, './log': log, }) @@ -237,6 +240,39 @@ describe('dogstatsd', () => { sinon.assert.notCalled(log.debug) }) + it('calls the flush callback after UDP accepts the metrics', (done) => { + udp4.send = sinon.stub().callsFake((...args) => args.at(-1)()) + client = createDogStatsDClient() + + client.gauge('test.avg', 1) + client.flush(() => { + sinon.assert.calledOnce(udp4.send) + done() + }) + }) + + it('joins an already in-flight UDP flush', (done) => { + let completeFirstFlush + udp4.send = sinon.stub().callsFake((...args) => { + completeFirstFlush = args.at(-1) + }) + client = createDogStatsDClient() + client.gauge('test.avg', 1) + client.flush() + + client.flush(() => { + try { + sinon.assert.calledOnce(udp4.send) + done() + } catch (error) { + done(error) + } + }) + + assert.strictEqual(completeFirstFlush instanceof Function, true) + completeFirstFlush() + }) + it('logs the metric count and the UDP transport on a non-empty flush', () => { client = createDogStatsDClient() @@ -388,6 +424,22 @@ describe('dogstatsd', () => { client.flush() }) + it('calls the flush callback after the HTTP proxy responds', (done) => { + client = createDogStatsDClient({ + metricsProxyUrl: `http://localhost:${httpPort}`, + }) + + client.gauge('test.avg', 1) + client.flush(() => { + try { + assert.strictEqual(Buffer.concat(httpData).toString(), 'test.avg:1|g\n') + done() + } catch (error) { + done(error) + } + }) + }) + it('should support HTTP via URL object', (done) => { assertData = () => { try { @@ -444,13 +496,23 @@ describe('dogstatsd', () => { } }) - statusCode = null + const request = sinon.stub().callsFake((buffer, options, callback) => { + callback(new Error('connection refused')) + }) + const { DogStatsDClient: FailingDogStatsDClient } = proxyquire.noPreserveCache().noCallThru()('../src/dogstatsd', { + dgram, + '../../datadog-core': datadogCore, + './exporters/common/docker': docker, + './exporters/common/request': request, + './log': log, + }) - // host exists but port does not, ECONNREFUSED - client = createDogStatsDClient({ - metricsProxyUrl: 'http://localhost:32700', + client = new FailingDogStatsDClient({ host: 'localhost', + lookup: dns.lookup, + metricsProxyUrl: 'http://localhost:8126', port: 8125, + tags: [], }) client.increment('test.foo', 10) @@ -459,6 +521,18 @@ describe('dogstatsd', () => { }) describe('CustomMetrics', () => { + it('registers its aggregated metrics flush with the telemetry lifecycle', () => { + udp4.send = sinon.stub().callsFake((_buffer, _offset, _length, _port, _host, done) => done()) + client = createCustomMetrics() + client.gauge('test.avg', 10) + const done = sinon.spy() + + registerTelemetryFlusher.firstCall.args[0](done) + + sinon.assert.calledOnce(done) + sinon.assert.calledOnce(udp4.send) + }) + it('.gauge()', () => { client = createCustomMetrics() diff --git a/packages/dd-trace/test/exporters/agent/exporter.spec.js b/packages/dd-trace/test/exporters/agent/exporter.spec.js index e9b38fa44b9..ee9a04d0c06 100644 --- a/packages/dd-trace/test/exporters/agent/exporter.spec.js +++ b/packages/dd-trace/test/exporters/agent/exporter.spec.js @@ -18,6 +18,7 @@ describe('Exporter', () => { let writer let prioritySampler let span + let writerOptions beforeEach(() => { url = 'http://www.example.com:8126' @@ -29,7 +30,10 @@ describe('Exporter', () => { setUrl: sinon.spy(), } prioritySampler = {} - Writer = sinon.stub().returns(writer) + Writer = sinon.stub().callsFake(options => { + writerOptions = options + return writer + }) Exporter = proxyquire('../../../src/exporters/agent', { './writer': Writer, @@ -112,6 +116,70 @@ describe('Exporter', () => { }) }) + describe('flush', () => { + it('waits for trace exports already in flight', () => { + const callbacks = [] + writer.flush = sinon.spy(done => callbacks.push(done)) + exporter = new Exporter({ url, flushInterval: 0 }, prioritySampler) + const flushed = sinon.spy() + + exporter.export([span]) + exporter.flush(flushed) + + callbacks[1]() + sinon.assert.notCalled(flushed) + callbacks[0]() + sinon.assert.calledOnce(flushed) + }) + + it('waits for an encoder-triggered writer flush already in flight', () => { + const callbacks = [] + const flushDirect = sinon.spy(done => callbacks.push(done)) + writer.flushDirect = flushDirect + writer.flush = sinon.spy(done => writerOptions.onFlush(flushDirect, done)) + exporter = new Exporter({ url, flushInterval: 0 }, prioritySampler) + const flushed = sinon.spy() + + // This is the path the encoder uses when it crosses its soft limit. + exporter._writer.flush() + exporter.flush(flushed) + + callbacks[1]() + sinon.assert.notCalled(flushed) + callbacks[0]() + sinon.assert.calledOnce(flushed) + }) + + it('does not retain a failed writer flush', () => { + writer.flush = sinon.stub() + writer.flush.onFirstCall().throws(new Error('encode failed')) + writer.flush.onSecondCall().callsFake(done => done()) + exporter = new Exporter({ url, flushInterval: 0 }, prioritySampler) + const flushed = sinon.spy() + + assert.throws(() => exporter.export([span]), /encode failed/) + exporter.flush(flushed) + + sinon.assert.calledOnce(flushed) + }) + + it('waits for an earlier export when the boundary flush fails', () => { + let inFlightDone + writer.flush = sinon.stub() + writer.flush.onFirstCall().callsFake(done => { inFlightDone = done }) + writer.flush.onSecondCall().throws(new Error('encode failed')) + exporter = new Exporter({ url, flushInterval: 0 }, prioritySampler) + const flushed = sinon.spy() + + exporter.export([span]) + exporter.flush(flushed) + + sinon.assert.notCalled(flushed) + inFlightDone() + sinon.assert.calledOnce(flushed) + }) + }) + describe('setUrl', () => { beforeEach(() => { exporter = new Exporter({ url }) diff --git a/packages/dd-trace/test/exporters/agent/writer.spec.js b/packages/dd-trace/test/exporters/agent/writer.spec.js index 7ed70d40476..f92a48d6866 100644 --- a/packages/dd-trace/test/exporters/agent/writer.spec.js +++ b/packages/dd-trace/test/exporters/agent/writer.spec.js @@ -109,6 +109,20 @@ function describeWriter (protocolVersion) { writer.flush(done) }) + it('routes flushes through the configured lifecycle hook', (done) => { + const onFlush = sinon.spy((flush, done) => flush(done)) + writer = new Writer({ url, prioritySampler, protocolVersion, onFlush }) + + writer.flush(() => { + try { + sinon.assert.calledOnce(onFlush) + done() + } catch (error) { + done(error) + } + }) + }) + it('should flush its traces to the agent, and call callback', (done) => { const expectedData = Buffer.from('prefixed') diff --git a/packages/dd-trace/test/exporters/span-stats/exporter.spec.js b/packages/dd-trace/test/exporters/span-stats/exporter.spec.js index 30c431d7e03..34b01df4e15 100644 --- a/packages/dd-trace/test/exporters/span-stats/exporter.spec.js +++ b/packages/dd-trace/test/exporters/span-stats/exporter.spec.js @@ -15,6 +15,7 @@ describe('span-stats exporter', () => { let exporter let Writer let writer + let log beforeEach(() => { url = new URL('http://www.example.com:8126') @@ -23,9 +24,11 @@ describe('span-stats exporter', () => { flush: sinon.spy(), } Writer = sinon.stub().returns(writer) + log = { error: sinon.spy() } Exporter = proxyquire('../../../src/exporters/span-stats', { './writer': { Writer }, + '../../log': log, }).SpanStatsExporter }) @@ -41,13 +44,74 @@ describe('span-stats exporter', () => { sinon.assert.called(writer.flush) }) + it('waits for an in-flight export during flush', () => { + exporter = new Exporter({ url }) + let inFlightDone + writer.flush = sinon.stub() + writer.flush.onFirstCall().callsFake(done => { inFlightDone = done }) + writer.flush.onSecondCall().callsFake(done => done()) + const done = sinon.spy() + + exporter.export('in flight') + exporter.flush(done) + + sinon.assert.notCalled(done) + inFlightDone() + sinon.assert.calledOnce(done) + }) + + it('waits for an encoder-triggered export during flush', () => { + exporter = new Exporter({ url }) + const onFlush = Writer.firstCall.args[0].onFlush + let automaticDone + onFlush(done => { automaticDone = done }) + writer.flush = sinon.stub().callsFake(done => done()) + const done = sinon.spy() + + exporter.flush(done) + + sinon.assert.notCalled(done) + automaticDone() + sinon.assert.calledOnce(done) + }) + + it('does not retain a failed writer flush', () => { + writer.flush = sinon.stub() + writer.flush.onFirstCall().throws(new Error('encode failed')) + writer.flush.onSecondCall().callsFake(done => done()) + exporter = new Exporter({ url }) + const done = sinon.spy() + + assert.throws(() => exporter.export('failed export'), /encode failed/) + exporter.flush(done) + + sinon.assert.calledOnce(done) + }) + + it('waits for an in-flight export when the boundary flush fails', () => { + writer.flush = sinon.stub() + let inFlightDone + writer.flush.onFirstCall().callsFake(done => { inFlightDone = done }) + writer.flush.onSecondCall().throws(new Error('encode failed')) + exporter = new Exporter({ url }) + const done = sinon.spy() + + exporter.export('in flight') + exporter.export('failed boundary', done) + + sinon.assert.notCalled(done) + inFlightDone() + sinon.assert.calledOnce(done) + sinon.assert.calledOnceWithExactly(log.error, 'Failed to flush span stats: %s', 'encode failed') + }) + it('should set url from config', () => { const url = new URL('http://0.0.0.0:1234') exporter = new Exporter({ url }) assert.strictEqual(exporter._url.toString(), url.toString()) - sinon.assert.calledWith(Writer, { + sinon.assert.calledWithMatch(Writer, { url: exporter._url, }) }) diff --git a/packages/dd-trace/test/exporters/span-stats/writer.spec.js b/packages/dd-trace/test/exporters/span-stats/writer.spec.js index 23919126944..202612c5dcc 100644 --- a/packages/dd-trace/test/exporters/span-stats/writer.spec.js +++ b/packages/dd-trace/test/exporters/span-stats/writer.spec.js @@ -77,6 +77,17 @@ describe('span-stats writer', () => { writer.flush(done) }) + it('routes encoder-triggered flushes through the configured lifecycle hook', () => { + const onFlush = sinon.stub().callsFake((flush, done) => flush(done)) + writer = new Writer({ url, onFlush }) + encoder.count.returns(1) + encoder.encode.callsFake(() => writer.flush()) + + writer.append([span]) + + sinon.assert.calledOnce(onFlush) + }) + it('should flush to the agent, and call callback', (done) => { const expectedData = Buffer.from('prefixed') diff --git a/packages/dd-trace/test/opentelemetry/logs.spec.js b/packages/dd-trace/test/opentelemetry/logs.spec.js index 3928372c6c9..2ec004d14b0 100644 --- a/packages/dd-trace/test/opentelemetry/logs.spec.js +++ b/packages/dd-trace/test/opentelemetry/logs.spec.js @@ -15,6 +15,7 @@ require('../setup/core') const { protoLogsService } = require('../../src/opentelemetry/otlp/protobuf_loader').getProtobufTypes() const { getConfigFresh } = require('../helpers/config') const { assertObjectContains } = require('../../../../integration-tests/helpers') +const BatchLogRecordProcessor = require('../../src/opentelemetry/logs/batch_log_processor') /** * @param {object} type protobufjs Type instance for the OTLP service message @@ -142,6 +143,87 @@ describe('OpenTelemetry Logs', () => { }) describe('Logs Export', () => { + it('waits for an in-flight export during forceFlush', () => { + let exportDone + let flushDone + const processor = new BatchLogRecordProcessor({ + export: (records, done) => { exportDone = done }, + flush: (done) => { flushDone = done }, + }, 60_000, 1) + const done = sinon.spy() + + processor.onEmit({ body: 'in flight' }, { name: 'test' }) + processor.forceFlush(done) + + sinon.assert.notCalled(done) + exportDone({ code: 0 }) + sinon.assert.notCalled(done) + flushDone() + sinon.assert.calledOnce(done) + }) + + it('drains queued batches and waits for earlier size-triggered exports', () => { + const batches = [] + const callbacks = [] + const flushCallbacks = [] + let activeExports = 0 + const completeFlushes = () => { + if (activeExports !== 0) return + while (flushCallbacks.length > 0) flushCallbacks.shift()() + } + const processor = new BatchLogRecordProcessor({ + export: (records, done) => { + batches.push(records) + activeExports++ + callbacks.push(() => { + activeExports-- + done({ code: 0 }) + completeFlushes() + }) + }, + flush: (done) => { + if (activeExports === 0) done() + else flushCallbacks.push(done) + }, + }, 60_000, 2) + const done = sinon.spy() + + for (let index = 0; index < 5; index++) { + processor.onEmit({ body: index }, { name: 'test' }) + } + processor.forceFlush(done) + + assert.deepStrictEqual(batches.map(batch => batch.map(record => record.body)), [ + [0, 1], [2, 3], [4], + ]) + callbacks.shift()() + callbacks.shift()() + callbacks.shift()() + + sinon.assert.calledOnce(done) + }) + + it('does not wait for records emitted after the flush boundary', () => { + const exports = [] + let firstExportDone + const processor = new BatchLogRecordProcessor({ + export: (records, done) => { + exports.push(records.map(record => record.body)) + if (records[0].body === 'before') firstExportDone = done + }, + flush: done => done(), + }, 60_000, 2) + const done = sinon.spy() + + processor.onEmit({ body: 'before' }, { name: 'test' }) + processor.forceFlush(done) + processor.onEmit({ body: 'after' }, { name: 'test' }) + firstExportDone({ code: 0 }) + + assert.deepStrictEqual(exports, [['before']]) + sinon.assert.calledOnce(done) + }) + it('exports logs with complete OTLP structure, trace correlation, and instrumentation info', () => { mockOtlpExport((decoded, capturedHeaders) => { const { resource } = decoded.resourceLogs[0] diff --git a/packages/dd-trace/test/opentelemetry/metrics.spec.js b/packages/dd-trace/test/opentelemetry/metrics.spec.js index 0ad70fbdc51..4f12bba0004 100644 --- a/packages/dd-trace/test/opentelemetry/metrics.spec.js +++ b/packages/dd-trace/test/opentelemetry/metrics.spec.js @@ -13,6 +13,8 @@ require('../setup/core') const { protoMetricsService } = require('../../src/opentelemetry/otlp/protobuf_loader').getProtobufTypes() const { getConfigFresh } = require('../helpers/config') const { DEFAULT_MAX_MEASUREMENT_QUEUE_SIZE } = require('../../src/opentelemetry/metrics/constants') +const MeterProvider = require('../../src/opentelemetry/metrics/meter_provider') +const PeriodicMetricReader = require('../../src/opentelemetry/metrics/periodic_metric_reader') /** * @param {object} type protobufjs Type instance for the OTLP service message @@ -661,6 +663,33 @@ describe('OpenTelemetry Meter Provider', () => { }) describe('Lifecycle', () => { + it('waits for an in-flight export during forceFlush', () => { + const exports = [] + const flushes = [] + const reader = new PeriodicMetricReader({ + export: (metrics, done) => { exports.push(done) }, + flush: (done) => { flushes.push(done) }, + }, 60_000, 'DELTA', 1024) + const meter = new MeterProvider({ reader }).getMeter('test') + const firstDone = sinon.spy() + const done = sinon.spy() + + meter.createCounter('in-flight').add(1) + reader.forceFlush(firstDone) + flushes.shift()() + meter.createCounter('boundary').add(1) + reader.forceFlush(done) + + sinon.assert.notCalled(done) + assert.strictEqual(exports.length, 2) + exports[1]({ code: 0 }) + sinon.assert.notCalled(done) + flushes[0]() + sinon.assert.calledOnce(done) + exports[0]({ code: 0 }) + reader.shutdown() + }) + it('handles shutdown gracefully', async () => { setupMetrics() const provider = metrics.getMeterProvider() 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 ccceb847fd6..84e292bf278 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 @@ -218,4 +218,54 @@ describe('OtlpStatsExporter', () => { exporter.export(drained, BUCKET_SIZE_NS) assert.ok(httpStub.calledOnce) }) + + it('flushes after an in-flight HTTP export completes', () => { + let onEnd + httpStub.callsFake((options, callback) => { + const mockRes = { + statusCode: 200, + on: sinon.stub(), + once: (event, handler) => { + if (event === 'end') onEnd = handler + return mockRes + }, + } + callback(mockRes) + return mockReq + }) + const flushed = sinon.spy() + + exporter.export(makeDrained([makeSpan()]), BUCKET_SIZE_NS) + exporter.flush(flushed) + + sinon.assert.notCalled(flushed) + onEnd() + sinon.assert.calledOnce(flushed) + }) + + it('does not wait for exports started after the flush boundary', () => { + const onEnd = [] + httpStub.callsFake((options, callback) => { + const mockRes = { + statusCode: 200, + on: sinon.stub(), + once: (event, handler) => { + if (event === 'end') onEnd.push(handler) + return mockRes + }, + } + callback(mockRes) + return mockReq + }) + const flushed = sinon.spy() + + exporter.export(makeDrained([makeSpan()]), BUCKET_SIZE_NS) + exporter.flush(flushed) + exporter.export(makeDrained([makeSpan()]), BUCKET_SIZE_NS) + + onEnd[0]() + sinon.assert.calledOnce(flushed) + onEnd[1]() + sinon.assert.calledOnce(flushed) + }) }) diff --git a/packages/dd-trace/test/proxy.spec.js b/packages/dd-trace/test/proxy.spec.js index bd24e4454c0..c72e4a2ae15 100644 --- a/packages/dd-trace/test/proxy.spec.js +++ b/packages/dd-trace/test/proxy.spec.js @@ -47,6 +47,9 @@ describe('TracerProxy', () => { let NoopDogStatsDClient let OpenFeatureProvider let openfeatureProvider + let registerTelemetryFlusher + let initializeServerlessTelemetry + let flushAll beforeEach(() => { process.env.DD_TRACE_MOCHA_ENABLED = 'false' @@ -180,8 +183,13 @@ describe('TracerProxy', () => { runtimeMetrics = { start: sinon.spy(), + flush: sinon.spy(), } + registerTelemetryFlusher = sinon.stub().returns(() => {}) + initializeServerlessTelemetry = sinon.spy() + flushAll = sinon.spy() + profiler = { start: sinon.spy(), } @@ -269,6 +277,11 @@ describe('TracerProxy', () => { './flare': flare, './openfeature': openfeature, './openfeature/flagging_provider': OpenFeatureProvider, + './serverless': { + IS_SERVERLESS: false, + initializeServerlessTelemetry, + }, + './flush': { flushAll, registerTelemetryFlusher }, }) proxy = new ProxyClass() @@ -586,6 +599,29 @@ describe('TracerProxy', () => { sinon.assert.called(runtimeMetrics.start) }) + it('registers the runtime metrics flush with the serverless lifecycle', () => { + config.runtimeMetrics.enabled = true + const done = sinon.spy() + + proxy.init() + registerTelemetryFlusher.firstCall.args[0](done) + + sinon.assert.calledOnceWithExactly(runtimeMetrics.flush, done) + }) + + it('registers Vercel telemetry retention when tracing is disabled', () => { + config.DD_TRACE_ENABLED = false + + proxy.init() + + sinon.assert.calledOnce(initializeServerlessTelemetry) + const telemetry = initializeServerlessTelemetry.firstCall.args[0] + assert.strictEqual(typeof telemetry.flushAll, 'function') + const done = sinon.spy() + telemetry.flushAll(done) + sinon.assert.calledOnceWithExactly(flushAll, proxy._tracer, done, undefined) + }) + it('should expose noop metrics methods prior to initialization', () => { proxy.dogstatsd.increment('foo') }) diff --git a/packages/dd-trace/test/runtime_metrics.spec.js b/packages/dd-trace/test/runtime_metrics.spec.js index 94309d1d4e5..25317b18285 100644 --- a/packages/dd-trace/test/runtime_metrics.spec.js +++ b/packages/dd-trace/test/runtime_metrics.spec.js @@ -74,6 +74,7 @@ NATIVE_METRICS_VARIANTS.forEach((nativeMetrics) => { gauge () {}, increment () {}, decrement () {}, + flush (done) { done?.() }, }) proxy = proxyquire('../src/runtime_metrics', { @@ -152,6 +153,20 @@ NATIVE_METRICS_VARIANTS.forEach((nativeMetrics) => { sinon.assert.notCalled(runtimeMetrics.decrement) sinon.assert.calledOnce(runtimeMetrics.stop) }) + + it('flushes when enabled and is noop when disabled', () => { + const done = sinon.spy() + + proxy.start() + proxy.flush(done) + sinon.assert.notCalled(runtimeMetrics.flush) + sinon.assert.calledOnce(done) + + config.runtimeMetrics.enabled = true + proxy.start(config) + proxy.flush(done) + sinon.assert.calledOnceWithExactly(runtimeMetrics.flush, done) + }) }) describe('runtimeMetrics', () => { @@ -186,7 +201,7 @@ NATIVE_METRICS_VARIANTS.forEach((nativeMetrics) => { gauge: sinon.spy(), increment: sinon.spy(), histogram: sinon.spy(), - flush: sinon.spy(), + flush: sinon.stub().callsFake(done => done?.()), } const proxiedObject = { @@ -246,6 +261,21 @@ NATIVE_METRICS_VARIANTS.forEach((nativeMetrics) => { runtimeMetrics.stop() }) + it('captures and waits for the final runtime metrics flush', (done) => { + client.flush.resetHistory() + client.gauge.resetHistory() + + runtimeMetrics.flush(() => { + try { + sinon.assert.calledOnce(client.flush) + sinon.assert.called(client.gauge) + done() + } catch (error) { + done(error) + } + }) + }) + describe('start', () => { it('it should initialize the Dogstatsd client with the correct options', function () { runtimeMetrics.stop() diff --git a/packages/dd-trace/test/serverless.spec.js b/packages/dd-trace/test/serverless.spec.js index d33923f932b..043f825b9ee 100644 --- a/packages/dd-trace/test/serverless.spec.js +++ b/packages/dd-trace/test/serverless.spec.js @@ -1,13 +1,28 @@ 'use strict' const assert = require('node:assert/strict') +const http = require('node:http') const { describe, it, afterEach } = require('mocha') +const { logs } = require('@opentelemetry/api-logs') +const { metrics } = require('@opentelemetry/api') +const { channel } = require('dc-polyfill') require('./setup/core') -const { getServerlessPlatformTags, enableGCPPubSubPushSubscription } = require('../src/serverless') +const { + getServerlessPlatformTags, + getServerlessPlatform, + enableGCPPubSubPushSubscription, + initializeServerlessTelemetry, +} = require('../src/serverless') +const { registerVercelTelemetryRetention } = require('../src/serverless/vercel') +const { flushAll, registerTelemetryFlusher } = require('../src/flush') +const Tracer = require('../src/tracer') +const { initializeOpenTelemetryLogs } = require('../src/opentelemetry/logs') +const { initializeOpenTelemetryMetrics } = require('../src/opentelemetry/metrics') const agent = require('./plugins/agent') +const { getConfigFresh } = require('./helpers/config') describe('enableGCPPubSubPushSubscription', () => { const originalKService = process.env.K_SERVICE @@ -121,4 +136,253 @@ describe('Vercel span metadata', () => { 'vercel.environment', 'preview', ]) }) + + it('records the Vercel environment in configuration', () => { + process.env = { ...environment, VERCEL: '1' } + + assert.strictEqual(getServerlessPlatform().isVercel, true) + }) +}) + +describe('Vercel telemetry retention', () => { + const requestContext = Symbol.for('@vercel/request-context') + const originalContext = globalThis[requestContext] + const endpointVariables = [ + 'VERCEL', + 'OTEL_TRACES_EXPORTER', + 'DD_LOGS_OTEL_ENABLED', + 'DD_METRICS_OTEL_ENABLED', + 'OTEL_EXPORTER_OTLP_TRACES_ENDPOINT', + 'OTEL_EXPORTER_OTLP_LOGS_ENDPOINT', + 'OTEL_EXPORTER_OTLP_METRICS_ENDPOINT', + ] + const originalEndpoints = Object.fromEntries(endpointVariables.map(name => [name, process.env[name]])) + + afterEach(() => { + if (originalContext === undefined) delete globalThis[requestContext] + else globalThis[requestContext] = originalContext + for (const name of endpointVariables) { + if (originalEndpoints[name] === undefined) delete process.env[name] + else process.env[name] = originalEndpoints[name] + } + logs.disable() + metrics.disable() + }) + + it('retains trace, log, and metric payloads until their intake responses complete', async () => { + process.env.OTEL_TRACES_EXPORTER = 'otlp' + process.env.DD_LOGS_OTEL_ENABLED = 'true' + process.env.DD_METRICS_OTEL_ENABLED = 'true' + const received = new Set() + let intakeReceived + let metricPayloads = 0 + const intake = http.createServer((req, res) => { + req.resume() + req.once('end', () => { + if (req.url === '/v1/logs') received.add('logs') + if (req.url === '/v1/metrics') { + received.add('metrics') + metricPayloads++ + } + if (req.url === '/v1/traces') received.add('traces') + res.end() + if (received.size === 3 && metricPayloads === 2) intakeReceived() + }) + }) + await new Promise(resolve => intake.listen(0, '127.0.0.1', resolve)) + const { port } = intake.address() + const endpoint = `http://127.0.0.1:${port}` + process.env.OTEL_EXPORTER_OTLP_TRACES_ENDPOINT = `${endpoint}/v1/traces` + process.env.OTEL_EXPORTER_OTLP_LOGS_ENDPOINT = `${endpoint}/v1/logs` + process.env.OTEL_EXPORTER_OTLP_METRICS_ENDPOINT = `${endpoint}/v1/metrics` + + let retained + globalThis[requestContext] = { + get: () => ({ waitUntil: promise => { retained = promise } }), + } + const intakeRequests = new Promise(resolve => { intakeReceived = resolve }) + + let unregister + try { + const config = getConfigFresh({ service: 'serverless-flush' }) + const tracer = new Tracer(config) + initializeOpenTelemetryLogs(config) + initializeOpenTelemetryMetrics(config) + + tracer.trace('serverless.flush', {}, () => {}) + logs.getLogger('serverless-flush').emit({ body: 'flush me' }) + metrics.getMeter('serverless-flush').createCounter('flush.me').add(1) + + unregister = registerVercelTelemetryRetention(tracer) + channel('apm:http:server:request:finish').publish({}) + await intakeRequests + await retained + assert.deepStrictEqual(received, new Set(['traces', 'logs', 'metrics'])) + assert.strictEqual(metricPayloads, 2) + } finally { + unregister?.() + metrics.getMeterProvider()?.reader?.shutdown() + logs.getLoggerProvider()?.shutdown?.() + await new Promise(resolve => intake.close(resolve)) + } + }) + + it('waits for HTTP response completion after Next request finish', async () => { + process.env.VERCEL = '1' + let retained + globalThis[requestContext] = { + get: () => ({ waitUntil: promise => { retained = promise } }), + } + const nextFinishChannel = channel('apm:next:request:finish') + const httpFinishChannel = channel('apm:http:server:request:finish') + let finished = false + const tracer = { + flushAll (done) { + assert.ok(finished) + done() + }, + } + + const unregister = initializeServerlessTelemetry(tracer) + try { + nextFinishChannel.publish({}) + assert.strictEqual(retained, undefined) + finished = true + httpFinishChannel.publish({}) + await retained + } finally { + unregister() + } + }) + + it('retains telemetry for an ordinary HTTP Vercel request only once', async () => { + const retained = [] + let flushes = 0 + const context = { waitUntil: promise => { retained.push(promise) } } + globalThis[requestContext] = { get: () => context } + + const unregister = registerVercelTelemetryRetention({ + flushAll (done) { + flushes++ + done() + }, + }) + try { + channel('apm:http:server:request:finish').publish({}) + await Promise.all(retained) + assert.strictEqual(flushes, 1) + } finally { + unregister() + } + }) + + it('retains a configured telemetry-only pipeline without a trace exporter', async () => { + let retained + const completeTelemetry = [] + let flushes = 0 + globalThis[requestContext] = { + get: () => ({ waitUntil: promise => { retained = promise } }), + } + const telemetryFlusher = done => { + flushes++ + completeTelemetry.push(done) + } + const unregisterTelemetry = registerTelemetryFlusher(telemetryFlusher) + const unregister = registerVercelTelemetryRetention({ + flushAll: (done, options) => flushAll(undefined, done, options), + }) + try { + channel('apm:http:server:request:finish').publish({}) + await new Promise(resolve => setImmediate(resolve)) + + assert.ok(flushes >= 1) + assert.ok(completeTelemetry.every(done => typeof done === 'function')) + let settled = false + retained.then(() => { settled = true }) + await new Promise(resolve => setImmediate(resolve)) + assert.strictEqual(settled, false) + + for (const done of completeTelemetry) done() + await retained + } finally { + unregister() + unregisterTelemetry() + } + }) + + it('retains telemetry again when an outer Vercel response follows a nested request', async () => { + const retained = [] + const flushes = [] + globalThis[requestContext] = { + get: () => ({ waitUntil: promise => { retained.push(promise) } }), + } + const unregister = registerVercelTelemetryRetention({ + flushAll (done) { + flushes.push(done) + }, + }) + try { + channel('apm:http:server:request:finish').publish({ req: {} }) + await new Promise(resolve => setImmediate(resolve)) + channel('apm:http:server:request:finish').publish({ req: {} }) + await new Promise(resolve => setImmediate(resolve)) + + // Other tracers initialized by this file can share this request context; + // the callback count below isolates this test's tracer. + assert.ok(retained.length >= 2) + assert.strictEqual(flushes.length, 2) + flushes[0]() + flushes[1]() + } finally { + unregister() + } + }) + + it('passes Vercel retention timeout to the telemetry flush barrier', async () => { + let retained + let options + globalThis[requestContext] = { + get: () => ({ waitUntil: promise => { retained = promise } }), + } + + let unregister + try { + unregister = registerVercelTelemetryRetention({ + flushAll (done, flushOptions) { + options = flushOptions + done() + }, + }) + channel('apm:http:server:request:finish').publish({}) + await retained + + assert.deepStrictEqual(options, { timeout: 2_000 }) + } finally { + unregister?.() + } + }) + + it('retains telemetry at HTTP/2 response completion', async () => { + let retained + globalThis[requestContext] = { + get: () => ({ waitUntil: promise => { retained = promise } }), + } + + let flushes = 0 + const unregister = registerVercelTelemetryRetention({ + flushAll (done) { + flushes++ + done() + }, + }) + try { + channel('apm:http2:server:response:emit').publish({ eventName: 'finish' }) + assert.strictEqual(retained, undefined) + channel('apm:http2:server:response:emit').publish({ eventName: 'close' }) + await retained + assert.strictEqual(flushes, 1) + } finally { + unregister() + } + }) }) diff --git a/packages/dd-trace/test/span_stats.spec.js b/packages/dd-trace/test/span_stats.spec.js index 3652c8fe168..4152d109854 100644 --- a/packages/dd-trace/test/span_stats.spec.js +++ b/packages/dd-trace/test/span_stats.spec.js @@ -81,12 +81,14 @@ const syntheticSpan = { const exporter = { export: sinon.stub(), + flush: sinon.stub(), } const SpanStatsExporter = sinon.stub().returns(exporter) const otlpExporter = { export: sinon.stub(), + flush: sinon.stub(), } const { @@ -643,6 +645,83 @@ describe('SpanStatsProcessor', () => { assert.ok(otlpExporter.export.calledOnce) }) + it('force flushes pending OTLP span statistics', () => { + const exporter = { + export: sinon.stub().callsFake((_drained, _bucketSizeNs, done) => done()), + flush: sinon.stub().callsFake(done => done()), + } + const p = new SpanStatsProcessor(config, exporter) + clearTimeout(p.timer) + p.onSpanFinished(topLevelSpan) + + let flushed = false + p.forceFlush(() => { flushed = true }) + + assert.ok(exporter.export.calledOnce) + assert.ok(exporter.flush.calledOnce) + assert.ok(flushed) + assert.strictEqual(p.buckets.size, 0) + }) + + it('snapshots prior OTLP exports before starting the boundary export', () => { + let priorDone + let exportDone + const exporter = { + flush: sinon.stub().callsFake(done => { priorDone = done }), + export: sinon.stub().callsFake((_drained, _bucketSizeNs, done) => { exportDone = done }), + } + const p = new SpanStatsProcessor(config, exporter) + clearTimeout(p.timer) + p.onSpanFinished(topLevelSpan) + const done = sinon.spy() + + p.forceFlush(done) + + sinon.assert.callOrder(exporter.flush, exporter.export) + exportDone() + sinon.assert.notCalled(done) + priorDone() + sinon.assert.calledOnce(done) + }) + + it('force flushes pending agent span statistics', () => { + exporter.export.resetHistory() + exporter.flush.resetHistory() + exporter.export.callsFake((_payload, done) => done()) + const p = new SpanStatsProcessor(config) + clearTimeout(p.timer) + p.onSpanFinished(topLevelSpan) + + let flushed = false + p.forceFlush(() => { flushed = true }) + + assert.ok(exporter.export.calledOnce) + assert.ok(exporter.flush.notCalled) + assert.ok(flushed) + assert.strictEqual(p.buckets.size, 0) + exporter.export.resetBehavior() + }) + + it('joins an in-flight agent span statistics export during force flush', () => { + exporter.export.resetHistory() + exporter.flush.resetHistory() + const p = new SpanStatsProcessor(config) + clearTimeout(p.timer) + p.onSpanFinished(topLevelSpan) + + let flushDone + exporter.export.callsFake((_payload, done) => { flushDone = done }) + let flushed = false + p.forceFlush(() => { flushed = true }) + + assert.ok(exporter.export.calledOnce) + assert.ok(exporter.flush.notCalled) + assert.strictEqual(flushed, false) + flushDone() + assert.strictEqual(flushed, true) + exporter.export.resetBehavior() + }) + it('should record spans when only OTLP is enabled', () => { otlpExporter.export.resetHistory() const p = new SpanStatsProcessor({ diff --git a/packages/dd-trace/test/tracer.spec.js b/packages/dd-trace/test/tracer.spec.js index ee9c4d5ec05..8b7a5062502 100644 --- a/packages/dd-trace/test/tracer.spec.js +++ b/packages/dd-trace/test/tracer.spec.js @@ -50,6 +50,78 @@ describe('Tracer', () => { }) }) + describe('flushAll', () => { + it('flushes registered telemetry pipelines with the configured trace exporter', () => { + const { flushAll, registerTelemetryFlusher } = require('../src/flush') + const tracer = { + _exporter: { + flush: sinon.stub().callsFake(done => done()), + }, + } + const telemetryFlusher = sinon.stub().callsFake(done => done()) + const unregister = registerTelemetryFlusher(telemetryFlusher) + let completed = false + + flushAll(tracer, () => { completed = true }) + + sinon.assert.calledOnce(tracer._exporter.flush) + sinon.assert.calledOnce(telemetryFlusher) + assert.strictEqual(completed, true) + unregister() + }) + + it('flushes post-trace telemetry after the trace exporter completes', () => { + const { flushAll, registerTelemetryFlusher } = require('../src/flush') + let traceDone + const tracer = { _exporter: { flush: sinon.stub().callsFake(done => { traceDone = done }) } } + const runtimeMetricsFlusher = sinon.stub().callsFake(done => done()) + const unregister = registerTelemetryFlusher(runtimeMetricsFlusher, { afterTrace: true }) + const done = sinon.spy() + + flushAll(tracer, done) + + sinon.assert.notCalled(runtimeMetricsFlusher) + traceDone() + sinon.assert.calledOnce(runtimeMetricsFlusher) + sinon.assert.calledOnce(done) + unregister() + }) + + it('flushes registered telemetry pipelines without a trace exporter', () => { + const { flushAll, registerTelemetryFlusher } = require('../src/flush') + const telemetryFlusher = sinon.stub().callsFake(done => done()) + const unregister = registerTelemetryFlusher(telemetryFlusher) + const done = sinon.spy() + + flushAll(undefined, done) + + sinon.assert.calledOnce(telemetryFlusher) + sinon.assert.calledOnce(done) + unregister() + }) + + it('bounds configured telemetry flushing', () => { + const { flushAll, registerTelemetryFlusher } = require('../src/flush') + const timeout = sinon.stub(global, 'setTimeout') + const clearTimeout = sinon.stub(global, 'clearTimeout') + const done = sinon.spy() + const unregister = registerTelemetryFlusher(() => {}) + + try { + flushAll({}, done, { timeout: 2_000 }) + + sinon.assert.calledWith(timeout, sinon.match.func, 2_000) + timeout.firstCall.args[0]() + sinon.assert.calledOnce(done) + sinon.assert.called(clearTimeout) + } finally { + unregister() + timeout.restore() + clearTimeout.restore() + } + }) + }) + describe('trace', () => { it('should run the callback with a new span', () => { tracer.trace('name', {}, span => {