Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
195 changes: 131 additions & 64 deletions core/status.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,7 @@
// framework logic shared by every wrapper.

import { randomBytes } from "node:crypto";
import { open, readdir, readFile, rm, stat, unlink, utimes } from "node:fs/promises";
import { open, readdir, readFile, rename, rm, stat, unlink, utimes } from "node:fs/promises";
import { basename, dirname, join } from "node:path";
import type {
CarryForwardEntry,
Expand Down Expand Up @@ -457,6 +457,8 @@ export interface AcquireLockOptions {
* hung; default {@link STALE_LOCK_MS}. A dead holder is removed at once.
*/
staleMs?: number;
/** @internal File operations seam for deterministic lock race tests. */
fsOps?: { stat: typeof stat; rename: typeof rename; unlink: typeof unlink };
}

/**
Expand Down Expand Up @@ -490,6 +492,7 @@ export interface AcquireLockOptions {
export async function acquireLock(lockPath: string, options: AcquireLockOptions = {}): Promise<LockHandle> {
const timeoutMs = options.timeoutMs ?? LOCK_TIMEOUT_MS;
const staleMs = options.staleMs ?? STALE_LOCK_MS;
const fsOps = options.fsOps ?? { stat, rename, unlink };
const startedAt = Date.now();
const dir = dirname(lockPath);
const base = basename(lockPath);
Expand Down Expand Up @@ -524,88 +527,152 @@ export async function acquireLock(lockPath: string, options: AcquireLockOptions
}, Math.max(50, Math.floor(staleMs / 4)));
heartbeat.unref?.();

const giveUp = async (): Promise<never> => {
clearInterval(heartbeat);
await rm(ticketPath, { force: true }).catch(() => undefined);
const giveUp = (): never => {
throw new Error(`Timed out waiting for lock: ${lockPath}`);
};

while (true) {
const names = await readdir(dir).catch(() => [] as string[]);
let blocked = false;

// A pre-#355 lock file: honour it while fresh, remove it when stale.
if (names.includes(base)) {
const legacyStat = await stat(lockPath).catch(() => null);
if (legacyStat && Date.now() - legacyStat.mtimeMs <= staleMs) blocked = true;
else if (legacyStat) {
const holder = await describeLockHolder(lockPath);
if (await removeIfPresent(lockPath)) brokeStale = holder;
let acquired = false;
try {
while (true) {
const names = await readdir(dir).catch(() => [] as string[]);
let blocked = false;
// Failed cleanup must not leave permanent tombstones or claim markers.
for (const name of names) {
const claimPrefix = `${base}.claim.`;
if (!name.startsWith(`${base}.reaped.`) && !name.startsWith(claimPrefix)) continue;
const path = join(dir, name);
const fileStat = await fsOps.stat(path).catch(() => null);
const claimSourceGone = name.startsWith(claimPrefix)
&& !(await pathExists(join(dir, name.slice(claimPrefix.length))));
if (claimSourceGone || (fileStat && Date.now() - fileStat.mtimeMs > staleMs)) {
await rm(path, { force: true }).catch(() => undefined);
}
}
}

// Someone is between taking a number and writing their ticket: their
// number may be earlier than ours. Wait, unless they died in the door.
for (const name of names) {
if (!name.startsWith(`${base}.c.`) || name === basename(choosingPath)) continue;
const path = join(dir, name);
if (await ownerIsGone(path, staleMs)) {
await rm(path, { force: true }).catch(() => undefined);
continue;
// A pre-#355 lock file: honour it while fresh, remove it when stale.
if (names.includes(base)) {
const legacyStat = await fsOps.stat(lockPath).catch((error: NodeJS.ErrnoException) => {
if (error.code !== "ENOENT") blocked = true;
return null;
});
if (legacyStat && Date.now() - legacyStat.mtimeMs <= staleMs) blocked = true;
else if (legacyStat) {
const holder = await describeLockHolder(lockPath);
const claim = await removeIfPresent(lockPath, lockPath, fsOps);
if (claim === "claimed") brokeStale = holder;
if (claim === "busy") blocked = true;
}
}
blocked = true;
}

// Every ticket ahead of ours whose owner is alive blocks us; a dead or
// hung owner's ticket is removed by its own name.
if (!blocked) {
const ahead = names.filter((name) => name.startsWith(`${base}.t.`) && name < ticketName).sort();
for (const name of ahead) {
// Someone is between taking a number and writing their ticket: their
// number may be earlier than ours. Wait, unless they died in the door.
for (const name of names) {
if (!name.startsWith(`${base}.c.`) || name === basename(choosingPath)) continue;
const path = join(dir, name);
if (await ownerIsGone(path, staleMs)) {
// Recorded only by the waiter whose rm actually removed it; the
// others saw the same dead ticket and removed nothing.
const holder = await describeLockHolder(path);
if (await removeIfPresent(path)) brokeStale = holder;
if (await ownerIsGone(path, staleMs, fsOps)) {
await rm(path, { force: true }).catch(() => undefined);
continue;
}
blocked = true;
break;
}
}

if (!blocked) {
// Our own ticket must still be there: a waiter that judged us hung
// (the machine slept past staleMs) has already let someone in.
if (!(await pathExists(ticketPath))) return giveUp();
return {
release: async () => {
clearInterval(heartbeat);
await rm(ticketPath, { force: true }).catch(() => undefined);
},
...(brokeStale && { brokeStale }),
};
}
// Every ticket ahead of ours whose owner is alive blocks us; a dead or
// hung owner's ticket is removed by its own name.
if (!blocked) {
const ahead = names.filter((name) => name.startsWith(`${base}.t.`) && name < ticketName).sort();
for (const name of ahead) {
const path = join(dir, name);
if (await ownerIsGone(path, staleMs, fsOps)) {
// Only the waiter that claims the dead ticket records its holder;
// other waiters cannot claim the same path.
const holder = await describeLockHolder(path);
const claim = await removeIfPresent(path, lockPath, fsOps);
if (claim === "claimed") brokeStale = holder;
if (claim === "busy") {
blocked = true;
break;
}
continue;
}
blocked = true;
break;
}
}

if (Date.now() - startedAt > timeoutMs) return giveUp();
await sleep(LOCK_RETRY_MS);
if (!blocked) {
// Our own ticket must still be there: a waiter that judged us hung
// (the machine slept past staleMs) has already let someone in.
if (!(await pathExists(ticketPath))) giveUp();
acquired = true;
return {
release: async () => {
clearInterval(heartbeat);
await rm(ticketPath, { force: true }).catch(() => undefined);
},
...(brokeStale && { brokeStale }),
};
}

if (Date.now() - startedAt > timeoutMs) giveUp();
await sleep(LOCK_RETRY_MS);
}
} finally {
if (!acquired) {
clearInterval(heartbeat);
await rm(ticketPath, { force: true }).catch(() => undefined);
}
}
}

/**
* Remove `path`; true when this call removed it, false when it was already
* gone. `unlink`, not `rm`: `fs.promises.rm` reports success to every one
* of several concurrent callers, and the point here is to know which one
* actually took the file away.
* Atomically move a stale claim out of the queue. Only the waiter whose
* rename succeeds reports the break. The tombstone is outside the ticket
* namespace, so even a Windows unlink failure cannot revive the claim.
*/
async function removeIfPresent(path: string): Promise<boolean> {
async function removeIfPresent(
path: string,
lockPath: string,
fsOps: NonNullable<AcquireLockOptions["fsOps"]>,
): Promise<"claimed" | "gone" | "busy"> {
const tombstone = `${lockPath}.reaped.${randomBytes(8).toString("hex")}`;
// Windows can let concurrent renames of the same source both report
// success. This stable, exclusive marker decides who may claim that name.
const claimPath = `${lockPath}.claim.${basename(path)}`;
try {
await unlink(path);
return true;
await writeExclusive(claimPath, `${process.pid}\n`);
} catch (error) {
if ((error as NodeJS.ErrnoException).code === "ENOENT") return false;
if ((error as NodeJS.ErrnoException).code === "EEXIST") return "busy";
throw error;
}
try {
await fsOps.rename(path, tombstone);
} catch (error) {
await rm(claimPath, { force: true }).catch(() => undefined);
if ((error as NodeJS.ErrnoException).code === "ENOENT") return "gone";
// Windows can report contention as a permission or sharing error.
// An existing source blocks this round; a missing one lost the race.
if (["EPERM", "EBUSY", "EACCES"].includes((error as NodeJS.ErrnoException).code ?? "")) {
try {
await fsOps.stat(path);
} catch (statError) {
if ((statError as NodeJS.ErrnoException).code === "ENOENT") return "gone";
throw statError;
}
return "busy";
}
throw error;
}
// Keep the claim marker until stale cleanup: Windows can briefly expose
// the source path after rename reports success. Neither file is a ticket.
await fsOps.unlink(tombstone).catch(() => undefined);
try {
await fsOps.stat(path);
} catch (error) {
if ((error as NodeJS.ErrnoException).code === "ENOENT") {
await rm(claimPath, { force: true }).catch(() => undefined);
}
}
return "claimed";
}

/** Create `path` exclusively with `content`; the descriptor is closed either way (#131). */
Expand All @@ -624,12 +691,12 @@ async function writeExclusive(path: string, content: string): Promise<void> {
* cannot be signalled (another user's process) counts as alive. A file that
* vanished while we looked is gone, and so is its owner's claim.
*/
async function ownerIsGone(path: string, staleMs: number): Promise<boolean> {
async function ownerIsGone(path: string, staleMs: number, fsOps: NonNullable<AcquireLockOptions["fsOps"]>): Promise<boolean> {
let fileStat;
try {
fileStat = await stat(path);
} catch {
return true;
fileStat = await fsOps.stat(path);
} catch (error) {
return (error as NodeJS.ErrnoException).code === "ENOENT";
}
if (Date.now() - fileStat.mtimeMs > staleMs) return true;
const { pid } = await describeLockHolder(path);
Expand Down
Loading
Loading