diff --git a/hypaware-core/plugins-workspace/format-jsonl/src/index.js b/hypaware-core/plugins-workspace/format-jsonl/src/index.js index 5cde1887..397aafd7 100644 --- a/hypaware-core/plugins-workspace/format-jsonl/src/index.js +++ b/hypaware-core/plugins-workspace/format-jsonl/src/index.js @@ -3,7 +3,7 @@ import { Buffer } from 'node:buffer' import zlib from 'node:zlib' -import { SpanStatusCode, trace } from '@opentelemetry/api' +import { getTracer, SpanStatusCode } from '../../../../src/core/observability/index.js' /** @typedef {import('../../../../collectivus-plugin-kernel-types').PluginActivationContext} PluginActivationContext */ /** @typedef {import('../../../../collectivus-plugin-kernel-types').QueryPartition} QueryPartition */ @@ -48,7 +48,7 @@ export async function activate(ctx) { * @returns {Promise} */ async function encodePartition(partition, ctx) { - const tracer = trace.getTracer(PLUGIN_NAME, PLUGIN_VERSION) + const tracer = getTracer('plugin.format-jsonl') return tracer.startActiveSpan( 'encoder.encode_jsonl', { diff --git a/hypaware-core/plugins-workspace/format-parquet/src/index.js b/hypaware-core/plugins-workspace/format-parquet/src/index.js index f819d294..352ce32e 100644 --- a/hypaware-core/plugins-workspace/format-parquet/src/index.js +++ b/hypaware-core/plugins-workspace/format-parquet/src/index.js @@ -1,8 +1,7 @@ // @ts-check -import { SpanStatusCode, trace } from '@opentelemetry/api' - import { rowsToColumnSources } from './columns.js' +import { getTracer, SpanStatusCode } from '../../../../src/core/observability/index.js' /** @typedef {import('../../../../collectivus-plugin-kernel-types').ColumnSpec} ColumnSpec */ /** @typedef {import('../../../../collectivus-plugin-kernel-types').PluginActivationContext} PluginActivationContext */ @@ -48,7 +47,7 @@ export async function activate(ctx) { * @returns {Promise} */ async function encodePartition(partition, ctx) { - const tracer = trace.getTracer(PLUGIN_NAME, PLUGIN_VERSION) + const tracer = getTracer('plugin.format-parquet') return tracer.startActiveSpan( 'encoder.encode_parquet', { diff --git a/hypaware-core/plugins-workspace/otel/src/collector.js b/hypaware-core/plugins-workspace/otel/src/collector.js index ec3e8fb0..aef6e85e 100644 --- a/hypaware-core/plugins-workspace/otel/src/collector.js +++ b/hypaware-core/plugins-workspace/otel/src/collector.js @@ -1,9 +1,8 @@ // @ts-check -import { trace } from '@opentelemetry/api' - import { Attr, + getActiveSpan, withSpan, } from '../../../../src/core/observability/index.js' import { columnsFor, otelTablePath, PLUGIN_NAME } from './datasets.js' @@ -162,7 +161,7 @@ function asObject(value) { * @param {number} port */ export function stampBoundAddress(host, port) { - const span = trace.getActiveSpan() + const span = getActiveSpan() if (!span) return span.setAttribute('listen_host', host) span.setAttribute('listen_port', port) diff --git a/package-lock.json b/package-lock.json index 673d3fd1..3c00f63c 100644 --- a/package-lock.json +++ b/package-lock.json @@ -8,18 +8,6 @@ "name": "hypaware", "version": "0.2.0", "dependencies": { - "@opentelemetry/api": "^1.9.1", - "@opentelemetry/api-logs": "^0.218.0", - "@opentelemetry/core": "^2.7.1", - "@opentelemetry/exporter-logs-otlp-http": "^0.218.0", - "@opentelemetry/exporter-metrics-otlp-http": "^0.218.0", - "@opentelemetry/exporter-trace-otlp-http": "^0.218.0", - "@opentelemetry/resources": "^2.7.1", - "@opentelemetry/sdk-logs": "^0.218.0", - "@opentelemetry/sdk-metrics": "^2.7.1", - "@opentelemetry/sdk-trace-base": "^2.7.1", - "@opentelemetry/sdk-trace-node": "^2.7.1", - "@opentelemetry/semantic-conventions": "^1.41.1", "hyparquet": "^1.25.8", "hyparquet-compressors": "^1.1.1", "icebird": "^0.7.0", @@ -36,240 +24,6 @@ "hyparquet-writer": "^0.15.1" } }, - "node_modules/@opentelemetry/api": { - "version": "1.9.1", - "resolved": "https://registry.npmjs.org/@opentelemetry/api/-/api-1.9.1.tgz", - "integrity": "sha512-gLyJlPHPZYdAk1JENA9LeHejZe1Ti77/pTeFm/nMXmQH/HFZlcS/O2XJB+L8fkbrNSqhdtlvjBVjxwUYanNH5Q==", - "license": "Apache-2.0", - "engines": { - "node": ">=8.0.0" - } - }, - "node_modules/@opentelemetry/api-logs": { - "version": "0.218.0", - "resolved": "https://registry.npmjs.org/@opentelemetry/api-logs/-/api-logs-0.218.0.tgz", - "integrity": "sha512-fmEWp5kXlGEc3i/lR698Hz41DfGyN4Tbe4g7L1AxSc7fF8Xeh/FQ9Quqpa9dVA413Q1Ad43QOLzU4JoXgbFPWw==", - "license": "Apache-2.0", - "dependencies": { - "@opentelemetry/api": "^1.3.0" - }, - "engines": { - "node": ">=8.0.0" - } - }, - "node_modules/@opentelemetry/context-async-hooks": { - "version": "2.7.1", - "resolved": "https://registry.npmjs.org/@opentelemetry/context-async-hooks/-/context-async-hooks-2.7.1.tgz", - "integrity": "sha512-OPFBYuXEn1E4ja3Y6eeA7O+ZnLBNcXTV5Cgsn1VaqBZ6hC5FnpZPLBNme1LJY8ZtF4aOujPKFoeWN4ik487KuQ==", - "license": "Apache-2.0", - "engines": { - "node": "^18.19.0 || >=20.6.0" - }, - "peerDependencies": { - "@opentelemetry/api": ">=1.0.0 <1.10.0" - } - }, - "node_modules/@opentelemetry/core": { - "version": "2.7.1", - "resolved": "https://registry.npmjs.org/@opentelemetry/core/-/core-2.7.1.tgz", - "integrity": "sha512-QAqIj32AtK6+pEVNG7EOVxHdE06RP+FM5qpiEJ4RtDcFIqKUZHYhl7/7UY5efhwmwNAg7j8QbJVBLxMerc0+gw==", - "license": "Apache-2.0", - "dependencies": { - "@opentelemetry/semantic-conventions": "^1.29.0" - }, - "engines": { - "node": "^18.19.0 || >=20.6.0" - }, - "peerDependencies": { - "@opentelemetry/api": ">=1.0.0 <1.10.0" - } - }, - "node_modules/@opentelemetry/exporter-logs-otlp-http": { - "version": "0.218.0", - "resolved": "https://registry.npmjs.org/@opentelemetry/exporter-logs-otlp-http/-/exporter-logs-otlp-http-0.218.0.tgz", - "integrity": "sha512-Qx+4rpVHzgg89dawcWRHyt+XRXeLnhFz/qBtvggmjkcgPUdr+NAB0/u/eIPA8yAeJV0J80Vz43JZCh/XFvZFGw==", - "license": "Apache-2.0", - "dependencies": { - "@opentelemetry/api-logs": "0.218.0", - "@opentelemetry/core": "2.7.1", - "@opentelemetry/otlp-exporter-base": "0.218.0", - "@opentelemetry/otlp-transformer": "0.218.0", - "@opentelemetry/sdk-logs": "0.218.0" - }, - "engines": { - "node": "^18.19.0 || >=20.6.0" - }, - "peerDependencies": { - "@opentelemetry/api": "^1.3.0" - } - }, - "node_modules/@opentelemetry/exporter-metrics-otlp-http": { - "version": "0.218.0", - "resolved": "https://registry.npmjs.org/@opentelemetry/exporter-metrics-otlp-http/-/exporter-metrics-otlp-http-0.218.0.tgz", - "integrity": "sha512-bV7d2OuMpZu2+gAaxUAhzfZ0h3WVZk8ETQUEE3DNSntbTaMpuITjtm8I0rNyHFdm7Ax57K6ty7SgFXlBmOLIvQ==", - "license": "Apache-2.0", - "dependencies": { - "@opentelemetry/core": "2.7.1", - "@opentelemetry/otlp-exporter-base": "0.218.0", - "@opentelemetry/otlp-transformer": "0.218.0", - "@opentelemetry/resources": "2.7.1", - "@opentelemetry/sdk-metrics": "2.7.1" - }, - "engines": { - "node": "^18.19.0 || >=20.6.0" - }, - "peerDependencies": { - "@opentelemetry/api": "^1.3.0" - } - }, - "node_modules/@opentelemetry/exporter-trace-otlp-http": { - "version": "0.218.0", - "resolved": "https://registry.npmjs.org/@opentelemetry/exporter-trace-otlp-http/-/exporter-trace-otlp-http-0.218.0.tgz", - "integrity": "sha512-8dqezsmPhtKitIK/eTipZhYl9EX2/gNQ5zUMhaz3uxEURwfkNf8IPvo6yNfrzbxdtpAOybS/+h7wmIWYqFSpiw==", - "license": "Apache-2.0", - "dependencies": { - "@opentelemetry/core": "2.7.1", - "@opentelemetry/otlp-exporter-base": "0.218.0", - "@opentelemetry/otlp-transformer": "0.218.0", - "@opentelemetry/resources": "2.7.1", - "@opentelemetry/sdk-trace-base": "2.7.1" - }, - "engines": { - "node": "^18.19.0 || >=20.6.0" - }, - "peerDependencies": { - "@opentelemetry/api": "^1.3.0" - } - }, - "node_modules/@opentelemetry/otlp-exporter-base": { - "version": "0.218.0", - "resolved": "https://registry.npmjs.org/@opentelemetry/otlp-exporter-base/-/otlp-exporter-base-0.218.0.tgz", - "integrity": "sha512-ZwqpkNL5W7RyGJPDZ9g06DvKp8KFTWPJPN12anpMQYSKpTSU0z3EIZuPq9vPGpS8siFyOqDYDAuCwlNO9FqgbA==", - "license": "Apache-2.0", - "dependencies": { - "@opentelemetry/core": "2.7.1", - "@opentelemetry/otlp-transformer": "0.218.0" - }, - "engines": { - "node": "^18.19.0 || >=20.6.0" - }, - "peerDependencies": { - "@opentelemetry/api": "^1.3.0" - } - }, - "node_modules/@opentelemetry/otlp-transformer": { - "version": "0.218.0", - "resolved": "https://registry.npmjs.org/@opentelemetry/otlp-transformer/-/otlp-transformer-0.218.0.tgz", - "integrity": "sha512-CFaKH87WAzjuJ4awowTTLzUvMfaRfiOFG5+qm5S5ncyalRtN4ecQ+YmuANJSCrVPuvZFEkUgKhBPBndxi3rHsQ==", - "license": "Apache-2.0", - "dependencies": { - "@opentelemetry/api-logs": "0.218.0", - "@opentelemetry/core": "2.7.1", - "@opentelemetry/resources": "2.7.1", - "@opentelemetry/sdk-logs": "0.218.0", - "@opentelemetry/sdk-metrics": "2.7.1", - "@opentelemetry/sdk-trace-base": "2.7.1" - }, - "engines": { - "node": "^18.19.0 || >=20.6.0" - }, - "peerDependencies": { - "@opentelemetry/api": "^1.3.0" - } - }, - "node_modules/@opentelemetry/resources": { - "version": "2.7.1", - "resolved": "https://registry.npmjs.org/@opentelemetry/resources/-/resources-2.7.1.tgz", - "integrity": "sha512-DeT6KKolmC4e/dRQvMQ/RwlnzhaqeiFOXY5ngoOPJ07GgVVKxZOg9EcrNZb5aTzUn+iCrJldAgOfQm1O/QfPAQ==", - "license": "Apache-2.0", - "dependencies": { - "@opentelemetry/core": "2.7.1", - "@opentelemetry/semantic-conventions": "^1.29.0" - }, - "engines": { - "node": "^18.19.0 || >=20.6.0" - }, - "peerDependencies": { - "@opentelemetry/api": ">=1.3.0 <1.10.0" - } - }, - "node_modules/@opentelemetry/sdk-logs": { - "version": "0.218.0", - "resolved": "https://registry.npmjs.org/@opentelemetry/sdk-logs/-/sdk-logs-0.218.0.tgz", - "integrity": "sha512-QvnNdugatFTVCJXH0Mcu7GOOJSylA9j127kIezOE4YwTI4YbowRons2K4WZTv5FMS8T4q9P0NdaRHdkSmeAIag==", - "license": "Apache-2.0", - "dependencies": { - "@opentelemetry/api-logs": "0.218.0", - "@opentelemetry/core": "2.7.1", - "@opentelemetry/resources": "2.7.1", - "@opentelemetry/semantic-conventions": "^1.29.0" - }, - "engines": { - "node": "^18.19.0 || >=20.6.0" - }, - "peerDependencies": { - "@opentelemetry/api": ">=1.4.0 <1.10.0" - } - }, - "node_modules/@opentelemetry/sdk-metrics": { - "version": "2.7.1", - "resolved": "https://registry.npmjs.org/@opentelemetry/sdk-metrics/-/sdk-metrics-2.7.1.tgz", - "integrity": "sha512-MpDJdkiFDs3Pm1RHO3KByuZbuBdJEXEAkiC0+yJdsZGVCdf1RpHR6n+LHDcS7ffmfrt5kVCzJSCfm4z2C7v0uQ==", - "license": "Apache-2.0", - "dependencies": { - "@opentelemetry/core": "2.7.1", - "@opentelemetry/resources": "2.7.1" - }, - "engines": { - "node": "^18.19.0 || >=20.6.0" - }, - "peerDependencies": { - "@opentelemetry/api": ">=1.9.0 <1.10.0" - } - }, - "node_modules/@opentelemetry/sdk-trace-base": { - "version": "2.7.1", - "resolved": "https://registry.npmjs.org/@opentelemetry/sdk-trace-base/-/sdk-trace-base-2.7.1.tgz", - "integrity": "sha512-NAYIlsF8MPUsKqJMiDQJTMPOmlbawC1Iz/omMLygZ1C9am8fTKYjTaI+OZM+WTY3t3Glo0wnOg/6/pac6RGPPw==", - "license": "Apache-2.0", - "dependencies": { - "@opentelemetry/core": "2.7.1", - "@opentelemetry/resources": "2.7.1", - "@opentelemetry/semantic-conventions": "^1.29.0" - }, - "engines": { - "node": "^18.19.0 || >=20.6.0" - }, - "peerDependencies": { - "@opentelemetry/api": ">=1.3.0 <1.10.0" - } - }, - "node_modules/@opentelemetry/sdk-trace-node": { - "version": "2.7.1", - "resolved": "https://registry.npmjs.org/@opentelemetry/sdk-trace-node/-/sdk-trace-node-2.7.1.tgz", - "integrity": "sha512-pCpQxU68lV+I9s9svqMyVu5iHdDDUnqUpSxqwyCU8A9ejEsSnMPCbearwsUO4yk08ZJzAIUCFuReMdVQvHrdvg==", - "license": "Apache-2.0", - "dependencies": { - "@opentelemetry/context-async-hooks": "2.7.1", - "@opentelemetry/core": "2.7.1", - "@opentelemetry/sdk-trace-base": "2.7.1" - }, - "engines": { - "node": "^18.19.0 || >=20.6.0" - }, - "peerDependencies": { - "@opentelemetry/api": ">=1.0.0 <1.10.0" - } - }, - "node_modules/@opentelemetry/semantic-conventions": { - "version": "1.41.1", - "resolved": "https://registry.npmjs.org/@opentelemetry/semantic-conventions/-/semantic-conventions-1.41.1.tgz", - "integrity": "sha512-/UhIkaZgPutTFmQ7RnIJGgDXZmtEJ7Dvi86xNTFWcnRxVRNk/aotsqDJYeEvDP+FSMB2SdW+pQzNMcWP0rwuNA==", - "license": "Apache-2.0", - "engines": { - "node": ">=14" - } - }, "node_modules/fzstd": { "version": "0.1.1", "resolved": "https://registry.npmjs.org/fzstd/-/fzstd-0.1.1.tgz", diff --git a/package.json b/package.json index f0ffcc2f..9c93609c 100644 --- a/package.json +++ b/package.json @@ -29,18 +29,6 @@ "smoke": "node ./hypaware-core/smoke/index.js" }, "dependencies": { - "@opentelemetry/api": "^1.9.1", - "@opentelemetry/api-logs": "^0.218.0", - "@opentelemetry/core": "^2.7.1", - "@opentelemetry/exporter-logs-otlp-http": "^0.218.0", - "@opentelemetry/exporter-metrics-otlp-http": "^0.218.0", - "@opentelemetry/exporter-trace-otlp-http": "^0.218.0", - "@opentelemetry/resources": "^2.7.1", - "@opentelemetry/sdk-logs": "^0.218.0", - "@opentelemetry/sdk-metrics": "^2.7.1", - "@opentelemetry/sdk-trace-base": "^2.7.1", - "@opentelemetry/sdk-trace-node": "^2.7.1", - "@opentelemetry/semantic-conventions": "^1.41.1", "hyparquet": "^1.25.8", "hyparquet-compressors": "^1.1.1", "icebird": "^0.7.0", diff --git a/src/core/cli/dispatch.js b/src/core/cli/dispatch.js index 93bbcc5f..876c5882 100644 --- a/src/core/cli/dispatch.js +++ b/src/core/cli/dispatch.js @@ -3,15 +3,17 @@ import process from 'node:process' import path from 'node:path' import { performance } from 'node:perf_hooks' -import { context, ROOT_CONTEXT, SpanStatusCode } from '@opentelemetry/api' import { Attr, buildAttrs, + context, getKernelInstruments, getLogger, getTracer, installObservability, + ROOT_CONTEXT, + SpanStatusCode, } from '../observability/index.js' import { createCommandRegistry } from '../registry/commands.js' import { createKernelRuntime } from '../runtime/activation.js' diff --git a/src/core/observability/index.js b/src/core/observability/index.js index 23dcb298..561ecf21 100644 --- a/src/core/observability/index.js +++ b/src/core/observability/index.js @@ -33,10 +33,10 @@ export function installObservability(opts = {}) { /** * @param {{ * env: import('./env.js').ObservabilityEnv, - * resource: import('@opentelemetry/resources').Resource, - * tracer: { provider: import('@opentelemetry/sdk-trace-node').NodeTracerProvider|null }, - * logger: { provider: import('@opentelemetry/sdk-logs').LoggerProvider|null }, - * meter: { provider: import('@opentelemetry/sdk-metrics').MeterProvider|null, readers: import('@opentelemetry/sdk-metrics').MetricReader[] } + * resource: { attributes: Record }, + * tracer: { provider: import('./runtime.js').TracerProvider|null }, + * logger: { provider: import('./runtime.js').LoggerProvider|null }, + * meter: { provider: import('./runtime.js').MeterProvider|null, readers: object[] } * }} parts */ function buildHandle({ env, resource, tracer, logger, meter }) { @@ -91,3 +91,4 @@ export { getLogger } from './logger.js' export { getMeter, getKernelInstruments } from './meter.js' export { withSpan, runRoot } from './span_helpers.js' export { buildAttrs, normalizeKey, Attr } from './attrs.js' +export { context, ROOT_CONTEXT, SpanStatusCode, getActiveSpan } from './runtime.js' diff --git a/src/core/observability/jsonl_exporters.js b/src/core/observability/jsonl_exporters.js index 99648fcc..7921d519 100644 --- a/src/core/observability/jsonl_exporters.js +++ b/src/core/observability/jsonl_exporters.js @@ -3,7 +3,12 @@ import fs from 'node:fs' import path from 'node:path' -import { ExportResultCode } from '@opentelemetry/core' +import { hrTimeToIso } from './runtime.js' + +const ExportResultCode = Object.freeze({ + SUCCESS: 0, + FAILED: 1, +}) /** * Append-only JSONL writer. One file per signal per pid; the file is @@ -75,12 +80,12 @@ class JsonlWriter { } /** - * @param {import('@opentelemetry/sdk-trace-base').ReadableSpan} span + * @param {import('./runtime.js').Span} span */ function spanToJsonl(span) { const ctx = span.spanContext() - const startMs = span.startTime[0] * 1000 + span.startTime[1] / 1_000_000 - const endMs = span.endTime[0] * 1000 + span.endTime[1] / 1_000_000 + const startMs = hrtimeToMs(span.startTime) + const endMs = hrtimeToMs(span.endTime) return { serviceName: span.resource.attributes['service.name'] ?? 'unknown', name: span.name, @@ -88,15 +93,15 @@ function spanToJsonl(span) { spanId: ctx.spanId, parentSpanId: span.parentSpanContext?.spanId ?? null, kind: span.kind, - startTimestamp: hrtimeToIso(span.startTime), - endTimestamp: hrtimeToIso(span.endTime), + startTimestamp: hrTimeToIso(span.startTime), + endTimestamp: hrTimeToIso(span.endTime), durationMs: endMs - startMs, status: spanStatusName(span.status.code), statusMessage: span.status.message, attributes: span.attributes, events: span.events.map((e) => ({ name: e.name, - time: hrtimeToIso(e.time), + time: hrTimeToIso(e.time), attributes: e.attributes, })), resource: span.resource.attributes, @@ -114,18 +119,8 @@ function spanStatusName(code) { return 'unset' } -/** - * @param {[number, number]} hrtime - */ -function hrtimeToIso(hrtime) { - const ms = hrtime[0] * 1000 + hrtime[1] / 1_000_000 - return new Date(ms).toISOString() -} - /** * SpanExporter implementation that writes each batch as JSONL. - * - * @implements {import('@opentelemetry/sdk-trace-base').SpanExporter} */ export class JsonlSpanExporter { /** @@ -138,7 +133,7 @@ export class JsonlSpanExporter { } /** - * @param {import('@opentelemetry/sdk-trace-base').ReadableSpan[]} spans + * @param {import('./runtime.js').Span[]} spans * @param {(result: { code: number, error?: Error }) => void} resultCallback */ export(spans, resultCallback) { @@ -153,6 +148,11 @@ export class JsonlSpanExporter { } } + /** @param {import('./runtime.js').Span[]} spans */ + exportBatch(spans) { + this.export(spans, () => {}) + } + async shutdown() { await this.writer.close() } @@ -163,14 +163,14 @@ export class JsonlSpanExporter { } /** - * @param {import('@opentelemetry/sdk-logs').ReadableLogRecord} record + * @param {import('./runtime.js').LogRecord} record */ function logRecordToJsonl(record) { const hr = record.hrTime || record.hrTimeObserved || [0, 0] return { serviceName: record.resource.attributes['service.name'] ?? 'unknown', - timestamp: hrtimeToIso(hr), - observedTimestamp: hrtimeToIso(record.hrTimeObserved || hr), + timestamp: hrTimeToIso(hr), + observedTimestamp: hrTimeToIso(record.hrTimeObserved || hr), severityNumber: record.severityNumber ?? 0, severityText: record.severityText ?? '', body: serializeBody(record.body), @@ -194,8 +194,6 @@ function serializeBody(body) { /** * LogRecordExporter that writes JSONL. - * - * @implements {import('@opentelemetry/sdk-logs').LogRecordExporter} */ export class JsonlLogRecordExporter { /** @@ -208,7 +206,7 @@ export class JsonlLogRecordExporter { } /** - * @param {import('@opentelemetry/sdk-logs').ReadableLogRecord[]} records + * @param {import('./runtime.js').LogRecord[]} records * @param {(result: { code: number, error?: Error }) => void} resultCallback */ export(records, resultCallback) { @@ -223,6 +221,11 @@ export class JsonlLogRecordExporter { } } + /** @param {import('./runtime.js').LogRecord[]} records */ + exportBatch(records) { + this.export(records, () => {}) + } + async shutdown() { await this.writer.close() } @@ -236,50 +239,24 @@ export class JsonlLogRecordExporter { * PushMetricExporter that writes JSONL. Each export call emits one * record per data point, flattened so smoke assertions can query a * single named metric without unpacking the OTel resource metrics tree. - * - * @implements {import('@opentelemetry/sdk-metrics').PushMetricExporter} */ export class JsonlMetricExporter { /** * @param {object} opts * @param {string} opts.dir * @param {number} [opts.pid] - * @param {import('@opentelemetry/sdk-metrics').AggregationTemporality} [opts.temporality] */ - constructor({ dir, pid = process.pid, temporality }) { + constructor({ dir, pid = process.pid }) { this.writer = new JsonlWriter(dir, `metrics-${pid}.jsonl`) - this._temporality = temporality } /** - * @param {import('@opentelemetry/sdk-metrics').ResourceMetrics} metrics + * @param {import('./runtime.js').MetricRecord[]} records * @param {(result: { code: number, error?: Error }) => void} resultCallback */ - export(metrics, resultCallback) { + export(records, resultCallback) { try { - /** @type {object[]} */ - const records = [] - const resourceAttrs = metrics.resource.attributes - const serviceName = resourceAttrs['service.name'] ?? 'unknown' - for (const scopeMetrics of metrics.scopeMetrics) { - for (const metric of scopeMetrics.metrics) { - for (const dataPoint of metric.dataPoints) { - records.push({ - serviceName, - name: metric.descriptor.name, - description: metric.descriptor.description, - unit: metric.descriptor.unit, - type: metric.dataPointType, - attributes: dataPoint.attributes, - value: serializeMetricValue(metric.dataPointType, dataPoint.value), - startTimestamp: hrtimeToIso(dataPoint.startTime), - endTimestamp: hrtimeToIso(dataPoint.endTime), - resource: resourceAttrs, - }) - } - } - } - this.writer.writeBatch(records) + this.writer.writeBatch(records.map(metricRecordToJsonl)) resultCallback({ code: ExportResultCode.SUCCESS }) } catch (error) { resultCallback({ @@ -289,12 +266,9 @@ export class JsonlMetricExporter { } } - selectAggregationTemporality() { - // Cumulative matches the default for Sum/Histogram aggregations and - // keeps Sum data points monotonic across exports. - return /** @type {import('@opentelemetry/sdk-metrics').AggregationTemporality} */ ( - this._temporality ?? 1 - ) + /** @param {import('./runtime.js').MetricRecord[]} records */ + exportBatch(records) { + this.export(records, () => {}) } async shutdown() { @@ -307,23 +281,27 @@ export class JsonlMetricExporter { } /** - * @param {number} type - * @param {unknown} value + * @param {import('./runtime.js').MetricRecord} record */ -function serializeMetricValue(type, value) { - if (value && typeof value === 'object' && 'buckets' in value) { - const v = /** @type {{count:number,sum:number,min?:number,max?:number,buckets:{boundaries:number[],counts:number[]}}} */ (value) - return { - count: v.count, - sum: v.sum, - min: v.min, - max: v.max, - boundaries: v.buckets.boundaries, - counts: v.buckets.counts, - } - } - if (typeof value === 'number' || typeof value === 'bigint') { - return Number(value) +function metricRecordToJsonl(record) { + const resourceAttrs = record.resource.attributes + return { + serviceName: resourceAttrs['service.name'] ?? 'unknown', + name: record.name, + description: record.description, + unit: record.unit, + type: record.kind, + attributes: record.attributes, + value: record.value, + startTimestamp: hrTimeToIso(record.startTime), + endTimestamp: hrTimeToIso(record.endTime), + resource: resourceAttrs, } - return value +} + +/** + * @param {[number, number]} hrtime + */ +function hrtimeToMs(hrtime) { + return hrtime[0] * 1000 + hrtime[1] / 1_000_000 } diff --git a/src/core/observability/logger.js b/src/core/observability/logger.js index a8ffbaf4..f36cdb3b 100644 --- a/src/core/observability/logger.js +++ b/src/core/observability/logger.js @@ -1,12 +1,10 @@ // @ts-check -import { logs, SeverityNumber } from '@opentelemetry/api-logs' -import { LoggerProvider, SimpleLogRecordProcessor } from '@opentelemetry/sdk-logs' -import { OTLPLogExporter } from '@opentelemetry/exporter-logs-otlp-http' - import { JsonlLogRecordExporter } from './jsonl_exporters.js' import { devTelemetryDir } from './env.js' import { Attr, buildAttrs } from './attrs.js' +import { logs, LoggerProvider, SeverityNumber } from './runtime.js' +import { OtlpLogExporter } from './otlp_exporters.js' const OTLP_EXPORT_TIMEOUT_MS = 1_000 @@ -31,38 +29,34 @@ const SEVERITY_TEXT = Object.freeze({ * * @param {object} args * @param {import('./env.js').ObservabilityEnv} args.env - * @param {import('@opentelemetry/resources').Resource} args.resource + * @param {{ attributes: Record }} args.resource * @returns {{ provider: LoggerProvider|null, exporters: object[] }} */ export function installLoggerProvider({ env, resource }) { /** @type {object[]} */ const exporters = [] - /** @type {import('@opentelemetry/sdk-logs').LogRecordProcessor[]} */ - const processors = [] if (env.devTelemetry) { const dir = devTelemetryDir(env.stateDir) const jsonlExporter = new JsonlLogRecordExporter({ dir }) - processors.push(new SimpleLogRecordProcessor(jsonlExporter)) exporters.push(jsonlExporter) } if (!env.devTelemetry && env.otlpEndpoint) { - const otlpExporter = new OTLPLogExporter({ + const otlpExporter = new OtlpLogExporter({ url: env.otlpEndpoint.replace(/\/$/, '') + '/v1/logs', timeoutMillis: OTLP_EXPORT_TIMEOUT_MS, }) - processors.push(new SimpleLogRecordProcessor(otlpExporter)) exporters.push(otlpExporter) } - if (processors.length === 0) { + if (exporters.length === 0) { return { provider: null, exporters: [] } } const provider = new LoggerProvider({ resource, - processors, + exporters, }) logs.setGlobalLoggerProvider(provider) return { provider, exporters } diff --git a/src/core/observability/meter.js b/src/core/observability/meter.js index 13a8266d..752fac0d 100644 --- a/src/core/observability/meter.js +++ b/src/core/observability/meter.js @@ -1,11 +1,9 @@ // @ts-check -import { metrics } from '@opentelemetry/api' -import { MeterProvider, PeriodicExportingMetricReader } from '@opentelemetry/sdk-metrics' -import { OTLPMetricExporter } from '@opentelemetry/exporter-metrics-otlp-http' - import { JsonlMetricExporter } from './jsonl_exporters.js' import { devTelemetryDir } from './env.js' +import { MeterProvider, metrics } from './runtime.js' +import { OtlpMetricExporter } from './otlp_exporters.js' const OTLP_EXPORT_TIMEOUT_MS = 1_000 @@ -17,45 +15,34 @@ const OTLP_EXPORT_TIMEOUT_MS = 1_000 * * @param {object} args * @param {import('./env.js').ObservabilityEnv} args.env - * @param {import('@opentelemetry/resources').Resource} args.resource - * @returns {{ provider: MeterProvider|null, exporters: object[], readers: import('@opentelemetry/sdk-metrics').MetricReader[] }} + * @param {{ attributes: Record }} args.resource + * @returns {{ provider: MeterProvider|null, exporters: object[], readers: object[] }} */ export function installMeterProvider({ env, resource }) { /** @type {object[]} */ const exporters = [] - /** @type {import('@opentelemetry/sdk-metrics').MetricReader[]} */ - const readers = [] if (env.devTelemetry) { const dir = devTelemetryDir(env.stateDir) const jsonlExporter = new JsonlMetricExporter({ dir }) - readers.push(new PeriodicExportingMetricReader({ - exporter: jsonlExporter, - exportIntervalMillis: 250, - })) exporters.push(jsonlExporter) } if (!env.devTelemetry && env.otlpEndpoint) { - const otlpExporter = new OTLPMetricExporter({ + const otlpExporter = new OtlpMetricExporter({ url: env.otlpEndpoint.replace(/\/$/, '') + '/v1/metrics', timeoutMillis: OTLP_EXPORT_TIMEOUT_MS, }) - readers.push(new PeriodicExportingMetricReader({ - exporter: otlpExporter, - exportIntervalMillis: 30_000, - exportTimeoutMillis: OTLP_EXPORT_TIMEOUT_MS, - })) exporters.push(otlpExporter) } - if (readers.length === 0) { + if (exporters.length === 0) { return { provider: null, exporters: [], readers: [] } } - const provider = new MeterProvider({ resource, readers }) + const provider = new MeterProvider({ resource, exporters }) metrics.setGlobalMeterProvider(provider) - return { provider, exporters, readers } + return { provider, exporters, readers: [] } } /** @@ -63,7 +50,7 @@ export function installMeterProvider({ env, resource }) { * contract. Plugins declare their own meters; this set is reserved * for things the kernel itself emits. * - * @param {import('@opentelemetry/api').Meter} meter + * @param {{ createCounter(name: string, opts?: object): object, createUpDownCounter(name: string, opts?: object): object, createGauge(name: string, opts?: object): object, createHistogram(name: string, opts?: object): object }} meter */ function buildKernelInstruments(meter) { return { diff --git a/src/core/observability/otlp_exporters.js b/src/core/observability/otlp_exporters.js new file mode 100644 index 00000000..37281f8d --- /dev/null +++ b/src/core/observability/otlp_exporters.js @@ -0,0 +1,251 @@ +// @ts-check + +import { hrTimeToUnixNano, SpanStatusCode } from './runtime.js' + +const OTLP_AGGREGATION_TEMPORALITY_CUMULATIVE = 2 + +class OtlpHttpJsonExporter { + /** + * @param {object} opts + * @param {string} opts.url + * @param {number} opts.timeoutMillis + */ + constructor({ url, timeoutMillis }) { + this.url = url + this.timeoutMillis = timeoutMillis + /** @type {Promise[]} */ + this.pending = [] + } + + /** @param {unknown} payload */ + post(payload) { + const controller = new AbortController() + const timer = setTimeout(() => controller.abort(), this.timeoutMillis) + if (typeof timer.unref === 'function') timer.unref() + const request = fetch(this.url, { + method: 'POST', + headers: { 'content-type': 'application/json' }, + body: JSON.stringify(payload), + signal: controller.signal, + }).catch(() => undefined).finally(() => clearTimeout(timer)) + this.pending.push(request) + } + + async forceFlush() { + const pending = this.pending.splice(0) + await Promise.allSettled(pending) + } + + async shutdown() { + await this.forceFlush() + } +} + +export class OtlpSpanExporter extends OtlpHttpJsonExporter { + /** @param {import('./runtime.js').Span[]} spans */ + exportBatch(spans) { + if (spans.length === 0) return + this.post({ resourceSpans: groupByResourceAndScope(spans, spanToOtlp) }) + } +} + +export class OtlpLogExporter extends OtlpHttpJsonExporter { + /** @param {import('./runtime.js').LogRecord[]} records */ + exportBatch(records) { + if (records.length === 0) return + this.post({ resourceLogs: groupLogsByResourceAndScope(records) }) + } +} + +export class OtlpMetricExporter extends OtlpHttpJsonExporter { + /** @param {import('./runtime.js').MetricRecord[]} records */ + exportBatch(records) { + if (records.length === 0) return + this.post({ resourceMetrics: groupMetricsByResourceAndScope(records) }) + } +} + +/** + * @param {import('./runtime.js').Span[]} spans + * @param {(span: import('./runtime.js').Span) => object} mapSpan + */ +function groupByResourceAndScope(spans, mapSpan) { + /** @type {Map }>} */ + const byResource = new Map() + for (const span of spans) { + const resourceKey = JSON.stringify(span.resource.attributes) + let group = byResource.get(resourceKey) + if (!group) { + group = { resource: { attributes: attrsToOtlp(span.resource.attributes) }, scopeSpans: [] } + byResource.set(resourceKey, group) + } + let scope = group.scopeSpans.find((s) => s.scope.name === span.tracerName && s.scope.version === span.tracerVersion) + if (!scope) { + scope = { scope: cleanObject({ name: span.tracerName, version: span.tracerVersion }), spans: [] } + group.scopeSpans.push(scope) + } + scope.spans.push(mapSpan(span)) + } + return [...byResource.values()] +} + +/** @param {import('./runtime.js').LogRecord[]} records */ +function groupLogsByResourceAndScope(records) { + /** @type {Map }>} */ + const byResource = new Map() + for (const record of records) { + const resourceKey = JSON.stringify(record.resource.attributes) + let group = byResource.get(resourceKey) + if (!group) { + group = { resource: { attributes: attrsToOtlp(record.resource.attributes) }, scopeLogs: [] } + byResource.set(resourceKey, group) + } + let scope = group.scopeLogs.find((s) => s.scope.name === record.loggerName && s.scope.version === record.loggerVersion) + if (!scope) { + scope = { scope: cleanObject({ name: record.loggerName, version: record.loggerVersion }), logRecords: [] } + group.scopeLogs.push(scope) + } + scope.logRecords.push(logToOtlp(record)) + } + return [...byResource.values()] +} + +/** @param {import('./runtime.js').MetricRecord[]} records */ +function groupMetricsByResourceAndScope(records) { + /** @type {Map }>} */ + const byResource = new Map() + for (const record of records) { + const resourceKey = JSON.stringify(record.resource.attributes) + let group = byResource.get(resourceKey) + if (!group) { + group = { resource: { attributes: attrsToOtlp(record.resource.attributes) }, scopeMetrics: [] } + byResource.set(resourceKey, group) + } + let scope = group.scopeMetrics.find((s) => s.scope.name === record.meterName && s.scope.version === record.meterVersion) + if (!scope) { + scope = { scope: cleanObject({ name: record.meterName, version: record.meterVersion }), metrics: [] } + group.scopeMetrics.push(scope) + } + scope.metrics.push(metricToOtlp(record)) + } + return [...byResource.values()] +} + +/** @param {import('./runtime.js').Span} span */ +function spanToOtlp(span) { + const ctx = span.spanContext() + return cleanObject({ + traceId: ctx.traceId, + spanId: ctx.spanId, + parentSpanId: span.parentSpanContext?.spanId, + name: span.name, + kind: span.kind, + startTimeUnixNano: String(hrTimeToUnixNano(span.startTime)), + endTimeUnixNano: String(hrTimeToUnixNano(span.endTime)), + attributes: attrsToOtlp(span.attributes), + events: span.events.map((event) => ({ + timeUnixNano: String(hrTimeToUnixNano(event.time)), + name: event.name, + attributes: attrsToOtlp(event.attributes), + })), + status: spanStatusToOtlp(span.status), + }) +} + +/** @param {import('./runtime.js').LogRecord} record */ +function logToOtlp(record) { + return cleanObject({ + timeUnixNano: String(hrTimeToUnixNano(record.hrTime)), + observedTimeUnixNano: String(hrTimeToUnixNano(record.hrTimeObserved)), + severityNumber: record.severityNumber, + severityText: record.severityText, + body: anyValueToOtlp(record.body), + traceId: record.spanContext?.traceId, + spanId: record.spanContext?.spanId, + flags: record.spanContext?.traceFlags, + attributes: attrsToOtlp(record.attributes), + }) +} + +/** @param {import('./runtime.js').MetricRecord} record */ +function metricToOtlp(record) { + const pointBase = { + startTimeUnixNano: String(hrTimeToUnixNano(record.startTime)), + timeUnixNano: String(hrTimeToUnixNano(record.endTime)), + attributes: attrsToOtlp(record.attributes), + } + if (record.kind === 'histogram') { + return cleanObject({ + name: record.name, + description: record.description, + unit: record.unit, + histogram: { + aggregationTemporality: OTLP_AGGREGATION_TEMPORALITY_CUMULATIVE, + dataPoints: [{ + ...pointBase, + count: '1', + sum: record.value, + bucketCounts: ['1'], + explicitBounds: [], + }], + }, + }) + } + const containerName = record.kind === 'gauge' ? 'gauge' : 'sum' + return cleanObject({ + name: record.name, + description: record.description, + unit: record.unit, + [containerName]: { + ...(containerName === 'sum' + ? { + aggregationTemporality: OTLP_AGGREGATION_TEMPORALITY_CUMULATIVE, + isMonotonic: record.monotonic, + } + : {}), + dataPoints: [{ ...pointBase, asDouble: record.value }], + }, + }) +} + +/** + * @param {{ code: number, message?: string }} status + */ +function spanStatusToOtlp(status) { + return cleanObject({ + code: status.code === SpanStatusCode.ERROR ? 2 : status.code === SpanStatusCode.OK ? 1 : 0, + message: status.message, + }) +} + +/** @param {Record} attrs */ +function attrsToOtlp(attrs) { + return Object.entries(attrs ?? {}) + .filter(([key, value]) => key.length > 0 && value !== undefined) + .map(([key, value]) => ({ key, value: anyValueToOtlp(value) })) +} + +/** @param {unknown} value */ +function anyValueToOtlp(value) { + if (typeof value === 'string') return { stringValue: value } + if (typeof value === 'boolean') return { boolValue: value } + if (typeof value === 'number') { + return Number.isInteger(value) ? { intValue: String(value) } : { doubleValue: value } + } + if (typeof value === 'bigint') return { intValue: String(value) } + if (Array.isArray(value)) return { arrayValue: { values: value.map(anyValueToOtlp) } } + if (value && typeof value === 'object') { + return { + kvlistValue: { + values: Object.entries(/** @type {Record} */ (value)) + .map(([key, item]) => ({ key, value: anyValueToOtlp(item) })), + }, + } + } + return { stringValue: '' } +} + +/** @param {Record} obj */ +function cleanObject(obj) { + return Object.fromEntries(Object.entries(obj).filter(([, value]) => value !== undefined)) +} diff --git a/src/core/observability/resource.js b/src/core/observability/resource.js index 9c6d1865..7155560c 100644 --- a/src/core/observability/resource.js +++ b/src/core/observability/resource.js @@ -1,16 +1,14 @@ // @ts-check -import { resourceFromAttributes } from '@opentelemetry/resources' - /** - * Build the OTel Resource that the tracer, logger, and meter providers + * Build the resource metadata that the tracer, logger, and meter providers * share. `service.name` comes from env; `dev_run_id` is mirrored onto * the resource so every signal carries it even when a caller forgets * to set it as a span attribute. Any `OTEL_RESOURCE_ATTRIBUTES` value * (a=b,c=d shape) is merged last so user-supplied keys win. * * @param {import('./env.js').ObservabilityEnv} env - * @returns {import('@opentelemetry/resources').Resource} + * @returns {{ attributes: Record }} */ export function buildResource(env) { /** @type {Record} */ @@ -30,5 +28,5 @@ export function buildResource(env) { if (key) attrs[key] = value } } - return resourceFromAttributes(attrs) + return { attributes: attrs } } diff --git a/src/core/observability/runtime.js b/src/core/observability/runtime.js new file mode 100644 index 00000000..799a94cf --- /dev/null +++ b/src/core/observability/runtime.js @@ -0,0 +1,493 @@ +// @ts-check + +import { AsyncLocalStorage } from 'node:async_hooks' +import { performance } from 'node:perf_hooks' +import crypto from 'node:crypto' + +export const SpanStatusCode = Object.freeze({ + UNSET: 0, + OK: 1, + ERROR: 2, +}) + +export const SeverityNumber = Object.freeze({ + DEBUG: 5, + INFO: 9, + WARN: 13, + ERROR: 17, +}) + +export const ROOT_CONTEXT = Object.freeze({ span: null }) + +/** @type {AsyncLocalStorage<{ span: Span|null }>} */ +const activeContext = new AsyncLocalStorage() + +/** @type {TracerProvider|null} */ +let globalTracerProvider = null +/** @type {LoggerProvider|null} */ +let globalLoggerProvider = null +/** @type {MeterProvider|null} */ +let globalMeterProvider = null + +export const context = Object.freeze({ + /** @param {{ span: Span|null }} ctx @param {() => unknown} fn */ + with(ctx, fn) { + return activeContext.run(ctx ?? ROOT_CONTEXT, fn) + }, + active() { + return activeContext.getStore() ?? ROOT_CONTEXT + }, +}) + +export const trace = Object.freeze({ + /** + * @param {string} name + * @param {string} [version] + */ + getTracer(name, version) { + return new Tracer(name, version) + }, + getActiveSpan() { + return activeContext.getStore()?.span ?? null + }, + getTracerProvider() { + return globalTracerProvider ?? NOOP_TRACER_PROVIDER + }, +}) + +export const logs = Object.freeze({ + /** @param {LoggerProvider} provider */ + setGlobalLoggerProvider(provider) { + globalLoggerProvider = provider + }, + /** + * @param {string} name + * @param {string} [version] + */ + getLogger(name, version) { + return new Logger(name, version) + }, +}) + +export const metrics = Object.freeze({ + /** @param {MeterProvider} provider */ + setGlobalMeterProvider(provider) { + globalMeterProvider = provider + }, + /** + * @param {string} name + * @param {string} [version] + */ + getMeter(name, version) { + return new Meter(name, version) + }, +}) + +export class TracerProvider { + /** + * @param {object} opts + * @param {{ attributes: Record }} opts.resource + * @param {Array<{ exportBatch(spans: Span[]): unknown, forceFlush?: () => Promise|void, shutdown?: () => Promise|void }>} [opts.exporters] + */ + constructor({ resource, exporters = [] }) { + this.resource = resource + this.exporters = exporters + } + + register() { + globalTracerProvider = this + } + + /** @param {Span} span */ + exportSpan(span) { + if (this.exporters.length === 0) return + for (const exporter of this.exporters) exporter.exportBatch([span]) + } + + async forceFlush() { + await flushExporters(this.exporters) + } + + async shutdown() { + await shutdownExporters(this.exporters) + if (globalTracerProvider === this) globalTracerProvider = null + } +} + +export class LoggerProvider { + /** + * @param {object} opts + * @param {{ attributes: Record }} opts.resource + * @param {Array<{ exportBatch(records: LogRecord[]): unknown, forceFlush?: () => Promise|void, shutdown?: () => Promise|void }>} [opts.exporters] + */ + constructor({ resource, exporters = [] }) { + this.resource = resource + this.exporters = exporters + } + + /** @param {LogRecord} record */ + exportRecord(record) { + if (this.exporters.length === 0) return + for (const exporter of this.exporters) exporter.exportBatch([record]) + } + + async forceFlush() { + await flushExporters(this.exporters) + } + + async shutdown() { + await shutdownExporters(this.exporters) + if (globalLoggerProvider === this) globalLoggerProvider = null + } +} + +export class MeterProvider { + /** + * @param {object} opts + * @param {{ attributes: Record }} opts.resource + * @param {Array<{ exportBatch(records: MetricRecord[]): unknown, forceFlush?: () => Promise|void, shutdown?: () => Promise|void }>} [opts.exporters] + */ + constructor({ resource, exporters = [] }) { + this.resource = resource + this.exporters = exporters + } + + /** @param {MetricRecord} record */ + exportRecord(record) { + if (this.exporters.length === 0) return + for (const exporter of this.exporters) exporter.exportBatch([record]) + } + + async forceFlush() { + await flushExporters(this.exporters) + } + + async shutdown() { + await shutdownExporters(this.exporters) + if (globalMeterProvider === this) globalMeterProvider = null + } +} + +class Tracer { + /** @param {string} name @param {string} [version] */ + constructor(name, version) { + this.name = name + this.version = version + } + + /** + * @param {string} name + * @param {object|((span: Span) => unknown)} [options] + * @param {(span: Span) => unknown} [fn] + */ + startActiveSpan(name, options, fn) { + const callback = typeof options === 'function' ? options : fn + const spanOptions = typeof options === 'object' && options !== null ? options : {} + const provider = globalTracerProvider + const parent = spanOptions.root ? null : (activeContext.getStore()?.span ?? null) + const span = new Span({ + name, + tracerName: this.name, + tracerVersion: this.version, + provider, + resource: provider?.resource ?? EMPTY_RESOURCE, + parent, + attributes: normalizeAttributes(Reflect.get(spanOptions, 'attributes')), + }) + if (!callback) return span + return activeContext.run({ span }, () => callback(span)) + } +} + +export class Span { + /** + * @param {object} opts + * @param {string} opts.name + * @param {string} opts.tracerName + * @param {string} [opts.tracerVersion] + * @param {TracerProvider|null} opts.provider + * @param {{ attributes: Record }} opts.resource + * @param {Span|null} opts.parent + * @param {Record} opts.attributes + */ + constructor({ name, tracerName, tracerVersion, provider, resource, parent, attributes }) { + this.name = name + this.tracerName = tracerName + this.tracerVersion = tracerVersion + this.provider = provider + this.resource = resource + this.parentSpanContext = parent ? parent.spanContext() : undefined + this.kind = 0 + this.attributes = { ...attributes } + this.events = [] + this.status = { code: SpanStatusCode.UNSET } + this.startTime = nowHrTime() + this.endTime = this.startTime + this._ended = false + this._context = { + traceId: parent ? parent.spanContext().traceId : randomHex(16), + spanId: randomHex(8), + traceFlags: 1, + } + } + + spanContext() { + return this._context + } + + /** @param {string} key @param {unknown} value */ + setAttribute(key, value) { + if (value !== undefined) this.attributes[key] = value + return this + } + + /** @param {Record} attrs */ + setAttributes(attrs) { + for (const [key, value] of Object.entries(attrs ?? {})) this.setAttribute(key, value) + return this + } + + /** @param {{ code: number, message?: string }} status */ + setStatus(status) { + this.status = { ...status } + return this + } + + /** @param {Error} error */ + recordException(error) { + this.addEvent('exception', { + 'exception.type': error.name, + 'exception.message': error.message, + ...(error.stack ? { 'exception.stacktrace': error.stack } : {}), + }) + } + + /** @param {string} name @param {Record} [attributes] */ + addEvent(name, attributes = {}) { + this.events.push({ name, time: nowHrTime(), attributes }) + } + + end() { + if (this._ended) return + this._ended = true + this.endTime = nowHrTime() + if (compareHrTime(this.endTime, this.startTime) <= 0) { + this.endTime = addNanos(this.startTime, 1_000_000) + } + this.provider?.exportSpan(this) + } +} + +class Logger { + /** @param {string} name @param {string} [version] */ + constructor(name, version) { + this.name = name + this.version = version + } + + /** + * @param {{ + * severityNumber?: number, + * severityText?: string, + * body?: unknown, + * attributes?: Record, + * }} record + */ + emit(record) { + const provider = globalLoggerProvider + if (!provider) return + const now = nowHrTime() + const activeSpan = trace.getActiveSpan() + provider.exportRecord({ + loggerName: this.name, + loggerVersion: this.version, + resource: provider.resource, + hrTime: now, + hrTimeObserved: now, + spanContext: activeSpan?.spanContext(), + severityNumber: record.severityNumber, + severityText: record.severityText, + body: record.body, + attributes: normalizeAttributes(record.attributes), + }) + } +} + +class Meter { + /** @param {string} name @param {string} [version] */ + constructor(name, version) { + this.name = name + this.version = version + } + + /** @param {string} name @param {{ description?: string, unit?: string }} [opts] */ + createCounter(name, opts = {}) { + return new Instrument({ meter: this, name, kind: 'counter', monotonic: true, ...opts }) + } + + /** @param {string} name @param {{ description?: string, unit?: string }} [opts] */ + createUpDownCounter(name, opts = {}) { + return new Instrument({ meter: this, name, kind: 'upDownCounter', monotonic: false, ...opts }) + } + + /** @param {string} name @param {{ description?: string, unit?: string }} [opts] */ + createGauge(name, opts = {}) { + return new Instrument({ meter: this, name, kind: 'gauge', monotonic: false, ...opts }) + } + + /** @param {string} name @param {{ description?: string, unit?: string }} [opts] */ + createHistogram(name, opts = {}) { + return new Instrument({ meter: this, name, kind: 'histogram', monotonic: false, ...opts }) + } +} + +class Instrument { + /** + * @param {object} opts + * @param {Meter} opts.meter + * @param {string} opts.name + * @param {'counter'|'upDownCounter'|'gauge'|'histogram'} opts.kind + * @param {boolean} opts.monotonic + * @param {string} [opts.description] + * @param {string} [opts.unit] + */ + constructor(opts) { + this.meter = opts.meter + this.name = opts.name + this.kind = opts.kind + this.description = opts.description + this.unit = opts.unit + this.monotonic = opts.monotonic + } + + /** @param {number} value @param {Record} [attributes] */ + add(value, attributes = {}) { + this._record(value, attributes) + } + + /** @param {number} value @param {Record} [attributes] */ + record(value, attributes = {}) { + this._record(value, attributes) + } + + /** @param {number} value @param {Record} attributes */ + _record(value, attributes) { + const provider = globalMeterProvider + if (!provider) return + const now = nowHrTime() + provider.exportRecord({ + meterName: this.meter.name, + meterVersion: this.meter.version, + resource: provider.resource, + name: this.name, + description: this.description, + unit: this.unit, + kind: this.kind, + monotonic: this.monotonic, + value, + attributes: normalizeAttributes(attributes), + startTime: now, + endTime: now, + }) + } +} + +const EMPTY_RESOURCE = Object.freeze({ attributes: Object.freeze({}) }) +const NOOP_TRACER_PROVIDER = Object.freeze({ resource: EMPTY_RESOURCE }) + +/** @param {unknown} value */ +export function normalizeAttributes(value) { + if (!value || typeof value !== 'object' || Array.isArray(value)) return {} + /** @type {Record} */ + const out = {} + for (const [key, attr] of Object.entries(/** @type {Record} */ (value))) { + if (attr !== undefined) out[key] = attr + } + return out +} + +export function getActiveSpan() { + return trace.getActiveSpan() +} + +/** @param {Array<{ forceFlush?: () => Promise|void }>} exporters */ +async function flushExporters(exporters) { + await Promise.allSettled(exporters.map((exporter) => exporter.forceFlush?.())) +} + +/** @param {Array<{ forceFlush?: () => Promise|void, shutdown?: () => Promise|void }>} exporters */ +async function shutdownExporters(exporters) { + await flushExporters(exporters) + await Promise.allSettled(exporters.map((exporter) => exporter.shutdown?.())) +} + +function nowHrTime() { + return nsToHrTime(nowUnixNano()) +} + +export function nowUnixNano() { + return BigInt(Math.round((performance.timeOrigin + performance.now()) * 1_000_000)) +} + +/** @param {bigint} ns */ +export function nsToHrTime(ns) { + const sec = ns / 1_000_000_000n + const nanos = ns % 1_000_000_000n + return [Number(sec), Number(nanos)] +} + +/** @param {[number, number]} hr */ +export function hrTimeToUnixNano(hr) { + return BigInt(hr[0]) * 1_000_000_000n + BigInt(hr[1]) +} + +/** @param {[number, number]} hr */ +export function hrTimeToIso(hr) { + return new Date(Number(hrTimeToUnixNano(hr) / 1_000_000n)).toISOString() +} + +/** @param {[number, number]} a @param {[number, number]} b */ +function compareHrTime(a, b) { + if (a[0] !== b[0]) return a[0] - b[0] + return a[1] - b[1] +} + +/** @param {[number, number]} hr @param {number} nanos */ +function addNanos(hr, nanos) { + return nsToHrTime(hrTimeToUnixNano(hr) + BigInt(nanos)) +} + +/** @param {number} bytes */ +function randomHex(bytes) { + return crypto.randomBytes(bytes).toString('hex') +} + +/** + * @typedef {Object} LogRecord + * @property {string} loggerName + * @property {string|undefined} loggerVersion + * @property {{ attributes: Record }} resource + * @property {[number, number]} hrTime + * @property {[number, number]} hrTimeObserved + * @property {{ traceId: string, spanId: string, traceFlags?: number }|undefined} spanContext + * @property {number|undefined} severityNumber + * @property {string|undefined} severityText + * @property {unknown} body + * @property {Record} attributes + */ + +/** + * @typedef {Object} MetricRecord + * @property {string} meterName + * @property {string|undefined} meterVersion + * @property {{ attributes: Record }} resource + * @property {string} name + * @property {string|undefined} description + * @property {string|undefined} unit + * @property {'counter'|'upDownCounter'|'gauge'|'histogram'} kind + * @property {boolean} monotonic + * @property {number} value + * @property {Record} attributes + * @property {[number, number]} startTime + * @property {[number, number]} endTime + */ diff --git a/src/core/observability/span_helpers.js b/src/core/observability/span_helpers.js index adf999d5..95ae8488 100644 --- a/src/core/observability/span_helpers.js +++ b/src/core/observability/span_helpers.js @@ -1,9 +1,8 @@ // @ts-check -import { context, SpanStatusCode, ROOT_CONTEXT } from '@opentelemetry/api' - import { buildAttrs } from './attrs.js' import { getTracer } from './tracer.js' +import { context, ROOT_CONTEXT, SpanStatusCode } from './runtime.js' /** * Run `fn` inside a span. Records the result on the span (status + any @@ -16,7 +15,7 @@ import { getTracer } from './tracer.js' * @template T * @param {string} name * @param {Record} attrs - * @param {(span: import('@opentelemetry/api').Span) => T|Promise} fn + * @param {(span: import('./runtime.js').Span) => T|Promise} fn * @param {{ component?: string }} [opts] * @returns {Promise} */ @@ -53,7 +52,7 @@ export async function withSpan(name, attrs, fn, opts = {}) { * @template T * @param {string} name * @param {Record} attrs - * @param {(span: import('@opentelemetry/api').Span) => T|Promise} fn + * @param {(span: import('./runtime.js').Span) => T|Promise} fn * @param {{ component?: string }} [opts] * @returns {Promise} */ diff --git a/src/core/observability/tracer.js b/src/core/observability/tracer.js index 8c0637a4..a76fbbfc 100644 --- a/src/core/observability/tracer.js +++ b/src/core/observability/tracer.js @@ -1,12 +1,9 @@ // @ts-check -import { trace, ProxyTracerProvider } from '@opentelemetry/api' -import { NodeTracerProvider } from '@opentelemetry/sdk-trace-node' -import { SimpleSpanProcessor } from '@opentelemetry/sdk-trace-base' -import { OTLPTraceExporter } from '@opentelemetry/exporter-trace-otlp-http' - import { JsonlSpanExporter } from './jsonl_exporters.js' import { devTelemetryDir } from './env.js' +import { OtlpSpanExporter } from './otlp_exporters.js' +import { trace, TracerProvider } from './runtime.js' const OTLP_EXPORT_TIMEOUT_MS = 1_000 @@ -24,38 +21,34 @@ const OTLP_EXPORT_TIMEOUT_MS = 1_000 * * @param {object} args * @param {import('./env.js').ObservabilityEnv} args.env - * @param {import('@opentelemetry/resources').Resource} args.resource - * @returns {{ provider: NodeTracerProvider|null, exporters: object[] }} + * @param {{ attributes: Record }} args.resource + * @returns {{ provider: TracerProvider|null, exporters: object[] }} */ export function installTracerProvider({ env, resource }) { /** @type {object[]} */ const exporters = [] - /** @type {import('@opentelemetry/sdk-trace-base').SpanProcessor[]} */ - const processors = [] if (env.devTelemetry) { const dir = devTelemetryDir(env.stateDir) const jsonlExporter = new JsonlSpanExporter({ dir }) - processors.push(new SimpleSpanProcessor(jsonlExporter)) exporters.push(jsonlExporter) } if (!env.devTelemetry && env.otlpEndpoint) { - const otlpExporter = new OTLPTraceExporter({ + const otlpExporter = new OtlpSpanExporter({ url: env.otlpEndpoint.replace(/\/$/, '') + '/v1/traces', timeoutMillis: OTLP_EXPORT_TIMEOUT_MS, }) - processors.push(new SimpleSpanProcessor(otlpExporter)) exporters.push(otlpExporter) } - if (processors.length === 0) { + if (exporters.length === 0) { return { provider: null, exporters: [] } } - const provider = new NodeTracerProvider({ + const provider = new TracerProvider({ resource, - spanProcessors: processors, + exporters, }) provider.register() return { provider, exporters } @@ -70,8 +63,7 @@ export function getTracer(component) { return trace.getTracer(`hypaware.${component}`) } -/** @returns {ProxyTracerProvider} */ +/** @returns {object} */ export function getActiveProvider() { - const provider = trace.getTracerProvider() - return /** @type {ProxyTracerProvider} */ (provider) + return trace.getTracerProvider() } diff --git a/src/core/plugin_install/install.js b/src/core/plugin_install/install.js index 785763e5..f55df1e1 100644 --- a/src/core/plugin_install/install.js +++ b/src/core/plugin_install/install.js @@ -2,12 +2,11 @@ import fs from 'node:fs/promises' -import { SpanStatusCode } from '@opentelemetry/api' - import { Attr, getKernelInstruments, getLogger, + SpanStatusCode, withSpan, } from '../observability/index.js' diff --git a/src/core/plugin_install/update_check.js b/src/core/plugin_install/update_check.js index 76654924..ec35055e 100644 --- a/src/core/plugin_install/update_check.js +++ b/src/core/plugin_install/update_check.js @@ -1,8 +1,6 @@ // @ts-check -import { SpanStatusCode } from '@opentelemetry/api' - -import { Attr, getKernelInstruments, withSpan } from '../observability/index.js' +import { Attr, getKernelInstruments, SpanStatusCode, withSpan } from '../observability/index.js' /** @typedef {import('../../../collectivus-plugin-kernel-types').PluginLockEntry} PluginLockEntry */ /** @typedef {import('../../../collectivus-plugin-kernel-types').PluginUpdateState} PluginUpdateState */ diff --git a/src/core/sinks/driver.js b/src/core/sinks/driver.js index a6bc8c7f..eef07644 100644 --- a/src/core/sinks/driver.js +++ b/src/core/sinks/driver.js @@ -228,7 +228,7 @@ export function createSinkDriver(opts) { * @param {string} batchId * @param {number} partitionsCount * @param {string} message - * @param {import('@opentelemetry/api').Span} span + * @param {import('../observability/runtime.js').Span} span */ function recordFailure(handle, batchId, partitionsCount, message, span) { instruments.sinkExportFailuresTotal.add(1, {