From 3fc08601a0095bcf53850e53e75c86026f63a314 Mon Sep 17 00:00:00 2001 From: Theo Browne Date: Fri, 25 Sep 2026 22:16:36 -0700 Subject: [PATCH 1/2] perf(server): opening Diagnostics no longer loads the whole trace ring into memory Settings > Diagnostics read every rotated trace file into memory as one string (about 110 MB) and parsed all of it in one sync pass. On a real 107 MB ring that added about 320 MB of RSS and blocked the event loop for about 380 ms, right when a slow, swapping user opens the page. Diagnostics now streams each file with decodeText and splitLines and feeds lines into an incremental aggregator, like `t3 trace summary` does. The CLI now uses the same shared helper. Output is unchanged: on the same real ring the old and new results are deep-equal. Peak RSS growth drops to about 50 MB and the longest event loop block to under 10 ms. Co-Authored-By: Claude Opus 5.5 (1M context) --- apps/server/src/cli/trace.ts | 23 +- .../src/diagnostics/TraceDiagnostics.test.ts | 318 +++++++++------ .../src/diagnostics/TraceDiagnostics.ts | 382 +++++++++--------- 3 files changed, 372 insertions(+), 351 deletions(-) diff --git a/apps/server/src/cli/trace.ts b/apps/server/src/cli/trace.ts index 7b8a83592ef8..61a209471d6e 100644 --- a/apps/server/src/cli/trace.ts +++ b/apps/server/src/cli/trace.ts @@ -13,11 +13,10 @@ import * as Effect from "effect/Effect"; import * as FileSystem from "effect/FileSystem"; import * as Option from "effect/Option"; import * as Schema from "effect/Schema"; -import * as Stream from "effect/Stream"; import { Command, Flag } from "effect/unstable/cli"; import * as ServerConfig from "../config.ts"; -import { toRotatedTracePaths, TraceFileReadError } from "../diagnostics/TraceDiagnostics.ts"; +import { streamTraceFileLines, toRotatedTracePaths } from "../diagnostics/TraceDiagnostics.ts"; import { resolveBaseDir } from "../os-jank.ts"; import { baseDirFlag, DurationFromString, traceFileConfig, traceMaxFilesConfig } from "./config.ts"; @@ -187,27 +186,9 @@ const traceSummaryCommand = Command.make("summary", { ? (yield* Clock.currentTimeMillis) - Duration.toMillis(flags.since.value) : undefined; const summarizer = makeTraceSpanSummary(sinceMs); - // Stream each file so only one chunk of text is in memory at a time. yield* Effect.forEach( toRotatedTracePaths(traceFilePath, yield* traceMaxFilesConfig), - (path) => - fs.stream(path).pipe( - Stream.decodeText, - Stream.splitLines, - Stream.runForEachArray((lines) => Effect.sync(() => lines.forEach(summarizer.addLine))), - Effect.catchTags({ - PlatformError: (cause) => - cause.reason._tag === "NotFound" - ? Effect.void - : Effect.fail( - new TraceFileReadError({ - traceFilePath: path, - causeTag: cause.reason._tag, - cause, - }), - ), - }), - ), + (path) => streamTraceFileLines(fs, path, summarizer.addLine), { discard: true }, ); const summary = summarizer.finish(); diff --git a/apps/server/src/diagnostics/TraceDiagnostics.test.ts b/apps/server/src/diagnostics/TraceDiagnostics.test.ts index 70bb4dc815c3..8ea72b95b9e0 100644 --- a/apps/server/src/diagnostics/TraceDiagnostics.test.ts +++ b/apps/server/src/diagnostics/TraceDiagnostics.test.ts @@ -7,6 +7,7 @@ import * as Logger from "effect/Logger"; import * as Option from "effect/Option"; import * as PlatformError from "effect/PlatformError"; import * as References from "effect/References"; +import * as Stream from "effect/Stream"; import * as TraceDiagnostics from "./TraceDiagnostics.ts"; @@ -40,72 +41,78 @@ function record(input: { }); } +const traceFilePath = "/tmp/server.trace.ndjson"; +const readAt = DateTime.makeUnsafe("2026-05-05T10:00:00.000Z"); + +/** Aggregates whole lines in memory, the way diagnostics read traces before streaming. */ +function aggregateLines(lines: ReadonlyArray) { + const aggregator = TraceDiagnostics.makeTraceDiagnosticsAggregator(); + lines.forEach(aggregator.addLine); + return aggregator.finish({ + traceFilePath, + scannedFilePaths: TraceDiagnostics.toRotatedTracePaths(traceFilePath, 1), + readAt, + }); +} + +/** Reads the trace file and one rotated backup through a fake file system. */ +function readTraces(fileSystem: Partial) { + return TraceDiagnostics.readTraceDiagnostics({ traceFilePath, maxFiles: 1, readAt }).pipe( + Effect.provide(TraceDiagnostics.layer.pipe(Layer.provide(FileSystem.layerNoop(fileSystem)))), + ); +} + describe("TraceDiagnostics", () => { it.effect("aggregates failures, slow spans, log levels, and parse errors", () => Effect.sync(() => { - const diagnostics = TraceDiagnostics.aggregateTraceDiagnostics({ - traceFilePath: "/tmp/server.trace.ndjson", - readAt: DateTime.makeUnsafe("2026-05-05T10:00:00.000Z"), - slowSpanThresholdMs: 1_000, - files: [ - { - path: "/tmp/server.trace.ndjson.1", - text: [ - record({ - name: "server.getConfig", - traceId: "trace-a", - spanId: "span-a", - startMs: 1_000, - durationMs: 50, - }), - "not-json", - ].join("\n"), - }, - { - path: "/tmp/server.trace.ndjson", - text: [ - record({ - name: "orchestration.dispatch", - traceId: "trace-b", - spanId: "span-b", - startMs: 2_000, - durationMs: 1_500, - exit: { _tag: "Failure", cause: "Provider crashed" }, - events: [ - { - name: "provider failed", - timeUnixNano: ns(3_400), - attributes: { "effect.logLevel": "Error" }, - }, - ], - }), - record({ - name: "orchestration.dispatch", - traceId: "trace-c", - spanId: "span-c", - startMs: 4_000, - durationMs: 250, - exit: { _tag: "Failure", cause: "Provider crashed" }, - }), - record({ - name: "git.status", - traceId: "trace-d", - spanId: "span-d", - startMs: 5_000, - durationMs: 25, - exit: { _tag: "Interrupted", cause: "Interrupted" }, - events: [ - { - name: "status delayed", - timeUnixNano: ns(5_010), - attributes: { "effect.logLevel": "Warning" }, - }, - ], - }), - ].join("\n"), - }, - ], - }); + const diagnostics = aggregateLines([ + record({ + name: "server.getConfig", + traceId: "trace-a", + spanId: "span-a", + startMs: 1_000, + durationMs: 50, + }), + "not-json", + record({ + name: "orchestration.dispatch", + traceId: "trace-b", + spanId: "span-b", + startMs: 2_000, + durationMs: 1_500, + exit: { _tag: "Failure", cause: "Provider crashed" }, + events: [ + { + name: "provider failed", + timeUnixNano: ns(3_400), + attributes: { "effect.logLevel": "Error" }, + }, + ], + }), + record({ + name: "orchestration.dispatch", + traceId: "trace-c", + spanId: "span-c", + startMs: 4_000, + durationMs: 250, + exit: { _tag: "Failure", cause: "Provider crashed" }, + }), + record({ + name: "git.status", + traceId: "trace-d", + spanId: "span-d", + startMs: 5_000, + durationMs: 25, + exit: { _tag: "Interrupted", cause: "Interrupted" }, + events: [ + { + name: "status delayed", + timeUnixNano: ns(5_010), + attributes: { "effect.logLevel": "Warning" }, + }, + ], + }), + ]); assert.equal(diagnostics.recordCount, 4); assert.equal(DateTime.formatIso(diagnostics.readAt), "2026-05-05T10:00:00.000Z"); @@ -139,47 +146,110 @@ describe("TraceDiagnostics", () => { ); it.effect("returns a not-found diagnostic when no files are available", () => - Effect.sync(() => { - const diagnostics = TraceDiagnostics.aggregateTraceDiagnostics({ - traceFilePath: "/tmp/missing.trace.ndjson", - readAt: DateTime.makeUnsafe("2026-05-05T10:00:00.000Z"), - files: [], - }); + Effect.gen(function* () { + const diagnostics = yield* readTraces({}); assert.equal(diagnostics.recordCount, 0); assert.equal(Option.getOrUndefined(diagnostics.error)?.kind, "trace-file-not-found"); }), ); - it.effect("preserves full failure causes and log messages", () => - Effect.sync(() => { - const longCause = `VcsProcessSpawnError: ${"missing executable ".repeat(80)}`.trim(); - const longMessage = `provider warning: ${"retrying command ".repeat(80)}`.trim(); - const diagnostics = TraceDiagnostics.aggregateTraceDiagnostics({ - traceFilePath: "/tmp/server.trace.ndjson", - readAt: DateTime.makeUnsafe("2026-05-05T10:00:00.000Z"), - files: [ - { - path: "/tmp/server.trace.ndjson", - text: record({ - name: "VcsProcess.run", - traceId: "trace-long", - spanId: "span-long", + it.effect("streams rotated files into the same result as reading them whole", () => + Effect.gen(function* () { + // CRLF and LF endings plus multi-byte text, served one byte per chunk so + // chunks split lines, line endings, and characters. + const files = new Map([ + [ + `${traceFilePath}.1`, + [ + record({ + name: "server.getConfig", + traceId: "trace-a", + spanId: "span-a", startMs: 1_000, + durationMs: 50, + }), + "not-json", + record({ + name: "orchestration.dispatch", + traceId: "trace-b", + spanId: "span-b", + startMs: 2_000, + durationMs: 1_500, + exit: { _tag: "Failure", cause: "Provider crashed: café 🔥" }, + }), + "", + ].join("\r\n"), + ], + [ + traceFilePath, + [ + record({ + name: "git.status", + traceId: "trace-c", + spanId: "span-c", + startMs: 3_000, durationMs: 25, - exit: { _tag: "Failure", cause: longCause }, + exit: { _tag: "Interrupted", cause: "Interrupted" }, events: [ { - name: longMessage, - timeUnixNano: ns(1_010), + name: "status delayed ⏳", + timeUnixNano: ns(3_010), attributes: { "effect.logLevel": "Warning" }, }, ], }), - }, + "", + record({ + name: "orchestration.dispatch", + traceId: "trace-d", + spanId: "span-d", + startMs: 4_000, + durationMs: 250, + exit: { _tag: "Failure", cause: "Provider crashed: café 🔥" }, + }), + ].join("\n"), ], + ]); + const encoder = new TextEncoder(); + + const diagnostics = yield* readTraces({ + stream: (path) => + Stream.fromIterable( + Array.from(encoder.encode(files.get(path)), (byte) => Uint8Array.of(byte)), + ), }); + assert.equal(diagnostics.recordCount, 4); + assert.deepStrictEqual( + diagnostics, + aggregateLines([...files.values()].flatMap((text) => text.split(/\r?\n/))), + ); + }), + ); + + it.effect("preserves full failure causes and log messages", () => + Effect.sync(() => { + const longCause = `VcsProcessSpawnError: ${"missing executable ".repeat(80)}`.trim(); + const longMessage = `provider warning: ${"retrying command ".repeat(80)}`.trim(); + const diagnostics = aggregateLines([ + record({ + name: "VcsProcess.run", + traceId: "trace-long", + spanId: "span-long", + startMs: 1_000, + durationMs: 25, + exit: { _tag: "Failure", cause: longCause }, + events: [ + { + name: longMessage, + timeUnixNano: ns(1_010), + attributes: { "effect.logLevel": "Warning" }, + }, + ], + }), + ]); + assert.equal(diagnostics.latestFailures[0]?.cause, longCause); assert.equal(diagnostics.commonFailures[0]?.cause, longCause); assert.equal(diagnostics.latestWarningAndErrorLogs[0]?.message, longMessage); @@ -188,45 +258,34 @@ describe("TraceDiagnostics", () => { it.effect("keeps loaded trace data when one rotated trace file fails to read", () => Effect.gen(function* () { - const traceFilePath = "/tmp/server.trace.ndjson"; const readFailure = PlatformError.systemError({ _tag: "PermissionDenied", module: "FileSystem", - method: "readFileString", + method: "open", description: "permission denied", pathOrDescriptor: `${traceFilePath}.1`, }); - const fileSystemLayer = FileSystem.layerNoop({ - readFileString: (path) => - path === `${traceFilePath}.1` - ? Effect.fail(readFailure) - : Effect.succeed( - record({ - name: "server.getConfig", - traceId: "trace-a", - spanId: "span-a", - startMs: 1_000, - durationMs: 50, - }), - ), - }); const logAnnotations: Array> = []; const logger = Logger.make((options) => { logAnnotations.push({ ...options.fiber.getRef(References.CurrentLogAnnotations) }); }); - const diagnostics = yield* TraceDiagnostics.readTraceDiagnostics({ - traceFilePath, - maxFiles: 1, - readAt: DateTime.makeUnsafe("2026-05-05T10:00:00.000Z"), - }).pipe( - Effect.provide( - Layer.mergeAll( - TraceDiagnostics.layer.pipe(Layer.provide(fileSystemLayer)), - Logger.layer([logger], { mergeWithExisting: false }), - ), - ), - ); + const diagnostics = yield* readTraces({ + stream: (path) => + path === `${traceFilePath}.1` + ? Stream.fail(readFailure) + : Stream.make( + new TextEncoder().encode( + record({ + name: "server.getConfig", + traceId: "trace-a", + spanId: "span-a", + startMs: 1_000, + durationMs: 50, + }), + ), + ), + }).pipe(Effect.provide(Logger.layer([logger], { mergeWithExisting: false }))); assert.equal(diagnostics.recordCount, 1); assert.equal( @@ -253,24 +312,17 @@ describe("TraceDiagnostics", () => { it.effect("keeps only the slowest span occurrences while aggregating large inputs", () => Effect.sync(() => { - const diagnostics = TraceDiagnostics.aggregateTraceDiagnostics({ - traceFilePath: "/tmp/server.trace.ndjson", - readAt: DateTime.makeUnsafe("2026-05-05T10:00:00.000Z"), - files: [ - { - path: "/tmp/server.trace.ndjson", - text: Array.from({ length: 25 }, (_, index) => - record({ - name: `span-${index}`, - traceId: `trace-${index}`, - spanId: `span-${index}`, - startMs: index * 1_000, - durationMs: index, - }), - ).join("\n"), - }, - ], - }); + const diagnostics = aggregateLines( + Array.from({ length: 25 }, (_, index) => + record({ + name: `span-${index}`, + traceId: `trace-${index}`, + spanId: `span-${index}`, + startMs: index * 1_000, + durationMs: index, + }), + ), + ); assert.equal(diagnostics.recordCount, 25); assert.equal(diagnostics.slowestSpans.length, 10); diff --git a/apps/server/src/diagnostics/TraceDiagnostics.ts b/apps/server/src/diagnostics/TraceDiagnostics.ts index afb7b4e923bb..79d62567ea6a 100644 --- a/apps/server/src/diagnostics/TraceDiagnostics.ts +++ b/apps/server/src/diagnostics/TraceDiagnostics.ts @@ -16,6 +16,7 @@ import * as Option from "effect/Option"; import * as PlatformError from "effect/PlatformError"; import * as Result from "effect/Result"; import * as Schema from "effect/Schema"; +import * as Stream from "effect/Stream"; interface TraceRecordLike { readonly name?: unknown; @@ -63,16 +64,6 @@ export class TraceDiagnostics extends Context.Service< } >()("t3/diagnostics/TraceDiagnostics") {} -interface TraceDiagnosticsInput { - readonly traceFilePath: string; - readonly files: ReadonlyArray<{ readonly path: string; readonly text: string }>; - readonly scannedFilePaths?: ReadonlyArray; - readonly slowSpanThresholdMs?: number; - readonly readAt: DateTime.Utc; - readonly error?: TraceDiagnosticsErrorSummary; - readonly partialFailure?: boolean; -} - interface TraceDiagnosticsErrorSummary { readonly kind: ServerTraceDiagnosticsErrorKind; readonly message: string; @@ -191,26 +182,13 @@ function insertBoundedSlowestSpan( } } -export function aggregateTraceDiagnostics( - input: TraceDiagnosticsInput, -): ServerTraceDiagnosticsResult { - const readAt = input.readAt; - const slowSpanThresholdMs = input.slowSpanThresholdMs ?? DEFAULT_SLOW_SPAN_THRESHOLD_MS; - const scannedFilePaths = input.scannedFilePaths ?? input.files.map((file) => file.path); - if (input.files.length === 0) { - return makeEmptyDiagnostics({ - traceFilePath: input.traceFilePath, - scannedFilePaths, - readAt, - slowSpanThresholdMs, - error: input.error ?? { - kind: "trace-file-not-found", - message: "No local trace files were found.", - }, - ...(input.partialFailure ? { partialFailure: true } : {}), - }); - } - +/** + * Folds trace NDJSON into diagnostics. Call `addLine` once per line as the + * rotated files stream in, then `finish` for the result. + */ +export function makeTraceDiagnosticsAggregator( + slowSpanThresholdMs = DEFAULT_SLOW_SPAN_THRESHOLD_MS, +) { let parseErrorCount = 0; let recordCount = 0; let failureCount = 0; @@ -229,182 +207,196 @@ export function aggregateTraceDiagnostics( const latestWarningAndErrorLogs: ServerTraceDiagnosticsLogEvent[] = []; const logLevelCounts: Record = {}; - for (const file of input.files) { - const lines = file.text.split(/\r?\n/); - for (const line of lines) { - if (line.trim().length === 0) continue; - - let parsed: unknown; - try { - parsed = JSON.parse(line); - } catch { - parseErrorCount += 1; - continue; - } + const addLine = (line: string) => { + if (line.trim().length === 0) return; - if (!isRecordObject(parsed)) { - parseErrorCount += 1; - continue; - } + let parsed: unknown; + try { + parsed = JSON.parse(line); + } catch { + parseErrorCount += 1; + return; + } - const name = toStringValue(parsed.name); - const traceId = toStringValue(parsed.traceId); - const spanId = toStringValue(parsed.spanId); - const durationMs = toNumberValue(parsed.durationMs); - const endedAt = unixNanoToDateTime(parsed.endTimeUnixNano); - const startedAt = unixNanoToDateTime(parsed.startTimeUnixNano); + if (!isRecordObject(parsed)) { + parseErrorCount += 1; + return; + } - if (!name || !traceId || !spanId || durationMs === null || !endedAt) { - parseErrorCount += 1; - continue; - } + const name = toStringValue(parsed.name); + const traceId = toStringValue(parsed.traceId); + const spanId = toStringValue(parsed.spanId); + const durationMs = toNumberValue(parsed.durationMs); + const endedAt = unixNanoToDateTime(parsed.endTimeUnixNano); + const startedAt = unixNanoToDateTime(parsed.startTimeUnixNano); - recordCount += 1; - firstSpanAt = - startedAt && (firstSpanAt === null || DateTime.isLessThan(startedAt, firstSpanAt)) - ? startedAt - : firstSpanAt; - lastSpanAt = - lastSpanAt === null || DateTime.isGreaterThan(endedAt, lastSpanAt) ? endedAt : lastSpanAt; - - const exitTag = readExitTag(parsed.exit); - const isFailure = exitTag === "Failure"; - const isInterrupted = exitTag === "Interrupted"; - if (isFailure) failureCount += 1; - if (isInterrupted) interruptionCount += 1; - - const spanSummary = spansByName.get(name) ?? { - count: 0, - failureCount: 0, - totalDurationMs: 0, - maxDurationMs: 0, - }; - spanSummary.count += 1; - spanSummary.totalDurationMs += durationMs; - spanSummary.maxDurationMs = Math.max(spanSummary.maxDurationMs, durationMs); - if (isFailure) spanSummary.failureCount += 1; - spansByName.set(name, spanSummary); - - const spanItem = { name, durationMs, endedAt, traceId, spanId }; - if (durationMs >= slowSpanThresholdMs) { - slowSpanCount += 1; - } - insertBoundedSlowestSpan(slowestSpans, spanItem); - - if (isFailure) { - const cause = readExitCause(parsed.exit); - latestFailures.push({ ...spanItem, cause }); - - const failureKey = `${name}\0${cause}`; - const existing = failuresByKey.get(failureKey); - const isLatestFailure = !existing || DateTime.isGreaterThan(endedAt, existing.lastSeenAt); - failuresByKey.set(failureKey, { - name, - cause, - count: (existing?.count ?? 0) + 1, - lastSeenAt: isLatestFailure ? endedAt : existing!.lastSeenAt, - traceId: isLatestFailure ? traceId : existing!.traceId, - spanId: isLatestFailure ? spanId : existing!.spanId, - }); - } + if (!name || !traceId || !spanId || durationMs === null || !endedAt) { + parseErrorCount += 1; + return; + } + + recordCount += 1; + firstSpanAt = + startedAt && (firstSpanAt === null || DateTime.isLessThan(startedAt, firstSpanAt)) + ? startedAt + : firstSpanAt; + lastSpanAt = + lastSpanAt === null || DateTime.isGreaterThan(endedAt, lastSpanAt) ? endedAt : lastSpanAt; + + const exitTag = readExitTag(parsed.exit); + const isFailure = exitTag === "Failure"; + const isInterrupted = exitTag === "Interrupted"; + if (isFailure) failureCount += 1; + if (isInterrupted) interruptionCount += 1; + + const spanSummary = spansByName.get(name) ?? { + count: 0, + failureCount: 0, + totalDurationMs: 0, + maxDurationMs: 0, + }; + spanSummary.count += 1; + spanSummary.totalDurationMs += durationMs; + spanSummary.maxDurationMs = Math.max(spanSummary.maxDurationMs, durationMs); + if (isFailure) spanSummary.failureCount += 1; + spansByName.set(name, spanSummary); + + const spanItem = { name, durationMs, endedAt, traceId, spanId }; + if (durationMs >= slowSpanThresholdMs) { + slowSpanCount += 1; + } + insertBoundedSlowestSpan(slowestSpans, spanItem); + + if (isFailure) { + const cause = readExitCause(parsed.exit); + latestFailures.push({ ...spanItem, cause }); + + const failureKey = `${name}\0${cause}`; + const existing = failuresByKey.get(failureKey); + const isLatestFailure = !existing || DateTime.isGreaterThan(endedAt, existing.lastSeenAt); + failuresByKey.set(failureKey, { + name, + cause, + count: (existing?.count ?? 0) + 1, + lastSeenAt: isLatestFailure ? endedAt : existing!.lastSeenAt, + traceId: isLatestFailure ? traceId : existing!.traceId, + spanId: isLatestFailure ? spanId : existing!.spanId, + }); + } - if (Array.isArray(parsed.events)) { - for (const rawEvent of parsed.events) { - if (!isTraceEvent(rawEvent)) continue; - const attributes = readEventAttributes(rawEvent); - const level = toStringValue(attributes["effect.logLevel"]); - if (!level) continue; - - logLevelCounts[level] = (logLevelCounts[level] ?? 0) + 1; - const normalizedLevel = level.toLowerCase(); - if ( - normalizedLevel !== "warning" && - normalizedLevel !== "warn" && - normalizedLevel !== "error" && - normalizedLevel !== "fatal" - ) { - continue; - } - - const seenAt = unixNanoToDateTime(rawEvent.timeUnixNano) ?? endedAt; - const message = toStringValue(rawEvent.name)?.trim() ?? "Log event"; - latestWarningAndErrorLogs.push({ - spanName: name, - level, - message, - seenAt, - traceId, - spanId, - }); + if (Array.isArray(parsed.events)) { + for (const rawEvent of parsed.events) { + if (!isTraceEvent(rawEvent)) continue; + const attributes = readEventAttributes(rawEvent); + const level = toStringValue(attributes["effect.logLevel"]); + if (!level) continue; + + logLevelCounts[level] = (logLevelCounts[level] ?? 0) + 1; + const normalizedLevel = level.toLowerCase(); + if ( + normalizedLevel !== "warning" && + normalizedLevel !== "warn" && + normalizedLevel !== "error" && + normalizedLevel !== "fatal" + ) { + continue; } + + const seenAt = unixNanoToDateTime(rawEvent.timeUnixNano) ?? endedAt; + const message = toStringValue(rawEvent.name)?.trim() ?? "Log event"; + latestWarningAndErrorLogs.push({ + spanName: name, + level, + message, + seenAt, + traceId, + spanId, + }); } } - } - - const topSpansByCount: ServerTraceDiagnosticsSpanSummary[] = [...spansByName.entries()] - .map(([name, span]) => ({ - name, - count: span.count, - failureCount: span.failureCount, - totalDurationMs: span.totalDurationMs, - averageDurationMs: span.count > 0 ? span.totalDurationMs / span.count : 0, - maxDurationMs: span.maxDurationMs, - })) - .toSorted((left, right) => right.count - left.count || right.maxDurationMs - left.maxDurationMs) - .slice(0, TOP_LIMIT); + }; - return { - traceFilePath: input.traceFilePath, - scannedFilePaths, - readAt, - recordCount, - parseErrorCount, - firstSpanAt: Option.fromNullishOr(firstSpanAt), - lastSpanAt: Option.fromNullishOr(lastSpanAt), - failureCount, - interruptionCount, - slowSpanThresholdMs, - slowSpanCount, - logLevelCounts, - topSpansByCount, - slowestSpans, - commonFailures: [...failuresByKey.values()] - .toSorted( - (left, right) => - right.count - left.count || - DateTime.toEpochMillis(right.lastSeenAt) - DateTime.toEpochMillis(left.lastSeenAt), - ) - .slice(0, TOP_LIMIT), - latestFailures: latestFailures - .toSorted( - (left, right) => - DateTime.toEpochMillis(right.endedAt) - DateTime.toEpochMillis(left.endedAt), - ) - .slice(0, RECENT_LIMIT), - latestWarningAndErrorLogs: latestWarningAndErrorLogs + const finish = (input: { + readonly traceFilePath: string; + readonly scannedFilePaths: ReadonlyArray; + readonly readAt: DateTime.Utc; + readonly error?: TraceDiagnosticsErrorSummary; + readonly partialFailure?: boolean; + }): ServerTraceDiagnosticsResult => { + const topSpansByCount: ServerTraceDiagnosticsSpanSummary[] = [...spansByName.entries()] + .map(([name, span]) => ({ + name, + count: span.count, + failureCount: span.failureCount, + totalDurationMs: span.totalDurationMs, + averageDurationMs: span.count > 0 ? span.totalDurationMs / span.count : 0, + maxDurationMs: span.maxDurationMs, + })) .toSorted( - (left, right) => DateTime.toEpochMillis(right.seenAt) - DateTime.toEpochMillis(left.seenAt), + (left, right) => right.count - left.count || right.maxDurationMs - left.maxDurationMs, ) - .slice(0, RECENT_LIMIT), - partialFailure: input.partialFailure ? Option.some(true) : Option.none(), - error: Option.fromNullishOr(input.error), + .slice(0, TOP_LIMIT); + + return { + traceFilePath: input.traceFilePath, + scannedFilePaths: input.scannedFilePaths, + readAt: input.readAt, + recordCount, + parseErrorCount, + firstSpanAt: Option.fromNullishOr(firstSpanAt), + lastSpanAt: Option.fromNullishOr(lastSpanAt), + failureCount, + interruptionCount, + slowSpanThresholdMs, + slowSpanCount, + logLevelCounts, + topSpansByCount, + slowestSpans, + commonFailures: [...failuresByKey.values()] + .toSorted( + (left, right) => + right.count - left.count || + DateTime.toEpochMillis(right.lastSeenAt) - DateTime.toEpochMillis(left.lastSeenAt), + ) + .slice(0, TOP_LIMIT), + latestFailures: latestFailures + .toSorted( + (left, right) => + DateTime.toEpochMillis(right.endedAt) - DateTime.toEpochMillis(left.endedAt), + ) + .slice(0, RECENT_LIMIT), + latestWarningAndErrorLogs: latestWarningAndErrorLogs + .toSorted( + (left, right) => + DateTime.toEpochMillis(right.seenAt) - DateTime.toEpochMillis(left.seenAt), + ) + .slice(0, RECENT_LIMIT), + partialFailure: input.partialFailure ? Option.some(true) : Option.none(), + error: Option.fromNullishOr(input.error), + }; }; -} -type TraceFileReadResult = - | { readonly _tag: "Loaded"; readonly path: string; readonly text: string } - | { readonly _tag: "Missing"; readonly path: string }; + return { addLine, finish }; +} -function readTraceFile( +/** + * Feeds each line of one trace file to `onLine`, streaming so only one chunk of + * text is in memory at a time. Succeeds with false when the file does not exist. + */ +export function streamTraceFileLines( fileSystem: FileSystem.FileSystem, path: string, -): Effect.Effect { - return fileSystem.readFileString(path).pipe( - Effect.map((text): TraceFileReadResult => ({ _tag: "Loaded", path, text })), + onLine: (line: string) => void, +): Effect.Effect { + return fileSystem.stream(path).pipe( + Stream.decodeText, + Stream.splitLines, + Stream.runForEachArray((lines) => Effect.sync(() => lines.forEach(onLine))), + Effect.as(true), Effect.catchTags({ PlatformError: (cause) => isNotFoundError(cause) - ? Effect.succeed({ _tag: "Missing", path }) + ? Effect.succeed(false) : Effect.fail( new TraceFileReadError({ traceFilePath: path, @@ -425,10 +417,12 @@ export const make = Effect.gen(function* () { const readAt = options.readAt ?? (yield* DateTime.now); const slowSpanThresholdMs = options.slowSpanThresholdMs ?? DEFAULT_SLOW_SPAN_THRESHOLD_MS; const paths = toRotatedTracePaths(options.traceFilePath, options.maxFiles); + const aggregator = makeTraceDiagnosticsAggregator(slowSpanThresholdMs); + // One aggregator reads every file, so keep them in order, oldest first. const results = yield* Effect.forEach( paths, (path) => - readTraceFile(fileSystem, path).pipe( + streamTraceFileLines(fileSystem, path, aggregator.addLine).pipe( Effect.tapError((cause) => Effect.logWarning("Failed to read local trace file.").pipe( Effect.annotateLogs({ @@ -444,11 +438,7 @@ export const make = Effect.gen(function* () { concurrency: 1, }, ); - const files = results.flatMap((result) => - Result.isSuccess(result) && result.success._tag === "Loaded" - ? [{ path: result.success.path, text: result.success.text }] - : [], - ); + const foundFile = results.some((result) => Result.isSuccess(result) && result.success); const readFailure = results.find(Result.isFailure); const readFailureError = readFailure ? ({ @@ -457,7 +447,7 @@ export const make = Effect.gen(function* () { } satisfies TraceDiagnosticsErrorSummary) : undefined; - if (files.length === 0) { + if (!foundFile) { return makeEmptyDiagnostics({ traceFilePath: options.traceFilePath, scannedFilePaths: paths, @@ -472,12 +462,10 @@ export const make = Effect.gen(function* () { }); } - return aggregateTraceDiagnostics({ + return aggregator.finish({ traceFilePath: options.traceFilePath, - files, scannedFilePaths: paths, readAt, - slowSpanThresholdMs, ...(readFailureError ? { partialFailure: true, error: readFailureError } : {}), }); }, From f26b76caa5845b042401b1cdcce4f6ede92e42b8 Mon Sep 17 00:00:00 2001 From: Theo Browne Date: Fri, 25 Sep 2026 22:24:18 -0700 Subject: [PATCH 2/2] perf(server): cap latest failures and warning logs while aggregating traces The aggregator kept every failure and every warning or error log, then sorted and sliced them to 20 at the end. On a ring with many warnings that held tens of thousands of objects. Now all three top lists use one bounded insert, so each holds at most its limit. Output is unchanged. Co-Authored-By: Claude Opus 5.5 (1M context) --- .../src/diagnostics/TraceDiagnostics.test.ts | 28 +++++-- .../src/diagnostics/TraceDiagnostics.ts | 79 ++++++++++--------- 2 files changed, 65 insertions(+), 42 deletions(-) diff --git a/apps/server/src/diagnostics/TraceDiagnostics.test.ts b/apps/server/src/diagnostics/TraceDiagnostics.test.ts index 8ea72b95b9e0..3b86841b99ba 100644 --- a/apps/server/src/diagnostics/TraceDiagnostics.test.ts +++ b/apps/server/src/diagnostics/TraceDiagnostics.test.ts @@ -310,25 +310,43 @@ describe("TraceDiagnostics", () => { }), ); - it.effect("keeps only the slowest span occurrences while aggregating large inputs", () => + it.effect("keeps only the top spans, failures, and warning logs from large inputs", () => Effect.sync(() => { + // Shuffled, so some older records arrive after the lists are full. + const indexes = Array.from({ length: 30 }, (_, step) => (step * 7) % 30); const diagnostics = aggregateLines( - Array.from({ length: 25 }, (_, index) => + indexes.map((index) => record({ name: `span-${index}`, traceId: `trace-${index}`, spanId: `span-${index}`, startMs: index * 1_000, durationMs: index, + exit: { _tag: "Failure", cause: "Provider crashed" }, + events: [ + { + name: `warning ${index}`, + timeUnixNano: ns(index * 1_000), + attributes: { "effect.logLevel": "Warning" }, + }, + ], }), ), ); + const newestTwenty = Array.from({ length: 20 }, (_, rank) => `trace-${29 - rank}`); - assert.equal(diagnostics.recordCount, 25); - assert.equal(diagnostics.slowestSpans.length, 10); + assert.equal(diagnostics.recordCount, 30); assert.deepStrictEqual( diagnostics.slowestSpans.map((span) => span.durationMs), - [24, 23, 22, 21, 20, 19, 18, 17, 16, 15], + [29, 28, 27, 26, 25, 24, 23, 22, 21, 20], + ); + assert.deepStrictEqual( + diagnostics.latestFailures.map((failure) => failure.traceId), + newestTwenty, + ); + assert.deepStrictEqual( + diagnostics.latestWarningAndErrorLogs.map((log) => log.traceId), + newestTwenty, ); }), ); diff --git a/apps/server/src/diagnostics/TraceDiagnostics.ts b/apps/server/src/diagnostics/TraceDiagnostics.ts index 79d62567ea6a..6dbc4afdec93 100644 --- a/apps/server/src/diagnostics/TraceDiagnostics.ts +++ b/apps/server/src/diagnostics/TraceDiagnostics.ts @@ -164,24 +164,43 @@ function isNotFoundError(error: PlatformError.PlatformError): boolean { return error.reason._tag === "NotFound"; } -function insertBoundedSlowestSpan( - slowestSpans: ServerTraceDiagnosticsSpanOccurrence[], - span: ServerTraceDiagnosticsSpanOccurrence, +/** + * Adds `item` to `items`, which stays sorted by `order` and holds at most + * `limit` entries. Same result as a stable sort and slice over every item, but + * memory stays bounded however many items stream in. + */ +function insertBounded( + items: A[], + item: A, + limit: number, + order: (left: A, right: A) => number, ): void { - if ( - slowestSpans.length >= TOP_LIMIT && - span.durationMs <= slowestSpans[slowestSpans.length - 1]!.durationMs - ) { + if (items.length >= limit && order(item, items[items.length - 1]!) >= 0) { return; } - slowestSpans.push(span); - slowestSpans.sort((left, right) => right.durationMs - left.durationMs); - if (slowestSpans.length > TOP_LIMIT) { - slowestSpans.length = TOP_LIMIT; + items.push(item); + items.sort(order); + if (items.length > limit) { + items.length = limit; } } +const slowestFirst = ( + left: ServerTraceDiagnosticsSpanOccurrence, + right: ServerTraceDiagnosticsSpanOccurrence, +) => right.durationMs - left.durationMs; + +const latestEndedFirst = ( + left: ServerTraceDiagnosticsRecentFailure, + right: ServerTraceDiagnosticsRecentFailure, +) => DateTime.toEpochMillis(right.endedAt) - DateTime.toEpochMillis(left.endedAt); + +const latestSeenFirst = ( + left: ServerTraceDiagnosticsLogEvent, + right: ServerTraceDiagnosticsLogEvent, +) => DateTime.toEpochMillis(right.seenAt) - DateTime.toEpochMillis(left.seenAt); + /** * Folds trace NDJSON into diagnostics. Call `addLine` once per line as the * rotated files stream in, then `finish` for the result. @@ -265,11 +284,11 @@ export function makeTraceDiagnosticsAggregator( if (durationMs >= slowSpanThresholdMs) { slowSpanCount += 1; } - insertBoundedSlowestSpan(slowestSpans, spanItem); + insertBounded(slowestSpans, spanItem, TOP_LIMIT, slowestFirst); if (isFailure) { const cause = readExitCause(parsed.exit); - latestFailures.push({ ...spanItem, cause }); + insertBounded(latestFailures, { ...spanItem, cause }, RECENT_LIMIT, latestEndedFirst); const failureKey = `${name}\0${cause}`; const existing = failuresByKey.get(failureKey); @@ -304,14 +323,12 @@ export function makeTraceDiagnosticsAggregator( const seenAt = unixNanoToDateTime(rawEvent.timeUnixNano) ?? endedAt; const message = toStringValue(rawEvent.name)?.trim() ?? "Log event"; - latestWarningAndErrorLogs.push({ - spanName: name, - level, - message, - seenAt, - traceId, - spanId, - }); + insertBounded( + latestWarningAndErrorLogs, + { spanName: name, level, message, seenAt, traceId, spanId }, + RECENT_LIMIT, + latestSeenFirst, + ); } } }; @@ -359,18 +376,8 @@ export function makeTraceDiagnosticsAggregator( DateTime.toEpochMillis(right.lastSeenAt) - DateTime.toEpochMillis(left.lastSeenAt), ) .slice(0, TOP_LIMIT), - latestFailures: latestFailures - .toSorted( - (left, right) => - DateTime.toEpochMillis(right.endedAt) - DateTime.toEpochMillis(left.endedAt), - ) - .slice(0, RECENT_LIMIT), - latestWarningAndErrorLogs: latestWarningAndErrorLogs - .toSorted( - (left, right) => - DateTime.toEpochMillis(right.seenAt) - DateTime.toEpochMillis(left.seenAt), - ) - .slice(0, RECENT_LIMIT), + latestFailures, + latestWarningAndErrorLogs, partialFailure: input.partialFailure ? Option.some(true) : Option.none(), error: Option.fromNullishOr(input.error), }; @@ -418,7 +425,6 @@ export const make = Effect.gen(function* () { const slowSpanThresholdMs = options.slowSpanThresholdMs ?? DEFAULT_SLOW_SPAN_THRESHOLD_MS; const paths = toRotatedTracePaths(options.traceFilePath, options.maxFiles); const aggregator = makeTraceDiagnosticsAggregator(slowSpanThresholdMs); - // One aggregator reads every file, so keep them in order, oldest first. const results = yield* Effect.forEach( paths, (path) => @@ -434,9 +440,8 @@ export const make = Effect.gen(function* () { ), Effect.result, ), - { - concurrency: 1, - }, + // Every file feeds one aggregator, so read them one at a time, oldest first. + { concurrency: 1 }, ); const foundFile = results.some((result) => Result.isSuccess(result) && result.success); const readFailure = results.find(Result.isFailure);