From 97d59406f6544a6a9be7065dce449daa6f6d41d8 Mon Sep 17 00:00:00 2001 From: Duy Nguyen Date: Thu, 8 Oct 2026 05:07:57 +0700 Subject: [PATCH] fix(runtime): retry worker and widget-dev cleanup with the promise form of rm On Windows, rmSync reports a held directory as EBUSY or EPERM at once and never runs its maxRetries, so the worker's run directory and a superseded widget-dev snapshot were left behind whenever something still held them. Both callers are already async, so they now await rm from node:fs/promises, which does retry, without blocking the event loop. The widget-dev prune runs on the session's chain, so nothing else of that session interleaves. Tests hold each directory with a child process's working directory until 300 ms after its removal starts; they fail against the old rmSync calls. --- .../src/application/widget-dev-sessions.ts | 16 ++-- apps/runtime/src/worker-process.ts | 9 ++- apps/runtime/test/hold-directory.ts | 48 ++++++++++++ apps/runtime/test/widget-dev-sessions.spec.ts | 53 ++++++++++++- apps/runtime/test/worker-process.spec.ts | 75 ++++++++++++++++++- 5 files changed, 188 insertions(+), 13 deletions(-) create mode 100644 apps/runtime/test/hold-directory.ts diff --git a/apps/runtime/src/application/widget-dev-sessions.ts b/apps/runtime/src/application/widget-dev-sessions.ts index edd497bd1..644bf5e66 100644 --- a/apps/runtime/src/application/widget-dev-sessions.ts +++ b/apps/runtime/src/application/widget-dev-sessions.ts @@ -1,4 +1,5 @@ -import { mkdirSync, realpathSync, rmSync, statSync } from "node:fs"; +import { mkdirSync, realpathSync, statSync } from "node:fs"; +import { rm } from "node:fs/promises"; import { isAbsolute, join, resolve } from "node:path"; import { @@ -393,7 +394,7 @@ export function createWidgetDevSessions( }); if (session !== undefined && active?.snapshotDigest === asked.generation.digest) { session.baseline = manifestAt(asked.listing.source.kind === "local" ? asked.listing.source.path : ""); - prune(sessionId); + await prune(sessionId); } } else { const denied = approval?.decision === "denied"; @@ -556,7 +557,7 @@ export function createWidgetDevSessions( if (session !== undefined) session.baseline = latest.manifest; // Only when nobody was asked: a build the person approved in the inbox was shown to them with what it adds. if (granted === undefined) sayWidened(stored, latest); - prune(sessionId); + await prune(sessionId); return; } if (outcome.kind === "approval-required") { @@ -582,9 +583,10 @@ export function createWidgetDevSessions( * Remove what superseded generations of a session left behind: their generation records, except the newest superseded * one (the generation a rollback returns to), and their snapshots in the package cache, except the ones something * still runs, waits on or can roll back to. Only what this session made is touched. A snapshot that cannot be removed - * now (a file still open on Windows) is kept on the list and tried again after the next install. + * now (a file still open on Windows past the retries) is kept on the list and tried again after the next install. It + * runs on the session's chain, so nothing else of the session interleaves with it. */ - const prune = (sessionId: string): void => { + const prune = async (sessionId: string): Promise => { const stored = read(sessionId); if (stored === undefined) return; const made = new Set(stored.snapshots ?? []); @@ -642,7 +644,9 @@ export function createWidgetDevSessions( continue; } try { - rmSync(path, { recursive: true, force: true, maxRetries: 2 }); + // The promise form on purpose: on Windows `rmSync` reports a held file as `EBUSY` or `EPERM` at once and never + // runs its retries. These wait up to about 1.5 s on this session's chain, never on the event loop. + await rm(path, { recursive: true, force: true, maxRetries: 5, retryDelay: 100 }); removed.push(digest); } catch (cause) { process.stderr.write(`widget dev: could not remove the superseded snapshot ${digest} yet: ${messageOf(cause)}\n`); diff --git a/apps/runtime/src/worker-process.ts b/apps/runtime/src/worker-process.ts index a042f6ed8..4e32c0163 100644 --- a/apps/runtime/src/worker-process.ts +++ b/apps/runtime/src/worker-process.ts @@ -1,5 +1,6 @@ import { spawn, type ChildProcess } from "node:child_process"; -import { mkdtempSync, rmSync, writeFileSync } from "node:fs"; +import { mkdtempSync, writeFileSync } from "node:fs"; +import { rm } from "node:fs/promises"; import { tmpdir } from "node:os"; import { isAbsolute, join, resolve as resolvePath } from "node:path"; import { fileURLToPath } from "node:url"; @@ -340,8 +341,10 @@ export async function runWorkerProcess(options: WorkerProcessOptions): Promise boolean; + /** Resolves once the holder runs in the directory. */ + readonly ready: Promise; + /** Let go of the directory; resolves once the holder has exited. */ + readonly release: () => Promise; +} + +/** + * Hold a directory the way Windows holds one for a moment after a process exits, or while an antivirus or indexer + * looks at it. + * + * A child process runs with the directory as its working directory. On Windows, removing the directory then fails with + * `EBUSY` or `EPERM` until the child exits; elsewhere the hold changes nothing. The hold starts once the child is + * running, not when it is spawned, so a test waits for `ready` (or checks `holding`) before it relies on it. + */ +export function holdDirectory(path: string): DirectoryHold { + const holder = spawn( + process.execPath, + ["-e", "process.stdout.write('ready'); process.stdin.resume(); process.stdin.on('end', () => process.exit(0));"], + { cwd: path, stdio: ["pipe", "pipe", "ignore"], windowsHide: true }, + ); + let running = false; + const ready = new Promise((resolve, reject) => { + holder.stdout?.once("data", () => { + running = true; + resolve(); + }); + holder.once("error", reject); + }); + // A caller that only checks `holding` never awaits `ready`; a holder that failed to start is then simply not holding. + ready.catch(() => undefined); + const exited = new Promise((resolve) => { + holder.once("exit", () => resolve()); + holder.once("error", () => resolve()); + }); + return { + holding: () => running, + ready, + release: async () => { + holder.stdin?.end(); + await exited; + }, + }; +} diff --git a/apps/runtime/test/widget-dev-sessions.spec.ts b/apps/runtime/test/widget-dev-sessions.spec.ts index 60bec7e9d..1c47e2f4b 100644 --- a/apps/runtime/test/widget-dev-sessions.spec.ts +++ b/apps/runtime/test/widget-dev-sessions.spec.ts @@ -1,4 +1,5 @@ import type * as fs from "node:fs"; +import type * as fsPromises from "node:fs/promises"; import { existsSync, mkdirSync, mkdtempSync, readFileSync, readdirSync, rmSync, writeFileSync } from "node:fs"; import { createServer, type Server } from "node:http"; import { tmpdir } from "node:os"; @@ -31,10 +32,14 @@ import { createDevelopWidgetTool } from "../src/develop-widget-tool.ts"; import { hostText } from "../src/host-text.ts"; import { handleRequest, type GatewayDeps, type GatewayResponse } from "../src/gateway.ts"; import { bootNodeServices, type NodeServices } from "../src/services.ts"; +import { holdDirectory } from "./hold-directory.ts"; /** A folder whose `stat` fails with `EPERM`, as an antivirus or indexer holding it on Windows makes it fail. */ const statFailure = vi.hoisted(() => ({ path: undefined as string | undefined })); +/** Seen as the runtime starts to remove a path, in either form, before the removal itself runs. */ +const removal = vi.hoisted(() => ({ starting: undefined as ((path: string) => void) | undefined })); + vi.mock("node:fs", async (importOriginal) => { const actual = await importOriginal(); const statSync = ((path: fs.PathLike, options?: fs.StatSyncOptions) => { @@ -43,7 +48,20 @@ vi.mock("node:fs", async (importOriginal) => { } return actual.statSync(path, options); }) as typeof actual.statSync; - return { ...actual, statSync, default: { ...actual, statSync } }; + const rmSync = ((path: fs.PathLike, options?: fs.RmOptions) => { + removal.starting?.(resolve(String(path))); + actual.rmSync(path, options); + }) as typeof actual.rmSync; + return { ...actual, statSync, rmSync, default: { ...actual, statSync, rmSync } }; +}); + +vi.mock("node:fs/promises", async (importOriginal) => { + const actual = await importOriginal(); + const rm = (async (path: fs.PathLike, options?: fs.RmOptions) => { + removal.starting?.(resolve(String(path))); + await actual.rm(path, options); + }) as typeof actual.rm; + return { ...actual, rm, default: { ...actual, rm } }; }); /** @@ -584,6 +602,39 @@ describe("a widget dev session", () => { expect(readDevSessions(join(dir, "node"))[0]?.snapshots).toHaveLength(2); }); + it("removes a superseded snapshot that is held for a moment, as a file still open on Windows holds it", async () => { + const started = session(await call("POST", "/widget-dev/sessions", { root })); + const local = join(dir, "node", "package-cache", "local"); + const [first] = readdirSync(local); + if (first === undefined) throw new Error("the first build left no snapshot"); + const firstPath = resolve(local, first); + // On Windows a directory that is someone's working directory cannot be removed until they let go; elsewhere the + // hold changes nothing. It ends 300 ms after the removal starts, so only a removal that really retries finds it free. + const hold = holdDirectory(firstPath); + await hold.ready; + let released: Promise | undefined; + removal.starting = (path) => { + if (path !== firstPath || released !== undefined) return; + released = new Promise((done) => setTimeout(done, 300)).then(() => hold.release()); + }; + try { + // The second build supersedes the first, which stays as the generation a rollback returns to; the third prunes it. + for (const text of ["second", "third"]) { + writePackage(`

${text}

\n`); + expect(session(await call("POST", `/widget-dev/sessions/${started.sessionId}/rebuild`)).running?.generation).toBeGreaterThan(1); + } + expect(released).toBeDefined(); + expect(existsSync(firstPath)).toBe(false); + // Removed, so no longer on the session's list of snapshots to try again. + const listed = readDevSessions(join(dir, "node"))[0]?.snapshots ?? []; + expect(listed).toHaveLength(2); + expect(listed.some((digest) => digest.endsWith(first))).toBe(false); + } finally { + removal.starting = undefined; + await (released ?? hold.release()); + } + }); + it("forgets the oldest stopped sessions rather than failing when the store is full", async () => { const at = new Date(Date.UTC(2026, 0, 1)).toISOString(); writeDevSessions( diff --git a/apps/runtime/test/worker-process.spec.ts b/apps/runtime/test/worker-process.spec.ts index 0c508d2e7..cf74a901f 100644 --- a/apps/runtime/test/worker-process.spec.ts +++ b/apps/runtime/test/worker-process.spec.ts @@ -1,11 +1,47 @@ -import { mkdtempSync, rmSync, writeFileSync } from "node:fs"; +import type * as fs from "node:fs"; +import type * as fsPromises from "node:fs/promises"; +import { existsSync, mkdtempSync, rmSync, writeFileSync } from "node:fs"; import { tmpdir } from "node:os"; -import { join } from "node:path"; +import { basename, join } from "node:path"; import type { WorkerBriefEnvelope } from "@clarkcant/app-worker"; -import { afterAll, beforeAll, describe, expect, it } from "vitest"; +import { afterAll, beforeAll, describe, expect, it, vi } from "vitest"; import { runWorkerProcess, workerEnvironment } from "../src/worker-process.ts"; +import { holdDirectory, type DirectoryHold } from "./hold-directory.ts"; + +/** + * Seen as the runtime makes a temporary directory (`created`) and as it starts to remove one, in either form, before + * the removal itself runs (`removing`). A test uses them to hold the worker's directory past the worker's exit, as a + * worker that was just stopped holds it on Windows. + */ +const watched = vi.hoisted(() => ({ + created: undefined as ((path: string) => void) | undefined, + removing: undefined as ((path: string) => void) | undefined, +})); + +vi.mock("node:fs", async (importOriginal) => { + const actual = await importOriginal(); + const mkdtempSync = ((prefix: string, options?: fs.EncodingOption) => { + const made = actual.mkdtempSync(prefix, options); + watched.created?.(String(made)); + return made; + }) as typeof actual.mkdtempSync; + const rmSync = ((path: fs.PathLike, options?: fs.RmOptions) => { + watched.removing?.(String(path)); + actual.rmSync(path, options); + }) as typeof actual.rmSync; + return { ...actual, mkdtempSync, rmSync, default: { ...actual, mkdtempSync, rmSync } }; +}); + +vi.mock("node:fs/promises", async (importOriginal) => { + const actual = await importOriginal(); + const rm = (async (path: fs.PathLike, options?: fs.RmOptions) => { + watched.removing?.(String(path)); + await actual.rm(path, options); + }) as typeof actual.rm; + return { ...actual, rm, default: { ...actual, rm } }; +}); /** * Dispatching a worker as a process. @@ -93,6 +129,39 @@ describe("the runtime starts a worker of its own", () => { ).rejects.toThrow(/no model was given to this worker/); }); + it("removes its run directory even when the directory is still held for a moment after the worker exits", async () => { + // On Windows a directory that is someone's working directory cannot be removed until they let go; elsewhere the + // hold changes nothing. The hold starts with the directory and ends 300 ms after its removal starts, so only a + // removal that really retries finds the directory free. + let directory: string | undefined; + let hold: DirectoryHold | undefined; + let heldAtRemoval = false; + let released: Promise | undefined; + watched.created = (path) => { + if (directory !== undefined || !basename(path).startsWith("clarkcant-worker-run-")) return; + directory = path; + hold = holdDirectory(path); + }; + watched.removing = (path) => { + if (path !== directory || hold === undefined || released !== undefined) return; + heldAtRemoval = hold.holding(); + const holding = hold; + released = new Promise((resolve) => setTimeout(resolve, 300)).then(() => holding.release()); + }; + try { + const result = await runWorkerProcess({ nodeId: "node_test", brief: brief() }); + expect(result.record.runId).toBe("run_1"); + // The hold was in place when the removal started; otherwise this test would prove nothing. + expect(heldAtRemoval).toBe(true); + expect(existsSync(directory ?? "")).toBe(false); + } finally { + watched.created = undefined; + watched.removing = undefined; + await (released ?? hold?.release()); + if (directory !== undefined) rmSync(directory, { recursive: true, force: true }); + } + }); + it("starts a real-model worker without any provider key in its environment, whatever this process holds", () => { const source = { PATH: "/usr/bin",