diff --git a/packages/cli/src/commands/preview.test.ts b/packages/cli/src/commands/preview.test.ts index 4d6c7528fd..6af3f53a7d 100644 --- a/packages/cli/src/commands/preview.test.ts +++ b/packages/cli/src/commands/preview.test.ts @@ -6,6 +6,7 @@ import { afterEach, describe, expect, it, vi } from "vitest"; import { runCommand } from "citty"; import { default as previewCommand, + drainEmbeddedPreviewResources, foregroundPreviewReadyPayload, handlePreviewKillAll, handlePreviewList, @@ -493,3 +494,33 @@ describe("waitForStudioChildClose", () => { expect(signalTarget.off).toHaveBeenCalledTimes(2); }); }); + +describe("embedded preview resource drain", () => { + it("closes child registration before cancelling and awaiting renders", async () => { + const order: string[] = []; + let finishRenderDrain!: () => void; + const renderDrain = new Promise((resolve) => { + finishRenderDrain = resolve; + }); + const draining = drainEmbeddedPreviewResources({ + beginProcessDrain: () => order.push("processes"), + disposeRenders: async () => { + order.push("renders:start"); + await renderDrain; + order.push("renders:end"); + }, + closeThumbnailBrowser: async () => { + order.push("thumbnail"); + }, + drainBrowserPool: async () => { + order.push("browsers"); + }, + }); + await Promise.resolve(); + + expect(order).toEqual(["processes", "renders:start", "thumbnail"]); + finishRenderDrain(); + await draining; + expect(order).toEqual(["processes", "renders:start", "thumbnail", "renders:end", "browsers"]); + }); +}); diff --git a/packages/cli/src/commands/preview.ts b/packages/cli/src/commands/preview.ts index 5203d7eead..fc143e95ee 100644 --- a/packages/cli/src/commands/preview.ts +++ b/packages/cli/src/commands/preview.ts @@ -90,6 +90,17 @@ interface BrowserLaunchOptions { browserNoGpu?: boolean; } +export async function drainEmbeddedPreviewResources(input: { + beginProcessDrain: () => void; + disposeRenders: () => Promise; + closeThumbnailBrowser: () => Promise; + drainBrowserPool: () => Promise; +}): Promise { + input.beginProcessDrain(); + await Promise.allSettled([input.disposeRenders(), input.closeThumbnailBrowser()]); + await input.drainBrowserPool().catch(() => {}); +} + interface StudioLaunchOptions extends BrowserLaunchOptions { projectName?: string; autoProxy?: boolean; @@ -1566,7 +1577,7 @@ async function runEmbeddedMode( // Compute everything that may throw before acquiring the fs.watch handle. // Once createStudioServer returns, every subsequent exit path must close it. const serverBuildSignature = await loadPreviewServerBuildSignature(); - const { app, watcher } = createStudioServer({ + const { app, watcher, dispose } = createStudioServer({ projectDir: dir, projectName: pName, autoProxy: options?.autoProxy, @@ -1584,6 +1595,7 @@ async function runEmbeddedMode( options?.browserGpuMode, ); } catch (err: unknown) { + await dispose().catch(() => {}); watcher.close(); reportPreviewFailure( Boolean(options?.json), @@ -1598,6 +1610,7 @@ async function runEmbeddedMode( // createStudioServer acquires an fs.watch handle before port discovery. // Reuse owns no local server, so release that handle before returning or // the otherwise-finished CLI process remains alive indefinitely. + await dispose().catch(() => {}); watcher.close(); const url = `http://localhost:${result.port}`; if (options?.json) { @@ -1685,10 +1698,13 @@ async function runEmbeddedMode( // Kill ffmpeg first (sync, fast), then drain browsers (async, slower). const cleanup = async () => { const { closeThumbnailBrowser } = await import("../server/studioServer.js"); - const { drainBrowserPool, killTrackedProcesses } = await import("@hyperframes/engine"); - killTrackedProcesses(); - await closeThumbnailBrowser().catch(() => {}); - await drainBrowserPool().catch(() => {}); + const { beginTrackedProcessDrain, drainBrowserPool } = await import("@hyperframes/engine"); + await drainEmbeddedPreviewResources({ + beginProcessDrain: beginTrackedProcessDrain, + disposeRenders: dispose, + closeThumbnailBrowser, + drainBrowserPool, + }); }; cleanup() diff --git a/packages/cli/src/server/studioServer.ts b/packages/cli/src/server/studioServer.ts index acc2e53d8b..d2bd0bb4c5 100644 --- a/packages/cli/src/server/studioServer.ts +++ b/packages/cli/src/server/studioServer.ts @@ -306,6 +306,8 @@ export interface StudioServerOptions { export interface StudioServer { app: Hono; watcher: ProjectWatcher; + /** Cancel and await every render owned by this server. */ + dispose(): Promise; /** Exposed for tests: the adapter handed to the shared studio API (carries * the resolved `autoProxy` flag the preview routes read). */ adapter: PreviewApiAdapter; @@ -474,7 +476,7 @@ export function createStudioServer(options: StudioServerOptions): StudioServer { // Run render asynchronously, mutating the state object const startTime = Date.now(); - (async () => { + const completion = (async () => { let renderJob: RenderJob | undefined; const removeCancelledOutput = () => { // User-initiated cancel: not a failure. Remove any output so the @@ -567,6 +569,7 @@ export function createStudioServer(options: StudioServerOptions): StudioServer { } } })(); + state.completion = completion; return state; }, @@ -983,5 +986,5 @@ export function createStudioServer(options: StudioServerOptions): StudioServer { return c.html(html); }); - return { app, watcher, adapter }; + return { app, watcher, adapter, dispose: () => api.dispose() }; } diff --git a/packages/cli/src/utils/orphanCleanup.test.ts b/packages/cli/src/utils/orphanCleanup.test.ts index a6e0cca924..756ebdecd5 100644 --- a/packages/cli/src/utils/orphanCleanup.test.ts +++ b/packages/cli/src/utils/orphanCleanup.test.ts @@ -1,7 +1,8 @@ -import { describe, it, expect } from "vitest"; +import { describe, it, expect, vi } from "vitest"; import { spawn } from "node:child_process"; import { isProcessDescendant, + killOwnedOrphanedFfmpegProcesses, killProcessTree, killOrphanedProcesses, processIdentity, @@ -72,6 +73,39 @@ describe("process-tree ownership", () => { }); }); +describe("owned FFmpeg orphan cleanup", () => { + it("kills only the ownership-verified PID list", () => { + const kill = vi.fn(); + const records = [ + { pid: 41, identity: "linux:one" }, + { pid: 42, identity: "linux:two" }, + ]; + + expect( + killOwnedOrphanedFfmpegProcesses( + records, + kill, + (pid) => records.find((record) => record.pid === pid)?.identity ?? null, + ), + ).toBe(2); + expect(kill.mock.calls.map(([pid]) => pid)).toEqual([41, 42]); + expect(kill.mock.calls.every(([, , stillOwned]) => stillOwned())).toBe(true); + }); + + it("does not kill when the PID birth identity changed after discovery", () => { + const kill = vi.fn(); + + expect( + killOwnedOrphanedFfmpegProcesses( + [{ pid: 41, identity: "linux:original" }], + kill, + () => "linux:reused", + ), + ).toBe(0); + expect(kill).not.toHaveBeenCalled(); + }); +}); + describe.skipIf(!IS_UNIX)("killProcessTree", () => { it("kills a process and all its children", async () => { // Spawn a parent that spawns two sleeping children diff --git a/packages/cli/src/utils/orphanCleanup.ts b/packages/cli/src/utils/orphanCleanup.ts index 66dcad2991..df40e1c8a0 100644 --- a/packages/cli/src/utils/orphanCleanup.ts +++ b/packages/cli/src/utils/orphanCleanup.ts @@ -1,13 +1,23 @@ import { execFileSync, execSync } from "node:child_process"; -import { readFileSync } from "node:fs"; +import { + findOwnedOrphanedFfmpegProcesses, + type OwnedFfmpegProcess, + processIdentity, + processParentPid, +} from "@hyperframes/engine/process-tracker"; + +export { processIdentity }; /** - * Find and kill orphaned Chrome processes from previous crashed sessions. + * Find and kill orphaned Chrome and HyperFrames-owned FFmpeg processes from + * previous crashed sessions. * Targets both chrome-headless-shell (production/CI) and Google Chrome * launched by Puppeteer (dev mode). Puppeteer Chrome is identified by the * `puppeteer_dev_chrome_profile` marker in its user-data-dir argument. * - * An orphan is a process whose PPID=1 (reparented to init/launchd after + * FFmpeg recovery additionally requires a private process-tracker record with + * a matching birth identity, so an unrelated same-user encoder is never + * selected by name. An orphan is a process whose PPID=1 (reparented to init/launchd after * its parent died). We kill the orphan's entire subtree so child helper * processes (GPU, renderer, network, etc.) are also cleaned up. * @@ -26,7 +36,27 @@ export function killOrphanedProcesses(): number { } killed += killOrphansByName("puppeteer_dev_chrome_profile"); + killed += killOwnedOrphanedFfmpegProcesses(); + + return killed; +} +export function killOwnedOrphanedFfmpegProcesses( + records: OwnedFfmpegProcess[] = findOwnedOrphanedFfmpegProcesses(), + kill: ( + pid: number, + signal?: NodeJS.Signals, + stillOwned?: () => boolean, + ) => void = killProcessTree, + identityForPid: (pid: number) => string | null = processIdentity, +): number { + let killed = 0; + for (const record of records) { + const stillOwned = () => identityForPid(record.pid) === record.identity; + if (!stillOwned()) continue; + kill(record.pid, "SIGTERM", stillOwned); + killed++; + } return killed; } @@ -44,7 +74,12 @@ export function killOrphanedProcesses(): number { * ignore, and leaving a preview server alive is the worse failure here. Do not * pass SIGTERM expecting a clean shutdown on Windows. */ -export function killProcessTree(pid: number, signal: NodeJS.Signals = "SIGTERM"): void { +export function killProcessTree( + pid: number, + signal: NodeJS.Signals = "SIGTERM", + stillOwned: () => boolean = () => true, +): void { + if (!stillOwned()) return; if (process.platform === "win32") { try { execFileSync("taskkill", windowsProcessTreeKillArgs(pid), { @@ -60,8 +95,10 @@ export function killProcessTree(pid: number, signal: NodeJS.Signals = "SIGTERM") const descendants = getDescendants(pid); const allPids = [...descendants.reverse(), pid]; + const identities = new Map(allPids.map((candidate) => [candidate, processIdentity(candidate)])); for (const p of allPids) { + if (!stillOwned()) return; try { process.kill(p, signal); } catch { @@ -72,7 +109,10 @@ export function killProcessTree(pid: number, signal: NodeJS.Signals = "SIGTERM") // Escalate to SIGKILL after a short grace period for any survivors. if (signal !== "SIGKILL") { setTimeout(() => { + if (!stillOwned()) return; for (const p of allPids) { + const identity = identities.get(p); + if (!identity || processIdentity(p) !== identity) continue; try { process.kill(p, "SIGKILL"); } catch { @@ -87,85 +127,8 @@ export function windowsProcessTreeKillArgs(pid: number): string[] { return ["/PID", String(pid), "/T", "/F"]; } -/** - * Return a process birth token suitable for detecting PID reuse. The token is - * diagnostic state only: callers must still prove the live server is a - * descendant before treating a saved wrapper as the owned process-tree root. - */ -export function processIdentity(pid: number): string | null { - if (!Number.isInteger(pid) || pid <= 0) return null; - try { - if (process.platform === "win32") { - const created = execFileSync( - "powershell.exe", - [ - "-NoProfile", - "-NonInteractive", - "-Command", - `$p = Get-CimInstance Win32_Process -Filter 'ProcessId = ${pid}' -ErrorAction SilentlyContinue; if ($p) { $p.CreationDate.ToFileTimeUtc() }`, - ], - { - encoding: "utf8", - timeout: 2000, - stdio: ["pipe", "pipe", "ignore"], - windowsHide: true, - }, - ).trim(); - return created ? `windows:${created}` : null; - } - - if (process.platform === "linux") { - const stat = readFileSync(`/proc/${pid}/stat`, "utf8"); - const fields = stat - .slice(stat.lastIndexOf(") ") + 2) - .trim() - .split(/\s+/); - const startTicks = fields[19]; // field 22 overall; fields starts at process state (3) - return startTicks ? `linux:${startTicks}` : null; - } - - const started = execFileSync("ps", ["-o", "lstart=", "-p", String(pid)], { - encoding: "utf8", - timeout: 2000, - }).trim(); - return started ? `posix:${started}` : null; - } catch { - return null; - } -} - type ParentPidLookup = (pid: number) => number | null; -function processParentPid(pid: number): number | null { - try { - const output = - process.platform === "win32" - ? execFileSync( - "powershell.exe", - [ - "-NoProfile", - "-NonInteractive", - "-Command", - `$p = Get-CimInstance Win32_Process -Filter 'ProcessId = ${pid}' -ErrorAction SilentlyContinue; if ($p) { $p.ParentProcessId }`, - ], - { - encoding: "utf8", - timeout: 2000, - stdio: ["pipe", "pipe", "ignore"], - windowsHide: true, - }, - ) - : execFileSync("ps", ["-o", "ppid=", "-p", String(pid)], { - encoding: "utf8", - timeout: 2000, - }); - const parentPid = Number(output.trim()); - return Number.isInteger(parentPid) && parentPid > 0 ? parentPid : null; - } catch { - return null; - } -} - /** * Prove that `childPid` currently belongs to the process tree rooted at * `ancestorPid`. The walk fails closed on missing, invalid, or cyclic process @@ -254,13 +217,5 @@ function getUid(): string | null { } function isOrphan(pid: number): boolean { - try { - const ppid = execSync(`ps -p ${pid} -o ppid=`, { - encoding: "utf-8", - timeout: 2000, - }).trim(); - return ppid === "1"; - } catch { - return false; - } + return processParentPid(pid) === 1; } diff --git a/packages/engine/package-subpaths.json b/packages/engine/package-subpaths.json index 285d44f41d..471670985e 100644 --- a/packages/engine/package-subpaths.json +++ b/packages/engine/package-subpaths.json @@ -19,6 +19,12 @@ "types": "./dist/utils/shaderTransitions.d.ts", "environments": ["browser", "bun", "node"] }, + "./process-tracker": { + "source": "./src/utils/processTracker.ts", + "runtime": "./dist/utils/processTracker.js", + "types": "./dist/utils/processTracker.d.ts", + "environments": ["bun", "node"] + }, "./package.json": { "source": "./package.json", "runtime": "./package.json", diff --git a/packages/engine/package.json b/packages/engine/package.json index 0d0892452e..2c98eebed8 100644 --- a/packages/engine/package.json +++ b/packages/engine/package.json @@ -30,6 +30,11 @@ "import": "./src/utils/shaderTransitions.ts", "types": "./src/utils/shaderTransitions.ts" }, + "./process-tracker": { + "bun": "./src/utils/processTracker.ts", + "import": "./src/utils/processTracker.ts", + "types": "./src/utils/processTracker.ts" + }, "./package.json": "./package.json" }, "publishConfig": { @@ -47,6 +52,10 @@ "import": "./dist/utils/shaderTransitions.js", "types": "./dist/utils/shaderTransitions.d.ts" }, + "./process-tracker": { + "import": "./dist/utils/processTracker.js", + "types": "./dist/utils/processTracker.d.ts" + }, "./package.json": "./package.json" }, "main": "./dist/index.js", diff --git a/packages/engine/src/index.ts b/packages/engine/src/index.ts index 6137591552..61ccf9bc79 100644 --- a/packages/engine/src/index.ts +++ b/packages/engine/src/index.ts @@ -332,7 +332,11 @@ export { FFPROBE_PATH_ENV, } from "./utils/ffmpegBinaries.js"; -export { trackChildProcess, killTrackedProcesses } from "./utils/processTracker.js"; +export { + beginTrackedProcessDrain, + killTrackedProcesses, + trackChildProcess, +} from "./utils/processTracker.js"; // drawElement self-verify comparison — shared by the streaming drain // (producer) and the parallel disk-path verify (parallelCoordinator). diff --git a/packages/engine/src/services/streamingEncoder.ts b/packages/engine/src/services/streamingEncoder.ts index e57946e0a8..e6ef7d96f8 100644 --- a/packages/engine/src/services/streamingEncoder.ts +++ b/packages/engine/src/services/streamingEncoder.ts @@ -464,7 +464,7 @@ export async function spawnStreamingEncoder( // See runFfmpeg.ts: keeps a console window off the user's desktop on Windows. windowsHide: true, }); - trackChildProcess(ffmpeg); + trackChildProcess(ffmpeg, { kind: "ffmpeg" }); let exitStatus: "running" | "success" | "error" = "running"; let stderr = ""; diff --git a/packages/engine/src/utils/gpuEncoder.ts b/packages/engine/src/utils/gpuEncoder.ts index 03b006fb68..4a1ce8f5ad 100644 --- a/packages/engine/src/utils/gpuEncoder.ts +++ b/packages/engine/src/utils/gpuEncoder.ts @@ -67,7 +67,7 @@ export async function detectGpuEncoder(): Promise { // See runFfmpeg.ts: keeps a console window off the user's desktop on Windows. windowsHide: true, }); - trackChildProcess(ffmpeg); + trackChildProcess(ffmpeg, { kind: "ffmpeg" }); let stdout = ""; ffmpeg.stdout.on("data", (data) => { stdout += data.toString(); @@ -150,7 +150,7 @@ async function canUseGpuEncoder(encoder: ConcreteGpuEncoder): Promise { stdio: ["ignore", "ignore", "pipe"], windowsHide: true, }); - trackChildProcess(ffmpeg); + trackChildProcess(ffmpeg, { kind: "ffmpeg" }); const outcome = await new ManagedChildProcess(ffmpeg, { deadlineAtMs: Date.now() + GPU_PROBE_TIMEOUT_MS, terminationGraceMs: GPU_PROBE_KILL_GRACE_MS, diff --git a/packages/engine/src/utils/processTracker.test.ts b/packages/engine/src/utils/processTracker.test.ts index b87a286159..6a795bfd08 100644 --- a/packages/engine/src/utils/processTracker.test.ts +++ b/packages/engine/src/utils/processTracker.test.ts @@ -1,6 +1,14 @@ import { describe, it, expect, beforeEach, vi } from "vitest"; import { spawn } from "node:child_process"; -import { trackChildProcess, killTrackedProcesses } from "./processTracker.js"; +import { mkdtempSync, readdirSync, rmSync, writeFileSync } from "node:fs"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import { + beginTrackedProcessDrain, + findOwnedOrphanedFfmpegProcesses, + trackChildProcess, + killTrackedProcesses, +} from "./processTracker.js"; // Reset tracked set between tests by killing everything beforeEach(() => { @@ -102,4 +110,61 @@ describe("killTrackedProcesses", () => { killTrackedProcesses(); killTrackedProcesses(); }); + + it("registers owned FFmpeg identity and removes it on clean exit", async () => { + const registryDir = mkdtempSync(join(tmpdir(), "hf-owned-ffmpeg-")); + const proc = spawn("sleep", ["60"], { stdio: "ignore" }); + const exitPromise = new Promise((resolve) => proc.on("close", () => resolve())); + try { + trackChildProcess(proc, { kind: "ffmpeg", registryDir }); + expect(readdirSync(registryDir)).toHaveLength(1); + + proc.kill("SIGTERM"); + await exitPromise; + expect(readdirSync(registryDir)).toHaveLength(0); + } finally { + proc.kill("SIGKILL"); + await exitPromise; + rmSync(registryDir, { recursive: true, force: true }); + } + }); + + it("recovers only identity-matched FFmpeg records reparented to init", () => { + const registryDir = mkdtempSync(join(tmpdir(), "hf-owned-ffmpeg-scan-")); + try { + for (const [pid, identity] of [ + [101, "linux:one"], + [102, "linux:two"], + [103, "linux:stale"], + ] as const) { + writeFileSync( + join(registryDir, `${pid}.json`), + JSON.stringify({ version: 1, kind: "ffmpeg", pid, identity }), + ); + } + + expect( + findOwnedOrphanedFfmpegProcesses({ + registryDir, + identityForPid: (pid) => + ({ 101: "linux:one", 102: "linux:two", 103: "linux:reused" })[pid] ?? null, + parentPidForPid: (pid) => (pid === 101 ? 1 : 77), + }), + ).toEqual([{ pid: 101, identity: "linux:one" }]); + expect(readdirSync(registryDir)).not.toContain("103.json"); + expect(readdirSync(registryDir)).toContain("102.json"); + } finally { + rmSync(registryDir, { recursive: true, force: true }); + } + }); + + it("kills a child registered after the terminal drain begins", async () => { + beginTrackedProcessDrain(); + const proc = spawn("sleep", ["60"], { stdio: "ignore" }); + const exitPromise = new Promise((resolve) => proc.on("close", resolve)); + + trackChildProcess(proc); + + expect(await exitPromise).toBeNull(); + }); }); diff --git a/packages/engine/src/utils/processTracker.ts b/packages/engine/src/utils/processTracker.ts index 877834b4f1..da7ab56706 100644 --- a/packages/engine/src/utils/processTracker.ts +++ b/packages/engine/src/utils/processTracker.ts @@ -1,12 +1,202 @@ -import type { ChildProcess } from "node:child_process"; +import { execFileSync, type ChildProcess } from "node:child_process"; +import { mkdirSync, readFileSync, readdirSync, unlinkSync, writeFileSync } from "node:fs"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; const tracked = new Set(); +let draining = false; -export function trackChildProcess(proc: ChildProcess): void { - tracked.add(proc); - const remove = () => tracked.delete(proc); +export interface TrackChildProcessOptions { + kind?: "ffmpeg"; + registryDir?: string; +} + +export function processIdentity(pid: number): string | null { + if (!Number.isInteger(pid) || pid <= 0) return null; + try { + if (process.platform === "win32") { + const created = execFileSync( + "powershell.exe", + [ + "-NoProfile", + "-NonInteractive", + "-Command", + `$p = Get-CimInstance Win32_Process -Filter 'ProcessId = ${pid}' -ErrorAction SilentlyContinue; if ($p) { $p.CreationDate.ToFileTimeUtc() }`, + ], + { + encoding: "utf8", + timeout: 2000, + stdio: ["pipe", "pipe", "ignore"], + windowsHide: true, + }, + ).trim(); + return created ? `windows:${created}` : null; + } + if (process.platform === "linux") { + const stat = readFileSync(`/proc/${pid}/stat`, "utf8"); + const fields = stat + .slice(stat.lastIndexOf(") ") + 2) + .trim() + .split(/\s+/); + const startTicks = fields[19]; + return startTicks ? `linux:${startTicks}` : null; + } + const started = execFileSync("ps", ["-o", "lstart=", "-p", String(pid)], { + encoding: "utf8", + timeout: 2000, + }).trim(); + return started ? `posix:${started}` : null; + } catch { + return null; + } +} + +export function ownedProcessRegistryDir(): string { + const uid = typeof process.getuid === "function" ? process.getuid() : "user"; + return join(tmpdir(), `hyperframes-owned-processes-${uid}`); +} + +export function processParentPid(pid: number): number | null { + if (!Number.isInteger(pid) || pid <= 0) return null; + try { + const output = + process.platform === "win32" + ? execFileSync( + "powershell.exe", + [ + "-NoProfile", + "-NonInteractive", + "-Command", + `$p = Get-CimInstance Win32_Process -Filter 'ProcessId = ${pid}' -ErrorAction SilentlyContinue; if ($p) { $p.ParentProcessId }`, + ], + { + encoding: "utf8", + timeout: 2000, + stdio: ["pipe", "pipe", "ignore"], + windowsHide: true, + }, + ) + : execFileSync("ps", ["-o", "ppid=", "-p", String(pid)], { + encoding: "utf8", + timeout: 2000, + }); + const parent = Number(output.trim()); + return Number.isInteger(parent) && parent > 0 ? parent : null; + } catch { + return null; + } +} + +export interface OwnedFfmpegProcess { + pid: number; + identity: string; +} + +export function findOwnedOrphanedFfmpegProcesses( + options: { + registryDir?: string; + identityForPid?: (pid: number) => string | null; + parentPidForPid?: (pid: number) => number | null; + } = {}, +): OwnedFfmpegProcess[] { + if (process.platform === "win32") return []; + const registryDir = options.registryDir ?? ownedProcessRegistryDir(); + const identityForPid = options.identityForPid ?? processIdentity; + const parentPidForPid = options.parentPidForPid ?? processParentPid; + let files: string[]; + try { + files = readdirSync(registryDir).filter((file) => file.endsWith(".json")); + } catch { + return []; + } + + const orphans: OwnedFfmpegProcess[] = []; + for (const file of files) { + const path = join(registryDir, file); + try { + const record = JSON.parse(readFileSync(path, "utf8")) as { + version?: unknown; + kind?: unknown; + pid?: unknown; + identity?: unknown; + }; + if ( + record.version !== 1 || + record.kind !== "ffmpeg" || + !Number.isInteger(record.pid) || + (record.pid as number) <= 0 || + typeof record.identity !== "string" + ) { + unlinkSync(path); + continue; + } + const pid = record.pid as number; + if (identityForPid(pid) !== record.identity) { + unlinkSync(path); + continue; + } + if (parentPidForPid(pid) === 1) orphans.push({ pid, identity: record.identity }); + } catch { + try { + unlinkSync(path); + } catch { + // Stale record already removed. + } + } + } + return orphans.sort((a, b) => a.pid - b.pid); +} + +function registerOwnedFfmpeg( + proc: ChildProcess, + registryDir = ownedProcessRegistryDir(), +): string | null { + if (!proc.pid || process.platform === "win32") return null; + const identity = processIdentity(proc.pid); + if (!identity) return null; + try { + mkdirSync(registryDir, { recursive: true, mode: 0o700 }); + const path = join(registryDir, `${proc.pid}.json`); + try { + unlinkSync(path); + } catch { + // No stale record for this reused PID. + } + writeFileSync(path, JSON.stringify({ version: 1, kind: "ffmpeg", pid: proc.pid, identity }), { + flag: "wx", + mode: 0o600, + }); + return path; + } catch { + return null; + } +} + +export function trackChildProcess( + proc: ChildProcess, + options: TrackChildProcessOptions = {}, +): void { + let ownershipPath = + options.kind === "ffmpeg" ? registerOwnedFfmpeg(proc, options.registryDir) : null; + const remove = () => { + tracked.delete(proc); + const path = ownershipPath; + ownershipPath = null; + if (path) { + try { + unlinkSync(path); + } catch { + // Already removed or unavailable. + } + } + }; proc.once("exit", remove); proc.once("close", remove); + if (draining) { + terminateProcesses([proc]); + return; + } + tracked.add(proc); } /** @@ -14,8 +204,20 @@ export function trackChildProcess(proc: ChildProcess): void { * after a short grace period. */ export function killTrackedProcesses(): void { + const processes = [...tracked]; + tracked.clear(); + terminateProcesses(processes); +} + +/** Permanently close this process's child-registration boundary during shutdown. */ +export function beginTrackedProcessDrain(): void { + draining = true; + killTrackedProcesses(); +} + +function terminateProcesses(processes: ChildProcess[]): void { const alive: ChildProcess[] = []; - for (const proc of tracked) { + for (const proc of processes) { if (!proc.killed) { try { proc.kill("SIGTERM"); @@ -25,8 +227,6 @@ export function killTrackedProcesses(): void { } } } - tracked.clear(); - if (alive.length === 0) return; setTimeout(() => { diff --git a/packages/engine/src/utils/runFfmpeg.ts b/packages/engine/src/utils/runFfmpeg.ts index d82772e1ae..30f5a85fb9 100644 --- a/packages/engine/src/utils/runFfmpeg.ts +++ b/packages/engine/src/utils/runFfmpeg.ts @@ -117,7 +117,7 @@ export async function runFfmpeg(args: string[], opts?: RunFfmpegOptions): Promis // shells out dozens of times across parallel workers, which flashes a burst // of windows across the user's desktop. No-op on macOS and Linux. const ffmpeg = spawn(getFfmpegBinary(), args, { windowsHide: true }); - trackChildProcess(ffmpeg); + trackChildProcess(ffmpeg, { kind: "ffmpeg" }); const managed = new ManagedChildProcess(ffmpeg, { signal: opts?.signal, deadlineAtMs: Date.now() + timeout, diff --git a/packages/producer/src/services/audioExtractor.ts b/packages/producer/src/services/audioExtractor.ts index 42208f1a11..3e651f023f 100644 --- a/packages/producer/src/services/audioExtractor.ts +++ b/packages/producer/src/services/audioExtractor.ts @@ -90,7 +90,7 @@ function runFFmpeg(args: string[]): Promise { return new Promise((resolve, reject) => { // See runFfmpeg.ts: keeps a console window off the user's desktop on Windows. const ffmpeg = spawn(getFfmpegBinary(), args, { windowsHide: true }); - trackChildProcess(ffmpeg); + trackChildProcess(ffmpeg, { kind: "ffmpeg" }); let stderr = ""; ffmpeg.stderr.on("data", (data) => { diff --git a/packages/producer/src/services/render/audioPadTrim.ts b/packages/producer/src/services/render/audioPadTrim.ts index e4ecb4e5d1..8c687cfe6d 100644 --- a/packages/producer/src/services/render/audioPadTrim.ts +++ b/packages/producer/src/services/render/audioPadTrim.ts @@ -635,7 +635,7 @@ async function runFfprobeJson(args: string[], signal?: AbortSignal): Promise< throw new Error('[audioPadTrim] ffprobe args must terminate options with "--".'); } const proc = spawn(getFfprobeBinary(), args, { stdio: ["ignore", "pipe", "pipe"] }); - trackChildProcess(proc); + trackChildProcess(proc, { kind: "ffmpeg" }); let stdout = ""; proc.stdout.on("data", (data: Buffer) => { stdout += data.toString(); diff --git a/packages/studio-server/src/createStudioApi.ts b/packages/studio-server/src/createStudioApi.ts index 5976ced33a..5153a5c04d 100644 --- a/packages/studio-server/src/createStudioApi.ts +++ b/packages/studio-server/src/createStudioApi.ts @@ -5,7 +5,7 @@ import { registerStoryboardRoutes } from "./routes/storyboard.js"; import { registerFileRoutes } from "./routes/files.js"; import { registerPreviewRoutes } from "./routes/preview.js"; import { registerLintRoutes } from "./routes/lint.js"; -import { registerRenderRoutes } from "./routes/render.js"; +import { registerRenderRoutes, type RenderRoutesHandle } from "./routes/render.js"; import { registerThumbnailRoutes } from "./routes/thumbnail.js"; import { registerWaveformRoutes } from "./routes/waveform.js"; import { registerFontRoutes } from "./routes/fonts.js"; @@ -14,13 +14,15 @@ import { registerSelectionRoutes } from "./routes/selection.js"; import { registerMediaRoutes } from "./routes/media.js"; import { registerGlobalAssetRoutes } from "./routes/globalAssets.js"; +export type StudioApi = Hono & RenderRoutesHandle; + /** * Create a Hono sub-app with all studio API routes. * * Both the vite dev server and CLI embedded server mount this app * under /api, each providing their own adapter for host-specific behavior. */ -export function createStudioApi(adapter: StudioApiAdapter): Hono { +export function createStudioApi(adapter: StudioApiAdapter): StudioApi { const api = new Hono(); registerProjectRoutes(api, adapter); @@ -28,7 +30,7 @@ export function createStudioApi(adapter: StudioApiAdapter): Hono { registerFileRoutes(api, adapter); registerPreviewRoutes(api, adapter); registerLintRoutes(api, adapter); - registerRenderRoutes(api, adapter); + const renderRoutes = registerRenderRoutes(api, adapter); registerThumbnailRoutes(api, adapter); registerSelectionRoutes(api, adapter); registerMediaRoutes(api, adapter); @@ -37,5 +39,5 @@ export function createStudioApi(adapter: StudioApiAdapter): Hono { registerRegistryRoutes(api, adapter); registerGlobalAssetRoutes(api); - return api; + return Object.assign(api, { dispose: () => renderRoutes.dispose() }); } diff --git a/packages/studio-server/src/index.ts b/packages/studio-server/src/index.ts index 7802ad9758..bd3abdf485 100644 --- a/packages/studio-server/src/index.ts +++ b/packages/studio-server/src/index.ts @@ -1,4 +1,4 @@ -export { createStudioApi } from "./createStudioApi.js"; +export { createStudioApi, type StudioApi } from "./createStudioApi.js"; export { createProjectSignature, affectsProjectSignature } from "./helpers/projectSignature.js"; export type { StudioApiAdapter, diff --git a/packages/studio-server/src/routes/render.test.ts b/packages/studio-server/src/routes/render.test.ts index 40183da199..14fa2d1f36 100644 --- a/packages/studio-server/src/routes/render.test.ts +++ b/packages/studio-server/src/routes/render.test.ts @@ -728,3 +728,112 @@ describe("POST /projects/:id/render — variables forwarding", () => { } }); }); + +describe("render route disposal", () => { + it("cancels and awaits every active render before resolving", async () => { + let finishRender!: () => void; + const completion = new Promise((resolve) => { + finishRender = resolve; + }); + const cancel = vi.fn(); + const spy = vi.fn(); + const { adapter, rendersDir } = createAdapter(spy); + adapter.startRender = (opts) => + ({ + id: opts.jobId, + status: "rendering", + progress: 0, + outputPath: opts.outputPath, + cancel, + completion, + }) as ReturnType; + const app = new Hono(); + const routes = registerRenderRoutes(app, adapter); + try { + const started = await app.request("http://localhost/projects/demo/render", { + method: "POST", + headers: { "content-type": "application/json" }, + body: JSON.stringify({ format: "mp4" }), + }); + expect(started.status).toBe(200); + + let disposed = false; + const disposal = routes.dispose().then(() => { + disposed = true; + }); + await Promise.resolve(); + + expect(cancel).toHaveBeenCalledOnce(); + expect(disposed).toBe(false); + + finishRender(); + await disposal; + expect(disposed).toBe(true); + } finally { + rmSync(rendersDir, { recursive: true, force: true }); + } + }); + + it("rejects a render request after disposal begins", async () => { + const spy = vi.fn(); + const { adapter, rendersDir } = createAdapter(spy); + const app = new Hono(); + const routes = registerRenderRoutes(app, adapter); + try { + await routes.dispose(); + const response = await app.request("http://localhost/projects/demo/render", { + method: "POST", + headers: { "content-type": "application/json" }, + body: JSON.stringify({ format: "mp4" }), + }); + + expect(response.status).toBe(503); + expect(spy).not.toHaveBeenCalled(); + } finally { + rmSync(rendersDir, { recursive: true, force: true }); + } + }); + + it("awaits a cancelled job whose resource cleanup is still running", async () => { + let finishRender!: () => void; + const completion = new Promise((resolve) => { + finishRender = resolve; + }); + const cancel = vi.fn(); + const spy = vi.fn(); + const { adapter, rendersDir } = createAdapter(spy); + adapter.startRender = (opts) => + ({ + id: opts.jobId, + status: "cancelled", + progress: 0, + outputPath: opts.outputPath, + cancel, + completion, + }) as ReturnType; + const app = new Hono(); + const routes = registerRenderRoutes(app, adapter); + try { + const started = await app.request("http://localhost/projects/demo/render", { + method: "POST", + headers: { "content-type": "application/json" }, + body: JSON.stringify({ format: "mp4" }), + }); + expect(started.status).toBe(200); + + let disposed = false; + const disposal = routes.dispose().then(() => { + disposed = true; + }); + await new Promise((resolve) => setTimeout(resolve, 0)); + + expect(cancel).not.toHaveBeenCalled(); + expect(disposed).toBe(false); + finishRender(); + await disposal; + expect(disposed).toBe(true); + } finally { + rmSync(rendersDir, { recursive: true, force: true }); + } + }); +}); diff --git a/packages/studio-server/src/routes/render.ts b/packages/studio-server/src/routes/render.ts index bc0c6de270..90fc65bc72 100644 --- a/packages/studio-server/src/routes/render.ts +++ b/packages/studio-server/src/routes/render.ts @@ -10,14 +10,26 @@ import { isVariablesPayload, VARIABLES_PAYLOAD_ERROR } from "../helpers/variable const VALID_RESOLUTIONS = new Set(VALID_CANVAS_RESOLUTIONS); -export function registerRenderRoutes(api: Hono, adapter: StudioApiAdapter): void { +export interface RenderRoutesHandle { + dispose(): Promise; +} + +interface StoredRenderJob { + state: RenderJobState; + createdAt: number; + finishedAt?: number; + pendingCompletion?: Promise; +} + +export function registerRenderRoutes(api: Hono, adapter: StudioApiAdapter): RenderRoutesHandle { // Scoped job store — not shared across createStudioApi() calls - const renderJobs = new Map(); + const renderJobs = new Map(); // TTL cleanup for completed jobs (5 minutes) const TTL_MS = 300_000; const CLEANUP_INTERVAL_MS = 60_000; let cleanupTimer: ReturnType | null = null; + let disposalPromise: Promise | null = null; const cleanupEnabled = () => typeof process !== "undefined" && @@ -27,9 +39,9 @@ export function registerRenderRoutes(api: Hono, adapter: StudioApiAdapter): void const cleanupFinishedJobs = () => { const now = Date.now(); for (const [key, job] of renderJobs) { - if (job.status !== "rendering" && now - job.createdAt > TTL_MS) { - renderJobs.delete(key); - } + if (job.state.status === "rendering" || job.pendingCompletion) continue; + job.finishedAt ??= now; + if (now - job.finishedAt > TTL_MS) renderJobs.delete(key); } if (renderJobs.size === 0 && cleanupTimer) { clearInterval(cleanupTimer); @@ -45,10 +57,17 @@ export function registerRenderRoutes(api: Hono, adapter: StudioApiAdapter): void } }; + const cancelRenderJob = (job: RenderJobState) => { + if (job.status !== "rendering") return; + job.status = "cancelled"; + job.cancel?.(); + }; + ensureCleanupTimer(); // Start a render api.post("/projects/:id/render", async (c) => { + if (disposalPromise) return c.json({ error: "studio is shutting down" }, 503); const project = await adapter.resolveProject(c.req.param("id")); if (!project) return c.json({ error: "not found" }, 404); @@ -118,6 +137,7 @@ export function registerRenderRoutes(api: Hono, adapter: StudioApiAdapter): void const ext = FORMAT_EXT[format] ?? ".mp4"; const outputPath = join(rendersDir, `${jobId}${ext}`); + if (disposalPromise) return c.json({ error: "studio is shutting down" }, 503); const jobState = adapter.startRender({ project, outputPath, @@ -132,8 +152,19 @@ export function registerRenderRoutes(api: Hono, adapter: StudioApiAdapter): void typeof body.telemetryDistinctId === "string" ? body.telemetryDistinctId : undefined, telemetryOptOut: body.telemetryOptOut === true, }); - (jobState as RenderJobState & { createdAt: number }).createdAt = Date.now(); - renderJobs.set(jobId, jobState as RenderJobState & { createdAt: number }); + const stored: StoredRenderJob = { state: jobState, createdAt: Date.now() }; + if (jobState.completion) { + const completion = jobState.completion; + const clearCompletion = () => { + if (stored.pendingCompletion === completion) { + stored.pendingCompletion = undefined; + stored.finishedAt = Date.now(); + } + }; + stored.pendingCompletion = completion; + void completion.then(clearCompletion, clearCompletion); + } + renderJobs.set(jobId, stored); ensureCleanupTimer(); @@ -143,12 +174,12 @@ export function registerRenderRoutes(api: Hono, adapter: StudioApiAdapter): void // SSE progress stream api.get("/render/:jobId/progress", (c) => { const { jobId } = c.req.param(); - const job = renderJobs.get(jobId); + const job = renderJobs.get(jobId)?.state; if (!job) return c.json({ error: "not found" }, 404); return streamSSE(c, async (stream) => { while (true) { - const current = renderJobs.get(jobId); + const current = renderJobs.get(jobId)?.state; if (!current) break; await stream.writeSSE({ event: "progress", @@ -169,12 +200,9 @@ export function registerRenderRoutes(api: Hono, adapter: StudioApiAdapter): void // SSE stream terminates) and invokes the adapter's abort hook when present. api.post("/render/:jobId/cancel", (c) => { const { jobId } = c.req.param(); - const job = renderJobs.get(jobId); + const job = renderJobs.get(jobId)?.state; if (!job) return c.json({ error: "not found" }, 404); - if (job.status === "rendering") { - job.status = "cancelled"; - job.cancel?.(); - } + cancelRenderJob(job); return c.json({ status: job.status }); }); @@ -194,7 +222,7 @@ export function registerRenderRoutes(api: Hono, adapter: StudioApiAdapter): void // fallow-ignore-next-line code-duplication api.get("/render/:jobId/view", (c) => { const { jobId } = c.req.param(); - const job = renderJobs.get(jobId); + const job = renderJobs.get(jobId)?.state; if (!job?.outputPath || !existsSync(job.outputPath)) { return c.json({ error: "not found" }, 404); } @@ -215,7 +243,7 @@ export function registerRenderRoutes(api: Hono, adapter: StudioApiAdapter): void // fallow-ignore-next-line code-duplication api.get("/render/:jobId/download", (c) => { const { jobId } = c.req.param(); - const job = renderJobs.get(jobId); + const job = renderJobs.get(jobId)?.state; if (!job?.outputPath || !existsSync(job.outputPath)) { return c.json({ error: "not found" }, 404); } @@ -233,7 +261,7 @@ export function registerRenderRoutes(api: Hono, adapter: StudioApiAdapter): void // Delete render api.delete("/render/:jobId", (c) => { const { jobId } = c.req.param(); - for (const [, state] of renderJobs) { + for (const { state } of renderJobs.values()) { if (state.id === jobId && state.outputPath) { const dir = state.outputPath.replace(/\/[^/]+$/, ""); for (const ext of [".mp4", ".webm", ".mov", ".meta.json"]) { @@ -317,14 +345,48 @@ export function registerRenderRoutes(api: Hono, adapter: StudioApiAdapter): void for (const file of files) { if (!renderJobs.has(file.id)) { renderJobs.set(file.id, { - id: file.id, - status: file.status, - progress: 100, - outputPath: join(rendersDir, file.filename), + state: { + id: file.id, + status: file.status, + progress: 100, + outputPath: join(rendersDir, file.filename), + }, createdAt: file.createdAt, - } as RenderJobState & { createdAt: number }); + finishedAt: file.createdAt, + }); } } return c.json({ renders: files }); }); + + const performDispose = async () => { + if (cleanupTimer) { + clearInterval(cleanupTimer); + cleanupTimer = null; + } + const jobs = [...renderJobs.values()]; + const active = jobs.filter((job) => job.state.status === "rendering"); + for (const job of active) { + try { + cancelRenderJob(job.state); + } catch { + // Continue cancelling and awaiting the remaining owned jobs. + } + } + await Promise.allSettled( + jobs + .map((job) => job.pendingCompletion) + .filter((completion): completion is Promise => completion !== undefined), + ); + }; + + return { + dispose() { + if (!disposalPromise) { + // Publish idempotency before a cancel hook can re-enter disposal. + disposalPromise = Promise.resolve().then(performDispose); + } + return disposalPromise; + }, + }; } diff --git a/packages/studio-server/src/types.ts b/packages/studio-server/src/types.ts index 1344b90701..982da92dad 100644 --- a/packages/studio-server/src/types.ts +++ b/packages/studio-server/src/types.ts @@ -23,6 +23,8 @@ export interface RenderJobState { * route still marks the job cancelled so the SSE stream terminates). */ cancel?: () => void; + /** Resolves after the adapter has released all resources owned by this render. */ + completion?: Promise; } export interface MediaProcessingJobState {