Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
146 changes: 146 additions & 0 deletions integration-tests/debugger/coordinated-sampling.spec.js
Original file line number Diff line number Diff line change
@@ -0,0 +1,146 @@
'use strict'

const assert = require('node:assert/strict')
const { setTimeout: delay } = require('node:timers/promises')

const { setup } = require('./utils')

/**
* @typedef {object} ReceivedSnapshot
* @property {string} probeId
* @property {string | undefined} traceId
*/

// How long to keep listening for snapshots that should not arrive, once the expected ones have
const GRACE_PERIOD_MS = 500

describe('Dynamic Instrumentation', function () {
const t = setup({ dependencies: ['fastify'] })

describe('coordinated sampling', function () {
it('should emit the snapshots of all probes hit in a trace, or none of them', async function () {
const [first, second, third] = t.breakpoints
const probes = {
// Sampled at most once per 10 seconds, so the second trace it starts is dropped
first: first.generateRemoteConfig({ captureSnapshot: true, sampling: { snapshotsPerSecond: 0.1 } }),
// Hit just before the first trace, so it would be rate limited in that trace if sampled on its own
second: second.generateRemoteConfig({ captureSnapshot: true, sampling: { snapshotsPerSecond: 1 } }),
// Practically not rate limited, so it would emit in the dropped trace if sampled on its own
third: third.generateRemoteConfig({ captureSnapshot: true, sampling: { snapshotsPerSecond: 1000 } }),
}
await installProbes(Object.values(probes))
const { snapshots, waitForCount } = collectSnapshots()

const { body: { traceId: secondTraceId } } = await t.request('/second')
const { body: { traceId: sampledTraceId } } = await t.request('/chain')
const { body: { traceId: droppedTraceId } } = await t.request('/chain')
await waitForCount(4)
await delay(GRACE_PERIOD_MS)

assert.notStrictEqual(sampledTraceId, droppedTraceId)
assert.deepStrictEqual(groupProbeNamesByTrace(snapshots, probes), new Map([
[secondTraceId, ['second']],
[sampledTraceId, ['first', 'second', 'third']],
]))
})

it('should emit one snapshot per probe per trace', async function () {
const [,,,, loopBody, afterLoop] = t.breakpoints
const probes = {
loopBody: loopBody.generateRemoteConfig({ captureSnapshot: true, sampling: { snapshotsPerSecond: 1000 } }),
afterLoop: afterLoop.generateRemoteConfig({ captureSnapshot: true, sampling: { snapshotsPerSecond: 1000 } }),
}
await installProbes(Object.values(probes))
const { snapshots, waitForCount } = collectSnapshots()

const { body: { traceId } } = await t.request('/loop')
await waitForCount(2)
await delay(GRACE_PERIOD_MS)

assert.deepStrictEqual(groupProbeNamesByTrace(snapshots, probes), new Map([
[traceId, ['afterLoop', 'loopBody']],
]))
})

it('should sample probes independently when hit in a trace that has finished', async function () {
const [,,, leaked] = t.breakpoints
// Without the fallback to independent sampling, the probe would emit only once in the finished trace
const probe = leaked.generateRemoteConfig({ captureSnapshot: true, sampling: { snapshotsPerSecond: 10 } })
await installProbes([probe])
const { snapshots, waitForCount } = collectSnapshots()

const { body: { traceId } } = await t.request('/leak')
await waitForCount(2)

for (const snapshot of snapshots) {
assert.strictEqual(snapshot.probeId, probe.config.id)
assert.strictEqual(snapshot.traceId, traceId)
}
})
})

/**
* Install probes and wait until all of them are installed.
*
* @param {Array<{ product: string, id: string, config: { id: string } }>} rcConfigs - The remote configs of the
* probes to install.
*/
async function installProbes (rcConfigs) {
const probesInstalled = t.waitForProbeStatus(rcConfigs.map(({ config }) => config.id), 'INSTALLED')
for (const rcConfig of rcConfigs) {
t.agent.addRemoteConfig(rcConfig)
}
await probesInstalled
}

/**
* Collect the snapshots received by the agent from now on.
*
* @returns {{ snapshots: ReceivedSnapshot[], waitForCount: (count: number) => Promise<void> }}
*/
function collectSnapshots () {
/** @type {ReceivedSnapshot[]} */
const snapshots = []
/** @type {Array<{ count: number, resolve: () => void }>} */
const waiters = []

t.agent.on('debugger-input', ({ payload }) => {
for (const { dd, debugger: { snapshot } } of payload) {
snapshots.push({ probeId: snapshot.probe.id, traceId: dd?.trace_id })
}
for (const waiter of waiters) {
if (snapshots.length >= waiter.count) waiter.resolve()
}
})

return {
snapshots,
waitForCount (count) {
return new Promise((resolve) => {
if (snapshots.length >= count) return resolve()
waiters.push({ count, resolve })
})
},
}
}
})

