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..2d9f6d4549e --- /dev/null +++ b/packages/coding-agent/src/telemetry/authorized-exporters.ts @@ -0,0 +1,183 @@ +/** + * 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. + * + * `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 { + 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; + signal?: AbortSignal; + }, + ): Promise; +} + +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 async function sendAuthorized( + options: AuthorizedExporterOptions, + body: Uint8Array | undefined, + resultCallback: (result: ExportResult) => void, +): Promise { + if (body === undefined || body.byteLength === 0) { + resultCallback({ code: ExportResultCode.SUCCESS }); + return; + } + const eligibleOrigin = new URL(options.url).origin; + 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 }); + return; + } + if (!RETRYABLE_STATUSES.has(response.status) || lastAttempt) { + resultCallback({ code: ExportResultCode.FAILED, error: new Error(`collector status ${response.status}`) }); + 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 { + 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/src/telemetry/init.ts b/packages/coding-agent/src/telemetry/init.ts index a191760bdd5..63a53860d2d 100644 --- a/packages/coding-agent/src/telemetry/init.ts +++ b/packages/coding-agent/src/telemetry/init.ts @@ -40,10 +40,16 @@ 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"; -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 @@ -52,10 +58,31 @@ 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). */ 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 +376,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 +487,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 +499,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,10 +514,13 @@ 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 })], + processors: [new BatchLogRecordProcessor({ exporter, maxExportBatchSize: MAX_LOG_EXPORT_BATCH_SIZE })], }); logs.setGlobalLoggerProvider(logProvider); otelLogger = logProvider.getLogger("aura"); @@ -620,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/src/telemetry/sink-otlp.ts b/packages/coding-agent/src/telemetry/sink-otlp.ts index 260f8313221..1bf1b20d7db 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,41 @@ 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; + // `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; +} + export interface OtlpSinkDeps { recorder?: AuraMetricRecorder; emitLog: EmitOtelLog; @@ -65,9 +101,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/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-authorized-exporters.test.ts b/packages/coding-agent/test/telemetry-authorized-exporters.test.ts new file mode 100644 index 00000000000..0f2e9f588f1 --- /dev/null +++ b/packages/coding-agent/test/telemetry-authorized-exporters.test.ts @@ -0,0 +1,231 @@ +import { describe, expect, it } from "bun:test"; +import { ExportResultCode } from "@opentelemetry/core"; +import { LoggerProvider, SimpleLogRecordProcessor } from "@opentelemetry/sdk-logs"; +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: FetchCall[] = []; + return { + calls, + transport: { + authorizedFetch: async (url: string, init: FetchCall["init"]) => { + calls.push({ url, init }); + return new Response(new Uint8Array(), { status }); + }, + }, + }; +} + +/** 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) { + // 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 }> = []; + await sendAuthorized({ url: "https://telemetry.elide.cloud/v1/logs", transport }, new Uint8Array([1]), result => + results.push(result), + ); + 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); + }); +}); + +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(); + }); +}); diff --git a/packages/coding-agent/test/telemetry-export.test.ts b/packages/coding-agent/test/telemetry-export.test.ts index a6198e01c8e..a0bf91cdd23 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"; @@ -143,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 @@ -257,3 +276,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(); + }); +}); diff --git a/packages/coding-agent/test/telemetry-sink.test.ts b/packages/coding-agent/test/telemetry-sink.test.ts index 78491358a5f..19acb902831 100644 --- a/packages/coding-agent/test/telemetry-sink.test.ts +++ b/packages/coding-agent/test/telemetry-sink.test.ts @@ -287,3 +287,75 @@ 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); + }); + + 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", + }); + }); +});