From 82970ea159d6a08d5ba3b133ede436879a600c26 Mon Sep 17 00:00:00 2001 From: Theo Browne Date: Fri, 25 Sep 2026 22:14:34 -0700 Subject: [PATCH 1/2] fix(observability): span events no longer pile up without limit A fiber that outlives its span kept logging into the ended span. Each log added an event that was never written, and it stayed in memory for as long as the fiber lived. Open spans also kept every log event. - A span ignores events after it ends. - A span keeps its first 128 events (the OpenTelemetry SDK default) and records the rest in span.dropped_events_count. - Event names keep 500 characters, like attribute strings. - A failure cause keeps 8,000 characters. The longest cause in local traces is 4,346, so this only stops pathological ones. Co-Authored-By: Claude Opus 5.5 (1M context) --- docs/operations/observability.md | 10 ++- packages/shared/src/observability.test.ts | 98 ++++++++++++++++++++++- packages/shared/src/observability.ts | 41 ++++++++-- 3 files changed, 137 insertions(+), 12 deletions(-) diff --git a/docs/operations/observability.md b/docs/operations/observability.md index c6537f6eac71..f26783a5dc83 100644 --- a/docs/operations/observability.md +++ b/docs/operations/observability.md @@ -50,8 +50,10 @@ Important fields common to both record types: - `attributes`: structured context - `events`: embedded logs and custom events -`effect-span` records also contain `exit` with `Success`, `Failure`, or `Interrupted`. `otlp-span` -records instead carry OTLP resource, scope, and optional status fields. +`effect-span` records also contain `exit` with `Success`, `Failure`, or `Interrupted`. Their +attribute strings and event names keep the first 500 characters (`db.query.text` keeps 200). A +failure's `cause` keeps the first 8,000 characters, which fits a normal stack and `[cause]` chain. +`otlp-span` records instead carry OTLP resource, scope, and optional status fields. The `TraceRecord`, `EffectTraceRecord`, and `OtlpTraceRecord` schemas live in `packages/shared/src/observability.ts`. @@ -533,7 +535,9 @@ yield * Effect.logInfo("starting provider turn"); yield * Effect.logDebug("waiting for approval response"); ``` -Those messages show up as span events because `Logger.tracerLogger` is installed. +Those messages show up as span events because `Logger.tracerLogger` is installed. A span keeps +its first 128 events and counts the rest in its `span.dropped_events_count` attribute. Logs from a +fiber that outlives its span are not written. ### Use The Pipeable Metrics API diff --git a/packages/shared/src/observability.test.ts b/packages/shared/src/observability.test.ts index b9bf4ef1375d..51c9f999fc5e 100644 --- a/packages/shared/src/observability.test.ts +++ b/packages/shared/src/observability.test.ts @@ -62,6 +62,7 @@ const TraceRecordLine = Schema.Struct({ exit: Schema.optional( Schema.Struct({ _tag: Schema.String, + cause: Schema.optional(Schema.String), }), ), }); @@ -97,7 +98,7 @@ const readTraceRecords = Effect.fn("readTraceRecords")(function* (tracePath: str .map((line) => decodeTraceRecordLine(line)); }); -const makeTestLayer = (tracePath: string) => +const makeTestLayer = (tracePath: string, delegate?: Tracer.Tracer) => Layer.mergeAll( Layer.effect( Tracer.Tracer, @@ -106,12 +107,27 @@ const makeTestLayer = (tracePath: string) => maxBytes: 1024 * 1024, maxFiles: 2, batchWindowMs: 10_000, + ...(delegate ? { delegate } : {}), }), ), Logger.layer([Logger.tracerLogger], { mergeWithExisting: false }), Layer.succeed(References.MinimumLogLevel, "Info"), ); +// A delegate tracer that keeps its spans, to see what the local tracer +// forwards to it. +const makeRecordingDelegate = () => { + const spans: Array = []; + const delegate = Tracer.make({ + span: (options) => { + const span = new Tracer.NativeSpan(options); + spans.push(span); + return span; + }, + }); + return { delegate, spans }; +}; + const nodeServicesIt = it.layer(NodeServices.layer); describe("truncateTraceAttributes", () => { @@ -506,6 +522,86 @@ describe("observability", () => { }), ), ); + + it.effect("clamps oversized failure causes and log event names", () => + Effect.scoped( + Effect.gen(function* () { + const fileSystem = yield* FileSystem.FileSystem; + const path = yield* Path.Path; + const tempDir = yield* fileSystem.makeTempDirectoryScoped({ prefix: "t3-local-tracer-" }); + const tracePath = path.join(tempDir, "shared.trace.ndjson"); + + yield* Effect.scoped( + Effect.exit( + Effect.logWarning("w".repeat(5_000)).pipe( + Effect.andThen(Effect.fail("f".repeat(20_000))), + Effect.withSpan("failing-span"), + Effect.provide(makeTestLayer(tracePath)), + ), + ), + ); + + const records = yield* readTraceRecords(tracePath); + assert.equal(records[0]?.events[0]?.name, `${"w".repeat(500)}…[truncated]`); + const cause = records[0]?.exit?.cause ?? ""; + assert.equal(cause.length, 8_000 + "…[truncated]".length); + assert.isTrue(cause.endsWith("…[truncated]")); + }), + ), + ); + + it.effect("ignores events added after a span has ended", () => + Effect.scoped( + Effect.gen(function* () { + const fileSystem = yield* FileSystem.FileSystem; + const path = yield* Path.Path; + const tempDir = yield* fileSystem.makeTempDirectoryScoped({ prefix: "t3-local-tracer-" }); + const tracePath = path.join(tempDir, "shared.trace.ndjson"); + const { delegate, spans } = makeRecordingDelegate(); + + const span = yield* Effect.scoped( + Effect.logInfo("during").pipe( + Effect.andThen(Effect.currentSpan), + Effect.withSpan("ended-span"), + Effect.provide(makeTestLayer(tracePath, delegate)), + ), + ); + span.event("after end", 0n); + + assert.deepEqual( + spans[0]?.events.map(([name]) => name), + ["during"], + ); + }), + ), + ); + + it.effect("caps the events kept per span and records how many were dropped", () => + Effect.scoped( + Effect.gen(function* () { + const fileSystem = yield* FileSystem.FileSystem; + const path = yield* Path.Path; + const tempDir = yield* fileSystem.makeTempDirectoryScoped({ prefix: "t3-local-tracer-" }); + const tracePath = path.join(tempDir, "shared.trace.ndjson"); + const { delegate, spans } = makeRecordingDelegate(); + + yield* Effect.scoped( + Effect.forEach(Arr.range(1, 200), (index) => Effect.logInfo(`event ${index}`), { + discard: true, + }).pipe( + Effect.withSpan("chatty-span"), + Effect.provide(makeTestLayer(tracePath, delegate)), + ), + ); + + const records = yield* readTraceRecords(tracePath); + assert.equal(records[0]?.events.length, 128); + assert.equal(records[0]?.events.at(-1)?.name, "event 128"); + assert.equal(records[0]?.attributes["span.dropped_events_count"], 72); + assert.equal(spans[0]?.events.length, 128); + }), + ), + ); }); }); diff --git a/packages/shared/src/observability.ts b/packages/shared/src/observability.ts index b4dbdf88651a..eee374128d3a 100644 --- a/packages/shared/src/observability.ts +++ b/packages/shared/src/observability.ts @@ -259,19 +259,25 @@ export function compactTraceAttributes( return Object.fromEntries(entries); } +// A failure repeats its pretty cause on every failing ancestor span. Normal +// causes (message, stack, and [cause] chain) are well under the cap; it only +// stops a pathological one, such as an error that embeds a large payload. +const TRACE_CAUSE_MAX_LENGTH = 8_000; + function formatTraceExit(exit: Exit.Exit): EffectTraceRecord["exit"] { if (ExitRuntime.isSuccess(exit)) { return { _tag: "Success" }; } + const cause = truncateTraceString(Cause.pretty(exit.cause), TRACE_CAUSE_MAX_LENGTH); if (Cause.hasInterruptsOnly(exit.cause)) { return { _tag: "Interrupted", - cause: Cause.pretty(exit.cause), + cause, }; } return { _tag: "Failure", - cause: Cause.pretty(exit.cause), + cause, }; } @@ -280,14 +286,18 @@ const TRACE_ATTRIBUTE_TRUNCATED_LENGTH = 200; const TRACE_ATTRIBUTE_TRUNCATION_SUFFIX = "…[truncated]"; const ALWAYS_TRUNCATED_TRACE_ATTRIBUTES: ReadonlySet = new Set(["db.query.text"]); +function truncateTraceString(value: string, maxLength: number): string { + return value.length <= maxLength + ? value + : `${value.slice(0, maxLength)}${TRACE_ATTRIBUTE_TRUNCATION_SUFFIX}`; +} + // Clamps strings nested inside already-normalized attribute values (arrays and // plain objects from normalizeJsonValue, e.g. an Error's `stack`). Returns the // input reference when nothing was clamped. function truncateNestedValue(value: unknown): unknown { if (typeof value === "string") { - return value.length <= TRACE_ATTRIBUTE_MAX_LENGTH - ? value - : `${value.slice(0, TRACE_ATTRIBUTE_MAX_LENGTH)}${TRACE_ATTRIBUTE_TRUNCATION_SUFFIX}`; + return truncateTraceString(value, TRACE_ATTRIBUTE_MAX_LENGTH); } if (Array.isArray(value)) { const truncated = value.map(truncateNestedValue); @@ -318,8 +328,7 @@ export function truncateTraceAttributes(attributes: TraceAttributes): TraceAttri if (typeof value === "string" && ALWAYS_TRUNCATED_TRACE_ATTRIBUTES.has(key)) { if (value.length <= TRACE_ATTRIBUTE_TRUNCATED_LENGTH) continue; truncated ??= { ...attributes }; - truncated[key] = - `${value.slice(0, TRACE_ATTRIBUTE_TRUNCATED_LENGTH)}${TRACE_ATTRIBUTE_TRUNCATION_SUFFIX}`; + truncated[key] = truncateTraceString(value, TRACE_ATTRIBUTE_TRUNCATED_LENGTH); continue; } const next = truncateNestedValue(value); @@ -349,7 +358,8 @@ function spanToTraceRecord(span: SerializableSpan): EffectTraceRecord { compactTraceAttributes(Object.fromEntries(span.attributes)), ), events: span.events.map(([name, startTime, attributes]) => ({ - name, + // A log event is named after its whole formatted message. + name: truncateTraceString(name, TRACE_ATTRIBUTE_MAX_LENGTH), timeUnixNano: String(startTime), attributes: truncateTraceAttributes(compactTraceAttributes(attributes)), })), @@ -457,6 +467,11 @@ export const makeTraceSink = Effect.fn("makeTraceSink")(function* (options: Trac } satisfies TraceSink; }); +// Long-lived spans (a whole RPC stream, for example) would otherwise keep +// every log event in memory until they end. 128 is the OpenTelemetry SDK +// default event count limit. +const TRACE_SPAN_MAX_EVENTS = 128; + class LocalFileSpan implements Tracer.Span { readonly _tag = "Span"; readonly name: string; @@ -471,6 +486,7 @@ class LocalFileSpan implements Tracer.Span { status: Tracer.SpanStatus; attributes: Map; events: Array<[name: string, startTime: bigint, attributes: Record]>; + private droppedEventCount = 0; private readonly delegate: Tracer.Span; private readonly push: (record: EffectTraceRecord) => void; @@ -498,6 +514,9 @@ class LocalFileSpan implements Tracer.Span { } end(endTime: bigint, exit: Exit.Exit): void { + if (this.droppedEventCount > 0) { + this.attribute("span.dropped_events_count", this.droppedEventCount); + } this.status = { _tag: "Ended", startTime: this.status.startTime, @@ -516,7 +535,13 @@ class LocalFileSpan implements Tracer.Span { this.delegate.attribute(key, value); } + // Events on an ended span are never written, so they are not kept either. event(name: string, startTime: bigint, attributes?: Record): void { + if (this.status._tag === "Ended") return; + if (this.events.length >= TRACE_SPAN_MAX_EVENTS) { + this.droppedEventCount += 1; + return; + } const nextAttributes = attributes ?? {}; this.events.push([name, startTime, nextAttributes]); this.delegate.event(name, startTime, nextAttributes); From e760d2e003e87efcba1581fa6ee6796afca3008f Mon Sep 17 00:00:00 2001 From: Theo Browne Date: Fri, 25 Sep 2026 23:09:45 -0700 Subject: [PATCH 2/2] fix(observability): keep a span's newest events, not its first A long-lived span, such as an RPC stream span, kept its first 128 events and dropped the rest, so a warning or error logged late in its life was lost. It now drops the oldest event instead, like the OpenTelemetry SDK. The delegate span gets the kept events at end, so both hold the same set. Co-Authored-By: Claude Opus 5.5 (1M context) --- docs/operations/observability.md | 4 ++-- packages/shared/src/observability.test.ts | 29 +++++++++++++---------- packages/shared/src/observability.ts | 15 +++++++----- 3 files changed, 27 insertions(+), 21 deletions(-) diff --git a/docs/operations/observability.md b/docs/operations/observability.md index f26783a5dc83..b2f0a2940991 100644 --- a/docs/operations/observability.md +++ b/docs/operations/observability.md @@ -536,8 +536,8 @@ yield * Effect.logDebug("waiting for approval response"); ``` Those messages show up as span events because `Logger.tracerLogger` is installed. A span keeps -its first 128 events and counts the rest in its `span.dropped_events_count` attribute. Logs from a -fiber that outlives its span are not written. +its newest 128 events and counts the older ones it drops in its `span.dropped_events_count` +attribute. Logs from a fiber that outlives its span are not written. ### Use The Pipeable Metrics API diff --git a/packages/shared/src/observability.test.ts b/packages/shared/src/observability.test.ts index 51c9f999fc5e..137f0020f6f5 100644 --- a/packages/shared/src/observability.test.ts +++ b/packages/shared/src/observability.test.ts @@ -557,26 +557,23 @@ describe("observability", () => { const path = yield* Path.Path; const tempDir = yield* fileSystem.makeTempDirectoryScoped({ prefix: "t3-local-tracer-" }); const tracePath = path.join(tempDir, "shared.trace.ndjson"); - const { delegate, spans } = makeRecordingDelegate(); const span = yield* Effect.scoped( - Effect.logInfo("during").pipe( - Effect.andThen(Effect.currentSpan), + Effect.currentSpan.pipe( + Effect.tap((span) => Effect.sync(() => span.event("during", 1n))), Effect.withSpan("ended-span"), - Effect.provide(makeTestLayer(tracePath, delegate)), + Effect.provide(makeTestLayer(tracePath)), ), ); - span.event("after end", 0n); + span.event("after end", 2n); - assert.deepEqual( - spans[0]?.events.map(([name]) => name), - ["during"], - ); + // A fiber that outlives its span would grow this list. + assert.deepPropertyVal(span, "events", [["during", 1n, {}]]); }), ), ); - it.effect("caps the events kept per span and records how many were dropped", () => + it.effect("keeps the newest events per span and records how many were dropped", () => Effect.scoped( Effect.gen(function* () { const fileSystem = yield* FileSystem.FileSystem; @@ -595,10 +592,16 @@ describe("observability", () => { ); const records = yield* readTraceRecords(tracePath); - assert.equal(records[0]?.events.length, 128); - assert.equal(records[0]?.events.at(-1)?.name, "event 128"); + const names = records[0]?.events.map((event) => event.name); + assert.deepEqual( + names, + Arr.range(73, 200).map((index) => `event ${index}`), + ); assert.equal(records[0]?.attributes["span.dropped_events_count"], 72); - assert.equal(spans[0]?.events.length, 128); + assert.deepEqual( + spans[0]?.events.map(([name]) => name), + names, + ); }), ), ); diff --git a/packages/shared/src/observability.ts b/packages/shared/src/observability.ts index eee374128d3a..5e4432167092 100644 --- a/packages/shared/src/observability.ts +++ b/packages/shared/src/observability.ts @@ -468,8 +468,8 @@ export const makeTraceSink = Effect.fn("makeTraceSink")(function* (options: Trac }); // Long-lived spans (a whole RPC stream, for example) would otherwise keep -// every log event in memory until they end. 128 is the OpenTelemetry SDK -// default event count limit. +// every log event in memory until they end. Like the OpenTelemetry SDK, a +// span keeps its newest 128 events and drops the oldest. const TRACE_SPAN_MAX_EVENTS = 128; class LocalFileSpan implements Tracer.Span { @@ -523,6 +523,11 @@ class LocalFileSpan implements Tracer.Span { endTime, exit, }; + // The delegate gets the kept events only now, so it holds the same bounded + // set. The OTLP delegate reads a span's events only at end. + for (const [name, startTime, attributes] of this.events) { + this.delegate.event(name, startTime, attributes); + } this.delegate.end(endTime, exit); if (this.sampled) { @@ -539,12 +544,10 @@ class LocalFileSpan implements Tracer.Span { event(name: string, startTime: bigint, attributes?: Record): void { if (this.status._tag === "Ended") return; if (this.events.length >= TRACE_SPAN_MAX_EVENTS) { + this.events.shift(); this.droppedEventCount += 1; - return; } - const nextAttributes = attributes ?? {}; - this.events.push([name, startTime, nextAttributes]); - this.delegate.event(name, startTime, nextAttributes); + this.events.push([name, startTime, attributes ?? {}]); } addLinks(links: ReadonlyArray): void {