From 77b3181fbeb44638c5a2a9640c23b2ca2f5121ac Mon Sep 17 00:00:00 2001 From: Theo Browne Date: Fri, 25 Sep 2026 22:13:07 -0700 Subject: [PATCH 1/5] fix(observability): a failing trace disk no longer stalls the server When a trace file write failed, the sink put the whole backlog back in its buffer. Every later push then retried the full backlog, which took about 40% of the event loop after a few minutes of a failed disk. At about 120k records, the unshift threw RangeError, lost the backlog, and stopped the timed flush. Now a failed write drops the rest of its batch and counts it. The next flush logs the count once as a warning. The buffer stays at or below one batch, and writes continue when the disk recovers. Also fix the documented trace batch window default: it is 1000 ms, not 200 ms. Co-Authored-By: Claude Opus 5.5 (1M context) --- docs/operations/observability.md | 2 +- packages/shared/src/observability.test.ts | 56 +++++++++++++++++++++++ packages/shared/src/observability.ts | 28 ++++++++---- 3 files changed, 76 insertions(+), 10 deletions(-) diff --git a/docs/operations/observability.md b/docs/operations/observability.md index c6537f6eac71..1faa085a4542 100644 --- a/docs/operations/observability.md +++ b/docs/operations/observability.md @@ -584,7 +584,7 @@ Local trace file: - `T3CODE_TRACE_FILE`: override trace file path - `T3CODE_TRACE_MAX_BYTES`: per-file rotation size, default `10485760` - `T3CODE_TRACE_MAX_FILES`: rotated file count, default `10` -- `T3CODE_TRACE_BATCH_WINDOW_MS`: flush window, default `200` +- `T3CODE_TRACE_BATCH_WINDOW_MS`: flush window, default `1000` - `T3CODE_TRACE_MIN_LEVEL`: minimum trace level, default `Info` - `T3CODE_TRACE_TIMING_ENABLED`: enable timing metadata, default `true` diff --git a/packages/shared/src/observability.test.ts b/packages/shared/src/observability.test.ts index b9bf4ef1375d..a11760c26cfd 100644 --- a/packages/shared/src/observability.test.ts +++ b/packages/shared/src/observability.test.ts @@ -12,7 +12,9 @@ import * as Ref from "effect/Ref"; import * as References from "effect/References"; import * as Schema from "effect/Schema"; import * as Tracer from "effect/Tracer"; +import { vi } from "vite-plus/test"; +import { RotatingFileSink } from "./logging.ts"; import { causeErrorTag, compactTraceAttributes, @@ -424,6 +426,60 @@ describe("observability", () => { ), ); + it.effect("drops records after a failed write instead of retrying them on every push", () => { + const warnings: Array = []; + const captureWarnings = Logger.make(({ logLevel, message }) => { + if (logLevel === "Warn") warnings.push(message); + }); + + return Effect.gen(function* () { + const fileSystem = yield* FileSystem.FileSystem; + const path = yield* Path.Path; + const tempDir = yield* fileSystem.makeTempDirectoryScoped({ prefix: "t3-trace-sink-" }); + const tracePath = path.join(tempDir, "shared.trace.ndjson"); + // A directory at the trace path fails every append, like a full disk. + yield* fileSystem.makeDirectory(tracePath); + const write = vi.spyOn(RotatingFileSink.prototype, "write"); + yield* Effect.addFinalizer(() => Effect.sync(() => write.mockRestore())); + + const sink = yield* makeTraceSink({ + filePath: tracePath, + maxBytes: 1024 * 1024, + maxFiles: 2, + batchWindowMs: 10_000, + }); + + for (let index = 0; index < 1_024; index += 1) { + sink.push(makeRecord("lost", String(index))); + } + yield* sink.flush; + + // One write per full batch of 256, never a growing backlog. + assert.deepStrictEqual( + write.mock.calls.map(([chunk]) => String(chunk).split("\n").length - 1), + [256, 256, 256, 256], + ); + expect(warnings).toEqual([ + [expect.any(String), { filePath: tracePath, droppedCount: 1_024 }], + ]); + + // Once the disk recovers, new records are written again. + yield* fileSystem.remove(tracePath, { recursive: true }); + sink.push(makeRecord("recovered")); + yield* sink.flush; + + const records = yield* readTraceRecords(tracePath); + assert.deepStrictEqual( + records.map((record) => record.name), + ["recovered"], + ); + assert.equal(warnings.length, 1); + }).pipe( + Effect.scoped, + Effect.provide(Logger.layer([captureWarnings], { mergeWithExisting: false })), + ); + }); + it.effect("writes nested spans to disk and captures log messages as span events", () => Effect.scoped( Effect.gen(function* () { diff --git a/packages/shared/src/observability.ts b/packages/shared/src/observability.ts index b4dbdf88651a..c880f8a08343 100644 --- a/packages/shared/src/observability.ts +++ b/packages/shared/src/observability.ts @@ -371,6 +371,8 @@ export const makeTraceSink = Effect.fn("makeTraceSink")(function* (options: Trac }); let buffer: Array = []; + // Records lost to failed writes since the last flush reported them. + let droppedCount = 0; let pendingFlushStats: TraceSinkFlushStats = { logicalWriteBytes: 0, count: 0, @@ -407,7 +409,10 @@ export const makeTraceSink = Effect.fn("makeTraceSink")(function* (options: Trac try { sink.write(chunk); } catch { - buffer.unshift(...records.slice(persistedCount)); + // A failing disk (ENOSPC, EACCES, EIO) drops the rest of the batch. + // Retrying it would grow the backlog, and every later push would + // retry all of it. + droppedCount += records.length - persistedCount; return; } pendingFlushStats = { @@ -419,21 +424,26 @@ export const makeTraceSink = Effect.fn("makeTraceSink")(function* (options: Trac } }; - const flush = Effect.sync(() => { + const flush = Effect.gen(function* () { flushUnsafe(); const stats = pendingFlushStats; + const dropped = droppedCount; pendingFlushStats = { logicalWriteBytes: 0, count: 0, durationMs: 0, }; - return stats; - }).pipe( - Effect.flatMap((stats) => - stats.count > 0 && options.onFlush ? options.onFlush(stats).pipe(Effect.ignore) : Effect.void, - ), - Effect.withTracerEnabled(false), - ); + droppedCount = 0; + if (stats.count > 0 && options.onFlush) { + yield* options.onFlush(stats).pipe(Effect.ignore); + } + if (dropped > 0) { + yield* Effect.logWarning("Dropped trace records after a failed write", { + filePath: options.filePath, + droppedCount: dropped, + }); + } + }).pipe(Effect.withTracerEnabled(false)); yield* Effect.addFinalizer(() => flush.pipe(Effect.ignore)); yield* Effect.forkScoped( From f510c7e0ccaa3a191c2ef4dec4a30fde61e1fde7 Mon Sep 17 00:00:00 2001 From: Theo Browne Date: Fri, 25 Sep 2026 22:56:19 -0700 Subject: [PATCH 2/5] fix(observability): log trace write failures once per episode The drop warning ran on every flush while the disk stayed broken, and the tracer logger added each one as an event on the ended makeTraceSink span that the timed flush fiber keeps alive. Now the sink warns once when writes start failing and logs the total dropped once they recover. The flush also runs without the inherited parent span. Co-Authored-By: Claude Opus 5.5 (1M context) --- packages/shared/src/observability.test.ts | 64 +++++++++++++++++------ packages/shared/src/observability.ts | 37 ++++++++++--- 2 files changed, 79 insertions(+), 22 deletions(-) diff --git a/packages/shared/src/observability.test.ts b/packages/shared/src/observability.test.ts index a11760c26cfd..cef3fe8de459 100644 --- a/packages/shared/src/observability.test.ts +++ b/packages/shared/src/observability.test.ts @@ -12,6 +12,7 @@ import * as Ref from "effect/Ref"; import * as References from "effect/References"; import * as Schema from "effect/Schema"; import * as Tracer from "effect/Tracer"; +import * as TestClock from "effect/testing/TestClock"; import { vi } from "vite-plus/test"; import { RotatingFileSink } from "./logging.ts"; @@ -426,10 +427,18 @@ describe("observability", () => { ), ); - it.effect("drops records after a failed write instead of retrying them on every push", () => { - const warnings: Array = []; - const captureWarnings = Logger.make(({ logLevel, message }) => { - if (logLevel === "Warn") warnings.push(message); + it.effect("drops records after a failed write and logs once per failure episode", () => { + const logs: Array<{ readonly logLevel: string; readonly message: unknown }> = []; + const captureLogs = Logger.make(({ logLevel, message }) => { + logs.push({ logLevel, message }); + }); + const spans: Array = []; + const recordingTracer = Tracer.make({ + span: (options) => { + const span = new Tracer.NativeSpan(options); + spans.push(span); + return span; + }, }); return Effect.gen(function* () { @@ -446,37 +455,62 @@ describe("observability", () => { filePath: tracePath, maxBytes: 1024 * 1024, maxFiles: 2, - batchWindowMs: 10_000, - }); + batchWindowMs: 1_000, + }).pipe(Effect.withTracer(recordingTracer)); for (let index = 0; index < 1_024; index += 1) { sink.push(makeRecord("lost", String(index))); } - yield* sink.flush; + // Timed flushes run in the fiber forked inside the makeTraceSink span. + for (let index = 0; index < 5; index += 1) { + sink.push(makeRecord("lost")); + yield* TestClock.adjust("1 second"); + } - // One write per full batch of 256, never a growing backlog. + // One write per batch, never a growing backlog. assert.deepStrictEqual( write.mock.calls.map(([chunk]) => String(chunk).split("\n").length - 1), - [256, 256, 256, 256], + [256, 256, 256, 256, 1, 1, 1, 1, 1], ); - expect(warnings).toEqual([ - [expect.any(String), { filePath: tracePath, droppedCount: 1_024 }], + expect(logs).toEqual([ + { logLevel: "Warn", message: [expect.any(String), { filePath: tracePath }] }, ]); - // Once the disk recovers, new records are written again. + // Once the disk recovers, new records are written and the loss is reported. yield* fileSystem.remove(tracePath, { recursive: true }); sink.push(makeRecord("recovered")); - yield* sink.flush; + yield* TestClock.adjust("1 second"); const records = yield* readTraceRecords(tracePath); assert.deepStrictEqual( records.map((record) => record.name), ["recovered"], ); - assert.equal(warnings.length, 1); + expect(logs[1]).toEqual({ + logLevel: "Info", + message: [expect.any(String), { filePath: tracePath, droppedCount: 1_029 }], + }); + + // A new failure episode warns again. + yield* fileSystem.remove(tracePath); + yield* fileSystem.makeDirectory(tracePath); + sink.push(makeRecord("lost-again")); + yield* TestClock.adjust("1 second"); + assert.deepStrictEqual( + logs.map((log) => log.logLevel), + ["Warn", "Info", "Warn"], + ); + + // The ended makeTraceSink span is never released, so it must not collect log events. + assert.deepStrictEqual( + spans.map((span) => [span.name, span.events.length]), + [["makeTraceSink", 0]], + ); }).pipe( Effect.scoped, - Effect.provide(Logger.layer([captureWarnings], { mergeWithExisting: false })), + Effect.provide( + Logger.layer([captureLogs, Logger.tracerLogger], { mergeWithExisting: false }), + ), ); }); diff --git a/packages/shared/src/observability.ts b/packages/shared/src/observability.ts index c880f8a08343..08e618c2556d 100644 --- a/packages/shared/src/observability.ts +++ b/packages/shared/src/observability.ts @@ -1,4 +1,5 @@ import * as Cause from "effect/Cause"; +import * as Context from "effect/Context"; import * as Effect from "effect/Effect"; import type * as Exit from "effect/Exit"; import * as ExitRuntime from "effect/Exit"; @@ -371,8 +372,12 @@ export const makeTraceSink = Effect.fn("makeTraceSink")(function* (options: Trac }); let buffer: Array = []; - // Records lost to failed writes since the last flush reported them. + // Failure episode state. The latest write result says if the disk is + // failing. Records lost since the episode started are counted until a write + // succeeds and the flush reports the recovery. + let writeFailing = false; let droppedCount = 0; + let failureReported = false; let pendingFlushStats: TraceSinkFlushStats = { logicalWriteBytes: 0, count: 0, @@ -412,9 +417,11 @@ export const makeTraceSink = Effect.fn("makeTraceSink")(function* (options: Trac // A failing disk (ENOSPC, EACCES, EIO) drops the rest of the batch. // Retrying it would grow the backlog, and every later push would // retry all of it. + writeFailing = true; droppedCount += records.length - persistedCount; return; } + writeFailing = false; pendingFlushStats = { logicalWriteBytes: pendingFlushStats.logicalWriteBytes + chunkBytes, count: pendingFlushStats.count + nextIndex - persistedCount, @@ -424,26 +431,42 @@ export const makeTraceSink = Effect.fn("makeTraceSink")(function* (options: Trac } }; + // Logs only when writes start failing and when they recover, so a disk that + // stays broken costs one line, not one per flush. It drops the parent span: + // the timed fiber inherits the makeTraceSink span, which has ended but stays + // referenced for the life of the sink, and the tracer logger would add every + // log to it as an event. const flush = Effect.gen(function* () { flushUnsafe(); const stats = pendingFlushStats; - const dropped = droppedCount; pendingFlushStats = { logicalWriteBytes: 0, count: 0, durationMs: 0, }; - droppedCount = 0; if (stats.count > 0 && options.onFlush) { yield* options.onFlush(stats).pipe(Effect.ignore); } - if (dropped > 0) { - yield* Effect.logWarning("Dropped trace records after a failed write", { + if (droppedCount > 0 && !failureReported) { + failureReported = true; + yield* Effect.logWarning("Trace writes are failing, dropping records until they recover", { filePath: options.filePath, - droppedCount: dropped, }); } - }).pipe(Effect.withTracerEnabled(false)); + if (failureReported && !writeFailing) { + yield* Effect.logInfo("Trace writes recovered", { + filePath: options.filePath, + droppedCount, + }); + failureReported = false; + droppedCount = 0; + } + }).pipe( + Effect.withTracerEnabled(false), + Effect.updateContext((context: Context.Context) => + Context.omit(Tracer.ParentSpan)(context), + ), + ); yield* Effect.addFinalizer(() => flush.pipe(Effect.ignore)); yield* Effect.forkScoped( From b0281956ff57d66b5485c074c25691339b6e4f71 Mon Sep 17 00:00:00 2001 From: Theo Browne Date: Fri, 25 Sep 2026 23:02:33 -0700 Subject: [PATCH 3/5] docs(observability): note trace failure episodes are judged per flush Co-Authored-By: Claude Opus 5.5 (1M context) --- packages/shared/src/observability.ts | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/packages/shared/src/observability.ts b/packages/shared/src/observability.ts index 08e618c2556d..1a50dffd7aa2 100644 --- a/packages/shared/src/observability.ts +++ b/packages/shared/src/observability.ts @@ -373,8 +373,9 @@ export const makeTraceSink = Effect.fn("makeTraceSink")(function* (options: Trac let buffer: Array = []; // Failure episode state. The latest write result says if the disk is - // failing. Records lost since the episode started are counted until a write - // succeeds and the flush reports the recovery. + // failing. Flush checks it once per window, so a disk that fails and + // recovers more than once inside one window stays one episode. Records lost + // since the episode started are counted until a flush sees a good write. let writeFailing = false; let droppedCount = 0; let failureReported = false; From b61b040580d03d883b8367e15bbbe26f30c5f384 Mon Sep 17 00:00:00 2001 From: Theo Browne Date: Fri, 25 Sep 2026 23:19:39 -0700 Subject: [PATCH 4/5] test(observability): healthy flushes after trace recovery log nothing Co-Authored-By: Claude Opus 5.5 (1M context) --- packages/shared/src/observability.test.ts | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/packages/shared/src/observability.test.ts b/packages/shared/src/observability.test.ts index cef3fe8de459..a31770bfa4ba 100644 --- a/packages/shared/src/observability.test.ts +++ b/packages/shared/src/observability.test.ts @@ -491,6 +491,11 @@ describe("observability", () => { message: [expect.any(String), { filePath: tracePath, droppedCount: 1_029 }], }); + // Healthy flushes after the recovery log nothing. + sink.push(makeRecord("healthy")); + yield* TestClock.adjust("1 second"); + expect(logs).toHaveLength(2); + // A new failure episode warns again. yield* fileSystem.remove(tracePath); yield* fileSystem.makeDirectory(tracePath); From dd32fb61f22078ffa714cdb8d738bceaec3d4cda Mon Sep 17 00:00:00 2001 From: Theo Browne Date: Sat, 26 Sep 2026 01:46:06 -0700 Subject: [PATCH 5/5] refactor(observability): tighten trace sink comments and parent span removal Co-Authored-By: Claude Opus 5.5 (1M context) --- packages/shared/src/observability.ts | 20 ++++++++------------ 1 file changed, 8 insertions(+), 12 deletions(-) diff --git a/packages/shared/src/observability.ts b/packages/shared/src/observability.ts index 1a50dffd7aa2..9f92cdc1dfdc 100644 --- a/packages/shared/src/observability.ts +++ b/packages/shared/src/observability.ts @@ -372,10 +372,9 @@ export const makeTraceSink = Effect.fn("makeTraceSink")(function* (options: Trac }); let buffer: Array = []; - // Failure episode state. The latest write result says if the disk is - // failing. Flush checks it once per window, so a disk that fails and - // recovers more than once inside one window stays one episode. Records lost - // since the episode started are counted until a flush sees a good write. + // A failure episode starts at the first dropped record and ends when a flush + // sees that the latest write succeeded. Flush checks once per window, so a + // disk that flaps inside one window stays one episode. let writeFailing = false; let droppedCount = 0; let failureReported = false; @@ -432,11 +431,8 @@ export const makeTraceSink = Effect.fn("makeTraceSink")(function* (options: Trac } }; - // Logs only when writes start failing and when they recover, so a disk that - // stays broken costs one line, not one per flush. It drops the parent span: - // the timed fiber inherits the makeTraceSink span, which has ended but stays - // referenced for the life of the sink, and the tracer logger would add every - // log to it as an event. + // Logs once when writes start failing and once when they recover, so a disk + // that stays broken costs one line, not one per flush. const flush = Effect.gen(function* () { flushUnsafe(); const stats = pendingFlushStats; @@ -464,9 +460,9 @@ export const makeTraceSink = Effect.fn("makeTraceSink")(function* (options: Trac } }).pipe( Effect.withTracerEnabled(false), - Effect.updateContext((context: Context.Context) => - Context.omit(Tracer.ParentSpan)(context), - ), + // The timed fiber inherits the makeTraceSink span. That span has ended but + // lives as long as the sink, so the tracer logger must not add logs to it. + Effect.updateContext(Context.omit(Tracer.ParentSpan)), ); yield* Effect.addFinalizer(() => flush.pipe(Effect.ignore));