Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
12 changes: 12 additions & 0 deletions .changeset/otlp-logs-signal.md
Original file line number Diff line number Diff line change
@@ -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.
3 changes: 3 additions & 0 deletions packages/core/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down
14 changes: 11 additions & 3 deletions packages/core/src/bootstrap/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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)
*
Expand Down Expand Up @@ -106,6 +107,7 @@ export async function bootstrapObservability(overrides: Partial<BootstrapEnv> =
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,
Expand Down Expand Up @@ -163,6 +165,7 @@ export async function bootstrapObservability(overrides: Partial<BootstrapEnv> =

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
Expand All @@ -174,13 +177,17 @@ export async function bootstrapObservability(overrides: Partial<BootstrapEnv> =
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',
environment: env.environment,
release: env.release,
otlpEndpoint: tracesEndpoint,
otlpMetricsEndpoint: metricsEndpoint,
otlpLogsEndpoint: logsEndpoint,
otlpHeaders: sharedHeaders,
tokenProvider,
};
Expand Down Expand Up @@ -212,6 +219,7 @@ export interface BootstrapEnv {
endpoint?: string;
tracesEndpoint?: string;
metricsEndpoint?: string;
logsEndpoint?: string;
token?: string;
authUrl?: string;
clientId?: string;
Expand Down
Original file line number Diff line number Diff line change
@@ -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<typeof vi.fn> } {
const invalidate = vi.fn();
return {
getAccessToken: vi.fn().mockResolvedValue('tok-123'),
invalidate,
} as unknown as TokenProvider & { invalidate: ReturnType<typeof vi.fn> };
}

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<string, string>).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');
});
});
63 changes: 63 additions & 0 deletions packages/core/src/otel/__tests__/logs-correlation.test.ts
Original file line number Diff line number Diff line change
@@ -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');
});
});
10 changes: 10 additions & 0 deletions packages/core/src/otel/__tests__/setup-otel-sdk.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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();
});
});
13 changes: 12 additions & 1 deletion packages/core/src/otel/auth-injecting-exporter.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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';

/**
Expand Down Expand Up @@ -208,3 +209,13 @@ export class AuthInjectingMetricExporter extends BaseAuthInjectingExporter<Resou
selectAggregation = undefined;
selectAggregationTemporality = undefined;
}

export class AuthInjectingLogExporter extends BaseAuthInjectingExporter<ReadableLogRecord> 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);
}
}
37 changes: 36 additions & 1 deletion packages/core/src/otel/setup-otel-sdk.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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'). */
Expand Down Expand Up @@ -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`.
Expand All @@ -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
Expand Down Expand Up @@ -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
? []
: [
Expand All @@ -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.
Expand All @@ -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 */
}
Expand All @@ -178,6 +212,7 @@ export function setupOtelSdk(options: SetupOtelOptions): OtelSdkHandle {
async shutdown() {
try {
await sdk.shutdown();
await loggerProvider?.shutdown();
} catch {
/* swallow */
} finally {
Expand Down
Loading
Loading