From f25083c8417e17c3e048526568bedcf8d9700af5 Mon Sep 17 00:00:00 2001 From: "exe.dev user" Date: Thu, 24 Sep 2026 14:55:50 +0000 Subject: [PATCH 1/3] fix: bound diagnostic span and export memory --- bun.lock | 2 + packages/bcode-laminar/package.json | 2 + packages/bcode-laminar/src/budget.ts | 109 ++++++++++++++ packages/bcode-laminar/src/exporter.ts | 20 +-- packages/bcode-laminar/src/processor.ts | 75 +++++++--- .../bcode-laminar/test/budget-safety.test.ts | 134 ++++++++++++++++++ .../bcode-laminar/test/memory-budget.test.ts | 79 +++++++++++ 7 files changed, 385 insertions(+), 36 deletions(-) create mode 100644 packages/bcode-laminar/src/budget.ts create mode 100644 packages/bcode-laminar/test/budget-safety.test.ts create mode 100644 packages/bcode-laminar/test/memory-budget.test.ts diff --git a/bun.lock b/bun.lock index fec4318cc6..2606761deb 100644 --- a/bun.lock +++ b/bun.lock @@ -113,6 +113,8 @@ "@opencode-ai/plugin": "workspace:*", "@opentelemetry/api": "1.9.0", "@opentelemetry/core": "2.9.0", + "@opentelemetry/otlp-transformer": "0.217.0", + "@opentelemetry/resources": "2.7.1", "@opentelemetry/exporter-trace-otlp-grpc": "0.217.0", "@opentelemetry/exporter-trace-otlp-proto": "0.217.0", "@opentelemetry/sdk-node": "0.217.0", diff --git a/packages/bcode-laminar/package.json b/packages/bcode-laminar/package.json index 7fafff77ba..813de904e8 100644 --- a/packages/bcode-laminar/package.json +++ b/packages/bcode-laminar/package.json @@ -19,6 +19,8 @@ "@opentelemetry/sdk-node": "0.217.0", "@opentelemetry/sdk-trace-base": "2.9.0", "@opentelemetry/core": "2.9.0", + "@opentelemetry/otlp-transformer": "0.217.0", + "@opentelemetry/resources": "2.7.1", "@opentelemetry/exporter-trace-otlp-grpc": "0.217.0", "@opentelemetry/exporter-trace-otlp-proto": "0.217.0", "@grpc/grpc-js": "1.14.4" diff --git a/packages/bcode-laminar/src/budget.ts b/packages/bcode-laminar/src/budget.ts new file mode 100644 index 0000000000..5fe5162729 --- /dev/null +++ b/packages/bcode-laminar/src/budget.ts @@ -0,0 +1,109 @@ +import { resourceFromAttributes } from "@opentelemetry/resources" +import { ProtobufTraceSerializer } from "@opentelemetry/otlp-transformer" +import type { ReadableSpan } from "@opentelemetry/sdk-trace-base" +import type { Attributes, AttributeValue } from "@opentelemetry/api" + +export const RECORD_BYTES = 64 * 1024 +export const QUEUE_BYTES = 4 * 1024 * 1024 +export const QUEUE_RECORDS = 512 +const FIELD_BYTES = 16 * 1024 +const encoder = new TextEncoder() +const decoder = new TextDecoder("utf-8", { fatal: true }) + +export function encodedBytes(spans: ReadableSpan[]): number { + return ProtobufTraceSerializer.serializeRequest(spans)?.byteLength ?? 0 +} + +function text(value: string): string { + // Slice before encoding so a huge diagnostic string cannot create another huge buffer. + const bytes = encoder.encode(value.slice(0, FIELD_BYTES)) + if (bytes.length <= FIELD_BYTES) return decoder.decode(bytes) + let end = FIELD_BYTES + while ((bytes[end]! & 0xc0) === 0x80) end-- + return decoder.decode(bytes.subarray(0, end)) +} + +function attributes(source: Attributes, mark: () => void): Attributes { + const result: Attributes = {} + for (const [key, value] of Object.entries(source)) { + if (typeof value === "string") { + result[text(key)] = text(value) + if (result[text(key)] !== value) mark() + } else if (Array.isArray(value)) { + const values: AttributeValue[] = [] + let remaining = FIELD_BYTES + for (const item of value) { + const bounded = typeof item === "string" ? text(item) : item + const size = typeof bounded === "string" ? encoder.encode(bounded).length : 8 + if (size > remaining) { + mark() + break + } + if (bounded !== item) mark() + remaining -= size + values.push(bounded as AttributeValue) + } + result[text(key)] = values as AttributeValue + } else result[text(key)] = value + if (text(key) !== key) mark() + } + return result +} + +export function boundedSpan(source: ReadableSpan): ReadableSpan | undefined { + let changed = false + const mark = () => { + changed = true + } + const boundedText = (value: string) => { + const result = text(value) + if (result !== value) mark() + return result + } + const context = { ...source.spanContext() } + // A plain snapshot must not retain the original span and its unbounded strings. + const span: ReadableSpan = { + name: boundedText(source.name), + kind: source.kind, + spanContext: () => context, + parentSpanContext: source.parentSpanContext, + startTime: source.startTime, + endTime: source.endTime, + duration: source.duration, + ended: source.ended, + resource: resourceFromAttributes(attributes(source.resource.attributes, mark)), + instrumentationScope: source.instrumentationScope, + droppedAttributesCount: source.droppedAttributesCount, + droppedEventsCount: source.droppedEventsCount, + droppedLinksCount: source.droppedLinksCount, + attributes: attributes(source.attributes, mark), + status: { ...source.status, message: source.status.message && boundedText(source.status.message) }, + events: source.events.map((event) => ({ + ...event, + name: boundedText(event.name), + attributes: attributes(event.attributes ?? {}, mark), + })), + links: source.links.map((link) => ({ ...link, attributes: attributes(link.attributes ?? {}, mark) })), + } + if (changed) span.attributes["bcode.telemetry.truncated"] = true + // Keep identity, timings, numeric usage and error status before diagnostic text. + for (const key of Object.keys(span.attributes).sort( + (a, b) => + (typeof span.attributes[b] === "string" ? span.attributes[b].length : 0) - + (typeof span.attributes[a] === "string" ? span.attributes[a].length : 0), + )) { + if (encodedBytes([span]) <= RECORD_BYTES) return span + if (typeof span.attributes[key] === "number" || typeof span.attributes[key] === "boolean") continue + delete span.attributes[key] + span.attributes["bcode.telemetry.truncated"] = true + } + while (encodedBytes([span]) > RECORD_BYTES && span.events.length) { + span.events.pop() + span.attributes["bcode.telemetry.truncated"] = true + } + while (encodedBytes([span]) > RECORD_BYTES && span.links.length) { + span.links.pop() + span.attributes["bcode.telemetry.truncated"] = true + } + return encodedBytes([span]) <= RECORD_BYTES ? span : undefined +} diff --git a/packages/bcode-laminar/src/exporter.ts b/packages/bcode-laminar/src/exporter.ts index 21da7ed15e..9faf113476 100644 --- a/packages/bcode-laminar/src/exporter.ts +++ b/packages/bcode-laminar/src/exporter.ts @@ -14,19 +14,14 @@ import { makeSpanOtelV2Compatible } from "./compat" export class LaminarSpanExporter implements SpanExporter { private exporter: SpanExporter - constructor(options: { - apiKey: string - baseUrl: string - port: number - timeoutMillis?: number - }) { + constructor(options: { apiKey: string; baseUrl: string; port: number; timeoutMillis?: number }) { const url = options.baseUrl.replace(/\/$/, "").replace(/:\d{1,5}$/g, "") const metadata = new Metadata() metadata.set("authorization", `Bearer ${options.apiKey}`) this.exporter = new ExporterGrpc({ url: `${url}:${options.port}`, metadata, - timeoutMillis: options.timeoutMillis ?? 30000, + timeoutMillis: options.timeoutMillis ?? 10000, }) } @@ -54,14 +49,9 @@ export class LaminarSpanExporter implements SpanExporter { // the runtime never needs LMNR_PROJECT_API_KEY. // // Default path is unchanged: gRPC to Laminar with bearer auth. -export const createSpanExporter = ( - laminar: { apiKey: string; baseUrl: string; port: number }, -): SpanExporter => { - if ( - process.env.OTEL_EXPORTER_OTLP_TRACES_ENDPOINT || - process.env.OTEL_EXPORTER_OTLP_ENDPOINT - ) { - return new ExporterHttpProto() +export const createSpanExporter = (laminar: { apiKey: string; baseUrl: string; port: number }): SpanExporter => { + if (process.env.OTEL_EXPORTER_OTLP_TRACES_ENDPOINT || process.env.OTEL_EXPORTER_OTLP_ENDPOINT) { + return new ExporterHttpProto({ timeoutMillis: 10000 }) } return new LaminarSpanExporter(laminar) } diff --git a/packages/bcode-laminar/src/processor.ts b/packages/bcode-laminar/src/processor.ts index f958d9ef52..0336fe172d 100644 --- a/packages/bcode-laminar/src/processor.ts +++ b/packages/bcode-laminar/src/processor.ts @@ -17,7 +17,8 @@ // - `pino` logger — opencode plugins log via `client.app.log`; the plugin passes // in a logger callback. -import { type Context, type Span, SpanStatusCode, trace } from "@opentelemetry/api" +import { type Context, type Span, SpanStatusCode, TraceFlags, trace } from "@opentelemetry/api" +import { ExportResultCode } from "@opentelemetry/core" import { BatchSpanProcessor, type ReadableSpan, @@ -34,6 +35,7 @@ import { SPAN_SDK_VERSION, } from "./attributes" import { getParentSpanId, makeSpanOtelV2Compatible, type OTelSpanCompat } from "./compat" +import { boundedSpan, encodedBytes, QUEUE_BYTES, QUEUE_RECORDS } from "./budget" import { sessionCurrentTurnSpan } from "./state" import { otelSpanIdToUUID, type StringUUID } from "./utils" @@ -44,16 +46,45 @@ type LogFn = (level: "debug" | "info" | "warn" | "error", message: string) => vo export class OpenCodeLaminarSpanProcessor implements SpanProcessor { private inner: BatchSpanProcessor + private pendingBytes = 0 + private pendingRecords = 0 + private sizes = new WeakMap() private readonly spanIdToPath = new Map() private readonly spanIdLists = new Map() private readonly spawningSpanIdToToolUseId: Record = {} private readonly log: LogFn constructor(options: { exporter: SpanExporter; log?: LogFn }) { - this.inner = new BatchSpanProcessor(options.exporter, { - maxExportBatchSize: 512, - exportTimeoutMillis: 30000, - }) + this.inner = new BatchSpanProcessor( + { + export: (spans, callback) => { + let completed = false + const finish: typeof callback = (result) => { + if (completed) return + completed = true + for (const span of spans) { + this.pendingBytes -= this.sizes.get(span) ?? 0 + this.pendingRecords-- + this.sizes.delete(span) + } + callback(result) + } + try { + options.exporter.export(spans, finish) + } catch { + finish({ code: ExportResultCode.FAILED, error: new Error("Diagnostic export failed") }) + } + }, + shutdown: () => options.exporter.shutdown(), + forceFlush: () => options.exporter.forceFlush?.() ?? Promise.resolve(), + }, + { + // Sixteen 64 KiB records fit in 1 MiB, including each record's OTLP envelope. + maxExportBatchSize: 16, + maxQueueSize: QUEUE_RECORDS, + exportTimeoutMillis: 10000, + }, + ) this.log = options.log ?? (() => {}) } @@ -71,9 +102,7 @@ export class OpenCodeLaminarSpanProcessor implements SpanProcessor { // 1. Re-parent AI-SDK spans onto the live "turn" span for this opencode // session, so each turn becomes its own Laminar trace instead of a // forest of orphan traces. - const sessionId = span.attributes?.["ai.telemetry.metadata.sessionId"] as - | string - | undefined + const sessionId = span.attributes?.["ai.telemetry.metadata.sessionId"] as string | undefined let ctx = parentContext if (sessionId && typeof sessionId === "string") { const parentSpanContext = sessionCurrentTurnSpan[sessionId]?.spanContext() @@ -96,9 +125,7 @@ export class OpenCodeLaminarSpanProcessor implements SpanProcessor { const toolCallNameAttr = span.attributes?.["ai.toolCall.name"] as string | undefined if ( SPAWNING_TOOL_NAMES.includes(span.name) || - (span.name === "ai.toolCall" && - toolCallNameAttr && - SPAWNING_TOOL_NAMES.includes(toolCallNameAttr)) + (span.name === "ai.toolCall" && toolCallNameAttr && SPAWNING_TOOL_NAMES.includes(toolCallNameAttr)) ) { this.spawningSpanIdToToolUseId[otelSpanIdToUUID(span.spanContext().spanId)] = toolCallId } @@ -107,16 +134,12 @@ export class OpenCodeLaminarSpanProcessor implements SpanProcessor { // 3. Stamp Laminar's path attributes. The UI nests by these, NOT by // OTel parentSpanId — must run for every span. const parentPathFromAttribute = span.attributes?.[PARENT_SPAN_PATH] as string[] | undefined - const parentIdsPathFromAttribute = span.attributes?.[PARENT_SPAN_IDS_PATH] as - | StringUUID[] - | undefined + const parentIdsPathFromAttribute = span.attributes?.[PARENT_SPAN_IDS_PATH] as StringUUID[] | undefined const parentSpanId = getParentSpanId(span) const parentSpanPath = - parentPathFromAttribute ?? - (parentSpanId !== undefined ? this.spanIdToPath.get(parentSpanId) : undefined) + parentPathFromAttribute ?? (parentSpanId !== undefined ? this.spanIdToPath.get(parentSpanId) : undefined) const parentSpanIdsPath = - parentIdsPathFromAttribute ?? - (parentSpanId !== undefined ? this.spanIdLists.get(parentSpanId) : []) + parentIdsPathFromAttribute ?? (parentSpanId !== undefined ? this.spanIdLists.get(parentSpanId) : []) const spanId = span.spanContext().spanId const spanPath = parentSpanPath ? [...parentSpanPath, span.name] : [span.name] @@ -143,8 +166,7 @@ export class OpenCodeLaminarSpanProcessor implements SpanProcessor { if (spawningToolCallSpanId) { span.setAttributes({ "lmnr.spawning_subagent.span_id": spawningToolCallSpanId, - "lmnr.spawning_subagent.tool_use_id": - this.spawningSpanIdToToolUseId[spawningToolCallSpanId]!, + "lmnr.spawning_subagent.tool_use_id": this.spawningSpanIdToToolUseId[spawningToolCallSpanId]!, }) } } @@ -176,6 +198,17 @@ export class OpenCodeLaminarSpanProcessor implements SpanProcessor { // entries that the descendant scan iterates), not by the raw hex span id. delete this.spawningSpanIdToToolUseId[otelSpanIdToUUID(spanId)] makeSpanOtelV2Compatible(span) - this.inner.onEnd(span) + if (!(span.spanContext().traceFlags & TraceFlags.SAMPLED)) return + const bounded = boundedSpan(span) + const size = bounded ? encodedBytes([bounded]) : 0 + if (!bounded || this.pendingBytes + size > QUEUE_BYTES || this.pendingRecords >= QUEUE_RECORDS) { + this.log("warn", "Diagnostic span dropped: telemetry memory budget") + return + } + makeSpanOtelV2Compatible(bounded) + this.sizes.set(bounded, size) + this.pendingBytes += size + this.pendingRecords++ + this.inner.onEnd(bounded) } } diff --git a/packages/bcode-laminar/test/budget-safety.test.ts b/packages/bcode-laminar/test/budget-safety.test.ts new file mode 100644 index 0000000000..023a8234f4 --- /dev/null +++ b/packages/bcode-laminar/test/budget-safety.test.ts @@ -0,0 +1,134 @@ +import { expect, test } from "bun:test" +import { context, trace, SpanStatusCode } from "@opentelemetry/api" +import { ExportResultCode, type ExportResult } from "@opentelemetry/core" +import { BasicTracerProvider, type ReadableSpan } from "@opentelemetry/sdk-trace-base" +import { OTLPTraceExporter } from "@opentelemetry/exporter-trace-otlp-proto" +import { OpenCodeLaminarSpanProcessor } from "../src/processor" +import { encodedBytes } from "../src/budget" + +function setup() { + const spans: ReadableSpan[] = [] + const processor = new OpenCodeLaminarSpanProcessor({ + exporter: { + export(items, done) { + spans.push(...items) + done({ code: ExportResultCode.SUCCESS }) + }, + async shutdown() {}, + }, + }) + const provider = new BasicTracerProvider({ spanProcessors: [processor] }) + return { provider, spans, tracer: provider.getTracer("safety") } +} + +test("truncation preserves identity, parenting, usage, status and the original data", async () => { + const { provider, tracer, spans } = setup() + const parent = tracer.startSpan("parent") + const child = tracer.startSpan("child", {}, trace.setSpan(context.active(), parent)) + const original = "😺".repeat(20000) + child.setAttribute("lmnr.span.input", original) + child.setAttribute("gen_ai.usage.input_tokens", 123) + child.setStatus({ code: SpanStatusCode.ERROR, message: "test failure" }) + child.end() + parent.end() + await provider.forceFlush() + await provider.shutdown() + const output = spans.find((span) => span.name === "child")! + expect(output.spanContext().spanId).toBe(child.spanContext().spanId) + expect(output.parentSpanContext?.spanId).toBe(parent.spanContext().spanId) + expect(output.attributes["gen_ai.usage.input_tokens"]).toBe(123) + expect(output.status).toEqual({ code: SpanStatusCode.ERROR, message: "test failure" }) + expect(output.attributes["bcode.telemetry.truncated"]).toBe(true) + expect(String(output.attributes["lmnr.span.input"])).not.toContain("�") + expect((child as unknown as ReadableSpan).attributes["lmnr.span.input"]).toBe(original) + expect(Object.getPrototypeOf(output)).toBe(Object.prototype) + expect(encodedBytes([output])).toBeLessThanOrEqual(65536) +}) + +test("oversized events and links stay bounded", async () => { + const { provider, tracer, spans } = setup() + const span = tracer.startSpan("event", { + links: [ + { + context: { traceId: "a".repeat(32), spanId: "b".repeat(16), traceFlags: 1 }, + attributes: { text: "x".repeat(100000) }, + }, + ], + }) + for (let i = 0; i < 12; i++) span.addEvent("error", { message: "x".repeat(100000) }) + span.end() + await provider.forceFlush() + await provider.shutdown() + expect(spans.length).toBe(1) + expect(encodedBytes(spans)).toBeLessThanOrEqual(65536) + expect(spans[0].attributes["bcode.telemetry.truncated"]).toBe(true) +}) + +test("a stalled exporter cannot grow the queue beyond byte or record budgets and recovers", async () => { + const callbacks: ((result: ExportResult) => void)[] = [] + let exported = 0 + const processor = new OpenCodeLaminarSpanProcessor({ + exporter: { + export(items, done) { + exported += items.length + callbacks.push(done) + }, + async shutdown() {}, + }, + }) + const provider = new BasicTracerProvider({ spanProcessors: [processor] }) + const tracer = provider.getTracer("stalled") + for (let i = 0; i < 1000; i++) { + const span = tracer.startSpan("step") + span.setAttribute("input", "x".repeat(15000)) + span.end() + } + const state = processor as unknown as { pendingBytes: number; pendingRecords: number } + expect(state.pendingBytes).toBeLessThanOrEqual(4 * 1024 * 1024) + expect(state.pendingRecords).toBeLessThanOrEqual(512) + const admitted = state.pendingRecords + const flushing = provider.forceFlush() + await new Promise((resolve) => setTimeout(resolve, 0)) + for (const done of callbacks.splice(0)) done({ code: ExportResultCode.SUCCESS }) + await flushing + expect(exported).toBe(admitted) + expect(state.pendingBytes).toBe(0) + tracer.startSpan("recovered").end() + const recovered = provider.forceFlush() + await new Promise((resolve) => setTimeout(resolve, 0)) + for (const done of callbacks.splice(0)) done({ code: ExportResultCode.SUCCESS }) + await recovered + await provider.shutdown() + expect(exported).toBe(admitted + 1) +}) + +for (const status of [200, 413]) { + test(`real OTLP HTTP transport: ${status} response produces one attempt with a bounded body`, async () => { + const bodies: number[] = [] + const server = Bun.serve({ + hostname: "127.0.0.1", + port: 0, + async fetch(request) { + bodies.push((await request.arrayBuffer()).byteLength) + return new Response(status === 200 ? new Uint8Array() : "too large", { + status, + headers: { "content-type": "application/x-protobuf" }, + }) + }, + }) + const exporter = new OTLPTraceExporter({ url: `http://127.0.0.1:${server.port}/v1/traces`, timeoutMillis: 10000 }) + const { provider, tracer, spans } = setup() + tracer.startSpan("wire").end() + await provider.forceFlush() + await provider.shutdown() + try { + const result = await new Promise((done) => exporter.export(spans, done)) + expect(result.code).toBe(status === 200 ? ExportResultCode.SUCCESS : ExportResultCode.FAILED) + expect(bodies.length).toBe(1) + expect(bodies[0]).toBeLessThanOrEqual(1024 * 1024) + } finally { + await exporter.shutdown() + server.stop(true) + } + }) +} diff --git a/packages/bcode-laminar/test/memory-budget.test.ts b/packages/bcode-laminar/test/memory-budget.test.ts new file mode 100644 index 0000000000..9727bbd672 --- /dev/null +++ b/packages/bcode-laminar/test/memory-budget.test.ts @@ -0,0 +1,79 @@ +import { describe, expect, test } from "bun:test" +import { BasicTracerProvider, type ReadableSpan, type SpanExporter } from "@opentelemetry/sdk-trace-base" +import { ExportResultCode } from "@opentelemetry/core" +import { ProtobufTraceSerializer } from "@opentelemetry/otlp-transformer" +import { OpenCodeLaminarSpanProcessor } from "../src/processor" + +function setup() { + const batches: ReadableSpan[][] = [] + const exporter: SpanExporter = { + export(spans, done) { batches.push([...spans]); done({ code: ExportResultCode.SUCCESS }) }, + async shutdown() {}, + } + const processor = new OpenCodeLaminarSpanProcessor({ exporter }) + const provider = new BasicTracerProvider({ spanProcessors: [processor] }) + return { batches, processor, provider, tracer: provider.getTracer("memory-budget") } +} +function bytes(spans: ReadableSpan[]) { + return ProtobufTraceSerializer.serializeRequest(spans)!.byteLength +} + +describe("diagnostic byte budgets", () => { + test("red: diagnostic text is capped at 16 KiB in UTF-8", async () => { + const ctx = setup() + const span = ctx.tracer.startSpan("step") + span.setAttribute("lmnr.span.input", "😺".repeat(20000)) + span.end() + await ctx.provider.forceFlush() + await ctx.provider.shutdown() + expect(ctx.batches.flat().length).toBe(1) + expect(Buffer.byteLength(String(ctx.batches[0][0].attributes["lmnr.span.input"]))).toBeLessThanOrEqual(16 * 1024) + }) + test("red: one record stays below 64 KiB", async () => { + const ctx = setup() + const span = ctx.tracer.startSpan("step") + for (let i = 0; i < 8; i++) span.setAttribute(`diagnostic.${i}`, "x".repeat(16000)) + span.end() + await ctx.provider.forceFlush() + await ctx.provider.shutdown() + expect(ctx.batches.flat().length).toBe(1) + expect(bytes(ctx.batches.flat())).toBeLessThanOrEqual(64 * 1024) + }) + test("red: uploads split at 1 MiB including wire encoding", async () => { + const ctx = setup() + for (let i = 0; i < 100; i++) { + const span = ctx.tracer.startSpan(`step-${i}`) + span.setAttribute("lmnr.span.input", "x".repeat(15000)) + span.end() + } + await ctx.provider.forceFlush() + await ctx.provider.shutdown() + expect(ctx.batches.flat().length).toBe(100) + expect(Math.max(...ctx.batches.map(bytes))).toBeLessThanOrEqual(1024 * 1024) + }) + test("red: queued diagnostics stay below 4 MiB", async () => { + const ctx = setup() + for (let i = 0; i < 128; i++) { + const span = ctx.tracer.startSpan(`step-${i}`) + for (let j = 0; j < 3; j++) span.setAttribute(`diagnostic.${j}`, "x".repeat(15000)) + span.end() + } + const queued = (ctx.processor as unknown as { inner: { _finishedSpans: ReadableSpan[] } }).inner._finishedSpans + const queuedBytes = bytes(queued) + await ctx.provider.forceFlush() + await ctx.provider.shutdown() + expect(queuedBytes).toBeLessThanOrEqual(4 * 1024 * 1024) + }) + test("control: small diagnostic retains name, timing and input", async () => { + const ctx = setup() + const span = ctx.tracer.startSpan("read page") + span.setAttribute("lmnr.span.input", "cats") + span.end() + await ctx.provider.forceFlush() + await ctx.provider.shutdown() + expect(ctx.batches.flat().length).toBe(1) + expect(ctx.batches[0][0].name).toBe("read page") + expect(ctx.batches[0][0].attributes["lmnr.span.input"]).toBe("cats") + expect(ctx.batches[0][0].duration[0]).toBeGreaterThanOrEqual(0) + }) +}) From d5ead5d748e5a9f9ddfbc853c4809980356763f2 Mon Sep 17 00:00:00 2001 From: "exe.dev user" Date: Thu, 24 Sep 2026 15:06:12 +0000 Subject: [PATCH 2/3] fix: avoid collisions when bounding diagnostic attribute keys --- packages/bcode-laminar/src/budget.ts | 15 +++++++++------ packages/bcode-laminar/test/budget-safety.test.ts | 13 +++++++++++++ 2 files changed, 22 insertions(+), 6 deletions(-) diff --git a/packages/bcode-laminar/src/budget.ts b/packages/bcode-laminar/src/budget.ts index 5fe5162729..a97f0abbbe 100644 --- a/packages/bcode-laminar/src/budget.ts +++ b/packages/bcode-laminar/src/budget.ts @@ -24,11 +24,15 @@ function text(value: string): string { } function attributes(source: Attributes, mark: () => void): Attributes { - const result: Attributes = {} + const result: Attributes = Object.create(null) for (const [key, value] of Object.entries(source)) { + if (text(key) !== key) { + mark() + continue + } if (typeof value === "string") { - result[text(key)] = text(value) - if (result[text(key)] !== value) mark() + result[key] = text(value) + if (result[key] !== value) mark() } else if (Array.isArray(value)) { const values: AttributeValue[] = [] let remaining = FIELD_BYTES @@ -43,9 +47,8 @@ function attributes(source: Attributes, mark: () => void): Attributes { remaining -= size values.push(bounded as AttributeValue) } - result[text(key)] = values as AttributeValue - } else result[text(key)] = value - if (text(key) !== key) mark() + result[key] = values as AttributeValue + } else result[key] = value } return result } diff --git a/packages/bcode-laminar/test/budget-safety.test.ts b/packages/bcode-laminar/test/budget-safety.test.ts index 023a8234f4..cea756fa33 100644 --- a/packages/bcode-laminar/test/budget-safety.test.ts +++ b/packages/bcode-laminar/test/budget-safety.test.ts @@ -132,3 +132,16 @@ for (const status of [200, 413]) { } }) } + +test("overlong attribute keys cannot overwrite a valid key", async () => { + const { provider, tracer, spans } = setup() + const key = "x".repeat(16384) + const span = tracer.startSpan("keys") + span.setAttribute(key, "keep") + span.setAttribute(key + "suffix", "overwrite") + span.end() + await provider.forceFlush() + await provider.shutdown() + expect(spans[0].attributes[key]).toBe("keep") + expect(spans[0].attributes["bcode.telemetry.truncated"]).toBe(true) +}) From 12a406f250c6d248eba284d1a60ef6a9c98fa91b Mon Sep 17 00:00:00 2001 From: "exe.dev user" Date: Thu, 24 Sep 2026 17:12:46 +0000 Subject: [PATCH 3/3] fix(laminar): recover export budget after timeouts --- bun.lock | 2 +- packages/bcode-laminar/package.json | 2 +- packages/bcode-laminar/src/budget.ts | 4 +- packages/bcode-laminar/src/processor.ts | 10 ++++ .../bcode-laminar/test/budget-safety.test.ts | 54 ++++++++++++++++++- .../bcode-laminar/test/memory-budget.test.ts | 12 +++-- 6 files changed, 76 insertions(+), 8 deletions(-) diff --git a/bun.lock b/bun.lock index 2606761deb..f56d638d04 100644 --- a/bun.lock +++ b/bun.lock @@ -114,7 +114,7 @@ "@opentelemetry/api": "1.9.0", "@opentelemetry/core": "2.9.0", "@opentelemetry/otlp-transformer": "0.217.0", - "@opentelemetry/resources": "2.7.1", + "@opentelemetry/resources": "2.9.0", "@opentelemetry/exporter-trace-otlp-grpc": "0.217.0", "@opentelemetry/exporter-trace-otlp-proto": "0.217.0", "@opentelemetry/sdk-node": "0.217.0", diff --git a/packages/bcode-laminar/package.json b/packages/bcode-laminar/package.json index 813de904e8..4147536dbf 100644 --- a/packages/bcode-laminar/package.json +++ b/packages/bcode-laminar/package.json @@ -20,7 +20,7 @@ "@opentelemetry/sdk-trace-base": "2.9.0", "@opentelemetry/core": "2.9.0", "@opentelemetry/otlp-transformer": "0.217.0", - "@opentelemetry/resources": "2.7.1", + "@opentelemetry/resources": "2.9.0", "@opentelemetry/exporter-trace-otlp-grpc": "0.217.0", "@opentelemetry/exporter-trace-otlp-proto": "0.217.0", "@grpc/grpc-js": "1.14.4" diff --git a/packages/bcode-laminar/src/budget.ts b/packages/bcode-laminar/src/budget.ts index a97f0abbbe..f6a175bf18 100644 --- a/packages/bcode-laminar/src/budget.ts +++ b/packages/bcode-laminar/src/budget.ts @@ -74,7 +74,9 @@ export function boundedSpan(source: ReadableSpan): ReadableSpan | undefined { endTime: source.endTime, duration: source.duration, ended: source.ended, - resource: resourceFromAttributes(attributes(source.resource.attributes, mark)), + resource: resourceFromAttributes(attributes(source.resource.attributes, mark), { + schemaUrl: source.resource.schemaUrl, + }), instrumentationScope: source.instrumentationScope, droppedAttributesCount: source.droppedAttributesCount, droppedEventsCount: source.droppedEventsCount, diff --git a/packages/bcode-laminar/src/processor.ts b/packages/bcode-laminar/src/processor.ts index 0336fe172d..b1c4d66441 100644 --- a/packages/bcode-laminar/src/processor.ts +++ b/packages/bcode-laminar/src/processor.ts @@ -62,6 +62,7 @@ export class OpenCodeLaminarSpanProcessor implements SpanProcessor { const finish: typeof callback = (result) => { if (completed) return completed = true + clearTimeout(timer) for (const span of spans) { this.pendingBytes -= this.sizes.get(span) ?? 0 this.pendingRecords-- @@ -69,6 +70,15 @@ export class OpenCodeLaminarSpanProcessor implements SpanProcessor { } callback(result) } + // SDK timeout does not invoke the exporter callback; release our budget too. + const timer = setTimeout( + () => + finish({ + code: ExportResultCode.FAILED, + error: new Error("Diagnostic export timed out"), + }), + 10000, + ) try { options.exporter.export(spans, finish) } catch { diff --git a/packages/bcode-laminar/test/budget-safety.test.ts b/packages/bcode-laminar/test/budget-safety.test.ts index cea756fa33..ec3785bd98 100644 --- a/packages/bcode-laminar/test/budget-safety.test.ts +++ b/packages/bcode-laminar/test/budget-safety.test.ts @@ -4,7 +4,8 @@ import { ExportResultCode, type ExportResult } from "@opentelemetry/core" import { BasicTracerProvider, type ReadableSpan } from "@opentelemetry/sdk-trace-base" import { OTLPTraceExporter } from "@opentelemetry/exporter-trace-otlp-proto" import { OpenCodeLaminarSpanProcessor } from "../src/processor" -import { encodedBytes } from "../src/budget" +import { resourceFromAttributes } from "@opentelemetry/resources" +import { boundedSpan, encodedBytes } from "../src/budget" function setup() { const spans: ReadableSpan[] = [] @@ -145,3 +146,54 @@ test("overlong attribute keys cannot overwrite a valid key", async () => { expect(spans[0].attributes[key]).toBe("keep") expect(spans[0].attributes["bcode.telemetry.truncated"]).toBe(true) }) + +test("export timeout releases budget and late callbacks cannot release it twice", async () => { + const callbacks: ((result: ExportResult) => void)[] = [] + const processor = new OpenCodeLaminarSpanProcessor({ + exporter: { + export(_items, done) { + callbacks.push(done) + }, + async shutdown() {}, + }, + }) + const provider = new BasicTracerProvider({ spanProcessors: [processor] }) + const tracer = provider.getTracer("timeout") + tracer.startSpan("stalled").end() + await provider.forceFlush().catch(() => {}) + await Bun.sleep(50) + const state = processor as unknown as { pendingBytes: number; pendingRecords: number } + expect(state.pendingRecords).toBe(0) + expect(state.pendingBytes).toBe(0) + callbacks.shift()!({ code: ExportResultCode.SUCCESS }) + expect(state.pendingRecords).toBe(0) + tracer.startSpan("recovered").end() + const recovered = provider.forceFlush() + await Bun.sleep(0) + callbacks.shift()!({ code: ExportResultCode.SUCCESS }) + await recovered + expect(state.pendingRecords).toBe(0) + await provider.shutdown() +}, 15000) + +test("bounded resource keeps its schema URL", async () => { + const { provider, tracer, spans } = setup() + tracer.startSpan("schema").end() + await provider.forceFlush() + await provider.shutdown() + const source = { + ...spans[0], + resource: resourceFromAttributes({ service: "test" }, { schemaUrl: "https://opentelemetry.io/schemas/1.24.0" }), + } + expect(boundedSpan(source)?.resource.schemaUrl).toBe(source.resource.schemaUrl) +}) + +test("UTF-8 truncation does not split a surrogate pair at the initial slice", async () => { + const { provider, tracer, spans } = setup() + const span = tracer.startSpan("unicode") + span.setAttribute("text", "x".repeat(16383) + "😺") + span.end() + await provider.forceFlush() + await provider.shutdown() + expect(spans[0].attributes.text).toBe("x".repeat(16383)) +}) diff --git a/packages/bcode-laminar/test/memory-budget.test.ts b/packages/bcode-laminar/test/memory-budget.test.ts index 9727bbd672..82b0861aa0 100644 --- a/packages/bcode-laminar/test/memory-budget.test.ts +++ b/packages/bcode-laminar/test/memory-budget.test.ts @@ -7,7 +7,10 @@ import { OpenCodeLaminarSpanProcessor } from "../src/processor" function setup() { const batches: ReadableSpan[][] = [] const exporter: SpanExporter = { - export(spans, done) { batches.push([...spans]); done({ code: ExportResultCode.SUCCESS }) }, + export(spans, done) { + batches.push([...spans]) + done({ code: ExportResultCode.SUCCESS }) + }, async shutdown() {}, } const processor = new OpenCodeLaminarSpanProcessor({ exporter }) @@ -41,14 +44,15 @@ describe("diagnostic byte budgets", () => { }) test("red: uploads split at 1 MiB including wire encoding", async () => { const ctx = setup() - for (let i = 0; i < 100; i++) { + for (let i = 0; i < 32; i++) { const span = ctx.tracer.startSpan(`step-${i}`) - span.setAttribute("lmnr.span.input", "x".repeat(15000)) + for (let j = 0; j < 4; j++) span.setAttribute(`diagnostic.${j}`, "x".repeat(16200)) span.end() } await ctx.provider.forceFlush() await ctx.provider.shutdown() - expect(ctx.batches.flat().length).toBe(100) + expect(ctx.batches.flat().length).toBe(32) + expect(Math.max(...ctx.batches.map(bytes))).toBeGreaterThan(1000000) expect(Math.max(...ctx.batches.map(bytes))).toBeLessThanOrEqual(1024 * 1024) }) test("red: queued diagnostics stay below 4 MiB", async () => {