import { isSensitiveTelemetryFieldName, redactTelemetryText, } from './redaction'; export { isSensitiveTelemetryFieldName } from './redaction'; export const TELEMETRY_SCHEMA_VERSION = 1 as const; export type TelemetryScalar = string | number | boolean | null | undefined; export type TelemetryValue = | Exclude | readonly TelemetryValue[] | { readonly [key: string]: TelemetryValue }; export type TelemetryFields = Readonly>; export type TelemetryContext = { environment?: string | null; release?: string | null; orgId?: string | null; orgName?: string | null; userId?: string | null; apiKeyId?: string | null; authType?: string | null; requestId?: string | null; runId?: string | null; workflowId?: string | null; workflowName?: string | null; workflowOrigin?: string | null; playName?: string | null; machineId?: string | null; }; export type TelemetryLevel = 'debug' | 'info' | 'warn' | 'error'; export type TelemetryMetricKind = 'counter' | 'gauge' | 'distribution'; export type TelemetrySpanOutcome = 'ok' | 'error'; export const TELEMETRY_TRIAGE_SCHEMA_VERSION = 1 as const; export const DEFAULT_ERROR_TRIAGE_CLASS = 'unclassified_error' as const; const TELEMETRY_TRIAGE_CLASS_PATTERN = /^[a-z][a-z0-9_]{0,63}$/; export type TelemetryTriage = { schema_version: typeof TELEMETRY_TRIAGE_SCHEMA_VERSION; reportable: boolean; class: string; }; type TelemetryEventBase = { telemetrySchemaVersion: typeof TELEMETRY_SCHEMA_VERSION; timestamp: string; service: string; component: string | null; event: string; tag: string; triage?: TelemetryTriage; } & Omit & Record; export type TelemetryLogEvent = TelemetryEventBase & { telemetryKind: 'log'; telemetryLevel: TelemetryLevel; errorType?: string; errorMessage?: string; }; export type TelemetryMetricEvent = TelemetryEventBase & { telemetryKind: 'metric'; telemetryMetricKind: TelemetryMetricKind; telemetryMetricValue: number; telemetryMetricUnit: string | null; }; export type TelemetrySpanEvent = TelemetryEventBase & { telemetryKind: 'span'; telemetryLevel: 'info' | 'warn' | 'error'; spanDurationMs: number; spanOutcome: TelemetrySpanOutcome; errorType?: string; errorMessage?: string; }; export type TelemetryEvent = | TelemetryLogEvent | TelemetryMetricEvent | TelemetrySpanEvent; export type TelemetryAdapter = { emit(event: TelemetryEvent): void | Promise; }; export type TelemetryEventOptions = { tag?: string }; export type TelemetryTimeErrorOptions = { telemetryLevel: 'warn' | 'error'; fields?: TelemetryFields; }; export type Telemetry = { child(input: { component?: string; context?: TelemetryContext }): Telemetry; debug( event: string, fields?: TelemetryFields, options?: TelemetryEventOptions, ): Promise; info( event: string, fields?: TelemetryFields, options?: TelemetryEventOptions, ): Promise; warn( event: string, fields?: TelemetryFields, options?: TelemetryEventOptions, ): Promise; error( event: string, error?: unknown, fields?: TelemetryFields, options?: TelemetryEventOptions, ): Promise; counter( metric: string, value?: number, fields?: TelemetryFields, ): Promise; gauge(metric: string, value: number, fields?: TelemetryFields): Promise; distribution( metric: string, value: number, options?: { unit?: string; fields?: TelemetryFields }, ): Promise; time( span: string, fields: TelemetryFields | undefined, fn: () => Promise | T, options?: { onError?: (error: unknown) => TelemetryTimeErrorOptions }, ): Promise; }; type CreateTelemetryInput = { service: string; component?: string; context?: TelemetryContext; adapter?: TelemetryAdapter; now?: () => Date; }; const MAX_FIELD_STRING_LENGTH = 8_000; const MAX_FIELD_DEPTH = 8; const MAX_ARRAY_ITEMS = 100; const MAX_OBJECT_KEYS = 100; const RESERVED_FIELDS = new Set([ 'telemetrySchemaVersion', 'timestamp', 'service', 'component', 'event', 'tag', 'telemetryKind', 'telemetryLevel', 'telemetryMetricKind', 'telemetryMetricValue', 'telemetryMetricUnit', 'spanDurationMs', 'spanOutcome', 'environment', 'release', 'orgId', 'orgName', 'userId', 'apiKeyId', 'authType', 'requestId', 'runId', 'workflowId', 'workflowName', 'workflowOrigin', 'playName', 'machineId', ]); function errorTriage(fields: TelemetryFields | undefined): TelemetryTriage { const candidate = fields?.triage; const triageClass = candidate && typeof candidate === 'object' && !Array.isArray(candidate) && typeof (candidate as { class?: unknown }).class === 'string' && TELEMETRY_TRIAGE_CLASS_PATTERN.test( (candidate as { class: string }).class.trim(), ) ? (candidate as { class: string }).class.trim() : DEFAULT_ERROR_TRIAGE_CLASS; return { schema_version: TELEMETRY_TRIAGE_SCHEMA_VERSION, reportable: true, class: triageClass, }; } function normalizeName(value: string, label: string): string { const normalized = value.trim(); if (!normalized) throw new Error(`${label} must not be empty.`); return normalized; } export function redactTelemetryString(value: string): string { return redactTelemetryText(value); } function normalizeFieldValue( key: string, value: unknown, depth = 0, seen = new WeakSet(), ): TelemetryValue | undefined { if (isSensitiveTelemetryFieldName(key)) return '[REDACTED]'; if (value === null) return null; if (typeof value === 'boolean') return value; if (typeof value === 'number') { return Number.isFinite(value) ? value : undefined; } if (typeof value === 'string') { const redacted = redactTelemetryString(value); if (redacted.length <= MAX_FIELD_STRING_LENGTH) return redacted; return `${redacted.slice(0, MAX_FIELD_STRING_LENGTH)}…`; } if (typeof value === 'bigint') return String(value); if (value instanceof Date) return value.toISOString(); if (value instanceof Error) { return { name: value.name || 'Error', message: normalizeFieldValue( 'message', value.message, depth + 1, seen, ) as string, }; } if (typeof value !== 'object' || value === undefined) return undefined; if (depth >= MAX_FIELD_DEPTH) return '[max-depth]'; if (seen.has(value)) return '[circular]'; seen.add(value); if (Array.isArray(value)) { const normalized = value .slice(0, MAX_ARRAY_ITEMS) .map( (item) => normalizeFieldValue('item', item, depth + 1, seen) ?? null, ); if (value.length > MAX_ARRAY_ITEMS) normalized.push('[truncated]'); seen.delete(value); return normalized; } const normalized: Record = {}; const entries = Object.entries(value).slice(0, MAX_OBJECT_KEYS); for (const [childKey, childValue] of entries) { const item = normalizeFieldValue(childKey, childValue, depth + 1, seen); if (item !== undefined) normalized[childKey] = item; } if (Object.keys(value).length > MAX_OBJECT_KEYS) { normalized.__truncated__ = true; } seen.delete(value); return normalized; } function normalizeFields(fields: TelemetryFields | undefined) { const normalized: Record = {}; for (const [key, value] of Object.entries(fields ?? {})) { if (RESERVED_FIELDS.has(key) || value === undefined) continue; const normalizedValue = normalizeFieldValue(key, value); if (normalizedValue !== undefined) normalized[key] = normalizedValue; } return normalized; } function normalizeContext(context: TelemetryContext | undefined) { return Object.fromEntries( Object.entries(context ?? {}).filter(([, value]) => value !== undefined), ) as TelemetryContext; } function errorFields( error: unknown, ): Partial> { if (error === undefined || error === null) return {}; if (error instanceof Error) { return { errorType: error.name || 'Error', errorMessage: redactTelemetryString(error.message), }; } return { errorType: 'Error', errorMessage: redactTelemetryString(String(error)), }; } export function createConsoleTelemetryAdapter(input?: { console?: Pick; }): TelemetryAdapter { const target = input?.console ?? console; return { emit(event) { const line = JSON.stringify(event); if (event.telemetryKind === 'log' || event.telemetryKind === 'span') target[event.telemetryLevel](line); else target.info(line); }, }; } export function composeTelemetryAdapters( ...adapters: readonly TelemetryAdapter[] ): TelemetryAdapter { return { async emit(event) { await Promise.allSettled( adapters.map((adapter) => Promise.resolve().then(() => adapter.emit(event)), ), ); }, }; } export function createTelemetry(input: CreateTelemetryInput): Telemetry { const service = normalizeName(input.service, 'Telemetry service'); const component = input.component?.trim() || null; const context = normalizeContext(input.context); const adapter = input.adapter ?? createConsoleTelemetryAdapter(); const now = input.now ?? (() => new Date()); function base( event: string, options?: TelemetryEventOptions, ): TelemetryEventBase { const eventName = normalizeName(event, 'Telemetry event'); return { telemetrySchemaVersion: TELEMETRY_SCHEMA_VERSION, timestamp: now().toISOString(), service, component, event: eventName, tag: options?.tag?.trim() || `[${eventName}]`, ...context, }; } async function publish(event: TelemetryEvent): Promise { try { await adapter.emit(event); } catch { // Product execution must not depend on operational telemetry delivery. } } async function log( level: TelemetryLevel, event: string, fields?: TelemetryFields, options?: TelemetryEventOptions, ) { await publish({ ...base(event, options), ...normalizeFields(fields), telemetryKind: 'log', telemetryLevel: level, }); } async function metric( metricKind: TelemetryMetricKind, name: string, value: number, unit: string | undefined, fields: TelemetryFields | undefined, ) { if (!Number.isFinite(value)) { throw new Error(`Telemetry metric ${name} must have a finite value.`); } await publish({ ...base(name), ...normalizeFields(fields), telemetryKind: 'metric', telemetryMetricKind: metricKind, telemetryMetricValue: value, telemetryMetricUnit: unit?.trim() || null, }); } return { child(childInput) { return createTelemetry({ service, component: childInput.component ?? component ?? undefined, context: { ...context, ...normalizeContext(childInput.context) }, adapter, now, }); }, debug: (event, fields, options) => log('debug', event, fields, options), info: (event, fields, options) => log('info', event, fields, options), warn: (event, fields, options) => log('warn', event, fields, options), async error(event, error, fields, options) { await publish({ ...base(event, options), ...normalizeFields(fields), triage: errorTriage(fields), telemetryKind: 'log', telemetryLevel: 'error', ...errorFields(error), }); }, counter: (name, value = 1, fields) => metric('counter', name, value, undefined, fields), gauge: (name, value, fields) => metric('gauge', name, value, undefined, fields), distribution: (name, value, options) => metric('distribution', name, value, options?.unit, options?.fields), async time(span, fields, fn, options) { const startedAt = performance.now(); try { const result = await fn(); await publish({ ...base(span), ...normalizeFields(fields), telemetryKind: 'span', telemetryLevel: 'info', spanDurationMs: performance.now() - startedAt, spanOutcome: 'ok', }); return result; } catch (error) { const errorOptions = options?.onError?.(error); const failureFields = { ...(fields ?? {}), ...(errorOptions?.fields ?? {}), }; await publish({ ...base(span), ...normalizeFields(failureFields), triage: errorTriage(failureFields), telemetryKind: 'span', telemetryLevel: errorOptions?.telemetryLevel ?? 'error', spanDurationMs: performance.now() - startedAt, spanOutcome: 'error', ...errorFields(error), }); throw error; } }, }; }