diff --git a/extensions/processes/hooks/event-bridge.test.ts b/extensions/processes/hooks/event-bridge.test.ts index 55873e1..3608163 100644 --- a/extensions/processes/hooks/event-bridge.test.ts +++ b/extensions/processes/hooks/event-bridge.test.ts @@ -77,10 +77,19 @@ describe("registerEventBridge", () => { const info = makeInfo({ status: "exited", success: true, exitCode: 0 }); const appendedText = [{ type: "stdout" as const, text: "ready" }]; fake.emit({ type: "process_ended", info }); - fake.emit({ type: "process_output_changed", id: "proc_1", appendedText }); + fake.emit({ + type: "process_output_changed", + id: "proc_1", + appendedText, + droppedLines: 3, + }); expect(ended).toHaveBeenCalledWith(info); - expect(output).toHaveBeenCalledWith({ id: "proc_1", appendedText }); + expect(output).toHaveBeenCalledWith({ + id: "proc_1", + appendedText, + droppedLines: 3, + }); }); it("bridges processes changed events", () => { diff --git a/extensions/processes/hooks/event-bridge.ts b/extensions/processes/hooks/event-bridge.ts index cd0db27..4a8a4e6 100644 --- a/extensions/processes/hooks/event-bridge.ts +++ b/extensions/processes/hooks/event-bridge.ts @@ -21,6 +21,7 @@ export function registerEventBridge( events.emit(CHANNELS.OUTPUT_CHANGED, { id: event.id, appendedText: event.appendedText, + droppedLines: event.droppedLines, }); break; case "processes_changed": diff --git a/src/manager/index.ts b/src/manager/index.ts index de63ccc..32c6249 100644 --- a/src/manager/index.ts +++ b/src/manager/index.ts @@ -98,6 +98,10 @@ export class ProcessManager { return this.logStore.getCombinedOutput(managed.combinedFile, tailLines); } + /** + * Return output capped at 16 MiB per stream. Use getLogFiles() and stream + * the returned paths when an exact, unbounded read is required. + */ getFullOutput(id: string): { stdout: string; stderr: string } | null { const managed = this.registry.getRecord(id); if (!managed) return null; diff --git a/src/manager/internal-types.ts b/src/manager/internal-types.ts index a653216..faa1b65 100644 --- a/src/manager/internal-types.ts +++ b/src/manager/internal-types.ts @@ -1,4 +1,3 @@ -import type { ChildProcess } from "node:child_process"; import type { Writable } from "node:stream"; import type { @@ -46,7 +45,6 @@ interface ProcessPublicState { * to expose because they allow direct mutation/control outside manager methods. */ interface ProcessRuntimeState { - process: ChildProcess; stdin: Writable | null; stdinClosed: boolean; lastSignalSent: NodeJS.Signals | null; @@ -70,8 +68,12 @@ interface ProcessLogState { * buffers let `ProcessOutput` emit only completed lines in events/logs. */ interface ProcessLineBufferState { - stdoutPendingLine: string; - stderrPendingLine: string; + stdoutPendingLine: Buffer; + stderrPendingLine: Buffer; + /** The line head was emitted; discard the tail through the next newline. */ + stdoutLineOverflowed: boolean; + /** The line head was emitted; discard the tail through the next newline. */ + stderrLineOverflowed: boolean; } /** @@ -82,6 +84,7 @@ interface ProcessLineBufferState { */ interface ProcessOutputEventBufferState { appendedLines: Array<{ type: "stdout" | "stderr"; text: string }>; + droppedLineCount: number; } /** diff --git a/src/manager/limits.ts b/src/manager/limits.ts new file mode 100644 index 0000000..48c8c69 --- /dev/null +++ b/src/manager/limits.ts @@ -0,0 +1,9 @@ +/** Max bytes retained for an unterminated line before force-flush. */ +export const MAX_PENDING_LINE_BYTES = 64 * 1024; +/** Max bytes retained for any single line in an event payload. */ +export const MAX_LINE_BYTES = 32 * 1024; +/** Max lines carried in one process_output_changed payload. */ +export const MAX_LINES_PER_EMIT = 2000; +/** Max bytes read when tailing a log file. */ +export const MAX_TAIL_READ_BYTES = 2 * 1024 * 1024; +export const TRUNCATION_SUFFIX = " … [line truncated]"; diff --git a/src/manager/process-log-store.test.ts b/src/manager/process-log-store.test.ts index f12da6e..984160e 100644 --- a/src/manager/process-log-store.test.ts +++ b/src/manager/process-log-store.test.ts @@ -1,6 +1,7 @@ import { fs, vol } from "memfs"; import { assert, beforeEach, describe, expect, it } from "vitest"; +import { MAX_TAIL_READ_BYTES } from "./limits"; import { ProcessLogStore } from "./process-log-store"; describe("ProcessLogStore", () => { @@ -99,6 +100,105 @@ describe("ProcessLogStore", () => { expect(store.readTailLines("/nonexistent", 10)).toEqual([]); }); + it("readTailLines handles files without a trailing newline", () => { + using store = new ProcessLogStore("/tmp/test-logs"); + const paths = store.createLogs("proc_1"); + + fs.writeFileSync(paths.stdoutFile, "line1\nline2\nline3"); + + expect(store.readTailLines(paths.stdoutFile, 2)).toEqual([ + "line2", + "line3", + ]); + }); + + it("readTailLines handles empty and zero-line requests", () => { + using store = new ProcessLogStore("/tmp/test-logs"); + const paths = store.createLogs("proc_1"); + + expect(store.readTailLines(paths.stdoutFile, 10)).toEqual([]); + fs.writeFileSync(paths.stdoutFile, "line\n"); + expect(store.readTailLines(paths.stdoutFile, 0)).toEqual([]); + }); + + it("readTailLines handles an exact chunk boundary", () => { + using store = new ProcessLogStore("/tmp/test-logs"); + const paths = store.createLogs("proc_1"); + const prefix = "x".repeat(64 * 1024 - "\nlast\n".length); + fs.writeFileSync(paths.stdoutFile, `${prefix}\nlast\n`); + + expect(store.readTailLines(paths.stdoutFile, 1)).toEqual(["last"]); + }); + + it("readTailLines reads the tail of a multi-megabyte numbered file", () => { + using store = new ProcessLogStore("/tmp/test-logs"); + const paths = store.createLogs("proc_1"); + const lines = Array.from( + { length: 500_000 }, + (_, index) => `line-${String(index).padStart(6, "0")}`, + ); + fs.writeFileSync(paths.stdoutFile, `${lines.join("\n")}\n`); + + expect(store.readTailLines(paths.stdoutFile, 10)).toEqual(lines.slice(-10)); + }); + + it("readTailLines bounds a tail containing one huge line", () => { + using store = new ProcessLogStore("/tmp/test-logs"); + const paths = store.createLogs("proc_1"); + fs.writeFileSync(paths.stdoutFile, "x".repeat(3 * 1024 * 1024)); + + const result = store.readTailLines(paths.stdoutFile, 10); + + expect(result).toHaveLength(1); + expect(result[0]?.startsWith("[… truncated] ")).toBe(true); + expect(Buffer.byteLength(result[0] ?? "")).toBeLessThanOrEqual( + MAX_TAIL_READ_BYTES + Buffer.byteLength("[… truncated] "), + ); + }); + + it("readTailLines retains a marked huge line ending in newline", () => { + using store = new ProcessLogStore("/tmp/test-logs"); + const paths = store.createLogs("proc_1"); + fs.writeFileSync(paths.stdoutFile, `${"x".repeat(3 * 1024 * 1024)}\n`); + + const result = store.readTailLines(paths.stdoutFile, 1); + + expect(result).toHaveLength(1); + expect(result[0]?.startsWith("[… truncated] ")).toBe(true); + }); + + it("readTailLines bounds decoded invalid UTF-8 output", () => { + using store = new ProcessLogStore("/tmp/test-logs"); + const paths = store.createLogs("proc_1"); + fs.writeFileSync( + paths.stdoutFile, + Buffer.alloc(MAX_TAIL_READ_BYTES + 1, 0xff), + ); + + const result = store.readTailLines(paths.stdoutFile, 1); + + expect(Buffer.byteLength(result[0] ?? "")).toBeLessThanOrEqual( + MAX_TAIL_READ_BYTES + Buffer.byteLength("[… truncated] "), + ); + }); + + it("readTailLines keeps final lines after invalid UTF-8 output", () => { + using store = new ProcessLogStore("/tmp/test-logs"); + const paths = store.createLogs("proc_1"); + fs.writeFileSync( + paths.stdoutFile, + Buffer.concat([ + Buffer.alloc(MAX_TAIL_READ_BYTES - 32, 0xff), + Buffer.from("\npenultimate\nlast\n"), + ]), + ); + + expect(store.readTailLines(paths.stdoutFile, 2)).toEqual([ + "penultimate", + "last", + ]); + }); + it("readFullFile returns entire content", () => { using store = new ProcessLogStore("/tmp/test-logs"); const paths = store.createLogs("proc_1"); @@ -114,6 +214,21 @@ describe("ProcessLogStore", () => { expect(store.readFullFile("/nonexistent")).toBe(""); }); + it("readFullFile returns a marked bounded suffix for large files", () => { + using store = new ProcessLogStore("/tmp/test-logs"); + const paths = store.createLogs("proc_1"); + fs.writeFileSync(paths.stdoutFile, "x".repeat(MAX_TAIL_READ_BYTES * 8 + 1)); + + const content = store.readFullFile(paths.stdoutFile); + + expect(content.startsWith(`[… truncated, see ${paths.stdoutFile}]\n`)).toBe( + true, + ); + expect(Buffer.byteLength(content)).toBeLessThanOrEqual( + MAX_TAIL_READ_BYTES * 8, + ); + }); + it("getCombinedOutput parses tagged lines", () => { using store = new ProcessLogStore("/tmp/test-logs"); const paths = store.createLogs("proc_1"); diff --git a/src/manager/process-log-store.ts b/src/manager/process-log-store.ts index 42a8676..ed09400 100644 --- a/src/manager/process-log-store.ts +++ b/src/manager/process-log-store.ts @@ -1,8 +1,12 @@ import { appendFileSync, + closeSync, + fstatSync, mkdirSync, mkdtempSync, + openSync, readFileSync, + readSync, rmSync, statSync, } from "node:fs"; @@ -10,6 +14,11 @@ import { tmpdir } from "node:os"; import { join } from "node:path"; import type { ProcessLogPaths } from "./internal-types"; +import { MAX_TAIL_READ_BYTES } from "./limits"; + +const TAIL_CHUNK_BYTES = 64 * 1024; +const MAX_FULL_FILE_BYTES = MAX_TAIL_READ_BYTES * 8; +const TAIL_TRUNCATION_PREFIX = "[… truncated] "; export class ProcessLogStore { private logDir: string; @@ -80,20 +89,99 @@ export class ProcessLogStore { } readTailLines(filePath: string, lines: number): string[] { + if (lines <= 0) return []; + + let fd: number | undefined; try { - const content = readFileSync(filePath, "utf-8"); - const allLines = content.split("\n"); - if (allLines.length > 0 && allLines[allLines.length - 1] === "") { - allLines.pop(); + fd = openSync(filePath, "r"); + const size = fstatSync(fd).size; + if (size === 0) return []; + + const buffer = Buffer.allocUnsafe(Math.min(size, MAX_TAIL_READ_BYTES)); + let bufferStart = buffer.length; + let position = size; + let bytesRead = 0; + let newlineCount = 0; + + while ( + position > 0 && + bytesRead < MAX_TAIL_READ_BYTES && + newlineCount <= lines + ) { + const length = Math.min( + TAIL_CHUNK_BYTES, + position, + MAX_TAIL_READ_BYTES - bytesRead, + ); + position -= length; + bufferStart -= length; + readExact(fd, buffer, bufferStart, length, position); + bytesRead += length; + newlineCount += countNewlines( + buffer.subarray(bufferStart, bufferStart + length), + ); + } + + let populated: Buffer = buffer.subarray(bufferStart); + const startedMidFile = position > 0; + let partialFirstLine: Buffer | undefined; + if (startedMidFile) { + const firstNewline = populated.indexOf(0x0a); + if (firstNewline === -1) { + populated = alignUtf8Start(populated); + const suffix = decodeUtf8Bounded(populated, MAX_TAIL_READ_BYTES); + return [`${TAIL_TRUNCATION_PREFIX}${suffix}`]; + } + // The read began in the middle of a logical line. Dropping through + // its newline also removes any partial UTF-8 code point at the start. + const partial = alignUtf8Start(populated.subarray(0, firstNewline)); + if (partial.length > 0) { + partialFirstLine = partial; + } + populated = populated.subarray(firstNewline + 1); } - return allLines.slice(-lines); + + const rawLines = tailLineBuffersNewestFirst(populated, lines); + if (partialFirstLine && rawLines.length < lines) { + rawLines.push(partialFirstLine); + } + + const decodedNewestFirst: string[] = []; + let remainingBytes = MAX_TAIL_READ_BYTES; + for ( + let index = 0; + index < rawLines.length && remainingBytes > 0; + index++ + ) { + const rawLine = rawLines[index]; + const decoded = decodeUtf8Bounded(rawLine, remainingBytes); + const marked = + partialFirstLine && rawLine === partialFirstLine + ? `${TAIL_TRUNCATION_PREFIX}${decoded}` + : decoded; + decodedNewestFirst.push(marked); + remainingBytes -= Buffer.byteLength(decoded); + } + return decodedNewestFirst.reverse(); } catch (_error) { return []; + } finally { + if (fd !== undefined) closeSync(fd); } } readFullFile(filePath: string): string { try { + const size = statSync(filePath).size; + if (size > MAX_FULL_FILE_BYTES) { + const marker = `[… truncated, see ${filePath}]\n`; + const markerBytes = Buffer.byteLength(marker); + const suffix = this.readFileSuffix( + filePath, + Math.max(0, MAX_FULL_FILE_BYTES - markerBytes), + ); + return marker + suffix; + } return readFileSync(filePath, "utf-8"); } catch (_error) { return ""; @@ -148,4 +236,123 @@ export class ProcessLogStore { [Symbol.dispose](): void { this.cleanup(); } + + private readFileSuffix(filePath: string, maxBytes: number): string { + let fd: number | undefined; + try { + fd = openSync(filePath, "r"); + const size = fstatSync(fd).size; + const length = Math.min(size, maxBytes); + const buffer = Buffer.allocUnsafe(length); + readExact(fd, buffer, 0, length, size - length); + return decodeUtf8Bounded(alignUtf8Start(buffer), maxBytes); + } catch (_error) { + return ""; + } finally { + if (fd !== undefined) closeSync(fd); + } + } +} + +function countNewlines(buffer: Buffer): number { + let count = 0; + for (const byte of buffer) { + if (byte === 0x0a) count++; + } + return count; +} + +function alignUtf8Start(buffer: Buffer): Buffer { + let start = 0; + while (start < buffer.length && (buffer[start] & 0xc0) === 0x80) start++; + return buffer.subarray(start); +} + +function readExact( + fd: number, + buffer: Buffer, + offset: number, + length: number, + position: number, +): void { + let total = 0; + while (total < length) { + const count = readSync( + fd, + buffer, + offset + total, + length - total, + position + total, + ); + if (count === 0) throw new Error("Unexpected end of log file"); + total += count; + } +} + +function tailLineBuffersNewestFirst(buffer: Buffer, count: number): Buffer[] { + if (buffer.length === 0 || count <= 0) return []; + + const result: Buffer[] = []; + let end = buffer.length; + if (buffer[end - 1] === 0x0a) end--; + + while (result.length < count && end >= 0) { + const newline = buffer.lastIndexOf(0x0a, end - 1); + const start = newline + 1; + let lineEnd = end; + if (lineEnd > start && buffer[lineEnd - 1] === 0x0d) lineEnd--; + result.push(buffer.subarray(start, lineEnd)); + if (newline === -1) break; + end = newline; + } + + return result; +} + +function decodeUtf8Bounded(buffer: Buffer, maxOutputBytes: number): string { + if (buffer.length === 0 || maxOutputBytes <= 0) return ""; + + const decoder = new TextDecoder("utf-8"); + const parts: string[] = []; + let remaining = maxOutputBytes; + + for (let offset = 0; offset < buffer.length && remaining > 0; ) { + const end = Math.min(buffer.length, offset + TAIL_CHUNK_BYTES); + const text = decoder.decode(buffer.subarray(offset, end), { + stream: end < buffer.length, + }); + const textBytes = Buffer.byteLength(text); + if (textBytes <= remaining) { + parts.push(text); + remaining -= textBytes; + } else { + const prefix = trimIncompleteUtf8Suffix( + Buffer.from(text).subarray(0, remaining), + ); + parts.push(prefix.toString("utf-8")); + remaining = 0; + } + offset = end; + } + + return parts.join(""); +} + +function trimIncompleteUtf8Suffix(buffer: Buffer): Buffer { + if (buffer.length === 0) return buffer; + let lead = buffer.length - 1; + while (lead >= 0 && (buffer[lead] & 0xc0) === 0x80) lead--; + if (lead < 0) return Buffer.alloc(0); + const byte = buffer[lead]; + const expected = + byte < 0x80 + ? 1 + : (byte & 0xe0) === 0xc0 + ? 2 + : (byte & 0xf0) === 0xe0 + ? 3 + : (byte & 0xf8) === 0xf0 + ? 4 + : 1; + return buffer.length - lead < expected ? buffer.subarray(0, lead) : buffer; } diff --git a/src/manager/process-output.test.ts b/src/manager/process-output.test.ts index 9ed2fb2..802118b 100644 --- a/src/manager/process-output.test.ts +++ b/src/manager/process-output.test.ts @@ -2,6 +2,12 @@ import { createMock, type PartialFuncReturn } from "@golevelup/ts-vitest"; import { beforeEach, describe, expect, it, vi } from "vitest"; import type { ManagerEvent } from "../types"; import type { ManagedProcessRecord } from "./internal-types"; +import { + MAX_LINE_BYTES, + MAX_LINES_PER_EMIT, + MAX_PENDING_LINE_BYTES, + TRUNCATION_SUFFIX, +} from "./limits"; import type { ProcessLogStore } from "./process-log-store"; import { ProcessOutput } from "./process-output"; @@ -25,9 +31,12 @@ const recordDefaults = { stdin: null, stdinClosed: false, lastSignalSent: null, - stdoutPendingLine: "", - stderrPendingLine: "", + stdoutPendingLine: Buffer.alloc(0), + stderrPendingLine: Buffer.alloc(0), + stdoutLineOverflowed: false, + stderrLineOverflowed: false, appendedLines: [], + droppedLineCount: 0, } satisfies PartialFuncReturn; describe("ProcessOutput", () => { @@ -112,7 +121,7 @@ describe("ProcessOutput", () => { output.onStdoutChunk(record, Buffer.from("partial")); expect(combinedLines).toEqual([]); - expect(record.stdoutPendingLine).toBe("partial"); + expect(record.stdoutPendingLine.toString()).toBe("partial"); expect(emitted).toEqual([{ type: "process_output_changed", id: "proc_1" }]); output.onStdoutChunk(record, Buffer.from(" done\nnext")); @@ -120,14 +129,14 @@ describe("ProcessOutput", () => { expect(combinedLines).toEqual([ { file: "/tmp/combined.log", source: "stdout", line: "partial done" }, ]); - expect(record.stdoutPendingLine).toBe("next"); + expect(record.stdoutPendingLine.toString()).toBe("next"); }); it("flushes pending stdout and stderr lines", () => { using output = createOutput(); const record = createRecord({ - stdoutPendingLine: "stdout tail", - stderrPendingLine: "stderr tail", + stdoutPendingLine: Buffer.from("stdout tail"), + stderrPendingLine: Buffer.from("stderr tail"), }); output.flush(record); @@ -136,8 +145,10 @@ describe("ProcessOutput", () => { { file: "/tmp/combined.log", source: "stdout", line: "stdout tail" }, { file: "/tmp/combined.log", source: "stderr", line: "stderr tail" }, ]); - expect(record.stdoutPendingLine).toBe(""); - expect(record.stderrPendingLine).toBe(""); + expect(record.stdoutPendingLine).toHaveLength(0); + expect(record.stderrPendingLine).toHaveLength(0); + expect(record.stdoutLineOverflowed).toBe(false); + expect(record.stderrLineOverflowed).toBe(false); expect(emitted).toEqual([ { type: "process_output_changed", @@ -193,6 +204,29 @@ describe("ProcessOutput", () => { vi.useRealTimers(); }); + it("caps buffered lines after adding final partial output during flush", () => { + using output = createOutput(100); + const record = createRecord({ + appendedLines: Array.from({ length: MAX_LINES_PER_EMIT }, (_, index) => ({ + type: "stdout" as const, + text: `line-${index}`, + })), + stdoutPendingLine: Buffer.from("final tail"), + }); + + output.flush(record); + + const event = emitted[0]; + expect(event).toMatchObject({ + type: "process_output_changed", + droppedLines: 1, + }); + if (event?.type !== "process_output_changed") return; + expect(event.appendedText).toHaveLength(MAX_LINES_PER_EMIT); + expect(event.appendedText?.[0]?.text).toBe("line-1"); + expect(event.appendedText?.at(-1)?.text).toBe("final tail"); + }); + it("clear removes pending timers", () => { vi.useFakeTimers(); using output = createOutput(100); @@ -208,4 +242,114 @@ describe("ProcessOutput", () => { vi.useRealTimers(); }); + + it("emits one bounded line and drops the remainder until newline", () => { + using output = createOutput(0); + const record = createRecord(); + const chunks = ["a".repeat(40_000), "b".repeat(40_000), "c".repeat(40_000)]; + + const lines = chunks.flatMap((chunk) => + output.onStdoutChunk(record, Buffer.from(chunk)), + ); + + expect(lines).toHaveLength(1); + expect(lines[0].length).toBeLessThanOrEqual( + MAX_PENDING_LINE_BYTES + TRUNCATION_SUFFIX.length, + ); + expect(lines[0]?.endsWith(TRUNCATION_SUFFIX)).toBe(true); + expect(record.stdoutLineOverflowed).toBe(true); + + expect(output.onStdoutChunk(record, Buffer.from("tail\nafter\n"))).toEqual([ + "after", + ]); + }); + + it("truncates a complete overlong line and keeps following lines", () => { + using output = createOutput(0); + const record = createRecord(); + + const lines = output.onStdoutChunk( + record, + Buffer.from(`${"a".repeat(100_000)}\nshort\n`), + ); + + expect(lines).toHaveLength(2); + expect(lines[0]?.endsWith(TRUNCATION_SUFFIX)).toBe(true); + expect(lines[1]).toBe("short"); + }); + + it("handles CRLF split across chunks", () => { + using output = createOutput(0); + const record = createRecord(); + + output.onStdoutChunk(record, Buffer.from("line\r")); + const lines = output.onStdoutChunk(record, Buffer.from("\nnext\r\n")); + + expect(lines).toEqual(["line", "next"]); + }); + + it("enforces byte limits without splitting UTF-8 characters", () => { + using output = createOutput(0); + const record = createRecord(); + + const lines = output.onStdoutChunk( + record, + Buffer.from(`${"€".repeat(30_000)}\n`), + ); + + expect(Buffer.byteLength(lines[0] ?? "")).toBeLessThanOrEqual( + MAX_PENDING_LINE_BYTES + Buffer.byteLength(TRUNCATION_SUFFIX), + ); + const event = emitted[0]; + if (event?.type !== "process_output_changed") return; + expect( + Buffer.byteLength(event.appendedText?.[0]?.text ?? ""), + ).toBeLessThanOrEqual(MAX_LINE_BYTES); + expect(event.appendedText?.[0]?.text).not.toContain("�"); + }); + + it("caps an output burst and reports the newest retained lines", () => { + vi.useFakeTimers(); + using output = createOutput(100); + const record = createRecord(); + + output.onStdoutChunk(record, Buffer.from("initial\n")); + const lines = Array.from( + { length: MAX_LINES_PER_EMIT + 500 }, + (_, index) => `line-${index}`, + ); + output.onStdoutChunk(record, Buffer.from(`${lines.join("\n")}\n`)); + vi.advanceTimersByTime(100); + + const event = emitted.at(-1); + expect(event).toMatchObject({ + type: "process_output_changed", + droppedLines: 500, + }); + if (event?.type !== "process_output_changed") return; + expect(event.appendedText).toHaveLength(MAX_LINES_PER_EMIT); + expect(event.appendedText?.[0]?.text).toBe("line-500"); + expect(event.appendedText?.at(-1)?.text).toBe( + `line-${MAX_LINES_PER_EMIT + 499}`, + ); + + vi.useRealTimers(); + }); + + it("clamps event lines while preserving the longer combined-log line", () => { + using output = createOutput(0); + const record = createRecord(); + + output.onStdoutChunk(record, Buffer.from(`${"x".repeat(100_000)}\n`)); + + expect(combinedLines[0]?.line.length).toBeGreaterThan(MAX_LINE_BYTES); + const event = emitted[0]; + if (event?.type !== "process_output_changed") return; + expect(Buffer.byteLength(event.appendedText?.[0]?.text ?? "")).toBe( + MAX_LINE_BYTES, + ); + expect(event.appendedText?.[0]?.text.endsWith(TRUNCATION_SUFFIX)).toBe( + true, + ); + }); }); diff --git a/src/manager/process-output.ts b/src/manager/process-output.ts index b8a5d08..75307e6 100644 --- a/src/manager/process-output.ts +++ b/src/manager/process-output.ts @@ -1,5 +1,11 @@ import type { ManagerEvent } from "../types"; import type { ManagedProcessRecord } from "./internal-types"; +import { + MAX_LINE_BYTES, + MAX_LINES_PER_EMIT, + MAX_PENDING_LINE_BYTES, + TRUNCATION_SUFFIX, +} from "./limits"; import type { ProcessLogStore } from "./process-log-store"; interface ProcessOutputDeps { @@ -23,27 +29,22 @@ export class ProcessOutput { } onStdoutChunk(record: ManagedProcessRecord, data: Buffer): string[] { - const lines = this.extractCompleteLines(record, "stdout", data); - for (const line of lines) { + return this.extractCompleteLines(record, "stdout", data, (line) => { this.logStore.appendCombinedLine(record.combinedFile, "stdout", line); - record.appendedLines.push({ type: "stdout", text: line }); - } - this.notify(record); - return lines; + this.appendEventLine(record, "stdout", line); + }); } onStderrChunk(record: ManagedProcessRecord, data: Buffer): string[] { - const lines = this.extractCompleteLines(record, "stderr", data); - for (const line of lines) { + return this.extractCompleteLines(record, "stderr", data, (line) => { this.logStore.appendCombinedLine(record.combinedFile, "stderr", line); - record.appendedLines.push({ type: "stderr", text: line }); - } - this.notify(record); - return lines; + this.appendEventLine(record, "stderr", line); + }); } flush(record: ManagedProcessRecord): void { this.flushPendingLines(record); + this.trimEventBuffer(record); const timeout = this.pendingOutputEmit.get(record.id); if (timeout) { @@ -51,14 +52,15 @@ export class ProcessOutput { this.pendingOutputEmit.delete(record.id); } - const appendedText = this.drainAppendedLines(record); - if (!timeout && !appendedText) return; + const drained = this.drainAppendedLines(record); + if (!timeout && !drained) return; this.lastOutputEmitAt.set(record.id, Date.now()); this.emit({ type: "process_output_changed", id: record.id, - ...(appendedText ? { appendedText } : {}), + ...(drained?.appendedText ? { appendedText: drained.appendedText } : {}), + ...(drained?.droppedLines ? { droppedLines: drained.droppedLines } : {}), }); } @@ -78,17 +80,23 @@ export class ProcessOutput { } private notify(record: ManagedProcessRecord): void { + this.trimEventBuffer(record); const now = Date.now(); const lastEmit = this.lastOutputEmitAt.get(record.id) ?? 0; const elapsed = now - lastEmit; if (elapsed >= this.throttleMs) { this.lastOutputEmitAt.set(record.id, now); - const appendedText = this.drainAppendedLines(record); + const drained = this.drainAppendedLines(record); this.emit({ type: "process_output_changed", id: record.id, - ...(appendedText ? { appendedText } : {}), + ...(drained?.appendedText + ? { appendedText: drained.appendedText } + : {}), + ...(drained?.droppedLines + ? { droppedLines: drained.droppedLines } + : {}), }); return; } @@ -97,13 +105,18 @@ export class ProcessOutput { const delay = this.throttleMs - elapsed; const timeout = setTimeout(() => { this.pendingOutputEmit.delete(record.id); - const appendedText = this.drainAppendedLines(record); - if (!appendedText) return; + const drained = this.drainAppendedLines(record); + if (!drained) return; this.lastOutputEmitAt.set(record.id, Date.now()); this.emit({ type: "process_output_changed", id: record.id, - appendedText, + ...(drained.appendedText + ? { appendedText: drained.appendedText } + : {}), + ...(drained.droppedLines + ? { droppedLines: drained.droppedLines } + : {}), }); }, delay); this.pendingOutputEmit.set(record.id, timeout); @@ -111,65 +124,222 @@ export class ProcessOutput { } private flushPendingLines(record: ManagedProcessRecord): void { - if (record.stdoutPendingLine) { - this.logStore.appendCombinedLine( - record.combinedFile, - "stdout", - record.stdoutPendingLine, - ); - record.appendedLines.push({ - type: "stdout", - text: record.stdoutPendingLine, - }); - record.stdoutPendingLine = ""; + if (record.stdoutPendingLine.length > 0) { + const line = record.stdoutPendingLine.toString("utf-8"); + this.logStore.appendCombinedLine(record.combinedFile, "stdout", line); + this.appendEventLine(record, "stdout", line); } + record.stdoutPendingLine = Buffer.alloc(0); + record.stdoutLineOverflowed = false; - if (record.stderrPendingLine) { - this.logStore.appendCombinedLine( - record.combinedFile, - "stderr", - record.stderrPendingLine, - ); - record.appendedLines.push({ - type: "stderr", - text: record.stderrPendingLine, - }); - record.stderrPendingLine = ""; + if (record.stderrPendingLine.length > 0) { + const line = record.stderrPendingLine.toString("utf-8"); + this.logStore.appendCombinedLine(record.combinedFile, "stderr", line); + this.appendEventLine(record, "stderr", line); } + record.stderrPendingLine = Buffer.alloc(0); + record.stderrLineOverflowed = false; } - private drainAppendedLines( - record: ManagedProcessRecord, - ): Array<{ type: "stdout" | "stderr"; text: string }> | undefined { - if (record.appendedLines.length === 0) return undefined; + private drainAppendedLines(record: ManagedProcessRecord): + | { + appendedText?: Array<{ + type: "stdout" | "stderr"; + text: string; + }>; + droppedLines?: number; + } + | undefined { + if (record.appendedLines.length === 0 && record.droppedLineCount === 0) { + return undefined; + } const lines = record.appendedLines; + const droppedLines = record.droppedLineCount; record.appendedLines = []; - return lines; + record.droppedLineCount = 0; + return { + ...(lines.length > 0 ? { appendedText: lines } : {}), + ...(droppedLines > 0 ? { droppedLines } : {}), + }; } private extractCompleteLines( record: ManagedProcessRecord, source: "stdout" | "stderr", data: Buffer, + onLine: (line: string) => void, ): string[] { - const chunk = data.toString(); - const pending = - source === "stdout" ? record.stdoutPendingLine : record.stderrPendingLine; - const merged = pending + chunk; - const parts = merged.split(/\r?\n/); - const completeLines = parts.slice(0, -1); - const nextPending = parts[parts.length - 1] ?? ""; - - if (source === "stdout") { - record.stdoutPendingLine = nextPending; + let pending = this.getPending(record, source); + let overflowed = this.getOverflowed(record, source); + const completeLines: string[] = []; + let cursor = 0; + + if (overflowed) { + const newline = data.indexOf(0x0a); + if (newline === -1) { + this.notify(record); + return []; + } + cursor = newline + 1; + overflowed = false; + pending = Buffer.alloc(0); + } + + let newline = data.indexOf(0x0a, cursor); + while (newline !== -1) { + const line = this.buildBoundedLine( + pending, + data.subarray(cursor, newline), + true, + ); + onLine(line); + if (completeLines.length < MAX_LINES_PER_EMIT) completeLines.push(line); + pending = Buffer.alloc(0); + cursor = newline + 1; + newline = data.indexOf(0x0a, cursor); + } + + const tail = data.subarray(cursor); + const tailLength = tail.length; + if (pending.length + tailLength > MAX_PENDING_LINE_BYTES) { + const line = this.buildBoundedLine(pending, tail, false); + onLine(line); + if (completeLines.length < MAX_LINES_PER_EMIT) completeLines.push(line); + pending = Buffer.alloc(0); + overflowed = true; + // The raw stdout/stderr log still has the complete bytes. Dropping the + // event-stream tail avoids turning one huge logical line into thousands + // of synthetic lines in UI buffers. } else { - record.stderrPendingLine = nextPending; + pending = + pending.length === 0 + ? Buffer.from(tail) + : Buffer.concat([pending, tail], pending.length + tail.length); } + this.setPending(record, source, pending); + this.setOverflowed(record, source, overflowed); + this.notify(record); return completeLines; } + private buildBoundedLine( + pending: Buffer, + segment: Buffer, + stripCarriageReturn: boolean, + ): string { + let pendingEnd = pending.length; + let segmentEnd = segment.length; + if (stripCarriageReturn) { + if (segmentEnd > 0 && segment[segmentEnd - 1] === 0x0d) segmentEnd--; + else if ( + segmentEnd === 0 && + pendingEnd > 0 && + pending[pendingEnd - 1] === 0x0d + ) { + pendingEnd--; + } + } + + const totalLength = pendingEnd + segmentEnd; + if (totalLength <= MAX_PENDING_LINE_BYTES) { + return Buffer.concat( + [pending.subarray(0, pendingEnd), segment.subarray(0, segmentEnd)], + totalLength, + ).toString("utf-8"); + } + + const prefix = Buffer.allocUnsafe(MAX_PENDING_LINE_BYTES); + const pendingLength = Math.min(pendingEnd, MAX_PENDING_LINE_BYTES); + pending.copy(prefix, 0, 0, pendingLength); + const segmentLength = MAX_PENDING_LINE_BYTES - pendingLength; + segment.copy(prefix, pendingLength, 0, segmentLength); + return `${trimIncompleteUtf8Suffix(prefix).toString("utf-8")}${TRUNCATION_SUFFIX}`; + } + + private appendEventLine( + record: ManagedProcessRecord, + type: "stdout" | "stderr", + text: string, + ): void { + record.appendedLines.push({ type, text: this.clampLine(text) }); + if (record.appendedLines.length <= MAX_LINES_PER_EMIT * 2) return; + this.trimEventBuffer(record); + } + + private trimEventBuffer(record: ManagedProcessRecord): void { + if (record.appendedLines.length <= MAX_LINES_PER_EMIT) return; + const overflow = record.appendedLines.length - MAX_LINES_PER_EMIT; + record.appendedLines.splice(0, overflow); + record.droppedLineCount += overflow; + } + + private clampLine(text: string): string { + if (Buffer.byteLength(text) <= MAX_LINE_BYTES) return text; + const suffixBytes = Buffer.byteLength(TRUNCATION_SUFFIX); + const prefix = Buffer.from(text).subarray(0, MAX_LINE_BYTES - suffixBytes); + return `${trimIncompleteUtf8Suffix(prefix).toString("utf-8")}${TRUNCATION_SUFFIX}`; + } + + private getPending( + record: ManagedProcessRecord, + source: "stdout" | "stderr", + ): Buffer { + return source === "stdout" + ? record.stdoutPendingLine + : record.stderrPendingLine; + } + + private setPending( + record: ManagedProcessRecord, + source: "stdout" | "stderr", + value: Buffer, + ): void { + if (source === "stdout") record.stdoutPendingLine = value; + else record.stderrPendingLine = value; + } + + private getOverflowed( + record: ManagedProcessRecord, + source: "stdout" | "stderr", + ): boolean { + return source === "stdout" + ? record.stdoutLineOverflowed + : record.stderrLineOverflowed; + } + + private setOverflowed( + record: ManagedProcessRecord, + source: "stdout" | "stderr", + value: boolean, + ): void { + if (source === "stdout") record.stdoutLineOverflowed = value; + else record.stderrLineOverflowed = value; + } + [Symbol.dispose](): void { this.clearAll(); } } + +function trimIncompleteUtf8Suffix(buffer: Buffer): Buffer { + if (buffer.length === 0) return buffer; + + let lead = buffer.length - 1; + while (lead >= 0 && (buffer[lead] & 0xc0) === 0x80) lead--; + if (lead < 0) return Buffer.alloc(0); + + const leadByte = buffer[lead]; + const expectedLength = + leadByte < 0x80 + ? 1 + : (leadByte & 0xe0) === 0xc0 + ? 2 + : (leadByte & 0xf0) === 0xe0 + ? 3 + : (leadByte & 0xf8) === 0xf0 + ? 4 + : 1; + const actualLength = buffer.length - lead; + return actualLength < expectedLength ? buffer.subarray(0, lead) : buffer; +} diff --git a/src/manager/process-registry.test.ts b/src/manager/process-registry.test.ts index b1e3be0..551ae2d 100644 --- a/src/manager/process-registry.test.ts +++ b/src/manager/process-registry.test.ts @@ -23,9 +23,12 @@ const managedDefaults = { stdin: null, stdinClosed: false, lastSignalSent: null, - stdoutPendingLine: "", - stderrPendingLine: "", + stdoutPendingLine: Buffer.alloc(0), + stderrPendingLine: Buffer.alloc(0), + stdoutLineOverflowed: false, + stderrLineOverflowed: false, appendedLines: [], + droppedLineCount: 0, } satisfies PartialFuncReturn; describe("ProcessRegistry", () => { diff --git a/src/manager/process-runtime-controller.ts b/src/manager/process-runtime-controller.ts index f75c387..e902c9c 100644 --- a/src/manager/process-runtime-controller.ts +++ b/src/manager/process-runtime-controller.ts @@ -64,13 +64,15 @@ export class ProcessRuntimeController { signal: null, errorMessage: null, combinedFile: logPaths.combinedFile, - process: child, stdin: child.stdin, stdinClosed: false, lastSignalSent: null, - stdoutPendingLine: "", - stderrPendingLine: "", + stdoutPendingLine: Buffer.alloc(0), + stderrPendingLine: Buffer.alloc(0), + stdoutLineOverflowed: false, + stderrLineOverflowed: false, appendedLines: [], + droppedLineCount: 0, }; this.registry.add(managed); diff --git a/src/protocol/broadcasts.ts b/src/protocol/broadcasts.ts index 95414e1..3bbd2c9 100644 --- a/src/protocol/broadcasts.ts +++ b/src/protocol/broadcasts.ts @@ -7,6 +7,7 @@ export type ProcessesEndedPayload = ProcessInfo; export type ProcessesOutputChangedPayload = { id: string; appendedText?: Array<{ type: "stdout" | "stderr"; text: string }>; + droppedLines?: number; }; export type ProcessesChangedPayload = { reason: "started" | "ended" | "cleared"; diff --git a/src/types.ts b/src/types.ts index ecb8bbd..5012426 100644 --- a/src/types.ts +++ b/src/types.ts @@ -50,6 +50,7 @@ export type ManagerEvent = type: "process_output_changed"; id: string; appendedText?: Array<{ type: "stdout" | "stderr"; text: string }>; + droppedLines?: number; } | { type: "processes_changed" };