From 7c58ea4367c38c3cc040e12af6b5b605ac375b18 Mon Sep 17 00:00:00 2001 From: Ali Fayed Date: Tue, 4 Aug 2026 14:34:46 -0700 Subject: [PATCH 1/7] feat(coding-agent): emit aura.usage.tokens billing log records per chat call --- .../coding-agent/src/telemetry/sink-otlp.ts | 38 ++++++++++++++- .../coding-agent/test/telemetry-sink.test.ts | 48 +++++++++++++++++++ 2 files changed, 85 insertions(+), 1 deletion(-) diff --git a/packages/coding-agent/src/telemetry/sink-otlp.ts b/packages/coding-agent/src/telemetry/sink-otlp.ts index 260f8313221..c3d2dc9cac7 100644 --- a/packages/coding-agent/src/telemetry/sink-otlp.ts +++ b/packages/coding-agent/src/telemetry/sink-otlp.ts @@ -3,6 +3,7 @@ * structured log records. Publishers stay OTel-free; this module is only * registered when an OTLP provider is live (init.ts). */ +import type { ChatUsageEvent } from "@oh-my-pi/pi-agent-core"; import type { logger } from "@oh-my-pi/pi-utils"; import { type Attributes, context, SpanStatusCode, trace } from "@opentelemetry/api"; import type { LogAttributes } from "@opentelemetry/api-logs"; @@ -11,6 +12,38 @@ import type { AuraMetricRecorder } from "./metrics"; export type EmitOtelLog = (level: logger.LogLevel, body: string, attributes: LogAttributes, eventName: string) => void; +/** + * Billing-record contract consumed by the cloud telemetry worker (elide + * cloud repo, workers/telemetry/meter.ts): one log record per chat call, + * one attribute per positive token count. Renaming the event or attribute + * keys silently stops usage metering — treat both as frozen. + */ +export const AURA_USAGE_TOKENS_EVENT = "aura.usage.tokens"; + +const USAGE_TOKEN_ATTRS = [ + ["input", (usage: ChatUsageEvent["usage"]) => usage.inputTokens], + ["output", (usage: ChatUsageEvent["usage"]) => usage.outputTokens], + ["cache_read", (usage: ChatUsageEvent["usage"]) => usage.cachedInputTokens], + ["cache_write", (usage: ChatUsageEvent["usage"]) => usage.cacheWriteTokens], + ["reasoning", (usage: ChatUsageEvent["usage"]) => usage.reasoningOutputTokens], +] as const satisfies ReadonlyArray number | undefined]>; + +/** Attributes for one billing record, or undefined when there is nothing billable. */ +export function usageTokenAttributes(event: ChatUsageEvent): LogAttributes | undefined { + const attributes: LogAttributes = {}; + let billable = false; + for (const [tokenType, read] of USAGE_TOKEN_ATTRS) { + const count = read(event.usage); + if (typeof count !== "number" || !Number.isFinite(count) || count <= 0) continue; + attributes[`${AURA_USAGE_TOKENS_EVENT}.${tokenType}`] = count; + billable = true; + } + if (!billable) return undefined; + attributes["gen_ai.provider.name"] = event.provider; + attributes["gen_ai.request.model"] = event.model; + return attributes; +} + export interface OtlpSinkDeps { recorder?: AuraMetricRecorder; emitLog: EmitOtelLog; @@ -65,9 +98,12 @@ function handle(deps: OtlpSinkDeps, event: TelemetryEvent): void { case "turn.completed": deps.recorder?.recordRun(event.summary, event.coverage); break; - case "chat.usage": + case "chat.usage": { deps.recorder?.recordChatUsage(event.event); + const usageAttributes = usageTokenAttributes(event.event); + if (usageAttributes) deps.emitLog("info", "chat usage", usageAttributes, AURA_USAGE_TOKENS_EVENT); break; + } case "runtime.call.completed": deps.recorder?.recordRuntimeCall(event); deps.emitLog( diff --git a/packages/coding-agent/test/telemetry-sink.test.ts b/packages/coding-agent/test/telemetry-sink.test.ts index 78491358a5f..40e7929f6db 100644 --- a/packages/coding-agent/test/telemetry-sink.test.ts +++ b/packages/coding-agent/test/telemetry-sink.test.ts @@ -287,3 +287,51 @@ describe("otlp sink", () => { unregister(); }); }); + +describe("usage token billing records", () => { + it("emits one aura.usage.tokens record per chat.usage event with positive counts", () => { + const logs: Array<{ eventName: string; attributes: Record }> = []; + const unregister = registerOtlpSink({ + emitLog: (_level, _body, attributes, eventName) => logs.push({ eventName, attributes }), + }); + emitTelemetryEvent({ + type: "chat.usage", + event: { + provider: "anthropic", + model: "claude-fable-5", + usage: { + inputTokens: 1200, + outputTokens: 340, + cachedInputTokens: 9000, + cacheWriteTokens: 0, + reasoningOutputTokens: 25, + totalTokens: 10565, + }, + } as never, + }); + unregister(); + const billing = logs.filter(entry => entry.eventName === "aura.usage.tokens"); + expect(billing).toHaveLength(1); + expect(billing[0]?.attributes).toEqual({ + "aura.usage.tokens.input": 1200, + "aura.usage.tokens.output": 340, + "aura.usage.tokens.cache_read": 9000, + "aura.usage.tokens.reasoning": 25, + "gen_ai.provider.name": "anthropic", + "gen_ai.request.model": "claude-fable-5", + }); + }); + + it("emits nothing when every count is zero or absent", () => { + const logs: Array<{ eventName: string }> = []; + const unregister = registerOtlpSink({ + emitLog: (_level, _body, _attributes, eventName) => logs.push({ eventName }), + }); + emitTelemetryEvent({ + type: "chat.usage", + event: { provider: "anthropic", model: "m", usage: { inputTokens: 0, outputTokens: 0 } } as never, + }); + unregister(); + expect(logs.filter(entry => entry.eventName === "aura.usage.tokens")).toHaveLength(0); + }); +}); From 03cd6f3e55f370d1ceae2aa57ddde5c05ae9837d Mon Sep 17 00:00:00 2001 From: Ali Fayed Date: Tue, 4 Aug 2026 14:41:22 -0700 Subject: [PATCH 2/7] fix(coding-agent): drop absent provider/model from usage token records --- .../coding-agent/src/telemetry/sink-otlp.ts | 5 +++- .../coding-agent/test/telemetry-sink.test.ts | 24 +++++++++++++++++++ 2 files changed, 28 insertions(+), 1 deletion(-) diff --git a/packages/coding-agent/src/telemetry/sink-otlp.ts b/packages/coding-agent/src/telemetry/sink-otlp.ts index c3d2dc9cac7..1bf1b20d7db 100644 --- a/packages/coding-agent/src/telemetry/sink-otlp.ts +++ b/packages/coding-agent/src/telemetry/sink-otlp.ts @@ -39,7 +39,10 @@ export function usageTokenAttributes(event: ChatUsageEvent): LogAttributes | und billable = true; } if (!billable) return undefined; - attributes["gen_ai.provider.name"] = event.provider; + // `provider` is optional on ChatUsageEvent; drop it rather than shipping an + // undefined value, matching otelAttributes/metricAttributes' convention for + // the same event fields. + if (event.provider !== undefined) attributes["gen_ai.provider.name"] = event.provider; attributes["gen_ai.request.model"] = event.model; return attributes; } diff --git a/packages/coding-agent/test/telemetry-sink.test.ts b/packages/coding-agent/test/telemetry-sink.test.ts index 40e7929f6db..19acb902831 100644 --- a/packages/coding-agent/test/telemetry-sink.test.ts +++ b/packages/coding-agent/test/telemetry-sink.test.ts @@ -334,4 +334,28 @@ describe("usage token billing records", () => { unregister(); expect(logs.filter(entry => entry.eventName === "aura.usage.tokens")).toHaveLength(0); }); + + it("drops gen_ai.provider.name rather than shipping it undefined when provider is absent", () => { + const logs: Array<{ eventName: string; attributes: Record }> = []; + const unregister = registerOtlpSink({ + emitLog: (_level, _body, attributes, eventName) => logs.push({ eventName, attributes }), + }); + emitTelemetryEvent({ + type: "chat.usage", + event: { + provider: undefined, + model: "claude-fable-5", + usage: { inputTokens: 100, outputTokens: 50 }, + } as never, + }); + unregister(); + const billing = logs.filter(entry => entry.eventName === "aura.usage.tokens"); + expect(billing).toHaveLength(1); + expect("gen_ai.provider.name" in (billing[0]?.attributes ?? {})).toBe(false); + expect(billing[0]?.attributes).toEqual({ + "aura.usage.tokens.input": 100, + "aura.usage.tokens.output": 50, + "gen_ai.request.model": "claude-fable-5", + }); + }); }); From 82de4a61c1f621fc230706e2cf8396699c0de1b9 Mon Sep 17 00:00:00 2001 From: Ali Fayed Date: Tue, 4 Aug 2026 14:49:59 -0700 Subject: [PATCH 3/7] feat(coding-agent): authorized OTLP exporters over TokenManager.authorizedFetch Step 1 probe: @opentelemetry/otlp-transformer@0.220.0 exports ProtobufLogsSerializer/ProtobufMetricsSerializer/ProtobufTraceSerializer exactly as expected, no fallback needed. --- bun.lock | 2 + package.json | 1 + packages/coding-agent/package.json | 1 + .../src/telemetry/authorized-exporters.ts | 115 ++++++++++++++++++ .../telemetry-authorized-exporters.test.ts | 64 ++++++++++ 5 files changed, 183 insertions(+) create mode 100644 packages/coding-agent/src/telemetry/authorized-exporters.ts create mode 100644 packages/coding-agent/test/telemetry-authorized-exporters.test.ts diff --git a/bun.lock b/bun.lock index 4730b16494c..5b7391af740 100644 --- a/bun.lock +++ b/bun.lock @@ -97,6 +97,7 @@ "@opentelemetry/exporter-logs-otlp-proto": "catalog:", "@opentelemetry/exporter-metrics-otlp-proto": "catalog:", "@opentelemetry/exporter-trace-otlp-proto": "catalog:", + "@opentelemetry/otlp-transformer": "catalog:", "@opentelemetry/resources": "catalog:", "@opentelemetry/sdk-logs": "catalog:", "@opentelemetry/sdk-metrics": "catalog:", @@ -394,6 +395,7 @@ "@opentelemetry/exporter-logs-otlp-proto": "^0.220.0", "@opentelemetry/exporter-metrics-otlp-proto": "^0.220.0", "@opentelemetry/exporter-trace-otlp-proto": "^0.220.0", + "@opentelemetry/otlp-transformer": "^0.220.0", "@opentelemetry/resources": "^2.9.0", "@opentelemetry/sdk-logs": "^0.220.0", "@opentelemetry/sdk-metrics": "^2.9.0", diff --git a/package.json b/package.json index a2c27f72c3f..3b398a5b9f4 100644 --- a/package.json +++ b/package.json @@ -44,6 +44,7 @@ "@opentelemetry/exporter-logs-otlp-proto": "^0.220.0", "@opentelemetry/exporter-metrics-otlp-proto": "^0.220.0", "@opentelemetry/exporter-trace-otlp-proto": "^0.220.0", + "@opentelemetry/otlp-transformer": "^0.220.0", "@opentelemetry/resources": "^2.9.0", "@opentelemetry/sdk-logs": "^0.220.0", "@opentelemetry/sdk-metrics": "^2.9.0", diff --git a/packages/coding-agent/package.json b/packages/coding-agent/package.json index 9f8e1935023..e283237dad0 100644 --- a/packages/coding-agent/package.json +++ b/packages/coding-agent/package.json @@ -72,6 +72,7 @@ "@opentelemetry/exporter-logs-otlp-proto": "catalog:", "@opentelemetry/exporter-metrics-otlp-proto": "catalog:", "@opentelemetry/exporter-trace-otlp-proto": "catalog:", + "@opentelemetry/otlp-transformer": "catalog:", "@opentelemetry/resources": "catalog:", "@opentelemetry/sdk-logs": "catalog:", "@opentelemetry/sdk-metrics": "catalog:", diff --git a/packages/coding-agent/src/telemetry/authorized-exporters.ts b/packages/coding-agent/src/telemetry/authorized-exporters.ts new file mode 100644 index 00000000000..7a9f0fdcc0d --- /dev/null +++ b/packages/coding-agent/src/telemetry/authorized-exporters.ts @@ -0,0 +1,115 @@ +/** + * OTLP exporters for the Aura telemetry tier. + * + * The stock OTLP exporters take static headers at construction; an Aura + * access token lives ≤15 minutes and may only ever be attached by + * `TokenManager.authorizedFetch` (the single bearer-attachment point). These + * exporters serialize with the same otlp-transformer the stock exporters + * use, then send through an injected `authorizedFetch` — fresh token per + * export, 401-refresh-retry and redirect guards included. Only the Aura + * tier constructs them (init.ts); every other destination keeps the stock + * exporters and never sees a credential. + */ +import { type ExportResult, ExportResultCode } from "@opentelemetry/core"; +import { + ProtobufLogsSerializer, + ProtobufMetricsSerializer, + ProtobufTraceSerializer, +} from "@opentelemetry/otlp-transformer"; +import type { LogRecordExporter, ReadableLogRecord } from "@opentelemetry/sdk-logs"; +import type { PushMetricExporter, ResourceMetrics } from "@opentelemetry/sdk-metrics"; +import { AggregationTemporality } from "@opentelemetry/sdk-metrics"; +import type { ReadableSpan, SpanExporter } from "@opentelemetry/sdk-trace-base"; + +export interface AuraTelemetryTransport { + authorizedFetch( + url: string, + init: { method?: string; headers?: Bun.HeadersInit; body?: Bun.BodyInit; eligibleOrigin?: string }, + ): Promise; +} + +export interface AuthorizedExporterOptions { + /** Full signal URL, e.g. https://telemetry.elide.cloud/v1/logs. */ + url: string; + transport: AuraTelemetryTransport; +} + +/** Exported for tests (like errorEventFromLog in init.ts). */ +export function sendAuthorized( + options: AuthorizedExporterOptions, + body: Uint8Array | undefined, + resultCallback: (result: ExportResult) => void, +): void { + if (body === undefined || body.byteLength === 0) { + resultCallback({ code: ExportResultCode.SUCCESS }); + return; + } + const eligibleOrigin = new URL(options.url).origin; + options.transport + .authorizedFetch(options.url, { + method: "POST", + headers: { "content-type": "application/x-protobuf" }, + body: body as unknown as Bun.BodyInit, + eligibleOrigin, + }) + .then(async response => { + await response.body?.cancel(); + if (response.ok) resultCallback({ code: ExportResultCode.SUCCESS }); + else + resultCallback({ code: ExportResultCode.FAILED, error: new Error(`collector status ${response.status}`) }); + }) + .catch(error => { + resultCallback({ + code: ExportResultCode.FAILED, + error: error instanceof Error ? error : new Error(String(error)), + }); + }); +} + +export class AuthorizedLogExporter implements LogRecordExporter { + constructor(private readonly options: AuthorizedExporterOptions) {} + + export(logs: ReadableLogRecord[], resultCallback: (result: ExportResult) => void): void { + sendAuthorized( + this.options, + logs.length === 0 ? undefined : ProtobufLogsSerializer.serializeRequest(logs), + resultCallback, + ); + } + + /** Nothing buffered — every export() call sends immediately. */ + async forceFlush(): Promise {} + + async shutdown(): Promise {} +} + +export class AuthorizedTraceExporter implements SpanExporter { + constructor(private readonly options: AuthorizedExporterOptions) {} + + export(spans: ReadableSpan[], resultCallback: (result: ExportResult) => void): void { + sendAuthorized( + this.options, + spans.length === 0 ? undefined : ProtobufTraceSerializer.serializeRequest(spans), + resultCallback, + ); + } + + async shutdown(): Promise {} +} + +export class AuthorizedMetricExporter implements PushMetricExporter { + constructor(private readonly options: AuthorizedExporterOptions) {} + + export(metrics: ResourceMetrics, resultCallback: (result: ExportResult) => void): void { + sendAuthorized(this.options, ProtobufMetricsSerializer.serializeRequest(metrics), resultCallback); + } + + /** Cumulative, matching the stock exporter default — billing reads logs, not metrics. */ + selectAggregationTemporality(): AggregationTemporality { + return AggregationTemporality.CUMULATIVE; + } + + async forceFlush(): Promise {} + + async shutdown(): Promise {} +} diff --git a/packages/coding-agent/test/telemetry-authorized-exporters.test.ts b/packages/coding-agent/test/telemetry-authorized-exporters.test.ts new file mode 100644 index 00000000000..b825fc899ed --- /dev/null +++ b/packages/coding-agent/test/telemetry-authorized-exporters.test.ts @@ -0,0 +1,64 @@ +import { describe, expect, it } from "bun:test"; +import { ExportResultCode } from "@opentelemetry/core"; +import { LoggerProvider, SimpleLogRecordProcessor } from "@opentelemetry/sdk-logs"; +import { AuthorizedLogExporter, sendAuthorized } from "../src/telemetry/authorized-exporters"; + +function fakeTransport(status = 200) { + const calls: Array<{ + url: string; + init: { method?: string; headers?: Bun.HeadersInit; body?: Bun.BodyInit; eligibleOrigin?: string }; + }> = []; + return { + calls, + transport: { + authorizedFetch: async (url: string, init: (typeof calls)[number]["init"]) => { + calls.push({ url, init }); + return new Response(new Uint8Array(), { status }); + }, + }, + }; +} + +// Minimal ReadableLogRecord: drive a real LoggerProvider with a +// SimpleLogRecordProcessor into the exporter instead of hand-building one. +function emitOne(exporter: AuthorizedLogExporter) { + // This pinned sdk-logs (0.220.0) takes { exporter }, not the exporter directly. + const provider = new LoggerProvider({ processors: [new SimpleLogRecordProcessor({ exporter })] }); + provider.getLogger("test").emit({ body: "chat usage", eventName: "aura.usage.tokens" }); + return provider.forceFlush(); +} + +describe("AuthorizedLogExporter", () => { + it("POSTs protobuf to the configured url with the eligible origin", async () => { + const { calls, transport } = fakeTransport(); + const exporter = new AuthorizedLogExporter({ url: "https://telemetry.elide.cloud/v1/logs", transport }); + await emitOne(exporter); + expect(calls).toHaveLength(1); + expect(calls[0]?.url).toBe("https://telemetry.elide.cloud/v1/logs"); + expect(calls[0]?.init.method).toBe("POST"); + expect(new Headers(calls[0]?.init.headers).get("content-type")).toBe("application/x-protobuf"); + expect(calls[0]?.init.eligibleOrigin).toBe("https://telemetry.elide.cloud"); + const body = calls[0]?.init.body as Uint8Array | undefined; + expect(body?.byteLength).toBeGreaterThan(0); + }); + + it("maps a collector rejection to a FAILED result", async () => { + const { transport } = fakeTransport(500); + const results: Array<{ code: number }> = []; + sendAuthorized({ url: "https://telemetry.elide.cloud/v1/logs", transport }, new Uint8Array([1]), result => + results.push(result), + ); + await new Promise(resolve => setTimeout(resolve, 0)); + expect(results[0]?.code).toBe(ExportResultCode.FAILED); + }); + + it("an empty batch short-circuits SUCCESS without touching the transport", async () => { + const { calls, transport } = fakeTransport(); + const exporter = new AuthorizedLogExporter({ url: "https://telemetry.elide.cloud/v1/logs", transport }); + const results: Array<{ code: number }> = []; + exporter.export([], result => results.push(result)); + await new Promise(resolve => setTimeout(resolve, 0)); + expect(results[0]?.code).toBe(ExportResultCode.SUCCESS); + expect(calls).toHaveLength(0); + }); +}); From 58738e8a1be735a579ae0a8809687f808f5a26a3 Mon Sep 17 00:00:00 2001 From: Ali Fayed Date: Tue, 4 Aug 2026 15:00:10 -0700 Subject: [PATCH 4/7] feat(coding-agent): route Aura-tier OTLP exports through authorized exporters --- packages/coding-agent/src/telemetry/init.ts | 53 +++++++++++++++++-- .../test/telemetry-export.test.ts | 39 ++++++++++++++ 2 files changed, 89 insertions(+), 3 deletions(-) diff --git a/packages/coding-agent/src/telemetry/init.ts b/packages/coding-agent/src/telemetry/init.ts index a191760bdd5..2f841d7fddf 100644 --- a/packages/coding-agent/src/telemetry/init.ts +++ b/packages/coding-agent/src/telemetry/init.ts @@ -40,6 +40,12 @@ import { BatchSpanProcessor } from "@opentelemetry/sdk-trace-base"; import { NodeTracerProvider } from "@opentelemetry/sdk-trace-node"; import { auraDeploymentFor, readCloudSwitches, resolveAuraDeployment } from "../cloud/deployment"; import type { Settings } from "../config/settings"; +import { + type AuraTelemetryTransport, + AuthorizedLogExporter, + AuthorizedMetricExporter, + AuthorizedTraceExporter, +} from "./authorized-exporters"; import { type ErrorReportedTelemetry, emitTelemetryEvent, getActiveTelemetrySessionId } from "./events"; import { buildResourceAttributes, getOrCreateInstallId } from "./identity"; import { AuraMetricRecorder } from "./metrics"; @@ -56,6 +62,16 @@ const FLUSH_INTERVAL_MS = 30_000; export interface InitTelemetryOptions { /** Settings instance for telemetry.* config (env always wins). */ settings?: Settings; + /** + * Cloud transport for the Aura telemetry tier. When present AND the Aura + * tier wins the destination, exports go through authorized exporters — + * fresh bearer per export via TokenManager.authorizedFetch. Absent (no + * signed-in cloud session wired up yet), the Aura tier exports + * unauthenticated and the collector's edge will 401 them; the built-in + * and operator tiers are unaffected either way. The CLI wires this once + * the cloud login command lands (cloud/auth.ts AuraAuthClient.manager). + */ + cloud?: { transport: AuraTelemetryTransport }; } export type TelemetrySignal = "trace" | "log" | "metric"; @@ -349,6 +365,28 @@ export function resolveExporterConfig( return config; } +/** + * The full signal URL to export through an authorized exporter, or + * undefined when the Aura tier is not the winning destination. Mirrors + * resolveExporterConfig's precedence exactly: env endpoints always win + * (and never carry the Aura credential), an explicit telemetry.endpoint + * outranks the Aura tier, and the built-in tier never authenticates. + */ +export function resolveAuraAuthorizedUrl( + signal: TelemetrySignal, + settings: Pick | undefined, + processEnv: Record = process.env, +): string | undefined { + if (!settings?.get("telemetry.enabled")) return undefined; + const { path, envInfix } = SIGNAL_OTLP[signal]; + const envEndpoint = processEnv[`OTEL_EXPORTER_OTLP_${envInfix}_ENDPOINT`] ?? processEnv.OTEL_EXPORTER_OTLP_ENDPOINT; + if (envEndpoint) return undefined; + if (settings.get("telemetry.endpoint")) return undefined; + const aura = auraTelemetryEndpoint(settings, processEnv); + if (!aura) return undefined; + return `${aura.replace(/\/+$/, "")}${path}`; +} + /** A destination and the headers that belong to it — resolved together, always. */ interface TelemetryDestination { readonly url: string; @@ -438,7 +476,10 @@ async function registerProviders(signalConfig: SignalConfig, options: InitTeleme const settings = options.settings; if (signalConfig.trace) { - const exporter = new OTLPTraceExporter(resolveExporterConfig("trace", settings)); + const authorizedUrl = options.cloud && resolveAuraAuthorizedUrl("trace", settings); + const exporter = authorizedUrl + ? new AuthorizedTraceExporter({ url: authorizedUrl, transport: options.cloud!.transport }) + : new OTLPTraceExporter(resolveExporterConfig("trace", settings)); traceProvider = new NodeTracerProvider({ resource, spanProcessors: [new BatchSpanProcessor(exporter)], @@ -447,7 +488,10 @@ async function registerProviders(signalConfig: SignalConfig, options: InitTeleme } if (signalConfig.metric) { - const exporter = new OTLPMetricExporter(resolveExporterConfig("metric", settings)); + const authorizedUrl = options.cloud && resolveAuraAuthorizedUrl("metric", settings); + const exporter = authorizedUrl + ? new AuthorizedMetricExporter({ url: authorizedUrl, transport: options.cloud!.transport }) + : new OTLPMetricExporter(resolveExporterConfig("metric", settings)); meterProvider = new MeterProvider({ resource, readers: [new PeriodicExportingMetricReader({ exporter })], @@ -459,7 +503,10 @@ async function registerProviders(signalConfig: SignalConfig, options: InitTeleme } if (signalConfig.log) { - const exporter = new OTLPLogExporter(resolveExporterConfig("log", settings)); + const authorizedUrl = options.cloud && resolveAuraAuthorizedUrl("log", settings); + const exporter = authorizedUrl + ? new AuthorizedLogExporter({ url: authorizedUrl, transport: options.cloud!.transport }) + : new OTLPLogExporter(resolveExporterConfig("log", settings)); logProvider = new LoggerProvider({ resource, processors: [new BatchLogRecordProcessor({ exporter })], diff --git a/packages/coding-agent/test/telemetry-export.test.ts b/packages/coding-agent/test/telemetry-export.test.ts index a6198e01c8e..233c89002b8 100644 --- a/packages/coding-agent/test/telemetry-export.test.ts +++ b/packages/coding-agent/test/telemetry-export.test.ts @@ -4,6 +4,7 @@ import { createTelemetryExportConfig, initTelemetryExport, isTelemetryExportEnabled, + resolveAuraAuthorizedUrl, subscribeTelemetry, } from "@oh-my-pi/pi-coding-agent/telemetry-export"; import { logger } from "@oh-my-pi/pi-utils"; @@ -257,3 +258,41 @@ describe("initTelemetryExport registration failures", () => { expect(warnings).toHaveLength(2); }); }); + +describe("resolveAuraAuthorizedUrl", () => { + const settings = { + get: (key: string) => + (({ "telemetry.enabled": true, "cloud.telemetry.enabled": true }) as Record)[key], + } as never; + + it("returns the signal URL when the Aura tier wins the destination", () => { + const url = resolveAuraAuthorizedUrl("log", settings, { AURA_DOMAIN: "elide.cloud" }); + expect(url).toBe("https://telemetry.elide.cloud/v1/logs"); + }); + + it("returns undefined when env owns the endpoint (credential never follows env)", () => { + const url = resolveAuraAuthorizedUrl("log", settings, { + AURA_DOMAIN: "elide.cloud", + OTEL_EXPORTER_OTLP_ENDPOINT: "https://operator.example", + }); + expect(url).toBeUndefined(); + }); + + it("returns undefined when an explicit telemetry.endpoint outranks the Aura tier", () => { + const withEndpoint = { + get: (key: string) => + ( + ({ + "telemetry.enabled": true, + "cloud.telemetry.enabled": true, + "telemetry.endpoint": "https://collector.example", + }) as Record + )[key], + } as never; + expect(resolveAuraAuthorizedUrl("log", withEndpoint, { AURA_DOMAIN: "elide.cloud" })).toBeUndefined(); + }); + + it("returns undefined for the built-in tier (no AURA_DOMAIN)", () => { + expect(resolveAuraAuthorizedUrl("log", settings, {})).toBeUndefined(); + }); +}); From 8eb8d268f77d21bee970a54ca0f9a9550ff2bd92 Mon Sep 17 00:00:00 2001 From: Ali Fayed Date: Tue, 4 Aug 2026 15:31:41 -0700 Subject: [PATCH 5/7] fix(coding-agent): cap log export batch size below the cloud relay's 112 KiB limit MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit BatchLogRecordProcessor's default maxExportBatchSize (512) can serialize to a protobuf payload larger than workers/telemetry's MAX_BODY_BYTES cap in the elide cloud repo, which 413s the whole batch — dropping every record in it, including any aura.usage.tokens billing records riding along with observability logs. Pin a conservative 64-record cap on the LoggerProvider's batch processor. --- packages/coding-agent/src/telemetry/init.ts | 13 ++++++++++++- 1 file changed, 12 insertions(+), 1 deletion(-) diff --git a/packages/coding-agent/src/telemetry/init.ts b/packages/coding-agent/src/telemetry/init.ts index 2f841d7fddf..d4f8a0f199b 100644 --- a/packages/coding-agent/src/telemetry/init.ts +++ b/packages/coding-agent/src/telemetry/init.ts @@ -58,6 +58,17 @@ import { registerOtlpSink } from "./sink-otlp"; */ const FLUSH_INTERVAL_MS = 30_000; +/** + * Max log records per export batch. The BatchLogRecordProcessor default + * (512) can serialize to a protobuf payload larger than the cloud relay's + * 112 KiB per-request cap (workers/telemetry's MAX_BODY_BYTES in the elide + * cloud repo); a batch that trips that cap 413s in full, dropping every + * record in it — including any aura.usage.tokens billing records riding + * along with observability logs. A conservative cap keeps batches well + * under that ceiling. + */ +const MAX_LOG_EXPORT_BATCH_SIZE = 64; + /** Options for {@link initTelemetryExport}. */ export interface InitTelemetryOptions { /** Settings instance for telemetry.* config (env always wins). */ @@ -509,7 +520,7 @@ async function registerProviders(signalConfig: SignalConfig, options: InitTeleme : new OTLPLogExporter(resolveExporterConfig("log", settings)); logProvider = new LoggerProvider({ resource, - processors: [new BatchLogRecordProcessor({ exporter })], + processors: [new BatchLogRecordProcessor({ exporter, maxExportBatchSize: MAX_LOG_EXPORT_BATCH_SIZE })], }); logs.setGlobalLoggerProvider(logProvider); otelLogger = logProvider.getLogger("aura"); From faea57986adc1117ccea614bf28df28a80b58e80 Mon Sep 17 00:00:00 2001 From: Ali Fayed Date: Tue, 4 Aug 2026 15:31:57 -0700 Subject: [PATCH 6/7] fix(coding-agent): exempt aura.usage.tokens from the OTEL_LOG_LEVEL gate MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit emitOtelLog (init.ts) filtered every record — including the aura.usage.tokens billing record — on OTEL_LOG_LEVEL, so an operator setting OTEL_LOG_LEVEL=warn (or lower) to quiet observability noise would silently stop usage metering too. The billing record is metering, not observability, so exempt it from the level gate by eventName. Covered by a new out-of-process probe (mirroring the existing otel-*-probe.ts pattern so the global LoggerProvider singleton never leaks into the test runner): with OTEL_LOG_LEVEL=error, an ordinary info-level bridged log is still suppressed (the gate still works), while a billable chat.usage event's aura.usage.tokens record still reaches the collector (the exemption holds). --- packages/coding-agent/src/telemetry/init.ts | 14 +- .../test/otel-log-level-billing-probe.ts | 187 ++++++++++++++++++ .../test/telemetry-export.test.ts | 18 ++ 3 files changed, 215 insertions(+), 4 deletions(-) create mode 100644 packages/coding-agent/test/otel-log-level-billing-probe.ts diff --git a/packages/coding-agent/src/telemetry/init.ts b/packages/coding-agent/src/telemetry/init.ts index d4f8a0f199b..63a53860d2d 100644 --- a/packages/coding-agent/src/telemetry/init.ts +++ b/packages/coding-agent/src/telemetry/init.ts @@ -49,7 +49,7 @@ import { import { type ErrorReportedTelemetry, emitTelemetryEvent, getActiveTelemetrySessionId } from "./events"; import { buildResourceAttributes, getOrCreateInstallId } from "./identity"; import { AuraMetricRecorder } from "./metrics"; -import { registerOtlpSink } from "./sink-otlp"; +import { AURA_USAGE_TOKENS_EVENT, registerOtlpSink } from "./sink-otlp"; /** * Periodic flush interval. A long-lived `omp` process (the ACP server is @@ -678,9 +678,15 @@ function emitOtelLog( timestamp = new Date(), ): void { if (!otelLogger) return; - const minLevel = parseOtelLogLevel(process.env.OTEL_LOG_LEVEL); - if (minLevel === "none") return; - if (LOG_LEVEL_WEIGHT[level] > LOG_LEVEL_WEIGHT[minLevel]) return; + // OTEL_LOG_LEVEL is an observability verbosity knob (suppress debug/info + // noise); the aura.usage.tokens billing record is metering, not + // observability, so it is exempt from the gate — an operator tuning log + // verbosity down must never silently stop usage metering. + if (eventName !== AURA_USAGE_TOKENS_EVENT) { + const minLevel = parseOtelLogLevel(process.env.OTEL_LOG_LEVEL); + if (minLevel === "none") return; + if (LOG_LEVEL_WEIGHT[level] > LOG_LEVEL_WEIGHT[minLevel]) return; + } otelLogger.emit({ eventName, timestamp, diff --git a/packages/coding-agent/test/otel-log-level-billing-probe.ts b/packages/coding-agent/test/otel-log-level-billing-probe.ts new file mode 100644 index 00000000000..cc8e055d868 --- /dev/null +++ b/packages/coding-agent/test/otel-log-level-billing-probe.ts @@ -0,0 +1,187 @@ +/** + * Positive-path probe for the OTEL_LOG_LEVEL billing exemption, run as a + * subprocess by telemetry-export.test.ts. Keeping it out-of-process means the + * global LoggerProvider singleton that initTelemetryExport() registers never + * leaks into the test runner, and OTEL_LOG_LEVEL is set only in the + * subprocess's own environment — nothing to restore in the parent. + * + * Regression for the OTEL_LOG_LEVEL gate silently stopping metering + * (init.ts's emitOtelLog): with OTEL_LOG_LEVEL=error, an ordinary info-level + * bridged log must still be suppressed (proving the gate still works), while + * a billable chat.usage event — which emits its aura.usage.tokens record at + * "info" — must still reach the collector (proving the billing record is + * exempt from the gate). + */ + +import * as fs from "node:fs"; +import * as os from "node:os"; +import * as path from "node:path"; +import type { ChatUsageEvent } from "@oh-my-pi/pi-agent-core"; +import { + AURA_USAGE_TOKENS_EVENT, + createTelemetryExportConfig, + flushTelemetryExport, + initTelemetryExport, + isTelemetryExportEnabled, +} from "@oh-my-pi/pi-coding-agent/telemetry-export"; +import { logger } from "@oh-my-pi/pi-utils"; + +const logPayloads: Uint8Array[] = []; + +interface ProtobufField { + readonly number: number; + readonly bytes?: Uint8Array; +} + +function readVarint(bytes: Uint8Array, offset: number): [number, number] { + let value = 0; + let shift = 0; + while (offset < bytes.length) { + const byte = bytes[offset++]; + value += (byte & 0x7f) * 2 ** shift; + if ((byte & 0x80) === 0) return [value, offset]; + shift += 7; + } + throw new Error("Truncated protobuf varint"); +} + +function protobufFields(bytes: Uint8Array): ProtobufField[] { + const fields: ProtobufField[] = []; + for (let offset = 0; offset < bytes.length; ) { + const [tag, nextOffset] = readVarint(bytes, offset); + offset = nextOffset; + const wireType = tag & 7; + const number = tag >>> 3; + if (wireType === 0) { + [, offset] = readVarint(bytes, offset); + fields.push({ number }); + } else if (wireType === 1) { + offset += 8; + fields.push({ number }); + } else if (wireType === 2) { + const [length, valueOffset] = readVarint(bytes, offset); + offset = valueOffset; + const end = offset + length; + if (end > bytes.length) throw new Error("Truncated protobuf field"); + fields.push({ number, bytes: bytes.slice(offset, end) }); + offset = end; + } else if (wireType === 5) { + offset += 4; + fields.push({ number }); + } else { + throw new Error(`Unsupported protobuf wire type ${wireType}`); + } + } + return fields; +} + +function text(bytes: Uint8Array): string { + return new TextDecoder().decode(bytes); +} + +/** `LogRecord.event_name` (field 12) via resource_logs(1) -> scope_logs(2) -> log_records(2). */ +function logEventNames(payload: Uint8Array): string[] { + const names: string[] = []; + for (const resourceLogs of protobufFields(payload)) { + if (resourceLogs.number !== 1 || !resourceLogs.bytes) continue; + for (const scopeLogs of protobufFields(resourceLogs.bytes)) { + if (scopeLogs.number !== 2 || !scopeLogs.bytes) continue; + for (const record of protobufFields(scopeLogs.bytes)) { + if (record.number !== 2 || !record.bytes) continue; + const eventName = protobufFields(record.bytes).find(field => field.number === 12)?.bytes; + if (eventName) names.push(text(eventName)); + } + } + } + return names; +} + +const server = Bun.serve({ + port: 0, + async fetch(req) { + const endpoint = new URL(req.url).pathname; + if (req.method === "POST" && endpoint.endsWith("/v1/logs")) { + const body = await req.arrayBuffer(); + if (body.byteLength > 0) logPayloads.push(new Uint8Array(body)); + return new Response('{"partialSuccess":{}}', { + status: 200, + headers: { "content-type": "application/json" }, + }); + } + return new Response("not found", { status: 404 }); + }, +}); + +process.env.OTEL_EXPORTER_OTLP_LOGS_ENDPOINT = `http://localhost:${server.port}/v1/logs`; +process.env.OTEL_SERVICE_NAME = "oh-my-pi-log-level-billing-probe"; +// The gate under test: raised well above "info" so an ordinary bridged log is +// suppressed, isolating the billing record's exemption from it. +process.env.OTEL_LOG_LEVEL = "error"; + +// Sandbox the config root so the probe's `aura.install.id` is minted into a temp +// directory instead of persisting into the developer's real config root. +const configRoot = fs.mkdtempSync(path.join(os.homedir(), ".aura-log-level-billing-probe-")); +process.on("exit", () => fs.rmSync(configRoot, { recursive: true, force: true })); +process.env.PI_CONFIG_DIR = path.basename(configRoot); +const { refreshDirsFromEnv } = await import("@oh-my-pi/pi-utils/dirs"); +refreshDirsFromEnv(); + +await initTelemetryExport(); +if (!isTelemetryExportEnabled()) { + console.error("PROBE: provider did not register"); + await server.stop(true); + process.exit(2); +} + +const config = createTelemetryExportConfig(undefined); +if (!config) { + console.error("PROBE: export config not produced"); + await server.stop(true); + process.exit(2); +} + +// An ordinary info-level bridged log: must be suppressed under OTEL_LOG_LEVEL=error. +logger.info("probe info log — should be suppressed", { code: "probe_info" }); + +// A billable chat.usage event: sink-otlp emits its aura.usage.tokens record at +// "info" (sink-otlp.ts), which must still get through despite the gate above. +const usage: ChatUsageEvent = { + span: undefined as never, + agent: { id: "main", name: "Main" }, + conversationId: "probe-session", + stepNumber: 0, + model: "claude-haiku-4-5", + provider: "anthropic", + serviceTier: undefined, + usage: { + inputTokens: 1000, + outputTokens: 200, + totalTokens: 1200, + cachedInputTokens: 0, + cacheWriteTokens: 0, + reasoningOutputTokens: 0, + }, + cost: { usd: 0.01 }, + attributes: undefined, + headers: undefined, +}; +await config.onChatUsage?.(usage); + +await flushTelemetryExport(); +await server.stop(true); + +const eventNames = logPayloads.flatMap(logEventNames); +const sawBilling = eventNames.includes(AURA_USAGE_TOKENS_EVENT); +const sawSuppressedInfo = eventNames.includes("aura.log"); + +if (!sawBilling) { + console.error(`PROBE: billing record missing (saw ${eventNames.join(",") || "none"})`); + process.exit(1); +} +if (sawSuppressedInfo) { + console.error("PROBE: OTEL_LOG_LEVEL=error failed to suppress the ordinary info log — gate is not selective"); + process.exit(1); +} + +console.log("PROBE: BILLING EXEMPT FROM LOG LEVEL GATE"); +process.exit(0); diff --git a/packages/coding-agent/test/telemetry-export.test.ts b/packages/coding-agent/test/telemetry-export.test.ts index 233c89002b8..a0bf91cdd23 100644 --- a/packages/coding-agent/test/telemetry-export.test.ts +++ b/packages/coding-agent/test/telemetry-export.test.ts @@ -144,6 +144,24 @@ describe("initTelemetryExport signals export path", () => { expect(code).toBe(0); }, 20_000); + it("exempts the aura.usage.tokens billing record from OTEL_LOG_LEVEL", async () => { + // Regression: emitOtelLog's level gate (init.ts) used to apply uniformly, + // so OTEL_LOG_LEVEL=error silently stopped usage metering. The probe sets + // OTEL_LOG_LEVEL=error, confirms an ordinary info-level bridged log is + // still suppressed (the gate still works), and confirms a billable + // chat.usage event's aura.usage.tokens record — emitted at "info" — still + // reaches the collector (the billing record is exempt). + const probe = fileURLToPath(new URL("./otel-log-level-billing-probe.ts", import.meta.url)); + const proc = Bun.spawn(["bun", probe], { stdout: "pipe", stderr: "pipe" }); + const [code, stdout, stderr] = await Promise.all([ + proc.exited, + new Response(proc.stdout).text(), + new Response(proc.stderr).text(), + ]); + expect(`${stdout}\n${stderr}`).toContain("PROBE: BILLING EXEMPT FROM LOG LEVEL GATE"); + expect(code).toBe(0); + }, 20_000); + it("merges OTEL_RESOURCE_ATTRIBUTES into the exported resource", async () => { // Regression for #7134: the resource only carried service.name, so // OTEL_RESOURCE_ATTRIBUTES entries never reached the collector. The probe From 3e6629705a782a4e98547d05a707e085317b556f Mon Sep 17 00:00:00 2001 From: Ali Fayed Date: Tue, 4 Aug 2026 15:32:08 -0700 Subject: [PATCH 7/7] fix(coding-agent): retry/backoff/timeout for the authorized OTLP send path sendAuthorized did one authorizedFetch and nothing else: a transient 429/5xx permanently dropped the batch (BatchLogRecordProcessor does not re-queue a FAILED export), and no signal meant an export could hang forever. Mirror the stock otlp-exporter-base behavior in spirit: retry on 429/502/503/504 and network rejection with jittered exponential backoff (honoring Retry-After when present) up to 3 attempts total, and bound every attempt with AbortSignal.timeout (10s, the stock default). Non-retryable statuses still fail immediately, with no retry. sleep is injectable via AuthorizedExporterOptions so the new retry tests run instantly instead of on real timers. New tests: retry-then-success on 503 and on a network rejection, exhaustion after persistent 503s -> FAILED, a non-retryable 400 fails immediately without retrying, Retry-After overrides the computed backoff, each attempt carries an AbortSignal, and the previously-missing direct happy-path tests for AuthorizedTraceExporter and AuthorizedMetricExporter through the shared send. --- .../src/telemetry/authorized-exporters.ts | 108 +++++++++-- .../telemetry-authorized-exporters.test.ts | 183 +++++++++++++++++- 2 files changed, 263 insertions(+), 28 deletions(-) diff --git a/packages/coding-agent/src/telemetry/authorized-exporters.ts b/packages/coding-agent/src/telemetry/authorized-exporters.ts index 7a9f0fdcc0d..2d9f6d4549e 100644 --- a/packages/coding-agent/src/telemetry/authorized-exporters.ts +++ b/packages/coding-agent/src/telemetry/authorized-exporters.ts @@ -9,6 +9,15 @@ * export, 401-refresh-retry and redirect guards included. Only the Aura * tier constructs them (init.ts); every other destination keeps the stock * exporters and never sees a credential. + * + * `sendAuthorized` mirrors, in spirit, the stock otlp-exporter-base retry + * behavior: a 429/502/503/504 or a network-level rejection is retried with + * jittered exponential backoff (honoring `Retry-After` when the collector + * sends one) up to {@link MAX_ATTEMPTS} attempts, and every attempt is bounded + * by {@link EXPORT_TIMEOUT_MS} — otherwise `BatchLogRecordProcessor` marks a + * hung export FAILED but never re-queues it, so a transient collector blip or + * a wedged connection would permanently drop the batch, billing records + * included. */ import { type ExportResult, ExportResultCode } from "@opentelemetry/core"; import { @@ -24,7 +33,13 @@ import type { ReadableSpan, SpanExporter } from "@opentelemetry/sdk-trace-base"; export interface AuraTelemetryTransport { authorizedFetch( url: string, - init: { method?: string; headers?: Bun.HeadersInit; body?: Bun.BodyInit; eligibleOrigin?: string }, + init: { + method?: string; + headers?: Bun.HeadersInit; + body?: Bun.BodyInit; + eligibleOrigin?: string; + signal?: AbortSignal; + }, ): Promise; } @@ -32,38 +47,91 @@ export interface AuthorizedExporterOptions { /** Full signal URL, e.g. https://telemetry.elide.cloud/v1/logs. */ url: string; transport: AuraTelemetryTransport; + /** + * Backoff delay override, for tests: given the computed delay in + * milliseconds, resolve whenever the test is ready to let the retry + * proceed. Defaults to a real `setTimeout`-based sleep. + */ + sleep?: (ms: number) => Promise; +} + +/** Up to 3 attempts total (the initial send plus 2 retries). */ +const MAX_ATTEMPTS = 3; +/** Collector statuses worth retrying — transient overload/unavailability, per the OTLP spec. */ +const RETRYABLE_STATUSES = new Set([429, 502, 503, 504]); +const INITIAL_BACKOFF_MS = 1000; +const MAX_BACKOFF_MS = 5000; +const BACKOFF_MULTIPLIER = 2; +/** Jitter fraction applied symmetrically around the computed backoff. */ +const JITTER = 0.2; +/** Per-attempt bound, matching the stock OTLP exporters' default export timeout. */ +const EXPORT_TIMEOUT_MS = 10_000; + +function jitteredBackoff(baseMs: number): number { + const jitter = Math.random() * (2 * JITTER) - JITTER; + return Math.max(0, Math.min(baseMs * (1 + jitter), MAX_BACKOFF_MS)); +} + +/** `Retry-After`, in ms: an integer is seconds, otherwise an HTTP-date. Undefined when absent or unparsable. */ +function parseRetryAfterMs(value: string | null): number | undefined { + if (!value) return undefined; + const seconds = Number.parseInt(value, 10); + if (Number.isInteger(seconds) && String(seconds) === value.trim()) return Math.max(0, seconds * 1000); + const delay = new Date(value).getTime() - Date.now(); + return Number.isNaN(delay) ? undefined : Math.max(0, delay); +} + +function defaultSleep(ms: number): Promise { + return new Promise(resolve => setTimeout(resolve, ms)); } /** Exported for tests (like errorEventFromLog in init.ts). */ -export function sendAuthorized( +export async function sendAuthorized( options: AuthorizedExporterOptions, body: Uint8Array | undefined, resultCallback: (result: ExportResult) => void, -): void { +): Promise { if (body === undefined || body.byteLength === 0) { resultCallback({ code: ExportResultCode.SUCCESS }); return; } const eligibleOrigin = new URL(options.url).origin; - options.transport - .authorizedFetch(options.url, { - method: "POST", - headers: { "content-type": "application/x-protobuf" }, - body: body as unknown as Bun.BodyInit, - eligibleOrigin, - }) - .then(async response => { + const sleep = options.sleep ?? defaultSleep; + let backoffMs = INITIAL_BACKOFF_MS; + + for (let attempt = 1; attempt <= MAX_ATTEMPTS; attempt++) { + const lastAttempt = attempt === MAX_ATTEMPTS; + try { + const response = await options.transport.authorizedFetch(options.url, { + method: "POST", + headers: { "content-type": "application/x-protobuf" }, + body: body as unknown as Bun.BodyInit, + eligibleOrigin, + signal: AbortSignal.timeout(EXPORT_TIMEOUT_MS), + }); await response.body?.cancel(); - if (response.ok) resultCallback({ code: ExportResultCode.SUCCESS }); - else + if (response.ok) { + resultCallback({ code: ExportResultCode.SUCCESS }); + return; + } + if (!RETRYABLE_STATUSES.has(response.status) || lastAttempt) { resultCallback({ code: ExportResultCode.FAILED, error: new Error(`collector status ${response.status}`) }); - }) - .catch(error => { - resultCallback({ - code: ExportResultCode.FAILED, - error: error instanceof Error ? error : new Error(String(error)), - }); - }); + return; + } + const retryAfterMs = parseRetryAfterMs(response.headers.get("retry-after")); + await sleep(retryAfterMs ?? jitteredBackoff(backoffMs)); + } catch (error) { + if (lastAttempt) { + resultCallback({ + code: ExportResultCode.FAILED, + error: error instanceof Error ? error : new Error(String(error)), + }); + return; + } + await sleep(jitteredBackoff(backoffMs)); + } + backoffMs *= BACKOFF_MULTIPLIER; + } } export class AuthorizedLogExporter implements LogRecordExporter { diff --git a/packages/coding-agent/test/telemetry-authorized-exporters.test.ts b/packages/coding-agent/test/telemetry-authorized-exporters.test.ts index b825fc899ed..0f2e9f588f1 100644 --- a/packages/coding-agent/test/telemetry-authorized-exporters.test.ts +++ b/packages/coding-agent/test/telemetry-authorized-exporters.test.ts @@ -1,17 +1,26 @@ import { describe, expect, it } from "bun:test"; import { ExportResultCode } from "@opentelemetry/core"; import { LoggerProvider, SimpleLogRecordProcessor } from "@opentelemetry/sdk-logs"; -import { AuthorizedLogExporter, sendAuthorized } from "../src/telemetry/authorized-exporters"; +import { MeterProvider, PeriodicExportingMetricReader } from "@opentelemetry/sdk-metrics"; +import { BasicTracerProvider, SimpleSpanProcessor } from "@opentelemetry/sdk-trace-base"; +import { + AuthorizedLogExporter, + AuthorizedMetricExporter, + AuthorizedTraceExporter, + sendAuthorized, +} from "../src/telemetry/authorized-exporters"; + +type FetchCall = { + url: string; + init: { method?: string; headers?: Bun.HeadersInit; body?: Bun.BodyInit; eligibleOrigin?: string }; +}; function fakeTransport(status = 200) { - const calls: Array<{ - url: string; - init: { method?: string; headers?: Bun.HeadersInit; body?: Bun.BodyInit; eligibleOrigin?: string }; - }> = []; + const calls: FetchCall[] = []; return { calls, transport: { - authorizedFetch: async (url: string, init: (typeof calls)[number]["init"]) => { + authorizedFetch: async (url: string, init: FetchCall["init"]) => { calls.push({ url, init }); return new Response(new Uint8Array(), { status }); }, @@ -19,6 +28,32 @@ function fakeTransport(status = 200) { }; } +/** One scripted step per call: a response status/headers, or "throw" for a network rejection. */ +type ScriptedStep = { status: number; headers?: Record } | "throw"; + +function scriptedTransport(steps: ScriptedStep[]) { + const calls: FetchCall[] = []; + let index = 0; + return { + calls, + transport: { + authorizedFetch: async (url: string, init: FetchCall["init"]) => { + calls.push({ url, init }); + const step = steps[Math.min(index, steps.length - 1)]; + index++; + if (step === "throw") throw new Error("network down"); + return new Response(new Uint8Array(), { status: step.status, headers: step.headers }); + }, + }, + }; +} + +/** No-op sleep that resolves immediately, recording every requested delay. */ +function fakeSleep() { + const delays: number[] = []; + return { delays, sleep: async (ms: number) => void delays.push(ms) }; +} + // Minimal ReadableLogRecord: drive a real LoggerProvider with a // SimpleLogRecordProcessor into the exporter instead of hand-building one. function emitOne(exporter: AuthorizedLogExporter) { @@ -45,10 +80,9 @@ describe("AuthorizedLogExporter", () => { it("maps a collector rejection to a FAILED result", async () => { const { transport } = fakeTransport(500); const results: Array<{ code: number }> = []; - sendAuthorized({ url: "https://telemetry.elide.cloud/v1/logs", transport }, new Uint8Array([1]), result => + await sendAuthorized({ url: "https://telemetry.elide.cloud/v1/logs", transport }, new Uint8Array([1]), result => results.push(result), ); - await new Promise(resolve => setTimeout(resolve, 0)); expect(results[0]?.code).toBe(ExportResultCode.FAILED); }); @@ -62,3 +96,136 @@ describe("AuthorizedLogExporter", () => { expect(calls).toHaveLength(0); }); }); + +describe("sendAuthorized retry/backoff/timeout", () => { + it("retries a 503 once and succeeds, backing off between attempts", async () => { + const { calls, transport } = scriptedTransport([{ status: 503 }, { status: 200 }]); + const { delays, sleep } = fakeSleep(); + const results: Array<{ code: number }> = []; + await sendAuthorized( + { url: "https://telemetry.elide.cloud/v1/logs", transport, sleep }, + new Uint8Array([1]), + result => results.push(result), + ); + expect(calls).toHaveLength(2); + expect(delays).toHaveLength(1); + expect(delays[0]).toBeGreaterThan(0); + expect(results[0]?.code).toBe(ExportResultCode.SUCCESS); + }); + + it("retries a network rejection and succeeds", async () => { + const { calls, transport } = scriptedTransport(["throw", { status: 200 }]); + const { sleep } = fakeSleep(); + const results: Array<{ code: number }> = []; + await sendAuthorized( + { url: "https://telemetry.elide.cloud/v1/logs", transport, sleep }, + new Uint8Array([1]), + result => results.push(result), + ); + expect(calls).toHaveLength(2); + expect(results[0]?.code).toBe(ExportResultCode.SUCCESS); + }); + + it("exhausts retries on persistent 503s and reports FAILED", async () => { + const { calls, transport } = scriptedTransport([ + { status: 503 }, + { status: 503 }, + { status: 503 }, + { status: 503 }, + ]); + const { sleep } = fakeSleep(); + const results: Array<{ code: number; error?: Error }> = []; + await sendAuthorized( + { url: "https://telemetry.elide.cloud/v1/logs", transport, sleep }, + new Uint8Array([1]), + result => results.push(result), + ); + // 3 attempts total: the initial send plus 2 retries — no more calls after that. + expect(calls).toHaveLength(3); + expect(results).toHaveLength(1); + expect(results[0]?.code).toBe(ExportResultCode.FAILED); + expect(results[0]?.error?.message).toContain("503"); + }); + + it("fails immediately on a non-retryable 400, without retrying", async () => { + const { calls, transport } = scriptedTransport([{ status: 400 }, { status: 200 }]); + const { sleep } = fakeSleep(); + const results: Array<{ code: number }> = []; + await sendAuthorized( + { url: "https://telemetry.elide.cloud/v1/logs", transport, sleep }, + new Uint8Array([1]), + result => results.push(result), + ); + expect(calls).toHaveLength(1); + expect(results[0]?.code).toBe(ExportResultCode.FAILED); + }); + + it("honors Retry-After over the computed backoff", async () => { + const { transport } = scriptedTransport([{ status: 429, headers: { "retry-after": "7" } }, { status: 200 }]); + const { delays, sleep } = fakeSleep(); + const results: Array<{ code: number }> = []; + await sendAuthorized( + { url: "https://telemetry.elide.cloud/v1/logs", transport, sleep }, + new Uint8Array([1]), + result => results.push(result), + ); + expect(delays[0]).toBe(7000); + expect(results[0]?.code).toBe(ExportResultCode.SUCCESS); + }); + + it("bounds each attempt with an AbortSignal timeout", async () => { + let sawSignal: AbortSignal | undefined; + const transport = { + authorizedFetch: async (_url: string, init: { signal?: AbortSignal }) => { + sawSignal = init.signal; + return new Response(new Uint8Array(), { status: 200 }); + }, + }; + const results: Array<{ code: number }> = []; + await sendAuthorized({ url: "https://telemetry.elide.cloud/v1/logs", transport }, new Uint8Array([1]), result => + results.push(result), + ); + expect(sawSignal).toBeInstanceOf(AbortSignal); + expect(sawSignal?.aborted).toBe(false); + }); +}); + +describe("AuthorizedTraceExporter", () => { + it("POSTs protobuf spans through the shared send, happy path", async () => { + const { calls, transport } = fakeTransport(); + const exporter = new AuthorizedTraceExporter({ url: "https://telemetry.elide.cloud/v1/traces", transport }); + const provider = new BasicTracerProvider({ spanProcessors: [new SimpleSpanProcessor(exporter)] }); + provider.getTracer("test").startSpan("probe").end(); + await provider.forceFlush(); + + expect(calls).toHaveLength(1); + expect(calls[0]?.url).toBe("https://telemetry.elide.cloud/v1/traces"); + expect(calls[0]?.init.method).toBe("POST"); + expect(new Headers(calls[0]?.init.headers).get("content-type")).toBe("application/x-protobuf"); + expect(calls[0]?.init.eligibleOrigin).toBe("https://telemetry.elide.cloud"); + const body = calls[0]?.init.body as Uint8Array | undefined; + expect(body?.byteLength).toBeGreaterThan(0); + }); +}); + +describe("AuthorizedMetricExporter", () => { + it("POSTs protobuf metrics through the shared send, happy path", async () => { + const { calls, transport } = fakeTransport(); + const exporter = new AuthorizedMetricExporter({ url: "https://telemetry.elide.cloud/v1/metrics", transport }); + const provider = new MeterProvider({ + readers: [new PeriodicExportingMetricReader({ exporter, exportIntervalMillis: 60_000 })], + }); + provider.getMeter("test").createCounter("probe.count").add(1); + await provider.forceFlush(); + + expect(calls).toHaveLength(1); + expect(calls[0]?.url).toBe("https://telemetry.elide.cloud/v1/metrics"); + expect(calls[0]?.init.method).toBe("POST"); + expect(new Headers(calls[0]?.init.headers).get("content-type")).toBe("application/x-protobuf"); + expect(calls[0]?.init.eligibleOrigin).toBe("https://telemetry.elide.cloud"); + const body = calls[0]?.init.body as Uint8Array | undefined; + expect(body?.byteLength).toBeGreaterThan(0); + + await provider.shutdown(); + }); +});