import { existsSync, mkdtempSync, mkdirSync, readFileSync, readdirSync, rmSync, statSync, writeFileSync, } from "node:fs"; import { tmpdir } from "node:os"; import { dirnamePortablePath, joinPortablePath } from "./path-utils.ts"; import { isValidOperationId } from "./operation-identity.ts"; import { isTransientFileError, OPERATION_RELEASE_MARKER_FILE_NAME, retryFileOperation, } from "./file-retry.ts"; /** Files needed by one ephemeral operation's protocol and launch. */ export type OperationArtifactKind = | "completion" | "launchPrompt" | "activity" | "cancellationRequest" | "cancellation"; const ARTIFACT_FILE_NAMES = { completion: "completion.json", launchPrompt: "launch-prompt.md", activity: "activity.json", cancellationRequest: "cancellation-request.json", cancellation: "cancellation.json", } as const; /** Paths for the active operation protocol. */ export interface OperationArtifactPaths { readonly completion: string; /** Optional when a child only needs to read the protocol files. */ readonly launchPrompt?: string; readonly activity: string; readonly cancellationRequest: string; readonly cancellation: string; } /** * Opaque namespace for one Completion operation's temporary protocol data. * Callers ask for a named artifact and never construct the directory layout. */ export interface OperationArtifacts { readonly operationId: string; path(kind: OperationArtifactKind): string; release(): void; } /** Adapter seam for allocating operation namespaces. */ export interface OperationArtifactsAllocator { allocate(operationId: string): OperationArtifacts; } export interface OperationArtifactsAllocatorOptions { /** Override only the parent directory; production defaults to os.tmpdir(). */ temporaryDirectory?: string; /** Prefix passed to the platform's temporary-directory allocator. */ prefix?: string; /** Age after which an abandoned operation directory may be removed. */ staleAfterMs?: number; /** Clock seam for deterministic stale-artifact cleanup tests. */ now?: () => number; /** Bound one cleanup pass so a large temporary directory stays responsive. */ maxCleanupEntries?: number; } const DEFAULT_STALE_AFTER_MS = 24 * 60 * 60 * 1_000; const DEFAULT_MAX_CLEANUP_ENTRIES = 100; const RELEASE_RETRY_ATTEMPTS = 5; const RELEASE_RETRY_INITIAL_DELAY_MS = 10; const RELEASE_RETRY_MAX_DELAY_MS = 100; const DEFERRED_RELEASE_ATTEMPTS = 8; const DEFERRED_RELEASE_DELAY_MS = 100; const OWNER_FILE_NAME = ".owner"; function requireOperationId(operationId: string): void { if (!isValidOperationId(operationId)) { throw new Error("Operation artifacts require a valid operation identity."); } } function requireArtifactPath(path: string, kind: OperationArtifactKind): void { if ( typeof path !== "string" || path.length === 0 || path.includes("\0") || /[\r\n]/.test(path) ) { throw new Error(`Operation artifact ${kind} requires a valid path.`); } } function completePaths(directory: string): OperationArtifactPaths { return Object.freeze({ completion: joinPortablePath(directory, ARTIFACT_FILE_NAMES.completion), launchPrompt: joinPortablePath(directory, ARTIFACT_FILE_NAMES.launchPrompt), activity: joinPortablePath(directory, ARTIFACT_FILE_NAMES.activity), cancellationRequest: joinPortablePath(directory, ARTIFACT_FILE_NAMES.cancellationRequest), cancellation: joinPortablePath(directory, ARTIFACT_FILE_NAMES.cancellation), }); } /** * Build an operation reference around already allocated paths. * * This is used at process seams where the child receives artifact paths through * its environment. The optional release callback belongs to the allocating * process; child references therefore default to a no-op release. */ export function createOperationArtifactsReference( operationId: string, paths: OperationArtifactPaths, release: () => void = () => {}, ): OperationArtifacts { requireOperationId(operationId); const resolvedPaths: Record = { completion: paths.completion, launchPrompt: paths.launchPrompt ?? joinPortablePath(dirnamePortablePath(paths.completion), ARTIFACT_FILE_NAMES.launchPrompt), activity: paths.activity, cancellationRequest: paths.cancellationRequest, cancellation: paths.cancellation, }; for (const kind of Object.keys(ARTIFACT_FILE_NAMES) as OperationArtifactKind[]) { requireArtifactPath(resolvedPaths[kind], kind); } let released = false; let deferredAttempts = 0; let deferredRelease: ReturnType | undefined; const scheduleDeferredRelease = (): void => { if (released || deferredRelease || deferredAttempts >= DEFERRED_RELEASE_ATTEMPTS) return; deferredAttempts += 1; deferredRelease = setTimeout(() => { deferredRelease = undefined; if (released) return; try { release(); released = true; } catch (error) { if (!isTransientFileError(error)) return; scheduleDeferredRelease(); } }, DEFERRED_RELEASE_DELAY_MS); (deferredRelease as any).unref?.(); }; return Object.freeze({ operationId, path(kind: OperationArtifactKind): string { return resolvedPaths[kind]; }, release(): void { if (released || deferredRelease) return; // Mark the namespace released only after the filesystem operation // succeeds. A transient Windows lock gets bounded synchronous retries, // followed by asynchronous retries so a parent-side catch cannot turn a // short-lived handle contention into a permanent leak. try { release(); released = true; } catch (error) { if (!isTransientFileError(error)) throw error; scheduleDeferredRelease(); } }, }); } function validateStaleCleanupPolicy(staleAfterMs: number, maxEntries: number): void { if (staleAfterMs < 0 || !Number.isFinite(staleAfterMs)) { throw new Error("Stale operation artifact age must be a finite non-negative number."); } if (maxEntries < 0 || !Number.isFinite(maxEntries)) { throw new Error("Stale operation artifact cleanup bound must be finite and non-negative."); } } function removeOperationDirectory(directory: string): void { retryFileOperation( () => rmSync(directory, { recursive: true, force: true, maxRetries: 3, retryDelay: 25, }), { attempts: RELEASE_RETRY_ATTEMPTS, initialDelayMs: RELEASE_RETRY_INITIAL_DELAY_MS, maxDelayMs: RELEASE_RETRY_MAX_DELAY_MS, }, ); } /** Mark a release before removal so a later cleanup pass can finish it. */ function markOperationDirectoryReleased(directory: string): void { try { retryFileOperation( () => writeFileSync( joinPortablePath(directory, OPERATION_RELEASE_MARKER_FILE_NAME), `${JSON.stringify({ pid: process.pid })}\n`, { encoding: "utf8", flag: "wx", mode: 0o600 }, ), ); } catch { // Removal remains the primary path. If the marker itself is contended or // already exists, the release retry still owns cleanup of this directory. } } function ownerProcessIsAlive(directory: string): boolean { try { const owner = JSON.parse(readFileSync(joinPortablePath(directory, OWNER_FILE_NAME), "utf8")) as { pid?: unknown; }; if (!Number.isSafeInteger(owner.pid) || (owner.pid as number) <= 0) return false; process.kill(owner.pid as number, 0); return true; } catch (error: unknown) { // A process that exists but cannot be inspected must be treated as live so // cleanup never removes another user's active operation. A transient lock // while reading the owner marker is likewise not proof of abandonment. return isTransientFileError(error) || (error != null && typeof error === "object" && (error as { code?: unknown }).code === "EACCES"); } } export function cleanupStaleOperationArtifacts(options: { temporaryDirectory: string; prefix: string; staleAfterMs: number; now?: () => number; maxEntries?: number; }): number { const now = options.now ?? Date.now; const maxEntries = options.maxEntries ?? DEFAULT_MAX_CLEANUP_ENTRIES; validateStaleCleanupPolicy(options.staleAfterMs, maxEntries); let entries: ReturnType; try { entries = readdirSync(options.temporaryDirectory, { withFileTypes: true }); } catch { return 0; } let removed = 0; let inspected = 0; for (const entry of entries) { if (inspected >= maxEntries) break; if (!entry.isDirectory() || !entry.name.startsWith(options.prefix)) continue; inspected += 1; const directory = joinPortablePath(options.temporaryDirectory, entry.name); try { const age = now() - statSync(directory).mtimeMs; const released = existsSync(joinPortablePath(directory, OPERATION_RELEASE_MARKER_FILE_NAME)); if (!released && (age <= options.staleAfterMs || ownerProcessIsAlive(directory))) continue; removeOperationDirectory(directory); removed += 1; } catch { // A live operation or a short-lived platform file lock is left for a // later bounded cleanup pass rather than disrupting a new allocation. } } return removed; } function allocateInDirectory( operationId: string, options: Required>, ): OperationArtifacts { requireOperationId(operationId); retryFileOperation(() => mkdirSync(options.temporaryDirectory, { recursive: true })); cleanupStaleOperationArtifacts({ temporaryDirectory: options.temporaryDirectory, prefix: options.prefix, staleAfterMs: options.staleAfterMs, now: options.now, maxEntries: options.maxCleanupEntries, }); const directory = retryFileOperation( () => mkdtempSync(joinPortablePath(options.temporaryDirectory, options.prefix)), ); try { retryFileOperation(() => writeFileSync( joinPortablePath(directory, OWNER_FILE_NAME), `${JSON.stringify({ pid: process.pid })}\n`, { encoding: "utf8", flag: "wx", mode: 0o600 }, )); } catch (error) { try { removeOperationDirectory(directory); } catch { // A failed owner marker leaves an unclaimed directory that age-based // cleanup can safely remove on a later allocation. } throw error; } const paths = completePaths(directory); return createOperationArtifactsReference(operationId, paths, () => { // Mark the namespace before removal. If a Windows handle outlives the // bounded release retries, the next allocation can safely remove this // explicitly released directory without mistaking a live operation for a // stale one. markOperationDirectoryReleased(directory); removeOperationDirectory(directory); }); } /** * Create the production allocator. Each allocation is a private directory * directly beneath the platform OS temporary directory. */ export function createOperationArtifactsAllocator( options: OperationArtifactsAllocatorOptions = {}, ): OperationArtifactsAllocator { const allocatorOptions = { temporaryDirectory: options.temporaryDirectory ?? tmpdir(), prefix: options.prefix ?? "pi-subagent-operation-", staleAfterMs: options.staleAfterMs ?? DEFAULT_STALE_AFTER_MS, now: options.now ?? Date.now, maxCleanupEntries: options.maxCleanupEntries ?? DEFAULT_MAX_CLEANUP_ENTRIES, }; if (!allocatorOptions.prefix) throw new Error("Operation artifact prefix cannot be empty."); if (allocatorOptions.prefix.includes("\0") || /[\\/]/.test(allocatorOptions.prefix)) { throw new Error("Operation artifact prefix must be a single path component."); } if ( typeof allocatorOptions.temporaryDirectory !== "string" || allocatorOptions.temporaryDirectory.length === 0 || allocatorOptions.temporaryDirectory.includes("\0") ) { throw new Error("Operation artifact temporary directory must be a valid path."); } validateStaleCleanupPolicy(allocatorOptions.staleAfterMs, allocatorOptions.maxCleanupEntries); return { allocate(operationId: string): OperationArtifacts { return allocateInDirectory(operationId, allocatorOptions); }, }; } /** Allocate one operation namespace with the production OS-temporary policy. */ export function allocateOperationArtifacts(operationId: string): OperationArtifacts { return createOperationArtifactsAllocator().allocate(operationId); }