diff --git a/packages/dd-trace/src/exporters/agent/index.js b/packages/dd-trace/src/exporters/agent/index.js index e0f1ae9a22b..bc039206c02 100644 --- a/packages/dd-trace/src/exporters/agent/index.js +++ b/packages/dd-trace/src/exporters/agent/index.js @@ -1,6 +1,7 @@ 'use strict' const { URL } = require('url') +const getFlushError = require('../../flush-error') const log = require('../../log') const { createServerlessDeliveryTracker } = require('../../serverless') const Writer = require('./writer') @@ -59,26 +60,42 @@ class AgentExporter { } } - flush (done) { + /** + * @param {(error?: Error) => void} [done] + * @param {{ reportErrors?: boolean }} [options] + */ + flush (done, options) { clearTimeout(this.#timer) this.#timer = undefined if (!this.#serverlessDeliveryTracker) { try { - return this._writer.flush(done) + return this._writer.flush(done, options) } catch (error) { log.error('Failed to flush traces: %s', error.message) - done?.() + done?.(options?.reportErrors ? error : undefined) return } } + let boundaryError + let waiting = false + const captureError = error => { + if (!waiting) boundaryError = error + } try { - this._writer.flush() + this._writer.flush(captureError, options) } catch (error) { log.error('Failed to flush traces: %s', error.message) + boundaryError = error } - this.#serverlessDeliveryTracker.waitForIdle(done) + waiting = true + if (!done) return + + this.#serverlessDeliveryTracker.waitForIdle(error => { + if (!options?.reportErrors || !boundaryError) return done(error) + done(getFlushError(error ? [boundaryError, error] : [boundaryError])) + }, options) } } diff --git a/packages/dd-trace/src/exporters/agent/writer.js b/packages/dd-trace/src/exporters/agent/writer.js index a815b45f625..ab0c71a1b6f 100644 --- a/packages/dd-trace/src/exporters/agent/writer.js +++ b/packages/dd-trace/src/exporters/agent/writer.js @@ -43,7 +43,7 @@ class AgentWriter extends BaseWriter { * Performs the writer flush without registering a serverless delivery. * Test Optimization owns its own request lifecycle tracking. * @param {(error?: Error) => void} [done] - * @param {{ deadline?: number }} [options] + * @param {{ deadline?: number, reportErrors?: boolean }} [options] * @returns {void} */ flushDirect (done, options) { @@ -73,7 +73,8 @@ class AgentWriter extends BaseWriter { if (err) { log.errorWithoutTelemetry('Error sending payload to the agent (status code: %s)', err.status, err) - done(flushOptions?.deadline === undefined ? undefined : err) + const reportError = flushOptions?.reportErrors || flushOptions?.deadline !== undefined + done(reportError ? err : undefined) return } diff --git a/packages/dd-trace/src/exporters/common/writer.js b/packages/dd-trace/src/exporters/common/writer.js index e52a950e5b8..6c9a293ccb7 100644 --- a/packages/dd-trace/src/exporters/common/writer.js +++ b/packages/dd-trace/src/exporters/common/writer.js @@ -26,7 +26,7 @@ class Writer { /** * Flushes queued telemetry, retaining delivery on supported serverless platforms. * @param {(error?: Error) => void} [done] - * @param {{ deadline?: number }} [options] + * @param {{ deadline?: number, reportErrors?: boolean }} [options] * @returns {void} */ flush (done, options) { @@ -39,7 +39,7 @@ class Writer { /** * Flushes queued telemetry without registering serverless delivery retention. * @param {(error?: Error) => void} [done] - * @param {{ deadline?: number }} [options] + * @param {{ deadline?: number, reportErrors?: boolean }} [options] * @returns {void} */ flushDirect (done = noop, options) { @@ -47,7 +47,13 @@ class Writer { if (!request.writable && options?.deadline === undefined && !this.#retainOnBackpressure) { this._encoder.reset() - done() + if (options?.reportErrors) { + const error = new log.NoTransmitError('Maximum active request buffer size reached: payload is discarded.') + error.code = 'ERR_DD_REQUEST_BUFFER_FULL' + done(error) + } else { + done() + } } else if (count > 0) { if (this.#isFirstFlush && firstFlushChannel.hasSubscribers && this._beforeFirstFlush) { this.#isFirstFlush = false @@ -66,7 +72,7 @@ class Writer { // the oversized payload at the network boundary anyway. this._encoder.reset() log.error('Writer dropped %d trace(s) that exceeded the %d byte chunk cap', count, MAX_CHUNK_SIZE) - done() + done(options?.reportErrors ? error : undefined) return } if (options === undefined) { diff --git a/packages/dd-trace/src/flush-error.js b/packages/dd-trace/src/flush-error.js new file mode 100644 index 00000000000..5f528c6979d --- /dev/null +++ b/packages/dd-trace/src/flush-error.js @@ -0,0 +1,34 @@ +'use strict' + +/** + * @param {unknown} reason + * @param {unknown[]} reasons + */ +function collectReason (reason, reasons) { + if (reason instanceof AggregateError) { + for (const error of reason.errors) collectReason(error, reasons) + return + } + + reasons.push(reason) +} + +/** + * Preserves a single rejection reason and aggregates independent failures. + * @param {unknown[]} flushReasons + * @returns {unknown} + */ +function getFlushError (flushReasons) { + let flushError + if (flushReasons.length > 0) { + const reasons = [] + for (const reason of flushReasons) collectReason(reason, reasons) + + flushError = reasons.length === 1 + ? reasons[0] + : new AggregateError(reasons, 'Multiple errors occurred while flushing') + } + return flushError +} + +module.exports = getFlushError 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..30930e5948b 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 @@ -5,7 +5,7 @@ const https = require('node:https') const { URL } = require('node:url') const { storage } = require('../../../../datadog-core') const log = require('../../log') -const { createServerlessDeliveryTracker } = require('../../serverless') +const TelemetryDeliveryTracker = require('../../serverless/telemetry-delivery-tracker') const telemetryMetrics = require('../../telemetry/metrics') const tracerMetrics = telemetryMetrics.manager.namespace('tracers') @@ -20,8 +20,8 @@ const legacyStorage = storage('legacy') * @class OtlpHttpExporterBase */ class OtlpHttpExporterBase { + #deliveryTracker = new TelemetryDeliveryTracker() #transport = https - #serverlessDeliveryTracker /** * Creates a new OtlpHttpExporterBase instance. @@ -34,7 +34,6 @@ class OtlpHttpExporterBase { * @param {string} signalType - Signal type for error messages (e.g., 'logs', 'metrics') */ constructor (url, headers, timeout, protocol, signalType) { - this.#serverlessDeliveryTracker = createServerlessDeliveryTracker() this.protocol = protocol this.signalType = signalType @@ -83,10 +82,7 @@ class OtlpHttpExporterBase { * @protected */ sendPayload (payload, resultCallback) { - if (this.#serverlessDeliveryTracker) { - return this.#serverlessDeliveryTracker.track(done => this.#sendPayload(payload, resultCallback, done)) - } - this.#sendPayload(payload, resultCallback) + this.#deliveryTracker.track(done => this.#sendPayload(payload, resultCallback, done)) } #sendPayload (payload, resultCallback, done) { @@ -103,7 +99,7 @@ class OtlpHttpExporterBase { if (completed) return completed = true resultCallback(result) - done?.() + done?.(result.error) } try { @@ -151,12 +147,12 @@ class OtlpHttpExporterBase { } /** - * Calls back once Vercel-tracked requests active at the flush boundary complete. - * @param {Function} [done] + * Calls back once requests active at the flush boundary complete. + * @param {(error?: Error) => void} [done] + * @param {{ reportErrors?: boolean }} [options] */ - flush (done) { - if (this.#serverlessDeliveryTracker) return this.#serverlessDeliveryTracker.waitForIdle(done) - done?.() + flush (done, options) { + this.#deliveryTracker.waitForIdle(done, options) } /** diff --git a/packages/dd-trace/src/opentelemetry/span_processor.js b/packages/dd-trace/src/opentelemetry/span_processor.js index 71bbc5df63b..8b3ee5bbd44 100644 --- a/packages/dd-trace/src/opentelemetry/span_processor.js +++ b/packages/dd-trace/src/opentelemetry/span_processor.js @@ -1,5 +1,33 @@ 'use strict' +const getFlushError = require('../flush-error') + +/** + * @typedef {{ status: 'fulfilled', value: unknown } | { status: 'rejected', reason: unknown }} FlushResult + */ + +/** + * @param {FlushResult[]} results + * @returns {Promise | undefined} + */ +function collectFlushErrors (results) { + const reasons = [] + for (const result of results) { + if (result.status !== 'rejected') continue + reasons.push(result.reason) + } + + if (reasons.length > 0) return Promise.reject(getFlushError(reasons)) +} + +/** + * @param {Promise[]} flushes + * @returns {Promise} + */ +function settleAllFlushes (flushes) { + return Promise.allSettled(flushes).then(collectFlushErrors) +} + class NoopSpanProcessor { forceFlush () { return Promise.resolve() @@ -22,9 +50,9 @@ class MultiSpanProcessor extends NoopSpanProcessor { } forceFlush () { - return Promise.all( - this.#processors.map(p => p.forceFlush()) - ) + const flushes = [] + for (const processor of this.#processors) flushes.push(processor.forceFlush()) + return settleAllFlushes(flushes) } onStart (span, context) { @@ -49,4 +77,5 @@ class MultiSpanProcessor extends NoopSpanProcessor { module.exports = { MultiSpanProcessor, NoopSpanProcessor, + settleAllFlushes, } diff --git a/packages/dd-trace/src/opentelemetry/tracer_provider.js b/packages/dd-trace/src/opentelemetry/tracer_provider.js index 910c464cc84..c61c4df8521 100644 --- a/packages/dd-trace/src/opentelemetry/tracer_provider.js +++ b/packages/dd-trace/src/opentelemetry/tracer_provider.js @@ -6,12 +6,45 @@ const { W3CTraceContextPropagator } = require('../../../../vendor/dist/@opentele const tracer = require('../../') const ContextManager = require('./context_manager') -const { MultiSpanProcessor, NoopSpanProcessor } = require('./span_processor') +const { MultiSpanProcessor, NoopSpanProcessor, settleAllFlushes } = require('./span_processor') const Tracer = require('./tracer') +/** + * @typedef {{ + * flush?: (done?: (error?: Error) => void, options?: { reportErrors?: boolean }) => void + * }} TraceExporter + */ + +/** + * @param {TraceExporter} exporter + * @returns {Promise} + */ +function flushExporter (exporter) { + if (typeof exporter.flush !== 'function') return Promise.resolve() + + /** + * @param {() => void} resolve + * @param {(reason?: unknown) => void} reject + */ + function flush (resolve, reject) { + /** + * @param {Error} [error] + */ + function done (error) { + if (error) reject(error) + else resolve() + } + + exporter.flush(done, { reportErrors: true }) + } + + return new Promise(flush) +} + class TracerProvider { #activeProcessor = new NoopSpanProcessor() #contextManager = new ContextManager() + #flush #processors = [] #tracers = new Map() @@ -84,8 +117,19 @@ class TracerProvider { return Promise.reject(new Error('Not started')) } - exporter._writer?.flush() - return this.#activeProcessor.forceFlush() + const flush = () => settleAllFlushes([ + flushExporter(exporter), + this.#activeProcessor.forceFlush(), + ]) + const pending = this.#flush ? this.#flush.then(flush, flush) : flush() + this.#flush = pending + + const clear = () => { + if (this.#flush === pending) this.#flush = undefined + } + pending.then(clear, clear) + + return pending } shutdown () { diff --git a/packages/dd-trace/src/serverless/telemetry-delivery-tracker.js b/packages/dd-trace/src/serverless/telemetry-delivery-tracker.js index 6b20a55f737..e4b654f91b8 100644 --- a/packages/dd-trace/src/serverless/telemetry-delivery-tracker.js +++ b/packages/dd-trace/src/serverless/telemetry-delivery-tracker.js @@ -1,48 +1,56 @@ 'use strict' +const getFlushError = require('../flush-error') + /** - * Tracks transport deliveries that must outlive a serverless request. + * Tracks transport deliveries until a completion boundary. * - * The tracker is created only for platforms with an invocation-retention - * boundary. Exporters keep their normal callback path when it is absent. + * Serverless exporters use the boundary to retain an invocation. OTLP + * exporters use it to make explicit flushes wait for active HTTP requests. */ class TelemetryDeliveryTracker { #deliveries = new Set() /** * Tracks one asynchronous transport delivery until its callback runs. - * @param {(done: () => void) => void} deliver - * @param {(() => void)|undefined} done + * @param {(done: (error?: Error) => void) => void} deliver + * @param {((error?: Error) => void)|undefined} done */ track (deliver, done) { const delivery = { callbacks: done ? [done] : [] } this.#deliveries.add(delivery) - const complete = () => { + const complete = (error) => { if (!this.#deliveries.delete(delivery)) return - for (const callback of delivery.callbacks) callback() + for (const callback of delivery.callbacks) callback(error) } try { deliver(complete) } catch (error) { - complete() + complete(error) throw error } } /** * Calls back after every delivery active at this boundary has completed. - * @param {(() => void)|undefined} done + * @param {((error?: Error) => void)|undefined} done + * @param {{ reportErrors?: boolean }} [options] */ - waitForIdle (done) { + waitForIdle (done, options) { if (!done) return + const errors = [] let pending = this.#deliveries.size - if (pending === 0) return done() + const finish = () => { + done(options?.reportErrors ? getFlushError(errors) : undefined) + } + if (pending === 0) return finish() - const complete = () => { - if (--pending === 0) done() + const complete = (error) => { + if (error) errors.push(error) + if (--pending === 0) finish() } for (const delivery of this.#deliveries) delivery.callbacks.push(complete) } diff --git a/packages/dd-trace/test/exporters/agent/exporter.spec.js b/packages/dd-trace/test/exporters/agent/exporter.spec.js index 8ba3fc86fb3..93900c7af5d 100644 --- a/packages/dd-trace/test/exporters/agent/exporter.spec.js +++ b/packages/dd-trace/test/exporters/agent/exporter.spec.js @@ -212,6 +212,87 @@ describe('Exporter', () => { complete() sinon.assert.calledOnce(flushed) }) + + it('reports a serverless delivery failure when requested', () => { + const error = new Error('agent failed') + writer.flush = sinon.spy(done => { + writerOptions.deliveryTracker.track(complete => complete(error), done) + }) + exporter = new Exporter({ url, flushInterval: 0 }, prioritySampler) + const flushed = sinon.spy() + + exporter.flush(flushed, { reportErrors: true }) + + sinon.assert.calledOnceWithExactly(flushed, error) + }) + + it('does not report a serverless delivery failure by default', () => { + writer.flush = sinon.spy(done => { + writerOptions.deliveryTracker.track(complete => complete(new Error('agent failed')), done) + }) + exporter = new Exporter({ url, flushInterval: 0 }, prioritySampler) + const flushed = sinon.spy() + + exporter.flush(flushed) + + sinon.assert.calledOnceWithExactly(flushed, undefined) + }) + + it('supports a serverless flush without a completion callback', () => { + exporter = new Exporter({ url, flushInterval: 0 }, prioritySampler) + + exporter.flush() + + sinon.assert.calledOnce(writer.flush) + }) + + it('aggregates a synchronous boundary failure with an in-flight failure', () => { + const inFlightError = new Error('in-flight request failed') + const boundaryError = new Error('boundary flush failed') + let completeInFlight + writer.flush = sinon.stub() + writer.flush.onFirstCall().callsFake(done => { + writerOptions.deliveryTracker.track(complete => { completeInFlight = complete }, done) + }) + writer.flush.onSecondCall().callsFake(done => { + writerOptions.deliveryTracker.track(() => { throw boundaryError }, done) + }) + exporter = new Exporter({ url, flushInterval: 0 }, prioritySampler) + const flushed = sinon.spy() + + exporter.export([span]) + exporter.flush(flushed, { reportErrors: true }) + sinon.assert.notCalled(flushed) + completeInFlight(inFlightError) + + sinon.assert.calledOnce(flushed) + const error = flushed.firstCall.firstArg + assert.ok(error instanceof AggregateError) + assert.deepStrictEqual(error.errors, [boundaryError, inFlightError]) + }) + + it('reports a boundary flush failure when requested', () => { + createServerlessDeliveryTracker.resetBehavior() + const error = new Error('encode failed') + writer.flush = sinon.stub().throws(error) + exporter = new Exporter({ url, flushInterval: 0 }, prioritySampler) + const flushed = sinon.spy() + + exporter.flush(flushed, { reportErrors: true }) + + sinon.assert.calledOnceWithExactly(flushed, error) + }) + + it('does not report a boundary flush failure by default', () => { + createServerlessDeliveryTracker.resetBehavior() + writer.flush = sinon.stub().throws(new Error('encode failed')) + exporter = new Exporter({ url, flushInterval: 0 }, prioritySampler) + const flushed = sinon.spy() + + exporter.flush(flushed) + + sinon.assert.calledOnceWithExactly(flushed, undefined) + }) }) describe('setUrl', () => { diff --git a/packages/dd-trace/test/exporters/common/writer.spec.js b/packages/dd-trace/test/exporters/common/writer.spec.js index 79fd04547b9..6217b3885d8 100644 --- a/packages/dd-trace/test/exporters/common/writer.spec.js +++ b/packages/dd-trace/test/exporters/common/writer.spec.js @@ -57,6 +57,16 @@ describe('common Writer', () => { sinon.assert.calledOnce(done) }) + it('reports a chunk overflow when requested', () => { + const error = new OverflowError(MAX_SIZE + 1) + encoder.makePayload.throws(error) + const done = sinon.stub() + + writer.flush(done, { reportErrors: true }) + + sinon.assert.calledOnceWithExactly(done, error) + }) + it('rethrows non-overflow makePayload errors', () => { encoder.makePayload.throws(new Error('not an overflow')) @@ -112,6 +122,18 @@ describe('common Writer', () => { sinon.assert.calledOnceWithExactly(done) }) + it('reports a dropped payload when the request buffer is full and errors are requested', () => { + request.writable = false + const done = sinon.stub() + + writer.flush(done, { reportErrors: true }) + + sinon.assert.calledOnceWithMatch(done, { + code: 'ERR_DD_REQUEST_BUFFER_FULL', + message: 'Maximum active request buffer size reached: payload is discarded.', + }) + }) + it('retains a non-final payload under backpressure when configured', () => { request.writable = false writer = new Writer({ url: 'http://localhost:8126', retainOnBackpressure: true }) diff --git a/packages/dd-trace/test/opentelemetry/tracer_provider.spec.js b/packages/dd-trace/test/opentelemetry/tracer_provider.spec.js index 8c758f6ff7e..d24abafa158 100644 --- a/packages/dd-trace/test/opentelemetry/tracer_provider.spec.js +++ b/packages/dd-trace/test/opentelemetry/tracer_provider.spec.js @@ -1,8 +1,10 @@ 'use strict' const assert = require('node:assert/strict') +const { EventEmitter, once } = require('node:events') +const http = require('node:http') -const { describe, it } = require('mocha') +const { after, before, describe, it } = require('mocha') const sinon = require('sinon') const { trace } = require('@opentelemetry/api') @@ -10,7 +12,62 @@ require('../setup/core') const TracerProvider = require('../../src/opentelemetry/tracer_provider') const Tracer = require('../../src/opentelemetry/tracer') const { MultiSpanProcessor, NoopSpanProcessor } = require('../../src/opentelemetry/span_processor') -require('../../index').init() + +const agentResponse = JSON.stringify({ rate_by_service: {} }) + +let traceRequestCount = 0 +let traceRequestEvents +let traceRequestResolve + +/** + * @param {(response: import('node:http').ServerResponse) => void} resolve + */ +function captureTraceRequest (resolve) { + traceRequestResolve = resolve +} + +/** + * @param {string[]} [events] + * @returns {Promise} + */ +function waitForTraceRequest (events) { + traceRequestEvents = events + return new Promise(captureTraceRequest) +} + +/** + * @param {import('node:http').IncomingMessage} request + * @param {import('node:http').ServerResponse} response + */ +function handleAgentRequest (request, response) { + request.resume() + + if (!request.url?.endsWith('/traces')) { + response.end('{}') + return + } + + traceRequestCount++ + if (!traceRequestResolve) { + response.end(agentResponse) + return + } + + const resolve = traceRequestResolve + traceRequestResolve = undefined + traceRequestEvents?.push('agent-received') + traceRequestEvents = undefined + resolve(response) +} + +/** + * @param {Error & { status?: number }} error + * @returns {boolean} + */ +function isAgentFailure (error) { + assert.strictEqual(error.status, 500) + return true +} describe('OTel TracerProvider', () => { it('should register with OTel API', () => { @@ -113,13 +170,203 @@ describe('OTel TracerProvider', () => { sinon.assert.calledOnce(processor.shutdown) }) - it('should delegate forceFlush to active span processor', () => { - const provider = new TracerProvider() - const processor = new NoopSpanProcessor() - provider.addSpanProcessor(processor) - processor.forceFlush = sinon.stub() + describe('forceFlush without an initialized tracer', () => { + it('rejects', async () => { + const provider = new TracerProvider() + + await assert.rejects(provider.forceFlush(), { message: 'Not started' }) + }) + }) + + describe('forceFlush with an initialized tracer', () => { + let agent + let originalRemoteConfigEnabled + + before(async () => { + originalRemoteConfigEnabled = process.env.DD_REMOTE_CONFIGURATION_ENABLED + process.env.DD_REMOTE_CONFIGURATION_ENABLED = 'false' + + agent = http.createServer(handleAgentRequest) + agent.listen(0, '127.0.0.1') + await once(agent, 'listening') + + const { port } = agent.address() + require('../../index').init({ + flushInterval: 60_000, + plugins: false, + startupLogs: false, + url: `http://127.0.0.1:${port}`, + }) + }) + + after(async () => { + if (originalRemoteConfigEnabled === undefined) { + delete process.env.DD_REMOTE_CONFIGURATION_ENABLED + } else { + process.env.DD_REMOTE_CONFIGURATION_ENABLED = originalRemoteConfigEnabled + } + + const closed = once(agent, 'close') + agent.close() + agent.closeAllConnections?.() + await closed + }) + + it('waits for Datadog delivery to complete', async () => { + const events = [] + const requestReceived = waitForTraceRequest(events) + const provider = new TracerProvider() + provider.getTracer().startSpan('otel.force_flush.delayed').end() + + const forceFlush = provider.forceFlush().then(() => events.push('forceFlush-resolved')) + const response = await requestReceived + + assert.deepStrictEqual(events, ['agent-received']) + + const responseFinished = once(response, 'finish') + response.end(agentResponse) + await responseFinished + events.push('agent-responded') + await forceFlush + + assert.deepStrictEqual(events, ['agent-received', 'agent-responded', 'forceFlush-resolved']) + }) + + it('serializes overlapping flush generations', async () => { + const provider = new TracerProvider() + const firstRequest = waitForTraceRequest() + provider.getTracer().startSpan('otel.force_flush.first_generation').end() + let firstSettled = false + const firstFlush = provider.forceFlush().finally(() => { firstSettled = true }) + const firstResponse = await firstRequest + + provider.getTracer().startSpan('otel.force_flush.second_generation').end() + let secondSettled = false + const secondFlush = provider.forceFlush().finally(() => { secondSettled = true }) + assert.strictEqual(firstSettled, false) + assert.strictEqual(secondSettled, false) + + const secondRequest = waitForTraceRequest() + firstResponse.end(agentResponse) + await firstFlush + const secondResponse = await secondRequest + assert.strictEqual(secondSettled, false) + + secondResponse.end(agentResponse) + await secondFlush + assert.strictEqual(secondSettled, true) + }) + + it('resolves without sending when the Datadog buffer is empty', async () => { + const requestsBeforeFlush = traceRequestCount + + await new TracerProvider().forceFlush() + + assert.strictEqual(traceRequestCount, requestsBeforeFlush) + }) + + it('waits for every configured span processor', async () => { + const firstSignal = new EventEmitter() + const secondSignal = new EventEmitter() + const firstDone = once(firstSignal, 'done') + const secondDone = once(secondSignal, 'done') + const first = new NoopSpanProcessor() + const second = new NoopSpanProcessor() + first.forceFlush = sinon.stub().returns(firstDone) + second.forceFlush = sinon.stub().returns(secondDone) + const provider = new TracerProvider({ spanProcessors: [first, second] }) + let settled = false + + const forceFlush = provider.forceFlush().finally(() => { settled = true }) + + sinon.assert.calledOnce(first.forceFlush) + sinon.assert.calledOnce(second.forceFlush) + firstSignal.emit('done') + await firstDone + assert.strictEqual(settled, false) + secondSignal.emit('done') + await forceFlush + assert.strictEqual(settled, true) + }) + + it('waits for delivery after a span processor fails, then rejects with the original error', async () => { + const processorError = new Error('processor failed') + const processor = new NoopSpanProcessor() + processor.forceFlush = sinon.stub().rejects(processorError) + const provider = new TracerProvider({ spanProcessors: [processor] }) + const requestReceived = waitForTraceRequest() + provider.getTracer().startSpan('otel.force_flush.processor_error').end() + let settled = false + + const forceFlush = provider.forceFlush().finally(() => { settled = true }) + await assert.rejects(processor.forceFlush.firstCall.returnValue, processorError) + assert.strictEqual(settled, false) + + const response = await requestReceived + response.end(agentResponse) + + await assert.rejects(forceFlush, error => { + assert.strictEqual(error, processorError) + return true + }) + assert.strictEqual(settled, true) + }) + + it('aggregates multiple span processor failures', async () => { + const firstError = new Error('first processor failed') + const secondError = new Error('second processor failed') + const first = new NoopSpanProcessor() + const second = new NoopSpanProcessor() + first.forceFlush = sinon.stub().rejects(firstError) + second.forceFlush = sinon.stub().rejects(secondError) + const provider = new TracerProvider({ spanProcessors: [first, second] }) + + await assert.rejects(provider.forceFlush(), error => { + assert.ok(error instanceof AggregateError) + assert.deepStrictEqual(error.errors, [firstError, secondError]) + return true + }) + }) + + it('rejects when Datadog delivery fails', async () => { + const requestReceived = waitForTraceRequest() + const provider = new TracerProvider() + provider.getTracer().startSpan('otel.force_flush.exporter_error').end() + + const forceFlush = provider.forceFlush() + const response = await requestReceived + response.statusCode = 500 + response.end('agent failed') + + await assert.rejects(forceFlush, isAgentFailure) + }) + + it('delegates to the active span processor', async () => { + const provider = new TracerProvider() + const processor = new NoopSpanProcessor() + provider.addSpanProcessor(processor) + processor.forceFlush = sinon.stub().resolves() + + await provider.forceFlush() + + sinon.assert.calledOnce(processor.forceFlush) + }) + + it('still flushes processors when the exporter has no flush method', async () => { + const datadogTracer = require('../../index')._tracer + const originalExporter = datadogTracer._exporter + datadogTracer._exporter = { export: sinon.stub() } + const processor = new NoopSpanProcessor() + processor.forceFlush = sinon.stub().resolves() + const provider = new TracerProvider({ spanProcessors: [processor] }) + + try { + await provider.forceFlush() + } finally { + datadogTracer._exporter = originalExporter + } - provider.forceFlush() - sinon.assert.calledOnce(processor.forceFlush) + sinon.assert.calledOnce(processor.forceFlush) + }) }) }) diff --git a/packages/dd-trace/test/opentelemetry/traces.spec.js b/packages/dd-trace/test/opentelemetry/traces.spec.js index d448b10c03e..9677e8d4c4c 100644 --- a/packages/dd-trace/test/opentelemetry/traces.spec.js +++ b/packages/dd-trace/test/opentelemetry/traces.spec.js @@ -1,6 +1,7 @@ 'use strict' const assert = require('assert') +const { once } = require('node:events') const http = require('node:http') const https = require('node:https') @@ -653,6 +654,38 @@ describe('OpenTelemetry Traces', () => { exporter.export([span]) }) + it('waits for an OTLP HTTP response and reports its failure', async () => { + let resolveRequest + const requestReceived = new Promise(resolve => { resolveRequest = resolve }) + const receiver = http.createServer((request, response) => { + request.resume() + resolveRequest(response) + }) + receiver.listen(0, '127.0.0.1') + await once(receiver, 'listening') + + try { + const { port } = receiver.address() + const exporter = new OtlpHttpTraceExporter(`http://127.0.0.1:${port}/v1/traces`, {}, 1000, {}) + exporter.export([createMockSpan()]) + let settled = false + const flush = new Promise((resolve, reject) => { + exporter.flush(error => error ? reject(error) : resolve(), { reportErrors: true }) + }).finally(() => { settled = true }) + const response = await requestReceived + + assert.strictEqual(settled, false) + response.statusCode = 500 + response.end('collector failed') + + await assert.rejects(flush, /HTTP 500: collector failed/) + assert.strictEqual(settled, true) + } finally { + receiver.close() + await once(receiver, 'close') + } + }) + it('sends JSON content-type header', () => { mockOtlpExport((decoded, headers) => { assert.strictEqual(headers['Content-Type'], 'application/json') diff --git a/packages/dd-trace/test/serverless.spec.js b/packages/dd-trace/test/serverless.spec.js index 270bf8e3705..0dda1094306 100644 --- a/packages/dd-trace/test/serverless.spec.js +++ b/packages/dd-trace/test/serverless.spec.js @@ -111,6 +111,56 @@ describe('TelemetryDeliveryTracker', () => { assert.strictEqual(done, 1) }) + + it('reports a synchronous delivery failure before rethrowing it', () => { + const tracker = new TelemetryDeliveryTracker() + const deliveryError = new Error('delivery failed') + let reportedError + + assert.throws(() => tracker.track( + () => { throw deliveryError }, + error => { reportedError = error } + ), deliveryError) + + assert.strictEqual(reportedError, deliveryError) + }) + + it('does not retain failures completed before the flush boundary', () => { + const tracker = new TelemetryDeliveryTracker() + let flushError + + tracker.track(complete => complete(new Error('delivery failed'))) + tracker.waitForIdle(error => { flushError = error }, { reportErrors: true }) + + assert.strictEqual(flushError, undefined) + }) + + it('aggregates failures from deliveries active at the flush boundary', () => { + const tracker = new TelemetryDeliveryTracker() + const firstError = new Error('first delivery failed') + const secondError = new Error('second delivery failed') + const complete = [] + let flushError + + tracker.track(callback => complete.push(callback)) + tracker.track(callback => complete.push(callback)) + tracker.waitForIdle(error => { flushError = error }, { reportErrors: true }) + complete[0](firstError) + complete[1](secondError) + + assert.ok(flushError instanceof AggregateError) + assert.deepStrictEqual(flushError.errors, [firstError, secondError]) + }) + + it('keeps delivery failures private unless requested', () => { + const tracker = new TelemetryDeliveryTracker() + + tracker.track(complete => complete(new Error('delivery failed'))) + + tracker.waitForIdle(error => { + assert.strictEqual(error, undefined) + }) + }) }) describe('flushServerlessTelemetry', () => {