Repository navigation
perf: check for subscribers before building the per-call context in wrappers #102
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
base: main
Are you sure you want to change the base?
Changes from all commits
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -266,7 +266,7 @@ function traceFunction (state, node, program) { | |
| async: node.async, | ||
| expression: false, | ||
| generator: node.generator, | ||
| }, program) | ||
| }, program, node.type === 'ArrowFunctionExpression') | ||
|
|
||
| node.params = wrapParams(params) | ||
|
|
||
|
|
@@ -368,16 +368,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 takes the call arguments as its parameter) 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, selfBinding } = state | ||
| function wrap (state, node, program, outerIsArrow = false) { | ||
| const { operator, moduleVersion, channelName, selfBinding } = state | ||
| const channelVariable = formatChannelVariable(channelName) | ||
| const { returnKind } = state.functionQuery | ||
|
|
||
| const iterPatch = returnKind ? generateIterPatch(state, returnKind, program) : '' | ||
|
|
@@ -396,18 +404,21 @@ function wrap (state, node, program) { | |
| const block = wrapper.body[0].body // Extract only block statement of function body. | ||
| // Declared in the outer function scope so that the `super()` call sites in | ||
| // the moved body can write to it and the `finally` block can read it back. | ||
| // It sits above the subscriber check because the fast path runs `super()` | ||
| // too, and assigning a `let` before its declaration runs throws. | ||
| const declareSelf = selfBinding ? `let ${selfBinding};` : '' | ||
| const common = parse(declareSelf + (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); | ||
|
|
@@ -416,11 +427,18 @@ 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 () { | ||
| ${declareSelf} | ||
| const __apm$traced = (__apm$callArgs) => { | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Hm, Custom transforms are public ( I checked the two downstream users that GitHub code search finds:
Suggestion: say in the PR body (and make sure to keep in the squash commit and changelog) that The
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Ah, also, the fast path still creates Every unsubscribed call still allocates the I don't think there needs to be any change for this. But in the follow-up, include the
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Couldn't we add backwards compatibility doing something like the following: Note: in this case buildArgumentArray is just a shorthand for the existing arguments copying logic and it's just meant for showcasing the above cleaner. This would keep __apm$traced() working for custom transforms, reuse subscriber-modified arguments, and still skip array construction on the normal fast path.
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Yeah, I think that would work. I think if we don't require any changes in the downstream consumers, then the contract has only expanded, not broken. 👍
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. I tried implementing the above proposal, but it gave back about 57% of the sync fast-path performance gain (it even made arrows and derived constructors slower for some reason). It looks like adding the default parameter interferes with some V8 optimizations, even when arguments are passed explicitly and the default never runs. I'll look into other possible ways we can keep both the performance gain and the backwards compatibility. Maybe we could keep the legacy wrapper when custom transforms are present and use the fast path otherwise.
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Update: I implemented the above locally, it involves retaining the legacy wrapper shape for files using custom transforms and the fast wrapper otherwise. Benchmarks retained the same gain for the fast path, while custom transforms stayed on par with the current state. What do you think, should we move forward with this as a temporary backward compatible solution? We could also look into version gating it and removing the legacy wrappers in the next major release. |
||
| const __apm$wrapped = () => {}; | ||
| return ${callWrapped}; | ||
| }; | ||
| if (!tr_ch_apm_hasSubscribers(${channelVariable})) return __apm$traced(${fastArgs}); | ||
| ${transport} | ||
| } | ||
| `).body[0].body.body | ||
|
|
||
| block.body.unshift(...common) | ||
|
|
||
|
|
@@ -619,9 +637,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); | ||
|
|
@@ -639,16 +657,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); | ||
|
|
@@ -685,11 +703,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} | ||
|
|
@@ -770,11 +786,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; | ||
|
|
@@ -836,8 +850,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); | ||
|
|
||
| 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 } |
| 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) | ||
| }) |
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,41 @@ | ||
| /** | ||
| * 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. | ||
| // | ||
| // Replacing the prototype methods is process-wide, so keep the window small: | ||
| // they are only swapped for one synchronous call and always restored, and each | ||
| // fixture runs in its own `node` process. Callers should only assert zero | ||
| // counts on the fast path and a non-zero count on the subscribed path, never | ||
| // exact numbers. The "fast path ordering" test in tests.test.mjs checks the | ||
| // same property on the generated code without patching anything. | ||
| 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 | ||
|
isaacs marked this conversation as resolved.
|
||
| 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 } | ||
| 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, | ||
| } |
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
The arrow fast path still copies the arguments into a new array, then spreads it.
Line 393 sets the fast-path argument to the array literal
[${args}], because an arrow has no ownarguments. Line 391 makes__apm$tracedcall an arrow original with__apm$wrapped(...__apm$callArgs). The generated arrow fast path is then:So each unsubscribed arrow call reads the rest array
__apm$args, builds a new array from it, and spreads that array again into the call. A non-arrow wrapper passesargumentsand calls.apply, which does not build a new array. (The constructor case also uses line 391, but its outer function is not an arrow, so it getsargumentsfrom line 393 and only does the one spread.) Arrow expressions are in the tests (observable: false, so allocation is not checked) but not in the benchmark table. The PR body is honest that arrows use[params], but the reader can not tell how much arrows gain.Suggestion: no code change needed here, but it would be good to add an arrow row to the benchmark table, or say that arrows gain less. A cheaper arrow fast path (for example, calling
__apm$wrappeddirectly) needs__apm$wrappedhoisted out of__apm$traced, so it belongs with the "hoist the original" follow-up, not this PR.