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..3b86841b99ba 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( @@ -251,32 +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(() => { - 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"), - }, - ], - }); + // Shuffled, so some older records arrive after the lists are full. + const indexes = Array.from({ length: 30 }, (_, step) => (step * 7) % 30); + const diagnostics = aggregateLines( + 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 afb7b4e923bb..6dbc4afdec93 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; @@ -173,44 +164,50 @@ 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; } } -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 } : {}), - }); - } - +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. + */ +export function makeTraceDiagnosticsAggregator( + slowSpanThresholdMs = DEFAULT_SLOW_SPAN_THRESHOLD_MS, +) { let parseErrorCount = 0; let recordCount = 0; let failureCount = 0; @@ -229,182 +226,184 @@ 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; + } - 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, - }); + 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; + } + insertBounded(slowestSpans, spanItem, TOP_LIMIT, slowestFirst); + + if (isFailure) { + const cause = readExitCause(parsed.exit); + insertBounded(latestFailures, { ...spanItem, cause }, RECENT_LIMIT, latestEndedFirst); + + 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"; + insertBounded( + latestWarningAndErrorLogs, + { spanName: name, level, message, seenAt, traceId, spanId }, + RECENT_LIMIT, + latestSeenFirst, + ); } } - } - - 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 + 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.endedAt) - DateTime.toEpochMillis(left.endedAt), + (left, right) => right.count - left.count || right.maxDurationMs - left.maxDurationMs, ) - .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), + .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, + latestWarningAndErrorLogs, + 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 +424,11 @@ 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); 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({ @@ -440,15 +440,10 @@ export const make = Effect.gen(function* () { ), Effect.result, ), - { - concurrency: 1, - }, - ); - const files = results.flatMap((result) => - Result.isSuccess(result) && result.success._tag === "Loaded" - ? [{ path: result.success.path, text: result.success.text }] - : [], + // 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); const readFailureError = readFailure ? ({ @@ -457,7 +452,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 +467,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 } : {}), }); },