From 3b0e0b065f4c5aa4684cc4c07b587645293d7e26 Mon Sep 17 00:00:00 2001 From: badcuban <108198679+badcuban@users.noreply.github.com> Date: Mon, 24 Aug 2026 14:46:13 -0400 Subject: [PATCH] perf(observability): bound log flushing and retries --- .../provider/Layers/EventNdjsonLogger.test.ts | 268 ++++++++++++++- .../src/provider/Layers/EventNdjsonLogger.ts | 287 ++++++++++++++-- packages/shared/src/observability.test.ts | 315 +++++++++++++++++- packages/shared/src/observability.ts | 162 ++++++++- 4 files changed, 989 insertions(+), 43 deletions(-) diff --git a/apps/server/src/provider/Layers/EventNdjsonLogger.test.ts b/apps/server/src/provider/Layers/EventNdjsonLogger.test.ts index 0e4e55e8d..1df0a3ae5 100644 --- a/apps/server/src/provider/Layers/EventNdjsonLogger.test.ts +++ b/apps/server/src/provider/Layers/EventNdjsonLogger.test.ts @@ -6,8 +6,14 @@ import path from "node:path"; import { ThreadId } from "@threadlines/contracts"; import { assert, describe, it } from "@effect/vitest"; import * as Effect from "effect/Effect"; +import * as TestClock from "effect/testing/TestClock"; +import { vi } from "vitest"; -import { cleanupProviderEventLogDirectory, makeEventNdjsonLogger } from "./EventNdjsonLogger.ts"; +import { + cleanupProviderEventLogDirectory, + type EventNdjsonLogger, + makeEventNdjsonLogger, +} from "./EventNdjsonLogger.ts"; function parseLogLine(line: string) { const match = /^\[([^\]]+)\] ([A-Z]+): (.+)$/.exec(line); @@ -147,6 +153,264 @@ describe("EventNdjsonLogger", () => { }), ); + it.effect("flushes after the batch window and re-arms for later writes", () => + Effect.gen(function* () { + const tempDir = fs.mkdtempSync(path.join(os.tmpdir(), "t3-provider-log-")); + const basePath = path.join(tempDir, "provider-canonical.ndjson"); + const globalPath = path.join(tempDir, "_global.log"); + let logger: EventNdjsonLogger | undefined; + + try { + logger = yield* makeEventNdjsonLogger(basePath, { + stream: "canonical", + batchWindowMs: 1_000, + }); + assert.notEqual(logger, undefined); + if (!logger) { + return; + } + + yield* logger.write({ id: "evt-first" }, null); + assert.equal(fs.existsSync(globalPath), false); + yield* TestClock.adjust("999 millis"); + yield* Effect.yieldNow; + assert.equal(fs.existsSync(globalPath), false); + yield* TestClock.adjust("1 millis"); + yield* Effect.yieldNow; + assert.equal(fs.existsSync(globalPath), true); + + yield* logger.write({ id: "evt-second" }, null); + assert.equal(fs.readFileSync(globalPath, "utf8").includes("evt-second"), false); + yield* TestClock.adjust("1 second"); + yield* Effect.yieldNow; + + const lines = fs + .readFileSync(globalPath, "utf8") + .trim() + .split("\n") + .map((line) => parseLogLine(line)); + assert.deepEqual( + lines.map((line) => line.payload), + ['{"id":"evt-first"}', '{"id":"evt-second"}'], + ); + } finally { + if (logger) { + yield* logger.close(); + } + fs.rmSync(tempDir, { recursive: true, force: true }); + } + }), + ); + + it.effect("flushes immediately when buffered bytes reach the configured cap", () => + Effect.gen(function* () { + const tempDir = fs.mkdtempSync(path.join(os.tmpdir(), "t3-provider-log-")); + const basePath = path.join(tempDir, "provider-native.ndjson"); + const globalPath = path.join(tempDir, "_global.log"); + let logger: EventNdjsonLogger | undefined; + + try { + logger = yield* makeEventNdjsonLogger(basePath, { + stream: "native", + batchWindowMs: 10_000, + maxBufferedBytes: 350, + }); + assert.notEqual(logger, undefined); + if (!logger) { + return; + } + + yield* logger.write({ id: "evt-byte-threshold", payload: "\u{1F642}".repeat(100) }, null); + assert.equal(fs.existsSync(globalPath), true); + assert.equal(fs.readFileSync(globalPath, "utf8").includes("evt-byte-threshold"), true); + } finally { + if (logger) { + yield* logger.close(); + } + fs.rmSync(tempDir, { recursive: true, force: true }); + } + }), + ); + + it.effect("drops oversized records before the rotating sink can append them", () => + Effect.gen(function* () { + const tempDir = fs.mkdtempSync(path.join(os.tmpdir(), "t3-provider-log-")); + const basePath = path.join(tempDir, "provider-canonical.ndjson"); + const globalPath = path.join(tempDir, "_global.log"); + const blockingBackup = `${globalPath}.2`; + let logger: EventNdjsonLogger | undefined; + + try { + yield* TestClock.setTime(0); + logger = yield* makeEventNdjsonLogger(basePath, { + stream: "canonical", + maxBytes: 180, + maxFiles: 2, + batchWindowMs: 0, + }); + assert.notEqual(logger, undefined); + if (!logger) { + return; + } + + fs.mkdirSync(path.join(blockingBackup, "child"), { recursive: true }); + yield* logger.write({ id: "evt-oversized", payload: "x".repeat(300) }, null); + fs.rmSync(blockingBackup, { recursive: true, force: true }); + yield* TestClock.adjust("200 millis"); + yield* Effect.yieldNow; + yield* logger.write({ id: "evt-small" }, null); + + const payloadIds = fs + .readdirSync(tempDir) + .filter((entry) => entry === "_global.log" || entry.startsWith("_global.log.")) + .flatMap((entry) => + fs + .readFileSync(path.join(tempDir, entry), "utf8") + .trim() + .split("\n") + .map((line) => JSON.parse(parseLogLine(line).payload).id as string), + ); + assert.deepEqual(payloadIds, ["evt-small"]); + } finally { + if (logger) { + yield* logger.close(); + } + fs.rmSync(tempDir, { recursive: true, force: true }); + } + }), + ); + + it.effect("waits for the retry timer and does not duplicate completed chunks", () => + Effect.gen(function* () { + const tempDir = fs.mkdtempSync(path.join(os.tmpdir(), "t3-provider-log-")); + const basePath = path.join(tempDir, "provider-canonical.ndjson"); + const realAppendFileSync = fs.appendFileSync.bind(fs); + let appendAttempts = 0; + let failWritesAfterFirst = true; + let logger: EventNdjsonLogger | undefined; + const appendSpy = vi.spyOn(fs, "appendFileSync").mockImplementation((...args) => { + appendAttempts += 1; + if (failWritesAfterFirst && appendAttempts > 1) { + throw new Error("simulated provider log write failure"); + } + Reflect.apply(realAppendFileSync, fs, args); + }); + + try { + yield* TestClock.setTime(0); + logger = yield* makeEventNdjsonLogger(basePath, { + stream: "canonical", + maxBytes: 180, + maxFiles: 10, + batchWindowMs: 1_000, + maxBufferedBytes: 2_000, + }); + assert.notEqual(logger, undefined); + if (!logger) { + return; + } + + yield* logger.write({ id: "evt-1", payload: "x".repeat(60) }, null); + yield* logger.write({ id: "evt-2", payload: "x".repeat(60) }, null); + yield* logger.write({ id: "evt-3", payload: "x".repeat(60) }, null); + yield* TestClock.adjust("1 second"); + yield* Effect.yieldNow; + assert.equal(appendAttempts, 2); + + yield* logger.write({ id: "evt-4", payload: "x".repeat(60) }, null); + assert.equal(appendAttempts, 2); + yield* TestClock.adjust("999 millis"); + yield* Effect.yieldNow; + assert.equal(appendAttempts, 2); + + failWritesAfterFirst = false; + yield* TestClock.adjust("1 millis"); + yield* Effect.yieldNow; + yield* TestClock.adjust("1 second"); + yield* Effect.yieldNow; + const payloadIds = fs + .readdirSync(tempDir) + .filter((entry) => entry === "_global.log" || entry.startsWith("_global.log.")) + .flatMap((entry) => + fs + .readFileSync(path.join(tempDir, entry), "utf8") + .trim() + .split("\n") + .map((line) => JSON.parse(parseLogLine(line).payload).id as string), + ) + .toSorted(); + assert.deepEqual(payloadIds, ["evt-1", "evt-2", "evt-3", "evt-4"]); + } finally { + if (logger) { + yield* logger.close(); + } + appendSpy.mockRestore(); + fs.rmSync(tempDir, { recursive: true, force: true }); + } + }), + ); + + it.effect("bounds buffered records globally across thread writers", () => + Effect.gen(function* () { + const tempDir = fs.mkdtempSync(path.join(os.tmpdir(), "t3-provider-log-")); + const basePath = path.join(tempDir, "provider-canonical.ndjson"); + const realAppendFileSync = fs.appendFileSync.bind(fs); + let failWrites = true; + let logger: EventNdjsonLogger | undefined; + const appendSpy = vi.spyOn(fs, "appendFileSync").mockImplementation((...args) => { + if (failWrites) { + throw new Error("simulated provider log write failure"); + } + Reflect.apply(realAppendFileSync, fs, args); + }); + + try { + yield* TestClock.setTime(0); + logger = yield* makeEventNdjsonLogger(basePath, { + stream: "canonical", + batchWindowMs: 1_000, + maxBufferedBytes: 500, + }); + assert.notEqual(logger, undefined); + if (!logger) { + return; + } + + for (let index = 0; index < 20; index += 1) { + yield* logger.write( + { id: `evt-${index}`, payload: "x".repeat(40) }, + ThreadId.make(`thread-${index}`), + ); + } + yield* TestClock.adjust("1 second"); + yield* Effect.yieldNow; + + failWrites = false; + yield* TestClock.adjust("1 second"); + yield* Effect.yieldNow; + const logFiles = fs.readdirSync(tempDir).filter((entry) => entry.endsWith(".log")); + const persistedBytes = logFiles.reduce( + (total, entry) => total + fs.statSync(path.join(tempDir, entry)).size, + 0, + ); + const persistedRecords = logFiles.reduce( + (total, entry) => + total + fs.readFileSync(path.join(tempDir, entry), "utf8").trim().split("\n").length, + 0, + ); + assert.equal(persistedBytes <= 500, true); + assert.equal(persistedRecords > 0, true); + assert.equal(persistedRecords < 20, true); + } finally { + if (logger) { + yield* logger.close(); + } + appendSpy.mockRestore(); + fs.rmSync(tempDir, { recursive: true, force: true }); + } + }), + ); + it.effect( "falls back to a global segment when orchestration thread id is missing or invalid", () => @@ -236,7 +500,7 @@ describe("EventNdjsonLogger", () => { try { const logger = yield* makeEventNdjsonLogger(basePath, { stream: "native", - maxBytes: 120, + maxBytes: 200, maxFiles: 2, }); assert.notEqual(logger, undefined); diff --git a/apps/server/src/provider/Layers/EventNdjsonLogger.ts b/apps/server/src/provider/Layers/EventNdjsonLogger.ts index a48fc11ed..1d9b8b2d4 100644 --- a/apps/server/src/provider/Layers/EventNdjsonLogger.ts +++ b/apps/server/src/provider/Layers/EventNdjsonLogger.ts @@ -22,7 +22,9 @@ import { toSafeThreadAttachmentSegment } from "../../attachmentStore.ts"; const DEFAULT_MAX_BYTES = 10 * 1024 * 1024; const DEFAULT_MAX_FILES = 10; const DEFAULT_BATCH_WINDOW_MS = 200; +const MAX_RETRY_DELAY_MS = 30_000; const FLUSH_BUFFER_THRESHOLD = 32; +const DEFAULT_MAX_BUFFERED_BYTES = 1024 * 1024; const GLOBAL_THREAD_SEGMENT = "_global"; const LOG_SCOPE = "provider-observability"; const DEFAULT_LOG_CLEANUP_MAX_AGE_MS = 30 * 24 * 60 * 60 * 1_000; @@ -43,6 +45,7 @@ export interface EventNdjsonLoggerOptions { readonly maxBytes?: number; readonly maxFiles?: number; readonly batchWindowMs?: number; + readonly maxBufferedBytes?: number; } export interface ProviderEventLogCleanupOptions { @@ -65,6 +68,16 @@ interface ThreadWriter { close: () => Effect.Effect; } +interface BufferBudget { + readonly maxBytes: number; + readonly usedBytes: number; + reserve: ( + bytes: number, + allowTemporaryOverflow: boolean, + ) => { readonly reserved: boolean; readonly reportDrop: boolean }; + release: (bytes: number) => void; +} + interface LoggerState { readonly threadWriters: Map; readonly failedSegments: Set; @@ -243,11 +256,39 @@ const toLogMessage = Effect.fnUntraced(function* ( return serialized.value; }); +function makeBufferBudget(maxBytes: number): BufferBudget { + let usedBytes = 0; + let dropReported = false; + + return { + maxBytes, + get usedBytes() { + return usedBytes; + }, + reserve(bytes, allowTemporaryOverflow) { + if (bytes > maxBytes - usedBytes && !(allowTemporaryOverflow && usedBytes === 0)) { + const reportDrop = !dropReported; + dropReported = true; + return { reserved: false, reportDrop }; + } + usedBytes += bytes; + return { reserved: true, reportDrop: false }; + }, + release(bytes) { + usedBytes = Math.max(0, usedBytes - bytes); + if (usedBytes < maxBytes) { + dropReported = false; + } + }, + }; +} + const makeThreadWriter = Effect.fnUntraced(function* (input: { readonly filePath: string; readonly maxBytes: number; readonly maxFiles: number; readonly batchWindowMs: number; + readonly bufferBudget: BufferBudget; readonly streamLabel: string; }): Effect.fn.Return { const sinkResult = yield* Effect.sync(() => { @@ -276,68 +317,256 @@ const makeThreadWriter = Effect.fnUntraced(function* (input: { const sink = sinkResult.sink; const scope = yield* Scope.make(); + const initialRetryDelayMs = + input.batchWindowMs > 0 ? input.batchWindowMs : DEFAULT_BATCH_WINDOW_MS; let closed = false; - let buffer: string[] = []; - - const flushUnsafe = (): { ok: true } | { ok: false; error: unknown } => { + let flushScheduled = false; + let recovering = false; + let failureReported = false; + let oversizedRecordWarningReported = false; + let nextRetryDelayMs = initialRetryDelayMs; + let bufferedBytes = 0; + let buffer: Array<{ readonly line: string; readonly bytes: number }> = []; + + const flushUnsafe = (): + | { ok: true; attempted: boolean } + | { + ok: false; + error: unknown; + droppedRecords: number; + droppedBytes: number; + } => { if (buffer.length === 0) { - return { ok: true }; + return { ok: true, attempted: false }; } const messages = buffer; buffer = []; + bufferedBytes = 0; + let persistedCount = 0; try { - for (const message of messages) { - sink.write(message); + while (persistedCount < messages.length) { + const firstMessage = messages[persistedCount]; + if (!firstMessage) { + break; + } + let nextIndex = persistedCount + 1; + let chunkBytes = firstMessage.bytes; + while (nextIndex < messages.length) { + const nextMessage = messages[nextIndex]; + if (!nextMessage || chunkBytes + nextMessage.bytes > input.maxBytes) { + break; + } + chunkBytes += nextMessage.bytes; + nextIndex += 1; + } + sink.write( + messages + .slice(persistedCount, nextIndex) + .map((message) => message.line) + .join(""), + ); + input.bufferBudget.release(chunkBytes); + persistedCount = nextIndex; } - return { ok: true }; + return { ok: true, attempted: true }; } catch (error) { - buffer = [...messages, ...buffer]; - return { ok: false, error }; + buffer = [...messages.slice(persistedCount), ...buffer]; + bufferedBytes = buffer.reduce((total, message) => total + message.bytes, 0); + let droppedRecords = 0; + let droppedBytes = 0; + while (input.bufferBudget.usedBytes > input.bufferBudget.maxBytes && buffer.length > 0) { + const dropped = buffer.pop(); + if (!dropped) { + break; + } + bufferedBytes -= dropped.bytes; + input.bufferBudget.release(dropped.bytes); + droppedRecords += 1; + droppedBytes += dropped.bytes; + } + return { ok: false, error, droppedRecords, droppedBytes }; + } + }; + + const runFlushUnsafe = () => { + const result = flushUnsafe(); + if (result.ok) { + if (result.attempted) { + recovering = false; + failureReported = false; + nextRetryDelayMs = initialRetryDelayMs; + } + return { result, reportFailure: false }; } + + nextRetryDelayMs = recovering + ? Math.min(MAX_RETRY_DELAY_MS, nextRetryDelayMs * 2) + : initialRetryDelayMs; + recovering = true; + const reportFailure = !failureReported; + failureReported = true; + return { result, reportFailure }; }; - const reportFlushResult = (result: ReturnType) => - result.ok + const reportFlushResult = (outcome: ReturnType) => + outcome.result.ok || !outcome.reportFailure ? Effect.void : logWarning("provider event log batch flush failed", { filePath: input.filePath, - error: result.error, + error: outcome.result.error, + droppedRecords: outcome.result.droppedRecords, + droppedBytes: outcome.result.droppedBytes, }); - const flush = Effect.sync(flushUnsafe).pipe( + const flush = Effect.sync(runFlushUnsafe).pipe( Effect.flatMap(reportFlushResult), Effect.withTracerEnabled(false), ); - if (input.batchWindowMs > 0) { - yield* Effect.sleep(`${input.batchWindowMs} millis`).pipe( - Effect.andThen(flush), - Effect.forever, - Effect.forkIn(scope), - ); + function scheduleFlush(): Effect.Effect { + const delayMs = recovering ? nextRetryDelayMs : initialRetryDelayMs; + return Effect.forkIn( + Effect.sleep(`${delayMs} millis`).pipe( + Effect.andThen( + Effect.gen(function* () { + const outcome = yield* Effect.sync(() => { + flushScheduled = false; + return runFlushUnsafe(); + }); + yield* reportFlushResult(outcome); + + const retry = yield* Effect.sync(() => { + if (closed || buffer.length === 0 || flushScheduled) { + return false; + } + flushScheduled = true; + return true; + }); + if (retry) { + yield* scheduleFlush(); + } + }), + ), + ), + scope, + { startImmediately: true }, + ).pipe(Effect.asVoid); } const writeMessage = (message: string) => Effect.gen(function* () { const observedAt = DateTime.formatIso(yield* DateTime.now); - return yield* Effect.sync(() => { + const action = yield* Effect.sync(() => { if (closed) { - return { ok: true as const }; + return { + outcome: undefined, + schedule: false, + droppedBytes: 0, + reportDrop: false, + reportOversized: false, + }; + } + const line = formatLogLine(input.streamLabel, observedAt, message); + const bytes = Buffer.byteLength(line); + if (bytes > input.maxBytes) { + const reportOversized = !oversizedRecordWarningReported; + oversizedRecordWarningReported = true; + return { + outcome: undefined, + schedule: false, + droppedBytes: bytes, + reportDrop: false, + reportOversized, + }; + } + const reservation = input.bufferBudget.reserve(bytes, !recovering); + if (!reservation.reserved) { + return { + outcome: undefined, + schedule: false, + droppedBytes: bytes, + reportDrop: reservation.reportDrop, + reportOversized: false, + }; } - buffer.push(formatLogLine(input.streamLabel, observedAt, message)); - return input.batchWindowMs <= 0 || buffer.length >= FLUSH_BUFFER_THRESHOLD - ? flushUnsafe() - : { ok: true as const }; + buffer.push({ line, bytes }); + bufferedBytes += bytes; + + if ( + !recovering && + (input.batchWindowMs <= 0 || + buffer.length >= FLUSH_BUFFER_THRESHOLD || + bufferedBytes >= input.bufferBudget.maxBytes) + ) { + const outcome = runFlushUnsafe(); + const schedule = !outcome.result.ok && buffer.length > 0 && !flushScheduled; + if (schedule) { + flushScheduled = true; + } + return { + outcome, + schedule, + droppedBytes: 0, + reportDrop: false, + reportOversized: false, + }; + } + + if (!flushScheduled && (input.batchWindowMs > 0 || recovering)) { + flushScheduled = true; + return { + outcome: undefined, + schedule: true, + droppedBytes: 0, + reportDrop: false, + reportOversized: false, + }; + } + return { + outcome: undefined, + schedule: false, + droppedBytes: 0, + reportDrop: false, + reportOversized: false, + }; }); - }).pipe(Effect.flatMap(reportFlushResult), Effect.withTracerEnabled(false)); + if (action.reportOversized) { + yield* logWarning("provider event log record exceeds file limit; dropping record", { + filePath: input.filePath, + maxBytes: input.maxBytes, + droppedRecords: 1, + droppedBytes: action.droppedBytes, + }); + } + if (action.reportDrop) { + yield* logWarning("provider event log buffer limit reached; dropping record", { + filePath: input.filePath, + maxBufferedBytes: input.bufferBudget.maxBytes, + bufferedBytes: input.bufferBudget.usedBytes, + droppedRecords: 1, + droppedBytes: action.droppedBytes, + }); + } + if (action.outcome) { + yield* reportFlushResult(action.outcome); + } + if (action.schedule) { + yield* scheduleFlush(); + } + }).pipe(Effect.uninterruptible, Effect.withTracerEnabled(false)); const close = Effect.gen(function* () { closed = true; yield* Scope.close(scope, Exit.void); yield* flush; + yield* Effect.sync(() => { + input.bufferBudget.release(bufferedBytes); + bufferedBytes = 0; + buffer = []; + }); }).pipe(Effect.withTracerEnabled(false)); return { @@ -353,6 +582,10 @@ export const makeEventNdjsonLogger = Effect.fnUntraced(function* ( const maxBytes = options.maxBytes ?? DEFAULT_MAX_BYTES; const maxFiles = options.maxFiles ?? DEFAULT_MAX_FILES; const batchWindowMs = options.batchWindowMs ?? DEFAULT_BATCH_WINDOW_MS; + const requestedMaxBufferedBytes = options.maxBufferedBytes ?? DEFAULT_MAX_BUFFERED_BYTES; + const maxBufferedBytes = Number.isFinite(requestedMaxBufferedBytes) + ? Math.max(1, Math.trunc(requestedMaxBufferedBytes)) + : DEFAULT_MAX_BUFFERED_BYTES; const streamLabel = resolveStreamLabel(options.stream); const directoryReady = yield* Effect.sync(() => { @@ -375,6 +608,7 @@ export const makeEventNdjsonLogger = Effect.fnUntraced(function* ( threadWriters: new Map(), failedSegments: new Set(), }); + const bufferBudget = makeBufferBudget(maxBufferedBytes); const resolveThreadWriter = Effect.fnUntraced(function* ( threadSegment: string, @@ -394,6 +628,7 @@ export const makeEventNdjsonLogger = Effect.fnUntraced(function* ( maxBytes, maxFiles, batchWindowMs, + bufferBudget, streamLabel, }).pipe( Effect.map((writer) => { diff --git a/packages/shared/src/observability.test.ts b/packages/shared/src/observability.test.ts index fb77c3f41..aa761e24a 100644 --- a/packages/shared/src/observability.test.ts +++ b/packages/shared/src/observability.test.ts @@ -1,3 +1,6 @@ +// @effect-diagnostics nodeBuiltinImport:off +import fs from "node:fs"; + import { assert, describe, it } from "@effect/vitest"; import * as NodeServices from "@effect/platform-node/NodeServices"; import * as Effect from "effect/Effect"; @@ -7,7 +10,9 @@ import * as Logger from "effect/Logger"; import * as Path from "effect/Path"; import * as References from "effect/References"; import * as Schema from "effect/Schema"; +import * as TestClock from "effect/testing/TestClock"; import * as Tracer from "effect/Tracer"; +import { vi } from "vitest"; import { compactTraceAttributes, @@ -149,6 +154,45 @@ describe("observability", () => { ), ); + it.effect("flushes traces after the batch window and re-arms for later records", () => + Effect.scoped( + 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"); + + const sink = yield* makeTraceSink({ + filePath: tracePath, + maxBytes: 1024, + maxFiles: 2, + batchWindowMs: 1_000, + }); + + sink.push(makeRecord("first")); + assert.equal(fs.existsSync(tracePath), false); + yield* TestClock.adjust("999 millis"); + yield* Effect.yieldNow; + assert.equal(fs.existsSync(tracePath), false); + yield* TestClock.adjust("1 millis"); + yield* Effect.yieldNow; + assert.equal(fs.existsSync(tracePath), true); + + sink.push(makeRecord("second")); + assert.equal(fs.readFileSync(tracePath, "utf8").includes('"name":"second"'), false); + yield* TestClock.adjust("1 second"); + yield* Effect.yieldNow; + + const records = yield* readTraceRecords(tracePath); + assert.deepStrictEqual( + records.map((record) => record.name), + ["first", "second"], + ); + yield* sink.close(); + }), + ), + ); + it.effect("applies local tracer record filters after a span ends", () => Effect.scoped( Effect.gen(function* () { @@ -193,9 +237,11 @@ describe("observability", () => { const tempDir = yield* fileSystem.makeTempDirectoryScoped({ prefix: "t3-trace-sink-" }); const tracePath = path.join(tempDir, "shared.trace.ndjson"); + const record = makeRecord("rotate", `0-${"x".repeat(48)}`); + const recordBytes = Buffer.byteLength(`${JSON.stringify(record)}\n`); const sink = yield* makeTraceSink({ filePath: tracePath, - maxBytes: 180, + maxBytes: recordBytes, maxFiles: 2, batchWindowMs: 10_000, }); @@ -221,6 +267,273 @@ describe("observability", () => { matchingFiles.some((entry) => entry === "shared.trace.ndjson.3"), false, ); + for (const entry of matchingFiles) { + assert.equal(fs.statSync(path.join(tempDir, entry)).size <= recordBytes, true); + } + }), + ), + ); + + it.effect("chunks buffered trace records without exceeding the rotation size", () => + Effect.scoped( + 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"); + const first = makeRecord("first", "\u{1F642}".repeat(32)); + const second = makeRecord("second", "\u{1F642}".repeat(32)); + const maxBytes = Math.max( + Buffer.byteLength(`${JSON.stringify(first)}\n`), + Buffer.byteLength(`${JSON.stringify(second)}\n`), + ); + + const sink = yield* makeTraceSink({ + filePath: tracePath, + maxBytes, + maxFiles: 2, + batchWindowMs: 10_000, + }); + + sink.push(first); + sink.push(second); + yield* sink.close(); + + const files = ["shared.trace.ndjson.1", "shared.trace.ndjson"].filter((entry) => + fs.existsSync(path.join(tempDir, entry)), + ); + assert.deepStrictEqual(files, ["shared.trace.ndjson.1", "shared.trace.ndjson"]); + for (const entry of files) { + assert.equal(fs.statSync(path.join(tempDir, entry)).size <= maxBytes, true); + } + const names = files.flatMap((entry) => + fs + .readFileSync(path.join(tempDir, entry), "utf8") + .trim() + .split("\n") + .map((line) => decodeTraceRecordLine(line).name), + ); + assert.deepStrictEqual(names, ["first", "second"]); + }), + ), + ); + + it.effect("drops an oversized trace record without dropping later records", () => + Effect.scoped( + 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"); + const valid = makeRecord("valid"); + const maxBytes = Buffer.byteLength(`${JSON.stringify(valid)}\n`); + + const sink = yield* makeTraceSink({ + filePath: tracePath, + maxBytes, + maxFiles: 2, + batchWindowMs: 10_000, + }); + + sink.push(makeRecord("oversized", "x".repeat(maxBytes))); + sink.push(valid); + yield* sink.close(); + + const lines = yield* readTraceRecords(tracePath); + assert.deepStrictEqual( + lines.map((line) => line.name), + ["valid"], + ); + assert.equal(fs.statSync(tracePath).size <= maxBytes, true); + }), + ), + ); + + it.effect("bounds failed trace retries without duplicating earlier records", () => + Effect.scoped( + 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"); + const first = makeRecord("first"); + const second = makeRecord("second"); + const third = makeRecord("third"); + const fourth = makeRecord("fourth"); + const recordBytes = (record: TraceRecord) => + Buffer.byteLength(`${JSON.stringify(record)}\n`); + const maxBytes = Math.max( + recordBytes(first), + recordBytes(second), + recordBytes(third), + recordBytes(fourth), + ); + const maxBufferedBytes = Math.max( + recordBytes(second) + recordBytes(third), + recordBytes(third) + recordBytes(fourth), + ); + + const sink = yield* makeTraceSink({ + filePath: tracePath, + maxBytes, + maxFiles: 2, + batchWindowMs: 10_000, + maxBufferedBytes, + }); + + const blockingBackup = `${tracePath}.2`; + fs.mkdirSync(path.join(blockingBackup, "child"), { recursive: true }); + sink.push(first); + sink.push(second); + yield* sink.flush; + sink.push(third); + sink.push(fourth); + fs.rmSync(blockingBackup, { recursive: true, force: true }); + yield* sink.flush; + + const names = [ + "shared.trace.ndjson.2", + "shared.trace.ndjson.1", + "shared.trace.ndjson", + ].flatMap((entry) => + fs + .readFileSync(path.join(tempDir, entry), "utf8") + .trim() + .split("\n") + .map((line) => decodeTraceRecordLine(line).name), + ); + assert.deepStrictEqual(names, ["first", "third", "fourth"]); + }), + ), + ); + + it.effect("leaves failed trace recovery to the retry timer", () => + Effect.scoped( + 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"); + const first = makeRecord("first"); + const second = makeRecord("second"); + const third = makeRecord("third", "x"); + const recordBytes = (record: TraceRecord) => + Buffer.byteLength(`${JSON.stringify(record)}\n`); + const maxBytes = Math.max(recordBytes(first), recordBytes(second), recordBytes(third)); + + const sink = yield* makeTraceSink({ + filePath: tracePath, + maxBytes, + maxFiles: 2, + batchWindowMs: 1_000, + maxBufferedBytes: recordBytes(second) + recordBytes(third), + }); + + const blockingBackup = `${tracePath}.2`; + fs.mkdirSync(path.join(blockingBackup, "child"), { recursive: true }); + sink.push(first); + sink.push(second); + yield* sink.flush; + + fs.rmSync(blockingBackup, { recursive: true, force: true }); + sink.push(third); + assert.equal(fs.existsSync(`${tracePath}.1`), false); + assert.deepStrictEqual( + (yield* readTraceRecords(tracePath)).map((record) => record.name), + ["first"], + ); + + yield* TestClock.adjust("1 second"); + yield* Effect.yieldNow; + + const names = ["shared.trace.ndjson.2", "shared.trace.ndjson.1", "shared.trace.ndjson"] + .filter((entry) => fs.existsSync(path.join(tempDir, entry))) + .flatMap((entry) => + fs + .readFileSync(path.join(tempDir, entry), "utf8") + .trim() + .split("\n") + .map((line) => decodeTraceRecordLine(line).name), + ); + assert.deepStrictEqual(names, ["first", "second", "third"]); + }), + ), + ); + + it.effect("backs off persistent trace failures with an immediate batch window", () => + Effect.scoped( + 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"); + const realAppendFileSync = fs.appendFileSync.bind(fs); + let appendAttempts = 0; + let failWrites = true; + const appendSpy = vi.spyOn(fs, "appendFileSync").mockImplementation((...args) => { + appendAttempts += 1; + if (failWrites) { + throw new Error("simulated trace write failure"); + } + Reflect.apply(realAppendFileSync, fs, args); + }); + yield* Effect.addFinalizer(() => Effect.sync(() => appendSpy.mockRestore())); + yield* TestClock.setTime(0); + + const sink = yield* makeTraceSink({ + filePath: tracePath, + maxBytes: 4_096, + maxFiles: 2, + batchWindowMs: 0, + maxBufferedBytes: 4_096, + }); + + sink.push(makeRecord("first")); + assert.equal(appendAttempts, 1); + sink.push(makeRecord("second")); + assert.equal(appendAttempts, 1); + + yield* TestClock.adjust("199 millis"); + yield* Effect.yieldNow; + assert.equal(appendAttempts, 1); + yield* TestClock.adjust("1 millis"); + yield* Effect.yieldNow; + assert.equal(appendAttempts, 2); + + sink.push(makeRecord("third")); + assert.equal(appendAttempts, 2); + yield* TestClock.adjust("399 millis"); + yield* Effect.yieldNow; + assert.equal(appendAttempts, 2); + yield* TestClock.adjust("1 millis"); + yield* Effect.yieldNow; + assert.equal(appendAttempts, 3); + + let expectedAttempts = 3; + for (const delayMs of [800, 1_600, 3_200, 6_400, 12_800, 25_600, 30_000]) { + yield* TestClock.adjust(`${delayMs} millis`); + yield* Effect.yieldNow; + expectedAttempts += 1; + assert.equal(appendAttempts, expectedAttempts); + } + assert.equal(appendAttempts, 10); + yield* TestClock.adjust("29999 millis"); + yield* Effect.yieldNow; + assert.equal(appendAttempts, 10); + yield* TestClock.adjust("1 millis"); + yield* Effect.yieldNow; + assert.equal(appendAttempts, 11); + yield* TestClock.adjust("30 seconds"); + yield* Effect.yieldNow; + assert.equal(appendAttempts, 12); + + failWrites = false; + yield* sink.close(); + const records = yield* readTraceRecords(tracePath); + assert.deepStrictEqual( + records.map((record) => record.name), + ["first", "second", "third"], + ); }), ), ); diff --git a/packages/shared/src/observability.ts b/packages/shared/src/observability.ts index a78bd8e39..63d56387b 100644 --- a/packages/shared/src/observability.ts +++ b/packages/shared/src/observability.ts @@ -2,6 +2,7 @@ import * as Cause from "effect/Cause"; import * as Effect from "effect/Effect"; import type * as Exit from "effect/Exit"; import * as ExitRuntime from "effect/Exit"; +import * as Fiber from "effect/Fiber"; import * as Option from "effect/Option"; import * as Tracer from "effect/Tracer"; import { OtlpResource, OtlpTracer } from "effect/unstable/observability"; @@ -9,6 +10,10 @@ import { OtlpResource, OtlpTracer } from "effect/unstable/observability"; import { RotatingFileSink } from "./logging.ts"; const FLUSH_BUFFER_THRESHOLD = 32; +const DEFAULT_MAX_BUFFERED_BYTES = 1024 * 1024; +const DEFAULT_RETRY_DELAY_MS = 200; +const MAX_RETRY_DELAY_MS = 30_000; +const traceTextEncoder = new TextEncoder(); export type TraceAttributes = Readonly>; @@ -78,6 +83,7 @@ export interface TraceSinkOptions { readonly maxBytes: number; readonly maxFiles: number; readonly batchWindowMs: number; + readonly maxBufferedBytes?: number; } export interface TraceSink { @@ -241,50 +247,175 @@ export function spanToTraceRecord(span: SerializableSpan): EffectTraceRecord { } export const makeTraceSink = Effect.fn("makeTraceSink")(function* (options: TraceSinkOptions) { + const batchWindowMs = Number.isFinite(options.batchWindowMs) + ? Math.max(0, options.batchWindowMs) + : 0; + const initialRetryDelayMs = + batchWindowMs > 0 + ? Math.max(1, Math.min(MAX_RETRY_DELAY_MS, batchWindowMs)) + : DEFAULT_RETRY_DELAY_MS; + const requestedMaxBufferedBytes = options.maxBufferedBytes ?? DEFAULT_MAX_BUFFERED_BYTES; + const maxBufferedBytes = Number.isFinite(requestedMaxBufferedBytes) + ? Math.max(1, Math.trunc(requestedMaxBufferedBytes)) + : DEFAULT_MAX_BUFFERED_BYTES; const sink = new RotatingFileSink({ filePath: options.filePath, maxBytes: options.maxBytes, maxFiles: options.maxFiles, + throwOnError: true, }); - let buffer: Array = []; + const runFork = Effect.runForkWith(yield* Effect.context()); + let closed = false; + let flushFiber: Fiber.Fiber | undefined; + let flushFailed = false; + let nextRetryDelayMs = initialRetryDelayMs; + let bufferedBytes = 0; + let buffer: Array<{ readonly line: string; readonly bytes: number }> = []; + + const trimRetryBuffer = (): void => { + while (bufferedBytes > maxBufferedBytes && buffer.length > 0) { + const dropped = buffer.shift(); + if (!dropped) { + break; + } + bufferedBytes -= dropped.bytes; + } + }; const flushUnsafe = () => { if (buffer.length === 0) { return; } - const chunk = buffer.join(""); + const records = buffer; buffer = []; + bufferedBytes = 0; + let persistedCount = 0; + + while (persistedCount < records.length) { + const firstRecord = records[persistedCount]; + if (!firstRecord) { + break; + } + if (firstRecord.bytes > options.maxBytes) { + persistedCount += 1; + continue; + } + + let nextIndex = persistedCount + 1; + let chunkBytes = firstRecord.bytes; + while (nextIndex < records.length) { + const nextRecord = records[nextIndex]; + if (!nextRecord || chunkBytes + nextRecord.bytes > options.maxBytes) { + break; + } + chunkBytes += nextRecord.bytes; + nextIndex += 1; + } - try { - sink.write(chunk); - } catch { - buffer.unshift(chunk); + const chunk = records + .slice(persistedCount, nextIndex) + .map((record) => record.line) + .join(""); + try { + sink.write(chunk); + } catch { + buffer = [...records.slice(persistedCount), ...buffer]; + bufferedBytes = buffer.reduce((total, record) => total + record.bytes, 0); + trimRetryBuffer(); + nextRetryDelayMs = flushFailed + ? Math.min(MAX_RETRY_DELAY_MS, nextRetryDelayMs * 2) + : initialRetryDelayMs; + flushFailed = buffer.length > 0; + return; + } + persistedCount = nextIndex; } + flushFailed = false; + nextRetryDelayMs = initialRetryDelayMs; }; - const flush = Effect.sync(flushUnsafe).pipe(Effect.withTracerEnabled(false)); + const scheduleFlush = (): void => { + if (closed || flushFiber !== undefined || buffer.length === 0) { + return; + } - yield* Effect.addFinalizer(() => flush.pipe(Effect.ignore)); - yield* Effect.forkScoped( - Effect.sleep(`${options.batchWindowMs} millis`).pipe(Effect.andThen(flush), Effect.forever), - ); + const delayMs = flushFailed ? nextRetryDelayMs : batchWindowMs; + if (delayMs <= 0) { + return; + } + + flushFiber = runFork( + Effect.sleep(`${delayMs} millis`).pipe( + Effect.andThen( + Effect.sync(() => { + flushFiber = undefined; + flushUnsafe(); + scheduleFlush(); + }), + ), + Effect.withTracerEnabled(false), + ), + ); + }; + + const flush = Effect.gen(function* () { + const scheduled = flushFiber; + flushFiber = undefined; + if (scheduled) { + yield* Fiber.interrupt(scheduled).pipe(Effect.ignore); + } + yield* Effect.sync(() => { + flushUnsafe(); + scheduleFlush(); + }); + }).pipe(Effect.withTracerEnabled(false)); + + const close = Effect.gen(function* () { + closed = true; + const scheduled = flushFiber; + flushFiber = undefined; + if (scheduled) { + yield* Fiber.interrupt(scheduled).pipe(Effect.ignore); + } + yield* Effect.sync(flushUnsafe); + }).pipe(Effect.withTracerEnabled(false)); + + yield* Effect.addFinalizer(() => close.pipe(Effect.ignore)); return { filePath: options.filePath, push(record) { + if (closed) { + return; + } try { - buffer.push(`${JSON.stringify(record)}\n`); - if (buffer.length >= FLUSH_BUFFER_THRESHOLD) { + const line = `${JSON.stringify(record)}\n`; + const bytes = traceTextEncoder.encode(line).byteLength; + if (bytes > options.maxBytes) { + return; + } + buffer.push({ line, bytes }); + bufferedBytes += bytes; + if (flushFailed) { + trimRetryBuffer(); + } + if ( + !flushFailed && + (batchWindowMs <= 0 || + buffer.length >= FLUSH_BUFFER_THRESHOLD || + bufferedBytes >= maxBufferedBytes) + ) { flushUnsafe(); } + scheduleFlush(); } catch { return; } }, flush, - close: () => flush, + close: () => close, } satisfies TraceSink; }); @@ -375,6 +506,9 @@ export const makeLocalFileTracer = Effect.fn("makeLocalFileTracer")(function* ( maxBytes: options.maxBytes, maxFiles: options.maxFiles, batchWindowMs: options.batchWindowMs, + ...(options.maxBufferedBytes === undefined + ? {} + : { maxBufferedBytes: options.maxBufferedBytes }), })); const delegate =