import { existsSync, linkSync, mkdirSync, readFileSync, rmSync, writeFileSync, } from "node:fs"; import { randomUUID } from "node:crypto"; import { dirname, join } from "node:path"; /** Marker left while an explicitly released operation awaits final cleanup. */ export const OPERATION_RELEASE_MARKER_FILE_NAME = ".released"; /** File errors used by Node for short-lived handle contention. */ const TRANSIENT_FILE_ERROR_CODES = new Set(["EBUSY", "EPERM"]); export interface FileRetryOptions { /** Total number of attempts, including the first call. */ attempts?: number; /** Delay before the second attempt. */ initialDelayMs?: number; /** Upper bound for the exponential delay between attempts. */ maxDelayMs?: number; /** Clock-free sleep seam used by tests and synchronous Node callers. */ sleep?: (milliseconds: number) => void; /** Predicate for errors that may become successful on a later attempt. */ shouldRetry?: (error: unknown) => boolean; } const DEFAULT_ATTEMPTS = 5; const DEFAULT_INITIAL_DELAY_MS = 10; const DEFAULT_MAX_DELAY_MS = 100; function sleepSync(milliseconds: number): void { if (milliseconds <= 0) return; const signal = new Int32Array(new SharedArrayBuffer(4)); Atomics.wait(signal, 0, 0, milliseconds); } function errorCode(error: unknown): string | undefined { if (error == null || typeof error !== "object") return undefined; const code = (error as { code?: unknown }).code; return typeof code === "string" ? code : undefined; } /** Whether an error represents short-lived filesystem contention. */ export function isTransientFileError(error: unknown): boolean { const code = errorCode(error); return code !== undefined && TRANSIENT_FILE_ERROR_CODES.has(code); } function normalizedAttempts(value: number | undefined): number { if (value === undefined) return DEFAULT_ATTEMPTS; if (!Number.isFinite(value)) return DEFAULT_ATTEMPTS; return Math.max(1, Math.floor(value)); } function normalizedDelay(value: number | undefined, fallback: number): number { if (value === undefined || !Number.isFinite(value)) return fallback; return Math.max(0, value); } /** Retry a synchronous filesystem operation across transient contention. */ export function retryFileOperation( operation: () => T, options: FileRetryOptions = {}, ): T { const attempts = normalizedAttempts(options.attempts); const sleep = options.sleep ?? sleepSync; const shouldRetry = options.shouldRetry ?? isTransientFileError; const maxDelayMs = normalizedDelay(options.maxDelayMs, DEFAULT_MAX_DELAY_MS); let delayMs = Math.min( normalizedDelay(options.initialDelayMs, DEFAULT_INITIAL_DELAY_MS), maxDelayMs, ); for (let attempt = 1; ; attempt += 1) { try { return operation(); } catch (error) { if (attempt >= attempts || !shouldRetry(error)) throw error; sleep(delayMs); delayMs = Math.min(Math.max(delayMs * 2, delayMs), maxDelayMs); } } } function isMissingFileError(error: unknown): boolean { return errorCode(error) === "ENOENT"; } /** Read and parse a JSON artifact, retrying only transient handle contention. */ export function readJsonFileWithRetry(path: string): unknown { try { return retryFileOperation( () => JSON.parse(readFileSync(path, "utf8")) as unknown, ); } catch (error) { if (isMissingFileError(error)) return undefined; throw error; } } export interface AtomicFilePublishOptions extends FileRetryOptions { /** Treat a destination created by another publisher as success. */ idempotent?: boolean; /** Validate an existing destination before an idempotent publish succeeds. */ isExistingDestinationAcceptable?: (path: string) => boolean; /** File mode for the private temporary artifact. */ mode?: number; /** Whether to create the parent; operation namespaces disable this after release. */ createParentDirectory?: boolean; } function bestEffortRemoveTemporaryFile(path: string): void { try { retryFileOperation( () => rmSync(path, { force: true, maxRetries: 3, retryDelay: 25 }), { attempts: 3, initialDelayMs: 10, maxDelayMs: 50 }, ); } catch { // The operation directory owns temporary files; a later namespace // cleanup can remove one that is briefly held by a scanner. } } /** * Publish one file without ever replacing an existing destination. * * When `createParentDirectory` is false, the caller owns an already-created * artifact namespace and this function never recreates it after release. The * payload is written to a private sibling and then hard-linked into place. * Hard-link creation is an exclusive, atomic publication on supported Node * filesystems, including Windows. A transient EPERM/EBUSY-style error retries * the complete write/link attempt with a fresh temporary sibling. */ export function publishAtomicFile( path: string, contents: string, options: AtomicFilePublishOptions = {}, ): string { const createParentDirectory = options.createParentDirectory ?? true; if (createParentDirectory) { retryFileOperation(() => mkdirSync(dirname(path), { recursive: true })); } const publishAttempt = (): string => { if (!createParentDirectory && existsSync(join(dirname(path), OPERATION_RELEASE_MARKER_FILE_NAME))) { throw new Error("Operation artifact namespace has already been released."); } const temporaryPath = `${path}.tmp-${process.pid}-${randomUUID()}`; try { writeFileSync(temporaryPath, contents, { encoding: "utf8", flag: "wx", mode: options.mode ?? 0o600, }); linkSync(temporaryPath, path); return path; } finally { bestEffortRemoveTemporaryFile(temporaryPath); } }; try { return retryFileOperation(publishAttempt, options); } catch (error) { // An idempotent publisher may lose the link race to another publisher. It // never replaces the destination; callers may validate that the existing // payload is the expected idempotent record before accepting the race. if (options.idempotent && existsSync(path)) { let acceptable = true; try { acceptable = options.isExistingDestinationAcceptable?.(path) ?? true; } catch { acceptable = false; } if (acceptable) return path; } throw error; } }