Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
16 changes: 10 additions & 6 deletions apps/runtime/src/application/widget-dev-sessions.ts
Original file line number Diff line number Diff line change
@@ -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 { homedir } from "node:os";
import { dirname, isAbsolute, join, resolve } from "node:path";

Expand Down Expand Up @@ -462,7 +463,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";
Expand Down Expand Up @@ -625,7 +626,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") {
Expand All @@ -651,9 +652,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<void> => {
const stored = read(sessionId);
if (stored === undefined) return;
const made = new Set(stored.snapshots ?? []);
Expand Down Expand Up @@ -711,7 +713,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`);
Expand Down
9 changes: 6 additions & 3 deletions apps/runtime/src/worker-process.ts
Original file line number Diff line number Diff line change
@@ -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";
Expand Down Expand Up @@ -340,8 +341,10 @@ export async function runWorkerProcess(options: WorkerProcessOptions): Promise<W
};
} finally {
try {
// Retried, because on Windows a worker being stopped still holds its working directory for a moment.
rmSync(directory, { recursive: true, force: true, maxRetries: 5, retryDelay: 100 });
// Retried, because on Windows a worker being stopped still holds its working directory for a moment. The promise
// form on purpose: on Windows `rmSync` reports a held directory as `EBUSY` or `EPERM` at once and never runs its
// retries, while `rm` waits `retryDelay` longer after each failed attempt (about 1.5 s over five), off the event loop.
await rm(directory, { recursive: true, force: true, maxRetries: 5, retryDelay: 100 });
} catch {
// A temporary directory left behind (it holds the brief, never a key) is not worth replacing the run's own result
// or error with.
Expand Down
48 changes: 48 additions & 0 deletions apps/runtime/test/hold-directory.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,48 @@
import { spawn } from "node:child_process";

export interface DirectoryHold {
/** Whether the holder runs in the directory yet; until it does, nothing is held. */
readonly holding: () => boolean;
/** Resolves once the holder runs in the directory. */
readonly ready: Promise<void>;
/** Let go of the directory; resolves once the holder has exited. */
readonly release: () => Promise<void>;
}

/**
* 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<void>((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<void>((resolve) => {
holder.once("exit", () => resolve());
holder.once("error", () => resolve());
});
return {
holding: () => running,
ready,
release: async () => {
holder.stdin?.end();
await exited;
},
};
}
53 changes: 52 additions & 1 deletion apps/runtime/test/widget-dev-sessions.spec.ts
Original file line number Diff line number Diff line change
@@ -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, realpathSync, renameSync, rmSync, symlinkSync, writeFileSync } from "node:fs";
import { createServer, type Server } from "node:http";
import type * as os from "node:os";
Expand Down Expand Up @@ -36,10 +37,14 @@ 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 { removeTestDirectory } from "../../../tools/test-cleanup.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<typeof fs>();
const statSync = ((path: fs.PathLike, options?: fs.StatSyncOptions) => {
Expand All @@ -48,7 +53,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<typeof fsPromises>();
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 } };
});

/** The home folder the node sees, when a test needs it to be one of the test's own folders. */
Expand Down Expand Up @@ -883,6 +901,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<void> | undefined;
removal.starting = (path) => {
if (path !== firstPath || released !== undefined) return;
released = new Promise<void>((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(`<!doctype html><p>${text}</p>\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(
Expand Down
75 changes: 72 additions & 3 deletions apps/runtime/test/worker-process.spec.ts
Original file line number Diff line number Diff line change
@@ -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<typeof fs>();
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<typeof fsPromises>();
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.
Expand Down Expand Up @@ -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<void> | 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<void>((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",
Expand Down
Loading