From b3d59995a8316cc20ca850701e8203b31f3b0fd4 Mon Sep 17 00:00:00 2001 From: Bryan English Date: Wed, 11 Feb 2026 16:08:28 -0500 Subject: [PATCH] lib: add built-in OpenTelemetry tracing Adds experimental OpenTelemetry tracing support, as described in the included docs. Tracing is activated via environment variables only and requires the --experimental-otel flag. There is no programmatic API yet. Custom instrumentation is not supported. Assisted-by: pi:glm-5.3 Signed-off-by: Bryan English --- benchmark/otel/http.js | 75 ++++ doc/api/cli.md | 85 +++++ doc/api/index.md | 1 + doc/api/otel.md | 187 ++++++++++ doc/node.1 | 45 +++ lib/internal/otel/core.js | 121 ++++++ lib/internal/otel/flush.js | 347 ++++++++++++++++++ lib/internal/otel/id.js | 56 +++ lib/internal/otel/instrumentations.js | 263 +++++++++++++ lib/internal/otel/span.js | 208 +++++++++++ lib/internal/process/pre_execution.js | 57 +++ src/node_options.cc | 4 + src/node_options.h | 5 + test/common/otel.js | 116 ++++++ test/doctool/test-doc-api-json.mjs | 1 + test/parallel/test-otel-beforeexit-flush.js | 61 +++ .../parallel/test-otel-context-propagation.js | 104 ++++++ test/parallel/test-otel-env.js | 283 ++++++++++++++ test/parallel/test-otel-exporter.js | 90 +++++ test/parallel/test-otel-filter.js | 43 +++ test/parallel/test-otel-flush-coverage.js | 118 ++++++ test/parallel/test-otel-http-client.js | 120 ++++++ test/parallel/test-otel-http-server.js | 52 +++ test/parallel/test-otel-https.js | 64 ++++ test/parallel/test-otel-id-refill.js | 30 ++ test/parallel/test-otel-otlp-compliance.js | 283 ++++++++++++++ test/parallel/test-otel-self-trace.js | 79 ++++ test/parallel/test-otel-server-close.js | 65 ++++ test/parallel/test-otel-span-coverage.js | 81 ++++ .../test-otel-traceparent-validation.js | 112 ++++++ test/parallel/test-otel-undici.js | 80 ++++ 31 files changed, 3236 insertions(+) create mode 100644 benchmark/otel/http.js create mode 100644 doc/api/otel.md create mode 100644 lib/internal/otel/core.js create mode 100644 lib/internal/otel/flush.js create mode 100644 lib/internal/otel/id.js create mode 100644 lib/internal/otel/instrumentations.js create mode 100644 lib/internal/otel/span.js create mode 100644 test/common/otel.js create mode 100644 test/parallel/test-otel-beforeexit-flush.js create mode 100644 test/parallel/test-otel-context-propagation.js create mode 100644 test/parallel/test-otel-env.js create mode 100644 test/parallel/test-otel-exporter.js create mode 100644 test/parallel/test-otel-filter.js create mode 100644 test/parallel/test-otel-flush-coverage.js create mode 100644 test/parallel/test-otel-http-client.js create mode 100644 test/parallel/test-otel-http-server.js create mode 100644 test/parallel/test-otel-https.js create mode 100644 test/parallel/test-otel-id-refill.js create mode 100644 test/parallel/test-otel-otlp-compliance.js create mode 100644 test/parallel/test-otel-self-trace.js create mode 100644 test/parallel/test-otel-server-close.js create mode 100644 test/parallel/test-otel-span-coverage.js create mode 100644 test/parallel/test-otel-traceparent-validation.js create mode 100644 test/parallel/test-otel-undici.js diff --git a/benchmark/otel/http.js b/benchmark/otel/http.js new file mode 100644 index 000000000000..4763bc0e688e --- /dev/null +++ b/benchmark/otel/http.js @@ -0,0 +1,75 @@ +'use strict'; + +// Measures the per-request overhead of the built-in OpenTelemetry tracing +// subsystem on the http hot path. Requests are issued sequentially over a +// single keep-alive connection so that the difference between the traced +// and untraced configurations reflects per-request work (two spans per +// request: one SERVER, one CLIENT), not scheduling or connection noise. +// +// When tracing is enabled, spans are exported to an unreachable endpoint: +// the measurement covers span creation and serialization but not collector +// latency. Export failures are reported as throttled warnings. +// +// Run with: node benchmark/otel/http.js + +const common = require('../common'); + +const bench = common.createBenchmark(main, { + tracing: [0, 1], + n: [1e5], +}, { + flags: ['--expose-internals'], +}); + +function main({ tracing, n }) { + if (tracing) { + const otel = require('internal/otel/core'); + otel.start({ endpoint: 'http://127.0.0.1:1' }); + } + + const http = require('http'); + + const server = http.createServer((req, res) => { + res.writeHead(200); + res.end('ok'); + }); + + server.listen(0, '127.0.0.1', () => { + const port = server.address().port; + const agent = new http.Agent({ keepAlive: true, maxSockets: 1 }); + + let completed = 0; + let measured = false; + const kWarmup = 100; + + function request(done) { + http.get({ host: '127.0.0.1', port, agent, path: '/bench' }, (res) => { + res.resume(); + res.on('end', done); + }); + } + + request(function done() { + if (!measured) { + completed++; + if (completed < kWarmup) { + request(done); + return; + } + // Warmup is done; the measured requests start now. + measured = true; + completed = 0; + bench.start(); + } else { + completed++; + } + if (completed >= n) { + bench.end(n); + server.close(); + agent.destroy(); + return; + } + request(done); + }); + }); +} diff --git a/doc/api/cli.md b/doc/api/cli.md index 3822685d3d50..42910ebae3e6 100644 --- a/doc/api/cli.md +++ b/doc/api/cli.md @@ -1545,6 +1545,19 @@ added: Enable experimental support for the network inspection with Chrome DevTools. +### `--experimental-otel` + + + +> Stability: 1 - Experimental + +Enable the experimental built-in OpenTelemetry tracing subsystem. When +enabled, tracing is activated by setting the [`NODE_OTEL`][] or +[`NODE_OTEL_ENDPOINT`][] environment variables. See the [OpenTelemetry][] +documentation for details. + ### `--experimental-package-map=` +### `NODE_OTEL=value` + + + +> Stability: 1 - Experimental + +When set to `1` while the [`--experimental-otel`][] flag is +enabled, activates the built-in OpenTelemetry tracing subsystem using the +default collector endpoint (`http://localhost:4318`). Values other than +`1` are ignored. If `NODE_OTEL_ENDPOINT` is also set, it takes precedence +for the endpoint. See the [OpenTelemetry][] documentation for details. + +### `NODE_OTEL_ENDPOINT=url` + + + +> Stability: 1 - Experimental + +When set to a non-empty value while the [`--experimental-otel`][] flag is +enabled, activates the built-in OpenTelemetry tracing subsystem and directs +spans to the specified OTLP/HTTP collector endpoint. The endpoint must be +the base URL of the collector: any path present in the endpoint is +replaced with `/v1/traces`, and an endpoint that already ends with +`/v1/traces` is used as is. When only `NODE_OTEL=1` is set, the default +collector endpoint (`http://localhost:4318`) is used. If `NODE_OTEL` is +also set, `NODE_OTEL_ENDPOINT` takes precedence for the endpoint. See the +[OpenTelemetry][] documentation for details. + +### `NODE_OTEL_FILTER=module[,…]` + + + +> Stability: 1 - Experimental + +Comma-separated list of core modules to instrument when OpenTelemetry tracing +is active. When not set, all supported modules are instrumented. Supported +values: `node:http`, `node:undici`, `node:fetch`. See the [OpenTelemetry][] +documentation for details. + +### `NODE_OTEL_FLUSH_INTERVAL=milliseconds` + + + +> Stability: 1 - Experimental + +Interval in milliseconds between periodic flushes of buffered spans to the +collector. Must be a positive integer. **Default:** `10000`. + +### `NODE_OTEL_MAX_BUFFER_SIZE=number` + + + +> Stability: 1 - Experimental + +Maximum number of spans buffered in memory before an immediate flush to the +collector is triggered. Must be a positive integer. **Default:** `100`. + ### `NODE_PATH=path[:…]` + + + + + +> Stability: 1 - Experimental + +Node.js includes an experimental built-in [OpenTelemetry][] tracing subsystem. +When activated, the subsystem automatically creates spans for HTTP server and +client operations and exports them using the [OTLP/HTTP JSON][] protocol. + +The subsystem is experimental and must be enabled with the +`--experimental-otel` flag. It is activated only via environment +variables, listed below. There is currently no programmatic API for +activating or configuring the subsystem, and no support for custom +instrumentations; these may be added in the future. + +## Limitations + +The subsystem is independent of the OpenTelemetry JavaScript packages. It +is not interoperable with `@opentelemetry/api`, and spans created by one +are not visible to the other. Running both in the same process produces +duplicate trace and span IDs. The built-in instrumentation also +overwrites the `traceparent` (and, when present, `tracestate`) header on +outgoing HTTP requests, discarding whatever a userland propagator may +have set. Users should run either the built-in subsystem or a userland +OpenTelemetry SDK, not both. + +The subsystem runs in the main thread only. Worker threads +(`node:worker_threads`) are not traced, even though they inherit the +environment variables. + +Buffered spans are exported periodically, when the internal buffer fills, +and when the event loop is about to drain. Spans that are still buffered +when `process.exit()` is called explicitly are lost, because explicit +exits do not run the exit-time flush. + +## Environment variables + +### `NODE_OTEL` + +When set to `1`, activates the tracing subsystem using the default +collector endpoint (`http://localhost:4318`). Values other than `1` are +ignored. If `NODE_OTEL_ENDPOINT` is also set, it takes precedence for the +endpoint. + +```bash +node --experimental-otel app.js # NODE_OTEL=1 set in the environment +``` + +### `NODE_OTEL_ENDPOINT` + +When set to a non-empty value, activates the tracing subsystem and directs +spans to the specified OTLP collector endpoint. The endpoint should be the +base URL of an OTLP/HTTP collector (e.g. `http://localhost:4318`) without +a path: any path present in the endpoint is replaced with `/v1/traces`, +and an endpoint that already ends with `/v1/traces` is used as is. When +only `NODE_OTEL=1` is set, the default collector endpoint +(`http://localhost:4318`) is used. + +```bash +NODE_OTEL_ENDPOINT=http://collector.example.com:4318 \ + node --experimental-otel app.js +``` + +### `NODE_OTEL_FILTER` + +Accepts a comma-separated list of core modules to instrument. When not set, all +supported modules are instrumented. For example, setting +`NODE_OTEL_FILTER=node:http` would enable tracing only for the `node:http` +module. + +Supported module filter values: + +* `node:http` — HTTP server and client operations +* `node:undici` — Undici HTTP client operations +* `node:fetch` — Fetch API operations (alias for undici) + +### `NODE_OTEL_MAX_BUFFER_SIZE` + +Maximum number of spans buffered in memory before an immediate flush to the +collector is triggered. Must be a positive integer. **Default:** `100`. + +### `NODE_OTEL_FLUSH_INTERVAL` + +Interval in milliseconds between periodic flushes of buffered spans to the +collector. Must be a positive integer. **Default:** `10000`. + +### `OTEL_SERVICE_NAME` + +Standard OpenTelemetry environment variable used to set the service name in +exported resource attributes. Defaults to `unknown_service:node`, per the +OpenTelemetry [semantic conventions][] for low-cardinality service names. + +## Instrumented operations + +When the subsystem is active, spans are automatically created for the +following operations. Per the OpenTelemetry [semantic conventions][], span +names must be low cardinality; Node.js core has no route concept, so spans +are named `{method}` and the request details live in the attributes. + +### HTTP server + +A span with kind `SERVER` is created for each incoming HTTP request. The span +starts when the request is received and ends when the response finishes. If the +client disconnects before the response completes, the span ends with an error +status. + +Server spans receive error status (`STATUS_ERROR`) for 5xx response codes. 4xx +responses are not treated as server errors per OpenTelemetry semantic +conventions. + +Attributes set on server spans: + +| Attribute | Description | Condition | +| --------------------------- | -------------------------------- | ----------------------------- | +| `http.request.method` | HTTP method (e.g. `GET`, `POST`) | Always | +| `url.path` | Request URL path (without query) | Always | +| `url.query` | Query string (without `?`) | When query string is present | +| `url.scheme` | `http` or `https` | Always | +| `server.address` | Host header value | When `Host` header is present | +| `network.protocol.version` | HTTP version (e.g. `1.1`) | Always | +| `http.response.status_code` | Response status code | When response finishes | +| `error.type` | HTTP status code as string | On 5xx responses | + +### HTTP client + +A span with kind `CLIENT` is created for each outgoing HTTP request made via +`node:http`. The span starts when the request is created and ends when the +response body completes or an error occurs. + +Client spans receive error status (`STATUS_ERROR`) for 4xx and 5xx response +codes. On connection errors, an `exception` event is added to the span with +`exception.type`, `exception.message`, and `exception.stacktrace` attributes. + +Attributes set on client spans: + +| Attribute | Description | Condition | +| --------------------------- | ------------------------- | ------------------------- | +| `http.request.method` | HTTP method | Always | +| `url.full` | Full request URL | Always | +| `server.address` | Target host | Always | +| `server.port` | Target port | When available | +| `http.response.status_code` | Response status code | When response is received | +| `network.protocol.version` | HTTP version | When response is received | +| `error.type` | Status code or error name | On 4xx/5xx or errors | + +### Undici/Fetch client + +A span with kind `CLIENT` is created for each outgoing request made via +`fetch()` or undici's `request()`. Error status, `exception` event behavior, +and span end timing are the same as for HTTP client spans above. + +Attributes set on undici/fetch client spans: + +| Attribute | Description | Condition | +| --------------------------- | ------------------------- | ------------------------- | +| `http.request.method` | HTTP method | Always | +| `url.full` | Full request URL | Always | +| `server.address` | Target origin | Always | +| `http.response.status_code` | Response status code | When response is received | +| `error.type` | Status code or error name | On 4xx/5xx or errors | + +## W3C Trace Context propagation + +The tracing subsystem automatically propagates [W3C Trace Context][] across HTTP +boundaries: + +* **Incoming requests**: The `traceparent` header is read from incoming HTTP + requests, and child spans created during request processing inherit the + trace ID. `traceparent` values using a version other than `00` are + rejected, and a fresh trace is started instead. The `tracestate` header, + when present, is attached to the span and forwarded verbatim to outgoing + requests. It is never parsed or modified. +* **Outgoing requests**: The `traceparent` header is injected into outgoing + HTTP and undici/fetch requests, together with the stored `tracestate` + when one is present, enabling distributed tracing across services. + +[OTLP/HTTP JSON]: https://opentelemetry.io/docs/specs/otlp/#otlphttp +[W3C Trace Context]: https://www.w3.org/TR/trace-context/ +[OpenTelemetry]: https://opentelemetry.io/ +[semantic conventions]: https://opentelemetry.io/docs/specs/semconv/http/http-spans/ diff --git a/doc/node.1 b/doc/node.1 index 863c468eed10..b607fa139db6 100644 --- a/doc/node.1 +++ b/doc/node.1 @@ -875,6 +875,12 @@ This feature requires \fB--allow-worker\fR if used with the Permission Model. .It Fl -experimental-network-inspection Enable experimental support for the network inspection with Chrome DevTools. . +.It Fl -experimental-otel +Enable the experimental built-in OpenTelemetry tracing subsystem. When +enabled, tracing is activated by setting the \fBNODE_OTEL\fR or +\fBNODE_OTEL_ENDPOINT\fR environment variables. See the OpenTelemetry +documentation for details. +. .It Fl -experimental-package-map Ns = Ns Ar Enable experimental package map resolution. The \fBpath\fR argument specifies the location of a JSON configuration file that defines package resolution mappings. @@ -2243,6 +2249,8 @@ one is included in the list below. .It \fB--experimental-modules\fR .It +\fB--experimental-otel\fR +.It \fB--experimental-package-map\fR .It \fB--experimental-print-required-tla\fR @@ -2544,6 +2552,43 @@ V8 options that are allowed are: \fB--perf-prof-unwinding-info\fR, and \fB--perf-prof\fR are only available on Linux. \fB--enable-etw-stack-walking\fR is only available on Windows. . +.It Ev NODE_OTEL Ar value + +When set to \fB1\fR while the \fB--experimental-otel\fR flag is +enabled, activates the built-in OpenTelemetry tracing subsystem using the +default collector endpoint (\fBhttp://localhost:4318\fR). Values other than +\fB1\fR are ignored. If \fBNODE_OTEL_ENDPOINT\fR is also set, it takes precedence +for the endpoint. See the OpenTelemetry documentation for details. +. +.It Ev NODE_OTEL_ENDPOINT Ar url + +When set to a non-empty value while the \fB--experimental-otel\fR flag is +enabled, activates the built-in OpenTelemetry tracing subsystem and directs +spans to the specified OTLP/HTTP collector endpoint. The endpoint must be +the base URL of the collector: any path present in the endpoint is +replaced with \fB/v1/traces\fR, and an endpoint that already ends with +\fB/v1/traces\fR is used as is. When only \fBNODE_OTEL=1\fR is set, the default +collector endpoint (\fBhttp://localhost:4318\fR) is used. If \fBNODE_OTEL\fR is +also set, \fBNODE_OTEL_ENDPOINT\fR takes precedence for the endpoint. See the +OpenTelemetry documentation for details. +. +.It Ev NODE_OTEL_FILTER Ar module[,…] + +Comma-separated list of core modules to instrument when OpenTelemetry tracing +is active. When not set, all supported modules are instrumented. Supported +values: \fBnode:http\fR, \fBnode:undici\fR, \fBnode:fetch\fR. See the OpenTelemetry +documentation for details. +. +.It Ev NODE_OTEL_FLUSH_INTERVAL Ar milliseconds + +Interval in milliseconds between periodic flushes of buffered spans to the +collector. Must be a positive integer. \fBDefault:\fR \fB10000\fR. +. +.It Ev NODE_OTEL_MAX_BUFFER_SIZE Ar number + +Maximum number of spans buffered in memory before an immediate flush to the +collector is triggered. Must be a positive integer. \fBDefault:\fR \fB100\fR. +. .It Ev NODE_PATH Ar path[:…] \fB':'\fR-separated list of directories prefixed to the module search path. On Windows, this is a \fB';'\fR-separated list instead. diff --git a/lib/internal/otel/core.js b/lib/internal/otel/core.js new file mode 100644 index 000000000000..40d9c04f3c32 --- /dev/null +++ b/lib/internal/otel/core.js @@ -0,0 +1,121 @@ +'use strict'; + +const { + SafeSet, + StringPrototypeSplit, + StringPrototypeTrim, +} = primordials; + +const { kEmptyObject } = require('internal/util'); +const { + codes: { + ERR_INVALID_ARG_TYPE, + ERR_INVALID_ARG_VALUE, + }, +} = require('internal/errors'); +const { + validateInteger, + validateObject, + validateString, +} = require('internal/validators'); + +const { AsyncLocalStorage } = require('async_hooks'); +const { URL } = require('internal/url'); + +const kDefaultEndpoint = 'http://localhost:4318'; + +let endpoint = null; +let active = false; +let filter = null; // null = all modules enabled; SafeSet = only listed modules. +let spanStorage = null; +let collectorHost = null; // Normalized host (e.g. "localhost:4318") for HTTP client filtering. + +function getSpanStorage() { + spanStorage ??= new AsyncLocalStorage(); + return spanStorage; +} + +function isActive() { + return active; +} + +function getEndpoint() { + return endpoint; +} + +function getCollectorHost() { + return collectorHost; +} + +function isModuleEnabled(moduleName) { + if (filter == null) return true; + return filter.has(moduleName); +} + +function parseFilter(filter) { + if (filter == null) return null; + if (typeof filter === 'string') { + const parts = StringPrototypeSplit(filter, ','); + const set = new SafeSet(); + for (let i = 0; i < parts.length; i++) { + const trimmed = StringPrototypeTrim(parts[i]); + if (trimmed) set.add(trimmed); + } + return set; + } + throw new ERR_INVALID_ARG_TYPE('options.filter', + ['string', 'null', 'undefined'], + filter); +} + +function start(options = kEmptyObject) { + // Tracing is activated at most once per process. There is no way to stop + // it again: disabling the flusher and instrumentations is not supported. + if (active) return; + + validateObject(options, 'options'); + + const endpointValue = options.endpoint ?? kDefaultEndpoint; + validateString(endpointValue, 'options.endpoint'); + if (!endpointValue) { + throw new ERR_INVALID_ARG_VALUE('options.endpoint', endpointValue, + 'must be a non-empty string'); + } + + let parsed; + try { + parsed = new URL(endpointValue); + } catch { + throw new ERR_INVALID_ARG_VALUE('options.endpoint', endpointValue, + 'must be a valid URL'); + } + + if (options.maxBufferSize !== undefined) { + validateInteger(options.maxBufferSize, 'options.maxBufferSize', 1); + } + if (options.flushInterval !== undefined) { + validateInteger(options.flushInterval, 'options.flushInterval', 1); + } + + const parsedFilter = parseFilter(options.filter); + + endpoint = endpointValue; + filter = parsedFilter; + active = true; + collectorHost = parsed.host; + + const { enableInstrumentations } = require('internal/otel/instrumentations'); + enableInstrumentations(); + + const { startFlusher } = require('internal/otel/flush'); + startFlusher(options); +} + +module.exports = { + start, + isActive, + getEndpoint, + getSpanStorage, + isModuleEnabled, + getCollectorHost, +}; diff --git a/lib/internal/otel/flush.js b/lib/internal/otel/flush.js new file mode 100644 index 000000000000..20987cfc2d75 --- /dev/null +++ b/lib/internal/otel/flush.js @@ -0,0 +1,347 @@ +'use strict'; + +const { + ArrayPrototypePush, + BigInt, + DateNow, + JSONStringify, + MathRound, + NumberIsInteger, + ObjectKeys, + String, +} = primordials; + +const { Buffer } = require('buffer'); +const http = require('http'); +const https = require('https'); +const { clearTimeout, setImmediate, setTimeout } = require('timers'); +const { URL } = require('internal/url'); +const { getEndpoint } = require('internal/otel/core'); + +const kDefaultMaxBufferSize = 100; +const kDefaultFlushIntervalMs = 10_000; +const kWarningThrottleMs = 30_000; +const kExportTimeoutMs = 10_000; + +// One keep-alive socket per protocol: connections are reused, and +// batches queue behind it, bounding in-flight exports to a slow collector. +const kExportAgents = { + '__proto__': null, + 'http:': new http.Agent({ keepAlive: true, maxSockets: 1 }), + 'https:': new https.Agent({ keepAlive: true, maxSockets: 1 }), +}; + +let spanBuffer = []; +let flushTimer = null; // One-shot timer; only set while spans are buffered. +let flushScheduled = false; // An off-thread flush is already scheduled. +let flushIntervalMs = kDefaultFlushIntervalMs; +let maxBufferSize = kDefaultMaxBufferSize; +let exportUrl = null; // Resolved once; the endpoint is fixed after start. + +// Read once at activation time so that later env mutations cannot change +// the service name of exports. The default follows the OpenTelemetry +// semantic conventions: a low-cardinality fallback value, deliberately +// without the pid. +const kServiceName = process.env.OTEL_SERVICE_NAME || + 'unknown_service:node'; + +let exportFailureCount = 0; +let lastExportWarningTime = 0; + +let cachedResource = null; +let cachedScope = null; + +function getResource() { + cachedResource ??= { + attributes: [ + { key: 'service.name', + value: { stringValue: kServiceName } }, + { key: 'telemetry.sdk.name', + value: { stringValue: 'nodejs-core' } }, + { key: 'telemetry.sdk.language', + value: { stringValue: 'nodejs' } }, + { key: 'telemetry.sdk.version', + value: { stringValue: process.version } }, + { key: 'process.runtime.name', + value: { stringValue: 'nodejs' } }, + { key: 'process.runtime.version', + value: { stringValue: process.version } }, + { key: 'process.pid', + value: { intValue: String(process.pid) } }, + ], + }; + return cachedResource; +} + +function getScope() { + cachedScope ??= { + name: 'nodejs-core', + version: process.version, + }; + return cachedScope; +} + +function encodeAttributeValue(value) { + if (typeof value === 'string') { + return { stringValue: value }; + } + if (typeof value === 'number') { + if (NumberIsInteger(value)) { + return { intValue: String(value) }; + } + return { doubleValue: value }; + } + if (typeof value === 'boolean') { + return { boolValue: value }; + } + return { stringValue: String(value) }; +} + +// Converts a wall-clock epoch millisecond timestamp into an OTLP +// nanosecond string. Rounding to whole milliseconds before multiplying +// keeps the product exactly representable in double precision. +function timeToUnixNano(t) { + return `${BigInt(MathRound(t)) * 1_000_000n}`; +} + +function spanToOtlp(span) { + // TODO(bengl): A lot of objects are created in here for all the atributes. + // As a future optimization, we could hand-write the JSON encoding. + const rawAttrs = span.getAttributes(); + const attrKeys = ObjectKeys(rawAttrs); + const attributes = []; + for (let i = 0; i < attrKeys.length; i++) { + // Attribute values may be deferred functions that only run at export + // time; values that resolve to undefined are omitted. + // Unused preallocated keys stay undefined and are skipped. + if (rawAttrs[attrKeys[i]] === undefined) continue; + ArrayPrototypePush(attributes, { + key: attrKeys[i], + value: encodeAttributeValue(rawAttrs[attrKeys[i]]), + }); + } + + const rawEvents = span.getEvents(); + const events = []; + for (let i = 0; i < rawEvents.length; i++) { + const event = rawEvents[i]; + const eventAttrs = []; + const eventAttrKeys = ObjectKeys(event.attributes); + for (let j = 0; j < eventAttrKeys.length; j++) { + if (event.attributes[eventAttrKeys[j]] === undefined) continue; + ArrayPrototypePush(eventAttrs, { + key: eventAttrKeys[j], + value: encodeAttributeValue(event.attributes[eventAttrKeys[j]]), + }); + } + const otlpEvent = { + name: event.name, + timeUnixNano: timeToUnixNano(event.time), + }; + if (eventAttrs.length > 0) { + otlpEvent.attributes = eventAttrs; + } + ArrayPrototypePush(events, otlpEvent); + } + + const otlpSpan = { + traceId: span.traceId, + spanId: span.spanId, + name: span.name, + kind: span.kind, + startTimeUnixNano: timeToUnixNano(span.startTime), + endTimeUnixNano: timeToUnixNano(span.endTime), + }; + + if (attributes.length > 0) { + otlpSpan.attributes = attributes; + } + + const status = span.status; + if (status.code !== 0 || status.message) { + otlpSpan.status = status; + } + + if (span.parentSpanId) { + otlpSpan.parentSpanId = span.parentSpanId; + } + + if (events.length > 0) { + otlpSpan.events = events; + } + + return otlpSpan; +} + +function addSpan(span) { + ArrayPrototypePush(spanBuffer, span); + if (spanBuffer.length >= maxBufferSize) { + // Defer the flush so serialization happens off the request path. + if (!flushScheduled) { + flushScheduled = true; + setImmediate(() => { + flushScheduled = false; + flush(); + }); + } + } else if (flushTimer == null) { + // One-shot timer; only exists while spans are buffered. + flushTimer = setTimeout(flush, flushIntervalMs); + flushTimer.unref(); + } +} + +function flush() { + if (flushTimer != null) { + clearTimeout(flushTimer); + flushTimer = null; + } + if (spanBuffer.length === 0) return; + if (getEndpoint() == null) return; + + const spans = spanBuffer; + spanBuffer = []; + + const otlpSpans = []; + for (let i = 0; i < spans.length; i++) { + try { + ArrayPrototypePush(otlpSpans, spanToOtlp(spans[i])); + } catch (err) { + warnExportFailure( + `Failed to serialize span "${spans[i].name}": ${err.message}`); + } + } + + if (otlpSpans.length === 0) return; + + let payload; + try { + payload = JSONStringify({ + resourceSpans: [{ + resource: getResource(), + scopeSpans: [{ + scope: getScope(), + spans: otlpSpans, + }], + }], + }); + } catch (err) { + warnExportFailure( + `Failed to serialize ${otlpSpans.length} spans: ${err.message}`); + return; + } + + sendToCollector(payload); +} + +// Emits a throttled OTelExportWarning and counts the failure. Warnings are +// throttled to at most one per kWarningThrottleMs. +function warnExportFailure(message) { + exportFailureCount++; + const now = DateNow(); + if (now - lastExportWarningTime >= kWarningThrottleMs) { + lastExportWarningTime = now; + const suffix = exportFailureCount > 1 ? + ` (${exportFailureCount} total failures)` : ''; + process.emitWarning( + `${message}${suffix}`, + 'OTelExportWarning', + ); + } +} + +function sendToCollector(body) { + const endpoint = getEndpoint(); + if (endpoint == null) return; + + try { + // The endpoint is fixed after start, so resolve the export URL once. + // /v1/traces is the standard OTLP/HTTP traces path; an endpoint + // already ending with it yields the same URL. + exportUrl ??= new URL('/v1/traces', endpoint); + const parsed = exportUrl; + + if (parsed.protocol !== 'https:' && parsed.protocol !== 'http:') { + warnExportFailure( + `Unsupported protocol "${parsed.protocol}" in OTLP endpoint; ` + + 'only http: and https: are supported'); + return; + } + + const transport = parsed.protocol === 'https:' ? https : http; + + const req = transport.request({ + hostname: parsed.hostname, + port: parsed.port, + path: parsed.pathname, + method: 'POST', + agent: kExportAgents[parsed.protocol], + timeout: kExportTimeoutMs, + headers: { + 'content-type': 'application/json', + 'content-length': Buffer.byteLength(body), + }, + }, (res) => { + // TODO(bengl): Once retry logic is added, parse the response body for + // ExportTraceServiceResponse.partial_success.rejected_spans. + res.resume(); + res.on('end', () => { + if (res.statusCode >= 400) { + warnExportFailure( + `OTLP collector responded with HTTP ${res.statusCode}`); + } + }); + res.on('error', (err) => { + warnExportFailure( + `OTLP export response stream error: ${err.message}`); + }); + }); + + req.on('timeout', () => { + // Destroying the request routes the failure into the (throttled) + // warning path and frees the socket for the next batch. + req.destroy(); + }); + + req.on('error', (err) => { + warnExportFailure( + `Failed to export spans to ${endpoint}: ${err.message}`); + }); + + req.end(body); + } catch (err) { + warnExportFailure( + `Failed to export spans to ${endpoint}: ${err.message}`); + } +} + +function startFlusher(options) { + maxBufferSize = options?.maxBufferSize ?? kDefaultMaxBufferSize; + flushIntervalMs = options?.flushInterval ?? kDefaultFlushIntervalMs; + + process.on('beforeExit', flush); +} + +// Resets per-run state. Used by the test suite to isolate scenarios; +// tracing itself never resets once started. +function resetCaches() { + if (flushTimer != null) { + clearTimeout(flushTimer); + flushTimer = null; + } + flushScheduled = false; + flushIntervalMs = kDefaultFlushIntervalMs; + cachedResource = null; + cachedScope = null; + exportFailureCount = 0; + lastExportWarningTime = 0; + maxBufferSize = kDefaultMaxBufferSize; + spanBuffer = []; +} + +module.exports = { + addSpan, + flush, + startFlusher, + resetCaches, +}; diff --git a/lib/internal/otel/id.js b/lib/internal/otel/id.js new file mode 100644 index 000000000000..29db05785ba2 --- /dev/null +++ b/lib/internal/otel/id.js @@ -0,0 +1,56 @@ +'use strict'; + +const { + Array, + NumberPrototypeToString, + StringPrototypePadStart, + Uint8Array, +} = primordials; + +const { randomFillSync } = require('internal/crypto/random'); + +// A 4KB buffer of random bytes, refilled when exhausted (once every ~170 +// spans: 16 bytes per trace ID, 8 per span ID), amortizes the CSPRNG cost. +// The CSPRNG stays: the W3C Trace Context spec recommends unpredictable +// IDs, and they are visible in outgoing requests. benchmark/otel/http.js +// tracks the per-request overhead. +const kBufferSize = 4096; +const randomBuffer = new Uint8Array(kBufferSize); +let randomOffset = kBufferSize; // Start at end to trigger first fill + +// Hex lookup table for fast byte-to-hex conversion. +const hexTable = new Array(256); +for (let i = 0; i < 256; i++) { + hexTable[i] = StringPrototypePadStart( + NumberPrototypeToString(i, 16), 2, '0', + ); +} + +function ensureRandomBytes(needed) { + if (randomOffset + needed > kBufferSize) { + randomFillSync(randomBuffer); + randomOffset = 0; + } +} + +function generateId(bytes) { + ensureRandomBytes(bytes); + let id = ''; + for (let i = 0; i < bytes; i++) { + id += hexTable[randomBuffer[randomOffset++]]; + } + return id; +} + +function generateTraceId() { + return generateId(16); +} + +function generateSpanId() { + return generateId(8); +} + +module.exports = { + generateTraceId, + generateSpanId, +}; diff --git a/lib/internal/otel/instrumentations.js b/lib/internal/otel/instrumentations.js new file mode 100644 index 000000000000..20f0036c12b2 --- /dev/null +++ b/lib/internal/otel/instrumentations.js @@ -0,0 +1,263 @@ +'use strict'; + +const { + StringPrototypeIndexOf, + StringPrototypeSlice, +} = primordials; + +const dc = require('diagnostics_channel'); +const { + Span, + SPAN_KIND_SERVER, + SPAN_KIND_CLIENT, + STATUS_ERROR, + kSpan, +} = require('internal/otel/span'); +const { + isModuleEnabled, + getSpanStorage, + getCollectorHost, +} = require('internal/otel/core'); + +// diagnostics_channel does not isolate subscriber exceptions: a throw +// here propagates into the instrumented module. Handlers must never throw. + +// Cached at subscription time to avoid a per-event AsyncLocalStorage +// lookup. +let spanStorage; + +// Shared response handling for client spans (http and undici alike). The +// caller decides when the span ends. +function applyClientResponse(span, statusCode) { + span.setAttribute('http.response.status_code', statusCode); + + if (statusCode >= 400) { + span.setStatus(STATUS_ERROR, `HTTP ${statusCode}`); + span.setAttribute('error.type', `${statusCode}`); + } +} + +// Shared error handling for client spans (http and undici alike). +function failClientSpan(span, error) { + span.setAttribute('error.type', error?.name || 'Error'); + span.setStatus(STATUS_ERROR, error?.message || 'unknown error'); + span.addEvent('exception', { + 'exception.type': error?.name || 'Error', + 'exception.message': error?.message || '', + 'exception.stacktrace': error?.stack || '', + }); + span.end(); +} + +function onHttpServerRequestStart({ request, socket }) { + // Skip the export request itself. Without this, tracing a collector + // hosted in the same process would feed itself indefinitely: every + // export would produce a server span that is exported again. + if (request.url === '/v1/traces' && + request.headers?.host === getCollectorHost()) { + return; + } + + const traceparent = request.headers?.traceparent; + const tracestate = request.headers?.tracestate; + const parent = traceparent != null ? + Span.fromTraceparent(traceparent) : undefined; + + const method = request.method || 'GET'; + const rawUrl = request.url || '/'; + const qIdx = StringPrototypeIndexOf(rawUrl, '?'); + const path = qIdx === -1 ? rawUrl : StringPrototypeSlice(rawUrl, 0, qIdx); + // Per the OpenTelemetry semantic conventions, span names must be low + // cardinality. Node.js core has no route concept, so only the method is + // used; url.path remains available as an attribute. + const span = new Span(method, SPAN_KIND_SERVER, { + parent, + // The tracestate is forwarded verbatim; it is never parsed. + tracestate, + }); + + span.setAttribute('http.request.method', method); + span.setAttribute('url.path', path); + if (qIdx !== -1) { + span.setAttribute('url.query', StringPrototypeSlice(rawUrl, qIdx + 1)); + } + + span.setAttribute('url.scheme', socket?.encrypted ? 'https' : 'http'); + span.setAttribute('network.protocol.version', + request.httpVersion || '1.1'); + + const host = request.headers?.host; + if (host) { + span.setAttribute('server.address', host); + } + + // The span is stashed on the request so that the response handlers do + // not need a second AsyncLocalStorage lookup to find it. + request[kSpan] = span; + + // The span deliberately stays in the AsyncLocalStorage after the request + // ends: work started while handling it (timers, outgoing requests from + // response callbacks) still links to the server span as its parent. + spanStorage.enterWith(span); + + request.on('close', () => { + if (span.endTime === undefined) { + span.setStatus(STATUS_ERROR, 'request closed before response'); + span.end(); + } + }); +} + +function onHttpServerResponseFinish({ request, response }) { + const span = request[kSpan]; + if (span == null) return; + + const statusCode = response.statusCode; + span.setAttribute('http.response.status_code', statusCode); + + if (statusCode >= 500) { + span.setStatus(STATUS_ERROR, `HTTP ${statusCode}`); + span.setAttribute('error.type', `${statusCode}`); + } + + span.end(); +} + +function onHttpClientRequestCreated({ request }) { + if (request.getHeader('host') === getCollectorHost()) return; + + const parent = spanStorage.getStore(); + + const method = request.method || 'GET'; + const span = new Span( + method, + SPAN_KIND_CLIENT, + { parent }, + ); + + span.setAttribute('http.request.method', method); + + const { protocol, host } = request; + const path = request.path; + if (host) { + span.setAttribute('server.address', host); + } + + const port = request.socket?.remotePort || request.port; + if (port != null) { + span.setAttribute('server.port', port); + } + + span.setAttribute('url.full', `${protocol}//${host}${path}`); + + request[kSpan] = span; + + request.setHeader('traceparent', span.toTraceparent()); + if (span.tracestate) { + request.setHeader('tracestate', span.tracestate); + } +} + +function onHttpClientResponseFinish({ request, response }) { + const span = request[kSpan]; + if (span == null) return; + + span.setAttribute('network.protocol.version', + response.httpVersion || '1.1'); + applyClientResponse(span, response.statusCode); + + // End the span when the response body completes. 'end' only fires when + // the body has been fully read; 'close' covers aborted responses. + const endSpan = () => span.end(); + response.once('end', endSpan); + response.once('close', endSpan); +} + +function onClientRequestError({ request, error }) { + const span = request[kSpan]; + if (span == null) return; + + failClientSpan(span, error); +} + +// Undici's addHeader appends rather than replaces; update an existing +// header in place so that a caller-supplied value is overridden, not +// duplicated. +function setUndiciHeader(request, name, value) { + const headers = request.headers; + for (let i = 0; i < headers.length; i += 2) { + if (headers[i] === name) { + headers[i + 1] = value; + return; + } + } + request.addHeader(name, value); +} + +function onUndiciRequestCreate({ request }) { + const parent = spanStorage.getStore(); + + const method = request.method || 'GET'; + const span = new Span( + method, + SPAN_KIND_CLIENT, + { parent }, + ); + + span.setAttribute('http.request.method', method); + + const { origin } = request; + const path = request.path; + if (origin) { + span.setAttribute('server.address', origin); + } + span.setAttribute('url.full', `${origin}${path}`); + + request[kSpan] = span; + + if (request.addHeader) { + setUndiciHeader(request, 'traceparent', span.toTraceparent()); + if (span.tracestate) { + setUndiciHeader(request, 'tracestate', span.tracestate); + } + } +} + +function onUndiciRequestHeaders({ request, response }) { + const span = request[kSpan]; + if (span == null) return; + + applyClientResponse(span, response.statusCode); +} + +// undici:request:trailers fires when the response completes, including its +// body, so the span duration covers the full transfer. +function onUndiciRequestTrailers({ request }) { + const span = request[kSpan]; + if (span == null) return; + + span.end(); +} + +function enableInstrumentations() { + spanStorage = getSpanStorage(); + + if (isModuleEnabled('node:http')) { + dc.subscribe('http.server.request.start', onHttpServerRequestStart); + dc.subscribe('http.server.response.finish', onHttpServerResponseFinish); + dc.subscribe('http.client.request.created', onHttpClientRequestCreated); + dc.subscribe('http.client.response.finish', onHttpClientResponseFinish); + dc.subscribe('http.client.request.error', onClientRequestError); + } + + if (isModuleEnabled('node:undici') || isModuleEnabled('node:fetch')) { + dc.subscribe('undici:request:create', onUndiciRequestCreate); + dc.subscribe('undici:request:headers', onUndiciRequestHeaders); + dc.subscribe('undici:request:trailers', onUndiciRequestTrailers); + dc.subscribe('undici:request:error', onClientRequestError); + } +} + +module.exports = { + enableInstrumentations, +}; diff --git a/lib/internal/otel/span.js b/lib/internal/otel/span.js new file mode 100644 index 000000000000..d28f7932e471 --- /dev/null +++ b/lib/internal/otel/span.js @@ -0,0 +1,208 @@ +'use strict'; + +const { + ArrayPrototypePush, + DateNow, + NumberParseInt, + NumberPrototypeToString, + RegExpPrototypeExec, + StringPrototypePadStart, + StringPrototypeSplit, + Symbol, +} = primordials; + +const { addSpan } = require('internal/otel/flush'); +const { generateTraceId, generateSpanId } = require('internal/otel/id'); +const { now } = require('internal/perf/utils'); + +// Span kind and status code values, per the SpanKind and StatusCode +// enums of the OpenTelemetry protocol: +// https://github.com/open-telemetry/opentelemetry-proto/blob/main/opentelemetry/proto/trace/v1/trace.proto +const SPAN_KIND_INTERNAL = 1; +const SPAN_KIND_SERVER = 2; +const SPAN_KIND_CLIENT = 3; + +// (Status codes: 0 UNSET, 1 OK, 2 ERROR.) +const STATUS_UNSET = 0; +const STATUS_ERROR = 2; + +const kSpan = Symbol('kOtelSpan'); + +const kHex32 = /^[0-9a-f]{32}$/; +const kHex16 = /^[0-9a-f]{16}$/; +const kHex2 = /^[0-9a-f]{2}$/; +const kAllZero32 = '00000000000000000000000000000000'; +const kAllZero16 = '0000000000000000'; + +class Span { + traceId; + spanId; + parentSpanId; + name; + kind; + // Raw W3C tracestate header value, forwarded verbatim to child + // requests. Empty when no tracestate was received. + tracestate; + // Timestamps are epoch milliseconds. The start is read from the wall + // clock (Date.now()); the end is derived as start + a monotonic duration + // (performance.now()), so that spans always have human-readable absolute + // timestamps while remaining safe from negative durations when the + // system clock steps backward mid-span. All timestamps are only converted + // to OTLP nanosecond strings at export time, so that unsampled or + // filtered spans never pay for the conversion. + startTime; + endTime; + #startMonotonic; + + #attributes; + #events; + #status; + #traceFlags; + + constructor(name, kind, options) { + const parent = options?.parent; + + this.name = name; + this.kind = kind; + this.spanId = generateSpanId(); + // The attribute keys below are the complete set written by the + // instrumentations (there is no public attribute API yet). + // Preallocating them gives every span's attribute object the same + // static shape, so writes never trigger hidden class transitions or + // dictionary mode. Unused keys stay undefined and are skipped at + // export time. + this.#attributes = { + '__proto__': null, + 'error.type': undefined, + 'http.request.method': undefined, + 'http.response.status_code': undefined, + 'network.protocol.version': undefined, + 'server.address': undefined, + 'server.port': undefined, + 'url.full': undefined, + 'url.path': undefined, + 'url.query': undefined, + 'url.scheme': undefined, + }; + this.#events = []; + this.#status = { code: STATUS_UNSET, message: '' }; + this.startTime = DateNow(); + this.#startMonotonic = now(); + + if (parent != null) { + this.traceId = parent.traceId; + this.parentSpanId = parent.spanId; + this.#traceFlags = parent.traceFlags; + } else { + this.traceId = generateTraceId(); + this.parentSpanId = ''; + this.#traceFlags = 0x01; // Sampled by default. + } + // An explicit tracestate wins; children inherit their parent's. + this.tracestate = options?.tracestate || parent?.tracestate || ''; + } + + get traceFlags() { + return this.#traceFlags; + } + + setAttribute(key, value) { + this.#attributes[key] = value; + return this; + } + + getAttributes() { + return this.#attributes; + } + + addEvent(name, attributes) { + // There is no public API for span events yet and the only internal + // producer adds at most one exception event per span. Once events are + // exposed through a public API, a per-span cap will be needed. + // TODO(bengl): Cap the number of events per span when addEvent becomes + // part of a public API. + ArrayPrototypePush(this.#events, { + name, + // Anchored to the span's start like endTime, for the same clock- + // step safety. + time: this.startTime + (now() - this.#startMonotonic), + attributes: attributes || { __proto__: null }, + }); + return this; + } + + getEvents() { + return this.#events; + } + + get status() { + return this.#status; + } + + setStatus(code, message) { + this.#status = { code, message: message || '' }; + return this; + } + + end() { + if (this.endTime !== undefined) return; // Already ended. + this.endTime = this.startTime + (now() - this.#startMonotonic); + + // Only export sampled spans. + if (this.#traceFlags & 0x01) { + addSpan(this); + } + } + + // The counterpart of static fromTraceparent(): produces the outgoing + // header value; the caller writes it into the request. + // Format: {version}-{trace-id}-{span-id}-{trace-flags} + toTraceparent() { + // The common case is a sampled root span (flags 0x01). + const flags = this.#traceFlags === 0x01 ? + '01' : + StringPrototypePadStart( + NumberPrototypeToString(this.#traceFlags, 16), 2, '0'); + return `00-${this.traceId}-${this.spanId}-${flags}`; + } + + // The counterpart of toTraceparent(): parses an incoming header value + // into a remote parent. Returns null when invalid per the W3C spec. + static fromTraceparent(traceparentHeader) { + if (typeof traceparentHeader !== 'string') { + return null; + } + + const parts = StringPrototypeSplit(traceparentHeader, '-'); + if (parts.length !== 4) return null; + + // Only version 00 is defined; reject everything else. + if (parts[0] !== '00') return null; + const traceId = parts[1]; + const spanId = parts[2]; + const flags = parts[3]; + if (RegExpPrototypeExec(kHex32, traceId) === null) return null; + if (RegExpPrototypeExec(kHex16, spanId) === null) return null; + if (RegExpPrototypeExec(kHex2, flags) === null) return null; + + if (traceId === kAllZero32) return null; + if (spanId === kAllZero16) return null; + + return { + __proto__: null, + traceId, + spanId, + traceFlags: NumberParseInt(flags, 16), + }; + } +} + +module.exports = { + Span, + SPAN_KIND_INTERNAL, + SPAN_KIND_SERVER, + SPAN_KIND_CLIENT, + STATUS_UNSET, + STATUS_ERROR, + kSpan, +}; diff --git a/lib/internal/process/pre_execution.js b/lib/internal/process/pre_execution.js index 0f5aaf857aeb..ca773ec61818 100644 --- a/lib/internal/process/pre_execution.js +++ b/lib/internal/process/pre_execution.js @@ -12,6 +12,7 @@ const { DatePrototypeGetMinutes, DatePrototypeGetMonth, DatePrototypeGetSeconds, + Number, NumberParseInt, ObjectDefineProperty, ObjectFreeze, @@ -124,6 +125,11 @@ function prepareExecution(options) { setupEventsource(); setupCodeCoverage(); setupDebugEnv(); + // Built-in OpenTelemetry tracing is main-thread only. Workers inherit + // the environment variables but do not activate the subsystem. + if (isMainThread) { + setupOtel(); + } // Process initial diagnostic reporting configuration, if present. initializeReport(); @@ -558,6 +564,57 @@ function setupDebugEnv() { } } +function setupOtel() { + // Tracing must never be baked into startup snapshots: activation belongs + // to the environment of the process that actually runs. + if (isBuildingSnapshot()) return; + + const endpoint = process.env.NODE_OTEL_ENDPOINT; + // Like other NODE_* environment variables, only the value 1 activates. + const enabled = process.env.NODE_OTEL === '1'; + const hasEnvVars = !!endpoint || !!process.env.NODE_OTEL; + + if (!getOptionValue('--experimental-otel')) { + if (hasEnvVars) { + process.emitWarning( + 'NODE_OTEL environment variables are set but --experimental-otel ' + + 'is not enabled; built-in OpenTelemetry tracing is not active', + 'OTelWarning', + ); + } + return; + } + + // Tracing is activated via environment variables only: + // NODE_OTEL_ENDPOINT directs spans to a specific OTLP collector, while + // NODE_OTEL=1 enables tracing with the default endpoint + // (http://localhost:4318). + if (!endpoint && !enabled) return; + + try { + const { start } = require('internal/otel/core'); + const opts = { __proto__: null }; + if (endpoint) opts.endpoint = endpoint; + const envFilter = process.env.NODE_OTEL_FILTER; + if (envFilter) opts.filter = envFilter; + const envMaxBufferSize = process.env.NODE_OTEL_MAX_BUFFER_SIZE; + if (envMaxBufferSize) opts.maxBufferSize = Number(envMaxBufferSize); + const envFlushInterval = process.env.NODE_OTEL_FLUSH_INTERVAL; + if (envFlushInterval) opts.flushInterval = Number(envFlushInterval); + start(opts); + process.emitWarning( + 'Built-in OpenTelemetry tracing is an experimental feature and might ' + + 'change at any time', + 'ExperimentalWarning', + ); + } catch (err) { + process.emitWarning( + `Failed to initialize OpenTelemetry tracing: ${err.message}`, + 'OTelWarning', + ); + } +} + // This has to be called after initializeReport() is called function initializeReportSignalHandlers() { if (getOptionValue('--report-on-signal')) { diff --git a/src/node_options.cc b/src/node_options.cc index acd98e6b9276..3b365e533768 100644 --- a/src/node_options.cc +++ b/src/node_options.cc @@ -815,6 +815,10 @@ EnvironmentOptionsParser::EnvironmentOptionsParser() { kDisallowedInEnvvar); AddOption("[vfs_load_set]", "", BOOL_FIELD(vfs_load)); Implies("--vfs-load", "[vfs_load_set]"); + AddOption("--experimental-otel", + "experimental built-in OpenTelemetry tracing", + BOOL_FIELD(experimental_otel), + kAllowedInEnvvar); AddOption("--experimental-quic", #ifndef OPENSSL_NO_QUIC "experimental QUIC support", diff --git a/src/node_options.h b/src/node_options.h index 37e05101f562..3583fa1cc206 100644 --- a/src/node_options.h +++ b/src/node_options.h @@ -219,6 +219,11 @@ class EnvironmentOptions : public Options { DEFINE_BOOL_FIELD(vfs_load) = false; DEFINE_BOOL_FIELD(webstorage) = HAVE_SQLITE; DEFINE_BOOL_FIELD(experimental_dtls) = EXPERIMENTALS_DEFAULT_VALUE; + // Uses EXPERIMENTALS_DEFAULT_VALUE (rather than a hard-coded false) so + // that builds which enable all experimental features by default opt + // into this flag as well. Tracing still requires explicit activation + // through environment variables either way. + DEFINE_BOOL_FIELD(experimental_otel) = EXPERIMENTALS_DEFAULT_VALUE; DEFINE_BOOL_FIELD(experimental_quic) = EXPERIMENTALS_DEFAULT_VALUE; DEFINE_BOOL_FIELD(experimental_global_navigator) = true; DEFINE_BOOL_FIELD(experimental_global_web_crypto) = true; diff --git a/test/common/otel.js b/test/common/otel.js new file mode 100644 index 000000000000..595236a95862 --- /dev/null +++ b/test/common/otel.js @@ -0,0 +1,116 @@ +'use strict'; + +const http = require('node:http'); + +// Helpers for tests of the built-in OpenTelemetry tracing subsystem +// (see doc/api/otel.md). They require no internals: tracing is either +// activated from scratch via environment variables in a child process, +// or started directly through `internal/otel/core` under +// `--expose-internals`. + +// Starts a local HTTP server that acts as an OTLP/HTTP JSON collector. +// +// Resolves with an object providing: +// * port: the port the collector listens on +// * next(): a promise that resolves with the next OTLP payload received, +// whether it arrives before or after next() is called +// * payloads: an array of every OTLP payload received so far +// * requests: an array of { url, headers } for each export request +// * close(): a promise that resolves once the collector has closed +async function startOTelCollector() { + const requests = []; + const payloads = []; + const queue = []; + let signal; + let signalPromise = new Promise((resolve) => { signal = resolve; }); + + const server = http.createServer((req, res) => { + let body = ''; + req.on('data', (chunk) => { body += chunk; }); + req.on('end', () => { + requests.push({ url: req.url, headers: req.headers }); + const payload = JSON.parse(body); + payloads.push(payload); + queue.push(payload); + res.writeHead(200); + res.end(); + signal(); + signalPromise = new Promise((resolve) => { signal = resolve; }); + }); + }); + + await new Promise((resolve) => server.listen(0, '127.0.0.1', resolve)); + + return { + port: server.address().port, + async next() { + while (queue.length === 0) { + await signalPromise; + } + return queue.shift(); + }, + payloads, + requests, + close: () => new Promise((resolve) => server.close(resolve)), + }; +} + +// Starts a plain HTTP server for tests to make traced requests against. +// The default handler responds with 200 OK. +// +// Resolves with an object providing: +// * port: the port the server listens on +// * close(): a promise that resolves once the server has closed +async function startHTTPServer( + handler = (req, res) => { + res.writeHead(200); + res.end('ok'); + }, +) { + const server = http.createServer(handler); + await new Promise((resolve) => server.listen(0, '127.0.0.1', resolve)); + + return { + port: server.address().port, + close: () => new Promise((resolve) => server.close(resolve)), + }; +} + +// Extracts the spans array from an OTLP/HTTP JSON export payload. +function getSpans(payload) { + return payload.resourceSpans[0].scopeSpans[0].spans; +} + +// Makes one traced GET request against the given server and waits for the +// response to finish. +function httpGet(port, path) { + return new Promise((resolve, reject) => { + http.get(`http://127.0.0.1:${port}${path}`, (res) => { + res.resume(); + res.on('end', resolve); + }).on('error', reject); + }); +} + +// Decodes a span's OTLP attributes into a plain key/value object. +function spanAttrs(span) { + const attrs = {}; + for (const a of span.attributes ?? []) { + attrs[a.key] = + a.value.stringValue || a.value.intValue || a.value.doubleValue; + } + return attrs; +} + +function delay(ms) { + return new Promise((resolve) => setTimeout(resolve, ms)); +} + +module.exports = { + startOTelCollector, + startHTTPServer, + getSpans, + httpGet, + spanAttrs, + delay, +}; diff --git a/test/doctool/test-doc-api-json.mjs b/test/doctool/test-doc-api-json.mjs index ff063e018d0e..a041b7bf7ffa 100644 --- a/test/doctool/test-doc-api-json.mjs +++ b/test/doctool/test-doc-api-json.mjs @@ -150,6 +150,7 @@ for await (const dirent of await fs.opendir(new URL('../../out/doc/api/', import 'index.md': [], 'intl.md': ['introduced_in', 'miscs'], 'n-api.md': ['introduced_in', 'stability', 'stabilityText', 'miscs'], + 'otel.md': ['introduced_in', 'meta', 'stability', 'stabilityText', 'miscs'], 'packages.md': ['introduced_in', 'meta', 'miscs'], 'process.md': ['globals'], 'report.md': ['introduced_in', 'meta', 'stability', 'stabilityText', 'miscs'], diff --git a/test/parallel/test-otel-beforeexit-flush.js b/test/parallel/test-otel-beforeexit-flush.js new file mode 100644 index 000000000000..f58adc2a36dc --- /dev/null +++ b/test/parallel/test-otel-beforeexit-flush.js @@ -0,0 +1,61 @@ +'use strict'; + +require('../common'); +const assert = require('node:assert'); +const { describe, it } = require('node:test'); +const { spawn } = require('node:child_process'); +const { + startOTelCollector, + getSpans, +} = require('../common/otel'); + +// Test that buffered spans are flushed via the beforeExit handler when the +// process exits while tracing is active. Tracing is activated in the child +// via environment variables. + +describe('otel beforeExit flush', () => { + it('flushes buffered spans on process beforeExit', async () => { + const collector = await startOTelCollector(); + + // Spawn a child with tracing activated via NODE_OTEL_ENDPOINT. The child + // makes a request and exits without any explicit shutdown. The beforeExit + // handler should flush the buffered spans. + const script = ` +const http = require("http"); +const server = http.createServer((req, res) => { + res.writeHead(200); + res.end("ok"); +}); +server.listen(0, () => { + const port = server.address().port; + http.get("http://127.0.0.1:" + port + "/test", (res) => { + res.resume(); + res.on("end", () => { + server.close(); + }); + }); +}); +`; + + const child = spawn(process.execPath, [ + '--experimental-otel', '-e', script, + ], { + stdio: 'pipe', + env: { + ...process.env, + NODE_OTEL_ENDPOINT: `http://127.0.0.1:${collector.port}`, + }, + }); + + const exitCode = await new Promise((resolve) => { + child.on('exit', (code) => resolve(code)); + }); + assert.strictEqual(exitCode, 0); + + const spans = getSpans(await collector.next()); + + await collector.close(); + + assert.ok(spans.length > 0, 'Expected spans to be flushed on beforeExit'); + }); +}); diff --git a/test/parallel/test-otel-context-propagation.js b/test/parallel/test-otel-context-propagation.js new file mode 100644 index 000000000000..104a764a91e8 --- /dev/null +++ b/test/parallel/test-otel-context-propagation.js @@ -0,0 +1,104 @@ +'use strict'; +// Flags: --expose-internals + +require('../common'); +const assert = require('node:assert'); +const http = require('node:http'); +const net = require('node:net'); +const { describe, it } = require('node:test'); + +const otel = require('internal/otel/core'); +const { flush } = require('internal/otel/flush'); +const { + startOTelCollector, + startHTTPServer, + getSpans, +} = require('../common/otel'); + +// This test verifies W3C trace context propagation: +// An incoming request with a traceparent header causes child spans to +// share the same traceId, and outgoing requests carry the traceparent. + +describe('otel context propagation', () => { + it('propagates trace context across HTTP hops', async () => { + const incomingTraceId = 'abcdef0123456789abcdef0123456789'; + const incomingSpanId = '0123456789abcdef'; + const incomingTraceparent = + `00-${incomingTraceId}-${incomingSpanId}-01`; + const incomingTracestate = 'vendor=value'; + + let outgoingTraceparent = null; + let outgoingTracestate = null; + + const collector = await startOTelCollector(); + + otel.start({ endpoint: `http://127.0.0.1:${collector.port}` }); + + // Backend server that records the outgoing trace context headers. + const backend = await startHTTPServer((req, res) => { + outgoingTraceparent = req.headers.traceparent || null; + outgoingTracestate = req.headers.tracestate || null; + res.writeHead(200); + res.end('backend-ok'); + }); + + // Frontend server: receives request, makes outgoing call to backend. + const frontend = await startHTTPServer((req, res) => { + http.get(`http://127.0.0.1:${backend.port}/backend`, (backendRes) => { + backendRes.resume(); + backendRes.on('end', () => { + res.writeHead(200); + res.end('frontend-ok'); + }); + }); + }); + + // Use a raw TCP socket to send the initial request with a custom + // traceparent header. This bypasses the HTTP client instrumentation + // which would overwrite the header. + await new Promise((resolve) => { + const socket = net.connect(frontend.port, '127.0.0.1', () => { + socket.write( + `GET /frontend HTTP/1.1\r\n` + + `Host: 127.0.0.1:${frontend.port}\r\n` + + `traceparent: ${incomingTraceparent}\r\n` + + `tracestate: ${incomingTracestate}\r\n` + + `Connection: close\r\n` + + `\r\n` + ); + socket.on('data', () => {}); + socket.on('end', resolve); + }); + }); + + flush(); + const spans = getSpans(await collector.next()); + + await collector.close(); + await frontend.close(); + await backend.close(); + + assert.ok(spans.length >= 1, `Expected at least 1 span, got: ${spans.length}`); + + const serverSpan = spans.find((s) => { + if (s.kind !== 2) return false; // SERVER + const pathAttr = s.attributes?.find((a) => a.key === 'url.path'); + return pathAttr?.value?.stringValue === '/frontend'; + }); + assert.ok(serverSpan, 'Expected a frontend server span'); + assert.strictEqual(serverSpan.traceId, incomingTraceId); + + assert.strictEqual(serverSpan.parentSpanId, incomingSpanId); + + const clientSpan = spans.find((s) => s.kind === 3); // CLIENT + assert.ok(clientSpan, 'Expected a client span for outgoing backend call'); + assert.strictEqual(clientSpan.traceId, incomingTraceId); + + assert.ok(outgoingTraceparent); + assert.ok(outgoingTraceparent.includes(incomingTraceId), + 'Outgoing traceparent should carry the original traceId'); + + // The tracestate must be forwarded verbatim to child requests. + assert.strictEqual(outgoingTracestate, incomingTracestate); + }); +}); diff --git a/test/parallel/test-otel-env.js b/test/parallel/test-otel-env.js new file mode 100644 index 000000000000..7e60f39a71ca --- /dev/null +++ b/test/parallel/test-otel-env.js @@ -0,0 +1,283 @@ +'use strict'; +const common = require('../common'); +const assert = require('node:assert'); +const { spawn } = require('node:child_process'); +const { describe, it } = require('node:test'); +const { + startOTelCollector, + getSpans, +} = require('../common/otel'); + +// Tracing is activated only via environment variables and only when the +// --experimental-otel flag is passed. These tests spawn child processes to +// exercise that activation path. + +const activationScript = ` +const { isActive, getEndpoint } = require('internal/otel/core'); +console.log('active:' + isActive()); +console.log('endpoint:' + getEndpoint()); +`; + +function spawnOtel(env, args = []) { + return common.spawnPromisified( + process.execPath, + ['--expose-internals', ...args, '-e', activationScript], + { env: { ...process.env, ...env } }, + ); +} + +describe('otel environment variable activation', () => { + it('activates tracing when NODE_OTEL_ENDPOINT is set with the flag', async () => { + const { code, stdout, stderr } = await spawnOtel({ + NODE_OTEL_ENDPOINT: 'http://127.0.0.1:9999', + }, ['--experimental-otel']); + + assert.strictEqual(code, 0); + assert.match(stdout, /active:true/); + assert.match(stdout, /endpoint:http:\/\/127\.0\.0\.1:9999/); + assert.match(stderr, /ExperimentalWarning/); + }); + + it('activates tracing with the default endpoint when NODE_OTEL is set with the flag', async () => { + const { code, stdout } = await spawnOtel({ + NODE_OTEL: '1', + }, ['--experimental-otel']); + + assert.strictEqual(code, 0); + assert.match(stdout, /active:true/); + assert.match(stdout, /endpoint:http:\/\/localhost:4318/); + }); + + it('NODE_OTEL_ENDPOINT takes precedence over NODE_OTEL', async () => { + const { code, stdout } = await spawnOtel({ + NODE_OTEL: '1', + NODE_OTEL_ENDPOINT: 'http://127.0.0.1:9999', + }, ['--experimental-otel']); + + assert.strictEqual(code, 0); + assert.match(stdout, /endpoint:http:\/\/127\.0\.0\.1:9999/); + }); + + it('ignores NODE_OTEL values other than 1', async () => { + const { code, stdout, stderr } = await spawnOtel({ + NODE_OTEL: '0', + }, ['--experimental-otel']); + + assert.strictEqual(code, 0); + assert.match(stdout, /active:false/); + assert.strictEqual(stderr, ''); + }); + + it('does not activate tracing without the --experimental-otel flag', async () => { + const { code, stdout, stderr } = await spawnOtel({ + NODE_OTEL: '1', + NODE_OTEL_ENDPOINT: 'http://127.0.0.1:9999', + }); + + assert.strictEqual(code, 0); + assert.match(stdout, /active:false/); + assert.match(stdout, /endpoint:null/); + assert.match(stderr, /OTelWarning/); + assert.match(stderr, /--experimental-otel/); + }); + + it('NODE_OTEL_FILTER accepts node:fetch as an alias for undici', async () => { + const { code, stdout } = await common.spawnPromisified(process.execPath, [ + '--experimental-otel', + '--expose-internals', + '-e', + ` + const { isModuleEnabled } = require('internal/otel/core'); + console.log('fetch:' + isModuleEnabled('node:fetch')); + console.log('http:' + isModuleEnabled('node:http')); + `, + ], { + env: { + ...process.env, + NODE_OTEL_ENDPOINT: 'http://127.0.0.1:1', + NODE_OTEL_FILTER: 'node:fetch', + }, + }); + + assert.strictEqual(code, 0); + const lines = stdout.trim().split('\n'); + assert.match(lines.find((l) => l.startsWith('fetch:')), /fetch:true/); + assert.match(lines.find((l) => l.startsWith('http:')), /http:false/); + }); + + it('flushes when the buffer fills, without waiting for exit', async () => { + const collector = await startOTelCollector(); + + // The child traces its own requests. Each request produces a server + // and a client span, so a buffer size of 10 flushes long before the + // child exits (it never exits on its own; the parent kills it). + const script = ` + const http = require('node:http'); + const server = http.createServer((req, res) => { + res.writeHead(200); + res.end('ok'); + }); + server.listen(0, '127.0.0.1', () => { + const port = server.address().port; + let remaining = 12; + function next() { + if (remaining-- === 0) return; + http.get('http://127.0.0.1:' + port + '/fill', () => next()).on('error', () => {}); + } + next(); + }); + `; + + const child = spawn(process.execPath, ['--experimental-otel', '-e', script], { + stdio: 'pipe', + env: { + ...process.env, + NODE_OTEL_ENDPOINT: `http://127.0.0.1:${collector.port}`, + NODE_OTEL_MAX_BUFFER_SIZE: '10', + }, + }); + + // The payload must arrive while the child is still alive; if the child + // exits first, the flush only happened at exit and the test fails. + const spans = await Promise.race([ + collector.next().then(getSpans), + new Promise((_, reject) => { + child.on('exit', () => reject(new Error('child exited before flush'))); + }), + ]); + + child.kill(); + await collector.close(); + + assert.ok(spans.length > 0, 'Expected spans flushed by the full buffer'); + }); + + it('does not activate tracing in worker threads', async () => { + const { code, stdout } = await common.spawnPromisified(process.execPath, [ + '--experimental-otel', + '--expose-internals', + '-e', + ` + const { Worker } = require('node:worker_threads'); + const { isActive } = require('internal/otel/core'); + console.log('main:' + isActive()); + const worker = new Worker( + "const { isActive } = require('internal/otel/core');" + + "console.log('worker:' + isActive());", + { eval: true }); + `, + ], { + env: { + ...process.env, + NODE_OTEL_ENDPOINT: 'http://127.0.0.1:1', + }, + }); + + assert.strictEqual(code, 0); + const lines = stdout.trim().split('\n'); + assert.match(lines.find((l) => l.startsWith('main:')), /main:true/); + assert.match(lines.find((l) => l.startsWith('worker:')), /worker:false/); + }); + + it('does not activate tracing when no environment variables are set', async () => { + const { code, stdout, stderr } = await spawnOtel({}, ['--experimental-otel']); + + assert.strictEqual(code, 0); + assert.match(stdout, /active:false/); + assert.strictEqual(stderr, ''); + }); + + it('NODE_OTEL_FILTER limits instrumented modules', async () => { + const { code, stdout } = await common.spawnPromisified(process.execPath, [ + '--experimental-otel', + '--expose-internals', + '-e', + ` + const { isModuleEnabled } = require('internal/otel/core'); + console.log('http:' + isModuleEnabled('node:http')); + console.log('net:' + isModuleEnabled('node:net')); + console.log('dns:' + isModuleEnabled('node:dns')); + `, + ], { + env: { + ...process.env, + NODE_OTEL_ENDPOINT: 'http://127.0.0.1:1', + NODE_OTEL_FILTER: 'node:http,node:net', + }, + }); + + assert.strictEqual(code, 0); + const lines = stdout.trim().split('\n'); + assert.match(lines.find((l) => l.startsWith('http:')), /http:true/); + assert.match(lines.find((l) => l.startsWith('net:')), /net:true/); + assert.match(lines.find((l) => l.startsWith('dns:')), /dns:false/); + }); + + it('accepts NODE_OTEL_MAX_BUFFER_SIZE and NODE_OTEL_FLUSH_INTERVAL', async () => { + const { code, stdout } = await spawnOtel({ + NODE_OTEL_ENDPOINT: 'http://127.0.0.1:9999', + NODE_OTEL_MAX_BUFFER_SIZE: '250', + NODE_OTEL_FLUSH_INTERVAL: '2500', + }, ['--experimental-otel']); + + assert.strictEqual(code, 0); + assert.match(stdout, /active:true/); + }); + + it('warns and does not activate on an invalid NODE_OTEL_MAX_BUFFER_SIZE', async () => { + const { code, stdout, stderr } = await spawnOtel({ + NODE_OTEL_ENDPOINT: 'http://127.0.0.1:9999', + NODE_OTEL_MAX_BUFFER_SIZE: 'abc', + }, ['--experimental-otel']); + + assert.strictEqual(code, 0); + assert.match(stdout, /active:false/); + assert.match(stderr, /OTelWarning/); + assert.match(stderr, /Failed to initialize OpenTelemetry tracing/); + }); + + it('warns and does not activate on an invalid NODE_OTEL_FLUSH_INTERVAL', async () => { + const { code, stdout, stderr } = await spawnOtel({ + NODE_OTEL_ENDPOINT: 'http://127.0.0.1:9999', + NODE_OTEL_FLUSH_INTERVAL: '0', + }, ['--experimental-otel']); + + assert.strictEqual(code, 0); + assert.match(stdout, /active:false/); + assert.match(stderr, /OTelWarning/); + assert.match(stderr, /Failed to initialize OpenTelemetry tracing/); + }); + + it('emits a warning and continues when the endpoint is invalid', async () => { + const { code, stdout, stderr } = await common.spawnPromisified(process.execPath, [ + '--experimental-otel', + '-e', + ` + const http = require('node:http'); + const server = http.createServer((req, res) => { + res.writeHead(200); + res.end('ok'); + }); + server.listen(0, () => { + http.get('http://127.0.0.1:' + server.address().port, (res) => { + res.resume(); + res.on('end', () => { + server.close(); + console.log('still running'); + }); + }); + }); + `, + ], { + env: { + ...process.env, + NODE_OTEL_ENDPOINT: 'not-a-valid-url', + }, + }); + + assert.strictEqual(code, 0); + assert.match(stdout, /still running/); + assert.match(stderr, /OTelWarning/); + assert.match(stderr, /Failed to initialize OpenTelemetry tracing/); + }); +}); diff --git a/test/parallel/test-otel-exporter.js b/test/parallel/test-otel-exporter.js new file mode 100644 index 000000000000..c5634b773d75 --- /dev/null +++ b/test/parallel/test-otel-exporter.js @@ -0,0 +1,90 @@ +'use strict'; +const common = require('../common'); +const assert = require('node:assert'); +const http = require('node:http'); +const { describe, it } = require('node:test'); + +// This test verifies exporter error handling and process lifecycle behavior. +// Tracing is activated in the child processes via environment variables. + +describe('otel exporter behavior', () => { + it('does not crash when collector is unreachable', async () => { + const { code, stdout } = await common.spawnPromisified(process.execPath, [ + '--experimental-otel', + '-e', + ` + const http = require('node:http'); + + const server = http.createServer((req, res) => { + res.writeHead(200); + res.end('ok'); + }); + + server.listen(0, () => { + http.get('http://127.0.0.1:' + server.address().port, (res) => { + res.resume(); + res.on('end', () => { + server.close(); + console.log('success'); + }); + }); + }); + `, + ], { + env: { + ...process.env, + NODE_OTEL_ENDPOINT: 'http://127.0.0.1:1', + }, + }); + + assert.strictEqual(code, 0); + assert.match(stdout, /success/); + }); + + it('gives up on an unresponsive collector and still exits', async () => { + // A collector that accepts the connection but never responds. The + // exporter request must time out, and the child must still exit. + const collector = http.createServer(() => {}); + await new Promise((resolve) => collector.listen(0, '127.0.0.1', resolve)); + + try { + const { code, stderr } = await common.spawnPromisified( + process.execPath, + [ + '--experimental-otel', + '-e', + 'http.get("http://127.0.0.1:1/", () => {}).on("error", () => {});', + ], + { + env: { + ...process.env, + NODE_OTEL_ENDPOINT: `http://127.0.0.1:${collector.address().port}`, + }, + timeout: 30_000, + }, + ); + + assert.strictEqual(code, 0); + assert.match(stderr, /OTelExportWarning/); + } finally { + collector.close(); + } + }); + + it('flush timer does not keep process alive', async () => { + const { code, stdout } = await common.spawnPromisified(process.execPath, [ + '--experimental-otel', + '-e', + 'console.log("exiting");', + ], { + env: { + ...process.env, + NODE_OTEL_ENDPOINT: 'http://127.0.0.1:1', + }, + timeout: 10_000, + }); + + assert.strictEqual(code, 0); + assert.match(stdout, /exiting/); + }); +}); diff --git a/test/parallel/test-otel-filter.js b/test/parallel/test-otel-filter.js new file mode 100644 index 000000000000..b7bdd5ae942f --- /dev/null +++ b/test/parallel/test-otel-filter.js @@ -0,0 +1,43 @@ +'use strict'; +// Flags: --expose-internals + +require('../common'); +const assert = require('node:assert'); +const { describe, it } = require('node:test'); + +const otel = require('internal/otel/core'); +const { flush } = require('internal/otel/flush'); +const { + startOTelCollector, + startHTTPServer, + httpGet, + delay, +} = require('../common/otel'); + +describe('otel filter', () => { + it('does not create spans for filtered-out modules', async () => { + const collector = await startOTelCollector(); + + // Filter to a module that will not be exercised. + otel.start({ + endpoint: `http://127.0.0.1:${collector.port}`, + filter: 'node:dns', + }); + + const server = await startHTTPServer(); + + await httpGet(server.port, '/'); + + // Flush explicitly — if any HTTP spans were created despite the filter, + // they would be sent to the collector. + flush(); + + // Wait for any potential request to the collector. + await delay(100); + + assert.strictEqual(collector.requests.length, 0); + + await collector.close(); + await server.close(); + }); +}); diff --git a/test/parallel/test-otel-flush-coverage.js b/test/parallel/test-otel-flush-coverage.js new file mode 100644 index 000000000000..b94a45ab02bb --- /dev/null +++ b/test/parallel/test-otel-flush-coverage.js @@ -0,0 +1,118 @@ +'use strict'; +// Flags: --expose-internals + +require('../common'); +const assert = require('node:assert'); +const { after, afterEach, before, describe, it } = require('node:test'); + +const otel = require('internal/otel/core'); +const { + Span, + SPAN_KIND_INTERNAL, +} = require('internal/otel/span'); +const { + addSpan, + flush, + resetCaches, +} = require('internal/otel/flush'); +const { startOTelCollector, delay } = require('../common/otel'); + +describe('flush.js coverage', () => { + let collector; + + before(async () => { + collector = await startOTelCollector(); + otel.start({ endpoint: `http://127.0.0.1:${collector.port}` }); + }); + + after(() => collector.close()); + + afterEach(() => { + // Clear the buffer and the export failure/warning counters so that + // tests do not interfere with each other. + resetCaches(); + }); + + it('flush is a no-op when buffer is empty', async () => { + flush(); + + // Wait for any potential HTTP request to arrive. + await delay(100); + + assert.strictEqual(collector.requests.length, 0); + }); + + it('flush handles spanToOtlp errors gracefully', async () => { + const warningPromise = new Promise((resolve) => { + process.on('warning', function onWarning(w) { + if (w.name === 'OTelExportWarning') { + process.removeListener('warning', onWarning); + resolve(w); + } + }); + }); + + // Create a poisoned span-like object that will throw during serialization. + const badSpan = { + name: 'bad-span', + getAttributes() { throw new Error('serialize boom'); }, + }; + + addSpan(badSpan); + flush(); + + const warning = await warningPromise; + + assert.ok(warning.message.includes('serialize boom')); + }); + + it('flush handles JSONStringify errors gracefully', async () => { + const warningPromise = new Promise((resolve) => { + process.on('warning', function onWarning(w) { + if (w.name === 'OTelExportWarning') { + process.removeListener('warning', onWarning); + resolve(w); + } + }); + }); + + // Create a fake span-like object that spanToOtlp can process, but whose + // name is a BigInt: it is copied into the OTLP structure untouched and + // makes JSONStringify throw. + const fakeSpan = { + name: 1n, + getAttributes() { return {}; }, + getEvents() { return []; }, + traceId: 'a'.repeat(32), + spanId: 'b'.repeat(16), + parentSpanId: null, + kind: 0, + startTime: 1, + endTime: 2, + status: { code: 0, message: '' }, + traceFlags: 0x01, + }; + + addSpan(fakeSpan); + flush(); + + const warning = await warningPromise; + + assert.ok(warning.message.includes('Failed to serialize')); + }); + + it('resetCaches clears the span buffer', async () => { + // Add a span so the buffer is non-empty. + const span = new Span('test', SPAN_KIND_INTERNAL); + span.end(); + + resetCaches(); + + flush(); + + // Wait for any potential HTTP request to arrive. + await delay(100); + + assert.strictEqual(collector.requests.length, 0); + }); +}); diff --git a/test/parallel/test-otel-http-client.js b/test/parallel/test-otel-http-client.js new file mode 100644 index 000000000000..c8bcf5ae4483 --- /dev/null +++ b/test/parallel/test-otel-http-client.js @@ -0,0 +1,120 @@ +'use strict'; +// Flags: --expose-internals + +require('../common'); +const assert = require('node:assert'); +const { after, before, describe, it } = require('node:test'); + +const otel = require('internal/otel/core'); +const { flush } = require('internal/otel/flush'); +const { + startOTelCollector, + startHTTPServer, + getSpans, + httpGet, + spanAttrs, +} = require('../common/otel'); + +describe('otel HTTP client spans', () => { + let collector; + + before(async () => { + collector = await startOTelCollector(); + otel.start({ endpoint: `http://127.0.0.1:${collector.port}` }); + }); + + after(() => collector.close()); + + it('creates client spans and injects traceparent', async () => { + let receivedTraceparent = null; + + // Target server records incoming traceparent header. + const target = await startHTTPServer((req, res) => { + receivedTraceparent = req.headers.traceparent || null; + res.writeHead(200); + res.end('ok'); + }); + + await httpGet(target.port, '/api/data'); + flush(); + const spans = getSpans(await collector.next()); + + await target.close(); + + const clientSpan = spans.find((s) => s.kind === 3); // SPAN_KIND_CLIENT + assert.ok(clientSpan, `Expected a client span in: ${JSON.stringify(spans)}`); + assert.strictEqual(clientSpan.name, 'GET'); + + const attrs = spanAttrs(clientSpan); + assert.strictEqual(attrs['http.request.method'], 'GET'); + assert.ok(attrs['url.full'], 'Expected url.full attribute'); + assert.strictEqual(attrs['http.response.status_code'], '200'); + + assert.ok(receivedTraceparent, 'Expected non-empty traceparent'); + assert.match(receivedTraceparent, /^00-[0-9a-f]{32}-[0-9a-f]{16}-0[01]$/); + + // The traceparent should reference the client span's IDs. + assert.ok(receivedTraceparent.includes(clientSpan.traceId), + 'traceparent should contain the client span traceId'); + assert.ok(receivedTraceparent.includes(clientSpan.spanId), + 'traceparent should contain the client span spanId'); + }); + + // Per OTel semantic conventions, client spans set STATUS_ERROR for >= 400, + // while server spans only set it for >= 500 (4xx is a client mistake). + it('sets error on client span but not server span for 404', async () => { + const server = await startHTTPServer((req, res) => { + res.writeHead(404); + res.end('not found'); + }); + + await httpGet(server.port, '/missing'); + flush(); + const spans = getSpans(await collector.next()); + + await server.close(); + + // Find server span (kind SERVER = 1, OTLP wire = 2) and client span + // (kind CLIENT = 2, OTLP wire = 3). + const serverSpan = spans.find((s) => { + return s.kind === 2 && s.attributes?.some( + (a) => a.key === 'url.path' && a.value.stringValue === '/missing'); + }); + const clientSpan = spans.find((s) => s.kind === 3); + + assert.ok(serverSpan, 'Expected a server span'); + assert.ok(clientSpan, 'Expected a client span'); + + // Server span: 404 should NOT have error status. + assert.strictEqual(serverSpan.status, undefined); + + // Client span: 404 should have error status. + assert.ok(clientSpan.status, 'Client span should have error status'); + assert.strictEqual(clientSpan.status.code, 2); // STATUS_ERROR + }); + + it('creates error span for HTTP client connection refused', async () => { + // Make a request to a port that is not listening — should fail. + await httpGet(1, '/fail').catch(() => {}); + + flush(); + const spans = getSpans(await collector.next()); + + const errorSpan = spans.find((s) => { + if (s.kind !== 3) return false; // CLIENT + const urlAttr = s.attributes?.find((a) => a.key === 'url.full'); + return urlAttr?.value?.stringValue?.includes('/fail'); + }); + + assert.ok(errorSpan, 'Expected an error client span for connection refused'); + + assert.ok(errorSpan.status, 'Error span should have status'); + assert.strictEqual(errorSpan.status.code, 2); // STATUS_ERROR + + assert.ok(errorSpan.events, 'Error span should have events'); + const exceptionEvent = errorSpan.events.find( + (e) => e.name === 'exception', + ); + assert.ok(exceptionEvent, 'Should have exception event'); + }); +}); diff --git a/test/parallel/test-otel-http-server.js b/test/parallel/test-otel-http-server.js new file mode 100644 index 000000000000..dc2885cf16aa --- /dev/null +++ b/test/parallel/test-otel-http-server.js @@ -0,0 +1,52 @@ +'use strict'; +// Flags: --expose-internals + +require('../common'); +const assert = require('node:assert'); +const { describe, it } = require('node:test'); + +const otel = require('internal/otel/core'); +const { flush } = require('internal/otel/flush'); +const { + startOTelCollector, + startHTTPServer, + getSpans, + httpGet, + spanAttrs, +} = require('../common/otel'); + +describe('otel HTTP server spans', () => { + it('creates server spans for incoming HTTP requests', async () => { + const collector = await startOTelCollector(); + + otel.start({ endpoint: `http://127.0.0.1:${collector.port}` }); + + const app = await startHTTPServer(); + + await httpGet(app.port, '/test-path'); + + flush(); + const spans = getSpans(await collector.next()); + + await collector.close(); + await app.close(); + + assert.ok(spans.length >= 1, `Expected at least 1 span, got ${spans.length}`); + + const serverSpan = spans.find((s) => s.kind === 2); // SPAN_KIND_SERVER + assert.ok(serverSpan, 'Expected a server span'); + assert.strictEqual(serverSpan.name, 'GET'); + + assert.ok(serverSpan.traceId); + assert.strictEqual(serverSpan.traceId.length, 32); + assert.ok(serverSpan.spanId); + assert.strictEqual(serverSpan.spanId.length, 16); + assert.ok(serverSpan.startTimeUnixNano); + assert.ok(serverSpan.endTimeUnixNano); + + const attrs = spanAttrs(serverSpan); + assert.strictEqual(attrs['http.request.method'], 'GET'); + assert.strictEqual(attrs['url.path'], '/test-path'); + assert.strictEqual(attrs['http.response.status_code'], '200'); + }); +}); diff --git a/test/parallel/test-otel-https.js b/test/parallel/test-otel-https.js new file mode 100644 index 000000000000..1891b9a92f35 --- /dev/null +++ b/test/parallel/test-otel-https.js @@ -0,0 +1,64 @@ +'use strict'; +// Flags: --expose-internals + +const common = require('../common'); +const assert = require('node:assert'); +const https = require('node:https'); +const fixtures = require('../common/fixtures'); +const { describe, it } = require('node:test'); + +const otel = require('internal/otel/core'); +const { flush } = require('internal/otel/flush'); +const { + startOTelCollector, + getSpans, + spanAttrs, +} = require('../common/otel'); + +describe('otel HTTPS server spans', () => { + it('marks the url.scheme attribute as https for TLS requests', async () => { + if (!common.hasCrypto) { + common.skip('missing crypto'); + return; + } + + const collector = await startOTelCollector(); + + otel.start({ endpoint: `http://127.0.0.1:${collector.port}` }); + + const server = https.createServer({ + key: fixtures.readKey('agent1-key.pem'), + cert: fixtures.readKey('agent1-cert.pem'), + }, (req, res) => { + res.writeHead(200); + res.end('ok'); + }); + await new Promise((resolve) => server.listen(0, '127.0.0.1', resolve)); + + await new Promise((resolve, reject) => { + https.get({ + host: '127.0.0.1', + port: server.address().port, + path: '/secure', + rejectUnauthorized: false, + agent: false, + }, (res) => { + res.resume(); + res.on('end', resolve); + }).on('error', reject); + }); + + flush(); + const spans = getSpans(await collector.next()); + + await collector.close(); + await new Promise((resolve) => server.close(resolve)); + + const serverSpan = spans.find((s) => s.kind === 2); // SPAN_KIND_SERVER + assert.ok(serverSpan, 'Expected a server span'); + + const attrs = spanAttrs(serverSpan); + assert.strictEqual(attrs['url.scheme'], 'https'); + assert.strictEqual(attrs['url.path'], '/secure'); + }); +}); diff --git a/test/parallel/test-otel-id-refill.js b/test/parallel/test-otel-id-refill.js new file mode 100644 index 000000000000..ed9fe40b07ca --- /dev/null +++ b/test/parallel/test-otel-id-refill.js @@ -0,0 +1,30 @@ +'use strict'; +// Flags: --expose-internals + +require('../common'); +const assert = require('node:assert'); +const { describe, it } = require('node:test'); + +const { generateTraceId, generateSpanId } = require('internal/otel/id'); + +describe('otel ID generation buffer refill', () => { + // The random buffer is 4096 bytes. Each trace ID uses 16 bytes and each + // span ID 8 bytes, so generating enough IDs forces at least one refill. + const cases = [ + ['trace', generateTraceId, 32, 300], + ['span', generateSpanId, 16, 600], + ]; + + for (const [name, generate, idLength, count] of cases) { + it(`generates valid ${name} IDs after exhausting the random buffer`, () => { + const ids = new Set(); + for (let i = 0; i < count; i++) { + const id = generate(); + assert.strictEqual(id.length, idLength); + assert.match(id, new RegExp(`^[0-9a-f]{${idLength}}$`)); + ids.add(id); + } + assert.strictEqual(ids.size, count); + }); + } +}); diff --git a/test/parallel/test-otel-otlp-compliance.js b/test/parallel/test-otel-otlp-compliance.js new file mode 100644 index 000000000000..7762fcfaa4d4 --- /dev/null +++ b/test/parallel/test-otel-otlp-compliance.js @@ -0,0 +1,283 @@ +'use strict'; +// Flags: --expose-internals + +// This test verifies that the OTLP/HTTP JSON export payload is compliant +// with the OpenTelemetry specification (opentelemetry-proto). + +require('../common'); +const assert = require('node:assert'); +const { after, before, describe, it } = require('node:test'); + +const otel = require('internal/otel/core'); +const { flush } = require('internal/otel/flush'); +const { + startOTelCollector, + startHTTPServer, + getSpans, + httpGet, +} = require('../common/otel'); + +// Verify that an attribute is a KeyValue holding exactly one AnyValue field. +function assertAnyValue(attr) { + assert.ok(typeof attr.key === 'string', 'attribute key must be a string'); + assert.ok(attr.value !== undefined, 'attribute value must be present'); + + const fields = Object.keys(attr.value); + assert.strictEqual(fields.length, 1, + `attribute value must have exactly one field, got: ${fields}`); + + const field = fields[0]; + assert.ok( + ['stringValue', 'boolValue', 'intValue', + 'doubleValue', 'arrayValue', 'kvlistValue', + 'bytesValue'].includes(field), + `unexpected AnyValue field: ${field}`, + ); + + // intValue must be a decimal string (int64 JSON encoding). + if (field === 'intValue') { + assert.ok(typeof attr.value.intValue === 'string', + 'intValue must be a string (int64 JSON encoding)'); + assert.match(attr.value.intValue, /^-?\d+$/); + } +} + +describe('OTLP/JSON spec compliance', () => { + let collector; + + before(async () => { + collector = await startOTelCollector(); + otel.start({ endpoint: `http://127.0.0.1:${collector.port}` }); + }); + + after(() => collector.close()); + + it('produces a spec-compliant ExportTraceServiceRequest', async () => { + const server = await startHTTPServer(); + + await httpGet(server.port, '/test'); + flush(); + const payload = await collector.next(); + const lastRequest = collector.requests.at(-1); + + await server.close(); + + // === Content-Type === + assert.strictEqual(lastRequest.headers['content-type'], 'application/json'); + + // === Top-level envelope === + assert.ok(Array.isArray(payload.resourceSpans), + 'resourceSpans must be an array'); + assert.strictEqual(payload.resourceSpans.length, 1); + + const resourceSpan = payload.resourceSpans[0]; + + // === Resource === + assert.ok(resourceSpan.resource, 'resource must be present'); + assert.ok(Array.isArray(resourceSpan.resource.attributes), + 'resource.attributes must be an array'); + + // Verify resource attributes are KeyValue format. + for (const attr of resourceSpan.resource.attributes) { + assertAnyValue(attr); + } + + // Verify required resource attributes. + const resourceAttrs = {}; + for (const a of resourceSpan.resource.attributes) { + resourceAttrs[a.key] = a.value; + } + assert.ok(resourceAttrs['service.name'], + 'resource must have service.name'); + assert.ok(resourceAttrs['service.name'].stringValue, + 'service.name must be a string value'); + + // === ScopeSpans === + assert.ok(Array.isArray(resourceSpan.scopeSpans), + 'scopeSpans must be an array'); + assert.strictEqual(resourceSpan.scopeSpans.length, 1); + + const scopeSpan = resourceSpan.scopeSpans[0]; + + // === InstrumentationScope === + assert.ok(scopeSpan.scope, 'scope must be present'); + assert.ok(typeof scopeSpan.scope.name === 'string', + 'scope.name must be a string'); + assert.ok(typeof scopeSpan.scope.version === 'string', + 'scope.version must be a string'); + + // === Spans === + assert.ok(Array.isArray(scopeSpan.spans), 'spans must be an array'); + assert.ok(scopeSpan.spans.length >= 1, 'must have at least 1 span'); + + for (const span of scopeSpan.spans) { + // --- traceId: 32-char lowercase hex string --- + assert.ok(typeof span.traceId === 'string', 'traceId must be a string'); + assert.strictEqual(span.traceId.length, 32); + assert.match(span.traceId, /^[0-9a-f]{32}$/, + 'traceId must be lowercase hex'); + // Must not be all zeros. + assert.notStrictEqual(span.traceId, '0'.repeat(32)); + + // --- spanId: 16-char lowercase hex string --- + assert.ok(typeof span.spanId === 'string', 'spanId must be a string'); + assert.strictEqual(span.spanId.length, 16); + assert.match(span.spanId, /^[0-9a-f]{16}$/, + 'spanId must be lowercase hex'); + assert.notStrictEqual(span.spanId, '0'.repeat(16)); + + // --- name: non-empty string --- + assert.ok(typeof span.name === 'string', 'name must be a string'); + assert.ok(span.name.length > 0, 'name must be non-empty'); + + // --- kind: integer 1-5 (OTLP SpanKind enum, no 0/UNSPECIFIED) --- + assert.ok(typeof span.kind === 'number', 'kind must be a number'); + assert.ok(Number.isInteger(span.kind), 'kind must be an integer'); + assert.ok(span.kind >= 1 && span.kind <= 5, + `kind must be 1-5, got: ${span.kind}`); + + // --- timestamps: decimal strings of nanoseconds --- + assert.ok(typeof span.startTimeUnixNano === 'string', + 'startTimeUnixNano must be a string'); + assert.match(span.startTimeUnixNano, /^\d+$/, + 'startTimeUnixNano must be a decimal string'); + assert.ok(typeof span.endTimeUnixNano === 'string', + 'endTimeUnixNano must be a string'); + assert.match(span.endTimeUnixNano, /^\d+$/, + 'endTimeUnixNano must be a decimal string'); + + // endTime >= startTime. + assert.ok(BigInt(span.endTimeUnixNano) >= BigInt(span.startTimeUnixNano), + 'endTimeUnixNano must be >= startTimeUnixNano'); + + // Timestamps should be plausible (after 2020, before 2100). + const startSec = Number(BigInt(span.startTimeUnixNano) / 1_000_000_000n); + assert.ok(startSec > 1577836800, 'timestamp too old'); // 2020-01-01 + assert.ok(startSec < 4102444800, 'timestamp too far in future'); // 2100-01-01 + + // --- parentSpanId: omitted for root, 16-char hex for child --- + if (span.parentSpanId !== undefined) { + assert.ok(typeof span.parentSpanId === 'string'); + assert.strictEqual(span.parentSpanId.length, 16); + assert.match(span.parentSpanId, /^[0-9a-f]{16}$/); + } + + // --- attributes: array of KeyValue (or omitted if empty) --- + if (span.attributes !== undefined) { + assert.ok(Array.isArray(span.attributes)); + for (const attr of span.attributes) { + assertAnyValue(attr); + } + } + + // --- status: omitted when unset, or {code, message} --- + if (span.status !== undefined) { + assert.ok(typeof span.status === 'object'); + assert.ok(typeof span.status.code === 'number'); + assert.ok(Number.isInteger(span.status.code)); + assert.ok(span.status.code >= 0 && span.status.code <= 2, + `status.code must be 0-2, got: ${span.status.code}`); + // If message is present, it must be a string. + if (span.status.message !== undefined) { + assert.ok(typeof span.status.message === 'string'); + } + } + + // --- events: omitted if empty, or array of Event --- + if (span.events !== undefined) { + assert.ok(Array.isArray(span.events)); + for (const event of span.events) { + assert.ok(typeof event.name === 'string'); + assert.ok(event.name.length > 0, 'event name must be non-empty'); + assert.ok(typeof event.timeUnixNano === 'string'); + assert.match(event.timeUnixNano, /^\d+$/); + if (event.attributes !== undefined) { + assert.ok(Array.isArray(event.attributes)); + } + } + } + + // --- No snake_case field names --- + for (const key of Object.keys(span)) { + assert.ok(!key.includes('_') || key === 'startTimeUnixNano' || + key === 'endTimeUnixNano' || key === 'timeUnixNano', + `unexpected snake_case-style field: ${key}`); + } + } + }); + + it('omits status when unset (200 OK response)', async () => { + const server = await startHTTPServer(); + + await httpGet(server.port, '/ok'); + flush(); + const payload = await collector.next(); + + await server.close(); + + const serverSpan = getSpans(payload).find((s) => { + return s.kind === 2 && s.attributes?.some( + (a) => a.key === 'url.path' && a.value.stringValue === '/ok'); + }); + assert.ok(serverSpan); + // For a 200 OK server span, status should be omitted (STATUS_UNSET). + assert.strictEqual(serverSpan.status, undefined); + }); + + it('includes status with code 2 for error responses', async () => { + const server = await startHTTPServer((req, res) => { + res.writeHead(500); + res.end('error'); + }); + + await httpGet(server.port, '/fail'); + flush(); + const payload = await collector.next(); + + await server.close(); + + const serverSpan = getSpans(payload).find((s) => { + return s.kind === 2 && s.attributes?.some( + (a) => a.key === 'url.path' && a.value.stringValue === '/fail'); + }); + assert.ok(serverSpan); + assert.ok(serverSpan.status, 'status must be present for error'); + assert.strictEqual(serverSpan.status.code, 2); // STATUS_CODE_ERROR + assert.ok(typeof serverSpan.status.message === 'string'); + }); + + it('omits empty attributes array', async () => { + const server = await startHTTPServer(); + + await httpGet(server.port, '/test'); + flush(); + const payload = await collector.next(); + + await server.close(); + + for (const span of getSpans(payload)) { + // If attributes is present, it must be non-empty. + if (span.attributes !== undefined) { + assert.ok(span.attributes.length > 0, + 'attributes must not be an empty array'); + } + // Events should not be present unless there are actual events. + if (span.events !== undefined) { + assert.ok(span.events.length > 0, + 'events must not be an empty array'); + } + } + }); + + it('posts to /v1/traces endpoint', async () => { + const server = await startHTTPServer(); + + await httpGet(server.port, '/x'); + flush(); + await collector.next(); + + await server.close(); + + assert.strictEqual(collector.requests.at(-1).url, '/v1/traces'); + }); +}); diff --git a/test/parallel/test-otel-self-trace.js b/test/parallel/test-otel-self-trace.js new file mode 100644 index 000000000000..1c3d2717e47a --- /dev/null +++ b/test/parallel/test-otel-self-trace.js @@ -0,0 +1,79 @@ +'use strict'; +// Flags: --expose-internals + +require('../common'); +const assert = require('node:assert'); +const { describe, it } = require('node:test'); + +const otel = require('internal/otel/core'); +const { flush } = require('internal/otel/flush'); +const { + startOTelCollector, + startHTTPServer, + getSpans, + httpGet, + delay, +} = require('../common/otel'); + +describe('otel self-trace prevention', () => { + it('does not create spans for export requests to the collector', async () => { + const collector = await startOTelCollector(); + + otel.start({ endpoint: `http://127.0.0.1:${collector.port}` }); + + const server = await startHTTPServer(); + + await httpGet(server.port, '/test'); + + // Flush while instrumentation is still active. This sends an HTTP + // request to the collector, which fires http.client.request.created. + // The self-trace check (getCollectorHost()) should prevent creating + // a client span for this export request. + flush(); + await collector.next(); + + // Allow time for any inadvertent export-request spans to be buffered, + // then flush again so that they would be exported and observable. + await delay(100); + flush(); + await delay(100); + + await collector.close(); + await server.close(); + + const allSpans = collector.payloads.flatMap(getSpans); + + // The export requests themselves must not produce server spans: a + // collector hosted in this process must not be traced by its own + // export traffic, or it would feed itself indefinitely. + for (const span of allSpans) { + const pathAttr = span.attributes?.find((a) => a.key === 'url.path'); + if (pathAttr) { + assert.notStrictEqual(pathAttr.value.stringValue, '/v1/traces'); + } + } + + // No span should target the collector. + for (const span of allSpans) { + const urlAttr = span.attributes?.find((a) => a.key === 'url.full'); + if (urlAttr) { + assert.ok( + !urlAttr.value.stringValue.includes(`:${collector.port}`), + `Span should not target collector: ${urlAttr.value.stringValue}`, + ); + } + } + + // Should have exactly the user request spans (server + client), + // not any client spans for the export requests. + const clientSpans = allSpans.filter((s) => s.kind === 3); // CLIENT + for (const cs of clientSpans) { + const urlAttr = cs.attributes?.find((a) => a.key === 'url.full'); + assert.ok(urlAttr, 'Client span should have url.full'); + assert.ok( + urlAttr.value.stringValue.includes('/test'), + `Client span should be for /test, got: ${urlAttr.value.stringValue}`, + ); + } + }); +}); diff --git a/test/parallel/test-otel-server-close.js b/test/parallel/test-otel-server-close.js new file mode 100644 index 000000000000..89f42ed67a15 --- /dev/null +++ b/test/parallel/test-otel-server-close.js @@ -0,0 +1,65 @@ +'use strict'; +// Flags: --expose-internals + +require('../common'); +const assert = require('node:assert'); +const http = require('node:http'); +const net = require('node:net'); +const { describe, it } = require('node:test'); + +const otel = require('internal/otel/core'); +const { flush } = require('internal/otel/flush'); +const { + startOTelCollector, + startHTTPServer, + getSpans, +} = require('../common/otel'); + +describe('otel server request close handling', () => { + it('ends span with error when client disconnects before response', async () => { + const collector = await startOTelCollector(); + + otel.start({ endpoint: `http://127.0.0.1:${collector.port}` }); + + // Create a server that delays its response for every path except + // /flush, which responds normally and is used to trigger an export. + const server = await startHTTPServer((req, res) => { + if (req.url === '/flush') { + res.writeHead(200); + res.end('ok'); + return; + } + // Don't respond — the client will disconnect first. + req.on('close', () => { + // After the client disconnects, make a normal request and flush. + http.get(`http://127.0.0.1:${server.port}/flush`, (r) => { + r.resume(); + r.on('end', () => flush()); + }); + }); + }); + + // Use a raw TCP socket to connect and send a partial HTTP request, + // then destroy the connection before the server responds. + const socket = net.connect(server.port, '127.0.0.1', () => { + socket.write('GET /disconnect HTTP/1.1\r\nHost: 127.0.0.1\r\n\r\n'); + // Destroy immediately — server never gets to respond. + setTimeout(() => socket.destroy(), 50); + }); + + const spans = getSpans(await collector.next()); + + await collector.close(); + await server.close(); + + const disconnectSpan = spans.find((s) => { + if (s.kind !== 2) return false; // SERVER + const pathAttr = s.attributes?.find((a) => a.key === 'url.path'); + return pathAttr?.value?.stringValue === '/disconnect'; + }); + + assert.ok(disconnectSpan, 'Expected a server span for /disconnect'); + assert.ok(disconnectSpan.status, 'Span should have error status'); + assert.strictEqual(disconnectSpan.status.code, 2); // STATUS_ERROR + }); +}); diff --git a/test/parallel/test-otel-span-coverage.js b/test/parallel/test-otel-span-coverage.js new file mode 100644 index 000000000000..6f59147c1f72 --- /dev/null +++ b/test/parallel/test-otel-span-coverage.js @@ -0,0 +1,81 @@ +'use strict'; +// Flags: --expose-internals + +require('../common'); +const assert = require('node:assert'); +const { after, before, describe, it } = require('node:test'); + +const otel = require('internal/otel/core'); +const { + Span, + SPAN_KIND_INTERNAL, + STATUS_UNSET, +} = require('internal/otel/span'); +const { flush } = require('internal/otel/flush'); +const { + startOTelCollector, + getSpans, +} = require('../common/otel'); + +describe('Span internals coverage', () => { + let collector; + + before(async () => { + collector = await startOTelCollector(); + otel.start({ endpoint: `http://127.0.0.1:${collector.port}` }); + }); + + after(() => collector.close()); + + it('unsampled span does not call addSpan', () => { + // Create a parent with traceFlags=0x00 (not sampled). + const parent = { + __proto__: null, + traceId: 'a'.repeat(32), + spanId: 'b'.repeat(16), + traceFlags: 0x00, + }; + + const span = new Span('unsampled', SPAN_KIND_INTERNAL, { parent }); + assert.strictEqual(span.traceFlags, 0x00); + + // end() should not throw even though addSpan won't buffer it. + span.end(); + assert.ok(span.endTime !== undefined); + }); + + it('addEvent without attributes uses empty object', () => { + const span = new Span('test', SPAN_KIND_INTERNAL); + span.addEvent('my-event'); + + const events = span.getEvents(); + assert.strictEqual(events.length, 1); + assert.strictEqual(events[0].name, 'my-event'); + assert.ok(events[0].attributes); + }); + + it('status defaults to UNSET', () => { + const span = new Span('test', SPAN_KIND_INTERNAL); + assert.strictEqual(span.status.code, STATUS_UNSET); + assert.strictEqual(span.status.message, ''); + }); + + it('a second end() call is a no-op and does not double-export', async () => { + const span = new Span('double-end', SPAN_KIND_INTERNAL); + assert.strictEqual(span.endTime, undefined); + + span.end(); + const firstEndTime = span.endTime; + assert.ok(firstEndTime !== undefined); + + span.end(); + assert.strictEqual(span.endTime, firstEndTime); + + flush(); + + const spans = getSpans(await collector.next()); + + const matches = spans.filter((s) => s.name === 'double-end'); + assert.strictEqual(matches.length, 1); + }); +}); diff --git a/test/parallel/test-otel-traceparent-validation.js b/test/parallel/test-otel-traceparent-validation.js new file mode 100644 index 000000000000..05b7c3e71ae5 --- /dev/null +++ b/test/parallel/test-otel-traceparent-validation.js @@ -0,0 +1,112 @@ +'use strict'; +// Flags: --expose-internals + +require('../common'); +const assert = require('node:assert'); +const net = require('node:net'); +const { after, before, describe, it } = require('node:test'); + +const otel = require('internal/otel/core'); +const { flush } = require('internal/otel/flush'); +const { + startOTelCollector, + startHTTPServer, + getSpans, +} = require('../common/otel'); + +// Sends a raw HTTP request with the given traceparent header to the server. +function requestWithTraceparent(port, traceparent) { + return new Promise((resolve) => { + const socket = net.connect(port, '127.0.0.1', () => { + socket.write( + `GET /test HTTP/1.1\r\n` + + `Host: 127.0.0.1:${port}\r\n` + + `traceparent: ${traceparent}\r\n` + + `Connection: close\r\n` + + `\r\n` + ); + socket.on('data', () => {}); + socket.on('end', resolve); + }); + }); +} + +describe('otel traceparent validation', () => { + let collector; + + before(async () => { + collector = await startOTelCollector(); + otel.start({ endpoint: `http://127.0.0.1:${collector.port}` }); + }); + + after(() => collector.close()); + + it('rejects invalid traceparent headers (non-hex chars)', async () => { + const server = await startHTTPServer(); + + await requestWithTraceparent(server.port, + '00-ZZZZZZZZZZZZZZZZZZZZZZZZZZZZZZZZ-' + + 'ZZZZZZZZZZZZZZZZ-01'); + + flush(); + const spans = getSpans(await collector.next()); + + await server.close(); + + const serverSpan = spans.find( + (s) => s.kind === 2 && s.attributes?.some( + (a) => a.key === 'url.path' && a.value.stringValue === '/test')); + assert.ok(serverSpan, 'Expected a server span'); + assert.match(serverSpan.traceId, /^[0-9a-f]{32}$/); + assert.notStrictEqual(serverSpan.traceId, + 'zzzzzzzzzzzzzzzzzzzzzzzzzzzzzzzz'); + // No parentSpanId since the invalid traceparent was rejected. + assert.ok(!serverSpan.parentSpanId, + 'Should not have parentSpanId from invalid traceparent'); + }); + + it('rejects all-zero traceId in traceparent', async () => { + const server = await startHTTPServer(); + + // All-zero traceId is invalid per W3C spec. + await requestWithTraceparent(server.port, + '00-00000000000000000000000000000000-' + + '0123456789abcdef-01'); + + flush(); + const spans = getSpans(await collector.next()); + + await server.close(); + + const serverSpan = spans.find( + (s) => s.kind === 2 && s.attributes?.some( + (a) => a.key === 'url.path' && a.value.stringValue === '/test')); + assert.ok(serverSpan); + assert.notStrictEqual(serverSpan.traceId, + '00000000000000000000000000000000'); + assert.ok(!serverSpan.parentSpanId); + }); + + it('rejects traceparent versions other than 00', async () => { + const server = await startHTTPServer(); + + // Version 01 is not understood; a fresh trace must be started. + await requestWithTraceparent(server.port, + '01-abcdef0123456789abcdef0123456789-' + + '0123456789abcdef-01'); + + flush(); + const spans = getSpans(await collector.next()); + + await server.close(); + + const serverSpan = spans.find( + (s) => s.kind === 2 && s.attributes?.some( + (a) => a.key === 'url.path' && a.value.stringValue === '/test')); + assert.ok(serverSpan); + assert.notStrictEqual(serverSpan.traceId, + 'abcdef0123456789abcdef0123456789'); + assert.ok(!serverSpan.parentSpanId, + 'Should not have a parent from an unknown version'); + }); +}); diff --git a/test/parallel/test-otel-undici.js b/test/parallel/test-otel-undici.js new file mode 100644 index 000000000000..a844e3c400a3 --- /dev/null +++ b/test/parallel/test-otel-undici.js @@ -0,0 +1,80 @@ +'use strict'; +// Flags: --expose-internals + +require('../common'); +const assert = require('node:assert'); +const { describe, it } = require('node:test'); + +const otel = require('internal/otel/core'); +const { flush } = require('internal/otel/flush'); +const { + startOTelCollector, + startHTTPServer, + getSpans, + spanAttrs, +} = require('../common/otel'); + +describe('otel undici spans', () => { + it('creates client spans for fetch requests and overrides existing ' + + 'trace context headers', async () => { + const collector = await startOTelCollector(); + + otel.start({ + endpoint: `http://127.0.0.1:${collector.port}`, + filter: 'node:undici,node:fetch', + }); + + const receivedTraceparent = {}; + const target = await startHTTPServer((req, res) => { + receivedTraceparent[req.url] = req.headers.traceparent; + res.writeHead(200); + res.end('ok'); + }); + + try { + const res = await fetch(`http://127.0.0.1:${target.port}/fetch-test`); + await res.text(); + } catch { + // Ignore fetch errors. + } + + // A caller-supplied traceparent must be overridden, not duplicated: + // undici's addHeader appends, so a duplicate header would arrive at + // the server as a single comma-separated value. + try { + const res = await fetch(`http://127.0.0.1:${target.port}/override-test`, { + headers: { + traceparent: '00-00000000000000000000000000000000-' + + '0000000000000000-00', + }, + }); + await res.text(); + } catch { + // Ignore fetch errors. + } + + flush(); + const spans = getSpans(await collector.next()); + + await collector.close(); + await target.close(); + + assert.ok(spans.length > 0, 'Expected at least one span from fetch'); + const clientSpan = spans.find((s) => { + return s.kind === 3 && s.attributes?.some( + (a) => a.key === 'url.full' && + a.value.stringValue.includes('/fetch-test')); + }); + assert.ok(clientSpan, 'Expected a CLIENT span from fetch'); + + const attrs = spanAttrs(clientSpan); + assert.strictEqual(attrs['http.request.method'], 'GET'); + assert.ok(attrs['url.full'], 'Expected url.full attribute'); + + const overridden = receivedTraceparent['/override-test']; + assert.ok(overridden, 'Expected a traceparent on the override request'); + assert.match(overridden, /^00-[0-9a-f]{32}-[0-9a-f]{16}-0[01]$/); + assert.ok(!overridden.includes(','), + 'traceparent must be overridden, not duplicated'); + }); +});