diff --git a/src/manager/index.test.ts b/src/manager/index.test.ts index 2dd5341..68f425b 100644 --- a/src/manager/index.test.ts +++ b/src/manager/index.test.ts @@ -1,7 +1,7 @@ import { EventEmitter } from "node:events"; import { PassThrough } from "node:stream"; -import { vol } from "memfs"; +import { fs, vol } from "memfs"; import { afterEach, assert, @@ -14,6 +14,7 @@ import { import type { ManagerEvent } from "../types"; import { LIVE_STATUSES } from "../types"; import { ProcessManager } from "."; +import { FINISHED_RECORD_GRACE_MS, MAX_FINISHED_RECORDS } from "./limits"; const fakeProcesses = new Map(); let nextPid = 10_000; @@ -741,6 +742,59 @@ describe("clearFinished", () => { }); }); +describe("finished record reaping", () => { + it("keeps recent records then reaps the oldest records after the grace period", async () => { + vi.useFakeTimers(); + vi.setSystemTime(new Date("2026-01-01T00:00:00Z")); + using manager = new ProcessManager(); + const infos = []; + const stdoutFiles: string[] = []; + + for (let index = 0; index < MAX_FINISHED_RECORDS + 10; index++) { + const info = manager.start(`task-${index}`, "true", "/tmp"); + const ended = waitForEnd(manager, info.id); + await ended; + infos.push(info); + stdoutFiles.push(info.stdoutFile); + vi.advanceTimersByTime(1); + } + + expect(manager.list()).toHaveLength(MAX_FINISHED_RECORDS + 10); + expect(fs.existsSync(stdoutFiles[0] ?? "")).toBe(true); + + await vi.advanceTimersByTimeAsync(FINISHED_RECORD_GRACE_MS + 10); + + const survivors = manager.list(); + expect(survivors).toHaveLength(MAX_FINISHED_RECORDS); + expect(survivors.map((info) => info.id)).toEqual( + infos.slice(-MAX_FINISHED_RECORDS).map((info) => info.id), + ); + expect(fs.existsSync(stdoutFiles[0] ?? "")).toBe(false); + }); + + it("uses completion order when end times are equal", async () => { + vi.useFakeTimers(); + vi.setSystemTime(new Date("2026-01-01T00:00:00Z")); + using manager = new ProcessManager(); + const infos = Array.from({ length: MAX_FINISHED_RECORDS + 2 }, (_, index) => + manager.start(`task-${index}`, "sleep 60", "/tmp"), + ); + const ended = waitForEndedCount(manager, infos.length); + + for (const info of [...infos].reverse()) { + const child = fakeProcesses.get(info.pid); + assert(child, "fake child should exist"); + child.finish(0); + } + await ended; + await vi.advanceTimersByTimeAsync(FINISHED_RECORD_GRACE_MS); + + expect(manager.list().map((info) => info.id)).toEqual( + infos.slice(0, MAX_FINISHED_RECORDS).map((info) => info.id), + ); + }); +}); + // --- Output retrieval --- describe("output retrieval", () => { diff --git a/src/manager/index.ts b/src/manager/index.ts index 32c6249..d39d1c5 100644 --- a/src/manager/index.ts +++ b/src/manager/index.ts @@ -161,6 +161,7 @@ export class ProcessManager { } cleanup(): void { + this.runtime.beginShutdown(); this.runtime.stopWatcher(); this.output.clearAll(); this.runtime.killAllLive(); diff --git a/src/manager/internal-types.ts b/src/manager/internal-types.ts index faa1b65..0e2e535 100644 --- a/src/manager/internal-types.ts +++ b/src/manager/internal-types.ts @@ -48,6 +48,7 @@ interface ProcessRuntimeState { stdin: Writable | null; stdinClosed: boolean; lastSignalSent: NodeJS.Signals | null; + completionSequence: number | null; } /** diff --git a/src/manager/limits.ts b/src/manager/limits.ts index 48c8c69..400f27a 100644 --- a/src/manager/limits.ts +++ b/src/manager/limits.ts @@ -7,3 +7,7 @@ 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]"; +/** Max finished process records retained after the grace period. */ +export const MAX_FINISHED_RECORDS = 50; +/** Keep recently finished records available to the UI before reaping. */ +export const FINISHED_RECORD_GRACE_MS = 5 * 60 * 1000; diff --git a/src/manager/process-output.test.ts b/src/manager/process-output.test.ts index 802118b..332fae5 100644 --- a/src/manager/process-output.test.ts +++ b/src/manager/process-output.test.ts @@ -31,6 +31,7 @@ const recordDefaults = { stdin: null, stdinClosed: false, lastSignalSent: null, + completionSequence: null, stdoutPendingLine: Buffer.alloc(0), stderrPendingLine: Buffer.alloc(0), stdoutLineOverflowed: false, diff --git a/src/manager/process-registry.test.ts b/src/manager/process-registry.test.ts index 551ae2d..7244fd8 100644 --- a/src/manager/process-registry.test.ts +++ b/src/manager/process-registry.test.ts @@ -23,6 +23,7 @@ const managedDefaults = { stdin: null, stdinClosed: false, lastSignalSent: null, + completionSequence: null, stdoutPendingLine: Buffer.alloc(0), stderrPendingLine: Buffer.alloc(0), stdoutLineOverflowed: false, diff --git a/src/manager/process-runtime-controller.ts b/src/manager/process-runtime-controller.ts index e902c9c..629aa22 100644 --- a/src/manager/process-runtime-controller.ts +++ b/src/manager/process-runtime-controller.ts @@ -7,6 +7,7 @@ import { spawnCommand } from "../utils/command-executor"; import { formatSignalInfo } from "../utils/signals"; import type { ManagedProcessRecord } from "./internal-types"; import { formatProcess } from "./internal-types"; +import { FINISHED_RECORD_GRACE_MS, MAX_FINISHED_RECORDS } from "./limits"; import type { ProcessLogStore } from "./process-log-store"; import type { ProcessOutput } from "./process-output"; import type { ProcessRegistry } from "./process-registry"; @@ -27,6 +28,9 @@ export class ProcessRuntimeController { private getConfiguredShellPath: () => string | undefined; private watcher: ReturnType | null = null; + private finishedReapTimer: ReturnType | null = null; + private completionSequence = 0; + private shuttingDown = false; constructor(deps: ProcessRuntimeControllerDeps) { this.registry = deps.registry; @@ -67,6 +71,7 @@ export class ProcessRuntimeController { stdin: child.stdin, stdinClosed: false, lastSignalSent: null, + completionSequence: null, stdoutPendingLine: Buffer.alloc(0), stderrPendingLine: Buffer.alloc(0), stdoutLineOverflowed: false, @@ -84,6 +89,7 @@ export class ProcessRuntimeController { managed.endReason = "missing_pid"; managed.errorMessage = "Spawn error: missing pid"; managed.endTime = Date.now(); + this.releaseRuntimeHandles(managed); this.transition(managed, "exited"); return managed; } @@ -101,7 +107,10 @@ export class ProcessRuntimeController { managed.status = next; if (next === "exited" || next === "killed") { + managed.completionSequence ??= ++this.completionSequence; this.emit({ type: "process_ended", info: formatProcess(managed) }); + this.reapOldestFinished(); + this.scheduleFinishedReap(); } this.ensureWatcherRunning(); @@ -197,6 +206,7 @@ export class ProcessRuntimeController { } this.output.flush(managed); + this.releaseRuntimeHandles(managed); this.transition(managed, "killed"); return { ok: true, info: formatProcess(managed) }; } @@ -258,19 +268,12 @@ export class ProcessRuntimeController { clearFinished(): number { let cleared = 0; - for (const [id, managed] of this.registry.entries()) { + for (const managed of this.registry.values()) { if (LIVE_STATUSES.has(managed.status)) { continue; } - this.logs.removeLogs({ - stdoutFile: managed.stdoutFile, - stderrFile: managed.stderrFile, - combinedFile: managed.combinedFile, - }); - - this.output.clear(id); - this.registry.delete(id); + this.removeFinishedRecord(managed); cleared++; } @@ -278,6 +281,7 @@ export class ProcessRuntimeController { this.emit({ type: "processes_changed" }); } + this.scheduleFinishedReap(); this.stopWatcherIfIdle(); return cleared; } @@ -289,6 +293,17 @@ export class ProcessRuntimeController { } } + stopFinishedReaper(): void { + if (!this.finishedReapTimer) return; + clearTimeout(this.finishedReapTimer); + this.finishedReapTimer = null; + } + + beginShutdown(): void { + this.shuttingDown = true; + this.stopFinishedReaper(); + } + /** * Kill all live processes (used by cleanup on actual pi exit). */ @@ -305,6 +320,7 @@ export class ProcessRuntimeController { [Symbol.dispose](): void { this.stopWatcher(); + this.beginShutdown(); this.killAllLive(); } @@ -323,6 +339,7 @@ export class ProcessRuntimeController { }); child.on("close", (code, signal) => { + this.releaseRuntimeHandles(managed); if (managed.endTime) return; managed.exitCode = code; @@ -350,6 +367,7 @@ export class ProcessRuntimeController { ); if (!managed.endTime) { + this.releaseRuntimeHandles(managed); managed.exitCode = -1; managed.success = false; managed.endReason = "spawn_error"; @@ -398,6 +416,7 @@ export class ProcessRuntimeController { } this.output.flush(managed); + this.releaseRuntimeHandles(managed); managed.success = false; managed.exitCode = null; @@ -412,4 +431,69 @@ export class ProcessRuntimeController { } } } + + private releaseRuntimeHandles(managed: ManagedProcessRecord): void { + managed.stdin = null; + managed.stdinClosed = true; + } + + private reapOldestFinished(): void { + const finished = this.finishedRecordsByAge(); + if (finished.length <= MAX_FINISHED_RECORDS) return; + + const cutoff = Date.now() - FINISHED_RECORD_GRACE_MS; + let remaining = finished.length; + let reaped = 0; + for (const record of finished) { + if (remaining <= MAX_FINISHED_RECORDS) break; + if ((record.endTime ?? Number.POSITIVE_INFINITY) > cutoff) break; + this.removeFinishedRecord(record); + remaining--; + reaped++; + } + + if (reaped > 0) this.emit({ type: "processes_changed" }); + } + + private scheduleFinishedReap(): void { + this.stopFinishedReaper(); + if (this.shuttingDown) return; + const finished = this.finishedRecordsByAge(); + if (finished.length <= MAX_FINISHED_RECORDS) return; + + const oldest = finished[0]; + const delay = Math.max( + 0, + (oldest.endTime ?? Date.now()) + FINISHED_RECORD_GRACE_MS - Date.now(), + ); + this.finishedReapTimer = setTimeout(() => { + this.finishedReapTimer = null; + this.reapOldestFinished(); + this.scheduleFinishedReap(); + }, delay); + this.finishedReapTimer.unref?.(); + } + + private finishedRecordsByAge(): ManagedProcessRecord[] { + return [...this.registry.values()] + .filter( + (record) => + !LIVE_STATUSES.has(record.status) && record.endTime !== null, + ) + .sort( + (a, b) => + (a.endTime ?? 0) - (b.endTime ?? 0) || + (a.completionSequence ?? 0) - (b.completionSequence ?? 0), + ); + } + + private removeFinishedRecord(record: ManagedProcessRecord): void { + this.logs.removeLogs({ + stdoutFile: record.stdoutFile, + stderrFile: record.stderrFile, + combinedFile: record.combinedFile, + }); + this.output.clear(record.id); + this.registry.delete(record.id); + } }