Skip to content
Merged
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.

2 changes: 2 additions & 0 deletions packages/bcode-laminar/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -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.9.0",
"@opentelemetry/exporter-trace-otlp-grpc": "0.217.0",
"@opentelemetry/exporter-trace-otlp-proto": "0.217.0",
"@grpc/grpc-js": "1.14.4"
Expand Down
114 changes: 114 additions & 0 deletions packages/bcode-laminar/src/budget.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,114 @@
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))
Comment thread
sarath-menon marked this conversation as resolved.
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 = Object.create(null)
for (const [key, value] of Object.entries(source)) {
if (text(key) !== key) {
mark()
continue
}
if (typeof value === "string") {
result[key] = text(value)
if (result[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[key] = values as AttributeValue
} else result[key] = value
}
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), {
schemaUrl: source.resource.schemaUrl,
}),
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
}
20 changes: 5 additions & 15 deletions packages/bcode-laminar/src/exporter.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
})
}

Expand Down Expand Up @@ -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)
}
85 changes: 64 additions & 21 deletions packages/bcode-laminar/src/processor.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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"

Expand All @@ -44,16 +46,55 @@ 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<ReadableSpan, number>()
private readonly spanIdToPath = new Map<string, string[]>()
private readonly spanIdLists = new Map<string, StringUUID[]>()
private readonly spawningSpanIdToToolUseId: Record<string, string> = {}
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
clearTimeout(timer)
for (const span of spans) {
this.pendingBytes -= this.sizes.get(span) ?? 0
this.pendingRecords--
this.sizes.delete(span)
}
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)
Comment thread
sarath-menon marked this conversation as resolved.
} 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 ?? (() => {})
}

Expand All @@ -71,9 +112,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()
Expand All @@ -96,9 +135,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
}
Expand All @@ -107,16 +144,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]
Expand All @@ -143,8 +176,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]!,
})
}
}
Expand Down Expand Up @@ -176,6 +208,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)
}
}
Loading
Loading