/**
* Group the names of the probes that emitted snapshots by the trace they were emitted in.
*
* @param {ReceivedSnapshot[]} snapshots - The received snapshots.
* @param {Record<string, { config: { id: string } }>} probes - The remote configs of the probes, by name.
* @returns {Map<string | undefined, string[]>} The sorted probe names, by trace id.
*/
function groupProbeNamesByTrace (snapshots, probes) {
const nameById = new Map(Object.entries(probes).map(([name, { config }]) => [config.id, name]))
/** @type {Map<string | undefined, string[]>} */
const namesByTrace = new Map()
for (const { probeId, traceId } of snapshots) {
const names = namesByTrace.get(traceId) ?? []
names.push(/** @type {string} */ (nameById.get(probeId)))
namesByTrace.set(traceId, names)
}
for (const names of namesByTrace.values()) names.sort()
return namesByTrace
}
70 changes: 70 additions & 0 deletions integration-tests/debugger/target-app/coordinated-sampling.js
Original file line number Diff line number Diff line change
@@ -0,0 +1,70 @@
'use strict'

// @ts-expect-error This code is running in a sandbox where dd-trace is available
require('dd-trace/init')
// @ts-expect-error This code is running in a sandbox where dd-trace is available
const tracer = require('dd-trace')
const { setTimeout: sleep } = require('node:timers/promises')

// @ts-expect-error This code is running in a sandbox where fastify is available
const Fastify = require('fastify')

const fastify = Fastify({ logger: { level: 'error' } })

function first (value) {
return value // BREAKPOINT: /chain
}

function second (value) {
return value // BREAKPOINT: /second
}

function third (value) {
return value // BREAKPOINT: /chain
}

function leaked () {
return 'leaked' // BREAKPOINT: /leak
}

fastify.get('/chain', async function chainHandler () {
first(1)
await sleep(10)
second(2)
await sleep(10)
third(3)
return { traceId: getActiveTraceId() }
})

fastify.get('/second', function secondHandler () {
second(2)
return { traceId: getActiveTraceId() }
})

fastify.get('/loop', async function loopHandler () {
let total = 0
for (let i = 0; i < 3; i++) {
total += i // BREAKPOINT: /loop
await sleep(10)
}
return { traceId: getActiveTraceId(), total } // BREAKPOINT: /loop
})

fastify.get('/leak', function leakHandler () {
// A timer created while handling the request runs its callbacks in the async context of the request, also after the
// request and its trace have finished.
setInterval(leaked, 20)
return { traceId: getActiveTraceId() }
})

function getActiveTraceId () {
return tracer.scope().active()?.context().toTraceId()
}

fastify.listen({ port: process.env.APP_PORT || 0 }, (err) => {
if (err) {
fastify.log.error(err)
process.exit(1)
}
process.send?.({ port: fastify.server.address().port })
})
119 changes: 108 additions & 11 deletions packages/dd-trace/src/debugger/probe_sampler.js
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@

const { types } = require('node:util')

