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
63 changes: 37 additions & 26 deletions lib/transforms.js
Original file line number Diff line number Diff line change
Expand Up @@ -257,7 +257,7 @@ function traceFunction (state, node, program) {
async: node.async,
expression: false,
generator: node.generator,
}, program)
}, program, node.type === 'ArrowFunctionExpression')

node.params = wrapParams(params)

Expand Down Expand Up @@ -352,16 +352,24 @@ function traceInstanceMethod (state, node, program) {
* Builds the replacement block statement for a function body.
*
* Selects the appropriate wrapper template (`wrapSync`, `wrapPromise`, etc.) and
* prepends the shared `__apm$ctx` / `__apm$traced` preamble before returning the
* prepends the shared `__apm$traced` / `__apm$ctx` preamble before returning the
* resulting block statement body.
*
* The no-subscriber check sits between the two halves of the preamble: after
* `__apm$traced` (which only captures the call arguments it is given) and
* before `__apm$arguments` / `__apm$ctx`, so an unsubscribed call does not
* build any of the per-call transport.
*
* @param {object} state
* @param {import('estree').Node} node - The original function (or identifier for instance methods).
* @param {import('estree').Program} program
* @param {boolean} [outerIsArrow] - Whether the wrapper itself is an arrow
* function, in which case it has no own `arguments`.
* @returns {import('estree').BlockStatement['body']}
*/
function wrap (state, node, program) {
const { operator, moduleVersion } = state
function wrap (state, node, program, outerIsArrow = false) {
const { operator, moduleVersion, channelName } = state
const channelVariable = formatChannelVariable(channelName)
const { returnKind } = state.functionQuery

const iterPatch = returnKind ? generateIterPatch(state, returnKind, program) : ''
Expand All @@ -378,17 +386,18 @@ function wrap (state, node, program) {
.map((_, i) => `__apm$arg${i}`).concat('...__apm$args').join(', ')

const block = wrapper.body[0].body // Extract only block statement of function body.
const common = parse(node.type === 'ArrowFunctionExpression'
const innerIsArrow = node.type === 'ArrowFunctionExpression'
const callWrapped = innerIsArrow
? '__apm$wrapped(...__apm$callArgs)'
: '__apm$wrapped.apply(this, __apm$callArgs)'
const fastArgs = outerIsArrow ? `[${args}]` : 'arguments'
const transport = innerIsArrow
? `
const __apm$arguments = [${args}];
const __apm$ctx = {
arguments: __apm$arguments,
moduleVersion: ${JSON.stringify(moduleVersion)}
};
const __apm$traced = () => {
const __apm$wrapped = () => {};
return __apm$wrapped(...__apm$arguments);
};
`
: `
const __apm$arguments = [${args}].slice(0, arguments.length);
Expand All @@ -397,11 +406,17 @@ function wrap (state, node, program) {
self: this,
moduleVersion: ${JSON.stringify(moduleVersion)}
};
const __apm$traced = () => {
const __apm$wrapped = () => {};
return __apm$wrapped.apply(this, __apm$arguments);
};
`).body
`
const common = parse(`
function wrapper () {
const __apm$traced = (__apm$callArgs) => {
const __apm$wrapped = () => {};
return ${callWrapped};
};
if (!tr_ch_apm_hasSubscribers(${channelVariable})) return __apm$traced(${fastArgs});
${transport}
}
`).body[0].body.body

block.body.unshift(...common)

Expand Down Expand Up @@ -531,9 +546,9 @@ function wrapCallback (state, node, iterPatch = '') {
function wrapper () {
const __apm$cb = Array.prototype.at.call(__apm$arguments, ${callbackIndex});

if (!${channelVariable}.start.hasSubscribers) return __apm$traced();
if (!${channelVariable}.start.hasSubscribers) return __apm$traced(__apm$arguments);

function __apm$wrappedCb(err, res) {
const __apm$wrappedCb = function (err, res) {
if (err) {
__apm$ctx.error = err;
${channelVariable}.error.publish(__apm$ctx);
Expand All @@ -551,16 +566,16 @@ function wrapCallback (state, node, iterPatch = '') {
${channelVariable}.asyncEnd.publish(__apm$ctx);
}
});
}
};

if (typeof __apm$cb !== 'function') {
return __apm$traced();
return __apm$traced(__apm$arguments);
}
Array.prototype.splice.call(__apm$arguments, ${callbackIndex}, 1, __apm$wrappedCb);

return ${channelVariable}.start.runStores(__apm$ctx, () => {
try {
return __apm$traced();
return __apm$traced(__apm$arguments);
} catch (err) {
__apm$ctx.error = err;
${channelVariable}.error.publish(__apm$ctx);
Expand Down Expand Up @@ -597,11 +612,9 @@ function wrapPromise (state, node, iterPatch = '') {
// mutated.
return parse(`
function wrapper () {
if (!tr_ch_apm_hasSubscribers(${channelVariable})) return __apm$traced();

return ${channelVariable}.start.runStores(__apm$ctx, () => {
try {
let promise = __apm$traced();
let promise = __apm$traced(__apm$arguments);
if (typeof promise?.then !== 'function') {
__apm$ctx.result = promise;
${iterPatch}
Expand Down Expand Up @@ -682,11 +695,9 @@ function wrapSync (state, node, iterPatch = '') {
// re-throws, so the trailing `return` is unreachable.
return parse(`
function wrapper () {
if (!tr_ch_apm_hasSubscribers(${channelVariable})) return __apm$traced();

return ${channelVariable}.start.runStores(__apm$ctx, () => {
try {
__apm$ctx.result = __apm$traced();
__apm$ctx.result = __apm$traced(__apm$arguments);
${iterPatch}
} catch (err) {
__apm$ctx.error = err;
Expand Down Expand Up @@ -748,8 +759,8 @@ function generateIterPatch (state, returnKind, program) {
const __apm$orig = __apm$iter[method];
if (typeof __apm$orig !== 'function') return;
__apm$iter[method] = function () {
if (!tr_ch_apm_hasSubscribers(${iterChannelVariable})) return __apm$orig.apply(this, arguments);
const __apm$iterArgs = Array.prototype.slice.call(arguments);
if (!tr_ch_apm_hasSubscribers(${iterChannelVariable})) return __apm$orig.apply(this, __apm$iterArgs);
__apm$ctx.method = method;
__apm$ctx.arguments = __apm$iterArgs;
return ${iterChannelVariable}.${traceMethod}(__apm$orig, __apm$ctx, this, ...__apm$iterArgs);
Expand Down
30 changes: 30 additions & 0 deletions tests/arguments_mutation_kinds_cjs/mod.js
Original file line number Diff line number Diff line change
@@ -0,0 +1,30 @@
'use strict'

function syncFn (a, ...rest) {
return { self: this, a, rest, len: arguments.length }
}

async function asyncFn (a, ...rest) {
return { self: this, a, rest, len: arguments.length }
}

function callbackFn (a, cb) {
cb(null, { self: this, a, len: arguments.length })
}

class Base {
constructor (a, ...rest) {
this.a = a
this.rest = rest
}
}

class Derived extends Base {
constructor (a, ...rest) {
super(a, ...rest)
this.derivedA = a
this.derivedRest = rest
}
}

module.exports = { syncFn, asyncFn, callbackFn, Base, Derived }
80 changes: 80 additions & 0 deletions tests/arguments_mutation_kinds_cjs/test.js
Original file line number Diff line number Diff line change
@@ -0,0 +1,80 @@
'use strict'

const { syncFn, asyncFn, callbackFn, Base, Derived } = require('./instrumented.js')
const assert = require('node:assert')
const { tracingChannel } = require('node:diagnostics_channel')

// A `start` subscriber replaces the first argument and appends one in place;
// the original must see both, because the wrapper hands it the same array
// that was published as `message.arguments`.
function replaceAndAppend (name) {
tracingChannel(`orchestrion:undici:${name}`).subscribe({
start (message) {
message.arguments[0] = 'replaced'
message.arguments.push('appended')
},
})
}

for (const name of ['syncFn', 'asyncFn', 'Base_ctor', 'Derived_ctor']) replaceAndAppend(name)

;(async () => {
const self = { self: true }

assert.deepStrictEqual(syncFn.call(self, 'original', 'extra'), {
self, a: 'replaced', rest: ['extra', 'appended'], len: 3,
})

assert.deepStrictEqual(await asyncFn.call(self, 'original'), {
self, a: 'replaced', rest: ['appended'], len: 2,
})

// Constructors forward `__apm$arguments` by spreading it into the moved
// body, so parameters see the mutation. (`arguments` inside the body is the
// constructor's own and is not affected, as before.)
const base = new Base('original')
assert.strictEqual(base.a, 'replaced')
assert.deepStrictEqual(base.rest, ['appended'])

// `super()` goes through the separately instrumented `Base`, whose own
// subscriber appends once more.
const derived = new Derived('original', 'extra')
assert.strictEqual(derived.derivedA, 'replaced')
assert.deepStrictEqual(derived.derivedRest, ['extra', 'appended'])
assert.strictEqual(derived.a, 'replaced')
assert.deepStrictEqual(derived.rest, ['extra', 'appended', 'appended'])

// Callback: by the time `start` fires the wrapper has already spliced its
// own callback into `message.arguments`. A subscriber can wrap that in turn,
// and both layers must run.
{
const events = []
tracingChannel('orchestrion:undici:callbackFn').subscribe({
start (message) {
events.push('start')
message.arguments[0] = 'replaced'
const wrappedCb = message.arguments[1]
assert.strictEqual(typeof wrappedCb, 'function')
message.arguments[1] = function (err, res) {
events.push('subscriber-cb')
return wrappedCb.call(this, err, res)
}
},
asyncStart () { events.push('asyncStart') },
asyncEnd () { events.push('asyncEnd') },
end () { events.push('end') },
})

const result = await new Promise((resolve, reject) => {
callbackFn.call(self, 'original', (err, res) => {
events.push('user-cb')
err ? reject(err) : resolve(res)
})
})
assert.deepStrictEqual(result, { self, a: 'replaced', len: 2 })
assert.deepStrictEqual(events, ['start', 'subscriber-cb', 'asyncStart', 'user-cb', 'asyncEnd', 'end'])
}
})().catch((err) => {
console.error(err)
process.exit(1)
})
34 changes: 34 additions & 0 deletions tests/common/transport_spy.js
Original file line number Diff line number Diff line change
@@ -0,0 +1,34 @@
/**
* Unless explicitly stated otherwise all files in this repository are licensed under the Apache-2.0 License.
* This product includes software developed at Datadog (https://www.datadoghq.com/). Copyright 2025 Datadog, Inc.
**/
'use strict'

// The wrappers build their per-call transport with `[...].slice(...)`
// (`__apm$arguments`), `Array.prototype.at.call(...)` (callback lookup) and
// `Array.prototype.slice.call(arguments)` (iterator `next` arguments). Counting
// those calls while a single wrapped call runs synchronously tells us whether
// the wrapper took the no-subscriber fast path or built the transport.
const { slice, at } = Array.prototype

function measure (fn) {
const counts = { slice: 0, at: 0 }
Array.prototype.slice = function (...args) { // eslint-disable-line no-extend-native
counts.slice++
return slice.apply(this, args)
}
Array.prototype.at = function (...args) { // eslint-disable-line no-extend-native
counts.at++
return at.apply(this, args)
}
let value
try {
value = fn()
} finally {
Array.prototype.slice = slice // eslint-disable-line no-extend-native
Array.prototype.at = at // eslint-disable-line no-extend-native
}
return { value, slice: counts.slice, at: counts.at }
}

module.exports = { measure }
82 changes: 82 additions & 0 deletions tests/fast_path_cjs/mod.js
Original file line number Diff line number Diff line change
@@ -0,0 +1,82 @@
'use strict'

// Every original records what it was actually called with, so the test can
// compare the unsubscribed fast path against the subscribed path.
function record (self, args, named) {
return { self, len: args.length, args: Array.from(args), ...named }
}

function syncFn (a, b = 'b-default', ...rest) {
return record(this, arguments, { a, b, rest })
}

async function asyncFn (a, b = 'b-default', ...rest) {
return record(this, arguments, { a, b, rest })
}

function callbackFn (a, cb) {
const cbArg = arguments[arguments.length - 1]
cbArg(null, record(this, arguments, { a, cb: cbArg }))
}

function autoFn (a, cb) {
const last = arguments[arguments.length - 1]
if (typeof last === 'function') {
last(null, record(this, arguments, { a, cb: last }))
} else {
return Promise.resolve(record(this, arguments, { a }))
}
}

function * iterFn (a) {
const sent = yield record(this, arguments, { a })
yield sent
}

async function * asyncIterFn (a) {
const sent = yield record(this, arguments, { a })
yield sent
}

const arrowFn = (a, b = 'b-default', ...rest) => ({ a, b, rest })

class Service {
method (a, b = 'b-default', ...rest) {
return record(this, arguments, { a, b, rest })
}
}

class Base {
constructor (a, b = 'b-default', ...rest) {
this.base = record(undefined, arguments, { a, b, rest })
this.newTarget = new.target
}
}

class Derived extends Base {
constructor (a, b = 'd-default', ...rest) {
super(a, b, ...rest)
this.derived = record(undefined, arguments, { a, b, rest })
}
}

// Not declared on the class body, so it is patched onto the instance from a
// synthesised constructor at runtime.
class Holder {}
Holder.prototype.run = function (a, b = 'b-default', ...rest) {
return record(this, arguments, { a, b, rest })
}

module.exports = {
syncFn,
asyncFn,
callbackFn,
autoFn,
iterFn,
asyncIterFn,
arrowFn,
Service,
Base,
Derived,
Holder,
}
Loading