diff --git a/packages/appkit/package.json b/packages/appkit/package.json index 1f0869209..4992229f9 100644 --- a/packages/appkit/package.json +++ b/packages/appkit/package.json @@ -73,8 +73,8 @@ "@opentelemetry/resources": "2.8.0", "@opentelemetry/sdk-logs": "0.219.0", "@opentelemetry/sdk-metrics": "2.8.0", - "@opentelemetry/sdk-node": "0.219.0", "@opentelemetry/sdk-trace-base": "2.8.0", + "@opentelemetry/sdk-trace-node": "2.8.0", "@opentelemetry/semantic-conventions": "1.38.0", "@types/semver": "7.7.1", "apache-arrow": "21.1.0", diff --git a/packages/appkit/src/core/appkit.ts b/packages/appkit/src/core/appkit.ts index 201acf190..6195836a4 100644 --- a/packages/appkit/src/core/appkit.ts +++ b/packages/appkit/src/core/appkit.ts @@ -222,6 +222,12 @@ export class AppKit { const instance = new AppKit(mergedConfig); await Promise.all(instance.#setupPromises); + + // Build the global tracer provider now that every plugin's setup() has run + // and contributed any span processors. Deferred to here so a single provider + // carries all processors (OTLP + plugin-contributed); see TelemetryManager. + TelemetryManager.start(); + await instance.#context.emitLifecycle("setup:complete"); const handle = instance as unknown as PluginMap; diff --git a/packages/appkit/src/core/tests/appkit-as-user-exports.test.ts b/packages/appkit/src/core/tests/appkit-as-user-exports.test.ts index 7cbadcd32..454a3d664 100644 --- a/packages/appkit/src/core/tests/appkit-as-user-exports.test.ts +++ b/packages/appkit/src/core/tests/appkit-as-user-exports.test.ts @@ -44,6 +44,8 @@ vi.mock("../../telemetry", async () => { ...actual, TelemetryManager: { initialize: vi.fn(), + start: vi.fn(), + registerSpanProcessor: vi.fn(), getProvider: () => ({ getTracer: () => ({ startActiveSpan: vi.fn((_name: string, fn: (span: any) => any) => diff --git a/packages/appkit/src/telemetry/telemetry-manager.ts b/packages/appkit/src/telemetry/telemetry-manager.ts index a2642f2ab..463e6c13e 100644 --- a/packages/appkit/src/telemetry/telemetry-manager.ts +++ b/packages/appkit/src/telemetry/telemetry-manager.ts @@ -1,3 +1,5 @@ +import { metrics } from "@opentelemetry/api"; +import { logs } from "@opentelemetry/api-logs"; import { getNodeAutoInstrumentations } from "@opentelemetry/auto-instrumentations-node"; import { OTLPLogExporter } from "@opentelemetry/exporter-logs-otlp-proto"; import { OTLPMetricExporter } from "@opentelemetry/exporter-metrics-otlp-proto"; @@ -14,9 +16,19 @@ import { type Resource, resourceFromAttributes, } from "@opentelemetry/resources"; -import { BatchLogRecordProcessor } from "@opentelemetry/sdk-logs"; -import { PeriodicExportingMetricReader } from "@opentelemetry/sdk-metrics"; -import { NodeSDK } from "@opentelemetry/sdk-node"; +import { + BatchLogRecordProcessor, + LoggerProvider, +} from "@opentelemetry/sdk-logs"; +import { + MeterProvider, + PeriodicExportingMetricReader, +} from "@opentelemetry/sdk-metrics"; +import { + BatchSpanProcessor, + type SpanProcessor, +} from "@opentelemetry/sdk-trace-base"; +import { NodeTracerProvider } from "@opentelemetry/sdk-trace-node"; import { ATTR_SERVICE_NAME, ATTR_SERVICE_VERSION, @@ -29,12 +41,34 @@ import type { TelemetryConfig } from "./types"; const logger = createLogger("telemetry"); +/** + * Owns the app's OpenTelemetry providers, split into two phases so plugins can + * contribute trace span processors before the tracer provider is built. + * + * - `initialize()` runs at app bootstrap, before plugin setup. It registers the + * meter and logger providers eagerly, because OTel's metrics API has no lazy + * proxy: a counter/histogram bound against the NoOp meter (as every connector + * and the cache do in their constructors) stays NoOp for the process lifetime. + * It does NOT register a tracer provider. + * - `registerSpanProcessor()` is called by plugins during `setup()` to add a + * span processor (e.g. an MLflow exporter) to the not-yet-built tracer. + * - `start()` runs after all plugin `setup()` completes. It builds the single + * global tracer provider with the OTLP processor (if configured) plus every + * contributed processor. Deferring is safe for traces: OTel's ProxyTracer + * rebinds tracers obtained before registration, and no span is emitted during + * setup. + */ export class TelemetryManager { private static readonly DEFAULT_EXPORT_INTERVAL_MS = 10000; private static readonly DEFAULT_FALLBACK_APP_NAME = "databricks-app"; private static instance?: TelemetryManager; - private sdk?: NodeSDK; + private resource?: Resource; + private meterProvider?: MeterProvider; + private loggerProvider?: LoggerProvider; + private tracerProvider?: NodeTracerProvider; + private readonly spanProcessors: SpanProcessor[] = []; + private started = false; private shutdownPromise?: Promise; /** @@ -66,20 +100,49 @@ export class TelemetryManager { instance._initialize(config); } + /** + * Contribute a span processor to the not-yet-built tracer provider. Called by + * plugins during `setup()`. No-op with a warning once `start()` has run, since + * a started provider's processors are immutable in OTel JS 2.x. + */ + static registerSpanProcessor(processor: SpanProcessor): void { + TelemetryManager.getInstance()._registerSpanProcessor(processor); + } + + private _registerSpanProcessor(processor: SpanProcessor): void { + if (this.started) { + logger.warn( + "registerSpanProcessor called after start(); processor ignored. " + + "Contribute span processors during plugin setup().", + ); + return; + } + this.spanProcessors.push(processor); + } + + /** + * Phase 1: register the meter and logger providers eagerly (before plugin + * setup), so metric instruments bound in connector/cache constructors attach + * to real meters. The tracer provider is deferred to `start()`. + * + * When no OTLP endpoint is configured, meter/logger registration is skipped; + * a contributed span processor can still bring up tracing in `start()`. + */ private _initialize(config: Partial): void { - if (this.sdk) return; + if (this.resource) return; + this.resource = this.createResource(config); + // OTLP exporters need an endpoint. Without one there is nothing to export + // metrics/logs to, so skip those providers — but still capture the resource + // and let `start()` bring up a tracer if a plugin contributed a processor. if (!process.env.OTEL_EXPORTER_OTLP_ENDPOINT) { return; } try { - this.sdk = new NodeSDK({ - resource: this.createResource(config), - autoDetectResources: false, - sampler: new AppKitSampler(), - traceExporter: new OTLPTraceExporter({ headers: config.headers }), - metricReaders: [ + this.meterProvider = new MeterProvider({ + resource: this.resource, + readers: [ new PeriodicExportingMetricReader({ exporter: new OTLPMetricExporter({ headers: config.headers }), exportIntervalMillis: @@ -87,21 +150,72 @@ export class TelemetryManager { TelemetryManager.DEFAULT_EXPORT_INTERVAL_MS, }), ], - logRecordProcessors: [ + }); + metrics.setGlobalMeterProvider(this.meterProvider); + + this.loggerProvider = new LoggerProvider({ + resource: this.resource, + processors: [ new BatchLogRecordProcessor( new OTLPLogExporter({ headers: config.headers }), ), ], - instrumentations: this.getDefaultInstrumentations(), }); + logs.setGlobalLoggerProvider(this.loggerProvider); + + // The OTLP trace exporter is the first span processor; contributed + // processors join it in `start()`. + this.spanProcessors.push( + new BatchSpanProcessor( + new OTLPTraceExporter({ headers: config.headers }), + ), + ); - this.sdk.start(); - logger.debug("Initialized successfully"); + this.registerInstrumentations(this.getDefaultInstrumentations()); + logger.debug("Meter/logger providers initialized"); } catch (error) { logger.error("Failed to initialize: %O", error); } } + /** + * Phase 2: build and register the global tracer provider. Called by core + * after every plugin's `setup()` completes, so all contributed span + * processors are known. No-op when nothing needs tracing (no OTLP endpoint + * and no contributed processor), preserving "no telemetry unless configured". + * + * `NodeTracerProvider.register()` installs the async-hooks context manager and + * W3C propagators — the same wiring `NodeSDK.start()` did — so span nesting + * across awaits is preserved. + */ + static start(): void { + TelemetryManager.getInstance()._start(); + } + + private _start(): void { + if (this.started) return; + this.started = true; + + if (this.spanProcessors.length === 0) { + return; + } + + try { + this.tracerProvider = new NodeTracerProvider({ + resource: this.resource, + sampler: new AppKitSampler(), + spanProcessors: this.spanProcessors, + }); + this.tracerProvider.register(); + logger.debug( + "Tracer provider started with %d span processor(s)", + this.spanProcessors.length, + ); + } catch (error) { + logger.error("Failed to start tracer provider: %O", error); + } + } + /** * Register OpenTelemetry instrumentations. * Can be called at any time, but recommended to call in plugin constructor. @@ -159,23 +273,34 @@ export class TelemetryManager { } /** - * Flush and shut down the OpenTelemetry SDK. + * Flush and shut down the tracer, meter, and logger providers. * - * Idempotent: the SDK reference is cleared synchronously and concurrent + * Idempotent: the provider references are cleared synchronously and concurrent * or repeated calls await the same in-flight flush. Awaited by the core * lifecycle manager during graceful shutdown — that manager owns the * process signal handlers, so telemetry no longer registers its own. */ async shutdown(): Promise { - if (this.sdk) { - const sdk = this.sdk; - this.sdk = undefined; + const providers = [ + this.tracerProvider, + this.meterProvider, + this.loggerProvider, + ].filter((p): p is NonNullable => p !== undefined); + + if (providers.length > 0) { + this.tracerProvider = undefined; + this.meterProvider = undefined; + this.loggerProvider = undefined; this.shutdownPromise = (async () => { - try { - await sdk.shutdown(); - } catch (error) { - logger.error("Error shutting down: %O", error); - } + await Promise.all( + providers.map(async (provider) => { + try { + await provider.shutdown(); + } catch (error) { + logger.error("Error shutting down: %O", error); + } + }), + ); })(); } diff --git a/packages/appkit/src/telemetry/tests/telemetry-manager.test.ts b/packages/appkit/src/telemetry/tests/telemetry-manager.test.ts index 11b85d9bf..187276d1d 100644 --- a/packages/appkit/src/telemetry/tests/telemetry-manager.test.ts +++ b/packages/appkit/src/telemetry/tests/telemetry-manager.test.ts @@ -1,3 +1,5 @@ +import { context, metrics, trace } from "@opentelemetry/api"; +import { logs } from "@opentelemetry/api-logs"; import { afterEach, beforeEach, describe, expect, test, vi } from "vitest"; import { TelemetryManager } from "../telemetry-manager"; @@ -54,12 +56,21 @@ describe("TelemetryManager", () => { vi.clearAllMocks(); // @ts-expect-error - accessing private static property for testing TelemetryManager.instance = undefined; - // @ts-expect-error - accessing private static property for testing - TelemetryManager.shutdownRegistered = false; + // OTel's registerGlobal is allowOverride=false: a global registered by one + // test would make the next test's registration a silent no-op. Reset all + // global providers so each test starts clean. + trace.disable(); + metrics.disable(); + logs.disable(); + context.disable(); }); afterEach(() => { process.env = originalEnv; + trace.disable(); + metrics.disable(); + logs.disable(); + context.disable(); }); test("getInstance() should return singleton instance", () => { @@ -89,6 +100,7 @@ describe("TelemetryManager", () => { serviceName: "integration-test", serviceVersion: "1.0.0", }); + TelemetryManager.start(); const telemetryProvider = TelemetryManager.getProvider("test-plugin"); const tracer = telemetryProvider.getTracer(); @@ -186,6 +198,7 @@ describe("TelemetryManager", () => { serviceName: "span-test", serviceVersion: "1.0.0", }); + TelemetryManager.start(); const telemetryProvider = TelemetryManager.getProvider("span-test-plugin"); @@ -211,6 +224,7 @@ describe("TelemetryManager", () => { serviceName: "error-test", serviceVersion: "1.0.0", }); + TelemetryManager.start(); const telemetryProvider = TelemetryManager.getProvider("error-test-plugin"); @@ -224,4 +238,90 @@ describe("TelemetryManager", () => { ).rejects.toThrow("Test error in span"); }); }); + + describe("two-phase init (registerSpanProcessor + start)", () => { + /** Minimal SpanProcessor that records the names of spans it sees start. */ + function recordingProcessor() { + const startedSpans: string[] = []; + return { + startedSpans, + onStart: (span: { name: string }) => { + startedSpans.push(span.name); + }, + onEnd: () => {}, + forceFlush: () => Promise.resolve(), + shutdown: () => Promise.resolve(), + }; + } + + test("routes spans to a contributed processor with no OTLP endpoint", async () => { + process.env.OTEL_EXPORTER_OTLP_ENDPOINT = ""; + const processor = recordingProcessor(); + + TelemetryManager.initialize({ serviceName: "contrib-only" }); + TelemetryManager.registerSpanProcessor(processor as any); + TelemetryManager.start(); + + const tracer = TelemetryManager.getProvider("contrib-plugin").getTracer(); + await tracer.startActiveSpan("contributed.span", {}, async (span) => { + span.end(); + }); + + expect(processor.startedSpans).toContain("contributed.span"); + }); + + test("start() is idempotent", () => { + process.env.OTEL_EXPORTER_OTLP_ENDPOINT = ""; + const processor = recordingProcessor(); + + TelemetryManager.initialize({ serviceName: "idempotent" }); + TelemetryManager.registerSpanProcessor(processor as any); + TelemetryManager.start(); + + expect(() => TelemetryManager.start()).not.toThrow(); + }); + + test("registerSpanProcessor after start() is ignored (not attached)", async () => { + process.env.OTEL_EXPORTER_OTLP_ENDPOINT = ""; + const early = recordingProcessor(); + const late = recordingProcessor(); + + TelemetryManager.initialize({ serviceName: "late-register" }); + TelemetryManager.registerSpanProcessor(early as any); + TelemetryManager.start(); + TelemetryManager.registerSpanProcessor(late as any); + + const tracer = TelemetryManager.getProvider("late-plugin").getTracer(); + await tracer.startActiveSpan("post.start.span", {}, async (span) => { + span.end(); + }); + + expect(early.startedSpans).toContain("post.start.span"); + expect(late.startedSpans).toHaveLength(0); + }); + + test("start() with no OTLP endpoint and no processors is a no-op", () => { + process.env.OTEL_EXPORTER_OTLP_ENDPOINT = ""; + + TelemetryManager.initialize({ serviceName: "no-telemetry" }); + + expect(() => TelemetryManager.start()).not.toThrow(); + }); + + test("metrics obtained after initialize() (before start()) still record", () => { + // Guards the eager-metrics invariant: the meter provider is registered in + // initialize(), not start(), because OTel's metrics API has no lazy proxy. + process.env.OTEL_EXPORTER_OTLP_ENDPOINT = "http://localhost:4318"; + + TelemetryManager.initialize({ serviceName: "eager-metrics" }); + + // Instrument bound BEFORE start() — mirrors connector/cache constructors. + const meter = TelemetryManager.getProvider("metrics-plugin").getMeter(); + const counter = meter.createCounter("eager.counter"); + + TelemetryManager.start(); + + expect(() => counter.add(1, { label: "value" })).not.toThrow(); + }); + }); }); diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index 86f24f334..db074bf01 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -291,12 +291,12 @@ importers: '@opentelemetry/sdk-metrics': specifier: 2.8.0 version: 2.8.0(@opentelemetry/api@1.9.0) - '@opentelemetry/sdk-node': - specifier: 0.219.0 - version: 0.219.0(@opentelemetry/api@1.9.0) '@opentelemetry/sdk-trace-base': specifier: 2.8.0 version: 2.8.0(@opentelemetry/api@1.9.0) + '@opentelemetry/sdk-trace-node': + specifier: 2.8.0 + version: 2.8.0(@opentelemetry/api@1.9.0) '@opentelemetry/semantic-conventions': specifier: 1.38.0 version: 1.38.0