diff --git a/docs/API.md b/docs/API.md index 0a3e31f37ff..799b8c0e4e0 100644 --- a/docs/API.md +++ b/docs/API.md @@ -536,6 +536,10 @@ cpuGauge.addCallback((result) => { }) ``` +Short-lived processes can use the callback-based `meterProvider.forceFlush(callback)` extension to wait for export. +Call `meterProvider.shutdown(callback)` after the final measurement to export once more and stop collection. The optional +callbacks receive an error when the operation fails. These methods aren't part of the OpenTelemetry Metrics API. + #### Supported Configuration The Datadog SDK supports many of the configurations supported by the OpenTelemetry SDK. The following environment variables are supported: diff --git a/docs/test.ts b/docs/test.ts index 49e80a27844..a08c1170279 100644 --- a/docs/test.ts +++ b/docs/test.ts @@ -566,6 +566,9 @@ const provider: opentelemetry.TracerProvider = new tracer.TracerProvider(); provider.register(); const otelTracer: opentelemetry.Tracer = provider.getTracer("name", "version") +const otelMeterProvider = {} as opentelemetry.MeterProvider +const otelForceFlush: (callback?: (error?: Error) => void) => void = otelMeterProvider.forceFlush +const otelShutdown: (callback?: (error?: Error) => void) => void = otelMeterProvider.shutdown // OTel supports several time input formats otelTracer.startSpan("name", { startTime: new Date() }) diff --git a/index.d.ts b/index.d.ts index 62fc054c974..9c89078165d 100644 --- a/index.d.ts +++ b/index.d.ts @@ -3279,6 +3279,11 @@ declare namespace tracer { } export namespace opentelemetry { + export interface MeterProvider extends otel.MeterProvider { + forceFlush(callback?: (error?: Error) => void): void; + shutdown(callback?: (error?: Error) => void): void; + } + /** * A registry for creating named {@link Tracer}s. */ diff --git a/index.d.v5.ts b/index.d.v5.ts index 888bc82237d..cb0412086ee 100644 --- a/index.d.v5.ts +++ b/index.d.v5.ts @@ -3465,6 +3465,11 @@ declare namespace tracer { } export namespace opentelemetry { + export interface MeterProvider extends otel.MeterProvider { + forceFlush(callback?: (error?: Error) => void): void; + shutdown(callback?: (error?: Error) => void): void; + } + /** * A registry for creating named {@link Tracer}s. */ diff --git a/packages/dd-trace/src/opentelemetry/metrics/meter_provider.js b/packages/dd-trace/src/opentelemetry/metrics/meter_provider.js index 53cfcb20a57..2782a81bd0e 100644 --- a/packages/dd-trace/src/opentelemetry/metrics/meter_provider.js +++ b/packages/dd-trace/src/opentelemetry/metrics/meter_provider.js @@ -57,6 +57,14 @@ class MeterProvider { if (this.reader) this.reader.forceFlush(done) else done?.() } + + /** + * @param {Function} [done] Called after shutdown completes + */ + shutdown (done) { + if (this.reader) this.reader.shutdown(done) + else done?.() + } } module.exports = MeterProvider 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 3bdf1408f80..27931dce31c 100644 --- a/packages/dd-trace/src/opentelemetry/metrics/periodic_metric_reader.js +++ b/packages/dd-trace/src/opentelemetry/metrics/periodic_metric_reader.js @@ -105,6 +105,10 @@ class PeriodicMetricReader { #exportInterval #aggregator #batchCallbacks = [] + #exportQueue = [] + #isExporting = false + #shutdownCallbacks = [] + #shutdownComplete = false /** * Creates a new PeriodicMetricReader instance. @@ -205,39 +209,25 @@ class PeriodicMetricReader { done?.() return } - let pending = 2 - const complete = () => { - if (--pending === 0) done?.() - } - - // Snapshot requests already active before starting this flush's export. - try { - if (typeof this.exporter.flush === 'function') this.exporter.flush(complete) - else complete() - } catch (error) { - log.error('Error flushing OTLP metrics:', error) - complete() - } - try { - this.#collectAndExport(complete) - } catch (error) { - log.error('Error exporting OTLP metrics:', error) - complete() - } + this.#enqueueExport(true, done) } /** * Shuts down the reader and stops periodic collection. - * @returns {void} + * @param {Function} [done] Called after the final export and exporter shutdown complete */ - shutdown () { + shutdown (done) { if (this.#isShutdown) { log.warn('PeriodicMetricReader is already shutdown') + if (this.#shutdownComplete) done?.() + else if (done) this.#shutdownCallbacks.push(done) return } + if (done) this.#shutdownCallbacks.push(done) + this.#isShutdown = true this.#clearTimer() - this.forceFlush() + this.#enqueueExport(true, error => this.#shutdownExporter(error)) } /** @@ -248,7 +238,7 @@ class PeriodicMetricReader { if (this.#timer) return this.#timer = setInterval(() => { - this.#collectAndExport() + if (!this.#isExporting && this.#exportQueue.length === 0) this.#enqueueExport(false) }, this.#exportInterval) this.#timer.unref?.() } @@ -264,6 +254,65 @@ class PeriodicMetricReader { } } + #enqueueExport (flushExporter, done) { + this.#exportQueue.push({ flushExporter, done }) + this.#drainExportQueue() + } + + #drainExportQueue () { + if (this.#isExporting || this.#exportQueue.length === 0) return + + this.#isExporting = true + const { flushExporter, done } = this.#exportQueue.shift() + let completed = false + const complete = error => { + if (completed) return + completed = true + this.#isExporting = false + queueMicrotask(() => this.#drainExportQueue()) + done?.(error) + } + let exportCompleted = false + const afterExport = error => { + if (exportCompleted) return + exportCompleted = true + if (!flushExporter || typeof this.exporter.flush !== 'function') return complete(error) + + try { + this.exporter.flush(flushError => complete(error || flushError)) + } catch (flushError) { + log.error('Error flushing OTLP metrics:', flushError) + complete(error || flushError) + } + } + + try { + this.#collectAndExport(afterExport) + } catch (error) { + log.error('Error exporting OTLP metrics:', error) + afterExport(error) + } + } + + #shutdownExporter (exportError) { + let completed = false + const complete = shutdownError => { + if (completed) return + completed = true + this.#shutdownComplete = true + const callbacks = this.#shutdownCallbacks.splice(0) + for (const callback of callbacks) callback(exportError || shutdownError) + } + + try { + if (typeof this.exporter.shutdown === 'function') this.exporter.shutdown(complete) + else complete() + } catch (error) { + log.error('Error shutting down OTLP metrics exporter:', error) + complete(error) + } + } + /** * Collects measurements and exports metrics. * @@ -322,7 +371,13 @@ class PeriodicMetricReader { this.#lastExportedState ) - this.exporter.export(metrics, callback) + this.exporter.export(metrics, result => { + if (result?.code === 1) { + callback?.(result.error || new Error('OTLP metrics export failed')) + } else { + callback?.() + } + }) } } 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 b5034396c7d..db17f69ecd3 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 @@ -172,7 +172,9 @@ class OtlpHttpExporterBase { this.telemetryTags[0] = `protocol:${this.#transport === https ? 'https' : 'http'}` } - shutdown () {} + shutdown (done) { + done?.() + } } module.exports = OtlpHttpExporterBase diff --git a/packages/dd-trace/test/opentelemetry/metrics.spec.js b/packages/dd-trace/test/opentelemetry/metrics.spec.js index fd11cc31618..f3ea94f99e9 100644 --- a/packages/dd-trace/test/opentelemetry/metrics.spec.js +++ b/packages/dd-trace/test/opentelemetry/metrics.spec.js @@ -667,12 +667,38 @@ describe('OpenTelemetry Meter Provider', () => { }) describe('Lifecycle', () => { - it('waits for an in-flight export during forceFlush', () => { + function completeNext (callbacks) { + const callback = callbacks.shift() + callback() + } + + it('exposes callback lifecycle methods on the global meter provider', () => { + const callbacks = {} + const reader = { + forceFlush: done => { callbacks.forceFlush = done }, + shutdown: done => { callbacks.shutdown = done }, + } + const provider = new MeterProvider({ reader }) + const forceFlushDone = sinon.spy() + const shutdownDone = sinon.spy() + + assert.strictEqual(provider.forceFlush(forceFlushDone), undefined) + sinon.assert.notCalled(forceFlushDone) + callbacks.forceFlush() + sinon.assert.calledOnce(forceFlushDone) + + assert.strictEqual(provider.shutdown(shutdownDone), undefined) + sinon.assert.notCalled(shutdownDone) + callbacks.shutdown() + sinon.assert.calledOnce(shutdownDone) + }) + + it('serializes forceFlush exports', async () => { const exports = [] const flushes = [] const reader = new PeriodicMetricReader({ export: (metrics, done) => { exports.push(done) }, - flush: (done) => { flushes.push(done) }, + flush: done => { flushes.push(done) }, }, 60_000, 'DELTA', 1024) const meter = new MeterProvider({ reader }).getMeter('test') const firstDone = sinon.spy() @@ -680,19 +706,23 @@ describe('OpenTelemetry Meter Provider', () => { meter.createCounter('in-flight').add(1) reader.forceFlush(firstDone) - const flush = flushes.shift() - flush() meter.createCounter('boundary').add(1) reader.forceFlush(done) sinon.assert.notCalled(done) + assert.strictEqual(exports.length, 1) + exports[0]({ code: 0 }) + sinon.assert.notCalled(firstDone) + completeNext(flushes) + sinon.assert.calledOnce(firstDone) + await Promise.resolve() assert.strictEqual(exports.length, 2) exports[1]({ code: 0 }) sinon.assert.notCalled(done) - flushes[0]() + completeNext(flushes) sinon.assert.calledOnce(done) - exports[0]({ code: 0 }) reader.shutdown() + completeNext(flushes) }) it('waits for an earlier export when the boundary export throws', () => { @@ -709,18 +739,54 @@ describe('OpenTelemetry Meter Provider', () => { sinon.assert.notCalled(done) priorDone() - sinon.assert.calledOnce(done) + sinon.assert.calledOnceWithMatch(done, { message: 'encode failed' }) reader.shutdown() + priorDone() }) - it('handles shutdown gracefully', async () => { + it('waits for an in-flight export before shutting down', async () => { + const exports = [] + const flushes = [] + const exporter = { + export: (metrics, done) => { exports.push(done) }, + flush: done => { flushes.push(done) }, + shutdown: sinon.spy(done => { done() }), + } + const reader = new PeriodicMetricReader(exporter, 60_000, 'DELTA', 1024) + const provider = new MeterProvider({ reader }) + const forceFlushDone = sinon.spy() + const done = sinon.spy() + + provider.getMeter('test').createCounter('requests').add(1) + provider.forceFlush(forceFlushDone) + provider.getMeter('test').createCounter('final').add(1) + assert.strictEqual(provider.shutdown(done), undefined) + + assert.strictEqual(exports.length, 1) + sinon.assert.notCalled(exporter.shutdown) + exports[0]({ code: 0 }) + completeNext(flushes) + sinon.assert.calledOnce(forceFlushDone) + await Promise.resolve() + + assert.strictEqual(exports.length, 2) + sinon.assert.notCalled(exporter.shutdown) + exports[1]({ code: 0 }) + completeNext(flushes) + sinon.assert.calledOnce(exporter.shutdown) + sinon.assert.calledOnce(done) + }) + + it('handles shutdown gracefully', (done) => { setupMetrics() const provider = metrics.getMeterProvider() - await provider.reader.shutdown() - await provider.reader.shutdown() // Second shutdown should be safe + provider.shutdown(error => { + assert.ifError(error) + provider.shutdown(done) + }) }) - it('handles forceFlush', async () => { + it('handles forceFlush', (done) => { const validator = mockOtlpExport((decoded) => { assert(decoded.resourceMetrics) }) @@ -730,8 +796,11 @@ describe('OpenTelemetry Meter Provider', () => { meter.createCounter('test').add(1) const provider = metrics.getMeterProvider() - await provider.reader.forceFlush() - validator() + provider.forceFlush(error => { + assert.ifError(error) + validator() + done() + }) }) it('removes callbacks from observable instruments', (done) => {