Files
stack/packages/queue/src/io.mjs
T
jason.woltjeandClaude Opus 5.5 6ca116b7ba feat(queue): queue as data A2, migration, render and dispatch (#1508)
Filbert approved round 1 (f167b85e). Manifest 782bcb62, 21 files, plus
the QUEUE.md markers and the TOOLS.md section. Lead decision 35.

Co-Authored-By: Claude Opus 5.5 <[email protected]>
2026-09-26 20:14:09 -05:00

108 lines
3.6 KiB
JavaScript

// The file layer every queue write goes through. `realIo` is the only layer
// the CLI uses. Tests pass their own layer through the API to inject one
// fault at a time (8.5); no flag or environment variable selects it.
import {
closeSync, constants, fsyncSync, linkSync, lstatSync, openSync, readFileSync, renameSync, statSync,
statfsSync, unlinkSync, writeSync,
} from "node:fs";
import { QueueError } from "./errors.mjs";
export const realIo = Object.freeze({
openExcl: (path, mode) => openSync(path, constants.O_CREAT | constants.O_EXCL | constants.O_WRONLY, mode),
openRead: (path) => openSync(path, constants.O_RDONLY),
write: (fd, buf, off, len) => writeSync(fd, buf, off, len),
fsync: (fd) => fsyncSync(fd),
close: (fd) => closeSync(fd),
readFile: (path) => readFileSync(path),
rename: (from, to) => renameSync(from, to),
unlink: (path) => unlinkSync(path),
link: (from, to) => linkSync(from, to),
stat: (path) => statSync(path, { bigint: true }),
lstat: (path) => lstatSync(path, { bigint: true }),
fsyncDir: (dir) => {
const fd = openSync(dir, constants.O_RDONLY | constants.O_DIRECTORY);
try { fsyncSync(fd); } finally { closeSync(fd); }
},
statfsType: (path) => statfsSync(path).type,
});
// ext4 (ext2/ext3 share the magic), xfs, btrfs. tmpfs only for a test that
// passes a layer with `allowTmpfs: true`; realIo never has it (N5).
const FS_TYPES = new Map([[0xef53, "ext4"], [0x58465342, "xfs"], [0x9123683e, "btrfs"]]);
export const TMPFS = 0x01021994;
export function checkPlatform(io, dirs) {
if (process.platform !== "linux") throw new QueueError(`unsupported platform ${process.platform}: the queue runs on Linux only`, 2);
for (const dir of dirs) {
const type = io.statfsType(dir);
if (type === TMPFS && io.allowTmpfs === true) continue;
if (!FS_TYPES.has(type)) throw new QueueError(`unsupported filesystem under ${dir} (type 0x${type.toString(16)}); ext4, xfs or btrfs only`, 2);
}
}
export function sleepMs(ms) {
Atomics.wait(new Int32Array(new SharedArrayBuffer(4)), 0, 0, ms);
}
// Loops on short writes and checks the total.
export function writeAll(io, fd, bytes) {
let off = 0;
while (off < bytes.length) {
const n = io.write(fd, bytes, off, bytes.length - off);
if (!Number.isInteger(n) || n <= 0) throw Object.assign(new Error(`short write: ${off} of ${bytes.length} bytes`), { code: "ESHORT" });
off += n;
}
if (off !== bytes.length) throw Object.assign(new Error(`wrote ${off} of ${bytes.length} bytes`), { code: "ESHORT" });
}
// O_EXCL temp file, every byte written, fsync, close. On any failure the
// temp file is removed and the error rethrown.
export function writeTemp(io, path, bytes, mode = 0o644) {
const fd = io.openExcl(path, mode);
try {
writeAll(io, fd, bytes);
io.fsync(fd);
} catch (err) {
try { io.close(fd); } catch { /* already failing */ }
unlinkQuiet(io, path);
throw err;
}
try {
io.close(fd);
} catch (err) {
unlinkQuiet(io, path);
throw err;
}
}
export function fsyncFile(io, path) {
const fd = io.openRead(path);
try { io.fsync(fd); } finally { io.close(fd); }
}
export function unlinkQuiet(io, path) {
try { io.unlink(path); } catch { /* absent or not ours to report */ }
}
export function readOrNull(io, path) {
try {
return io.readFile(path);
} catch (err) {
if (err.code === "ENOENT") return null;
throw err;
}
}
export function lstatOrNull(io, path) {
try {
return io.lstat(path);
} catch (err) {
if (err.code === "ENOENT") return null;
throw err;
}
}
export function errno(err) {
return err?.code ?? err?.message ?? String(err);
}