← Files ModRetro Chromatic PluginARCHIVED FILE

dist/chromatic-journal.js

19 KB · Oct 2, 2026 · 00:37 UTC

↓ Download file

import { closeSync, constants, fstatSync, fsyncSync, ftruncateSync, lstatSync, mkdirSync, openSync, opendirSync, readSync, readdirSync, realpathSync, renameSync, rmdirSync, unlinkSync, writeSync, } from "node:fs";
import fs from "node:fs";
import { createHash, randomUUID } from "node:crypto";
import os from "node:os";
import path from "node:path";
const OWNER = "codex-gb-studio-chromatic-operations-v1";
const RECORD_LIMIT = 8 * 1024 * 1024;
const STREAM_LIMIT = 8 * 1024 * 1024;
// Keep the slot in the existing directory namespace so older plugin versions
// can still admit explicit actions against the same journal after a rollback.
const OBSERVATION_DIRECTORY = createHash("sha256").update("chromatic-observation-slot-v1").digest("hex");
const OPERATION_ID = /^(?:[a-f0-9]{64}|[a-f0-9]{8}-[a-f0-9]{4}-[a-f0-9]{4}-[a-f0-9]{4}-[a-f0-9]{12})$/u;
const identity = (a, b) => a.dev === b.dev && a.ino === b.ino && a.uid === b.uid && a.gid === b.gid && a.mode === b.mode && a.nlink === b.nlink;
const unchanged = (a, b) => identity(a, b) && a.size === b.size && a.mtimeMs === b.mtimeMs && a.ctimeMs === b.ctimeMs;
const absent = (error) => error instanceof Error && "code" in error && error.code === "ENOENT";
function boundedRead(fd, size) {
    const bytes = Buffer.alloc(size);
    let offset = 0;
    while (offset < size) {
        const count = readSync(fd, bytes, offset, size - offset, offset);
        if (!count)
            throw new Error("Chromatic journal shrank while reading");
        offset += count;
    }
    if (readSync(fd, Buffer.alloc(1), 0, 1, offset))
        throw new Error("Chromatic journal grew while reading");
    return bytes;
}
function owned(filename, stat, directory) {
    if (stat.isSymbolicLink() || (directory ? !stat.isDirectory() : !stat.isFile() || stat.nlink !== 1)
        || realpathSync(filename) !== path.resolve(filename)
        || (process.getuid && stat.uid !== process.getuid())
        || (process.platform !== "win32" && (stat.mode & 0o777) !== (directory ? 0o700 : 0o600))) {
        throw new Error(`Unsafe Chromatic operation journal identity: ${filename}`);
    }
}
function readJson(filename, limit = RECORD_LIMIT) {
    let before;
    try {
        before = lstatSync(filename);
    }
    catch (error) {
        if (absent(error))
            return undefined;
        throw error;
    }
    owned(filename, before, false);
    if (before.size > limit)
        throw new Error("Chromatic operation journal exceeds its read bound");
    const fd = fs.openSync(filename, constants.O_RDONLY | (constants.O_NOFOLLOW ?? 0) | (constants.O_NONBLOCK ?? 0));
    try {
        if (!unchanged(before, fstatSync(fd)))
            throw new Error("Chromatic journal changed while opening");
        const raw = boundedRead(fd, before.size);
        if (raw.length !== before.size || !unchanged(before, fstatSync(fd)) || !unchanged(before, lstatSync(filename)))
            throw new Error("Chromatic journal changed while reading");
        const value = JSON.parse(raw.toString("utf8"));
        if (!value || typeof value !== "object" || Array.isArray(value) || !("owner" in value) || value.owner !== OWNER)
            throw new Error("Chromatic operation journal owner differs");
        return value;
    }
    finally {
        closeSync(fd);
    }
}
function readOutput(filename) {
    const before = lstatSync(filename);
    owned(filename, before, false);
    if (before.size > STREAM_LIMIT)
        throw new Error("Chromatic output journal exceeds its bound");
    const fd = fs.openSync(filename, constants.O_RDONLY | (constants.O_NOFOLLOW ?? 0) | (constants.O_NONBLOCK ?? 0));
    try {
        if (!unchanged(before, fstatSync(fd)))
            throw new Error("Chromatic output changed while opening");
        const bytes = boundedRead(fd, before.size);
        if (bytes.length !== before.size || !unchanged(before, fstatSync(fd)) || !unchanged(before, lstatSync(filename)))
            throw new Error("Chromatic output is changing; its current complete contents are not established");
        const text = bytes.toString("utf8");
        const preview = bytes.subarray(0, 512 * 1024);
        return { text: preview.toString("utf8"), bytes: bytes.length, sha256: createHash("sha256").update(bytes).digest("hex"),
            previewTruncated: preview.length !== bytes.length,
            ...(Buffer.from(text).equals(bytes) ? {} : { encoding: "invalid-utf8" }), complete: false };
    }
    finally {
        closeSync(fd);
    }
}
function createDirectory(parent, name) {
    const parentStat = lstatSync(parent);
    if (!parentStat.isDirectory() || parentStat.isSymbolicLink() || realpathSync(parent) !== parent
        || (process.getuid && parentStat.uid !== process.getuid())
        || (process.platform !== "win32" && (parentStat.mode & 0o022)))
        throw new Error("Chromatic state parent must be an owned canonical directory");
    const target = path.join(parent, name);
    try {
        mkdirSync(target, { mode: 0o700 });
    }
    catch (error) {
        if (!(error instanceof Error) || !("code" in error) || error.code !== "EEXIST")
            throw error;
    }
    owned(target, lstatSync(target), true);
    return target;
}
function writeAll(fd, bytes) {
    let offset = 0;
    while (offset < bytes.length) {
        const written = writeSync(fd, bytes, offset, bytes.length - offset, offset);
        if (written <= 0)
            throw new Error("Chromatic journal write made no progress");
        offset += written;
    }
    ftruncateSync(fd, bytes.length);
    fsyncSync(fd);
}
/** Small durable evidence journal, not a daemon or proof of child liveness. */
export class ChromaticJournal {
    root;
    #custom;
    constructor(root) {
        this.root = root ? path.resolve(root) : path.join(os.homedir(), ".codex-gb-studio", "chromatic");
        this.#custom = root !== undefined;
    }
    #root(create) {
        if (create) {
            if (!this.#custom)
                createDirectory(realpathSync(os.homedir()), ".codex-gb-studio");
            createDirectory(path.dirname(this.root), path.basename(this.root));
        }
        try {
            owned(this.root, lstatSync(this.root), true);
        }
        catch (error) {
            if (!create && absent(error))
                return false;
            throw error;
        }
        const marker = path.join(this.root, "owner.json");
        if (create && !readJson(marker, 4096)) {
            // A nonempty unmarked directory is never adopted as a journal.
            if (readdirSync(this.root).length)
                throw new Error("Chromatic journal directory is unowned");
            const fd = openSync(marker, constants.O_WRONLY | constants.O_CREAT | constants.O_EXCL, 0o600);
            try {
                writeAll(fd, Buffer.from(JSON.stringify({ owner: OWNER, schemaVersion: 1 }) + "\n"));
            }
            finally {
                closeSync(fd);
            }
        }
        if (!readJson(marker, 4096))
            throw new Error("Chromatic journal owner marker is missing");
        return true;
    }
    current() {
        if (!this.#root(false))
            return undefined;
        return readJson(path.join(this.root, "active.json"), 4096);
    }
    read(operationId) {
        if (!OPERATION_ID.test(operationId))
            throw new Error("Invalid Chromatic operation ID");
        if (!this.#root(false))
            return undefined;
        let directory = path.join(this.root, operationId);
        try {
            owned(directory, lstatSync(directory), true);
        }
        catch (error) {
            if (!absent(error))
                throw error;
            directory = path.join(this.root, OBSERVATION_DIRECTORY);
            try {
                owned(directory, lstatSync(directory), true);
            }
            catch (error) {
                if (absent(error))
                    return undefined;
                throw error;
            }
        }
        const saved = readJson(path.join(directory, "operation.json"));
        if (directory === path.join(this.root, OBSERVATION_DIRECTORY) && saved?.operationId !== operationId)
            return undefined;
        if (!saved || saved.operationId !== operationId || !saved.operation || typeof saved.operation !== "object")
            throw new Error("Chromatic operation record is incomplete");
        const operation = saved.operation;
        if (operation.state === "running" && Array.isArray(operation.commands)) {
            if (operation.commands.length > 2)
                throw new Error("Unexpected Chromatic command history");
            operation.commands = operation.commands.map((value, index) => {
                if (!value || typeof value !== "object" || Array.isArray(value))
                    throw new Error("Invalid Chromatic command record");
                const command = value;
                if (!command.streamFiles)
                    return command;
                return { ...command, retainedPartialOutput: {
                        stdout: readOutput(path.join(directory, `command-${index}.stdout`)),
                        stderr: readOutput(path.join(directory, `command-${index}.stderr`)),
                    } };
            });
        }
        return operation;
    }
    acquire(operationId, sessionId, initial) {
        if (!OPERATION_ID.test(operationId))
            throw new Error("Invalid Chromatic operation ID");
        this.#root(true);
        const current = this.current();
        if (current)
            throw Object.assign(new Error("A previous Chromatic operation still holds the journal. Its completion is unresolved here; do not retry or delete its record."), { code: "CHROMATIC_OWNERSHIP_UNRESOLVED", details: current });
        // Historical request IDs remain durable across sessions. Stream the names
        // without a lifetime-use quota or loading the full history into memory.
        const entries = opendirSync(this.root);
        try {
            for (let entry = entries.readSync(); entry; entry = entries.readSync()) {
                if (entry.name !== "owner.json" && !OPERATION_ID.test(entry.name))
                    throw new Error("Chromatic journal contains an unexpected entry; preserve it for review");
            }
        }
        finally {
            entries.closeSync();
        }
        const lockPath = path.join(this.root, "active.json");
        const lockFd = openSync(lockPath, constants.O_RDWR | constants.O_CREAT | constants.O_EXCL, 0o600);
        try {
            writeAll(lockFd, Buffer.from(JSON.stringify({ owner: OWNER, schemaVersion: 1, operationId, sessionId, pid: process.pid, createdAt: new Date().toISOString() }) + "\n"));
            const lockIdentity = fstatSync(lockFd);
            const observation = initial.command === "observe_devices";
            const directory = path.join(this.root, observation ? OBSERVATION_DIRECTORY : operationId);
            // Rotate only a settled read-only scan, while holding the shared claim.
            // Interrupted scans and every explicit action keep their original evidence.
            if (observation)
                this.#clearObservation(directory);
            mkdirSync(directory, { mode: 0o700 });
            owned(directory, lstatSync(directory), true);
            const lease = new ChromaticJournalLease(directory, operationId, lockPath, lockFd, lockIdentity);
            lease.write(initial);
            return lease;
        }
        catch (error) {
            closeSync(lockFd);
            // Preserve an incomplete reservation. Never infer that no child could exist
            // merely from a missing PID or partially written metadata after a crash.
            throw error;
        }
    }
    #clearObservation(directory) {
        let before;
        try {
            before = lstatSync(directory);
        }
        catch (error) {
            if (absent(error))
                return;
            throw error;
        }
        owned(directory, before, true);
        const saved = readJson(path.join(directory, "operation.json"));
        const operation = saved?.operation;
        if (!operation || operation.command !== "observe_devices" || !["succeeded", "failed"].includes(String(operation.state))) {
            throw new Error("The previous Chromatic observation is not settled; preserve its record");
        }
        const names = readdirSync(directory);
        for (const name of names) {
            if (!["operation.json", "command-0.stdout", "command-0.stderr"].includes(name))
                throw new Error("Unexpected Chromatic observation evidence; preserve it for review");
            const filename = path.join(directory, name);
            owned(filename, lstatSync(filename), false);
        }
        if (!unchanged(before, lstatSync(directory)))
            throw new Error("Chromatic observation changed before rotation");
        if (operation.state === "failed") {
            // Preserve a failed scan, including its original relative stream files.
            // Reserve the destination exclusively so a prior operation is never replaced.
            const id = saved.operationId;
            if (typeof id !== "string" || !OPERATION_ID.test(id) || id === OBSERVATION_DIRECTORY) {
                throw new Error("Invalid failed Chromatic observation identity");
            }
            const destination = path.join(this.root, id);
            mkdirSync(destination, { mode: 0o700 });
            const reserved = lstatSync(destination);
            owned(destination, reserved, true);
            if (!unchanged(before, lstatSync(directory)) || !unchanged(reserved, lstatSync(destination))) {
                throw new Error("Chromatic observation changed before preservation");
            }
            // Move into the exclusively created directory rather than replacing it:
            // Windows cannot rename a directory over an existing empty directory.
            for (const name of names)
                renameSync(path.join(directory, name), path.join(destination, name));
            rmdirSync(directory);
            return;
        }
        for (const name of names)
            unlinkSync(path.join(directory, name));
        rmdirSync(directory);
    }
}
export class ChromaticJournalLease {
    directory;
    operationId;
    lockPath;
    lockFd;
    lockIdentity;
    #streams = new Map();
    #recordIdentity;
    #checkpointFailure;
    #closed = false;
    constructor(directory, operationId, lockPath, lockFd, lockIdentity) {
        this.directory = directory;
        this.operationId = operationId;
        this.lockPath = lockPath;
        this.lockFd = lockFd;
        this.lockIdentity = lockIdentity;
    }
    write(operation) {
        if (this.#closed)
            throw new Error("Chromatic journal lease is closed");
        if (this.#checkpointFailure)
            throw this.#checkpointFailure;
        const bytes = Buffer.from(JSON.stringify({ owner: OWNER, schemaVersion: 1, operationId: this.operationId, operation }) + "\n");
        if (bytes.length > RECORD_LIMIT)
            throw new Error("Chromatic operation exceeds its durable record bound");
        const destination = path.join(this.directory, "operation.json");
        if (this.#recordIdentity) {
            const current = lstatSync(destination);
            owned(destination, current, false);
            if (!unchanged(current, this.#recordIdentity))
                throw new Error("Chromatic operation checkpoint changed outside its owner");
        }
        owned(this.directory, lstatSync(this.directory), true);
        const temporary = path.join(this.directory, `.checkpoint-${randomUUID()}.json`);
        try {
            const fd = openSync(temporary, constants.O_WRONLY | constants.O_CREAT | constants.O_EXCL, 0o600);
            try {
                writeAll(fd, bytes);
            }
            finally {
                closeSync(fd);
            }
            // The old complete checkpoint survives interruption until this atomic rename.
            renameSync(temporary, destination);
            this.#recordIdentity = lstatSync(destination);
        }
        catch (error) {
            this.#checkpointFailure = error;
            throw error;
        }
    }
    prepareCommand(index) {
        const names = { stdout: `command-${index}.stdout`, stderr: `command-${index}.stderr` };
        for (const stream of ["stdout", "stderr"]) {
            const fd = openSync(path.join(this.directory, names[stream]), constants.O_WRONLY | constants.O_CREAT | constants.O_EXCL, 0o600);
            this.#streams.set(`${index}:${stream}`, { fd, bytes: 0, observedBytes: 0, hash: createHash("sha256") });
        }
        return names;
    }
    output(index, stream, bytes) {
        const state = this.#streams.get(`${index}:${stream}`);
        if (!state || this.#closed)
            throw new Error("Chromatic output has no open journal stream");
        const keep = Math.min(bytes.length, STREAM_LIMIT - state.bytes);
        state.observedBytes += bytes.length;
        if (state.failure)
            throw state.failure;
        let offset = 0;
        try {
            while (offset < keep) {
                const written = fs.writeSync(state.fd, bytes, offset, keep - offset, state.bytes);
                if (written <= 0)
                    throw new Error("Chromatic output journal write made no progress");
                state.hash.update(bytes.subarray(offset, offset + written));
                state.bytes += written;
                offset += written;
            }
        }
        catch (error) {
            state.failure = error;
            throw error;
        }
        return { observedBytes: state.observedBytes, retainedBytes: state.bytes,
            complete: state.observedBytes === state.bytes, sha256: state.hash.copy().digest("hex") };
    }
    finish(operation, release) {
        if (this.#closed)
            return;
        let failure;
        for (const { fd } of this.#streams.values()) {
            try {
                fsyncSync(fd);
            }
            catch (error) {
                failure ??= error;
            }
            try {
                closeSync(fd);
            }
            catch (error) {
                failure ??= error;
            }
        }
        // Publish a terminal checkpoint only after the original stream files settle.
        if (!failure)
            try {
                this.write(operation);
            }
            catch (error) {
                failure = error;
            }
        try {
            if (release && !failure) {
                const actual = lstatSync(this.lockPath);
                owned(this.lockPath, actual, false);
                if (!identity(actual, this.lockIdentity) || !identity(fstatSync(this.lockFd), this.lockIdentity))
                    throw new Error("Chromatic ownership lock changed; it was not removed");
                unlinkSync(this.lockPath);
            }
        }
        catch (error) {
            failure ??= error;
        }
        try {
            closeSync(this.lockFd);
        }
        catch (error) {
            failure ??= error;
        }
        this.#closed = true;
        if (failure)
            throw failure;
    }
}
//# sourceMappingURL=chromatic-journal.js.map

SHA-256: f218a9c9cb7f31c6daf7f64d2b0cf9290eed8617200b10c5d815fc1ec4934c6e