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: 1 addition & 1 deletion docs/operations/observability.md
Original file line number Diff line number Diff line change
Expand Up @@ -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`

Expand Down
95 changes: 95 additions & 0 deletions packages/shared/src/observability.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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<Tracer.NativeSpan> = [];
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* () {
Expand Down
42 changes: 36 additions & 6 deletions packages/shared/src/observability.ts
Original file line number Diff line number Diff line change
@@ -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";
Expand Down Expand Up @@ -371,6 +372,12 @@ export const makeTraceSink = Effect.fn("makeTraceSink")(function* (options: Trac
});

let buffer: Array<string> = [];
// 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,
Expand Down Expand Up @@ -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,
Expand All @@ -419,20 +431,38 @@ 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 = {
logicalWriteBytes: 0,
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) {
Comment thread
coderabbitai[bot] marked this conversation as resolved.
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)<never>),
);

yield* Effect.addFinalizer(() => flush.pipe(Effect.ignore));
Expand Down
Loading