← Files Biological Sequence & Alignment ViewerARCHIVED FILE
src/workspace-atomic-publisher.ts
14.1 KB · Sep 30, 2026 · 23:01 UTC
import { randomUUID } from "node:crypto";
import { constants } from "node:fs";
import { link, lstat, open, rm, stat } from "node:fs/promises";
import path from "node:path";
const destinationTails = new Map<string, Promise<void>>();
export class SequenceWorkspaceCollisionError extends Error {
readonly code = "SEQUENCE_WORKSPACE_COLLISION";
constructor(
message = "The workspace export or provenance sidecar already exists.",
) {
super(message);
this.name = "SequenceWorkspaceCollisionError";
}
}
export function isSequenceWorkspaceCollisionError(
error: unknown,
): error is SequenceWorkspaceCollisionError {
return error instanceof SequenceWorkspaceCollisionError;
}
export type SequenceWorkspaceAtomicHooks = {
beforeArtifactLink?: () => Promise<void> | void;
beforeCommit?: () => Promise<void> | void;
beforeSidecarLink?: () => Promise<void> | void;
};
export async function publishSequenceWorkspacePair({
artifact,
artifactPath,
hooks = {},
sidecar,
sidecarPath,
signal,
}: {
artifact: Uint8Array;
artifactPath: string;
hooks?: SequenceWorkspaceAtomicHooks;
sidecar: string;
sidecarPath: string;
signal?: AbortSignal;
}): Promise<void> {
const release = await acquireDestination(artifactPath, sidecarPath, signal);
try {
await publishLocked({
artifact,
artifactPath,
hooks,
sidecar,
sidecarPath,
signal,
});
} finally {
release();
}
}
export type SequenceWorkspaceStagedIdentity = {
changedAtNanoseconds: string;
device: string;
inode: string;
modifiedAtNanoseconds: string;
size: string;
};
/** Publishes an already-written private file without reading it into memory. */
export async function publishSequenceWorkspaceStagedPair({
artifactPath,
hooks = {},
sidecar,
sidecarPath,
signal,
stagedArtifactPath,
stagedIdentity,
}: {
artifactPath: string;
hooks?: SequenceWorkspaceAtomicHooks;
sidecar: string;
sidecarPath: string;
signal?: AbortSignal;
stagedArtifactPath: string;
stagedIdentity: SequenceWorkspaceStagedIdentity;
}): Promise<void> {
if (
!sameCanonicalDirectory(
path.dirname(stagedArtifactPath),
path.dirname(artifactPath),
)
) {
throw new Error("Workspace staging must be local to its destination.");
}
const release = await acquireDestination(artifactPath, sidecarPath, signal);
try {
await publishStagedLocked({
artifactPath,
hooks,
sidecar,
sidecarPath,
signal,
stagedArtifactPath,
stagedIdentity,
});
} finally {
release();
}
}
async function publishLocked({
artifact,
artifactPath,
hooks,
sidecar,
sidecarPath,
signal,
}: {
artifact: Uint8Array;
artifactPath: string;
hooks: SequenceWorkspaceAtomicHooks;
sidecar: string;
sidecarPath: string;
signal?: AbortSignal;
}): Promise<void> {
const nonce = randomUUID();
const artifactTemp = temporaryPath(artifactPath, nonce, "artifact");
const sidecarTemp = temporaryPath(sidecarPath, nonce, "provenance");
let artifactPublished = false;
let sidecarPublished = false;
let committed = false;
const errors: unknown[] = [];
try {
throwIfAborted(signal);
await writePrivateExclusive(artifactTemp, artifact);
throwIfAborted(signal);
await writePrivateExclusive(sidecarTemp, sidecar);
throwIfAborted(signal);
await assertMissing(artifactPath);
await assertMissing(sidecarPath);
await hooks.beforeArtifactLink?.();
throwIfAborted(signal);
await link(artifactTemp, artifactPath);
artifactPublished = true;
await assertSameRegularFile(artifactPath, artifactTemp);
await hooks.beforeSidecarLink?.();
throwIfAborted(signal);
await assertSameRegularFile(artifactPath, artifactTemp);
await link(sidecarTemp, sidecarPath);
sidecarPublished = true;
await assertSameRegularFile(sidecarPath, sidecarTemp);
await hooks.beforeCommit?.();
throwIfAborted(signal);
await Promise.all([
assertSameRegularFile(artifactPath, artifactTemp),
assertSameRegularFile(sidecarPath, sidecarTemp),
]);
committed = true;
} catch (error) {
errors.push(normalizeCollision(error));
if (!committed) {
if (sidecarPublished) {
await unlinkIfSameFile(sidecarPath, sidecarTemp).catch((rollback) =>
errors.push(rollback),
);
}
if (artifactPublished) {
await unlinkIfSameFile(artifactPath, artifactTemp).catch((rollback) =>
errors.push(rollback),
);
}
}
}
await Promise.allSettled([
rm(artifactTemp, { force: true }),
rm(sidecarTemp, { force: true }),
]).then((results) => {
for (const result of results) {
if (result.status === "rejected") errors.push(result.reason);
}
});
if (committed) {
try {
await syncDirectory(path.dirname(artifactPath));
} catch (error) {
errors.push(error);
}
}
if (errors.length === 1) throw errors[0];
if (errors.length > 1) {
throw new AggregateError(
errors,
errors[0] instanceof Error
? errors[0].message
: "Workspace artifact publication failed.",
);
}
}
async function publishStagedLocked({
artifactPath,
hooks,
sidecar,
sidecarPath,
signal,
stagedArtifactPath,
stagedIdentity,
}: {
artifactPath: string;
hooks: SequenceWorkspaceAtomicHooks;
sidecar: string;
sidecarPath: string;
signal?: AbortSignal;
stagedArtifactPath: string;
stagedIdentity: SequenceWorkspaceStagedIdentity;
}): Promise<void> {
const nonce = randomUUID();
const sidecarTemp = temporaryPath(sidecarPath, nonce, "provenance");
let artifactPublished = false;
let sidecarPublished = false;
let committed = false;
const errors: unknown[] = [];
try {
throwIfAborted(signal);
await assertPrivateStagedFile(stagedArtifactPath, stagedIdentity);
await writePrivateExclusive(sidecarTemp, sidecar);
throwIfAborted(signal);
await assertMissing(artifactPath);
await assertMissing(sidecarPath);
await hooks.beforeArtifactLink?.();
throwIfAborted(signal);
await assertPrivateStagedFile(stagedArtifactPath, stagedIdentity);
await link(stagedArtifactPath, artifactPath);
artifactPublished = true;
await assertExpectedRegularFile(artifactPath, stagedIdentity);
await hooks.beforeSidecarLink?.();
throwIfAborted(signal);
await assertExpectedRegularFile(artifactPath, stagedIdentity);
await link(sidecarTemp, sidecarPath);
sidecarPublished = true;
await assertSameRegularFile(sidecarPath, sidecarTemp);
await hooks.beforeCommit?.();
throwIfAborted(signal);
await Promise.all([
assertExpectedRegularFile(artifactPath, stagedIdentity),
assertSameRegularFile(sidecarPath, sidecarTemp),
]);
committed = true;
} catch (error) {
errors.push(normalizeCollision(error));
if (!committed) {
if (sidecarPublished) {
await unlinkIfSameFile(sidecarPath, sidecarTemp).catch((rollback) =>
errors.push(rollback),
);
}
if (artifactPublished) {
await unlinkIfExpectedFile(artifactPath, stagedIdentity).catch(
(rollback) => errors.push(rollback),
);
}
}
}
await rm(sidecarTemp, { force: true }).catch((error) => errors.push(error));
if (committed) {
try {
await syncDirectory(path.dirname(artifactPath));
} catch (error) {
errors.push(error);
}
}
if (errors.length === 1) throw errors[0];
if (errors.length > 1) {
throw new AggregateError(
errors,
errors[0] instanceof Error
? errors[0].message
: "Workspace artifact publication failed.",
);
}
}
async function writePrivateExclusive(
filePath: string,
contents: Uint8Array | string,
): Promise<void> {
const handle = await open(
filePath,
constants.O_WRONLY |
constants.O_CREAT |
constants.O_EXCL |
constants.O_NOFOLLOW,
0o600,
);
try {
await handle.writeFile(contents);
await handle.sync();
const fileStat = await handle.stat({ bigint: true });
if (
!fileStat.isFile() ||
fileStat.nlink !== 1n ||
(process.platform !== "win32" && (fileStat.mode & 0o777n) !== 0o600n)
) {
throw new Error(
"Workspace staging did not create a private regular file.",
);
}
} finally {
await handle.close();
}
}
async function assertMissing(filePath: string): Promise<void> {
try {
await lstat(filePath);
} catch (error) {
if (isNodeError(error, "ENOENT")) return;
throw error;
}
throw new SequenceWorkspaceCollisionError();
}
async function unlinkIfSameFile(
destination: string,
staged: string,
): Promise<void> {
try {
const [destinationStat, stagedStat] = await Promise.all([
stat(destination, { bigint: true }),
stat(staged, { bigint: true }),
]);
if (
destinationStat.dev === stagedStat.dev &&
destinationStat.ino === stagedStat.ino
) {
await rm(destination);
}
} catch (error) {
if (!isNodeError(error, "ENOENT")) throw error;
}
}
async function assertSameRegularFile(
destination: string,
staged: string,
): Promise<void> {
const [destinationStat, stagedStat] = await Promise.all([
lstat(destination, { bigint: true }),
lstat(staged, { bigint: true }),
]);
if (
destinationStat.isSymbolicLink() ||
stagedStat.isSymbolicLink() ||
!destinationStat.isFile() ||
!stagedStat.isFile() ||
destinationStat.dev !== stagedStat.dev ||
destinationStat.ino !== stagedStat.ino
) {
throw new Error("The workspace publication target changed during commit.");
}
}
async function assertPrivateStagedFile(
staged: string,
expected: SequenceWorkspaceStagedIdentity,
): Promise<void> {
const fileStat = await lstat(staged, { bigint: true });
if (
fileStat.isSymbolicLink() ||
!fileStat.isFile() ||
fileStat.dev.toString() !== expected.device ||
fileStat.ino.toString() !== expected.inode ||
fileStat.ctimeNs.toString() !== expected.changedAtNanoseconds ||
fileStat.mtimeNs.toString() !== expected.modifiedAtNanoseconds ||
fileStat.size.toString() !== expected.size ||
fileStat.nlink !== 1n ||
(process.platform !== "win32" && (fileStat.mode & 0o777n) !== 0o600n)
) {
throw new Error("Workspace staging changed before publication.");
}
}
async function assertExpectedRegularFile(
destination: string,
expected: SequenceWorkspaceStagedIdentity,
): Promise<void> {
const fileStat = await lstat(destination, { bigint: true });
if (
fileStat.isSymbolicLink() ||
!fileStat.isFile() ||
fileStat.dev.toString() !== expected.device ||
fileStat.ino.toString() !== expected.inode ||
fileStat.size.toString() !== expected.size
) {
throw new Error("The workspace publication target changed during commit.");
}
}
async function unlinkIfExpectedFile(
destination: string,
expected: SequenceWorkspaceStagedIdentity,
): Promise<void> {
try {
await assertExpectedRegularFile(destination, expected);
await rm(destination);
} catch (error) {
if (!isNodeError(error, "ENOENT")) throw error;
}
}
async function syncDirectory(directory: string): Promise<void> {
if (process.platform === "win32") return;
const handle = await open(directory, constants.O_RDONLY);
try {
await handle.sync();
} catch (error) {
if (!isNodeError(error, "EINVAL") && !isNodeError(error, "ENOTSUP")) {
throw error;
}
} finally {
await handle.close();
}
}
function temporaryPath(
destination: string,
nonce: string,
suffix: string,
): string {
return path.join(
path.dirname(destination),
`.${path.basename(destination)}.${nonce}.${suffix}.tmp`,
);
}
function normalizeCollision(error: unknown): unknown {
return isNodeError(error, "EEXIST")
? new SequenceWorkspaceCollisionError()
: error;
}
function sameCanonicalDirectory(left: string, right: string): boolean {
return canonicalDestinationPath(left) === canonicalDestinationPath(right);
}
function canonicalDestinationPath(value: string): string {
const resolved = path.resolve(value);
return process.platform === "win32" ? resolved.toLowerCase() : resolved;
}
function isNodeError(
error: unknown,
code: string,
): error is NodeJS.ErrnoException {
return error instanceof Error && "code" in error && error.code === code;
}
function throwIfAborted(signal?: AbortSignal): void {
if (!signal?.aborted) return;
throw (
signal.reason ??
new DOMException(
"Workspace artifact publication was cancelled.",
"AbortError",
)
);
}
async function acquireDestination(
artifactPath: string,
sidecarPath: string,
signal?: AbortSignal,
): Promise<() => void> {
const key = JSON.stringify(
[artifactPath, sidecarPath].map(canonicalDestinationPath).sort(),
);
const previous = destinationTails.get(key) ?? Promise.resolve();
let openGate!: () => void;
const gate = new Promise<void>((resolve) => {
openGate = resolve;
});
const settled = previous.catch(() => undefined);
const tail = settled.then(() => gate);
destinationTails.set(key, tail);
try {
await waitForAbortable(settled, signal);
throwIfAborted(signal);
} catch (error) {
openGate();
throw error;
}
let released = false;
return () => {
if (released) return;
released = true;
openGate();
void tail.finally(() => {
if (destinationTails.get(key) === tail) destinationTails.delete(key);
});
};
}
async function waitForAbortable(
promise: Promise<void>,
signal?: AbortSignal,
): Promise<void> {
if (signal == null) {
await promise;
return;
}
throwIfAborted(signal);
await new Promise<void>((resolve, reject) => {
let finished = false;
const finish = (callback: () => void) => {
if (finished) return;
finished = true;
signal.removeEventListener("abort", onAbort);
callback();
};
const onAbort = () =>
finish(() =>
reject(
signal.reason ??
new DOMException(
"Workspace artifact publication was cancelled.",
"AbortError",
),
),
);
signal.addEventListener("abort", onAbort, { once: true });
void promise.then(
() => finish(resolve),
(error) => finish(() => reject(error)),
);
});
}
SHA-256: a31e817ee507af4a10e8934dfc3b2452374a32e58ab8630a7db8da1ea4d978d1