diff --git a/.changeset/otlp-logs-signal.md b/.changeset/otlp-logs-signal.md new file mode 100644 index 0000000..13148c7 --- /dev/null +++ b/.changeset/otlp-logs-signal.md @@ -0,0 +1,12 @@ +--- +'@smooai/observability': minor +--- + +Add the OTLP logs signal. The SDK now wires a `LoggerProvider` + +`BatchLogRecordProcessor` alongside traces and metrics — same endpoint +(`SMOOAI_OBSERVABILITY_ENDPOINT` → `/v1/logs`), same auth (static token or +M2M client_credentials via the per-request `AuthInjectingLogExporter`), same +enable path. App logs emitted through the standard `@opentelemetry/api-logs` +facade become OTLP log records correlated to the active span's trace_id / +span_id; when no logs endpoint resolves the global LoggerProvider stays the +api-logs no-op and stdout output is unchanged. diff --git a/packages/core/package.json b/packages/core/package.json index 5b5e27b..ef78c7f 100644 --- a/packages/core/package.json +++ b/packages/core/package.json @@ -70,12 +70,15 @@ "dependencies": { "@smooai/fetch": "^3.3.10", "@opentelemetry/api": "^1.9.0", + "@opentelemetry/api-logs": "^0.55.0", "@opentelemetry/auto-instrumentations-node": "^0.50.0", "@opentelemetry/core": "^1.30.0", + "@opentelemetry/exporter-logs-otlp-http": "^0.55.0", "@opentelemetry/exporter-metrics-otlp-http": "^0.55.0", "@opentelemetry/exporter-trace-otlp-http": "^0.55.0", "@opentelemetry/otlp-transformer": "^0.55.0", "@opentelemetry/resources": "^1.30.0", + "@opentelemetry/sdk-logs": "^0.55.0", "@opentelemetry/sdk-metrics": "^1.30.0", "@opentelemetry/sdk-node": "^0.55.0", "@opentelemetry/semantic-conventions": "^1.30.0" diff --git a/packages/core/src/bootstrap/index.ts b/packages/core/src/bootstrap/index.ts index a4492ad..1770a7b 100644 --- a/packages/core/src/bootstrap/index.ts +++ b/packages/core/src/bootstrap/index.ts @@ -16,11 +16,12 @@ * * SMOOAI_OBSERVABILITY_ENDPOINT — base URL of the ingest API (e.g. * "https://api.smoo.ai"). The SDK - * appends `/v1/traces` and - * `/v1/metrics`. May also be set per- + * appends `/v1/traces`, `/v1/metrics`, + * and `/v1/logs`. May also be set per- * signal via the standard * `OTEL_EXPORTER_OTLP_TRACES_ENDPOINT` - * / `_METRICS_ENDPOINT` env vars. + * / `_METRICS_ENDPOINT` / `_LOGS_ENDPOINT` + * env vars. * * ## Auth (pick ONE; pre-minted JWT wins if both are present) * @@ -106,6 +107,7 @@ export async function bootstrapObservability(overrides: Partial = endpoint: overrides.endpoint ?? process.env.SMOOAI_OBSERVABILITY_ENDPOINT, tracesEndpoint: overrides.tracesEndpoint ?? process.env.OTEL_EXPORTER_OTLP_TRACES_ENDPOINT, metricsEndpoint: overrides.metricsEndpoint ?? process.env.OTEL_EXPORTER_OTLP_METRICS_ENDPOINT, + logsEndpoint: overrides.logsEndpoint ?? process.env.OTEL_EXPORTER_OTLP_LOGS_ENDPOINT, token: overrides.token ?? process.env.SMOOAI_OBSERVABILITY_TOKEN, authUrl: overrides.authUrl ?? process.env.SMOOAI_OBSERVABILITY_AUTH_URL, clientId: overrides.clientId ?? process.env.SMOOAI_OBSERVABILITY_CLIENT_ID, @@ -163,6 +165,7 @@ export async function bootstrapObservability(overrides: Partial = const tracesEndpoint = env.tracesEndpoint ?? (env.endpoint ? `${stripTrailingSlash(env.endpoint)}/v1/traces` : undefined); const metricsEndpoint = env.metricsEndpoint ?? (env.endpoint ? `${stripTrailingSlash(env.endpoint)}/v1/metrics` : undefined); + const logsEndpoint = env.logsEndpoint ?? (env.endpoint ? `${stripTrailingSlash(env.endpoint)}/v1/logs` : undefined); // Set process.env so any *other* OTel-aware code in the process // (e.g. third-party libraries that read the env directly) sees the @@ -174,6 +177,9 @@ export async function bootstrapObservability(overrides: Partial = if (metricsEndpoint && !process.env.OTEL_EXPORTER_OTLP_METRICS_ENDPOINT) { process.env.OTEL_EXPORTER_OTLP_METRICS_ENDPOINT = metricsEndpoint; } + if (logsEndpoint && !process.env.OTEL_EXPORTER_OTLP_LOGS_ENDPOINT) { + process.env.OTEL_EXPORTER_OTLP_LOGS_ENDPOINT = logsEndpoint; + } const otelOptions: SetupOtelOptions = { serviceName: env.serviceName ?? 'smoo-service', @@ -181,6 +187,7 @@ export async function bootstrapObservability(overrides: Partial = release: env.release, otlpEndpoint: tracesEndpoint, otlpMetricsEndpoint: metricsEndpoint, + otlpLogsEndpoint: logsEndpoint, otlpHeaders: sharedHeaders, tokenProvider, }; @@ -212,6 +219,7 @@ export interface BootstrapEnv { endpoint?: string; tracesEndpoint?: string; metricsEndpoint?: string; + logsEndpoint?: string; token?: string; authUrl?: string; clientId?: string; diff --git a/packages/core/src/otel/__tests__/auth-injecting-log-exporter.test.ts b/packages/core/src/otel/__tests__/auth-injecting-log-exporter.test.ts new file mode 100644 index 0000000..309fb0c --- /dev/null +++ b/packages/core/src/otel/__tests__/auth-injecting-log-exporter.test.ts @@ -0,0 +1,85 @@ +import { SeverityNumber } from '@opentelemetry/api-logs'; +import { ExportResultCode } from '@opentelemetry/core'; +import { Resource } from '@opentelemetry/resources'; +import type { ReadableLogRecord } from '@opentelemetry/sdk-logs'; +import { describe, expect, it, vi } from 'vitest'; +import type { TokenProvider } from '../../auth/token-provider'; +import { AuthInjectingLogExporter } from '../auth-injecting-exporter'; + +// A serializable ReadableLogRecord stub. JsonLogsSerializer reads timing, +// severity, body, attributes, the (optional) span context, resource, and +// scope — supply all of them so a real serialize + export runs. The span +// context here proves a log carrying real W3C ids round-trips through the +// exporter unchanged. +const fakeLog = { + hrTime: [1609459200, 0] as [number, number], + hrTimeObserved: [1609459200, 0] as [number, number], + severityNumber: SeverityNumber.INFO, + severityText: 'info', + body: 'hello from a span', + attributes: { foo: 'bar' }, + droppedAttributesCount: 0, + spanContext: { traceId: '0af7651916cd43dd8448eb211c80319c', spanId: 'b7ad6b7169203331', traceFlags: 1 }, + resource: Resource.empty(), + instrumentationScope: { name: '@smooai/logger', version: '0.0.0' }, +} as unknown as ReadableLogRecord; + +function tokenProviderStub(): TokenProvider & { invalidate: ReturnType } { + const invalidate = vi.fn(); + return { + getAccessToken: vi.fn().mockResolvedValue('tok-123'), + invalidate, + } as unknown as TokenProvider & { invalidate: ReturnType }; +} + +describe('AuthInjectingLogExporter (transport via @smooai/fetch seam)', () => { + it('posts the serialized log with a fresh Bearer token and reports SUCCESS on 2xx', async () => { + const fetcher = vi.fn().mockResolvedValue(new Response('', { status: 200 })); + const tp = tokenProviderStub(); + const exporter = new AuthInjectingLogExporter({ url: 'https://api.smoo.ai/v1/logs', tokenProvider: tp, fetcher }); + + const result = await new Promise<{ code: ExportResultCode }>((resolve) => { + exporter.export([fakeLog], resolve); + }); + + expect(result.code).toBe(ExportResultCode.SUCCESS); + expect(fetcher).toHaveBeenCalledOnce(); + const [url, init] = fetcher.mock.calls[0]!; + expect(url).toBe('https://api.smoo.ai/v1/logs'); + expect(init.method).toBe('POST'); + expect((init.headers as Record).authorization).toBe('Bearer tok-123'); + // The serialized OTLP body carries the record's real W3C trace/span ids. + expect(init.body as string).toContain('0af7651916cd43dd8448eb211c80319c'); + expect(init.body as string).toContain('b7ad6b7169203331'); + }); + + it('invalidates the token and retries once on 401, then succeeds', async () => { + const fetcher = vi + .fn() + .mockResolvedValueOnce(new Response('', { status: 401 })) + .mockResolvedValueOnce(new Response('', { status: 200 })); + const tp = tokenProviderStub(); + const exporter = new AuthInjectingLogExporter({ url: 'https://api.smoo.ai/v1/logs', tokenProvider: tp, fetcher }); + + const result = await new Promise<{ code: ExportResultCode }>((resolve) => { + exporter.export([fakeLog], resolve); + }); + + expect(result.code).toBe(ExportResultCode.SUCCESS); + expect(tp.invalidate).toHaveBeenCalledOnce(); + expect(fetcher).toHaveBeenCalledTimes(2); + }); + + it('reports FAILED when the response is non-ok', async () => { + const fetcher = vi.fn().mockResolvedValue(new Response('boom', { status: 503 })); + const tp = tokenProviderStub(); + const exporter = new AuthInjectingLogExporter({ url: 'https://api.smoo.ai/v1/logs', tokenProvider: tp, fetcher }); + + const result = await new Promise<{ code: ExportResultCode; error?: Error }>((resolve) => { + exporter.export([fakeLog], resolve); + }); + + expect(result.code).toBe(ExportResultCode.FAILED); + expect(result.error?.message).toContain('503'); + }); +}); diff --git a/packages/core/src/otel/__tests__/logs-correlation.test.ts b/packages/core/src/otel/__tests__/logs-correlation.test.ts new file mode 100644 index 0000000..294e915 --- /dev/null +++ b/packages/core/src/otel/__tests__/logs-correlation.test.ts @@ -0,0 +1,63 @@ +import { context, trace } from '@opentelemetry/api'; +import { logs } from '@opentelemetry/api-logs'; +import { AsyncHooksContextManager } from '@opentelemetry/context-async-hooks'; +import { InMemoryLogRecordExporter, LoggerProvider, SimpleLogRecordProcessor } from '@opentelemetry/sdk-logs'; +import { BasicTracerProvider } from '@opentelemetry/sdk-trace-base'; +import { afterAll, afterEach, beforeAll, describe, expect, it } from 'vitest'; + +// End-to-end proof of the trace↔log correlation the whole logs signal rides +// on: a log emitted through the standard `@opentelemetry/api-logs` facade +// (the same facade @smooai/logger bridges to) while a span is active MUST +// carry that span's real W3C trace_id / span_id. If it doesn't, every logger +// line lands uncorrelated in the product and the feature is pointless — so we +// exercise the real LoggerProvider + a real active span, not a stub. +describe('logs ↔ trace correlation via the api-logs facade', () => { + const memoryExporter = new InMemoryLogRecordExporter(); + const loggerProvider = new LoggerProvider(); + const tracerProvider = new BasicTracerProvider(); + const contextManager = new AsyncHooksContextManager(); + + beforeAll(() => { + loggerProvider.addLogRecordProcessor(new SimpleLogRecordProcessor(memoryExporter)); + logs.setGlobalLoggerProvider(loggerProvider); + // Without a real ContextManager `context.with(...)` is a noop and no + // span is ever "active" — register async-hooks so the span propagates. + context.setGlobalContextManager(contextManager.enable()); + }); + + afterEach(() => { + memoryExporter.getFinishedLogRecords().length = 0; + }); + + afterAll(() => { + contextManager.disable(); + }); + + it('stamps the active span trace_id + span_id onto a log emitted inside it', () => { + const tracer = tracerProvider.getTracer('test'); + const span = tracer.startSpan('unit-of-work'); + const expected = span.spanContext(); + + context.with(trace.setSpan(context.active(), span), () => { + logs.getLogger('@smooai/logger').emit({ severityText: 'info', body: 'inside the span' }); + }); + span.end(); + + const records = memoryExporter.getFinishedLogRecords(); + expect(records).toHaveLength(1); + expect(records[0]!.spanContext?.traceId).toBe(expected.traceId); + expect(records[0]!.spanContext?.spanId).toBe(expected.spanId); + // Real W3C shapes — 32 / 16 lowercase hex. + expect(records[0]!.spanContext?.traceId).toMatch(/^[0-9a-f]{32}$/); + expect(records[0]!.spanContext?.spanId).toMatch(/^[0-9a-f]{16}$/); + }); + + it('emits without a span context when no span is active (graceful, still a record)', () => { + logs.getLogger('@smooai/logger').emit({ severityText: 'warn', body: 'no span here' }); + + const records = memoryExporter.getFinishedLogRecords(); + expect(records).toHaveLength(1); + expect(records[0]!.spanContext).toBeUndefined(); + expect(records[0]!.body).toBe('no span here'); + }); +}); diff --git a/packages/core/src/otel/__tests__/setup-otel-sdk.test.ts b/packages/core/src/otel/__tests__/setup-otel-sdk.test.ts index ec1b635..72fdda9 100644 --- a/packages/core/src/otel/__tests__/setup-otel-sdk.test.ts +++ b/packages/core/src/otel/__tests__/setup-otel-sdk.test.ts @@ -39,4 +39,14 @@ describe('setupOtelSdk', () => { const handle = setupOtelSdk({ serviceName: 'test', skipStart: true, disableAutoInstrumentations: true }); expect(handle.sdk).toBeDefined(); }); + + it('wires the logs signal (LoggerProvider) when a logs endpoint is given', () => { + const handle = setupOtelSdk({ serviceName: 'test', skipStart: true, otlpLogsEndpoint: 'https://api.smoo.ai/v1/logs' }); + expect(handle.loggerProvider).toBeDefined(); + }); + + it('leaves the logs signal unwired (no LoggerProvider) when no logs endpoint resolves', () => { + const handle = setupOtelSdk({ serviceName: 'test', skipStart: true }); + expect(handle.loggerProvider).toBeUndefined(); + }); }); diff --git a/packages/core/src/otel/auth-injecting-exporter.ts b/packages/core/src/otel/auth-injecting-exporter.ts index 3f27885..8174bfa 100644 --- a/packages/core/src/otel/auth-injecting-exporter.ts +++ b/packages/core/src/otel/auth-injecting-exporter.ts @@ -32,9 +32,10 @@ import smooFetch, { HTTPResponseError } from '@smooai/fetch'; import { ExportResult, ExportResultCode } from '@opentelemetry/core'; -import { JsonMetricsSerializer, JsonTraceSerializer } from '@opentelemetry/otlp-transformer'; +import { JsonLogsSerializer, JsonMetricsSerializer, JsonTraceSerializer } from '@opentelemetry/otlp-transformer'; import type { ResourceMetrics, PushMetricExporter } from '@opentelemetry/sdk-metrics'; import type { ReadableSpan, SpanExporter } from '@opentelemetry/sdk-trace-base'; +import type { LogRecordExporter, ReadableLogRecord } from '@opentelemetry/sdk-logs'; import type { TokenProvider } from '../auth/token-provider'; /** @@ -208,3 +209,13 @@ export class AuthInjectingMetricExporter extends BaseAuthInjectingExporter implements LogRecordExporter { + protected serialize(items: ReadableLogRecord[]): Uint8Array { + return JsonLogsSerializer.serializeRequest(items) ?? new Uint8Array(); + } + + export(logs: ReadableLogRecord[], resultCallback: (result: ExportResult) => void): void { + this.dispatch(logs, resultCallback); + } +} diff --git a/packages/core/src/otel/setup-otel-sdk.ts b/packages/core/src/otel/setup-otel-sdk.ts index ac43451..4df36e8 100644 --- a/packages/core/src/otel/setup-otel-sdk.ts +++ b/packages/core/src/otel/setup-otel-sdk.ts @@ -10,15 +10,18 @@ * tests and lazy boots don't accidentally double-register exporters. */ +import { logs } from '@opentelemetry/api-logs'; import { getNodeAutoInstrumentations } from '@opentelemetry/auto-instrumentations-node'; +import { OTLPLogExporter } from '@opentelemetry/exporter-logs-otlp-http'; import { OTLPMetricExporter } from '@opentelemetry/exporter-metrics-otlp-http'; import { OTLPTraceExporter } from '@opentelemetry/exporter-trace-otlp-http'; import { Resource } from '@opentelemetry/resources'; +import { BatchLogRecordProcessor, LoggerProvider } from '@opentelemetry/sdk-logs'; import { NodeSDK } from '@opentelemetry/sdk-node'; import { PeriodicExportingMetricReader } from '@opentelemetry/sdk-metrics'; import { ATTR_SERVICE_NAME, ATTR_SERVICE_VERSION } from '@opentelemetry/semantic-conventions'; import type { TokenProvider } from '../auth/token-provider'; -import { AuthInjectingMetricExporter, AuthInjectingTraceExporter } from './auth-injecting-exporter'; +import { AuthInjectingLogExporter, AuthInjectingMetricExporter, AuthInjectingTraceExporter } from './auth-injecting-exporter'; export interface SetupOtelOptions { /** Service name surfaced in spans (e.g. 'smoo-backend', 'smoo-web'). */ @@ -69,11 +72,21 @@ export interface SetupOtelOptions { * from the trace URL or env vars. */ otlpMetricsEndpoint?: string; + /** + * Endpoint for logs, if you want it different from the trace endpoint + * base. Defaults to `OTEL_EXPORTER_OTLP_LOGS_ENDPOINT` / `_ENDPOINT` env + * vars. When neither this nor the env vars resolve to a logs endpoint the + * logs signal is not wired (no LoggerProvider is registered), so app + * logger lines emitted through `@opentelemetry/api-logs` stay no-ops. + */ + otlpLogsEndpoint?: string; } export interface OtelSdkHandle { /** The underlying NodeSDK so callers can shutdown / flush in their own lifecycle. */ sdk: NodeSDK; + /** The logs LoggerProvider, if the logs signal was wired (a logs endpoint resolved). */ + loggerProvider?: LoggerProvider; /** * Force-flush spans now. Returns when the exporter has acknowledged or * the timeout elapses. Wired by the host into SIGTERM / `beforeExit`. @@ -92,6 +105,7 @@ export function setupOtelSdk(options: SetupOtelOptions): OtelSdkHandle { const traceEndpoint = options.otlpEndpoint ?? process.env.OTEL_EXPORTER_OTLP_TRACES_ENDPOINT ?? process.env.OTEL_EXPORTER_OTLP_ENDPOINT; const metricEndpoint = options.otlpMetricsEndpoint ?? process.env.OTEL_EXPORTER_OTLP_METRICS_ENDPOINT ?? process.env.OTEL_EXPORTER_OTLP_ENDPOINT; + const logEndpoint = options.otlpLogsEndpoint ?? process.env.OTEL_EXPORTER_OTLP_LOGS_ENDPOINT ?? process.env.OTEL_EXPORTER_OTLP_ENDPOINT; // SMOODEV-1206: when a TokenProvider is passed, route through the // auth-injecting exporters. They ask the TokenProvider for a fresh @@ -128,6 +142,23 @@ export function setupOtelSdk(options: SetupOtelOptions): OtelSdkHandle { ...(options.environment ? { 'deployment.environment.name': options.environment } : {}), }); + // Logs signal — same OTLP/HTTP transport + same auth path as traces. + // Only wired when a logs endpoint resolves; otherwise the global + // LoggerProvider stays the api-logs NoopLoggerProvider so app logger + // lines emitted through `@opentelemetry/api-logs` are cheap no-ops. + // The LogRecord SDK stamps trace_id/span_id from the active span context + // at emit time (W3C ids), giving trace↔log correlation for free. + let loggerProvider: LoggerProvider | undefined; + if (logEndpoint) { + const logExporter = + options.tokenProvider + ? new AuthInjectingLogExporter({ url: logEndpoint, tokenProvider: options.tokenProvider, staticHeaders: options.otlpHeaders }) + : new OTLPLogExporter({ url: logEndpoint, headers: options.otlpHeaders }); + loggerProvider = new LoggerProvider({ resource }); + loggerProvider.addLogRecordProcessor(new BatchLogRecordProcessor(logExporter)); + logs.setGlobalLoggerProvider(loggerProvider); + } + const instrumentations = options.disableAutoInstrumentations ? [] : [ @@ -152,6 +183,7 @@ export function setupOtelSdk(options: SetupOtelOptions): OtelSdkHandle { const handle: OtelSdkHandle = { sdk, + loggerProvider, async flush(timeoutMs = 2_000) { // NodeSDK doesn't expose a public flush; shutdown drains the exporter // queue. Wrap in Promise.race so a slow exporter doesn't stall SIGTERM. @@ -163,6 +195,8 @@ export function setupOtelSdk(options: SetupOtelOptions): OtelSdkHandle { if (typeof exporterWithFlush.forceFlush === 'function') { await exporterWithFlush.forceFlush(); } + // Drain the batched log records too so logs aren't lost on SIGTERM. + await loggerProvider?.forceFlush(); } catch { /* swallow — flush is best-effort */ } @@ -178,6 +212,7 @@ export function setupOtelSdk(options: SetupOtelOptions): OtelSdkHandle { async shutdown() { try { await sdk.shutdown(); + await loggerProvider?.shutdown(); } catch { /* swallow */ } finally { diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index cb12510..08fde39 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -33,12 +33,18 @@ importers: '@opentelemetry/api': specifier: ^1.9.0 version: 1.9.1 + '@opentelemetry/api-logs': + specifier: ^0.55.0 + version: 0.55.0 '@opentelemetry/auto-instrumentations-node': specifier: ^0.50.0 version: 0.50.2(@opentelemetry/api@1.9.1) '@opentelemetry/core': specifier: ^1.30.0 version: 1.30.1(@opentelemetry/api@1.9.1) + '@opentelemetry/exporter-logs-otlp-http': + specifier: ^0.55.0 + version: 0.55.0(@opentelemetry/api@1.9.1) '@opentelemetry/exporter-metrics-otlp-http': specifier: ^0.55.0 version: 0.55.0(@opentelemetry/api@1.9.1) @@ -51,6 +57,9 @@ importers: '@opentelemetry/resources': specifier: ^1.30.0 version: 1.30.1(@opentelemetry/api@1.9.1) + '@opentelemetry/sdk-logs': + specifier: ^0.55.0 + version: 0.55.0(@opentelemetry/api@1.9.1) '@opentelemetry/sdk-metrics': specifier: ^1.30.0 version: 1.30.1(@opentelemetry/api@1.9.1)