feat(studio): run nitro scans in the background with cancel-re-attach and nightly auto-clean
This commit is contained in:
1 parent
9d571e0c29
commit
6c53f4680c
14 files changed
+1547
-303
No files matched your search
@@ -88,10 +88,20 @@ vi.mock("node:fs", async (orig) => {
|
||||
import {
|
||||
autoCleanFakeNitros,
|
||||
deleteNitroCleanupFiles,
|
||||
getPersistedCleanupResult,
|
||||
type NitroRepairProgress,
|
||||
refreshScanCacheAfterMutation,
|
||||
repairBrokenNitros,
|
||||
scanFakeBrokenNitros,
|
||||
} from "@/lib/services/nitro-cleanup";
|
||||
import {
|
||||
cancelActiveScan,
|
||||
cancelScanSession,
|
||||
NitroScanConflictError,
|
||||
readScanSession,
|
||||
runScanInBackground,
|
||||
writeScanSession,
|
||||
} from "@/lib/services/nitro-scan-session";
|
||||
|
||||
beforeEach(() => {
|
||||
vi.clearAllMocks();
|
||||
@@ -128,6 +138,49 @@ beforeEach(() => {
|
||||
statFn.mockImplementation(async () => ({ size: 512, mtimeMs: 0 }));
|
||||
});
|
||||
|
||||
/**
|
||||
* File-backed state for the mocked `node:fs` promises: cache/session/history
|
||||
* writes land in the map and subsequent reads/exists checks consult it, while
|
||||
* binary reads (nitro bundles) fall back to placeholder bytes.
|
||||
*/
|
||||
function installFileState(
|
||||
initial: Record<string, string> = {},
|
||||
): Map<string, string> {
|
||||
const files = new Map<string, string>(Object.entries(initial));
|
||||
readFileFn.mockImplementation(async (filePath: string) => {
|
||||
const stored = files.get(filePath);
|
||||
if (stored !== undefined) return stored;
|
||||
return Buffer.from("nitro-data");
|
||||
});
|
||||
writeFileFn.mockImplementation(async (filePath: string, data: unknown) => {
|
||||
files.set(
|
||||
filePath,
|
||||
Buffer.isBuffer(data) ? data.toString("utf-8") : String(data),
|
||||
);
|
||||
});
|
||||
existsFn.mockImplementation((filePath: string) => files.has(filePath));
|
||||
return files;
|
||||
}
|
||||
|
||||
async function waitForSessionState(
|
||||
expected: string,
|
||||
timeoutMs = 2000,
|
||||
): Promise<Awaited<ReturnType<typeof readScanSession>>> {
|
||||
const deadline = Date.now() + timeoutMs;
|
||||
let session = await readScanSession();
|
||||
while ((session?.state ?? null) !== expected && Date.now() < deadline) {
|
||||
await new Promise((resolve) => setTimeout(resolve, 5));
|
||||
session = await readScanSession();
|
||||
}
|
||||
return session;
|
||||
}
|
||||
|
||||
/** Let a detached scan from the previous test settle before the next one. */
|
||||
async function settleDetachedScans(): Promise<void> {
|
||||
cancelActiveScan();
|
||||
await new Promise((resolve) => setTimeout(resolve, 20));
|
||||
}
|
||||
|
||||
function mockReaddir(byDir: Record<string, string[]>): void {
|
||||
readdirFn.mockImplementation(async (dir: string) => byDir[dir] ?? []);
|
||||
}
|
||||
@@ -429,6 +482,166 @@ describe("autoCleanFakeNitros", () => {
|
||||
});
|
||||
});
|
||||
|
||||
describe("refreshScanCacheAfterMutation", () => {
|
||||
it("drops removed files and recomputes digests so the cache stays valid", async () => {
|
||||
const dirs = {
|
||||
"nitro\u0000/assets/nitro": dirSignature(["chair.nitro", "ghost.nitro"]),
|
||||
"swf\u0000/assets/swf": dirSignature([]),
|
||||
"icon\u0000/assets/icons": dirSignature([]),
|
||||
};
|
||||
const files = installFileState({
|
||||
"/var/www/atom-nexst/storage/nitro-cleanup/scan-cache.json":
|
||||
JSON.stringify({
|
||||
dirs,
|
||||
result: {
|
||||
fake: [
|
||||
{
|
||||
fileName: "ghost.nitro",
|
||||
base: "ghost",
|
||||
size: 500,
|
||||
dirs: ["/assets/nitro"],
|
||||
lastModified: 0,
|
||||
},
|
||||
],
|
||||
broken: [],
|
||||
orphanedSwf: [],
|
||||
orphanedIcon: [],
|
||||
total: 2,
|
||||
},
|
||||
}),
|
||||
});
|
||||
// The directory changed after the cached scan (ghost.nitro deleted).
|
||||
readdirFn.mockImplementation(async (dir: string) =>
|
||||
dir.endsWith("/nitro") ? ["chair.nitro"] : [],
|
||||
);
|
||||
|
||||
await refreshScanCacheAfterMutation("nitro", ["ghost.nitro"], 1);
|
||||
|
||||
const written = JSON.parse(
|
||||
files.get(
|
||||
"/var/www/atom-nexst/storage/nitro-cleanup/scan-cache.json",
|
||||
) as string,
|
||||
) as {
|
||||
dirs: Record<string, string>;
|
||||
result: {
|
||||
fake: Array<{ fileName: string }>;
|
||||
broken: unknown[];
|
||||
orphanedSwf: unknown[];
|
||||
orphanedIcon: unknown[];
|
||||
total: number;
|
||||
};
|
||||
};
|
||||
expect(written.dirs["nitro\u0000/assets/nitro"]).toBe(
|
||||
dirSignature(["chair.nitro"]),
|
||||
);
|
||||
expect(written.result.fake).toEqual([]);
|
||||
expect(written.result.total).toBe(1);
|
||||
// The persisted result is now the refreshed one.
|
||||
await expect(getPersistedCleanupResult()).resolves.toMatchObject({
|
||||
fake: [],
|
||||
total: 1,
|
||||
});
|
||||
});
|
||||
|
||||
it("is a no-op when the file names list is empty or no cache exists", async () => {
|
||||
installFileState({});
|
||||
await refreshScanCacheAfterMutation("nitro", [], 0);
|
||||
installFileState({});
|
||||
await refreshScanCacheAfterMutation("nitro", ["ghost.nitro"], 1);
|
||||
expect(writeFileFn).not.toHaveBeenCalled();
|
||||
});
|
||||
});
|
||||
|
||||
describe("scan aborting", () => {
|
||||
it("throws an AbortError and skips the scan cache write when aborted", async () => {
|
||||
const controller = new AbortController();
|
||||
controller.abort();
|
||||
|
||||
await expect(
|
||||
scanFakeBrokenNitros({ signal: controller.signal }),
|
||||
).rejects.toMatchObject({ name: "AbortError" });
|
||||
expect(writeFileFn).not.toHaveBeenCalled();
|
||||
expect(readFileFn).not.toHaveBeenCalled();
|
||||
});
|
||||
});
|
||||
|
||||
describe("nitro-scan-session", () => {
|
||||
it("runs a scan in the background and persists the done session", async () => {
|
||||
await settleDetachedScans();
|
||||
installFileState({});
|
||||
readdirFn.mockImplementation(async (dir: string) =>
|
||||
dir.endsWith("/nitro") ? ["chair.nitro"] : [],
|
||||
);
|
||||
parseNitroBundle.mockReturnValue({});
|
||||
statFn.mockImplementation(async () => ({ size: 512, mtimeMs: 0 }));
|
||||
|
||||
const session = await runScanInBackground({ force: true });
|
||||
|
||||
expect(session.id).toBeTruthy();
|
||||
const done = await waitForSessionState("done");
|
||||
expect(done?.id).toBe(session.id);
|
||||
expect(done?.total).toBe(1);
|
||||
expect(done?.cached).toBe(false);
|
||||
});
|
||||
|
||||
it("rejects a second scan while a fresh one is running (single-flight)", async () => {
|
||||
await settleDetachedScans();
|
||||
installFileState({});
|
||||
readdirFn.mockImplementation(async (dir: string) =>
|
||||
dir.endsWith("/nitro") ? ["chair.nitro"] : [],
|
||||
);
|
||||
parseNitroBundle.mockReturnValue({});
|
||||
|
||||
await runScanInBackground({ force: false });
|
||||
await expect(runScanInBackground({ force: false })).rejects.toBeInstanceOf(
|
||||
NitroScanConflictError,
|
||||
);
|
||||
|
||||
await settleDetachedScans();
|
||||
});
|
||||
|
||||
it("lets a stale running session be taken over (crashed process recovery)", async () => {
|
||||
await settleDetachedScans();
|
||||
installFileState({});
|
||||
readdirFn.mockImplementation(async (dir: string) =>
|
||||
dir.endsWith("/nitro") ? ["chair.nitro"] : [],
|
||||
);
|
||||
parseNitroBundle.mockReturnValue({});
|
||||
const staleHourAgo = new Date(Date.now() - 60 * 60 * 1000).toISOString();
|
||||
await writeScanSession({
|
||||
id: "stale-session",
|
||||
state: "running",
|
||||
startedAt: staleHourAgo,
|
||||
updatedAt: staleHourAgo,
|
||||
force: false,
|
||||
});
|
||||
|
||||
const session = await runScanInBackground({ force: false });
|
||||
|
||||
expect(session.id).not.toBe("stale-session");
|
||||
await expect(waitForSessionState("done")).resolves.toMatchObject({
|
||||
id: session.id,
|
||||
});
|
||||
});
|
||||
|
||||
it("marks the session cancelled when the scan is aborted", async () => {
|
||||
await settleDetachedScans();
|
||||
installFileState({});
|
||||
readdirFn.mockImplementation(async (dir: string) =>
|
||||
dir.endsWith("/nitro") ? ["chair.nitro"] : [],
|
||||
);
|
||||
parseNitroBundle.mockReturnValue({});
|
||||
|
||||
const session = await runScanInBackground({ force: false });
|
||||
const cancelled = await cancelScanSession();
|
||||
|
||||
expect(cancelled?.state).toBe("cancelled");
|
||||
await expect(waitForSessionState("cancelled")).resolves.toMatchObject({
|
||||
id: session.id,
|
||||
});
|
||||
});
|
||||
});
|
||||
|
||||
describe("repairBrokenNitros", () => {
|
||||
it("re-downloads valid bundles and writes them to every nitro dir", async () => {
|
||||
getTargetsFn.mockResolvedValue({
|
||||
|
||||
@@ -74,6 +74,12 @@ export interface CleanupScanOptions {
|
||||
* leaving the tab looking frozen while a large directory is processed.
|
||||
*/
|
||||
onProgress?: (progress: CleanupProgress) => void;
|
||||
/**
|
||||
* Aborts the scan. Every CPU-bound stage (readdir, stat pool, header
|
||||
* validation) checks the signal and throws an `AbortError`-named error, so a
|
||||
* cancelled scan never writes a partial result to the on-disk cache.
|
||||
*/
|
||||
signal?: AbortSignal;
|
||||
}
|
||||
|
||||
export interface NitroCleanupDeleteResult {
|
||||
@@ -130,6 +136,16 @@ const KIND_FILE_RE: Record<CleanupAssetKind, RegExp> = {
|
||||
|
||||
const DAY_MS = 24 * 60 * 60 * 1000;
|
||||
|
||||
function abortSignalReason(): Error {
|
||||
const error = new Error("Scan cancelled");
|
||||
error.name = "AbortError";
|
||||
return error;
|
||||
}
|
||||
|
||||
function throwIfAborted(signal?: AbortSignal): void {
|
||||
if (signal?.aborted) throw abortSignalReason();
|
||||
}
|
||||
|
||||
function uniqueDirs(dirs: string[]): string[] {
|
||||
const seen = new Set<string>();
|
||||
const out: string[] = [];
|
||||
@@ -249,9 +265,11 @@ const DIR_STAT_CONCURRENCY = 32;
|
||||
async function buildDirEntries(
|
||||
dir: string,
|
||||
names: string[],
|
||||
signal?: AbortSignal,
|
||||
): Promise<DirectoryRead[]> {
|
||||
const entries: DirectoryRead[] = [];
|
||||
await runPool(names, DIR_STAT_CONCURRENCY, async (name) => {
|
||||
throwIfAborted(signal);
|
||||
let size = 0;
|
||||
let mtimeMs = 0;
|
||||
try {
|
||||
@@ -275,11 +293,13 @@ async function collectAssets(
|
||||
kind: CleanupAssetKind,
|
||||
dirs: string[],
|
||||
namesByKey: Map<string, string[]>,
|
||||
signal?: AbortSignal,
|
||||
): Promise<Map<string, AssetCleanupEntry>> {
|
||||
const byName = new Map<string, AssetCleanupEntry>();
|
||||
for (const dir of dirs) {
|
||||
throwIfAborted(signal);
|
||||
const names = namesByKey.get(kindKey(kind, dir)) ?? [];
|
||||
for (const read of await buildDirEntries(dir, names)) {
|
||||
for (const read of await buildDirEntries(dir, names, signal)) {
|
||||
const entry = byName.get(read.name) ?? {
|
||||
fileName: read.name,
|
||||
base: read.base,
|
||||
@@ -355,7 +375,7 @@ function cacheDirsMatch(
|
||||
export async function scanFakeBrokenNitros(
|
||||
options: CleanupScanOptions = {},
|
||||
): Promise<NitroCleanupScan> {
|
||||
const { force = false, onProgress } = options;
|
||||
const { force = false, onProgress, signal } = options;
|
||||
const targets = await getFurniAssetWriteTargets();
|
||||
const nitroDirs = uniqueDirs([
|
||||
targets.nitroDir,
|
||||
@@ -386,6 +406,7 @@ export async function scanFakeBrokenNitros(
|
||||
await Promise.all(
|
||||
directoryPlan.flatMap(({ kind, dirs }) =>
|
||||
dirs.map(async (dir) => {
|
||||
throwIfAborted(signal);
|
||||
namesByKey.set(
|
||||
kindKey(kind, dir),
|
||||
await readDirNames(dir, KIND_FILE_RE[kind]),
|
||||
@@ -398,6 +419,7 @@ export async function scanFakeBrokenNitros(
|
||||
currentDigests[key] = dirNameSignature(namesByKey.get(key) ?? []);
|
||||
}
|
||||
|
||||
throwIfAborted(signal);
|
||||
onProgress?.({ phase: "readdir", scanned: 0 });
|
||||
if (!force) {
|
||||
const cache = await readScanCache();
|
||||
@@ -407,14 +429,15 @@ export async function scanFakeBrokenNitros(
|
||||
}
|
||||
}
|
||||
|
||||
throwIfAborted(signal);
|
||||
// Needed only for classification — skip the DB query on a cache hit.
|
||||
onProgress?.({ phase: "stems", scanned: 0 });
|
||||
const validStems = await getCleanupValidStems();
|
||||
|
||||
const [byNitro, bySwf, byIcon] = await Promise.all([
|
||||
collectAssets("nitro", nitroDirs, namesByKey),
|
||||
collectAssets("swf", swfDirs, namesByKey),
|
||||
collectAssets("icon", iconDirs, namesByKey),
|
||||
collectAssets("nitro", nitroDirs, namesByKey, signal),
|
||||
collectAssets("swf", swfDirs, namesByKey, signal),
|
||||
collectAssets("icon", iconDirs, namesByKey, signal),
|
||||
]);
|
||||
|
||||
const entries: NitroCleanupEntry[] = [...byNitro.values()];
|
||||
@@ -424,6 +447,7 @@ export async function scanFakeBrokenNitros(
|
||||
const totalToValidate = entries.length;
|
||||
|
||||
await runPool(entries, 16, async (entry) => {
|
||||
throwIfAborted(signal);
|
||||
// A .nitro that no DB item maps to is a leftover / fake bundle.
|
||||
if (!validStems.has(entry.base.toLowerCase())) {
|
||||
fake.push(entry);
|
||||
@@ -477,6 +501,7 @@ export async function scanFakeBrokenNitros(
|
||||
orphanedIcon: byName(orphanedIcon),
|
||||
total: byNitro.size + bySwf.size + byIcon.size,
|
||||
};
|
||||
throwIfAborted(signal);
|
||||
await writeScanCache(currentDigests, result);
|
||||
onProgress?.({ phase: "done", scanned: result.total });
|
||||
return result;
|
||||
@@ -567,6 +592,147 @@ export async function autoCleanFakeNitros(
|
||||
};
|
||||
}
|
||||
|
||||
export async function getPersistedCleanupResult(): Promise<NitroCleanupScan | null> {
|
||||
const cache = await readScanCache();
|
||||
return cache?.result ?? null;
|
||||
}
|
||||
|
||||
/**
|
||||
* Keep the on-disk scan cache consistent after a mutation without re-scanning:
|
||||
* recompute the directory signatures for the touched asset kind (cheap readdir
|
||||
* pass — no stats, no header validation) and drop the removed file names from
|
||||
* the cached result lists. Without this, the next non-forced scan would see a
|
||||
* digest mismatch (e.g. after a delete) and re-validate the whole directory.
|
||||
*/
|
||||
export async function refreshScanCacheAfterMutation(
|
||||
kind: CleanupAssetKind,
|
||||
fileNames: string[],
|
||||
removedCount = fileNames.length,
|
||||
): Promise<void> {
|
||||
const remove = new Set(fileNames);
|
||||
if (remove.size === 0) return;
|
||||
const cache = await readScanCache();
|
||||
if (!cache) return;
|
||||
|
||||
const targets = await getFurniAssetWriteTargets();
|
||||
for (const dir of dirsForKind(targets, kind)) {
|
||||
cache.dirs[kindKey(kind, dir)] = dirNameSignature(
|
||||
await readDirNames(dir, KIND_FILE_RE[kind]),
|
||||
);
|
||||
}
|
||||
|
||||
const result = cache.result;
|
||||
if (kind === "nitro") {
|
||||
result.fake = result.fake.filter((e) => !remove.has(e.fileName));
|
||||
result.broken = result.broken.filter((e) => !remove.has(e.fileName));
|
||||
} else if (kind === "swf") {
|
||||
result.orphanedSwf = result.orphanedSwf.filter(
|
||||
(e) => !remove.has(e.fileName),
|
||||
);
|
||||
} else {
|
||||
result.orphanedIcon = result.orphanedIcon.filter(
|
||||
(e) => !remove.has(e.fileName),
|
||||
);
|
||||
}
|
||||
result.total = Math.max(0, result.total - removedCount);
|
||||
await writeScanCache(cache.dirs, result);
|
||||
}
|
||||
|
||||
// ── Auto-clean history ──────────────────────────────────────────────────────
|
||||
//
|
||||
// The nightly scheduled auto-clean and the manual button both record a history
|
||||
// entry so staff can see what the jobs worker removed and when.
|
||||
|
||||
export interface NitroAutoCleanHistoryEntry {
|
||||
triggeredAt: string;
|
||||
triggeredBy: "manual" | "scheduled";
|
||||
maxAgeDays: number;
|
||||
deleted: number;
|
||||
copiesRemoved: number;
|
||||
skippedRecent: number;
|
||||
skippedBroken: number;
|
||||
errors: string[];
|
||||
}
|
||||
|
||||
const HISTORY_MAX_ENTRIES = 50;
|
||||
const historyPath = () => path.join(SCAN_CACHE_DIR, "history.json");
|
||||
|
||||
export async function readAutoCleanHistory(): Promise<
|
||||
NitroAutoCleanHistoryEntry[]
|
||||
> {
|
||||
try {
|
||||
if (!existsSync(historyPath())) return [];
|
||||
const parsed = JSON.parse(await fs.readFile(historyPath(), "utf8"));
|
||||
if (!Array.isArray(parsed)) return [];
|
||||
return parsed as NitroAutoCleanHistoryEntry[];
|
||||
} catch {
|
||||
return [];
|
||||
}
|
||||
}
|
||||
|
||||
export async function appendAutoCleanHistory(
|
||||
entry: NitroAutoCleanHistoryEntry,
|
||||
): Promise<void> {
|
||||
try {
|
||||
const history = (await readAutoCleanHistory()).slice(
|
||||
0,
|
||||
HISTORY_MAX_ENTRIES - 1,
|
||||
);
|
||||
await fs.mkdir(SCAN_CACHE_DIR, { recursive: true });
|
||||
await fs.writeFile(historyPath(), JSON.stringify([entry, ...history]));
|
||||
} catch {
|
||||
/* best effort — history is never fatal */
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Refresh the scan cache + append a history entry after an auto-clean run.
|
||||
* Shared by the manual route and the nightly scheduled job.
|
||||
*/
|
||||
export async function finishAutoClean(options: {
|
||||
result: NitroAutoCleanResult;
|
||||
maxAgeDays: number;
|
||||
triggeredBy: "manual" | "scheduled";
|
||||
}): Promise<void> {
|
||||
const removed = options.result.files
|
||||
.filter((file) => file.deleted)
|
||||
.map((file) => file.fileName);
|
||||
await refreshScanCacheAfterMutation("nitro", removed, options.result.deleted);
|
||||
await appendAutoCleanHistory({
|
||||
triggeredAt: new Date().toISOString(),
|
||||
triggeredBy: options.triggeredBy,
|
||||
maxAgeDays: options.maxAgeDays,
|
||||
deleted: options.result.deleted,
|
||||
copiesRemoved: options.result.copiesRemoved,
|
||||
skippedRecent: options.result.skippedRecent,
|
||||
skippedBroken: options.result.skippedBroken,
|
||||
errors: options.result.errors,
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
* Entry point for the nightly jobs-worker cron. Never runs while an interactive
|
||||
* scan session is in flight, so scheduled deletion can't race the UI.
|
||||
*/
|
||||
export async function scheduledAutoCleanFakeNitros(
|
||||
maxAgeDays = 30,
|
||||
): Promise<{ skipped: boolean; deleted?: number }> {
|
||||
const { isSessionStale, readScanSession } = await import(
|
||||
"./nitro-scan-session"
|
||||
);
|
||||
const session = await readScanSession();
|
||||
if (session?.state === "running" && !isSessionStale(session)) {
|
||||
return { skipped: true };
|
||||
}
|
||||
const result = await autoCleanFakeNitros(maxAgeDays);
|
||||
await finishAutoClean({
|
||||
result,
|
||||
maxAgeDays,
|
||||
triggeredBy: "scheduled",
|
||||
});
|
||||
return { skipped: false, deleted: result.deleted };
|
||||
}
|
||||
|
||||
async function copyFileToDirs(
|
||||
buffer: Buffer,
|
||||
fileName: string,
|
||||
|
||||
@@ -0,0 +1,225 @@
|
||||
import { randomUUID } from "node:crypto";
|
||||
import { existsSync, promises as fs } from "node:fs";
|
||||
import path from "node:path";
|
||||
import { logger } from "@/lib/logger";
|
||||
import {
|
||||
type CleanupProgress,
|
||||
scanFakeBrokenNitros,
|
||||
} from "@/lib/services/nitro-cleanup";
|
||||
|
||||
/**
|
||||
* Durable nitro-cleanup scan session.
|
||||
*
|
||||
* A scan is expensive (hundreds of thousands of stat + header reads), so it is
|
||||
* better run as a detached in-process job than blocked inside one request:
|
||||
* closing the tab must not lose the scan and a second tab should re-attach
|
||||
* instead of starting a duplicate. The session is persisted on disk so the
|
||||
* Next.js process keeps working it after the POST that started it has already
|
||||
* returned, and a restart mid-scan is recoverable — a "running" session whose
|
||||
* heartbeat (`updatedAt`) has gone stale is treated as dead and taken over.
|
||||
*
|
||||
* The lightweight session record only stores metadata + progress; the actual
|
||||
* result lists live in the scan cache (`scan-cache.json`) written by the scan,
|
||||
* so a done session is read back without duplicating megabytes of JSON.
|
||||
*/
|
||||
|
||||
export const nitroScanRoot = () =>
|
||||
path.join(process.cwd(), "storage", "nitro-cleanup");
|
||||
const sessionPath = () => path.join(nitroScanRoot(), "session.json");
|
||||
|
||||
export type NitroScanState = "running" | "done" | "cancelled" | "error";
|
||||
|
||||
export interface NitroScanSession {
|
||||
id: string;
|
||||
state: NitroScanState;
|
||||
startedAt: string;
|
||||
updatedAt: string;
|
||||
force: boolean;
|
||||
progress?: { phase: string; scanned: number };
|
||||
cached?: boolean;
|
||||
total?: number;
|
||||
error?: string;
|
||||
}
|
||||
|
||||
/** A running session that has not been heard from for this long is stale. */
|
||||
export const SESSION_STALE_MS = 60_000;
|
||||
|
||||
/** Thrown when a fresh scan session already exists (single-flight). */
|
||||
export class NitroScanConflictError extends Error {
|
||||
readonly session: NitroScanSession;
|
||||
constructor(session: NitroScanSession) {
|
||||
super("A nitro scan is already running");
|
||||
this.name = "NitroScanConflictError";
|
||||
this.session = session;
|
||||
}
|
||||
}
|
||||
|
||||
export async function readScanSession(): Promise<NitroScanSession | null> {
|
||||
try {
|
||||
if (!existsSync(sessionPath())) return null;
|
||||
const parsed = JSON.parse(
|
||||
await fs.readFile(sessionPath(), "utf8"),
|
||||
) as NitroScanSession;
|
||||
return parsed && typeof parsed.id === "string" ? parsed : null;
|
||||
} catch {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
export async function writeScanSession(
|
||||
session: NitroScanSession,
|
||||
): Promise<void> {
|
||||
try {
|
||||
await fs.mkdir(nitroScanRoot(), { recursive: true });
|
||||
await fs.writeFile(sessionPath(), JSON.stringify(session));
|
||||
} catch {
|
||||
/* session file is best effort — never fatal */
|
||||
}
|
||||
}
|
||||
|
||||
export function isSessionStale(session: NitroScanSession): boolean {
|
||||
const age = Date.now() - new Date(session.updatedAt).getTime();
|
||||
return !Number.isFinite(age) || age > SESSION_STALE_MS;
|
||||
}
|
||||
|
||||
/**
|
||||
* Create a new scan session, or reject with `NitroScanConflictError` when a
|
||||
* fresh session is already running (single-flight across the whole server).
|
||||
*/
|
||||
export async function startScanSession(options: {
|
||||
force: boolean;
|
||||
}): Promise<NitroScanSession> {
|
||||
const current = await readScanSession();
|
||||
if (current?.state === "running" && !isSessionStale(current)) {
|
||||
throw new NitroScanConflictError(current);
|
||||
}
|
||||
const now = new Date().toISOString();
|
||||
const session: NitroScanSession = {
|
||||
id: randomUUID(),
|
||||
state: "running",
|
||||
startedAt: now,
|
||||
updatedAt: now,
|
||||
force: options.force,
|
||||
};
|
||||
await writeScanSession(session);
|
||||
return session;
|
||||
}
|
||||
|
||||
interface ActiveScan {
|
||||
sessionId: string;
|
||||
controller: AbortController;
|
||||
promise: Promise<void>;
|
||||
}
|
||||
|
||||
let activeScan: ActiveScan | null = null;
|
||||
|
||||
/**
|
||||
* Start (or re-attach to) a scan in the background. Resolves with the session
|
||||
* quickly; the scan itself keeps running detached in the process. Callers
|
||||
* poll `readScanSession()` to follow progress.
|
||||
*
|
||||
* @throws {NitroScanConflictError} when a fresh (non-stale) scan is already
|
||||
* running — the caller should attach to the returned session instead.
|
||||
*/
|
||||
export async function runScanInBackground(options: {
|
||||
force: boolean;
|
||||
}): Promise<NitroScanSession> {
|
||||
const session = await startScanSession(options);
|
||||
const controller = new AbortController();
|
||||
|
||||
const sessionEvent = async (patch?: {
|
||||
progress?: CleanupProgress;
|
||||
}): Promise<void> => {
|
||||
await writeScanSession({
|
||||
...session,
|
||||
state: "running",
|
||||
...(patch?.progress ? { progress: patch.progress } : {}),
|
||||
updatedAt: new Date().toISOString(),
|
||||
});
|
||||
};
|
||||
|
||||
const run = (async () => {
|
||||
try {
|
||||
const scan = await scanFakeBrokenNitros({
|
||||
force: session.force,
|
||||
signal: controller.signal,
|
||||
onProgress: (progress) => void sessionEvent({ progress }),
|
||||
});
|
||||
await writeScanSession({
|
||||
...session,
|
||||
state: "done",
|
||||
cached: scan.cached ?? false,
|
||||
total: scan.total,
|
||||
progress: { phase: "done", scanned: scan.total },
|
||||
updatedAt: new Date().toISOString(),
|
||||
});
|
||||
} catch (error) {
|
||||
if ((error as Error).name === "AbortError") {
|
||||
await writeScanSession({
|
||||
...session,
|
||||
state: "cancelled",
|
||||
updatedAt: new Date().toISOString(),
|
||||
});
|
||||
} else {
|
||||
logger.error("Nitro scan failed", {
|
||||
module: "nitro-cleanup-session",
|
||||
sessionId: session.id,
|
||||
error: error instanceof Error ? error.message : String(error),
|
||||
});
|
||||
await writeScanSession({
|
||||
...session,
|
||||
state: "error",
|
||||
error: error instanceof Error ? error.message : "Scan failed",
|
||||
updatedAt: new Date().toISOString(),
|
||||
});
|
||||
}
|
||||
} finally {
|
||||
if (activeScan?.sessionId === session.id) activeScan = null;
|
||||
}
|
||||
})();
|
||||
|
||||
activeScan = { sessionId: session.id, controller, promise: run };
|
||||
return session;
|
||||
}
|
||||
|
||||
/** Abort the in-process scan (if any). Returns whether one was aborted. */
|
||||
export function cancelActiveScan(): boolean {
|
||||
if (!activeScan) return false;
|
||||
activeScan.controller.abort();
|
||||
return true;
|
||||
}
|
||||
|
||||
/** Cancel the scan session regardless of who is running it. */
|
||||
export async function cancelScanSession(): Promise<NitroScanSession | null> {
|
||||
if (cancelActiveScan()) {
|
||||
// The detached run observes the abort and marks the session cancelled —
|
||||
// wait for that so the caller sees the definitive end state.
|
||||
const deadline = Date.now() + 2000;
|
||||
let session = await readScanSession();
|
||||
while (session?.state === "running" && Date.now() < deadline) {
|
||||
await new Promise((resolve) => setTimeout(resolve, 10));
|
||||
session = await readScanSession();
|
||||
}
|
||||
if (session?.state === "running" && session) {
|
||||
const cancelled: NitroScanSession = {
|
||||
...session,
|
||||
state: "cancelled",
|
||||
updatedAt: new Date().toISOString(),
|
||||
};
|
||||
await writeScanSession(cancelled);
|
||||
session = cancelled;
|
||||
}
|
||||
return session;
|
||||
}
|
||||
const session = await readScanSession();
|
||||
if (session?.state === "running") {
|
||||
const cancelled: NitroScanSession = {
|
||||
...session,
|
||||
state: "cancelled",
|
||||
updatedAt: new Date().toISOString(),
|
||||
};
|
||||
await writeScanSession(cancelled);
|
||||
return cancelled;
|
||||
}
|
||||
return session;
|
||||
}
|
||||
Reference in new issue
Block a user