import { execFile as execFileCallback, spawn, type ChildProcess } from "node:child_process"; import { createWriteStream } from "node:fs"; import { mkdtemp, rm, writeFile } from "node:fs/promises"; import net from "node:net"; import os from "node:os"; import path from "node:path"; import { setTimeout as sleep } from "node:timers/promises"; import { promisify } from "node:util"; const execFile = promisify(execFileCallback); export const repoRoot = path.resolve(import.meta.dirname, "../.."); const viteBinPath = path.join(repoRoot, "node_modules", "vite", "bin", "vite.js"); const wranglerBinPath = path.join(repoRoot, "node_modules", "wrangler", "bin", "wrangler.js"); const d1DatabaseName = "dab-control-plane"; const startAttempts = 5; const readyTimeoutMs = 30_000; export interface IsolatedWorkerOptions { name: string; tempPrefix: string; preserveEnv: string; oidcIssuer: string; d1DatabaseId: string; r2BucketName: string; resourceSuffix: string; devVars: Record; } export interface IsolatedWorker { port: number; tempDir: string; } export async function runIsolatedWorker( options: IsolatedWorkerOptions, callback: (worker: IsolatedWorker) => Promise, ): Promise { const tempDir = await mkdtemp(path.join(os.tmpdir(), options.tempPrefix)); const wranglerConfigPath = path.join(tempDir, "wrangler.json"); const devVarsPath = path.join(tempDir, ".dev.vars"); const preserve = process.env[options.preserveEnv] === "1"; let started: { child: ChildProcess; port: number } | undefined; await writeFile(wranglerConfigPath, `${localWranglerConfig(options)}\n`); await writeFile(devVarsPath, `${devVars(options.devVars)}\n`); await applyD1Migrations(tempDir, wranglerConfigPath); try { started = await startWorker(tempDir, wranglerConfigPath); await callback({ port: started.port, tempDir }); } finally { if (started) await stop(started.child); if (!preserve) { await rm(devVarsPath, { force: true }); await rm(tempDir, { recursive: true, force: true }); } else { console.log(`preserved ${tempDir}`); console.log(`preserved ${devVarsPath}`); console.log(`preserved ${wranglerConfigPath}`); } } } function devVars(values: Record): string { return Object.entries(values) .map(([key, value]) => `${key}=${value}`) .join("\n"); } async function startWorker( tempDir: string, wranglerConfigPath: string, ): Promise<{ child: ChildProcess; port: number }> { let lastError: unknown; for (let attempt = 1; attempt <= startAttempts; attempt += 1) { const port = await allocatePort(); try { const child = await startViteDev(tempDir, port, wranglerConfigPath); return { child, port }; } catch (error) { lastError = error; if (!isPortInUse(error) || attempt === startAttempts) throw error; } } throw lastError instanceof Error ? lastError : new Error("Unable to start Vite dev server"); } function allocatePort(): Promise { return new Promise((resolve, reject) => { const server = net.createServer(); server.on("error", reject); server.listen(0, "127.0.0.1", () => { const address = server.address(); if (!address || typeof address === "string") { reject(new Error("Unable to allocate loopback port")); return; } const port = address.port; server.close(() => resolve(port)); }); }); } async function startViteDev(tempDir: string, port: number, wranglerConfigPath: string): Promise { const stdout = createWriteStream(path.join(tempDir, "vite.stdout.log")); const stderr = createWriteStream(path.join(tempDir, "vite.stderr.log")); const child = spawn( process.execPath, [viteBinPath, "dev", "--host", "127.0.0.1", "--port", String(port), "--strictPort"], { cwd: repoRoot, env: { ...process.env, DAB_PERSIST_STATE_PATH: tempDir, DAB_WRANGLER_CONFIG_PATH: wranglerConfigPath, }, stdio: ["ignore", "pipe", "pipe"], detached: process.platform !== "win32", }, ); child.stdout?.pipe(stdout); child.stderr?.pipe(stderr); child.once("close", () => { stdout.close(); stderr.close(); }); try { await waitForHealth(port, child); return child; } catch (error) { await stop(child); throw error; } } async function waitForHealth(port: number, child: ChildProcess): Promise { const deadline = Date.now() + readyTimeoutMs; let lastError: unknown; while (Date.now() < deadline) { if (child.exitCode !== null) throw new Error(`Vite exited before readiness with code ${child.exitCode}`); try { const response = await fetch(`http://dab.localhost:${port}/healthz`); if (response.status === 200) return; lastError = new Error(`healthz returned ${response.status}`); } catch (error) { lastError = error; } await sleep(250); } throw lastError instanceof Error ? lastError : new Error("Timed out waiting for /healthz"); } async function applyD1Migrations(persistTo: string, configPath: string): Promise { await execFile( process.execPath, [ wranglerBinPath, "d1", "migrations", "apply", d1DatabaseName, "--local", "--persist-to", persistTo, "--config", configPath, ], { cwd: repoRoot, env: { ...process.env, CI: "1", NO_D1_WARNING: "true" }, timeout: 180_000, maxBuffer: 10 * 1024 * 1024, }, ); } async function stop(child: ChildProcess): Promise { if (child.exitCode !== null) return; const kill = (signal: NodeJS.Signals): void => { if (child.pid === undefined) return; if (process.platform === "win32") { child.kill(signal); return; } process.kill(-child.pid, signal); }; kill("SIGTERM"); await sleep(1_000); if (child.exitCode === null) kill("SIGKILL"); } function isPortInUse(error: unknown): boolean { const message = error instanceof Error ? error.message : String(error); return message.includes("is already in use") || message.includes("EADDRINUSE"); } function localWranglerConfig(options: IsolatedWorkerOptions): string { return JSON.stringify( { name: options.name, main: path.join(repoRoot, "src/worker/index.ts"), compatibility_date: "2026-05-20", compatibility_flags: ["nodejs_als"], vars: { LOG_LEVEL: "warn", CONTROL_PLANE_HOST: "dab.localhost", SUBJECT_HOST_SUFFIX: ".dab.localhost", TESSERA_OIDC_ISSUER: options.oidcIssuer, MAX_FILE_BYTES: "100000000", MAX_XML_BODY_BYTES: "1048576", MAX_CALENDAR_OBJECT_BYTES: "1048576", MAX_ADDRESS_OBJECT_BYTES: "1048576", MAX_DAV_DEPTH_INFINITY_NODES: "10000", MAX_REPORT_RESULTS: "5000", MAX_RECURRENCE_YEARS: "5", MAX_RECURRENCE_INSTANCES: "10000", DAB_ENABLE_TEST_FIXTURES: "1", }, secrets: { required: ["TESSERA_OIDC_CLIENT_ID", "TESSERA_OIDC_CLIENT_SECRET", "DAB_SESSION_SECRET"], }, d1_databases: [ { binding: "DAV_CONTROL_PLANE", database_name: d1DatabaseName, database_id: options.d1DatabaseId, migrations_dir: path.join(repoRoot, "drizzle/d1"), }, ], r2_buckets: [{ binding: "FILE_BLOBS", bucket_name: options.r2BucketName }], durable_objects: { bindings: [ { name: "AUTH", class_name: "AuthObject" }, { name: "FILE_DAV", class_name: "FileDavObject" }, { name: "CAL_DAV", class_name: "CalDavObject" }, { name: "CARD_DAV", class_name: "CardDavObject" }, ], }, migrations: [ { tag: "v1", new_sqlite_classes: ["AuthObject", "FileDavObject", "CalDavObject", "CardDavObject"], }, ], queues: { producers: [ { binding: "BLOB_GC", queue: `dab-blob-gc-${options.resourceSuffix}` }, { binding: "REPAIR_JOBS", queue: `dab-repair-jobs-${options.resourceSuffix}` }, ], consumers: [ { queue: `dab-blob-gc-${options.resourceSuffix}`, max_batch_size: 10, max_batch_timeout: 5 }, { queue: `dab-repair-jobs-${options.resourceSuffix}`, max_batch_size: 5, max_batch_timeout: 5 }, ], }, ratelimits: [ { name: "RL_AUTH", namespace_id: "83800", simple: { limit: 10, period: 60 } }, { name: "RL_DAV_AUTH", namespace_id: "83801", simple: { limit: 30, period: 60 } }, { name: "RL_REPORT", namespace_id: "83802", simple: { limit: 60, period: 60 } }, ], }, null, 2, ); }