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..a31770bfa4ba 100644 --- a/packages/shared/src/observability.test.ts +++ b/packages/shared/src/observability.test.ts @@ -12,7 +12,10 @@ 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"; import { causeErrorTag, compactTraceAttributes, @@ -424,6 +427,98 @@ describe("observability", () => { ), ); + 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* () { + 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: 1_000, + }).pipe(Effect.withTracer(recordingTracer)); + + for (let index = 0; index < 1_024; index += 1) { + sink.push(makeRecord("lost", String(index))); + } + // 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 batch, never a growing backlog. + assert.deepStrictEqual( + write.mock.calls.map(([chunk]) => String(chunk).split("\n").length - 1), + [256, 256, 256, 256, 1, 1, 1, 1, 1], + ); + expect(logs).toEqual([ + { logLevel: "Warn", message: [expect.any(String), { filePath: tracePath }] }, + ]); + + // Once the disk recovers, new records are written and the loss is reported. + yield* fileSystem.remove(tracePath, { recursive: true }); + sink.push(makeRecord("recovered")); + yield* TestClock.adjust("1 second"); + + const records = yield* readTraceRecords(tracePath); + assert.deepStrictEqual( + records.map((record) => record.name), + ["recovered"], + ); + expect(logs[1]).toEqual({ + logLevel: "Info", + 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); + 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([captureLogs, Logger.tracerLogger], { 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..9f92cdc1dfdc 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,6 +372,12 @@ export const makeTraceSink = Effect.fn("makeTraceSink")(function* (options: Trac }); let buffer: Array = []; + // 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; let pendingFlushStats: TraceSinkFlushStats = { logicalWriteBytes: 0, count: 0, @@ -407,9 +414,14 @@ 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. + writeFailing = true; + droppedCount += records.length - persistedCount; return; } + writeFailing = false; pendingFlushStats = { logicalWriteBytes: pendingFlushStats.logicalWriteBytes + chunkBytes, count: pendingFlushStats.count + nextIndex - persistedCount, @@ -419,7 +431,9 @@ export const makeTraceSink = Effect.fn("makeTraceSink")(function* (options: Trac } }; - const flush = Effect.sync(() => { + // 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; pendingFlushStats = { @@ -427,12 +441,28 @@ export const makeTraceSink = Effect.fn("makeTraceSink")(function* (options: Trac count: 0, durationMs: 0, }; - return stats; + if (stats.count > 0 && options.onFlush) { + yield* options.onFlush(stats).pipe(Effect.ignore); + } + if (droppedCount > 0 && !failureReported) { + failureReported = true; + yield* Effect.logWarning("Trace writes are failing, dropping records until they recover", { + filePath: options.filePath, + }); + } + if (failureReported && !writeFailing) { + yield* Effect.logInfo("Trace writes recovered", { + filePath: options.filePath, + droppedCount, + }); + failureReported = false; + droppedCount = 0; + } }).pipe( - Effect.flatMap((stats) => - stats.count > 0 && options.onFlush ? options.onFlush(stats).pipe(Effect.ignore) : Effect.void, - ), Effect.withTracerEnabled(false), + // 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));