From 1298a1793e8fac0d326573eb715cb113634d6c8e Mon Sep 17 00:00:00 2001 From: Daniel Strader Date: Fri, 21 Aug 2026 14:43:06 -0700 Subject: [PATCH] Restack bounded checkpoint readback after Phase 0 --- CHANGELOG.md | 1 + docs/local-checkpoints.md | 30 +- docs/shared-checkpoint-schema-v1.md | 4 +- src/checkpoints.ts | 1870 ++++++++++++++++++++++- test/fixtures/checkpoint-save-worker.js | 43 + test/inkcheck.test.js | 1004 +++++++++++- 6 files changed, 2866 insertions(+), 86 deletions(-) create mode 100644 test/fixtures/checkpoint-save-worker.js diff --git a/CHANGELOG.md b/CHANGELOG.md index a3d1835..67b2cb1 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,6 +2,7 @@ ## Unreleased +- Bound schema-v1 checkpoint readback by stored and decompressed bytes, classify reopen failures as corrupt, unsupported, or resource-limited, and add canonically self-bound private manifests so listing and retention use bounded metadata I/O without opening new frontiers. Full payload digests remain an open/resume boundary; no-clobber same-ID publication, payload-first crash recovery, and pair-inclusive quotas preserve existing v1 JSON/gzip reads, stable IDs, and exact resume pending framed schema v2. - Add an opt-in bounded NDJSON evidence stream for marathon-scale external consumers. It emits replayable numeric ending/runtime witnesses with global elapsed timestamps as they are retained, then a compact terminal summary without constructing the monolithic full-report JSON string. - Treat `--max-time` as a total CLI deadline and retain bounded time/heap headroom for clean report finalization. Explicit heap caps now expose the lower search watermark separately from the full process envelope. diff --git a/docs/local-checkpoints.md b/docs/local-checkpoints.md index ddae8bf..a9f486b 100644 --- a/docs/local-checkpoints.md +++ b/docs/local-checkpoints.md @@ -28,7 +28,9 @@ inkcheck checkpoints list --json inkcheck checkpoints show checkpoint-0123456789abcdef01234567 --json ``` -These commands return bounded metadata, not the frontier payload. `show` recompiles the project and reports: +These commands return metadata, not the frontier payload. New saves also write a private, at-most-64-KiB `checkpoint-.meta.json` sidecar. `list` and retention read only that canonical self-checksummed manifest plus the payload's file size; they never open, hash, inflate, or JSON-parse a manifested frontier. Earlier schema-v1 `.json` and `.json.gz` artifacts without a sidecar still work, using the bounded full reader as a compatibility fallback. + +The manifest self-checksum binds every field, including the stored-payload digest, so an isolated edit to creation time, grant, state count, or another field fails before it can affect listing or retention. This detects accidental/local metadata corruption; it is not authentication against someone with write access who deliberately rewrites both the manifest and its checksum. `open` and resume are the full-integrity boundary: they validate the manifest, hash the complete stored payload and compare its digest, then decompress, parse, verify the stable logical ID/configuration, and check source freshness. `show` reports: - `current`: compiled story and knot/source map match the saved checkpoint. - `stale`: source exists but no longer matches or compiles. @@ -36,19 +38,37 @@ These commands return bounded metadata, not the frontier payload. `show` recompi Resume requires `current`, the supported artifact and checkpoint schema versions, and every engine binding to match: story and knot hashes, depth, both seeds, hidden turn/random sensitivity, randomness detection, frontier envelopes, and external bindings. Corrupt content, metadata mismatch, unsupported versions, and a non-increasing grant fail closed. +Library callers can distinguish three `CheckpointReadError.kind` values instead of parsing messages: + +- `corrupt`: invalid gzip/JSON, checksum or stable-ID mismatch, or malformed metadata; +- `unsupported`: an artifact, manifest, or shared-checkpoint schema this Inkcheck cannot interpret; +- `resource_limit`: stored or decompressed bytes exceed the bounded readback envelope. This does not label the checkpoint corrupt and never deletes or replaces it. + +`listCheckpointArtifacts`, `openCheckpointArtifact`, and `loadCheckpointForResume` accept optional `maxStoredBytes` and `maxDecompressedBytes` read limits. Defaults cap stored input at 512 MiB and schema-v1 decompression at the smaller of 512 MiB and the runtime's maximum string length. Listing uses those limits only for a sidecar-free compatibility fallback. + ## Atomicity and retention -New checkpoints live under `.inkcheck/checkpoints/checkpoint-.json.gz`. Their stable ID derives from the exact logical checkpoint content plus the project-relative entrypoint, not from compression bytes, so repeating the same deterministic boundary reuses one artifact. Existing schema-v1 `.json` artifacts remain readable and resumable. +New checkpoints live under `.inkcheck/checkpoints/checkpoint-.json.gz`, with the bounded metadata sidecar beside them. Their stable ID derives from the exact logical checkpoint content plus the project-relative entrypoint, not from compression or manifest bytes, so repeating the same deterministic boundary reuses one artifact. Existing schema-v1 `.json` artifacts remain readable and resumable. + +The writer emits compact JSON through gzip directly into a private same-directory temporary file while computing the stored-byte digest in the same pass. It never constructs a second artifact-sized JSON string in memory, and it checks the final payload-plus-sidecar bytes before publication. The reader bounds stored input and gzip output before parsing and does not create a second decompressed string. Schema v1 is nevertheless still memory-heavy: the compressed buffer, decompressed buffer, one JSON string, and parsed frontier graph can overlap until garbage collection. A malformed gzip stream fails closed as storage corruption; a valid stream still passes the normal schema, stable-ID, source, and configuration checks after decompression. -The writer emits compact JSON through gzip directly into the private same-directory temporary file. It never constructs a second artifact-sized JSON string in memory, and it enforces the single-artifact ceiling against bytes that would actually remain on disk. A malformed gzip stream fails closed as storage corruption; a valid stream still passes the normal schema, stable-ID, source, and configuration checks after decompression. +Before exposing payload bytes, Inkcheck reserves one of 32 fixed, hidden recovery-manifest slots for that stable ID, writes the canonical at-most-64-KiB manifest there, flushes it, and directory-syncs the slot. A same-directory hard-link create is then the exclusive, no-clobber same-ID payload commit point on supported local filesystems, including ordinary NTFS/APFS/ext filesystems: concurrent losers reopen the winning bytes instead of overwriting them. The writer hashes the visible payload with a fixed-size buffer, selects only a recovery record whose size, digest, ID, and requested checkpoint summary match, and promotes that record to the canonical sidecar. Recovery therefore does not need to inflate or parse an artifact that exceeds the schema-v1 readback ceiling. -Inkcheck writes a same-directory temporary file with mode `0600`, flushes it, atomically renames it, and removes temporary files on failure. Only after the new file is durable does retention remove older artifacts. Defaults are hard safety ceilings: +The canonical sidecar is never pre-quarantined. A matching recovery record replaces an older orphan/corrupt sidecar only after this writer sees the published payload; POSIX uses atomic rename-over-existing, while the portable fallback retains bounded promotion and displaced companions so a crash in the replace gap can restore or complete the pair. Each slot has a nonce-bearing owner claim. Reservation and cleanup must first acquire the same fixed per-slot cleaning claim, which is removed last; a delayed cleaner therefore cannot delete a pathname after another writer has reused the slot. A transaction that returns before cleanup durably records that exact nonce as released, allowing a same-process or cross-process retry to reclaim it without confusing another live transaction for debris. Recovery records remain hidden from listing and retention. They are deleted only after the canonical manifest has been reread and matched against a second fixed-memory hash of the same visible payload; a crash before payload publication leaves ignored recovery-only metadata, and a crash after publication leaves enough metadata to finish without decoding. Sidecar reconstruction for legacy schema-v1 JSON also uses this fixed transaction namespace. The claim, release marker, cleaning claim, payload temporary, manifest, promotion, and displaced filenames are fixed and bounded per stable ID. A process that dies while holding a cleaning claim consumes that one slot until an operator removes it while no checkpoint save is active; it cannot expose or overwrite checkpoint evidence. + +This publication guarantee serializes writers for the same stable ID. Retention still validates the complete project set again after the pair is durable, but it is not a global transaction or recovery journal across simultaneous writers of different checkpoint IDs. Hidden files left by an abruptly terminated process are transaction debris, not listed checkpoint artifacts: they are excluded from retention and the project artifact byte ceiling, and the next successful save of that same stable ID cleans dead-owner slots. Repeated crashes across many distinct IDs can therefore leave additional hidden disk use; when no checkpoint save is active, an operator may remove those hidden recovery-slot files. Project-wide crash-debris accounting/recovery belongs with the framed checkpoint-v2 journal rather than this same-ID schema-v1 precursor. Only after the new pair is durable does retention remove older payloads and their sidecars. Defaults are hard safety ceilings for final checkpoint artifacts: - 512 MiB for one checkpoint; - 1 GiB across checkpoint artifacts in one project; - three generations per entrypoint. -An individually oversized compressed checkpoint is rejected. Once a new generation is durable, oldest generations for that entrypoint are removed first, then the oldest project checkpoints if needed to satisfy the project byte ceiling. The saved generation is protected from that cleanup. `checkpoints list/show` reports `storageEncoding` and the actual durable `sizeBytes`; this is storage cost, not an estimate of process heap or future search value. +An individually oversized payload-plus-sidecar pair is rejected. Once a new generation is durable, oldest generations for that entrypoint are removed first, then the oldest project checkpoints if needed to satisfy the project byte ceiling. The saved generation is protected from that cleanup. `checkpoints list/show` reports `payloadSizeBytes`, `metadataSizeBytes`, and their sum as the actual durable `sizeBytes`; this is storage cost, not an estimate of process heap or future search value. + +## Schema-v1 readback boundary + +This is the safe foundation for the observed 600,000-state boundary, not a claim that every such checkpoint can now resume. A gzip payload can be within the durable disk quota while its single logical JSON value is larger than V8 can represent. Inkcheck now returns `resource_limit` at that boundary, keeps the known-good bytes intact, and can still list/prune a manifested artifact without inflation. It does not misreport the file as corrupt or retry an unsafe allocation. A repeated same-ID save may recognize such an artifact only after its canonical manifest matches the requested checkpoint summary and its full stored-byte digest verifies; that preserves known bytes but does not claim the logical payload was decoded or resumable. + +Removing that format ceiling requires a framed artifact schema v2: independently bounded metadata and frontier frames, per-frame lengths/checksums, incremental decode into the resume structures, and a compatibility reader that leaves schema-v1 IDs and exact trajectories unchanged. Promotion should require split-run equality against uninterrupted search plus truncated-frame, oversized-frame, checksum, and mixed-v1/v2 retention tests. ## Privacy diff --git a/docs/shared-checkpoint-schema-v1.md b/docs/shared-checkpoint-schema-v1.md index 9da17ff..aed2f2b 100644 --- a/docs/shared-checkpoint-schema-v1.md +++ b/docs/shared-checkpoint-schema-v1.md @@ -44,8 +44,10 @@ See [local resumable checkpoints](local-checkpoints.md) for freshness, privacy, The logical checkpoint remains schema v1 JSON. New local artifacts stream that JSON through gzip as `.json.gz`; stable IDs still hash the same uncompressed logical checkpoint, and readers continue to accept earlier plain `.json` artifacts. Storage compression therefore changes neither deterministic frontier order nor split-run equivalence. +New artifacts have a small canonical self-checksummed metadata sidecar, allowing list and retention operations to read bounded metadata plus payload file size without opening the frontier. That checksum detects isolated field corruption but is not authentication against a hostile local writer who can recompute it. Opening and resuming remain full-integrity boundaries: manifest validation and the full stored-payload digest precede bounded decompression, then the original stable-ID, envelope, configuration, and freshness checks run. Read failures explicitly distinguish corruption, unsupported schemas, and an artifact that cannot fit the schema-v1 readback envelope. + ## Deliberate limits -Schema v1 supports only base `shared:deep-novelty-v1`; assertions, goals, variable-aware steering, goal-aware steering, and the default portfolio are rejected rather than resumed approximately. Hosted/MCP resume, frontier partitioning, and cross-version migration remain future work. +Schema v1 supports only base `shared:deep-novelty-v1`; assertions, goals, variable-aware steering, goal-aware steering, and the default portfolio are rejected rather than resumed approximately. Hosted/MCP resume, frontier partitioning, and cross-version migration remain future work. Because v1 is one JSON value, it also cannot safely reopen a payload above the runtime's maximum string length even when gzip keeps the stored file below quota. Its compressed buffer, decompressed buffer, JSON string, and parsed graph may overlap in memory. The reader reports an unsafe boundary as a resource limit and preserves the artifact; framed incremental schema v2 is required to remove both the string ceiling and this peak-memory shape. Checkpoint JSON can contain authored choice text, ending text, variable snapshots, serialized Ink runtime state, and exact witness paths. Treat it as sensitive project data and do not commit checkpoints by default. diff --git a/src/checkpoints.ts b/src/checkpoints.ts index 9eb80f0..79f34ee 100644 --- a/src/checkpoints.ts +++ b/src/checkpoints.ts @@ -1,4 +1,5 @@ import { createHash, randomUUID } from "crypto"; +import { constants as bufferConstants } from "buffer"; import * as fs from "fs"; import * as path from "path"; import { Readable, Transform } from "stream"; @@ -15,6 +16,16 @@ export const CHECKPOINT_ARTIFACT_SCHEMA_VERSION = 1; export const DEFAULT_MAX_CHECKPOINT_BYTES = 512 * 1024 * 1024; export const DEFAULT_MAX_PROJECT_CHECKPOINT_BYTES = 1024 * 1024 * 1024; export const DEFAULT_CHECKPOINT_GENERATIONS = 3; +export const CHECKPOINT_MANIFEST_SCHEMA_VERSION = 1; +export const DEFAULT_MAX_STORED_CHECKPOINT_READ_BYTES = DEFAULT_MAX_CHECKPOINT_BYTES; +export const DEFAULT_MAX_DECOMPRESSED_CHECKPOINT_READ_BYTES = Math.min( + DEFAULT_MAX_CHECKPOINT_BYTES, + bufferConstants.MAX_STRING_LENGTH +); + +const MAX_CHECKPOINT_MANIFEST_BYTES = 64 * 1024; +const MAX_CHECKPOINT_RECOVERY_MANIFESTS = 32; +const CHECKPOINT_RECOVERY_SLOT_WIDTH = String(MAX_CHECKPOINT_RECOVERY_MANIFESTS - 1).length; export type CheckpointFreshness = "current" | "stale" | "path_changed"; @@ -35,6 +46,8 @@ export interface CheckpointArtifactSummary { engine: string; totalGranted: number; statesExplored: number; + payloadSizeBytes: number; + metadataSizeBytes: number; sizeBytes: number; storageEncoding: "json" | "gzip"; } @@ -45,6 +58,29 @@ export interface CheckpointStorageLimits { maxGenerationsPerEntrypoint?: number; } +export interface CheckpointReadLimits { + maxStoredBytes?: number; + maxDecompressedBytes?: number; +} + +export type CheckpointReadErrorKind = "corrupt" | "resource_limit" | "unsupported"; +export type CheckpointReadStage = "manifest" | "storage" | "decompression" | "json" | "envelope"; + +export class CheckpointReadError extends Error { + payloadVerified = false; + + constructor( + public readonly kind: CheckpointReadErrorKind, + public readonly stage: CheckpointReadStage, + message: string, + public readonly observedBytes?: number, + public readonly limitBytes?: number + ) { + super(message); + this.name = "CheckpointReadError"; + } +} + interface CheckpointArtifact { artifactSchemaVersion: 1; artifactType: "shared-search-checkpoint"; @@ -59,6 +95,88 @@ interface CheckpointArtifact { checkpoint: SharedSearchCheckpoint; } +interface CheckpointArtifactManifest { + manifestSchemaVersion: 1; + artifactSchemaVersion: 1; + artifactType: "shared-search-checkpoint"; + id: string; + createdAt: string; + inkcheckVersion: string; + checkpointSchemaVersion: number; + entrypoint: string; + engine: string; + totalGranted: number; + statesExplored: number; + storageEncoding: "json" | "gzip"; + artifactSizeBytes: number; + artifactSha256: string; + manifestSha256: string; +} + +interface CheckpointRecord extends CheckpointArtifactSummary { + file: string; + manifestFile?: string; +} + +interface LoadedCheckpointArtifact { + artifact: CheckpointArtifact; + file: string; + storageEncoding: "json" | "gzip"; + payloadSizeBytes: number; + payloadSha256: string; +} + +interface StoredCheckpointManifest { + file: string; + raw: string; + manifest: CheckpointArtifactManifest; + sizeBytes: number; +} + +interface RecoveryManifest extends StoredCheckpointManifest { + slot: number; +} + +interface RecoveredCheckpointPair { + manifest: CheckpointArtifactManifest; + metadataSizeBytes: number; + payloadSizeBytes: number; + payloadSha256: string; +} + +interface CheckpointPayloadDigest { + sizeBytes: number; + sha256: string; + device: number; + inode: number; +} + +interface CheckpointRecoveryClaim { + schemaVersion: 1; + pid: number; + nonce?: string; +} + +interface CheckpointRecoveryRelease { + schemaVersion: 1; + nonce: string; +} + +interface CheckpointRecoveryCleaningClaim { + schemaVersion: 1; + pid: number; + nonce: string; +} + +interface CheckpointTransaction { + slot: number; + nonce: string; + temporary: string; +} + +const activeCheckpointTransactions = new Set(); +const retiredCheckpointTransactions = new Set(); + function checkpointsDirectory(projectRoot: string): string { return path.join(path.resolve(projectRoot), ".inkcheck", "checkpoints"); } @@ -78,6 +196,46 @@ function checkpointDestination(projectRoot: string, id: string): string { return path.join(checkpointsDirectory(projectRoot), `${id}.json.gz`); } +function checkpointManifestFile(projectRoot: string, id: string): string { + validateId(id); + return path.join(checkpointsDirectory(projectRoot), `${id}.meta.json`); +} + +function checkpointRecoveryManifestFile(projectRoot: string, id: string, slot: number): string { + validateId(id); + if (!Number.isSafeInteger(slot) || slot < 0 || slot >= MAX_CHECKPOINT_RECOVERY_MANIFESTS) { + throw new RangeError("checkpoint recovery manifest slot is out of range"); + } + return path.join( + checkpointsDirectory(projectRoot), + `.${id}.recovery-${String(slot).padStart(CHECKPOINT_RECOVERY_SLOT_WIDTH, "0")}.meta.json` + ); +} + +function checkpointRecoveryClaimFile(projectRoot: string, id: string, slot: number): string { + return `${checkpointRecoveryManifestFile(projectRoot, id, slot)}.claim`; +} + +function checkpointRecoveryReleaseFile(projectRoot: string, id: string, slot: number): string { + return `${checkpointRecoveryClaimFile(projectRoot, id, slot)}.released`; +} + +function checkpointRecoveryCleaningFile(projectRoot: string, id: string, slot: number): string { + return `${checkpointRecoveryClaimFile(projectRoot, id, slot)}.cleaning`; +} + +function checkpointRecoveryPromotionFile(projectRoot: string, id: string, slot: number): string { + return `${checkpointRecoveryManifestFile(projectRoot, id, slot)}.promote`; +} + +function checkpointRecoveryDisplacedFile(projectRoot: string, id: string, slot: number): string { + return `${checkpointRecoveryManifestFile(projectRoot, id, slot)}.displaced`; +} + +function checkpointRecoveryPayloadTemporaryFile(projectRoot: string, id: string, slot: number): string { + return `${checkpointRecoveryManifestFile(projectRoot, id, slot)}.payload.tmp`; +} + function checkpointFile(projectRoot: string, id: string): string { validateId(id); const compressed = checkpointDestination(projectRoot, id); @@ -162,17 +320,264 @@ function checkpointId(entrypoint: string, checkpoint: SharedSearchCheckpoint): s return `checkpoint-${hash.digest("hex").slice(0, 24)}`; } +function corrupt(stage: CheckpointReadStage, message: string): CheckpointReadError { + return new CheckpointReadError("corrupt", stage, message); +} + +function unsupported(stage: CheckpointReadStage, message: string): CheckpointReadError { + return new CheckpointReadError("unsupported", stage, message); +} + +function resourceLimit( + stage: CheckpointReadStage, + message: string, + observedBytes: number | undefined, + limitBytes: number +): CheckpointReadError { + return new CheckpointReadError("resource_limit", stage, message, observedBytes, limitBytes); +} + +function validateStoredEntrypoint(projectRoot: string, entrypoint: string, stage: "manifest" | "envelope"): void { + try { + sourcePath(projectRoot, entrypoint); + } catch { + throw corrupt(stage, "checkpoint metadata contains an invalid project-relative entrypoint"); + } +} + +function checkpointReadLimits(input: CheckpointReadLimits): Required { + const requestedStored = input.maxStoredBytes ?? DEFAULT_MAX_STORED_CHECKPOINT_READ_BYTES; + const requestedDecompressed = input.maxDecompressedBytes ?? DEFAULT_MAX_DECOMPRESSED_CHECKPOINT_READ_BYTES; + for (const [name, value] of Object.entries({ + maxStoredBytes: requestedStored, + maxDecompressedBytes: requestedDecompressed, + })) { + if (!Number.isSafeInteger(value) || value < 1) throw new RangeError(`${name} must be a positive safe integer`); + } + // Schema v1 is one JSON value. Even an explicitly larger caller limit cannot + // make V8 construct a string beyond its platform ceiling. + return { + maxStoredBytes: Math.min(requestedStored, bufferConstants.MAX_LENGTH), + maxDecompressedBytes: Math.min(requestedDecompressed, bufferConstants.MAX_STRING_LENGTH), + }; +} + +function readBoundedBuffer( + file: string, + limitBytes: number, + stage: CheckpointReadStage, + description: string +): Buffer { + const fd = fs.openSync(file, "r"); + try { + const stat = fs.fstatSync(fd); + if (!stat.isFile()) { + throw corrupt(stage, `${description} must be a regular file`); + } + const size = stat.size; + if (size > limitBytes) { + throw resourceLimit( + stage, + `${description} is ${size} bytes, above the ${limitBytes}-byte readback limit`, + size, + limitBytes + ); + } + const value = Buffer.allocUnsafe(size); + let offset = 0; + while (offset < size) { + const read = fs.readSync(fd, value, offset, size - offset, offset); + if (read === 0) break; + offset += read; + } + const extra = Buffer.allocUnsafe(1); + if (fs.readSync(fd, extra, 0, 1, offset) !== 0) { + throw corrupt(stage, `${description} changed while it was being read; retry from a stable copy`); + } + return offset === size ? value : value.subarray(0, offset); + } finally { + fs.closeSync(fd); + } +} + +function readLegacyJsonBuffer( + file: string, + limits: Required +): Buffer { + const fd = fs.openSync(file, "r"); + try { + const stat = fs.fstatSync(fd); + if (!stat.isFile()) { + throw corrupt("storage", "stored checkpoint artifact must be a regular file"); + } + const size = stat.size; + if (size > limits.maxStoredBytes) { + throw resourceLimit( + "storage", + `stored checkpoint artifact is ${size} bytes, above the ${limits.maxStoredBytes}-byte readback limit`, + size, + limits.maxStoredBytes + ); + } + if (size > limits.maxDecompressedBytes) { + throw resourceLimit( + "decompression", + `schema-v1 JSON checkpoint is ${size} bytes, above the ${limits.maxDecompressedBytes}-byte readback limit`, + size, + limits.maxDecompressedBytes + ); + } + const value = Buffer.allocUnsafe(size); + let offset = 0; + while (offset < size) { + const read = fs.readSync(fd, value, offset, size - offset, offset); + if (read === 0) break; + offset += read; + } + const extra = Buffer.allocUnsafe(1); + if (fs.readSync(fd, extra, 0, 1, offset) !== 0) { + const observed = Math.max(size + 1, fs.fstatSync(fd).size); + if (observed > limits.maxStoredBytes) { + throw resourceLimit( + "storage", + `stored checkpoint artifact grew to ${observed} bytes, above the ${limits.maxStoredBytes}-byte readback limit`, + observed, + limits.maxStoredBytes + ); + } + if (observed > limits.maxDecompressedBytes) { + throw resourceLimit( + "decompression", + `schema-v1 JSON checkpoint grew to ${observed} bytes, above the ${limits.maxDecompressedBytes}-byte readback limit`, + observed, + limits.maxDecompressedBytes + ); + } + throw corrupt( + "storage", + "stored checkpoint artifact changed while it was being read; retry from a stable copy" + ); + } + return offset === size ? value : value.subarray(0, offset); + } finally { + fs.closeSync(fd); + } +} + +function decompressCheckpointJson(compressed: Buffer, limits: Required): string { + let decompressed: Buffer; + try { + decompressed = gunzipSync(compressed, { maxOutputLength: limits.maxDecompressedBytes }); + } catch (error) { + if ((error as NodeJS.ErrnoException).code === "ERR_BUFFER_TOO_LARGE") { + throw resourceLimit( + "decompression", + `decompressed checkpoint exceeds the ${limits.maxDecompressedBytes}-byte schema-v1 readback limit; a framed checkpoint format is required to reopen it safely`, + undefined, + limits.maxDecompressedBytes + ); + } + throw corrupt("decompression", "checkpoint artifact is corrupt gzip; remove it or restore a valid copy before reopening it"); + } + try { + return decompressed.toString("utf8"); + } catch (error) { + if (error instanceof RangeError) { + throw resourceLimit( + "decompression", + `decompressed checkpoint exceeds the ${limits.maxDecompressedBytes}-byte schema-v1 string limit; a framed checkpoint format is required to reopen it safely`, + decompressed.length, + limits.maxDecompressedBytes + ); + } + throw error; + } +} + +function fileDigest( + file: string, + maxStoredBytes = DEFAULT_MAX_STORED_CHECKPOINT_READ_BYTES +): CheckpointPayloadDigest { + const fd = fs.openSync(file, "r"); + try { + const stat = fs.fstatSync(fd); + if (!stat.isFile()) { + throw corrupt("storage", "checkpoint artifact must be a regular file"); + } + const sizeBytes = stat.size; + if (sizeBytes > maxStoredBytes) { + throw resourceLimit( + "storage", + `stored checkpoint artifact is ${sizeBytes} bytes, above the ${maxStoredBytes}-byte checksum limit`, + sizeBytes, + maxStoredBytes + ); + } + const hash = createHash("sha256"); + const buffer = Buffer.allocUnsafe(Math.max(1, Math.min(1024 * 1024, sizeBytes))); + let offset = 0; + while (offset < sizeBytes) { + const read = fs.readSync(fd, buffer, 0, Math.min(buffer.length, sizeBytes - offset), offset); + if (read === 0) { + throw corrupt("storage", "checkpoint artifact changed while its metadata checksum was being read"); + } + hash.update(buffer.subarray(0, read)); + offset += read; + } + if (fs.readSync(fd, buffer, 0, 1, offset) !== 0) { + const observed = Math.max(sizeBytes + 1, fs.fstatSync(fd).size); + if (observed > maxStoredBytes) { + throw resourceLimit( + "storage", + `stored checkpoint artifact grew to ${observed} bytes, above the ${maxStoredBytes}-byte checksum limit`, + observed, + maxStoredBytes + ); + } + throw corrupt("storage", "checkpoint artifact changed while its metadata checksum was being read"); + } + return { + sizeBytes, + sha256: hash.digest("hex"), + device: stat.dev, + inode: stat.ino, + }; + } finally { + fs.closeSync(fd); + } +} + +function samePayloadDigest(left: CheckpointPayloadDigest, right: CheckpointPayloadDigest): boolean { + const identityMatches = left.inode === 0 || right.inode === 0 + || (left.device === right.device && left.inode === right.inode); + return identityMatches + && left.sizeBytes === right.sizeBytes + && left.sha256 === right.sha256; +} + function parseArtifact(raw: string, expectedId?: string): CheckpointArtifact { let value: unknown; try { value = JSON.parse(raw); } catch { - throw new Error("checkpoint artifact is corrupt JSON; remove it or restore a valid copy before reopening it"); + throw corrupt("json", "checkpoint artifact is corrupt JSON; remove it or restore a valid copy before reopening it"); + } + if (!value || typeof value !== "object") { + throw corrupt("envelope", "checkpoint artifact must be a JSON object"); } - if (!value || typeof value !== "object") throw new Error("checkpoint artifact must be a JSON object"); const artifact = value as Partial; if (artifact.artifactSchemaVersion !== CHECKPOINT_ARTIFACT_SCHEMA_VERSION) { - throw new Error(`unsupported checkpoint artifact schema ${String(artifact.artifactSchemaVersion)}; use a compatible Inkcheck version or migrate the artifact`); + throw unsupported( + "envelope", + `unsupported checkpoint artifact schema ${String(artifact.artifactSchemaVersion)}; use a compatible Inkcheck version or migrate the artifact` + ); + } + if (typeof artifact.checkpointSchemaVersion === "number" + && artifact.checkpointSchemaVersion !== SHARED_SEARCH_CHECKPOINT_SCHEMA_VERSION) { + throw unsupported( + "envelope", + `unsupported shared checkpoint schema ${String(artifact.checkpointSchemaVersion)}; use a compatible Inkcheck version or migrate the checkpoint` + ); } if (artifact.artifactType !== "shared-search-checkpoint" || typeof artifact.id !== "string" || typeof artifact.createdAt !== "string" || !Number.isFinite(Date.parse(artifact.createdAt)) @@ -187,45 +592,89 @@ function parseArtifact(raw: string, expectedId?: string): CheckpointArtifact { || !artifact.checkpoint.state || typeof artifact.checkpoint.state !== "object" || !Number.isSafeInteger(artifact.checkpoint.state.totalGranted) || !Number.isSafeInteger(artifact.checkpoint.state.statesExplored)) { - throw new Error("checkpoint artifact is missing required metadata; regenerate it with Inkcheck"); + throw corrupt("envelope", "checkpoint artifact is missing required metadata; regenerate it with Inkcheck"); } if (artifact.checkpoint.schemaVersion !== SHARED_SEARCH_CHECKPOINT_SCHEMA_VERSION) { - throw new Error(`unsupported shared checkpoint schema ${String(artifact.checkpoint.schemaVersion)}; use a compatible Inkcheck version or migrate the checkpoint`); + throw unsupported( + "envelope", + `unsupported shared checkpoint schema ${String(artifact.checkpoint.schemaVersion)}; use a compatible Inkcheck version or migrate the checkpoint` + ); } const actualId = checkpointId(artifact.source.entrypoint, artifact.checkpoint); if (artifact.id !== actualId || (expectedId !== undefined && artifact.id !== expectedId)) { - throw new Error("checkpoint artifact content does not match its stable ID; restore or regenerate the artifact"); + throw corrupt("envelope", "checkpoint artifact content does not match its stable ID; restore or regenerate the artifact"); } const configuration = artifact.checkpoint.configuration; if (artifact.storySha256 !== configuration.storySha256 || artifact.knotsSha256 !== configuration.knotsSha256 || JSON.stringify(artifact.configuration) !== JSON.stringify(configuration)) { - throw new Error("checkpoint artifact metadata does not match its saved frontier; restore or regenerate the artifact"); + throw corrupt("envelope", "checkpoint artifact metadata does not match its saved frontier; restore or regenerate the artifact"); } return artifact as CheckpointArtifact; } -function loadArtifact(projectRoot: string, id: string): CheckpointArtifact { +function loadArtifactDetailed( + projectRoot: string, + id: string, + inputLimits: CheckpointReadLimits = {} +): LoadedCheckpointArtifact { const file = checkpointFile(projectRoot, id); if (!fs.existsSync(file)) throw new Error(`checkpoint not found: ${id}`); - if (file.endsWith(".gz")) { - let raw: string; + const limits = checkpointReadLimits(inputLimits); + const storageEncoding = file.endsWith(".gz") ? "gzip" : "json"; + let storedManifest: CheckpointArtifactManifest | undefined; + const manifestFile = checkpointManifestFile(projectRoot, id); + if (fs.existsSync(manifestFile)) { + storedManifest = readCheckpointManifest(projectRoot, id).manifest; + validateManifestStorage(projectRoot, id, file, storedManifest); + } + let raw: string; + let payloadSizeBytes: number; + let payloadSha256: string; + if (storageEncoding === "gzip") { + const stored = readBoundedBuffer(file, limits.maxStoredBytes, "storage", "stored checkpoint artifact"); + payloadSizeBytes = stored.length; + payloadSha256 = createHash("sha256").update(stored).digest("hex"); + if (storedManifest && storedManifest.artifactSha256 !== payloadSha256) { + throw corrupt("storage", "checkpoint artifact bytes do not match its metadata checksum; restore or regenerate it"); + } try { - raw = gunzipSync(fs.readFileSync(file)).toString("utf8"); - } catch { - throw new Error("checkpoint artifact is corrupt gzip; remove it or restore a valid copy before reopening it"); + raw = decompressCheckpointJson(stored, limits); + } catch (error) { + if (storedManifest && error instanceof CheckpointReadError && error.kind === "resource_limit") { + error.payloadVerified = true; + } + throw error; + } + } else { + const stored = readLegacyJsonBuffer(file, limits); + raw = stored.toString("utf8"); + payloadSizeBytes = stored.length; + payloadSha256 = createHash("sha256").update(stored).digest("hex"); + if (storedManifest && storedManifest.artifactSha256 !== payloadSha256) { + throw corrupt("storage", "checkpoint artifact bytes do not match its metadata checksum; restore or regenerate it"); } - return parseArtifact(raw, id); } - return parseArtifact(fs.readFileSync(file, "utf8"), id); + const artifact = parseArtifact(raw, id); + validateStoredEntrypoint(projectRoot, artifact.source.entrypoint, "envelope"); + const loaded: LoadedCheckpointArtifact = { + artifact, + file, + storageEncoding, + payloadSizeBytes, + payloadSha256, + }; + validateManifestForLoaded(projectRoot, loaded, storedManifest); + return loaded; } -function summary(projectRoot: string, artifact: CheckpointArtifact): CheckpointArtifactSummary { - const file = checkpointFile(projectRoot, artifact.id); - const storageEncoding = file.endsWith(".gz") ? "gzip" : "json"; +function summary(projectRoot: string, loaded: LoadedCheckpointArtifact): CheckpointArtifactSummary { + const { artifact } = loaded; + const manifestFile = checkpointManifestFile(projectRoot, artifact.id); + const metadataSizeBytes = fs.existsSync(manifestFile) ? fs.statSync(manifestFile).size : 0; return { id: artifact.id, - path: checkpointRelativePath(artifact.id, storageEncoding), + path: checkpointRelativePath(artifact.id, loaded.storageEncoding), artifactType: "shared-search-checkpoint", createdAt: artifact.createdAt, inkcheckVersion: artifact.inkcheckVersion, @@ -234,23 +683,243 @@ function summary(projectRoot: string, artifact: CheckpointArtifact): CheckpointA engine: artifact.checkpoint.engine, totalGranted: artifact.checkpoint.state.totalGranted, statesExplored: artifact.checkpoint.state.statesExplored, - sizeBytes: fs.statSync(file).size, + payloadSizeBytes: loaded.payloadSizeBytes, + metadataSizeBytes, + sizeBytes: loaded.payloadSizeBytes + metadataSizeBytes, + storageEncoding: loaded.storageEncoding, + }; +} + +function manifestForArtifact( + artifact: CheckpointArtifact, + storageEncoding: "json" | "gzip", + artifactSizeBytes: number, + artifactSha256: string +): CheckpointArtifactManifest { + const body: Omit = { + manifestSchemaVersion: CHECKPOINT_MANIFEST_SCHEMA_VERSION, + artifactSchemaVersion: CHECKPOINT_ARTIFACT_SCHEMA_VERSION, + artifactType: "shared-search-checkpoint", + id: artifact.id, + createdAt: artifact.createdAt, + inkcheckVersion: artifact.inkcheckVersion, + checkpointSchemaVersion: artifact.checkpointSchemaVersion, + entrypoint: artifact.source.entrypoint, + engine: artifact.checkpoint.engine, + totalGranted: artifact.checkpoint.state.totalGranted, + statesExplored: artifact.checkpoint.state.statesExplored, + storageEncoding, + artifactSizeBytes, + artifactSha256, + }; + return { ...body, manifestSha256: manifestDigest(body) }; +} + +function manifestBody( + manifest: CheckpointArtifactManifest +): Omit { + // Property order is the manifest schema's canonical byte order. Keeping this + // explicit also ensures newly added fields cannot silently escape binding. + return { + manifestSchemaVersion: manifest.manifestSchemaVersion, + artifactSchemaVersion: manifest.artifactSchemaVersion, + artifactType: manifest.artifactType, + id: manifest.id, + createdAt: manifest.createdAt, + inkcheckVersion: manifest.inkcheckVersion, + checkpointSchemaVersion: manifest.checkpointSchemaVersion, + entrypoint: manifest.entrypoint, + engine: manifest.engine, + totalGranted: manifest.totalGranted, + statesExplored: manifest.statesExplored, + storageEncoding: manifest.storageEncoding, + artifactSizeBytes: manifest.artifactSizeBytes, + artifactSha256: manifest.artifactSha256, + }; +} + +function manifestDigest( + manifest: Omit +): string { + return createHash("sha256").update(JSON.stringify(manifest)).digest("hex"); +} + +function serializedManifest(manifest: CheckpointArtifactManifest): string { + return JSON.stringify(manifest); +} + +function parseManifest(raw: string, expectedId: string): CheckpointArtifactManifest { + let value: unknown; + try { + value = JSON.parse(raw); + } catch { + throw corrupt("manifest", "checkpoint metadata manifest is corrupt JSON; restore or regenerate it"); + } + if (!value || typeof value !== "object") { + throw corrupt("manifest", "checkpoint metadata manifest must be a JSON object"); + } + const manifest = value as Partial; + if (typeof manifest.manifestSchemaVersion === "number" + && manifest.manifestSchemaVersion !== CHECKPOINT_MANIFEST_SCHEMA_VERSION) { + throw unsupported( + "manifest", + `unsupported checkpoint metadata manifest schema ${String(manifest.manifestSchemaVersion)}; use a compatible Inkcheck version` + ); + } + const expectedKeys: Array = [ + "manifestSchemaVersion", "artifactSchemaVersion", "artifactType", "id", "createdAt", + "inkcheckVersion", "checkpointSchemaVersion", "entrypoint", "engine", "totalGranted", + "statesExplored", "storageEncoding", "artifactSizeBytes", "artifactSha256", "manifestSha256", + ]; + const actualKeys = Object.keys(manifest).sort(); + if (actualKeys.length !== expectedKeys.length + || actualKeys.some((key, index) => key !== [...expectedKeys].sort()[index])) { + throw corrupt("manifest", "checkpoint metadata manifest fields do not match its schema"); + } + if (manifest.manifestSchemaVersion !== CHECKPOINT_MANIFEST_SCHEMA_VERSION) { + throw corrupt("manifest", "checkpoint metadata manifest is missing its schema version"); + } + if (manifest.artifactSchemaVersion !== CHECKPOINT_ARTIFACT_SCHEMA_VERSION + || manifest.checkpointSchemaVersion !== SHARED_SEARCH_CHECKPOINT_SCHEMA_VERSION) { + throw unsupported( + "manifest", + `unsupported checkpoint schema in metadata manifest; use a compatible Inkcheck version or migrate the artifact` + ); + } + if (manifest.artifactType !== "shared-search-checkpoint" || manifest.id !== expectedId + || typeof manifest.createdAt !== "string" || !Number.isFinite(Date.parse(manifest.createdAt)) + || typeof manifest.inkcheckVersion !== "string" || typeof manifest.entrypoint !== "string" + || typeof manifest.engine !== "string" + || !Number.isSafeInteger(manifest.totalGranted) || !Number.isSafeInteger(manifest.statesExplored) + || (manifest.storageEncoding !== "json" && manifest.storageEncoding !== "gzip") + || !Number.isSafeInteger(manifest.artifactSizeBytes) || (manifest.artifactSizeBytes ?? 0) < 1 + || typeof manifest.artifactSha256 !== "string" || !/^[0-9a-f]{64}$/.test(manifest.artifactSha256) + || typeof manifest.manifestSha256 !== "string" || !/^[0-9a-f]{64}$/.test(manifest.manifestSha256)) { + throw corrupt("manifest", "checkpoint metadata manifest is missing required fields; restore or regenerate it"); + } + const complete = manifest as CheckpointArtifactManifest; + if (manifestDigest(manifestBody(complete)) !== complete.manifestSha256) { + throw corrupt("manifest", "checkpoint metadata manifest content does not match its canonical checksum"); + } + return complete; +} + +function readCheckpointManifest(projectRoot: string, id: string): { + manifest: CheckpointArtifactManifest; + sizeBytes: number; +} { + const manifestFile = checkpointManifestFile(projectRoot, id); + const stored = readBoundedBuffer( + manifestFile, + MAX_CHECKPOINT_MANIFEST_BYTES, + "manifest", + "checkpoint metadata manifest" + ); + return { manifest: parseManifest(stored.toString("utf8"), id), sizeBytes: stored.length }; +} + +function validateManifestStorage( + projectRoot: string, + id: string, + file: string, + manifest: CheckpointArtifactManifest +): void { + validateStoredEntrypoint(projectRoot, manifest.entrypoint, "manifest"); + const storageEncoding = file.endsWith(".gz") ? "gzip" : "json"; + if (manifest.storageEncoding !== storageEncoding) { + throw corrupt("manifest", `checkpoint ${id} metadata does not match its storage encoding`); + } + if (manifest.artifactSizeBytes !== fs.statSync(file).size) { + throw corrupt( + "storage", + `checkpoint artifact bytes do not match its metadata checksum; content does not match its stable ID or was modified` + ); + } +} + +function verifyManifestPayload( + file: string, + manifest: CheckpointArtifactManifest, + maxStoredBytes: number +): void { + const digest = fileDigest(file, maxStoredBytes); + if (manifest.artifactSha256 !== digest.sha256) { + throw corrupt( + "storage", + `checkpoint artifact bytes do not match its metadata checksum; content does not match its stable ID or was modified` + ); + } +} + +function validateManifestForLoaded( + projectRoot: string, + loaded: LoadedCheckpointArtifact, + alreadyRead?: CheckpointArtifactManifest +): void { + const manifestFile = checkpointManifestFile(projectRoot, loaded.artifact.id); + if (!alreadyRead && !fs.existsSync(manifestFile)) return; + const manifest = alreadyRead ?? readCheckpointManifest(projectRoot, loaded.artifact.id).manifest; + validateManifestStorage(projectRoot, loaded.artifact.id, loaded.file, manifest); + const artifact = loaded.artifact; + if (manifest.artifactSizeBytes !== loaded.payloadSizeBytes + || manifest.artifactSha256 !== loaded.payloadSha256 + || manifest.createdAt !== artifact.createdAt + || manifest.inkcheckVersion !== artifact.inkcheckVersion + || manifest.checkpointSchemaVersion !== artifact.checkpointSchemaVersion + || manifest.entrypoint !== artifact.source.entrypoint + || manifest.engine !== artifact.checkpoint.engine + || manifest.totalGranted !== artifact.checkpoint.state.totalGranted + || manifest.statesExplored !== artifact.checkpoint.state.statesExplored) { + throw corrupt("manifest", "checkpoint metadata manifest does not match its saved payload"); + } +} + +function recordFromManifest(projectRoot: string, id: string, file: string): CheckpointRecord { + const manifestFile = checkpointManifestFile(projectRoot, id); + const { manifest, sizeBytes: metadataSizeBytes } = readCheckpointManifest(projectRoot, id); + validateManifestStorage(projectRoot, id, file, manifest); + const storageEncoding = file.endsWith(".gz") ? "gzip" : "json"; + const payloadSizeBytes = manifest.artifactSizeBytes; + return { + id, + path: checkpointRelativePath(id, storageEncoding), + artifactType: "shared-search-checkpoint", + createdAt: manifest.createdAt, + inkcheckVersion: manifest.inkcheckVersion, + checkpointSchemaVersion: manifest.checkpointSchemaVersion, + entrypoint: manifest.entrypoint, + engine: manifest.engine, + totalGranted: manifest.totalGranted, + statesExplored: manifest.statesExplored, + payloadSizeBytes, + metadataSizeBytes, + sizeBytes: payloadSizeBytes + metadataSizeBytes, storageEncoding, + file, + manifestFile, }; } -function checkpointRecords(projectRoot: string): Array { +function checkpointRecords( + projectRoot: string, + inputLimits: CheckpointReadLimits = {} +): CheckpointRecord[] { const directory = checkpointsDirectory(projectRoot); if (!fs.existsSync(directory)) return []; const names = fs.readdirSync(directory) .filter((name) => /^checkpoint-[0-9a-f]{24}\.json(?:\.gz)?$/.test(name)); const ids = new Set(); return names.map((name) => { - const id = name.slice(0, name.indexOf(".json")); - if (ids.has(id)) throw new Error(`checkpoint ${id} has duplicate JSON and gzip artifacts; retain only one valid copy`); - ids.add(id); - return { ...summary(projectRoot, loadArtifact(projectRoot, id)), file: checkpointFile(projectRoot, id) }; - }); + const id = name.slice(0, name.indexOf(".json")); + if (ids.has(id)) { + throw corrupt("storage", `checkpoint ${id} has duplicate JSON and gzip artifacts; retain only one valid copy`); + } + ids.add(id); + const file = path.join(directory, name); + const manifestFile = checkpointManifestFile(projectRoot, id); + if (fs.existsSync(manifestFile)) return recordFromManifest(projectRoot, id, file); + return { ...summary(projectRoot, loadArtifactDetailed(projectRoot, id, inputLimits)), file }; + }); } export class CheckpointSizeLimitError extends Error { @@ -268,6 +937,7 @@ export class CheckpointSizeLimitError extends Error { class ByteLimitTransform extends Transform { bytes = 0; + private readonly hash = createHash("sha256"); constructor(private readonly kind: "single" | "project", private readonly limit: number) { super(); @@ -279,15 +949,21 @@ class ByteLimitTransform extends Transform { callback(new CheckpointSizeLimitError(this.kind, this.bytes, this.limit)); return; } + this.hash.update(chunk); callback(null, chunk); } + + digest(): string { + return this.hash.digest("hex"); + } } async function writeCompressedArtifact( temporary: string, artifact: CheckpointArtifact, - limits: Required -): Promise { + limits: Required, + precreated = false +): Promise<{ sizeBytes: number; sha256: string }> { const kind = limits.maxCheckpointBytes <= limits.maxProjectBytes ? "single" : "project"; const limit = Math.min(limits.maxCheckpointBytes, limits.maxProjectBytes); const limiter = new ByteLimitTransform(kind, limit); @@ -298,7 +974,7 @@ async function writeCompressedArtifact( // this durable result-window boundary. createGzip({ level: 1 }), limiter, - fs.createWriteStream(temporary, { flags: "wx", mode: 0o600 }) + fs.createWriteStream(temporary, { flags: precreated ? "r+" : "wx", mode: 0o600 }) ); // Windows requires a writable handle for fsync even after the stream has // closed; reopening r+ preserves the same durability step on every platform. @@ -308,6 +984,7 @@ async function writeCompressedArtifact( } finally { fs.closeSync(fd); } + return { sizeBytes: limiter.bytes, sha256: limiter.digest() }; } function oldestFirst(a: T, b: T): number { @@ -321,8 +998,9 @@ function pruneCheckpoints( ): string[] { let records = checkpointRecords(projectRoot); const removed: string[] = []; - const remove = (record: CheckpointArtifactSummary & { file: string }) => { - fs.rmSync(record.file); + const remove = (record: CheckpointRecord) => { + fs.rmSync(record.file, { force: true }); + if (record.manifestFile) fs.rmSync(record.manifestFile, { force: true }); removed.push(record.id); records = records.filter((candidate) => candidate.id !== record.id); }; @@ -370,39 +1048,1002 @@ function syncDirectory(directory: string): void { } } -export async function saveCheckpointArtifact( +function writePrivateManifestFile(file: string, raw: string): void { + if (Buffer.byteLength(raw) > MAX_CHECKPOINT_MANIFEST_BYTES) { + throw new Error("checkpoint metadata manifest exceeded its fixed 64-KiB limit"); + } + fs.writeFileSync(file, raw, { flag: "wx", mode: 0o600 }); + const fd = fs.openSync(file, "r+"); + try { + fs.fsyncSync(fd); + } finally { + fs.closeSync(fd); + } +} + +function recoveryManifestSourceFiles(projectRoot: string, id: string): Array<{ + slot: number; + file: string; +}> { + return Array.from({ length: MAX_CHECKPOINT_RECOVERY_MANIFESTS }, (_unused, slot) => [ + { slot, file: checkpointRecoveryManifestFile(projectRoot, id, slot) }, + { slot, file: checkpointRecoveryPromotionFile(projectRoot, id, slot) }, + ]).flat(); +} + +function hasRecoveryManifest(projectRoot: string, id: string): boolean { + return recoveryManifestSourceFiles(projectRoot, id).some(({ file }) => fs.existsSync(file)); +} + +function readRecoveryClaim(projectRoot: string, id: string, slot: number): CheckpointRecoveryClaim | undefined { + const file = checkpointRecoveryClaimFile(projectRoot, id, slot); + let raw: Buffer; + try { + raw = readBoundedBuffer(file, 1024, "manifest", "checkpoint recovery claim"); + } catch (error) { + if ((error as NodeJS.ErrnoException).code === "ENOENT") return undefined; + throw error; + } + let value: unknown; + try { + value = JSON.parse(raw.toString("utf8")); + } catch { + throw corrupt("manifest", "checkpoint recovery claim is corrupt JSON"); + } + if (!value || typeof value !== "object") { + throw corrupt("manifest", "checkpoint recovery claim must be a JSON object"); + } + const claim = value as Partial; + if (claim.schemaVersion !== 1 || !Number.isSafeInteger(claim.pid) || (claim.pid ?? 0) < 1 + || (claim.nonce !== undefined + && (typeof claim.nonce !== "string" + || !/^[0-9a-f]{8}-[0-9a-f]{4}-4[0-9a-f]{3}-[89ab][0-9a-f]{3}-[0-9a-f]{12}$/.test(claim.nonce)))) { + throw corrupt("manifest", "checkpoint recovery claim fields are invalid"); + } + return claim as CheckpointRecoveryClaim; +} + +function readRecoveryRelease( projectRoot: string, - entrypoint: string, + id: string, + slot: number +): CheckpointRecoveryRelease | undefined { + const file = checkpointRecoveryReleaseFile(projectRoot, id, slot); + let raw: Buffer; + try { + raw = readBoundedBuffer(file, 1024, "manifest", "checkpoint recovery release"); + } catch (error) { + if ((error as NodeJS.ErrnoException).code === "ENOENT") return undefined; + throw error; + } + let value: unknown; + try { + value = JSON.parse(raw.toString("utf8")); + } catch { + throw corrupt("manifest", "checkpoint recovery release is corrupt JSON"); + } + if (!value || typeof value !== "object") { + throw corrupt("manifest", "checkpoint recovery release must be a JSON object"); + } + const release = value as Partial; + if (release.schemaVersion !== 1 || typeof release.nonce !== "string" + || !/^[0-9a-f]{8}-[0-9a-f]{4}-4[0-9a-f]{3}-[89ab][0-9a-f]{3}-[0-9a-f]{12}$/.test(release.nonce)) { + throw corrupt("manifest", "checkpoint recovery release fields are invalid"); + } + return release as CheckpointRecoveryRelease; +} + +function checkpointTransactionKey( + projectRoot: string, + id: string, + slot: number, + nonce: string +): string { + return `${checkpointRecoveryClaimFile(projectRoot, id, slot)}\0${nonce}`; +} + +function processIsAlive(pid: number): boolean { + try { + process.kill(pid, 0); + return true; + } catch (error) { + const code = (error as NodeJS.ErrnoException).code; + if (code === "ESRCH") return false; + return true; + } +} + +function recoveryClaimIsActive( + projectRoot: string, + id: string, + slot: number, + claim: CheckpointRecoveryClaim +): boolean { + if (claim.nonce !== undefined) { + const key = checkpointTransactionKey(projectRoot, id, slot, claim.nonce); + if (retiredCheckpointTransactions.has(key)) return false; + let release: CheckpointRecoveryRelease | undefined; + try { + release = readRecoveryRelease(projectRoot, id, slot); + } catch (error) { + // A malformed unowned release marker cannot prove abandonment. Keep a + // live/unknown owner conservative, but a definitely dead foreign PID is + // still safe to recover and must not wedge portable replacement forever. + if (error instanceof CheckpointReadError) return processIsAlive(claim.pid); + throw error; + } + if (release?.nonce === claim.nonce) return false; + if (activeCheckpointTransactions.has(key)) return true; + // Worker threads share process.pid but have isolate-local module state. + // A current-PID nonce unknown to this isolate therefore remains active + // unless its exact durable release marker was observed above. + if (claim.pid === process.pid) return true; + } + return processIsAlive(claim.pid); +} + +function releaseCheckpointTransaction( + projectRoot: string, + id: string, + transaction: CheckpointTransaction +): void { + const key = checkpointTransactionKey(projectRoot, id, transaction.slot, transaction.nonce); + const claim = readRecoveryClaim(projectRoot, id, transaction.slot); + if (!claim || claim.pid !== process.pid || claim.nonce !== transaction.nonce) { + activeCheckpointTransactions.delete(key); + retiredCheckpointTransactions.delete(key); + return; + } + const releaseFile = checkpointRecoveryReleaseFile(projectRoot, id, transaction.slot); + const raw = JSON.stringify({ schemaVersion: 1, nonce: transaction.nonce }); + let durableRelease = false; + try { + try { + writePrivateManifestFile(releaseFile, raw); + syncDirectory(checkpointsDirectory(projectRoot)); + durableRelease = true; + } catch (error) { + if ((error as NodeJS.ErrnoException).code !== "EEXIST") throw error; + const current = readRecoveryRelease(projectRoot, id, transaction.slot); + if (current?.nonce !== transaction.nonce) { + throw corrupt("manifest", `checkpoint ${id} recovery slot was released by a different transaction`); + } + durableRelease = true; + } + } finally { + // Even a failed durable-release write ends this synchronous operation. + // Same-process retries may reclaim the exact nonce from the registry; + // other processes still conservatively honor the live PID until a marker + // is durable or that process exits. + activeCheckpointTransactions.delete(key); + if (durableRelease) retiredCheckpointTransactions.delete(key); + else retiredCheckpointTransactions.add(key); + } +} + +function acquireRecoverySlotCleaning( + projectRoot: string, + id: string, + slot: number +): string | undefined { + const file = checkpointRecoveryCleaningFile(projectRoot, id, slot); + const claim: CheckpointRecoveryCleaningClaim = { + schemaVersion: 1, + pid: process.pid, + nonce: randomUUID(), + }; + try { + writePrivateManifestFile(file, JSON.stringify(claim)); + syncDirectory(checkpointsDirectory(projectRoot)); + return file; + } catch (error) { + if ((error as NodeJS.ErrnoException).code === "EEXIST") return undefined; + throw error; + } +} + +function releaseRecoverySlotCleaning(projectRoot: string, file: string): void { + try { + fs.rmSync(file); + } catch (error) { + const code = (error as NodeJS.ErrnoException).code; + if (code === "ENOENT") return; + // Any failed release consumes this one fixed slot rather than throwing + // after reservation has already created an active transaction and losing + // the caller's only handle to release it safely. + return; + } + try { + syncDirectory(checkpointsDirectory(projectRoot)); + } catch { + // Namespace sync failure is likewise fail-closed at this slot. It must not + // override a completed reservation or cleanup operation. + } +} + +function recoverySlotHasPrivateState(projectRoot: string, id: string, slot: number): boolean { + return [ + checkpointRecoveryClaimFile(projectRoot, id, slot), + checkpointRecoveryReleaseFile(projectRoot, id, slot), + checkpointRecoveryManifestFile(projectRoot, id, slot), + checkpointRecoveryPromotionFile(projectRoot, id, slot), + checkpointRecoveryDisplacedFile(projectRoot, id, slot), + checkpointRecoveryPayloadTemporaryFile(projectRoot, id, slot), + ].some((file) => fs.existsSync(file)); +} + +function reserveCheckpointTransaction(projectRoot: string, id: string): CheckpointTransaction { + const directory = checkpointsDirectory(projectRoot); + for (let slot = 0; slot < MAX_CHECKPOINT_RECOVERY_MANIFESTS; slot++) { + const cleaning = acquireRecoverySlotCleaning(projectRoot, id, slot); + if (!cleaning) continue; + try { + // A returned failure is normally durably released even while its + // process continues running. Reclaim only that exact release (or an + // exact nonce retired by this isolate). Markerless crash records stay + // intact until a canonical pair has been verified. + let existingClaim: CheckpointRecoveryClaim | undefined; + try { + existingClaim = readRecoveryClaim(projectRoot, id, slot); + } catch (error) { + if (error instanceof CheckpointReadError) continue; + throw error; + } + if (existingClaim) { + let explicitlyReleased = false; + if (existingClaim.nonce !== undefined) { + const key = checkpointTransactionKey(projectRoot, id, slot, existingClaim.nonce); + explicitlyReleased = retiredCheckpointTransactions.has(key); + if (!explicitlyReleased) { + try { + explicitlyReleased = readRecoveryRelease(projectRoot, id, slot)?.nonce === existingClaim.nonce; + } catch (error) { + // A corrupt/torn release marker quarantines this one slot. + if (error instanceof CheckpointReadError) continue; + throw error; + } + } + } + if (!explicitlyReleased) continue; + cleanupRecoverySlotLocked(projectRoot, id, slot); + if (recoverySlotHasPrivateState(projectRoot, id, slot)) continue; + } else { + cleanupRecoverySlotLocked(projectRoot, id, slot); + if (recoverySlotHasPrivateState(projectRoot, id, slot)) continue; + } + + const claimFile = checkpointRecoveryClaimFile(projectRoot, id, slot); + const temporary = checkpointRecoveryPayloadTemporaryFile(projectRoot, id, slot); + const nonce = randomUUID(); + try { + writePrivateManifestFile(claimFile, JSON.stringify({ schemaVersion: 1, pid: process.pid, nonce })); + syncDirectory(directory); + } catch (error) { + if ((error as NodeJS.ErrnoException).code !== "EEXIST") throw error; + continue; + } + try { + fs.writeFileSync(temporary, Buffer.alloc(0), { flag: "wx", mode: 0o600 }); + syncDirectory(directory); + activeCheckpointTransactions.add(checkpointTransactionKey(projectRoot, id, slot, nonce)); + return { slot, nonce, temporary }; + } catch (error) { + fs.rmSync(claimFile, { force: true }); + syncDirectory(directory); + if ((error as NodeJS.ErrnoException).code !== "EEXIST") throw error; + } + } finally { + releaseRecoverySlotCleaning(projectRoot, cleaning); + } + } + throw new CheckpointReadError( + "resource_limit", + "manifest", + `checkpoint ${id} has reached its bounded ${MAX_CHECKPOINT_RECOVERY_MANIFESTS}-file recovery-manifest limit; complete or remove the stale transaction before retrying`, + MAX_CHECKPOINT_RECOVERY_MANIFESTS, + MAX_CHECKPOINT_RECOVERY_MANIFESTS + ); +} + +function writeRecoveryManifest( + projectRoot: string, + id: string, + slot: number, + raw: string +): void { + const file = checkpointRecoveryManifestFile(projectRoot, id, slot); + writePrivateManifestFile(file, raw); + syncDirectory(checkpointsDirectory(projectRoot)); +} + +function readRecoveryManifest(file: string, id: string, slot: number): RecoveryManifest { + const stored = readBoundedBuffer( + file, + MAX_CHECKPOINT_MANIFEST_BYTES, + "manifest", + "checkpoint recovery manifest" + ); + const raw = stored.toString("utf8"); + return { slot, file, raw, manifest: parseManifest(raw, id), sizeBytes: stored.length }; +} + +function manifestMatchesVisiblePayload( + projectRoot: string, + id: string, + file: string, + manifest: CheckpointArtifactManifest, + payload: { sizeBytes: number; sha256: string }, + relative: string, + checkpoint: SharedSearchCheckpoint +): boolean { + const storageEncoding = file.endsWith(".gz") ? "gzip" : "json"; + try { + validateStoredEntrypoint(projectRoot, manifest.entrypoint, "manifest"); + } catch { + return false; + } + return manifest.id === id + && manifest.storageEncoding === storageEncoding + && manifest.artifactSizeBytes === payload.sizeBytes + && manifest.artifactSha256 === payload.sha256 + && expectedManifestMatchesCheckpoint(manifest, relative, checkpoint); +} + +function matchingRecoveryManifest( + projectRoot: string, + id: string, + file: string, + payload: { sizeBytes: number; sha256: string }, + relative: string, + checkpoint: SharedSearchCheckpoint +): RecoveryManifest | undefined { + for (const candidate of recoveryManifestSourceFiles(projectRoot, id)) { + try { + const recovery = readRecoveryManifest(candidate.file, id, candidate.slot); + if (manifestMatchesVisiblePayload( + projectRoot, + id, + file, + recovery.manifest, + payload, + relative, + checkpoint + )) return recovery; + } catch (error) { + if ((error as NodeJS.ErrnoException).code === "ENOENT" || error instanceof CheckpointReadError) continue; + throw error; + } + } + return undefined; +} + +function readMatchingCanonicalManifest( + projectRoot: string, + id: string, + file: string, + payload: { sizeBytes: number; sha256: string }, + relative: string, checkpoint: SharedSearchCheckpoint, - inputLimits: CheckpointStorageLimits = {} + expectedRaw?: string +): StoredCheckpointManifest | undefined { + const manifestFile = checkpointManifestFile(projectRoot, id); + try { + const stored = readBoundedBuffer( + manifestFile, + MAX_CHECKPOINT_MANIFEST_BYTES, + "manifest", + "checkpoint metadata manifest" + ); + const raw = stored.toString("utf8"); + const manifest = parseManifest(raw, id); + // Parse/classify the canonical sidecar before comparing it with a v1 + // recovery record. A valid future schema is authoritative and must surface + // as unsupported rather than being overwritten as an apparent mismatch. + if (expectedRaw !== undefined && raw !== expectedRaw) return undefined; + if (!manifestMatchesVisiblePayload( + projectRoot, + id, + file, + manifest, + payload, + relative, + checkpoint + )) return undefined; + return { file: manifestFile, raw, manifest, sizeBytes: stored.length }; + } catch (error) { + if ((error as NodeJS.ErrnoException).code === "ENOENT") return undefined; + if (error instanceof CheckpointReadError) { + if (error.kind !== "corrupt") throw error; + return undefined; + } + throw error; + } +} + +function promoteRecoveryManifest( + projectRoot: string, + id: string, + file: string, + payload: { sizeBytes: number; sha256: string }, + relative: string, + checkpoint: SharedSearchCheckpoint, + recovery: RecoveryManifest +): StoredCheckpointManifest { + const directory = checkpointsDirectory(projectRoot); + const destination = checkpointManifestFile(projectRoot, id); + const promotion = checkpointRecoveryPromotionFile(projectRoot, id, recovery.slot); + const displaced = checkpointRecoveryDisplacedFile(projectRoot, id, recovery.slot); + const matches = () => readMatchingCanonicalManifest( + projectRoot, + id, + file, + payload, + relative, + checkpoint, + recovery.raw + ); + const alreadyPublished = matches(); + if (alreadyPublished) return alreadyPublished; + + try { + fs.linkSync(recovery.file, destination); + syncDirectory(directory); + } catch (error) { + const code = (error as NodeJS.ErrnoException).code; + const concurrent = matches(); + if (concurrent) return concurrent; + if (code !== "EEXIST" && code !== "ENOENT") throw error; + try { + if (recovery.file !== promotion) { + try { + fs.linkSync(recovery.file, promotion); + syncDirectory(directory); + } catch (linkError) { + const linkCode = (linkError as NodeJS.ErrnoException).code; + if (linkCode === "ENOENT") { + const completed = matches(); + if (completed) return completed; + } + if (linkCode !== "EEXIST") throw linkError; + const held = readRecoveryManifest(promotion, id, recovery.slot); + if (held.raw !== recovery.raw) { + throw corrupt("manifest", `checkpoint ${id} recovery promotion slot contains different metadata`); + } + } + } + try { + // POSIX atomically replaces an orphan/corrupt sidecar here. Some + // Windows filesystems reject rename-over-existing, so the fallback + // below preserves the durable recovery file across its replace gap. + if (process.env.NODE_ENV === "test" + && process.env.INKCHECK_TEST_CHECKPOINT_FORCE_WINDOWS_REPLACE === "1") { + const forced = new Error("simulated Windows rename-over-existing rejection") as NodeJS.ErrnoException; + forced.code = "EPERM"; + throw forced; + } + fs.renameSync(promotion, destination); + } catch (renameError) { + const completed = matches(); + if (completed) return completed; + const renameCode = (renameError as NodeJS.ErrnoException).code; + if (renameCode !== "EEXIST" && renameCode !== "EPERM" && renameCode !== "EACCES") { + throw renameError; + } + checkpointTestBarrier("before-manifest-displacement"); + if (fs.existsSync(displaced)) { + const ownerClaim = readRecoveryClaim(projectRoot, id, recovery.slot); + const ownerIsActive = ownerClaim !== undefined + && recoveryClaimIsActive(projectRoot, id, recovery.slot, ownerClaim); + if (ownerIsActive) { + if (process.env.NODE_ENV === "test" + && process.env.INKCHECK_TEST_CHECKPOINT_JOIN_MARKER) { + fs.writeFileSync(process.env.INKCHECK_TEST_CHECKPOINT_JOIN_MARKER, "joined"); + } + // Another portable promoter owns the same bounded replace gap. + // Join it by publishing from the shared durable promotion link; + // if it wins first, strict canonical reread accepts the same bytes. + try { + fs.linkSync(promotion, destination); + syncDirectory(directory); + } catch (joinError) { + const joinCode = (joinError as NodeJS.ErrnoException).code; + const joined = matches(); + if (joined) return joined; + if (joinCode !== "EEXIST" && joinCode !== "ENOENT") throw joinError; + for (let attempt = 0; attempt < 80; attempt++) { + Atomics.wait(new Int32Array(new SharedArrayBuffer(4)), 0, 0, 25); + const completedJoin = matches(); + if (completedJoin) return completedJoin; + } + throw corrupt( + "manifest", + `checkpoint ${id} bounded manifest replacement did not complete; retry from its recovery record` + ); + } + const joined = matches(); + if (joined) return joined; + throw corrupt( + "manifest", + `checkpoint ${id} joined manifest replacement does not match its visible payload` + ); + } + // A returned failure durably releases its nonce even if that process + // stays alive. Its restored canonical sidecar plus stale displaced + // link is therefore safe to take over instead of being mistaken for + // a live portable promoter forever. + try { + fs.rmSync(displaced); + syncDirectory(directory); + } catch (staleError) { + if ((staleError as NodeJS.ErrnoException).code !== "ENOENT") throw staleError; + } + if (process.env.NODE_ENV === "test" + && process.env.INKCHECK_TEST_CHECKPOINT_STALE_DISPLACED_MARKER) { + fs.writeFileSync(process.env.INKCHECK_TEST_CHECKPOINT_STALE_DISPLACED_MARKER, "taken-over"); + } + } + try { + fs.renameSync(destination, displaced); + syncDirectory(directory); + } catch (displaceError) { + if ((displaceError as NodeJS.ErrnoException).code !== "ENOENT") throw displaceError; + } + checkpointTestBarrier("after-manifest-displacement"); + checkpointTestCrash("after-manifest-displacement"); + try { + // Publish from the already-created private hard link. Another + // promoter may have removed the original recovery pathname while + // this writer held the replace gap. + fs.linkSync(promotion, destination); + syncDirectory(directory); + } catch (publishError) { + const raced = matches(); + if (!raced) { + if (!fs.existsSync(destination) && fs.existsSync(displaced)) { + try { + fs.linkSync(displaced, destination); + syncDirectory(directory); + fs.rmSync(displaced, { force: true }); + syncDirectory(directory); + } catch { + // Preserve both bounded files if restoration races or fails. + } + } + throw publishError; + } + } + } + syncDirectory(directory); + const published = matches(); + if (!published) { + throw corrupt( + "manifest", + `checkpoint ${id} recovery metadata could not be verified after publication` + ); + } + return published; + } catch (promotionError) { + const completed = matches(); + if (completed) return completed; + throw promotionError; + } + } + + const published = matches(); + if (!published) { + throw corrupt("manifest", `checkpoint ${id} recovery metadata does not match its visible payload`); + } + return published; +} + +function cleanupRecoverySlotLocked( + projectRoot: string, + id: string, + slot: number, + expectedTransaction?: CheckpointTransaction +): void { + const directory = checkpointsDirectory(projectRoot); + let claim: CheckpointRecoveryClaim | undefined; + try { + claim = readRecoveryClaim(projectRoot, id, slot); + } catch (error) { + // An unowned corrupt claim cannot be safely attributed or released. + if (!expectedTransaction) return; + throw error; + } + if (!claim) { + // The fixed cleaning claim excludes reservation, so claimless companions + // are abandoned state and cannot be replaced underneath this cleanup. + let removed = false; + for (const candidate of [ + checkpointRecoveryManifestFile(projectRoot, id, slot), + checkpointRecoveryPromotionFile(projectRoot, id, slot), + checkpointRecoveryDisplacedFile(projectRoot, id, slot), + checkpointRecoveryPayloadTemporaryFile(projectRoot, id, slot), + checkpointRecoveryReleaseFile(projectRoot, id, slot), + ]) { + try { + fs.rmSync(candidate); + removed = true; + } catch (error) { + const code = (error as NodeJS.ErrnoException).code; + if (code === "ENOENT" || code === "EPERM" || code === "EBUSY" || code === "EACCES") continue; + throw error; + } + } + if (expectedTransaction) { + const key = checkpointTransactionKey( + projectRoot, + id, + slot, + expectedTransaction.nonce + ); + activeCheckpointTransactions.delete(key); + retiredCheckpointTransactions.delete(key); + } + if (removed) syncDirectory(directory); + return; + } + if (expectedTransaction) { + if (claim.pid !== process.pid || claim.nonce !== expectedTransaction.nonce + || slot !== expectedTransaction.slot) return; + } else if (recoveryClaimIsActive(projectRoot, id, slot, claim)) { + return; + } + + let removed = false; + let complete = true; + for (const candidate of [ + checkpointRecoveryManifestFile(projectRoot, id, slot), + checkpointRecoveryPromotionFile(projectRoot, id, slot), + checkpointRecoveryDisplacedFile(projectRoot, id, slot), + checkpointRecoveryPayloadTemporaryFile(projectRoot, id, slot), + ]) { + try { + fs.rmSync(candidate); + removed = true; + } catch (error) { + const code = (error as NodeJS.ErrnoException).code; + if (code === "ENOENT") continue; + // Windows can refuse to unlink a still-open slot. Keep its claim so the + // pathname cannot be reused until the owner or a later recovery succeeds. + if (code === "EPERM" || code === "EBUSY" || code === "EACCES") { + complete = false; + continue; + } + throw error; + } + } + if (complete) { + try { + fs.rmSync(checkpointRecoveryClaimFile(projectRoot, id, slot)); + removed = true; + if (claim.nonce !== undefined) { + const key = checkpointTransactionKey(projectRoot, id, slot, claim.nonce); + activeCheckpointTransactions.delete(key); + retiredCheckpointTransactions.delete(key); + } + try { + fs.rmSync(checkpointRecoveryReleaseFile(projectRoot, id, slot)); + removed = true; + } catch (releaseError) { + const releaseCode = (releaseError as NodeJS.ErrnoException).code; + if (releaseCode !== "ENOENT" && releaseCode !== "EPERM" + && releaseCode !== "EBUSY" && releaseCode !== "EACCES") throw releaseError; + } + } catch (error) { + const code = (error as NodeJS.ErrnoException).code; + if (code !== "ENOENT" && code !== "EPERM" && code !== "EBUSY" && code !== "EACCES") throw error; + // The release marker, when present, stays beside a claim that could not + // be removed. A later cleanup can retry without misclassifying the live + // PID as a new transaction. + } + } + if (removed) syncDirectory(directory); +} + +function cleanupRecoverySlot( + projectRoot: string, + id: string, + slot: number, + expectedTransaction?: CheckpointTransaction +): void { + const cleaning = acquireRecoverySlotCleaning(projectRoot, id, slot); + if (!cleaning) return; + try { + cleanupRecoverySlotLocked(projectRoot, id, slot, expectedTransaction); + } finally { + releaseRecoverySlotCleaning(projectRoot, cleaning); + } +} + +function cleanupRecoveryManifests( + projectRoot: string, + id: string, + expectedTransaction?: CheckpointTransaction +): void { + for (let slot = 0; slot < MAX_CHECKPOINT_RECOVERY_MANIFESTS; slot++) { + cleanupRecoverySlot( + projectRoot, + id, + slot, + slot === expectedTransaction?.slot ? expectedTransaction : undefined + ); + } +} + +function recoverPublishedCheckpoint( + projectRoot: string, + relative: string, + id: string, + checkpoint: SharedSearchCheckpoint, + limits: Required, + expectedTransaction?: CheckpointTransaction +): RecoveredCheckpointPair | undefined { + // Healthy canonical pairs stay metadata-only here. Probe the fixed slot set + // before opening or hashing a potentially very large payload. + if (!hasRecoveryManifest(projectRoot, id)) return undefined; + const file = checkpointFile(projectRoot, id); + if (!fs.existsSync(file)) return undefined; + const digestLimit = Math.min( + limits.maxCheckpointBytes, + limits.maxProjectBytes + ); + const payload = fileDigest(file, digestLimit); + const recovery = matchingRecoveryManifest( + projectRoot, + id, + file, + payload, + relative, + checkpoint + ); + if (!recovery) return undefined; + promoteRecoveryManifest( + projectRoot, + id, + file, + payload, + relative, + checkpoint, + recovery + ); + // Reopen and hash the payload after manifest publication. This closes the + // prune/replacement window: recovery metadata is retained if the path now + // names different bytes (or a different regular file where inode identity + // is available), even when the first digest happened to match. + const verifiedPayload = fileDigest(file, digestLimit); + if (!samePayloadDigest(payload, verifiedPayload)) { + throw corrupt( + "storage", + `checkpoint ${id} payload changed while its recovery metadata was being published; retry from the durable recovery record` + ); + } + const verifiedManifest = readMatchingCanonicalManifest( + projectRoot, + id, + file, + verifiedPayload, + relative, + checkpoint, + recovery.raw + ); + if (!verifiedManifest) { + throw corrupt( + "manifest", + `checkpoint ${id} metadata changed before its recovered pair could be verified` + ); + } + // The canonical manifest has now been reread, parsed, checksum-checked, and + // matched against the rehashed visible payload. Only this point permits + // recovery cleanup; a crash before it leaves enough durable metadata to retry. + cleanupRecoveryManifests(projectRoot, id, expectedTransaction); + return { + manifest: verifiedManifest.manifest, + metadataSizeBytes: verifiedManifest.sizeBytes, + payloadSizeBytes: verifiedPayload.sizeBytes, + payloadSha256: verifiedPayload.sha256, + }; +} + +type CheckpointTestStage = + | "after-recovery-manifest" + | "after-payload-publication" + | "after-legacy-manifest-temporary" + | "before-manifest-displacement" + | "after-manifest-displacement"; + +function checkpointTestBarrier(stage: CheckpointTestStage): void { + const gate = process.env.INKCHECK_TEST_CHECKPOINT_GATE; + if (process.env.NODE_ENV !== "test" || !gate + || process.env.INKCHECK_TEST_CHECKPOINT_GATE_STAGE !== stage) return; + fs.writeFileSync(`${gate}.${process.pid}.ready`, stage); + const deadline = Date.now() + 15_000; + while (!fs.existsSync(gate)) { + if (Date.now() >= deadline) { + throw new Error(`timed out waiting for checkpoint test gate at ${stage}`); + } + Atomics.wait(new Int32Array(new SharedArrayBuffer(4)), 0, 0, 25); + } +} + +function checkpointTestCrash(stage: CheckpointTestStage): void { + if (process.env.NODE_ENV === "test" + && process.env.INKCHECK_TEST_CHECKPOINT_CRASH_STAGE === stage) { + process.exit(86); + } +} + +function publishCheckpointManifest( + projectRoot: string, + id: string, + raw: string +): void { + const directory = checkpointsDirectory(projectRoot); + const destination = checkpointManifestFile(projectRoot, id); + const bytes = Buffer.from(raw); + if (bytes.length > MAX_CHECKPOINT_MANIFEST_BYTES) { + throw new Error("checkpoint metadata manifest exceeded its fixed 64-KiB limit"); + } + // Sidecar-free v1 compatibility publication shares the same fixed 32-slot + // namespace as new checkpoint transactions. A crash can therefore leave at + // most one bounded private file per claimed slot, never UUID-named debris. + const transaction = reserveCheckpointTransaction(projectRoot, id); + try { + const fd = fs.openSync(transaction.temporary, "r+"); + try { + fs.ftruncateSync(fd, 0); + let offset = 0; + while (offset < bytes.length) { + offset += fs.writeSync(fd, bytes, offset, bytes.length - offset, offset); + } + fs.fsyncSync(fd); + } finally { + fs.closeSync(fd); + } + syncDirectory(directory); + checkpointTestCrash("after-legacy-manifest-temporary"); + try { + fs.linkSync(transaction.temporary, destination); + syncDirectory(directory); + } catch (error) { + if ((error as NodeJS.ErrnoException).code !== "EEXIST") throw error; + const current = readBoundedBuffer( + destination, + MAX_CHECKPOINT_MANIFEST_BYTES, + "manifest", + "checkpoint metadata manifest" + ).toString("utf8"); + if (current !== raw) { + throw corrupt("manifest", `checkpoint ${id} already has different metadata; reopen it before saving`); + } + } + } finally { + releaseCheckpointTransaction(projectRoot, id, transaction); + cleanupRecoverySlot(projectRoot, id, transaction.slot, transaction); + } +} + +function enforceDurableCheckpointLimits( + bytes: number, + limits: Required +): void { + if (bytes > limits.maxCheckpointBytes) { + throw new CheckpointSizeLimitError("single", bytes, limits.maxCheckpointBytes); + } + if (bytes > limits.maxProjectBytes) { + throw new CheckpointSizeLimitError("project", bytes, limits.maxProjectBytes); + } +} + +function expectedManifestMatchesCheckpoint( + manifest: CheckpointArtifactManifest, + relative: string, + checkpoint: SharedSearchCheckpoint +): boolean { + return manifest.entrypoint === relative + && manifest.engine === checkpoint.engine + && manifest.totalGranted === checkpoint.state.totalGranted + && manifest.statesExplored === checkpoint.state.statesExplored; +} + +async function reuseCheckpointArtifact( + root: string, + relative: string, + id: string, + checkpoint: SharedSearchCheckpoint, + limits: Required, + expectedTransaction?: CheckpointTransaction ): Promise { - const root = path.resolve(projectRoot); - const relative = relativeEntrypoint(root, entrypoint); - const id = checkpointId(relative, checkpoint); const directory = checkpointsDirectory(root); - const destination = checkpointDestination(root, id); - const limits = storageLimits(inputLimits); - fs.mkdirSync(directory, { recursive: true, mode: 0o700 }); - if (process.platform !== "win32") fs.chmodSync(directory, 0o700); const existing = checkpointFile(root, id); - if (fs.existsSync(existing)) { - loadArtifact(root, id); - if (process.platform !== "win32") fs.chmodSync(existing, 0o600); - const bytes = fs.statSync(existing).size; - if (bytes > limits.maxCheckpointBytes) { - throw new Error(`checkpoint is ${bytes} bytes, above the ${limits.maxCheckpointBytes}-byte single-checkpoint limit`); - } - if (bytes > limits.maxProjectBytes) { - throw new Error(`checkpoint is ${bytes} bytes, above the ${limits.maxProjectBytes}-byte project checkpoint quota`); + let sizeBytes: number; + const recovered = recoverPublishedCheckpoint(root, relative, id, checkpoint, limits, expectedTransaction); + if (recovered) { + sizeBytes = recovered.payloadSizeBytes + recovered.metadataSizeBytes; + } else { + try { + const loaded = loadArtifactDetailed(root, id); + const manifestFile = checkpointManifestFile(root, id); + if (!fs.existsSync(manifestFile)) { + const manifest = manifestForArtifact( + loaded.artifact, + loaded.storageEncoding, + loaded.payloadSizeBytes, + loaded.payloadSha256 + ); + const raw = serializedManifest(manifest); + enforceDurableCheckpointLimits(loaded.payloadSizeBytes + Buffer.byteLength(raw), limits); + publishCheckpointManifest(root, id, raw); + } + sizeBytes = summary(root, loaded).sizeBytes; + } catch (error) { + // Schema v1 may exceed the process string ceiling even though its stored + // bytes are durable. Reuse is allowed only when the canonical manifest + // matches this exact requested checkpoint and a fixed-memory full payload + // digest verifies the otherwise-unparseable bytes. + if (!(error instanceof CheckpointReadError) || error.kind !== "resource_limit") throw error; + if (!fs.existsSync(checkpointManifestFile(root, id))) throw error; + const { manifest, sizeBytes: metadataSizeBytes } = readCheckpointManifest(root, id); + validateManifestStorage(root, id, existing, manifest); + if (!expectedManifestMatchesCheckpoint(manifest, relative, checkpoint)) { + throw corrupt("manifest", "checkpoint metadata manifest does not match the requested saved frontier"); + } + sizeBytes = manifest.artifactSizeBytes + metadataSizeBytes; + // Refuse an already-oversized durable pair from metadata alone before + // spending I/O on its full fixed-memory checksum. + enforceDurableCheckpointLimits(sizeBytes, limits); + if (!error.payloadVerified) { + verifyManifestPayload( + existing, + manifest, + Math.min(limits.maxCheckpointBytes, limits.maxProjectBytes) + ); + } } - const pruned = pruneCheckpoints(root, id, limits); - if (pruned.length > 0) syncDirectory(directory); - const encoding = checkpointFile(root, id).endsWith(".gz") ? "gzip" : "json"; - return { id, path: checkpointRelativePath(id, encoding), pruned }; + // A complete canonical pair makes any older recovery-only records for the + // same stable ID redundant. They are invisible to list/prune throughout. + cleanupRecoveryManifests(root, id, expectedTransaction); + } + enforceDurableCheckpointLimits(sizeBytes, limits); + if (process.platform !== "win32") { + fs.chmodSync(existing, 0o600); + const manifestFile = checkpointManifestFile(root, id); + if (fs.existsSync(manifestFile)) fs.chmodSync(manifestFile, 0o600); } + const pruned = pruneCheckpoints(root, id, limits); + if (pruned.length > 0) syncDirectory(directory); + const encoding = checkpointFile(root, id).endsWith(".gz") ? "gzip" : "json"; + return { id, path: checkpointRelativePath(id, encoding), pruned }; +} + +async function saveCheckpointArtifactExclusive( + root: string, + relative: string, + id: string, + checkpoint: SharedSearchCheckpoint, + limits: Required +): Promise { + const directory = checkpointsDirectory(root); + const destination = checkpointDestination(root, id); + const existing = checkpointFile(root, id); + if (fs.existsSync(existing)) return reuseCheckpointArtifact(root, relative, id, checkpoint, limits); // Validate the existing retention set before creating a new durable file. // Corrupt old state must not turn a successful write into a partial cleanup. - checkpointRecords(root); + try { + checkpointRecords(root); + } catch (error) { + // A same-ID writer can publish between the initial existence check and + // retention validation. Reopen its durable transaction instead of making + // validation inflate a sidecar-free (and potentially huge) frontier. + if (fs.existsSync(checkpointFile(root, id))) { + return reuseCheckpointArtifact(root, relative, id, checkpoint, limits); + } + throw error; + } + if (fs.existsSync(checkpointFile(root, id))) { + return reuseCheckpointArtifact(root, relative, id, checkpoint, limits); + } const artifact: CheckpointArtifact = { artifactSchemaVersion: CHECKPOINT_ARTIFACT_SCHEMA_VERSION, artifactType: "shared-search-checkpoint", @@ -416,22 +2057,93 @@ export async function saveCheckpointArtifact( configuration: checkpoint.configuration, checkpoint, }; - const temporary = path.join(directory, `.${id}.${process.pid}.${randomUUID()}.tmp`); + const transaction = reserveCheckpointTransaction(root, id); + const temporary = transaction.temporary; + let publishedPayload = false; + let completedPair = false; try { - await writeCompressedArtifact(temporary, artifact, limits); - fs.renameSync(temporary, destination); - syncDirectory(directory); + const written = await writeCompressedArtifact(temporary, artifact, limits, true); + const manifest = manifestForArtifact(artifact, "gzip", written.sizeBytes, written.sha256); + const rawManifest = serializedManifest(manifest); + enforceDurableCheckpointLimits(written.sizeBytes + Buffer.byteLength(rawManifest), limits); + // The recovery manifest is the durable transaction intent. It is fsynced, + // then its directory is fsynced, before payload bytes can become visible. + writeRecoveryManifest(root, id, transaction.slot, rawManifest); + checkpointTestBarrier("after-recovery-manifest"); + checkpointTestCrash("after-recovery-manifest"); + try { + fs.linkSync(temporary, destination); + publishedPayload = true; + syncDirectory(directory); + checkpointTestCrash("after-payload-publication"); + } catch (error) { + if ((error as NodeJS.ErrnoException).code !== "EEXIST") throw error; + } + const reference = await reuseCheckpointArtifact( + root, + relative, + id, + checkpoint, + limits, + transaction + ); + completedPair = true; + return reference; + } catch (error) { + // Once another writer has completed the canonical pair, its verified + // cleanup may unlink this loser's private slot while an fd is still open. + // Treat an ensuing path-level ENOENT as loss of the no-clobber race and + // reopen the winner; a corrupt/orphan final still fails closed there. + if ((error as NodeJS.ErrnoException).code === "ENOENT" && fs.existsSync(destination)) { + const reference = await reuseCheckpointArtifact( + root, + relative, + id, + checkpoint, + limits, + transaction + ); + completedPair = true; + return reference; + } + throw error; } finally { - fs.rmSync(temporary, { force: true }); + // Publish a durable exact-nonce release before this operation returns. + // Other processes may then reclaim a failed transaction even though this + // Node process remains alive, while concurrent same-process transactions + // with different nonces stay protected. + releaseCheckpointTransaction(root, id, transaction); + // A losing candidate can never describe the visible payload and a fully + // verified pair no longer needs this writer's candidate. The winner's + // matching candidate remains durable across every earlier failure path. + if (!publishedPayload || completedPair) { + cleanupRecoverySlot(root, id, transaction.slot, transaction); + } } - const pruned = pruneCheckpoints(root, id, limits); - if (pruned.length > 0) syncDirectory(directory); - return { id, path: checkpointRelativePath(id), pruned }; } -export function listCheckpointArtifacts(projectRoot: string): CheckpointArtifactSummary[] { - return checkpointRecords(projectRoot) - .map(({ file: _file, ...record }) => record) +export async function saveCheckpointArtifact( + projectRoot: string, + entrypoint: string, + checkpoint: SharedSearchCheckpoint, + inputLimits: CheckpointStorageLimits = {} +): Promise { + const root = path.resolve(projectRoot); + const relative = relativeEntrypoint(root, entrypoint); + const id = checkpointId(relative, checkpoint); + const directory = checkpointsDirectory(root); + const limits = storageLimits(inputLimits); + fs.mkdirSync(directory, { recursive: true, mode: 0o700 }); + if (process.platform !== "win32") fs.chmodSync(directory, 0o700); + return saveCheckpointArtifactExclusive(root, relative, id, checkpoint, limits); +} + +export function listCheckpointArtifacts( + projectRoot: string, + readLimits: CheckpointReadLimits = {} +): CheckpointArtifactSummary[] { + return checkpointRecords(projectRoot, readLimits) + .map(({ file: _file, manifestFile: _manifestFile, ...record }) => record) .sort((a, b) => b.createdAt.localeCompare(a.createdAt) || a.id.localeCompare(b.id)); } @@ -451,27 +2163,35 @@ async function freshness( }; } -export async function openCheckpointArtifact(projectRoot: string, id: string): Promise<{ +export async function openCheckpointArtifact( + projectRoot: string, + id: string, + readLimits: CheckpointReadLimits = {} +): Promise<{ artifact: CheckpointArtifactSummary & { freshness: CheckpointFreshness }; }> { - const artifact = loadArtifact(projectRoot, id); - const current = await freshness(projectRoot, artifact); - return { artifact: { ...summary(projectRoot, artifact), freshness: current.freshness } }; + const loaded = loadArtifactDetailed(projectRoot, id, readLimits); + const current = await freshness(projectRoot, loaded.artifact); + return { artifact: { ...summary(projectRoot, loaded), freshness: current.freshness } }; } -export async function loadCheckpointForResume(projectRoot: string, id: string): Promise<{ +export async function loadCheckpointForResume( + projectRoot: string, + id: string, + readLimits: CheckpointReadLimits = {} +): Promise<{ artifact: CheckpointArtifactSummary & { freshness: "current" }; checkpoint: SharedSearchCheckpoint; entrypoint: string; }> { - const artifact = loadArtifact(projectRoot, id); - const current = await freshness(projectRoot, artifact); + const loaded = loadArtifactDetailed(projectRoot, id, readLimits); + const current = await freshness(projectRoot, loaded.artifact); if (current.freshness !== "current") { throw new Error(`checkpoint ${id} is ${current.freshness}; resume requires the exact source and knot map used to create it`); } return { - artifact: { ...summary(projectRoot, artifact), freshness: "current" }, - checkpoint: artifact.checkpoint, + artifact: { ...summary(projectRoot, loaded), freshness: "current" }, + checkpoint: loaded.artifact.checkpoint, entrypoint: current.entrypoint, }; } diff --git a/test/fixtures/checkpoint-save-worker.js b/test/fixtures/checkpoint-save-worker.js new file mode 100644 index 0000000..4620aae --- /dev/null +++ b/test/fixtures/checkpoint-save-worker.js @@ -0,0 +1,43 @@ +const fs = require("node:fs"); +const path = require("node:path"); + +if (process.env.INKCHECK_TEST_FORBID_CHECKPOINT_DECODE === "1") { + const zlib = require("node:zlib"); + const marker = process.env.INKCHECK_TEST_CHECKPOINT_DECODE_MARKER; + zlib.gunzipSync = () => { + if (marker) fs.writeFileSync(marker, "checkpoint payload decode attempted\n"); + const error = new RangeError("checkpoint recovery must not decode the published payload"); + error.code = "ERR_BUFFER_TOO_LARGE"; + throw error; + }; +} + +const { compile, scanKnots } = require("../../dist/inklecate"); +const { exploreSharedResumable } = require("../../dist/explore"); +const { saveCheckpointArtifact } = require("../../dist/checkpoints"); + +async function main() { + const projectRoot = path.resolve(process.argv[2]); + const story = path.join(projectRoot, "story.ink"); + const maxStates = Number(process.argv[3] ?? 20); + const compiled = await compile(story); + const options = { + maxStates, + preserveTurnState: false, + preserveRandomState: false, + }; + if (process.env.INKCHECK_TEST_CHECKPOINT_MAX_DEPTH) { + options.maxDepth = Number(process.env.INKCHECK_TEST_CHECKPOINT_MAX_DEPTH); + } + if (process.env.INKCHECK_TEST_CHECKPOINT_SEED) { + options.seed = Number(process.env.INKCHECK_TEST_CHECKPOINT_SEED); + } + const checkpoint = exploreSharedResumable(compiled.storyJson, scanKnots(story), [], options).checkpoint; + const reference = await saveCheckpointArtifact(projectRoot, story, checkpoint); + process.stdout.write(`${JSON.stringify(reference)}\n`); +} + +main().catch((error) => { + process.stderr.write(`${error.stack ?? error.message ?? String(error)}\n`); + process.exitCode = 1; +}); diff --git a/test/inkcheck.test.js b/test/inkcheck.test.js index e5217aa..016f74f 100644 --- a/test/inkcheck.test.js +++ b/test/inkcheck.test.js @@ -3,7 +3,7 @@ const assert = require("node:assert"); const fs = require("node:fs"); const os = require("node:os"); const path = require("node:path"); -const { spawnSync } = require("node:child_process"); +const { spawn, spawnSync } = require("node:child_process"); const { parseIssue, @@ -51,6 +51,8 @@ const { } = require("../dist/discovery"); const { CHECKPOINT_ARTIFACT_SCHEMA_VERSION, + CheckpointReadError, + CheckpointSizeLimitError, listCheckpointArtifacts, loadCheckpointForResume, openCheckpointArtifact, @@ -90,12 +92,30 @@ const EARLY_CHOICE_GRID = path.join(__dirname, "..", "examples", "early-choice-g const DEEP_CHAIN = path.join(__dirname, "..", "examples", "deep-chain.ink"); const EXTERNAL_STORY = path.join(__dirname, "..", "examples", "external-story.ink"); const CLI = path.join(__dirname, "..", "dist", "cli.js"); +const CHECKPOINT_SAVE_WORKER = path.join(__dirname, "fixtures", "checkpoint-save-worker.js"); const ROOT = path.join(__dirname, ".."); const SEARCH_FIXTURES = path.join(__dirname, "fixtures", "search"); const INSPECT_PROJECT = path.join(__dirname, "fixtures", "inspect", "project.ink"); const DUPLICATE_CHOICE_TEXT = path.join(__dirname, "fixtures", "duplicate-choice-text.ink"); const ASSERTION_STORY = path.join(__dirname, "fixtures", "assertions.ink"); const POLICY_LATE_ERROR = path.join(__dirname, "fixtures", "policy-late-error.ink"); + +function runCheckpointSaveWorker(projectRoot, maxStates = 20, environment = {}) { + return new Promise((resolve, reject) => { + const child = spawn(process.execPath, [CHECKPOINT_SAVE_WORKER, projectRoot, String(maxStates)], { + env: { ...process.env, ...environment }, + stdio: ["ignore", "pipe", "pipe"], + }); + let stdout = ""; + let stderr = ""; + child.stdout.setEncoding("utf8"); + child.stderr.setEncoding("utf8"); + child.stdout.on("data", (chunk) => { stdout += chunk; }); + child.stderr.on("data", (chunk) => { stderr += chunk; }); + child.on("error", reject); + child.on("close", (status, signal) => resolve({ status, signal, stdout, stderr })); + }); +} const STAGED_DISJOINT = path.join(__dirname, "fixtures", "staged-disjoint.ink"); const NO_DISCOVERY_BEFORE_DEPTH = path.join(__dirname, "fixtures", "no-discovery-before-depth.ink"); const LATE_RECOVERY = path.join(__dirname, "fixtures", "late-recovery.ink"); @@ -1660,6 +1680,9 @@ test("checkpoint artifacts are private, source-bound, idempotent, and determinis assert.deepStrictEqual(repeated.pruned, []); const firstFile = path.join(tmp, ...first.path.split("/")); if (process.platform !== "win32") assert.strictEqual(fs.statSync(firstFile).mode & 0o777, 0o600); + const firstManifest = firstFile.replace(/\.json(?:\.gz)?$/, ".meta.json"); + assert.strictEqual(fs.existsSync(firstManifest), true); + if (process.platform !== "win32") assert.strictEqual(fs.statSync(firstManifest).mode & 0o777, 0o600); assert.strictEqual((await openCheckpointArtifact(tmp, first.id)).artifact.freshness, "current"); assert.strictEqual((await loadCheckpointForResume(tmp, first.id)).checkpoint.state.totalGranted, 10); @@ -1678,20 +1701,43 @@ test("checkpoint artifacts are private, source-bound, idempotent, and determinis assert.match(newest.path, /\.json\.gz$/); const { gunzipSync } = require("node:zlib"); const compressedFile = path.join(tmp, ...newest.path.split("/")); + const compressedManifest = compressedFile.replace(/\.json\.gz$/, ".meta.json"); const legacyFile = compressedFile.slice(0, -3); const uncompressed = gunzipSync(fs.readFileSync(compressedFile)); assert.ok(fs.statSync(compressedFile).size < uncompressed.length); + assert.strictEqual(newest.payloadSizeBytes, fs.statSync(compressedFile).size); + assert.strictEqual(newest.metadataSizeBytes, fs.statSync(compressedManifest).size); + assert.strictEqual(newest.sizeBytes, newest.payloadSizeBytes + newest.metadataSizeBytes); - // Existing schema-v1 JSON artifacts remain readable after storage compression ships. + // Existing sidecar-free schema-v1 JSON artifacts remain readable after storage compression ships. fs.writeFileSync(legacyFile, uncompressed, { mode: 0o600 }); fs.rmSync(compressedFile); + fs.rmSync(compressedManifest); const legacy = await openCheckpointArtifact(tmp, references[3].id); assert.strictEqual(legacy.artifact.storageEncoding, "json"); assert.match(legacy.artifact.path, /\.json$/); assert.strictEqual((await loadCheckpointForResume(tmp, references[3].id)).checkpoint.state.totalGranted, 40); + for (let attempt = 0; attempt < 2; attempt++) { + const crashedLegacyPublish = await runCheckpointSaveWorker(tmp, 40, { + NODE_ENV: "test", + INKCHECK_TEST_CHECKPOINT_CRASH_STAGE: "after-legacy-manifest-temporary", + INKCHECK_TEST_CHECKPOINT_MAX_DEPTH: "150", + INKCHECK_TEST_CHECKPOINT_SEED: "7", + }); + assert.strictEqual(crashedLegacyPublish.status, 86, crashedLegacyPublish.stderr); + const legacyDebris = fs.readdirSync(path.dirname(legacyFile)) + .filter((name) => name.includes(references[3].id)); + assert.strictEqual(legacyDebris.filter((name) => name.endsWith(".claim")).length, attempt + 1); + assert.strictEqual(legacyDebris.filter((name) => name.endsWith(".payload.tmp")).length, attempt + 1); + assert.strictEqual(legacyDebris.some((name) => /^\.checkpoint-.*\.[0-9]+\.[0-9a-f-]+\.meta\.tmp$/.test(name)), false, + "legacy sidecar reconstruction never creates UUID-style temporary names"); + } const legacyRepeated = await saveCheckpointArtifact(tmp, story, makeCheckpoint(40)); assert.strictEqual(legacyRepeated.path, references[3].path.replace(/\.gz$/, "")); assert.strictEqual(fs.existsSync(`${legacyFile}.gz`), false, "legacy reuse does not create a duplicate encoding"); + assert.strictEqual(fs.readdirSync(path.dirname(legacyFile)) + .some((name) => name.includes(references[3].id) && name.startsWith(".")), false, + "successful legacy sidecar reconstruction cleans its fixed recovery slot"); const beforeQuotaFailure = fs.readdirSync(path.join(tmp, ".inkcheck", "checkpoints")); await assert.rejects( @@ -1737,15 +1783,963 @@ test("checkpoint artifacts fail closed on tampering and incompatible envelopes", const tampered = JSON.parse(original); tampered.checkpoint.state.totalGranted++; fs.writeFileSync(artifactFile, gzipSync(JSON.stringify(tampered))); - assert.throws(() => listCheckpointArtifacts(tmp), /content does not match its stable ID/); + assert.throws(() => listCheckpointArtifacts(tmp), (error) => { + assert.ok(error instanceof CheckpointReadError); + assert.strictEqual(error.kind, "corrupt"); + assert.match(error.message, /content does not match its stable ID/); + return true; + }); + + fs.writeFileSync(artifactFile, gzipSync("{")); + await assert.rejects(() => openCheckpointArtifact(tmp, reference.id), (error) => { + assert.ok(error instanceof CheckpointReadError); + assert.strictEqual(error.kind, "corrupt"); + assert.strictEqual(error.stage, "storage", "manifest payload digest fails before parsing"); + return true; + }); + fs.rmSync(artifactFile.replace(/\.json\.gz$/, ".meta.json")); + await assert.rejects(() => openCheckpointArtifact(tmp, reference.id), (error) => { + assert.ok(error instanceof CheckpointReadError); + assert.strictEqual(error.kind, "corrupt"); + assert.strictEqual(error.stage, "json", "sidecar-free v1 fallback still classifies corrupt JSON"); + return true; + }); const incompatible = JSON.parse(original); incompatible.artifactSchemaVersion = 999; fs.writeFileSync(artifactFile, gzipSync(JSON.stringify(incompatible))); - await assert.rejects(() => openCheckpointArtifact(tmp, reference.id), /compatible Inkcheck version or migrate/); + await assert.rejects(() => openCheckpointArtifact(tmp, reference.id), (error) => { + assert.ok(error instanceof CheckpointReadError); + assert.strictEqual(error.kind, "unsupported"); + assert.match(error.message, /compatible Inkcheck version or migrate/); + return true; + }); fs.writeFileSync(artifactFile, "not-gzip"); - await assert.rejects(() => openCheckpointArtifact(tmp, reference.id), /corrupt gzip/); + await assert.rejects(() => openCheckpointArtifact(tmp, reference.id), (error) => { + assert.ok(error instanceof CheckpointReadError); + assert.strictEqual(error.kind, "corrupt"); + assert.strictEqual(error.stage, "decompression"); + assert.match(error.message, /corrupt gzip/); + return true; + }); + } finally { + fs.rmSync(tmp, { recursive: true, force: true }); + } +}); + +test("checkpoint readback is resource-bounded while manifests keep listing metadata-only", async () => { + const tmp = fs.mkdtempSync(path.join(os.tmpdir(), "inkcheck-checkpoint-read-limit-")); + const story = path.join(tmp, "story.ink"); + try { + fs.copyFileSync(path.join(SEARCH_FIXTURES, "low-dedup-wide.ink"), story); + const compiled = await compile(story); + const checkpoint = exploreSharedResumable(compiled.storyJson, scanKnots(story), [], { + maxStates: 10, + preserveTurnState: false, + preserveRandomState: false, + }).checkpoint; + const reference = await saveCheckpointArtifact(tmp, story, checkpoint); + const artifactFile = path.join(tmp, ...reference.path.split("/")); + const manifestFile = artifactFile.replace(/\.json(?:\.gz)?$/, ".meta.json"); + + const listed = listCheckpointArtifacts(tmp, { maxDecompressedBytes: 1 }); + assert.strictEqual(listed.length, 1, "listing reads the sidecar instead of inflating the frontier"); + assert.strictEqual(Object.hasOwn(listed[0], "manifestFile"), false, "private sidecar paths are not exposed"); + + const originalOpenSync = fs.openSync; + fs.openSync = (candidate, ...args) => { + if (path.resolve(String(candidate)) === path.resolve(artifactFile)) { + throw new Error("checkpoint payload must not be opened during manifested list/prune"); + } + return originalOpenSync(candidate, ...args); + }; + try { + assert.strictEqual(listCheckpointArtifacts(tmp).length, 1); + const nextCheckpoint = exploreSharedResumable(compiled.storyJson, scanKnots(story), [], { + maxStates: 20, + preserveTurnState: false, + preserveRandomState: false, + }).checkpoint; + await saveCheckpointArtifact(tmp, story, nextCheckpoint); + } finally { + fs.openSync = originalOpenSync; + } + + const originalPayload = fs.readFileSync(artifactFile); + const sameSizeTamper = Buffer.from(originalPayload); + sameSizeTamper[Math.floor(sameSizeTamper.length / 2)] ^= 1; + fs.writeFileSync(artifactFile, sameSizeTamper); + assert.strictEqual(listCheckpointArtifacts(tmp).some((item) => item.id === reference.id), true, + "bounded listing deliberately defers same-size payload integrity to open/resume"); + await assert.rejects(() => openCheckpointArtifact(tmp, reference.id), (error) => { + assert.ok(error instanceof CheckpointReadError); + assert.strictEqual(error.kind, "corrupt"); + assert.strictEqual(error.stage, "storage"); + return true; + }); + fs.writeFileSync(artifactFile, originalPayload); + + await assert.rejects( + () => openCheckpointArtifact(tmp, reference.id, { maxDecompressedBytes: 1 }), + (error) => { + assert.ok(error instanceof CheckpointReadError); + assert.strictEqual(error.kind, "resource_limit"); + assert.strictEqual(error.stage, "decompression"); + assert.strictEqual(error.limitBytes, 1); + assert.match(error.message, /framed checkpoint format/); + return true; + } + ); + await assert.rejects( + () => loadCheckpointForResume(tmp, reference.id, { maxStoredBytes: 1 }), + (error) => { + assert.ok(error instanceof CheckpointReadError); + assert.strictEqual(error.kind, "resource_limit"); + assert.strictEqual(error.stage, "storage"); + assert.ok(error.observedBytes > error.limitBytes); + return true; + } + ); + + // A same-ID save first applies the schema-v1 default read envelope, then + // may fall back to a fixed-memory checksum under the caller's larger + // durable-storage cap. Use a synthetic fstat size so the test proves that + // cap is forwarded without reserving, reading, or allocating 512 MiB. + const sparseSize = (512 * 1024 * 1024) + 1; + const payloadStat = fs.statSync(artifactFile); + const sparseStat = new Proxy(payloadStat, { + get(target, property) { + if (property === "size") return sparseSize; + const value = Reflect.get(target, property, target); + return typeof value === "function" ? value.bind(target) : value; + }, + }); + const originalFstatSync = fs.fstatSync; + const originalCloseSync = fs.closeSync; + const originalReadSync = fs.readSync; + const payloadFds = new Map(); + let payloadOpenCount = 0; + const acceptedActiveCap = new Error("active checkpoint digest cap accepted sparse fstat"); + fs.openSync = (candidate, ...args) => { + const fd = originalOpenSync(candidate, ...args); + if (path.resolve(String(candidate)) === path.resolve(artifactFile)) { + payloadOpenCount++; + payloadFds.set(fd, payloadOpenCount); + } + return fd; + }; + fs.fstatSync = (fd, ...args) => payloadFds.has(fd) + ? sparseStat + : originalFstatSync(fd, ...args); + fs.closeSync = (fd) => { + payloadFds.delete(fd); + return originalCloseSync(fd); + }; + fs.readSync = (fd, ...args) => { + if (payloadFds.get(fd) === 2) throw acceptedActiveCap; + return originalReadSync(fd, ...args); + }; + try { + await assert.rejects( + () => saveCheckpointArtifact(tmp, story, checkpoint, { + maxCheckpointBytes: sparseSize + 1024, + maxProjectBytes: sparseSize + 1024, + }), + (error) => { + assert.strictEqual(error, acceptedActiveCap, + "fallback checksum reaches its fixed-buffer read under the caller's active cap"); + return true; + } + ); + } finally { + fs.openSync = originalOpenSync; + fs.fstatSync = originalFstatSync; + fs.closeSync = originalCloseSync; + fs.readSync = originalReadSync; + } + assert.strictEqual(payloadOpenCount, 2, + "default readback rejects from fstat before the active-cap verifier reopens once"); + + const heldManifest = `${manifestFile}.held`; + fs.renameSync(manifestFile, heldManifest); + try { + assert.throws(() => listCheckpointArtifacts(tmp, { maxDecompressedBytes: 1 }), (error) => { + assert.ok(error instanceof CheckpointReadError); + assert.strictEqual(error.kind, "resource_limit"); + return true; + }); + } finally { + fs.renameSync(heldManifest, manifestFile); + } + + const originalManifest = fs.readFileSync(manifestFile, "utf8"); + for (const [field, value] of [["createdAt", "1900-01-01T00:00:00.000Z"], ["totalGranted", 999]]) { + const changed = JSON.parse(originalManifest); + changed[field] = value; + fs.writeFileSync(manifestFile, JSON.stringify(changed)); + assert.throws(() => listCheckpointArtifacts(tmp), (error) => { + assert.ok(error instanceof CheckpointReadError); + assert.strictEqual(error.kind, "corrupt"); + assert.strictEqual(error.stage, "manifest"); + assert.match(error.message, /canonical checksum/); + return true; + }); + fs.writeFileSync(manifestFile, originalManifest); + } + + const incompatibleManifest = JSON.parse(originalManifest); + incompatibleManifest.manifestSchemaVersion = 2; + incompatibleManifest.futureSchemaField = { framedPayload: true }; + fs.writeFileSync(manifestFile, JSON.stringify(incompatibleManifest)); + assert.throws(() => listCheckpointArtifacts(tmp), (error) => { + assert.ok(error instanceof CheckpointReadError); + assert.strictEqual(error.kind, "unsupported"); + assert.strictEqual(error.stage, "manifest"); + return true; + }); + const futureManifest = fs.readFileSync(manifestFile); + const recoveryFile = path.join( + path.dirname(manifestFile), + `.${reference.id}.recovery-00.meta.json` + ); + fs.writeFileSync(recoveryFile, originalManifest); + fs.writeFileSync(`${recoveryFile}.claim`, JSON.stringify({ schemaVersion: 1, pid: 999999 })); + await assert.rejects(() => saveCheckpointArtifact(tmp, story, checkpoint), (error) => { + assert.ok(error instanceof CheckpointReadError); + assert.strictEqual(error.kind, "unsupported"); + assert.strictEqual(error.stage, "manifest"); + return true; + }); + assert.deepStrictEqual(fs.readFileSync(manifestFile), futureManifest, + "v1 recovery metadata never overwrites an authoritative future-schema sidecar"); + const oversizedCanonical = Buffer.alloc((64 * 1024) + 1, 0x20); + fs.writeFileSync(manifestFile, oversizedCanonical); + await assert.rejects(() => saveCheckpointArtifact(tmp, story, checkpoint), (error) => { + assert.ok(error instanceof CheckpointReadError); + assert.strictEqual(error.kind, "resource_limit"); + assert.strictEqual(error.stage, "manifest"); + return true; + }); + assert.deepStrictEqual(fs.readFileSync(manifestFile), oversizedCanonical, + "resource-limited canonical metadata is never replaced by a recovery record"); + } finally { + fs.rmSync(tmp, { recursive: true, force: true }); + } +}); + +test("legacy JSON checkpoint growth is bounded from one open descriptor", async () => { + const tmp = fs.mkdtempSync(path.join(os.tmpdir(), "inkcheck-checkpoint-legacy-growth-")); + const story = path.join(tmp, "story.ink"); + try { + fs.copyFileSync(path.join(SEARCH_FIXTURES, "low-dedup-wide.ink"), story); + const compiled = await compile(story); + const checkpoint = exploreSharedResumable(compiled.storyJson, scanKnots(story), [], { + maxStates: 10, + preserveTurnState: false, + preserveRandomState: false, + }).checkpoint; + const reference = await saveCheckpointArtifact(tmp, story, checkpoint); + const compressedFile = path.join(tmp, ...reference.path.split("/")); + const manifestFile = compressedFile.replace(/\.json\.gz$/, ".meta.json"); + const legacyFile = compressedFile.replace(/\.gz$/, ""); + const legacyBytes = require("node:zlib").gunzipSync(fs.readFileSync(compressedFile)); + fs.writeFileSync(legacyFile, legacyBytes); + fs.rmSync(compressedFile); + fs.rmSync(manifestFile); + + const maxStoredBytes = legacyBytes.length + 64; + const originalOpenSync = fs.openSync; + const originalReadSync = fs.readSync; + let payloadOpens = 0; + let grew = false; + fs.openSync = (candidate, flags, ...args) => { + if (path.resolve(String(candidate)) === path.resolve(legacyFile) && flags === "r") payloadOpens++; + return originalOpenSync(candidate, flags, ...args); + }; + fs.readSync = (fd, buffer, offset, length, position) => { + const read = originalReadSync(fd, buffer, offset, length, position); + if (!grew && position === 0) { + grew = true; + fs.appendFileSync(legacyFile, Buffer.alloc(128, 0x20)); + } + return read; + }; + try { + await assert.rejects( + () => openCheckpointArtifact(tmp, reference.id, { + maxStoredBytes, + maxDecompressedBytes: legacyBytes.length + 1024, + }), + (error) => { + assert.ok(error instanceof CheckpointReadError); + assert.strictEqual(error.kind, "resource_limit"); + assert.strictEqual(error.stage, "storage"); + assert.ok(error.observedBytes > maxStoredBytes); + return true; + } + ); + } finally { + fs.openSync = originalOpenSync; + fs.readSync = originalReadSync; + } + assert.strictEqual(payloadOpens, 1, "legacy bounds and bytes come from one opened file descriptor"); + } finally { + fs.rmSync(tmp, { recursive: true, force: true }); + } +}); + +test("checkpoint recovery manifests survive deterministic crash windows without payload decoding", async () => { + const tmp = fs.mkdtempSync(path.join(os.tmpdir(), "inkcheck-checkpoint-crash-recovery-")); + const story = path.join(tmp, "story.ink"); + try { + fs.copyFileSync(path.join(SEARCH_FIXTURES, "low-dedup-wide.ink"), story); + const beforePayload = await runCheckpointSaveWorker(tmp, 20, { + NODE_ENV: "test", + INKCHECK_TEST_CHECKPOINT_CRASH_STAGE: "after-recovery-manifest", + }); + assert.strictEqual(beforePayload.status, 86, beforePayload.stderr); + + const directory = path.join(tmp, ".inkcheck", "checkpoints"); + let names = fs.readdirSync(directory); + assert.strictEqual(names.filter((name) => /\.recovery-\d+\.meta\.json$/.test(name)).length, 1); + assert.strictEqual(names.some((name) => name.endsWith(".json.gz")), false, + "a recovery-only intent is never exposed as checkpoint evidence"); + + const afterPayload = await runCheckpointSaveWorker(tmp, 20, { + NODE_ENV: "test", + INKCHECK_TEST_CHECKPOINT_CRASH_STAGE: "after-payload-publication", + }); + assert.strictEqual(afterPayload.status, 86, afterPayload.stderr); + names = fs.readdirSync(directory); + const payloadName = names.find((name) => /^checkpoint-[0-9a-f]{24}\.json\.gz$/.test(name)); + assert.ok(payloadName); + const id = payloadName.slice(0, -".json.gz".length); + assert.strictEqual(fs.existsSync(path.join(directory, `${id}.meta.json`)), false); + assert.strictEqual(names.filter((name) => /\.recovery-\d+\.meta\.json$/.test(name)).length, 2, + "the matching durable record and an older recovery-only record both survive the crashes"); + assert.throws(() => listCheckpointArtifacts(tmp, { maxDecompressedBytes: 1 }), (error) => { + assert.ok(error instanceof CheckpointReadError); + assert.strictEqual(error.kind, "resource_limit"); + return true; + }); + + const decodeMarker = path.join(tmp, "decode-attempted"); + const recovered = await runCheckpointSaveWorker(tmp, 20, { + NODE_ENV: "test", + INKCHECK_TEST_FORBID_CHECKPOINT_DECODE: "1", + INKCHECK_TEST_CHECKPOINT_DECODE_MARKER: decodeMarker, + }); + assert.strictEqual(recovered.status, 0, recovered.stderr); + assert.strictEqual(fs.existsSync(decodeMarker), false, + "resource-limited recovery reconstructs metadata from the durable record without inflating JSON"); + names = fs.readdirSync(directory); + assert.deepStrictEqual(names.filter((name) => name.includes(id)).sort(), [ + `${id}.json.gz`, `${id}.meta.json`, + ]); + assert.strictEqual(names.some((name) => name.includes(".recovery-")), false, + "recovery-only debris is cleaned only after the canonical pair is reread and verified"); + assert.strictEqual(listCheckpointArtifacts(tmp, { maxDecompressedBytes: 1 })[0].id, id); + + fs.rmSync(path.join(directory, `${id}.json.gz`)); + const replaceGap = await runCheckpointSaveWorker(tmp, 20, { + NODE_ENV: "test", + INKCHECK_TEST_CHECKPOINT_FORCE_WINDOWS_REPLACE: "1", + INKCHECK_TEST_CHECKPOINT_CRASH_STAGE: "after-manifest-displacement", + }); + assert.strictEqual(replaceGap.status, 86, replaceGap.stderr); + names = fs.readdirSync(directory); + assert.strictEqual(fs.existsSync(path.join(directory, `${id}.meta.json`)), false, + "the deterministic failpoint stops inside the portable replace gap"); + assert.strictEqual(names.some((name) => name.endsWith(".promote")), true); + assert.strictEqual(names.some((name) => name.endsWith(".displaced")), true); + + const gapDecodeMarker = path.join(tmp, "replace-gap-decode-attempted"); + const gapRecovered = await runCheckpointSaveWorker(tmp, 20, { + NODE_ENV: "test", + INKCHECK_TEST_FORBID_CHECKPOINT_DECODE: "1", + INKCHECK_TEST_CHECKPOINT_DECODE_MARKER: gapDecodeMarker, + }); + assert.strictEqual(gapRecovered.status, 0, gapRecovered.stderr); + assert.strictEqual(fs.existsSync(gapDecodeMarker), false); + names = fs.readdirSync(directory); + assert.deepStrictEqual(names.filter((name) => name.includes(id)).sort(), [ + `${id}.json.gz`, `${id}.meta.json`, + ], "the bounded promotion and displaced companions are cleaned after gap recovery"); + } finally { + fs.rmSync(tmp, { recursive: true, force: true }); + } +}); + +test("checkpoint cleanup preserves a live same-ID writer's claimed slot", async () => { + const tmp = fs.mkdtempSync(path.join(os.tmpdir(), "inkcheck-checkpoint-live-owner-")); + const story = path.join(tmp, "story.ink"); + try { + fs.copyFileSync(path.join(SEARCH_FIXTURES, "low-dedup-wide.ink"), story); + const ownerGate = path.join(tmp, "live-owner-gate"); + const heldOwner = runCheckpointSaveWorker(tmp, 20, { + NODE_ENV: "test", + INKCHECK_TEST_CHECKPOINT_GATE: ownerGate, + INKCHECK_TEST_CHECKPOINT_GATE_STAGE: "after-recovery-manifest", + }); + const deadline = Date.now() + 10_000; + while (!fs.readdirSync(tmp).some((name) => name.startsWith("live-owner-gate.") + && name.endsWith(".ready"))) { + if (Date.now() >= deadline) throw new Error("timed out waiting for held checkpoint owner"); + await new Promise((resolve) => setTimeout(resolve, 25)); + } + + const winner = await runCheckpointSaveWorker(tmp, 20); + assert.strictEqual(winner.status, 0, winner.stderr); + const id = JSON.parse(winner.stdout).id; + const directory = path.join(tmp, ".inkcheck", "checkpoints"); + let names = fs.readdirSync(directory); + assert.strictEqual(names.some((name) => name === `${id}.json.gz`), true); + assert.strictEqual(names.some((name) => name === `${id}.meta.json`), true); + assert.strictEqual(names.filter((name) => name.endsWith(".claim")).length, 1); + assert.strictEqual(names.filter((name) => /\.recovery-\d+\.meta\.json$/.test(name)).length, 1); + assert.strictEqual(names.filter((name) => name.endsWith(".payload.tmp")).length, 1, + "the winner cannot unlink another live process's claimed transaction"); + + fs.writeFileSync(ownerGate, "release"); + const owner = await heldOwner; + assert.strictEqual(owner.status, 0, owner.stderr); + names = fs.readdirSync(directory); + assert.deepStrictEqual(names.filter((name) => name.includes(id)).sort(), [ + `${id}.json.gz`, `${id}.meta.json`, + ], "the slot owner releases its own claim after reopening the winner"); + } finally { + fs.rmSync(tmp, { recursive: true, force: true }); + } +}); + +test("checkpoint publication preserves concurrent same-process nonce owners", async () => { + const tmp = fs.mkdtempSync(path.join(os.tmpdir(), "inkcheck-checkpoint-same-process-owner-")); + const story = path.join(tmp, "story.ink"); + try { + fs.copyFileSync(path.join(SEARCH_FIXTURES, "low-dedup-wide.ink"), story); + const compiled = await compile(story); + const checkpoint = exploreSharedResumable(compiled.storyJson, scanKnots(story), [], { + maxStates: 20, + preserveTurnState: false, + preserveRandomState: false, + }).checkpoint; + const pending = Array.from({ length: 4 }, () => saveCheckpointArtifact(tmp, story, checkpoint)); + const directory = path.join(tmp, ".inkcheck", "checkpoints"); + let names = fs.readdirSync(directory); + assert.strictEqual(names.filter((name) => name.endsWith(".claim")).length, 4, + "all synchronous reservation phases own distinct nonce-bearing slots"); + assert.strictEqual(names.filter((name) => name.endsWith(".payload.tmp")).length, 4); + assert.strictEqual(names.some((name) => name.endsWith(".json.gz")), false); + + const references = await Promise.all(pending); + assert.strictEqual(new Set(references.map((reference) => reference.id)).size, 1); + const id = references[0].id; + names = fs.readdirSync(directory); + assert.deepStrictEqual(names.filter((name) => name.includes(id)).sort(), [ + `${id}.json.gz`, `${id}.meta.json`, + ], "a winner cannot delete live sibling transactions, and every owner cleans after joining it"); + } finally { + fs.rmSync(tmp, { recursive: true, force: true }); + } +}); + +test("checkpoint manifest tampering cannot steer deterministic retention", async () => { + const tmp = fs.mkdtempSync(path.join(os.tmpdir(), "inkcheck-checkpoint-manifest-retention-")); + const story = path.join(tmp, "story.ink"); + try { + fs.copyFileSync(path.join(SEARCH_FIXTURES, "low-dedup-wide.ink"), story); + const compiled = await compile(story); + const makeCheckpoint = (maxStates) => exploreSharedResumable(compiled.storyJson, scanKnots(story), [], { + maxStates, + preserveTurnState: false, + preserveRandomState: false, + }).checkpoint; + const references = []; + for (const maxStates of [10, 20, 30]) { + Atomics.wait(new Int32Array(new SharedArrayBuffer(4)), 0, 0, 2); + references.push(await saveCheckpointArtifact(tmp, story, makeCheckpoint(maxStates))); + } + const oldest = [...listCheckpointArtifacts(tmp)].sort((a, b) => a.createdAt.localeCompare(b.createdAt))[0]; + const manifestFile = path.join(tmp, ".inkcheck", "checkpoints", `${oldest.id}.meta.json`); + const originalManifest = fs.readFileSync(manifestFile, "utf8"); + const changed = JSON.parse(originalManifest); + changed.createdAt = "2999-01-01T00:00:00.000Z"; + fs.writeFileSync(manifestFile, JSON.stringify(changed)); + const before = fs.readdirSync(path.dirname(manifestFile)).sort(); + await assert.rejects(() => saveCheckpointArtifact(tmp, story, makeCheckpoint(40)), (error) => { + assert.ok(error instanceof CheckpointReadError); + assert.strictEqual(error.kind, "corrupt"); + assert.match(error.message, /canonical checksum/); + return true; + }); + assert.deepStrictEqual(fs.readdirSync(path.dirname(manifestFile)).sort(), before, + "retention does not delete or reorder around tampered metadata"); + + fs.writeFileSync(manifestFile, originalManifest); + const saved = await saveCheckpointArtifact(tmp, story, makeCheckpoint(40)); + assert.strictEqual(saved.pruned.includes(oldest.id), true, "untampered creation order remains authoritative"); + assert.strictEqual(references.length, 3); + } finally { + fs.rmSync(tmp, { recursive: true, force: true }); + } +}); + +test("checkpoint publication serializes cross-process same-ID writers and recovers one-file crash windows", async () => { + const tmp = fs.mkdtempSync(path.join(os.tmpdir(), "inkcheck-checkpoint-concurrent-save-")); + const story = path.join(tmp, "story.ink"); + try { + fs.copyFileSync(path.join(SEARCH_FIXTURES, "low-dedup-wide.ink"), story); + const compiled = await compile(story); + const checkpoint = exploreSharedResumable(compiled.storyJson, scanKnots(story), [], { + maxStates: 20, + preserveTurnState: false, + preserveRandomState: false, + }).checkpoint; + const publicationGate = path.join(tmp, "same-id-publication-gate"); + const workerPromises = Array.from({ length: 6 }, () => runCheckpointSaveWorker(tmp, 20, { + NODE_ENV: "test", + INKCHECK_TEST_CHECKPOINT_GATE: publicationGate, + INKCHECK_TEST_CHECKPOINT_GATE_STAGE: "after-recovery-manifest", + })); + const gateDeadline = Date.now() + 10_000; + while (fs.readdirSync(tmp).filter((name) => name.startsWith("same-id-publication-gate.") + && name.endsWith(".ready")).length < 6) { + if (Date.now() >= gateDeadline) throw new Error("timed out waiting for six checkpoint publication workers"); + await new Promise((resolve) => setTimeout(resolve, 25)); + } + const directory = path.join(tmp, ".inkcheck", "checkpoints"); + let names = fs.readdirSync(directory); + assert.strictEqual(names.filter((name) => /\.recovery-\d+\.meta\.json$/.test(name)).length, 6, + "all six processes hold durable recovery records before payload publication"); + assert.strictEqual(names.some((name) => name.endsWith(".json.gz")), false); + fs.writeFileSync(publicationGate, "release"); + const workers = await Promise.all(workerPromises); + for (const worker of workers) assert.strictEqual(worker.status, 0, worker.stderr); + const references = workers.map((worker) => JSON.parse(worker.stdout)); + assert.strictEqual(new Set(references.map((reference) => reference.id)).size, 1); + const id = references[0].id; + names = fs.readdirSync(directory); + assert.deepStrictEqual(names.filter((name) => name.includes(id)).sort(), [`${id}.json.gz`, `${id}.meta.json`]); + assert.strictEqual(names.some((name) => name.endsWith(".tmp") || name.endsWith(".lock")), false); + assert.strictEqual((await openCheckpointArtifact(tmp, id)).artifact.totalGranted, 20); + + const artifactFile = path.join(directory, `${id}.json.gz`); + const manifestFile = path.join(directory, `${id}.meta.json`); + const deadOwner = spawn(process.execPath, ["-e", ""]); + const deadOwnerPid = deadOwner.pid; + await new Promise((resolve, reject) => { + deadOwner.once("error", reject); + deadOwner.once("close", resolve); + }); + assert.ok(Number.isSafeInteger(deadOwnerPid) && deadOwnerPid > 0); + const oversizedRecovery = path.join(directory, `.${id}.recovery-00.meta.json`); + fs.writeFileSync(oversizedRecovery, Buffer.alloc((64 * 1024) + 1)); + fs.writeFileSync(`${oversizedRecovery}.claim`, JSON.stringify({ schemaVersion: 1, pid: deadOwnerPid })); + const originalOpenSync = fs.openSync; + const originalCloseSync = fs.closeSync; + const originalReadSync = fs.readSync; + const oversizedFds = new Set(); + let readOversizedRecovery = false; + fs.openSync = (candidate, ...args) => { + const fd = originalOpenSync(candidate, ...args); + if (path.resolve(String(candidate)) === path.resolve(oversizedRecovery)) oversizedFds.add(fd); + return fd; + }; + fs.closeSync = (fd) => { + oversizedFds.delete(fd); + return originalCloseSync(fd); + }; + fs.readSync = (fd, ...args) => { + if (oversizedFds.has(fd)) { + readOversizedRecovery = true; + throw new Error("oversized recovery metadata must be rejected from fstat before reading bytes"); + } + return originalReadSync(fd, ...args); + }; + try { + await saveCheckpointArtifact(tmp, story, checkpoint); + } finally { + fs.openSync = originalOpenSync; + fs.closeSync = originalCloseSync; + fs.readSync = originalReadSync; + } + assert.strictEqual(readOversizedRecovery, false); + assert.strictEqual(fs.existsSync(oversizedRecovery), false, + "a verified canonical pair safely cleans an oversized hidden recovery candidate"); + + const validManifest = fs.readFileSync(manifestFile); + fs.writeFileSync(manifestFile, "{corrupt"); + await assert.rejects(() => saveCheckpointArtifact(tmp, story, checkpoint), (error) => { + assert.ok(error instanceof CheckpointReadError); + assert.strictEqual(error.kind, "corrupt"); + assert.strictEqual(error.stage, "manifest"); + return true; + }); + assert.strictEqual(fs.existsSync(artifactFile), true); + names = fs.readdirSync(directory); + assert.strictEqual(names.some((name) => name.includes(".recovery-")), false, + "a corrupt final manifest without a recovery record fails closed"); + fs.writeFileSync(manifestFile, validManifest); + + fs.rmSync(manifestFile); + assert.strictEqual(listCheckpointArtifacts(tmp)[0].metadataSizeBytes, 0, + "a payload-first crash remains readable through the v1 compatibility path"); + await saveCheckpointArtifact(tmp, story, checkpoint); + assert.strictEqual(fs.existsSync(manifestFile), true, "idempotent reuse restores a missing sidecar"); + + fs.rmSync(artifactFile); + assert.deepStrictEqual(listCheckpointArtifacts(tmp), [], "an orphan sidecar is never listed as evidence"); + const beforeReplacementGate = path.join(tmp, "portable-replacement-before-gate"); + const afterReplacementGate = path.join(tmp, "portable-replacement-after-gate"); + const replacementJoinMarker = path.join(tmp, "portable-replacement-joined"); + const joiner = runCheckpointSaveWorker(tmp, 20, { + NODE_ENV: "test", + INKCHECK_TEST_CHECKPOINT_FORCE_WINDOWS_REPLACE: "1", + INKCHECK_TEST_CHECKPOINT_GATE: beforeReplacementGate, + INKCHECK_TEST_CHECKPOINT_GATE_STAGE: "before-manifest-displacement", + INKCHECK_TEST_CHECKPOINT_JOIN_MARKER: replacementJoinMarker, + }); + const replacementDeadline = Date.now() + 10_000; + while (!fs.readdirSync(tmp).some((name) => name.startsWith("portable-replacement-before-gate.") + && name.endsWith(".ready"))) { + if (Date.now() >= replacementDeadline) throw new Error("joiner did not reach the pre-displacement gate"); + await new Promise((resolve) => setTimeout(resolve, 25)); + } + const displacer = runCheckpointSaveWorker(tmp, 20, { + NODE_ENV: "test", + INKCHECK_TEST_CHECKPOINT_FORCE_WINDOWS_REPLACE: "1", + INKCHECK_TEST_CHECKPOINT_GATE: afterReplacementGate, + INKCHECK_TEST_CHECKPOINT_GATE_STAGE: "after-manifest-displacement", + }); + while (!fs.readdirSync(tmp).some((name) => name.startsWith("portable-replacement-after-gate.") + && name.endsWith(".ready"))) { + if (Date.now() >= replacementDeadline) throw new Error("displacer did not enter portable replace gap"); + await new Promise((resolve) => setTimeout(resolve, 25)); + } + assert.strictEqual(fs.existsSync(manifestFile), false); + fs.writeFileSync(beforeReplacementGate, "release joiner"); + while (!fs.existsSync(manifestFile)) { + if (Date.now() >= replacementDeadline) throw new Error("joiner did not publish from portable replace gap"); + await new Promise((resolve) => setTimeout(resolve, 25)); + } + assert.strictEqual(fs.existsSync(replacementJoinMarker), true, + "the second promoter executed the existing-displaced join branch"); + fs.writeFileSync(afterReplacementGate, "release displacer"); + const replacementResults = await Promise.all([joiner, displacer]); + for (const worker of replacementResults) assert.strictEqual(worker.status, 0, worker.stderr); + assert.strictEqual((await openCheckpointArtifact(tmp, id)).artifact.totalGranted, 20, + "a second portable promoter joins the bounded replace gap instead of failing corrupt"); + + const previousNodeEnv = process.env.NODE_ENV; + const previousForceWindows = process.env.INKCHECK_TEST_CHECKPOINT_FORCE_WINDOWS_REPLACE; + const previousStaleMarker = process.env.INKCHECK_TEST_CHECKPOINT_STALE_DISPLACED_MARKER; + const leaveReleasedPortableGap = async (marker, failRelease = false) => { + fs.rmSync(artifactFile); + const originalLinkSync = fs.linkSync; + const originalTransientRmSync = fs.rmSync; + const originalWriteFileSync = fs.writeFileSync; + let injectedPublishError = false; + let injectedRestoreBusy = false; + let injectedReleaseError = false; + fs.linkSync = (source, destination, ...args) => { + if (!injectedPublishError && String(source).endsWith(".promote") + && path.resolve(String(destination)) === path.resolve(manifestFile)) { + injectedPublishError = true; + const error = new Error("simulated portable manifest publish failure"); + error.code = "EIO"; + throw error; + } + return originalLinkSync(source, destination, ...args); + }; + fs.rmSync = (candidate, ...args) => { + if (injectedPublishError && !injectedRestoreBusy && String(candidate).endsWith(".displaced")) { + injectedRestoreBusy = true; + const error = new Error("simulated restored-displaced cleanup refusal"); + error.code = "EBUSY"; + throw error; + } + return originalTransientRmSync(candidate, ...args); + }; + fs.writeFileSync = (candidate, ...args) => { + if (failRelease && !injectedReleaseError && String(candidate).endsWith(".claim.released")) { + injectedReleaseError = true; + const error = new Error("simulated durable release failure"); + error.code = "EIO"; + throw error; + } + return originalWriteFileSync(candidate, ...args); + }; + process.env.NODE_ENV = "test"; + process.env.INKCHECK_TEST_CHECKPOINT_FORCE_WINDOWS_REPLACE = "1"; + process.env.INKCHECK_TEST_CHECKPOINT_STALE_DISPLACED_MARKER = marker; + try { + await assert.rejects(() => saveCheckpointArtifact(tmp, story, checkpoint), (error) => { + assert.strictEqual(error.code, "EIO"); + return true; + }); + } finally { + fs.linkSync = originalLinkSync; + fs.rmSync = originalTransientRmSync; + fs.writeFileSync = originalWriteFileSync; + } + assert.strictEqual(injectedPublishError, true); + assert.strictEqual(injectedRestoreBusy, true); + const releasedNames = fs.readdirSync(directory); + assert.strictEqual(releasedNames.some((name) => name.endsWith(".displaced")), true); + assert.strictEqual(injectedReleaseError, failRelease); + assert.strictEqual(releasedNames.some((name) => name.endsWith(".claim.released")), !failRelease, + failRelease + ? "the release seam failed before a durable marker became visible" + : "a returned transient error durably releases its exact transaction nonce"); + }; + try { + const sameProcessTakeover = path.join(tmp, "same-process-stale-displaced-takeover"); + await leaveReleasedPortableGap(sameProcessTakeover); + await saveCheckpointArtifact(tmp, story, checkpoint); + assert.strictEqual(fs.existsSync(sameProcessTakeover), true, + "a same-process retry takes over the inactive displaced companion"); + assert.deepStrictEqual(fs.readdirSync(directory).filter((name) => name.includes(id)).sort(), [ + `${id}.json.gz`, `${id}.meta.json`, + ], "same-process retry reclaims its released recovery namespace"); + + const crossProcessTakeover = path.join(tmp, "cross-process-stale-displaced-takeover"); + await leaveReleasedPortableGap(crossProcessTakeover); + const externalRetry = await runCheckpointSaveWorker(tmp, 20, { + NODE_ENV: "test", + INKCHECK_TEST_CHECKPOINT_FORCE_WINDOWS_REPLACE: "1", + INKCHECK_TEST_CHECKPOINT_STALE_DISPLACED_MARKER: crossProcessTakeover, + }); + assert.strictEqual(externalRetry.status, 0, externalRetry.stderr); + assert.strictEqual(fs.existsSync(crossProcessTakeover), true, + "another process distinguishes the durable release while its owner PID remains alive"); + assert.deepStrictEqual(fs.readdirSync(directory).filter((name) => name.includes(id)).sort(), [ + `${id}.json.gz`, `${id}.meta.json`, + ]); + + const inMemoryTakeover = path.join(tmp, "same-process-in-memory-stale-takeover"); + await leaveReleasedPortableGap(inMemoryTakeover, true); + await saveCheckpointArtifact(tmp, story, checkpoint); + assert.strictEqual(fs.existsSync(inMemoryTakeover), true, + "failed release I/O still clears the exact active nonce for a same-process retry"); + assert.deepStrictEqual(fs.readdirSync(directory).filter((name) => name.includes(id)).sort(), [ + `${id}.json.gz`, `${id}.meta.json`, + ]); + } finally { + if (previousNodeEnv === undefined) delete process.env.NODE_ENV; + else process.env.NODE_ENV = previousNodeEnv; + if (previousForceWindows === undefined) { + delete process.env.INKCHECK_TEST_CHECKPOINT_FORCE_WINDOWS_REPLACE; + } else { + process.env.INKCHECK_TEST_CHECKPOINT_FORCE_WINDOWS_REPLACE = previousForceWindows; + } + if (previousStaleMarker === undefined) { + delete process.env.INKCHECK_TEST_CHECKPOINT_STALE_DISPLACED_MARKER; + } else { + process.env.INKCHECK_TEST_CHECKPOINT_STALE_DISPLACED_MARKER = previousStaleMarker; + } + } + + fs.rmSync(artifactFile); + fs.rmSync(manifestFile); + const originalRmSync = fs.rmSync; + let injectedBusyCleanup = 0; + fs.rmSync = (candidate, ...args) => { + if (String(candidate).endsWith(".payload.tmp")) { + injectedBusyCleanup++; + const error = new Error("simulated Windows open-file cleanup refusal"); + error.code = "EBUSY"; + throw error; + } + return originalRmSync(candidate, ...args); + }; + try { + await saveCheckpointArtifact(tmp, story, checkpoint); + } finally { + fs.rmSync = originalRmSync; + } + assert.ok(injectedBusyCleanup >= 2, + "cleanup remains busy both before and after the transaction is durably released"); + assert.strictEqual((await openCheckpointArtifact(tmp, id)).artifact.totalGranted, 20, + "best-effort bounded cleanup cannot undo a verified pair on Windows"); + names = fs.readdirSync(directory); + const busyClaimName = names.find((name) => name.endsWith(".claim")); + assert.ok(busyClaimName); + const busyReleaseName = `${busyClaimName}.released`; + assert.strictEqual(names.includes(busyReleaseName), true, + "cleanup refusal preserves durable inactive state beside the claim"); + const busyClaim = JSON.parse(fs.readFileSync(path.join(directory, busyClaimName), "utf8")); + const busyRelease = JSON.parse(fs.readFileSync(path.join(directory, busyReleaseName), "utf8")); + assert.strictEqual(busyRelease.nonce, busyClaim.nonce); + await saveCheckpointArtifact(tmp, story, checkpoint); + assert.deepStrictEqual(fs.readdirSync(directory).filter((name) => name.includes(id)).sort(), [ + `${id}.json.gz`, `${id}.meta.json`, + ], "a clean retry reclaims the exact released transaction after EBUSY clears"); + + fs.rmSync(artifactFile); + fs.rmSync(manifestFile); + const orphanRelease = path.join(directory, `.${id}.recovery-00.meta.json.claim.released`); + const orphanNonce = "00000000-0000-4000-8000-000000000000"; + fs.writeFileSync(orphanRelease, JSON.stringify({ schemaVersion: 1, nonce: orphanNonce })); + const guardedRmSync = fs.rmSync; + let removedOrphanUnderClaim = false; + fs.rmSync = (candidate, ...args) => { + if (path.resolve(String(candidate)) === path.resolve(orphanRelease)) { + const claimFile = orphanRelease.slice(0, -".released".length); + assert.strictEqual(fs.existsSync(`${claimFile}.cleaning`), true, + "the atomic cleaning claim excludes reservation while orphan state is removed"); + removedOrphanUnderClaim = true; + } + return guardedRmSync(candidate, ...args); + }; + try { + await saveCheckpointArtifact(tmp, story, checkpoint); + } finally { + fs.rmSync = guardedRmSync; + } + assert.strictEqual(removedOrphanUnderClaim, true, + "a claimless marker is removed only while the fixed cleaning claim owns the slot"); + assert.deepStrictEqual(fs.readdirSync(directory).filter((name) => name.includes(id)).sort(), [ + `${id}.json.gz`, `${id}.meta.json`, + ]); + + const isolateClaimFile = path.join(directory, `.${id}.recovery-00.meta.json.claim`); + const isolateReleaseFile = `${isolateClaimFile}.released`; + const isolateTemporary = path.join(directory, `.${id}.recovery-00.meta.json.payload.tmp`); + const isolateNonce = "11111111-1111-4111-8111-111111111111"; + fs.writeFileSync(isolateClaimFile, JSON.stringify({ + schemaVersion: 1, + pid: process.pid, + nonce: isolateNonce, + })); + fs.writeFileSync(isolateTemporary, "held by another worker isolate"); + await saveCheckpointArtifact(tmp, story, checkpoint); + assert.strictEqual(fs.existsSync(isolateClaimFile), true); + assert.strictEqual(fs.existsSync(isolateTemporary), true, + "a markerless current-PID nonce unknown to this isolate remains live"); + fs.writeFileSync(isolateReleaseFile, JSON.stringify({ schemaVersion: 1, nonce: isolateNonce })); + await saveCheckpointArtifact(tmp, story, checkpoint); + assert.strictEqual(fs.existsSync(isolateClaimFile), false, + "the exact durable release permits cross-isolate cleanup"); + assert.strictEqual(fs.existsSync(isolateTemporary), false); + + const corruptReleaseNonce = "22222222-2222-4222-8222-222222222222"; + fs.writeFileSync(isolateClaimFile, JSON.stringify({ + schemaVersion: 1, + pid: process.pid, + nonce: corruptReleaseNonce, + })); + fs.writeFileSync(isolateTemporary, "quarantined"); + fs.writeFileSync(isolateReleaseFile, "{torn"); + await saveCheckpointArtifact(tmp, story, checkpoint); + assert.strictEqual(fs.existsSync(isolateClaimFile), true); + assert.strictEqual(fs.existsSync(isolateTemporary), true, + "a malformed unowned release marker quarantines only its fixed slot"); + fs.rmSync(isolateClaimFile); + fs.rmSync(isolateTemporary); + fs.rmSync(isolateReleaseFile); + + const deadCorruptReleaseNonce = "23232323-2323-4232-8232-232323232323"; + fs.writeFileSync(isolateClaimFile, JSON.stringify({ + schemaVersion: 1, + pid: deadOwnerPid, + nonce: deadCorruptReleaseNonce, + })); + fs.writeFileSync(isolateTemporary, "dead owner"); + fs.writeFileSync(isolateReleaseFile, "{torn"); + await saveCheckpointArtifact(tmp, story, checkpoint); + assert.strictEqual(fs.existsSync(isolateClaimFile), false); + assert.strictEqual(fs.existsSync(isolateTemporary), false, + "a corrupt release cannot wedge a definitely dead owner slot"); + assert.strictEqual(fs.existsSync(isolateReleaseFile), false); + + fs.writeFileSync(isolateTemporary, "claimless busy companion"); + const originalClaimlessRmSync = fs.rmSync; + fs.rmSync = (candidate, ...args) => { + if (path.resolve(String(candidate)) === path.resolve(isolateTemporary)) { + const error = new Error("simulated claimless cleanup refusal"); + error.code = "EBUSY"; + throw error; + } + return originalClaimlessRmSync(candidate, ...args); + }; + fs.rmSync(artifactFile); + fs.rmSync(manifestFile); + try { + await saveCheckpointArtifact(tmp, story, checkpoint); + } finally { + fs.rmSync = originalClaimlessRmSync; + } + assert.strictEqual(fs.existsSync(isolateTemporary), true, + "reservation skips a claimless slot when a companion cannot be removed"); + assert.strictEqual((await openCheckpointArtifact(tmp, id)).artifact.totalGranted, 20); + fs.rmSync(isolateTemporary); + + const staleCleaning = `${isolateClaimFile}.cleaning`; + fs.writeFileSync(staleCleaning, JSON.stringify({ + schemaVersion: 1, + pid: process.pid, + nonce: "33333333-3333-4333-8333-333333333333", + })); + fs.rmSync(artifactFile); + fs.rmSync(manifestFile); + await saveCheckpointArtifact(tmp, story, checkpoint); + assert.strictEqual(fs.existsSync(staleCleaning), true, + "a crashed cleaning owner consumes one bounded slot instead of permitting pathname reuse"); + assert.strictEqual((await openCheckpointArtifact(tmp, id)).artifact.totalGranted, 20); + fs.rmSync(staleCleaning); + + const recorded = listCheckpointArtifacts(tmp)[0]; + fs.rmSync(artifactFile); + fs.rmSync(manifestFile); + const pairLimit = recorded.payloadSizeBytes + Math.floor(recorded.metadataSizeBytes / 2); + await assert.rejects( + () => saveCheckpointArtifact(tmp, story, checkpoint, { maxCheckpointBytes: pairLimit }), + (error) => { + assert.ok(error instanceof CheckpointSizeLimitError); + assert.strictEqual(error.kind, "single"); + assert.ok(error.observedBytes > pairLimit, "quota includes the private metadata sidecar"); + return true; + } + ); + await assert.rejects( + () => saveCheckpointArtifact(tmp, story, checkpoint, { + maxCheckpointBytes: Number.MAX_SAFE_INTEGER, + maxProjectBytes: pairLimit, + }), + (error) => { + assert.ok(error instanceof CheckpointSizeLimitError); + assert.strictEqual(error.kind, "project"); + assert.ok(error.observedBytes > pairLimit, "project quota also includes sidecar bytes"); + return true; + } + ); + names = fs.readdirSync(directory); + assert.strictEqual(names.some((name) => name.endsWith(".tmp") || name.endsWith(".lock")), false); + assert.strictEqual(names.some((name) => name === `${id}.json.gz` || name === `${id}.meta.json`), false, + "pair-quota failure publishes neither file"); + + for (let slot = 0; slot < 32; slot++) { + fs.writeFileSync( + path.join(directory, `.${id}.recovery-${String(slot).padStart(2, "0")}.meta.json.claim.cleaning`), + JSON.stringify({ + schemaVersion: 1, + pid: process.pid, + nonce: "44444444-4444-4444-8444-444444444444", + }) + ); + } + const cappedNames = fs.readdirSync(directory).sort(); + await assert.rejects(() => saveCheckpointArtifact(tmp, story, checkpoint), (error) => { + assert.ok(error instanceof CheckpointReadError); + assert.strictEqual(error.kind, "resource_limit"); + assert.strictEqual(error.stage, "manifest"); + assert.strictEqual(error.observedBytes, 32); + assert.strictEqual(error.limitBytes, 32); + return true; + }); + assert.deepStrictEqual(fs.readdirSync(directory).sort(), cappedNames, + "32 fail-closed cleaning claims bound exhaustion and create no 33rd recovery namespace"); } finally { fs.rmSync(tmp, { recursive: true, force: true }); }