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

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions bun.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

1 change: 1 addition & 0 deletions package.json
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down
1 change: 1 addition & 0 deletions packages/coding-agent/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -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:",
Expand Down
183 changes: 183 additions & 0 deletions packages/coding-agent/src/telemetry/authorized-exporters.ts
Original file line number Diff line number Diff line change
@@ -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<Response>;
}

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<void>;
}

/** 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<void> {
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<void> {
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<void> {}

async shutdown(): Promise<void> {}
}

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<void> {}
}

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<void> {}

async shutdown(): Promise<void> {}
}
80 changes: 72 additions & 8 deletions packages/coding-agent/src/telemetry/init.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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";
Expand Down Expand Up @@ -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<Settings, "get"> | undefined,
processEnv: Record<string, string | undefined> = 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;
Expand Down Expand Up @@ -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)],
Expand All @@ -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 })],
Expand All @@ -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");
Expand Down Expand Up @@ -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,
Expand Down
Loading