const { storage } = require('../../../datadog-core')
const { MAX_SNAPSHOTS_PER_SECOND_GLOBALLY } = require('./devtools_client/defaults')
const { MAX_MESSAGE_LENGTH } = require('./constants')
const { EVENT_TYPE, SKIPPED_REASON } = require('./guardrail-metrics')
Expand Down Expand Up @@ -32,6 +33,20 @@ const ddTraceGlobal = /** @type {Record<symbol, SharedArrayBuffer | object | und
* @property {{ DD_DYNAMIC_INSTRUMENTATION_EVALUATION_TIMEOUT_MS: number }} dynamicInstrumentation
*/

/**
* The trace object shared by all spans of a trace in this process (`DatadogSpanContext#_trace`).
*
* @typedef {{ started: object[] }} LocalTrace
*/

/**
* The span active in the current async context, as stored in the legacy storage by the tracer.
*
* @typedef {{ context: () => { _trace: LocalTrace } }} ActiveSpan
*/

const legacyStorage = storage('legacy')

let evaluationTimeoutNs = 0n

module.exports = {
Expand Down Expand Up @@ -63,6 +78,14 @@ function installProbeSampler (guardrailMetrics, config) {
const buffer = createProbeSamplerBuffer()

const lastCaptureNsByProbeId = new Map()
/**
* The coordinated sampling decision of each trace in which a snapshot-producing probe was hit: `false` if the trace
* was dropped, otherwise the sampling indexes of the probes that already emitted a snapshot in it. Keyed weakly, so
* a decision is released together with its trace.
*
* @type {WeakMap<LocalTrace, Set<number> | false>}
*/
const emittedProbesByTrace = new WeakMap()
/**
* Probes that are skipped at entry because a recent evaluation failed or exceeded its time budget. One error result
* is reported per throttle window, so the probe stays visible without repeatedly paying for the evaluation.
Expand Down Expand Up @@ -219,7 +242,10 @@ function installProbeSampler (guardrailMetrics, config) {
}

/**
* Apply the per-probe and global rate limits and store the sampled probe index for the debugger worker.
* Decide if a probe hit should be sampled and store the sampled probe index for the debugger worker.
*
* Snapshot-producing probes hit within an active trace share one sampling decision per trace, so a trace emits the
* snapshots of all of its probes or none of them. Other hits are sampled independently.
*
* @param {number} probeIndex - The worker-side probe sampling index.
* @param {string} probeId - The probe id.
Expand All @@ -228,29 +254,100 @@ function installProbeSampler (guardrailMetrics, config) {
* @param {boolean} isSnapshotProducingProbe - Whether this probe counts toward the global snapshot sample limit.
*/
function sample (probeIndex, probeId, now, nsBetweenSampling, isSnapshotProducingProbe) {
if (isSnapshotProducingProbe === true) {
const span = /** @type {ActiveSpan | null | undefined} */ (legacyStorage.getStore()?.span)
const trace = span?.context()._trace
// The span processor empties `started` once all spans of a trace have finished. A probe hit under such a trace
// runs in a context that outlived it, e.g. a connection pool or event emitter callback bound to a request that
// has since ended. Coordinating that hit would let a decision made for the old request silence the probe, or
// limit it to a single snapshot, for as long as the context lives, so it's sampled independently instead.
if (trace !== undefined && trace.started.length !== 0) {
return sampleInTrace(trace, probeIndex, probeId, now, nsBetweenSampling)
}
}

if (isRateLimited(probeId, now, nsBetweenSampling, isSnapshotProducingProbe)) return false
return commitSample(probeIndex, probeId, now, isSnapshotProducingProbe)
}

/**
* Sample a snapshot-producing probe hit using the sampling decision of its trace. The first such probe hit in the
* trace makes the decision for all of them, subject to its per-probe rate limit and the global snapshot rate limit.
* If the trace is sampled, each probe then emits once in it without being subject to the rate limits, so the trace
* emits a complete set of snapshots. Those snapshots still count toward the rate limits, which keeps the overall
* volume close to the limits by making later traces less likely to be sampled.
*
* @param {LocalTrace} trace - The trace the probe was hit in.
* @param {number} probeIndex - The worker-side probe sampling index.
* @param {string} probeId - The probe id.
* @param {bigint} now - The current time.
* @param {bigint} nsBetweenSampling - Minimum nanoseconds between samples for this probe.
*/
function sampleInTrace (trace, probeIndex, probeId, now, nsBetweenSampling) {
const emittedProbes = emittedProbesByTrace.get(trace)

if (emittedProbes === undefined) {
if (isRateLimited(probeId, now, nsBetweenSampling, true)) {
emittedProbesByTrace.set(trace, false)
return false
}
// If the sampled probe can't be handed over, no decision is made, so the next probe hit in the trace makes it
if (!commitSample(probeIndex, probeId, now, true)) return false
emittedProbesByTrace.set(trace, new Set([probeIndex]))
return true
}

if (emittedProbes === false || emittedProbes.has(probeIndex)) {
guardrailMetrics.eventSkipped(SKIPPED_REASON.RATE_LIMIT_PROBE, EVENT_TYPE.SNAPSHOT)
return false
}

if (!commitSample(probeIndex, probeId, now, true)) return false
emittedProbes.add(probeIndex)
return true
}

/**
* Check the per-probe and global rate limits, recording the skip if a limit is reached.
*
* @param {string} probeId - The probe id.
* @param {bigint} now - The current time.
* @param {bigint} nsBetweenSampling - Minimum nanoseconds between samples for this probe.
* @param {boolean} isSnapshotProducingProbe - Whether this probe is subject to the global snapshot sample limit.
*/
function isRateLimited (probeId, now, nsBetweenSampling, isSnapshotProducingProbe) {
const lastCaptureNs = lastCaptureNsByProbeId.get(probeId)
if (lastCaptureNs !== undefined && now - lastCaptureNs < nsBetweenSampling) {
guardrailMetrics.eventSkipped(
SKIPPED_REASON.RATE_LIMIT_PROBE,
isSnapshotProducingProbe === true ? EVENT_TYPE.SNAPSHOT : EVENT_TYPE.LOG
)
return false
return true
}

let shouldResetGlobalSnapshotRateWindow = false
if (isSnapshotProducingProbe === true) {
if (now - globalSnapshotSamplingRateWindowStart > oneSecondNs) {
shouldResetGlobalSnapshotRateWindow = true
} else if (snapshotsSampledWithinTheLastSecond >= MAX_SNAPSHOTS_PER_SECOND_GLOBALLY) {
guardrailMetrics.eventSkipped(SKIPPED_REASON.RATE_LIMIT_GLOBAL, EVENT_TYPE.SNAPSHOT)
return false
}
if (isSnapshotProducingProbe === true &&
now - globalSnapshotSamplingRateWindowStart <= oneSecondNs &&
snapshotsSampledWithinTheLastSecond >= MAX_SNAPSHOTS_PER_SECOND_GLOBALLY) {
guardrailMetrics.eventSkipped(SKIPPED_REASON.RATE_LIMIT_GLOBAL, EVENT_TYPE.SNAPSHOT)
return true
}

return false
}

/**
* Store the sampled probe index for the debugger worker and count the sample toward the rate limits.
*
* @param {number} probeIndex - The worker-side probe sampling index.
* @param {string} probeId - The probe id.
* @param {bigint} now - The current time.
* @param {boolean} isSnapshotProducingProbe - Whether this probe counts toward the global snapshot sample limit.
*/
function commitSample (probeIndex, probeId, now, isSnapshotProducingProbe) {
if (!storeSampledProbeIndex(probeIndex)) return false

if (isSnapshotProducingProbe === true) {
if (shouldResetGlobalSnapshotRateWindow === true) {
if (now - globalSnapshotSamplingRateWindowStart > oneSecondNs) {
snapshotsSampledWithinTheLastSecond = 1
globalSnapshotSamplingRateWindowStart = now
} else {
Expand Down
Loading
Loading