From 0332586a99df327779ec49a658c5114ba4174d8e Mon Sep 17 00:00:00 2001 From: philipz Date: Mon, 14 Sep 2026 14:08:07 +0800 Subject: [PATCH 1/2] fix: resolve lock validity check, retry self-blocking, timer leak, and error masking --- src/index.ts | 122 ++++++--- src/regression.test.ts | 552 +++++++++++++++++++++++++++++++++++++++++ tsconfig.json | 10 +- 3 files changed, 651 insertions(+), 33 deletions(-) create mode 100644 src/regression.test.ts diff --git a/src/index.ts b/src/index.ts index e7a1152..8e9ee18 100644 --- a/src/index.ts +++ b/src/index.ts @@ -10,14 +10,16 @@ type Client = IORedisClient | IORedisCluster; // Define script constants. const ACQUIRE_SCRIPT = ` - -- Return 0 if an entry already exists. + -- Return 0 if an entry already exists with a different lock value. for i, key in ipairs(KEYS) do if redis.call("exists", key) == 1 then - return 0 + if redis.call("get", key) ~= ARGV[1] then + return 0 + end end end - -- Create an entry for each provided key. + -- Create or update the entry for each provided key. for i, key in ipairs(KEYS) do redis.call("set", key, ARGV[1], "PX", ARGV[2]) end @@ -156,8 +158,8 @@ export class Lock { return this.redlock.release(this); } - async extend(duration: number): Promise { - return this.redlock.extend(this, duration); + async extend(duration: number, settings?: Partial): Promise { + return this.redlock.extend(this, duration, settings); } } @@ -320,13 +322,15 @@ export default class Redlock extends EventEmitter { (settings?.driftFactor ?? this.settings.driftFactor) * duration ) + 2; - return new Lock( - this, - resources, - value, - attempts, - start + duration - drift - ); + const expiration = start + duration - drift; + if (Date.now() >= expiration) { + throw new ExecutionError( + "The lock validity time has elapsed before quorum was achieved.", + attempts + ); + } + + return new Lock(this, resources, value, attempts, expiration); } catch (error) { // If there was an error acquiring the lock, release any partial lock // state that may exist on a minority of clients. @@ -397,12 +401,29 @@ export default class Redlock extends EventEmitter { (settings?.driftFactor ?? this.settings.driftFactor) * duration ) + 2; + const expiration = start + duration - drift; + if (Date.now() >= expiration) { + await this._execute( + this.scripts.releaseScript, + existing.resources, + [existing.value], + { retryCount: 0 } + ).catch(() => { + // Any error here will be ignored. + }); + + throw new ExecutionError( + "The lock validity time has elapsed before extension was achieved.", + attempts + ); + } + const replacement = new Lock( this, existing.resources, existing.value, attempts, - start + duration - drift + expiration ); return replacement; @@ -713,6 +734,8 @@ export default class Redlock extends EventEmitter { const signal = controller.signal as RedlockAbortSignal; + let running = true; + function queue(): void { timeout = setTimeout( () => (extension = extend()), @@ -724,14 +747,16 @@ export default class Redlock extends EventEmitter { timeout = undefined; try { - lock = await lock.extend(duration); - queue(); + lock = await lock.extend(duration, settings); + if (running) { + queue(); + } } catch (error) { if (!(error instanceof Error)) { throw new Error(`Unexpected thrown ${typeof error}: ${error}.`); } - if (lock.expiration > Date.now()) { + if (running && lock.expiration > Date.now()) { return (extension = extend()); } @@ -745,26 +770,59 @@ export default class Redlock extends EventEmitter { let lock = await this.acquire(resources, duration, settings); queue(); + let routineResult: T | undefined = undefined; + let routineError: unknown; + let routineSucceeded = false; + try { - return await routine(signal); - } finally { - // Clean up the timer. - if (timeout) { - clearTimeout(timeout); - timeout = undefined; - } + routineResult = await routine(signal); + routineSucceeded = true; + } catch (error) { + routineError = error; + } - // Wait for an in-flight extension to finish. - if (extension) { - await extension.catch(() => { - // An error here doesn't matter at all, because the routine has - // already completed, and a release will be attempted regardless. The - // only reason for waiting here is to prevent possible contention - // between the extension and release. - }); - } + running = false; + + // Clean up the timer. + if (timeout) { + clearTimeout(timeout); + timeout = undefined; + } + + // Wait for an in-flight extension to finish. + if (extension) { + await extension.catch(() => { + // An error here doesn't matter at all, because the routine has + // already completed, and a release will be attempted regardless. The + // only reason for waiting here is to prevent possible contention + // between the extension and release. + }); + } + + if (timeout) { + clearTimeout(timeout); + timeout = undefined; + } + let releaseError: unknown; + try { await lock.release(); + } catch (error) { + releaseError = error; + } + + if (!routineSucceeded) { + throw routineError; } + + if (signal.aborted && signal.error) { + throw signal.error; + } + + if (releaseError) { + throw releaseError; + } + + return routineResult as T; } } diff --git a/src/regression.test.ts b/src/regression.test.ts new file mode 100644 index 0000000..98685d6 --- /dev/null +++ b/src/regression.test.ts @@ -0,0 +1,552 @@ +import test from "ava"; +import type { Redis as Client } from "ioredis"; +import Redlock, { ExecutionError, Lock } from "./index.js"; +import type { RedlockAbortSignal, Settings } from "./index.js"; + +const sleep = (ms: number): Promise => + new Promise((r) => setTimeout(r, ms)); + +interface KeyEntry { + value: string; + expiresAt: number; +} + +// --- Script discriminators --- +// +// Pitfall: The fixed ACQUIRE_SCRIPT **also contains** `redis.call("get", key) ~= ARGV[1]` +// (it now checks whether the key blocking us is our own). Any interceptor that identifies +// EXTEND solely using that string will falsely identify ACQUIRE as an extension, allowing +// tests to pass when **no extension ever occurred**. +// Only `redis.call("exists"` is unique to ACQUIRE, and exists both before and after the fix— +// allowing the same mock to work for both versions so red-green comparisons remain meaningful. +const ACQUIRE_MARKER = 'redis.call("exists"'; +const EXTEND_MARKER = 'redis.call("get", key) ~= ARGV[1]'; +const RELEASE_MARKER = 'redis.pcall("del"'; + +const isAcquireScript = (script: string): boolean => + script.includes(ACQUIRE_MARKER); +const isExtendScript = (script: string): boolean => + script.includes(EXTEND_MARKER) && !script.includes(ACQUIRE_MARKER); +const isReleaseScript = (script: string): boolean => + script.includes(RELEASE_MARKER); + +interface AcquireObservation { + granted: boolean; + blocker: string | null; + /** The key that blocked us carried our *own* value (the F2 defect). */ + selfBlocked: boolean; + /** We were granted while one of the keys already held our own value. */ + reacquiredOwn: boolean; +} + +class MockRedisClient { + public store = new Map(); + public delayMs = 0; + public failExtend = false; + public failRelease = false; + public acquireLog: AcquireObservation[] = []; + + private _entry(key: string): KeyEntry | undefined { + const e = this.store.get(key); + if (!e) return undefined; + if (e.expiresAt <= Date.now()) { + this.store.delete(key); + return undefined; + } + return e; + } + + public get(key: string): string | null { + const e = this._entry(key); + return e ? e.value : null; + } + + public set(key: string, value: string, ttlMs: number): void { + this.store.set(key, { value, expiresAt: Date.now() + ttlMs }); + } + + public evalsha(): Promise { + throw new Error("NOSCRIPT No matching script."); + } + + public async eval( + script: string, + numKeys: number, + args: (string | number)[] + ): Promise { + if (this.delayMs > 0) { + await sleep(this.delayMs); + } + + const keys = args.slice(0, numKeys).map(String); + const argv = args.slice(numKeys).map(String); + + if (isAcquireScript(script)) { + const lockValue = argv[0]; + const ttl = Number(argv[1]); + + // The blocking condition is determined by the script text: before the fix, + // it only checks exists (blocking even on identical values, i.e., self-blocking); + // after the fix, it blocks only if the value differs. Hardcoding either behavior + // would make the mock faithful to only one version of src, defeating the purpose + // of red-green comparisons. + const valueAware = script.includes(EXTEND_MARKER); + let reacquiredOwn = false; + + for (const key of keys) { + const e = this._entry(key); + if (!e) continue; + if (valueAware && e.value === lockValue) { + reacquiredOwn = true; + continue; + } + this.acquireLog.push({ + granted: false, + blocker: e.value, + selfBlocked: e.value === lockValue, + reacquiredOwn: false, + }); + return 0; + } + + for (const key of keys) { + this.set(key, lockValue, ttl); + } + this.acquireLog.push({ + granted: true, + blocker: null, + selfBlocked: false, + reacquiredOwn, + }); + return keys.length; + } + + if (isExtendScript(script)) { + if (this.failExtend) { + return 0; + } + const lockValue = argv[0]; + const ttl = Number(argv[1]); + + for (const key of keys) { + const e = this._entry(key); + if (!e || e.value !== lockValue) { + return 0; + } + } + + for (const key of keys) { + this.set(key, lockValue, ttl); + } + return keys.length; + } + + if (isReleaseScript(script)) { + if (this.failRelease) { + return 0; + } + const lockValue = argv[0]; + let deleted = 0; + for (const key of keys) { + const e = this._entry(key); + if (e && e.value === lockValue) { + this.store.delete(key); + deleted++; + } + } + return deleted; + } + + throw new Error("Unknown script in mock"); + } +} + +function makeClients(n: number): MockRedisClient[] { + return Array.from({ length: n }, () => new MockRedisClient()); +} + +test("F1 fix: acquire() throws ExecutionError when round-trip takes longer than validity", async (t) => { + const clients = makeClients(3); + for (const c of clients) { + c.delayMs = 100; + } + + const redlock = new Redlock(clients as unknown as Client[]); + const duration = 50; // duration 50ms, delay 100ms -> validity < 0 + + await t.throwsAsync( + async () => { + await redlock.acquire(["test-resource"], duration); + }, + { + instanceOf: ExecutionError, + message: + /The lock validity time has elapsed before quorum was achieved\./, + } + ); + + // Compensation release must have cleaned up the acquired keys + for (const c of clients) { + t.is( + c.get("test-resource"), + null, + "Partial key must be released on failure" + ); + } +}); + +test("F5 fix: extend() throws ExecutionError when extension round-trip exceeds validity", async (t) => { + const clients = makeClients(3); + const redlock = new Redlock(clients as unknown as Client[]); + + const lock = await redlock.acquire(["extend-resource"], 500); + t.true(lock.expiration > Date.now(), "Lock initially valid"); + + // Inject delay exceeding duration + for (const c of clients) { + c.delayMs = 120; + } + + await t.throwsAsync( + async () => { + await lock.extend(60); + }, + { + instanceOf: ExecutionError, + message: + /The lock validity time has elapsed before extension was achieved\./, + } + ); +}); + +test("F2 fix: retry is NOT blocked by keys left by the client's own previous attempt", async (t) => { + const clients = makeClients(3); + + // Client 2 and 3 occupied by a foreign lock for 50ms + clients[1].set("retry-resource", "FOREIGN", 50); + clients[2].set("retry-resource", "FOREIGN", 50); + + const redlock = new Redlock(clients as unknown as Client[]); + + // Attempt 1: succeeds on client 0, fails on 1 and 2 (no quorum). + // Attempt 2: retry after 60ms when foreign locks expire. + // Prior to fix, client 0 would be blocked by its own key from attempt 1! + const lock = await redlock.acquire(["retry-resource"], 2000, { + retryCount: 2, + retryDelay: 60, + retryJitter: 0, + }); + + t.truthy(lock, "Acquire succeeded on retry without self-blocking"); + t.is(clients[0].get("retry-resource"), lock.value); + t.is(clients[1].get("retry-resource"), lock.value); + t.is(clients[2].get("retry-resource"), lock.value); + + // The assertions above are NOT enough on their own: quorum here is 2 of 3, so + // even while client 0 was self-blocked the acquisition still succeeded via + // clients 1 and 2 — and client 0 still held the same value, left over from + // attempt 1. The defect is only observable on client 0's own votes. + const log = clients[0].acquireLog; + t.false( + log.some((e) => e.selfBlocked), + "Client 0 was never blocked by its own value" + ); + t.true( + log.some((e) => e.granted && !e.reacquiredOwn), + "Attempt 1 was a fresh acquisition on client 0" + ); + t.true( + log.some((e) => e.granted && e.reacquiredOwn), + "Attempt 2 re-acquired client 0 over its own leftover key (the fixed path)" + ); +}); + +test.serial( + "F6 fix: using() cleans up all timers and leaves no unhandled timeout after return", + async (t) => { + const clients = makeClients(3); + const redlock = new Redlock(clients as unknown as Client[], { + retryCount: 0, + retryDelay: 0, + retryJitter: 0, + automaticExtensionThreshold: 50, + }); + + const duration = 200; + let extendStartedResolve: () => void; + const extendStarted = new Promise((r) => (extendStartedResolve = r)); + let openExtendGate: () => void; + const extendGate = new Promise((r) => (openExtendGate = r)); + let extendIntercepted = 0; + let extendFinished = false; + let extendInFlightAtReturn = false; + + // Hold the extension at an explicit gate so the routine can return while the + // extension is genuinely in flight. + // + // The gate is a bare Promise on purpose: a sleep() here would create a + // setTimeout of its own, and this test counts lingering timers — the + // harness must not contribute any. + for (const c of clients) { + const origEval = c.eval.bind(c); + c.eval = async (script, numKeys, args) => { + if (isExtendScript(script)) { + extendIntercepted++; + extendStartedResolve(); + await extendGate; + const result = await origEval(script, numKeys, args); + extendFinished = true; + return result; + } + return origEval(script, numKeys, args); + }; + } + + // Track global setTimeout / clearTimeout + const activeTimers = new Set(); + const origSetTimeout = globalThis.setTimeout; + const origClearTimeout = globalThis.clearTimeout; + + globalThis.setTimeout = (( + fn: (...args: unknown[]) => void, + ms?: number, + ...args: unknown[] + ) => { + const handle = origSetTimeout(() => { + activeTimers.delete(handle); + fn(...args); + }, ms); + activeTimers.add(handle); + return handle; + }) as typeof setTimeout; + + globalThis.clearTimeout = ((handle?: NodeJS.Timeout) => { + if (handle) activeTimers.delete(handle); + return origClearTimeout(handle); + }) as typeof clearTimeout; + + const baseline = new Set(activeTimers); + + try { + const result = await redlock.using( + ["leak-resource"], + duration, + async () => { + await extendStarted; + extendInFlightAtReturn = !extendFinished; + // Let using() reach the "routine done, extension still in flight" + // state before the extension is allowed to complete. setImmediate is + // not a setTimeout, so it is not counted as a lingering timer. + setImmediate(openExtendGate); + return "ROUTINE_DONE"; + } + ); + + t.is(result, "ROUTINE_DONE"); + + // Preconditions: without these, "no lingering timers" would also hold in + // the case where no extension ever ran — a vacuous pass. + t.true(extendIntercepted > 0, "An automatic extension actually ran"); + t.true( + extendInFlightAtReturn, + "The routine returned while the extension was still in flight" + ); + + const leaked = [...activeTimers].filter((h) => !baseline.has(h)); + t.is(leaked.length, 0, "No lingering timers after using() completes"); + } finally { + globalThis.setTimeout = origSetTimeout; + globalThis.clearTimeout = origClearTimeout; + } + } +); + +test("F7 fix: using() propagates routine error without masking when release fails", async (t) => { + const clients = makeClients(3); + for (const c of clients) { + c.failRelease = true; // release will fail + } + + const redlock = new Redlock(clients as unknown as Client[], { + retryCount: 0, + automaticExtensionThreshold: 100, + }); + + const customError = new Error("Custom Business Logic Error"); + + await t.throwsAsync( + async () => { + await redlock.using(["err-resource"], 500, async () => { + throw customError; + }); + }, + { + is: customError, + message: "Custom Business Logic Error", + } + ); +}); + +test("F7 fix: using() propagates signal.error when abort occurs rather than release error", async (t) => { + const clients = makeClients(3); + for (const c of clients) { + c.failExtend = true; + c.failRelease = true; + } + + const redlock = new Redlock(clients as unknown as Client[], { + retryCount: 0, + retryDelay: 0, + retryJitter: 0, + automaticExtensionThreshold: 50, + }); + + // Lock duration 160ms, extension fails -> abort signal triggered + let captured: RedlockAbortSignal | undefined; + const error = await t.throwsAsync( + async () => { + await redlock.using(["abort-resource"], 160, async (signal) => { + captured = signal; + await sleep(250); // wait for extension to fail and lock to expire + t.true(signal.aborted, "Signal aborted"); + return "ROUTINE_FINISH"; + }); + }, + { + instanceOf: ExecutionError, + message: + /The operation was unable to achieve a quorum during its retry window\./, + } + ); + + // The message alone proves nothing: before the fix, the error thrown from the + // `finally` block (the failing release) carried exactly the same message. The + // discriminating assertion is *identity* — the thrown error must be the very + // object recorded on the signal, not a look-alike from the release path. + t.truthy(captured?.error, "signal.error was set when the lock was lost"); + t.is( + error, + captured?.error, + "using() rethrows signal.error itself, not the release error" + ); +}); + +// Counting eval calls cannot measure retryCount inside `using()`: when an +// extension fails, `using()` re-invokes extend() recursively for as long as the +// lock is still valid (src/index.ts, the `running && lock.expiration > Date.now()` +// branch). With retryDelay 0 that is a tight microtask loop — measured at 4149 +// eval calls in ~50ms. So the two properties are asserted separately, each in a +// deterministic way. + +test.serial( + "settings pass-through: using() forwards its per-call settings to lock.extend", + async (t) => { + const clients = makeClients(3); + const redlock = new Redlock(clients as unknown as Client[], { + retryCount: 5, + retryDelay: 200, + }); + + const perCall: Partial = { + retryCount: 0, + retryDelay: 0, + retryJitter: 0, + automaticExtensionThreshold: 50, + }; + + // Capture the second argument `using()` hands to Lock#extend. This is the + // property under test, observed directly rather than inferred from counts. + const captured: (Partial | undefined)[] = []; + const originalExtend = Lock.prototype.extend; + Lock.prototype.extend = function ( + this: Lock, + duration: number, + settings?: Partial + ): Promise { + captured.push(settings); + return originalExtend.call(this, duration, settings); + }; + + try { + // Extension succeeds here: no failure loop, exactly one extension. + const result = await redlock.using( + ["settings-resource"], + 200, + perCall, + async () => { + await sleep(250); + return "DONE"; + } + ); + t.is(result, "DONE"); + } finally { + Lock.prototype.extend = originalExtend; + } + + // Precondition: an automatic extension actually happened. + t.true(captured.length > 0, "Lock#extend was actually called"); + + // Before the fix, `using()` called `lock.extend(duration)` with no second + // argument, so every captured value was undefined. + t.not(captured[0], undefined, "extend received a settings argument"); + + // `using()` merges the per-call settings over the instance settings, so the + // captured object is the full merged set (it also carries driftFactor etc.). + // What matters is that the per-call overrides won. + t.like( + captured[0], + perCall, + "using() forwards its per-call settings to extend" + ); + t.is( + captured[0]?.retryCount, + 0, + "retryCount is the per-call 0, not the instance's 5" + ); + } +); + +test("settings pass-through: lock.extend honours the retryCount it is given", async (t) => { + const clients = makeClients(3); + // Instance is configured with retryCount 5 -> 6 attempts per client if the + // per-call settings are ignored. + const redlock = new Redlock(clients as unknown as Client[], { + retryCount: 5, + retryDelay: 0, + retryJitter: 0, + }); + + const lock = await redlock.acquire(["retrycount-resource"], 10_000); + + // Only now start failing extensions, so acquisition is unaffected. + let extendAttemptsCount = 0; + for (const c of clients) { + const origEval = c.eval.bind(c); + c.eval = async (script, numKeys, args) => { + if (isExtendScript(script)) { + extendAttemptsCount++; + return 0; // force failure + } + return origEval(script, numKeys, args); + }; + } + + // Called directly, so there is no `using()` retry loop in play: the attempt + // count is exactly (retryCount + 1) * clients. + await t.throwsAsync( + async () => { + await lock.extend(10_000, { + retryCount: 0, + retryDelay: 0, + retryJitter: 0, + }); + }, + { instanceOf: ExecutionError } + ); + + t.is( + extendAttemptsCount, + 3, + "retryCount: 0 means exactly 1 attempt on each of the 3 clients" + ); +}); diff --git a/tsconfig.json b/tsconfig.json index 32566b5..e3c17c6 100644 --- a/tsconfig.json +++ b/tsconfig.json @@ -11,7 +11,15 @@ "strict": true, "esModuleInterop": true, "skipLibCheck": false, - "sourceMap": true + "sourceMap": true, + // Only auto-load @types/node. **This is not a stylistic preference, but a build necessity**: + // devDependency `@informalsystems/quint@0.32.0` transitively pulls in + // `@types/lodash.clonedeep` → `@types/lodash@4.17.x`, whose .d.ts uses syntax that + // TypeScript 4.6 cannot parse (TS1005 '?' expected), causing `tsc` to exit with code 2. + // `@types/lodash` is never used by this project; it was merely pulled in by the default auto-inclusion of all `@types/*`. + // Restricting `types` still allows `import ... from "ioredis"` to resolve normally via `@types/ioredis` + // (`types` only controls **global auto-inclusion**, without affecting explicit module import resolution). + "types": ["node"] }, "include": ["./src/**/*"] } From 1d4b3b124cb08a23347a2636c0b8b2e95f81c74d Mon Sep 17 00:00:00 2001 From: philipz Date: Fri, 25 Sep 2026 15:42:22 +0800 Subject: [PATCH 2/2] fix: resolve attempt as against when even nodes tie to prevent hang In even-node configurations (e.g. N=2, quorumSize=2), a tie vote (1 for, 1 against) would result in all votes being collected without either side reaching quorumSize. Previously, only done() was invoked, leaving the outer attempt Promise permanently pending. This fix resolves the attempt with vote 'against' once all votes are in without reaching quorum. --- src/index.ts | 14 ++++++++++++++ src/regression.test.ts | 25 +++++++++++++++++++++++++ 2 files changed, 39 insertions(+) diff --git a/src/index.ts b/src/index.ts index 8e9ee18..d650440 100644 --- a/src/index.ts +++ b/src/index.ts @@ -556,6 +556,20 @@ export default class Redlock extends EventEmitter { stats.membershipSize ) { done(); + + // In an even-node configuration, a tie vote can result in all votes + // being collected without either side reaching quorumSize. Since a quorum + // in favor was not reached, the attempt has failed and must resolve as "against". + if ( + stats.votesFor.size < stats.quorumSize && + stats.votesAgainst.size < stats.quorumSize + ) { + resolve({ + vote: "against", + stats: statsPromise, + start, + }); + } } }; diff --git a/src/regression.test.ts b/src/regression.test.ts index 98685d6..17b6882 100644 --- a/src/regression.test.ts +++ b/src/regression.test.ts @@ -550,3 +550,28 @@ test("settings pass-through: lock.extend honours the retryCount it is given", as "retryCount: 0 means exactly 1 attempt on each of the 3 clients" ); }); + +test("Liveness defect: acquire() throws ExecutionError instead of hanging when even nodes tie", async (t) => { + // In an even-node configuration (N=2, quorumSize=2), a 1:1 tie vote + // (1 for, 1 against) must resolve as a failure ("against") rather than + // leaving the outer attempt Promise permanently pending. + const clients = makeClients(2); + clients[1].set("tie-resource", "OCCUPIED_BY_ANOTHER_CLIENT", 10_000); + + const redlock = new Redlock(clients as unknown as Client[], { + retryCount: 0, + retryDelay: 0, + retryJitter: 0, + }); + + await t.throwsAsync( + async () => { + await redlock.acquire(["tie-resource"], 1000); + }, + { + instanceOf: ExecutionError, + message: + /The operation was unable to achieve a quorum during its retry window\./, + } + ); +});