From 2d2fc5f2b6a917377e6dee1df0e336dd30204913 Mon Sep 17 00:00:00 2001 From: Harsh Rai Date: Wed, 26 Nov 2025 19:18:51 +0530 Subject: [PATCH 01/10] Sentinel scan iterator --- docs/sentinel.md | 24 +++ packages/client/lib/sentinel/index.spec.ts | 233 +++++++++++++++++++++ packages/client/lib/sentinel/index.ts | 51 ++++- 3 files changed, 307 insertions(+), 1 deletion(-) diff --git a/docs/sentinel.md b/docs/sentinel.md index 863d165b506..10f7dead1ad 100644 --- a/docs/sentinel.md +++ b/docs/sentinel.md @@ -160,3 +160,27 @@ try { clientLease.release(); } ``` + +## Scan Iterator + +The sentinel client supports `scanIterator` for iterating over keys on the master node: + +```javascript +for await (const keys of sentinel.scanIterator()) { + // ... +} +``` + +If a failover occurs during the scan, the iterator will automatically restart from the beginning on the new master to ensure all keys are covered. This may result in duplicate keys being yielded. If your application requires processing each key exactly once, you should implement a deduplication mechanism (like a `Set` or Bloom filter). + +```javascript +const processed = new Set(); +for await (const keys of sentinel.scanIterator()) { + for (const key of keys) { + if (processed.has(key)) continue; + processed.add(key); + + // process key + } +} +``` diff --git a/packages/client/lib/sentinel/index.spec.ts b/packages/client/lib/sentinel/index.spec.ts index 6dc1647a1e1..744beac9ea8 100644 --- a/packages/client/lib/sentinel/index.spec.ts +++ b/packages/client/lib/sentinel/index.spec.ts @@ -1379,4 +1379,237 @@ describe('legacy tests', () => { assert.equal(csc.stats().hitCount, 6); }) }); + + describe('scanIterator tests', () => { + testUtils.testWithClientSentinel('should iterate through all keys in normal operation', async sentinel => { + // Set up test data + const testKeys = new Set(); + const entries: Array = []; + + // Create 50 test keys to ensure we get multiple scan iterations + for (let i = 0; i < 50; i++) { + const key = `scantest:${i}`; + testKeys.add(key); + entries.push(key, `value${i}`); + } + + // Insert all test data + await sentinel.mSet(entries); + + // Collect all keys using scanIterator + const foundKeys = new Set(); + for await (const keyBatch of sentinel.scanIterator({ MATCH: 'scantest:*' })) { + for (const key of keyBatch) { + foundKeys.add(key); + } + } + + // Verify all keys were found + assert.deepEqual(testKeys, foundKeys); + }, GLOBAL.SENTINEL.OPEN); + + testUtils.testWithClientSentinel('should respect MATCH pattern', async sentinel => { + // Set up test data with different patterns + await sentinel.mSet([ + 'match:1', 'value1', + 'match:2', 'value2', + 'nomatch:1', 'value3', + 'nomatch:2', 'value4' + ]); + + const foundKeys = new Set(); + for await (const keyBatch of sentinel.scanIterator({ MATCH: 'match:*' })) { + for (const key of keyBatch) { + foundKeys.add(key); + } + } + + const expectedKeys = new Set(['match:1', 'match:2']); + assert.deepEqual(foundKeys, expectedKeys); + }, GLOBAL.SENTINEL.OPEN); + }); + + describe('scanIterator with master failover', () => { + const config: RedisSentinelConfig = { sentinelName: "test", numberOfNodes: 3, password: undefined }; + const frame = new SentinelFramework(config); + let sentinel: RedisSentinelType | undefined; + const tracer: Array = []; + + before(async function () { + this.timeout(60000); + await frame.spawnRedisSentinel(); + await steadyState(frame); + }); + + afterEach(async function () { + if (sentinel !== undefined) { + sentinel.destroy(); + sentinel = undefined; + } + }); + + after(async function () { + this.timeout(60000); + await frame.cleanup(); + }); + + it('should restart scan from beginning when master changes during iteration', async function () { + this.timeout(60000); + + sentinel = frame.getSentinelClient({ scanInterval: 1000 }); + sentinel.setTracer(tracer); + sentinel.on("error", () => {}); + await sentinel.connect(); + + // Set up test data + const testKeys = new Set(); + const entries: Array = []; + + for (let i = 0; i < 100; i++) { + const key = `failovertest:${i}`; + testKeys.add(key); + entries.push(key, `value${i}`); + } + + await sentinel.mSet(entries); + // Wait for addded keys to be replicated + await setTimeout(2000); + + let masterChangeDetected = false; + let iterationCount = 0; + const foundKeys = new Set(); + + // Listen for manifest change events + sentinel.on("topology-change", (event: RedisSentinelEvent) => { + if (event.type === "MASTER_CHANGE") { + masterChangeDetected = true; + tracer.push(`Master change detected during scan: ${event.node.port}`); + } + }); + + // Get the current master node before starting scan + const originalMaster = sentinel.getMasterNode(); + tracer.push(`Original master port: ${originalMaster?.port}`); + + // Start scanning with a small COUNT to ensure multiple iterations + const scanIterator = sentinel.scanIterator({ + MATCH: "failovertest:*", + COUNT: 10, + }); + + // Consume the scan iterator + try { + for await (const keyBatch of scanIterator) { + iterationCount++; + if (iterationCount === 1) { + tracer.push( + `Triggering master failover by stopping node ${originalMaster?.port}` + ); + await frame.stopNode(originalMaster!.port.toString()); + tracer.push(`Master node stopped`); + } + tracer.push( + `Scan iteration ${iterationCount}, got ${keyBatch.length} keys` + ); + + for (const key of keyBatch) { + foundKeys.add(key); + } + } + } catch (error) { + tracer.push(`Error during scan: ${error}`); + throw error; + } + + // Verify that master change was detected + assert.equal( + masterChangeDetected, + true, + "Master change should have been detected" + ); + + // Verify that we eventually got all keys despite the master change + assert.equal( + foundKeys.size, + testKeys.size, + "Should find all keys despite master failover" + ); + assert.deepEqual( + foundKeys, + testKeys, + "Found keys should match test keys" + ); + + // Verify that the master actually changed + const newMaster = sentinel.getMasterNode(); + tracer.push(`New master port: ${newMaster?.port}`); + assert.notEqual( + originalMaster?.port, + newMaster?.port, + "Master should have changed" + ); + + tracer.push( + `Test completed successfully with ${iterationCount} scan iterations` + ); + }); + + it('should handle master change at scan start', async function () { + this.timeout(60000); + + sentinel = frame.getSentinelClient({ scanInterval: 1000 }); + sentinel.setTracer(tracer); + sentinel.on("error", () => { }); + await sentinel.connect(); + + // Set up test data + const entries: Array = []; + for (let i = 0; i < 30; i++) { + entries.push(`startfailover:${i}`, `value${i}`); + } + await sentinel.mSet(entries); + + // Wait for addded keys to be replicated + await setTimeout(2000); + + // Get original master and trigger immediate failover + const originalMaster = sentinel.getMasterNode(); + + // Stop master immediately before starting scan + await frame.stopNode(originalMaster!.port.toString()); + + let masterChangeDetected = false; + let masterChangeResolve: () => void; + const masterChangePromise = new Promise((resolve) => { + masterChangeResolve = resolve; + }); + + // Listen for manifest change events + sentinel.on('topology-change', (event: RedisSentinelEvent) => { + if (event.type === "MASTER_CHANGE") { + masterChangeDetected = true; + tracer.push(`Master change detected during scan: ${event.node.port}`); + if (masterChangeResolve) masterChangeResolve(); + } + }); + + await masterChangePromise; + + // Now start scan - should work with new master + const foundKeys = new Set(); + for await (const keyBatch of sentinel.scanIterator({ MATCH: 'startfailover:*' })) { + for (const key of keyBatch) { + foundKeys.add(key); + } + } + + assert.equal(masterChangeDetected, true, 'Master change should have been detected'); + // Should find all keys even though master changed before scan started + assert.equal(foundKeys.size, 30); + + // Verify master actually changed + const newMaster = sentinel.getMasterNode(); + assert.notEqual(originalMaster?.port, newMaster?.port); + }); + }); }); diff --git a/packages/client/lib/sentinel/index.ts b/packages/client/lib/sentinel/index.ts index bb5eecd5ca1..e5cc486f9c6 100644 --- a/packages/client/lib/sentinel/index.ts +++ b/packages/client/lib/sentinel/index.ts @@ -1,5 +1,5 @@ import { EventEmitter } from 'node:events'; -import { CommandArguments, RedisFunctions, RedisModules, RedisScripts, ReplyUnion, RespVersions, TypeMapping, DEFAULT_RESP } from '../RESP/types'; +import { CommandArguments, RedisArgument, RedisFunctions, RedisModules, RedisScripts, ReplyUnion, RespVersions, TypeMapping, DEFAULT_RESP } from '../RESP/types'; import RedisClient, { AnyRedisClientOptions, RedisClientOptions, RedisClientType } from '../client'; import { CommandOptions } from '../client/commands-queue'; import { attachConfig } from '../commander'; @@ -19,11 +19,16 @@ import { RedisTcpSocketOptions } from '../client/socket'; import { BasicPooledClientSideCache, PooledClientSideCacheProvider } from '../client/cache'; import { ClientIdentity, ClientRole, generateClientId } from '../client/identity'; import { DEFAULT_COMMAND_TIMEOUT } from '../defaults'; +import { ScanOptions } from '../commands/SCAN'; interface ClientInfo { id: number; } +interface ScanIteratorOptions { + cursor?: RedisArgument; +} + export class RedisSentinelClient< M extends RedisModules, F extends RedisFunctions, @@ -654,6 +659,50 @@ export default class RedisSentinel< this._self.#internal.setTracer(tracer); } + + async *scanIterator( + this: RedisSentinelType, + options?: ScanOptions & ScanIteratorOptions + ) { + // Acquire a master client lease + const masterClient = await this.acquire(); + let cursor = options?.cursor ?? "0"; + let shouldRestart = false; + + // Set up topology change listener + const handleTopologyChange = (event: RedisSentinelEvent) => { + if (event.type === "MASTER_CHANGE") { + shouldRestart = true; + } + }; + + // Listen for master changes + this.on("topology-change", handleTopologyChange); + + try { + do { + // Check if we need to restart due to master change + if (shouldRestart) { + cursor = "0"; + shouldRestart = false; + } + + const reply = await masterClient.scan(cursor, options); + // If a topology change happened during the scan command (which caused a retry), + // the reply is from the new master using the old cursor. We should discard it + // and let the loop restart the scan from cursor "0". + if (shouldRestart) { + continue; + } + cursor = reply.cursor; + yield reply.keys; + } while (cursor !== "0"); + } finally { + // Clean up: remove event listener and release the client + this.removeListener("topology-change", handleTopologyChange); + masterClient.release(); + } + } } export class RedisSentinelInternal< From 074feecb6de78f58e07533d4b221775b96fd710f Mon Sep 17 00:00:00 2001 From: Harsh Rai Date: Fri, 28 Nov 2025 23:15:35 +0530 Subject: [PATCH 02/10] fix: cleaning up frame in between tests --- packages/client/lib/sentinel/index.spec.ts | 8 +++----- 1 file changed, 3 insertions(+), 5 deletions(-) diff --git a/packages/client/lib/sentinel/index.spec.ts b/packages/client/lib/sentinel/index.spec.ts index 744beac9ea8..b2c9d95e052 100644 --- a/packages/client/lib/sentinel/index.spec.ts +++ b/packages/client/lib/sentinel/index.spec.ts @@ -1435,21 +1435,19 @@ describe('legacy tests', () => { let sentinel: RedisSentinelType | undefined; const tracer: Array = []; - before(async function () { + beforeEach(async function () { this.timeout(60000); await frame.spawnRedisSentinel(); + await frame.getAllRunning(); await steadyState(frame); }); afterEach(async function () { + this.timeout(60000); if (sentinel !== undefined) { sentinel.destroy(); sentinel = undefined; } - }); - - after(async function () { - this.timeout(60000); await frame.cleanup(); }); From 27f3d53514ea661b0879ca2d909718e088d8fe9f Mon Sep 17 00:00:00 2001 From: Nikolay Karadzhov Date: Tue, 2 Jun 2026 16:08:35 +0300 Subject: [PATCH 03/10] fix(sentinel): acquire master lease per scan page, not for entire iteration The original implementation acquired the master client lease once and held it for the whole iteration. With the default `masterPoolSize` of 1, any command issued from inside the `for await` loop body (e.g. `sentinel.mGet(keys)`) would wait for a free slot that the iterator never released, hanging the caller indefinitely (reported by @TheIndra55). Move the acquire/release into the per-page loop so the pool is free between pages and consumer commands can run. The persistent `topology-change` listener still resets the cursor on `MASTER_CHANGE`, and if a failover occurs during a scan call the partial reply is discarded. Docs: - Clarify that the iterator inherits SCAN's standard guarantees plus an extra source of duplicates on failover. - Replace the "Set or Bloom filter" wording with a `Set` recommendation and a separate note that Bloom filters trade correctness for memory. - Document the per-page lease behaviour so users know the iterator is safe to combine with other commands inside the loop body. Tests: add a regression test that issues `mGet` from inside the loop on the default `masterPoolSize: 1` setup and asserts no deadlock. Co-Authored-By: Claude Opus 4.7 (1M context) --- docs/sentinel.md | 16 ++++++++++++- packages/client/lib/sentinel/index.spec.ts | 20 ++++++++++++++++ packages/client/lib/sentinel/index.ts | 27 ++++++++++++---------- 3 files changed, 50 insertions(+), 13 deletions(-) diff --git a/docs/sentinel.md b/docs/sentinel.md index 10f7dead1ad..e150ffe1de0 100644 --- a/docs/sentinel.md +++ b/docs/sentinel.md @@ -171,7 +171,15 @@ for await (const keys of sentinel.scanIterator()) { } ``` -If a failover occurs during the scan, the iterator will automatically restart from the beginning on the new master to ensure all keys are covered. This may result in duplicate keys being yielded. If your application requires processing each key exactly once, you should implement a deduplication mechanism (like a `Set` or Bloom filter). +### Semantics and differences from standalone `scanIterator` + +The standalone `RedisClient.scanIterator()` inherits SCAN's documented guarantees: a full iteration returns every key present from start to end, and may return a key multiple times. See the [SCAN guarantees](https://redis.io/docs/latest/commands/scan/#scan-guarantees) page. + +The sentinel iterator adds one extra source of duplicates: **master failover**. If the master changes mid-iteration (detected via the `topology-change` event with `type: "MASTER_CHANGE"`), the cursor is invalidated (SCAN cursors are node-local) and the iterator restarts from cursor `0` on the new master. Keys already yielded before the failover may be yielded again from the new master. + +Because Redis replication is asynchronous, the new master may also have a slightly different keyset than the old master at the moment of promotion — writes that had not yet replicated will be missing, and writes accepted on the new master after promotion will be present. + +If your processing must be exactly-once, deduplicate with a `Set`: ```javascript const processed = new Set(); @@ -184,3 +192,9 @@ for await (const keys of sentinel.scanIterator()) { } } ``` + +For very large keyspaces a `Set` may be memory-prohibitive. A Bloom filter is a lower-memory alternative but is **not** suitable for strict exactly-once processing: its false positives will cause some real keys to be skipped. Use it only when occasional skips are acceptable. + +### Pool behaviour + +The iterator acquires a master client lease only for the duration of each `SCAN` call and releases it before yielding to the consumer. This means commands issued from inside the `for await` loop body (e.g. `sentinel.mGet(keys)`) will not deadlock against the iterator, even with the default `masterPoolSize` of `1`. diff --git a/packages/client/lib/sentinel/index.spec.ts b/packages/client/lib/sentinel/index.spec.ts index b2c9d95e052..d31d689b259 100644 --- a/packages/client/lib/sentinel/index.spec.ts +++ b/packages/client/lib/sentinel/index.spec.ts @@ -1427,6 +1427,26 @@ describe('legacy tests', () => { const expectedKeys = new Set(['match:1', 'match:2']); assert.deepEqual(foundKeys, expectedKeys); }, GLOBAL.SENTINEL.OPEN); + + testUtils.testWithClientSentinel('should not deadlock when consumer runs commands inside the loop with masterPoolSize 1', async sentinel => { + const entries: Array = []; + for (let i = 0; i < 20; i++) { + entries.push(`deadlock:${i}`, `value${i}`); + } + await sentinel.mSet(entries); + + const collected = new Map(); + for await (const keyBatch of sentinel.scanIterator({ MATCH: 'deadlock:*', COUNT: 5 })) { + if (keyBatch.length === 0) continue; + const values = await sentinel.mGet(keyBatch); + keyBatch.forEach((key, i) => collected.set(key, values[i])); + } + + assert.equal(collected.size, 20); + for (let i = 0; i < 20; i++) { + assert.equal(collected.get(`deadlock:${i}`), `value${i}`); + } + }, GLOBAL.SENTINEL.OPEN); }); describe('scanIterator with master failover', () => { diff --git a/packages/client/lib/sentinel/index.ts b/packages/client/lib/sentinel/index.ts index e5cc486f9c6..4fa08bdf328 100644 --- a/packages/client/lib/sentinel/index.ts +++ b/packages/client/lib/sentinel/index.ts @@ -664,43 +664,46 @@ export default class RedisSentinel< this: RedisSentinelType, options?: ScanOptions & ScanIteratorOptions ) { - // Acquire a master client lease - const masterClient = await this.acquire(); let cursor = options?.cursor ?? "0"; let shouldRestart = false; - // Set up topology change listener const handleTopologyChange = (event: RedisSentinelEvent) => { if (event.type === "MASTER_CHANGE") { shouldRestart = true; } }; - - // Listen for master changes this.on("topology-change", handleTopologyChange); try { do { - // Check if we need to restart due to master change if (shouldRestart) { cursor = "0"; shouldRestart = false; } - const reply = await masterClient.scan(cursor, options); - // If a topology change happened during the scan command (which caused a retry), - // the reply is from the new master using the old cursor. We should discard it - // and let the loop restart the scan from cursor "0". + // Acquire the master lease only for the scan command itself, then + // release it before yielding so the consumer can run other commands + // (e.g. mGet) inside the for-await loop without exhausting the pool. + const masterClient = await this.acquire(); + let reply; + try { + reply = await masterClient.scan(cursor, options); + } finally { + const release = masterClient.release(); + if (release) await release; + } + + // If a topology change occurred during the scan command, the reply may + // mix keys from old and new masters. Discard and restart from cursor 0. if (shouldRestart) { continue; } + cursor = reply.cursor; yield reply.keys; } while (cursor !== "0"); } finally { - // Clean up: remove event listener and release the client this.removeListener("topology-change", handleTopologyChange); - masterClient.release(); } } } From ec1e2c31d5f63964d555638cd664e04c22e4a731 Mon Sep 17 00:00:00 2001 From: Nikolay Karadzhov Date: Wed, 3 Jun 2026 12:48:08 +0300 Subject: [PATCH 04/10] fix(sentinel): keep scanIterator running when MASTER_CHANGE fires at cursor "0" MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit When a topology change fired during the very first scan call (or right after a restart, where `cursor` is still `"0"`), the `continue` skipped the `cursor = reply.cursor` assignment, and the `do { ... } while (cursor !== "0")` condition then saw the unchanged `"0"` and exited — silently terminating the iterator without yielding any keys. Include `shouldRestart` in the loop condition so the loop survives the continue and runs another scan from the top of the body, where the cursor is reset to `"0"` and `shouldRestart` is cleared. Test: a regression case that emits a synthetic `topology-change` MASTER_CHANGE while the first `scan("0", ...)` call is in flight and asserts the iterator still yields the matching keys. Reported by Cursor Bugbot on PR #3141. Co-Authored-By: Claude Opus 4.7 (1M context) --- packages/client/lib/sentinel/index.spec.ts | 27 ++++++++++++++++++++++ packages/client/lib/sentinel/index.ts | 2 +- 2 files changed, 28 insertions(+), 1 deletion(-) diff --git a/packages/client/lib/sentinel/index.spec.ts b/packages/client/lib/sentinel/index.spec.ts index d31d689b259..c47b714a8e6 100644 --- a/packages/client/lib/sentinel/index.spec.ts +++ b/packages/client/lib/sentinel/index.spec.ts @@ -1447,6 +1447,33 @@ describe('legacy tests', () => { assert.equal(collected.get(`deadlock:${i}`), `value${i}`); } }, GLOBAL.SENTINEL.OPEN); + + testUtils.testWithClientSentinel('should restart and not silently terminate when MASTER_CHANGE fires while cursor is "0"', async sentinel => { + // Regression: a topology change during the first scan (cursor "0") used + // to `continue` past the cursor assignment, then the do/while condition + // saw the still-"0" cursor and exited, yielding nothing. + await sentinel.mSet([ + 'restart-on-zero:1', '1', + 'restart-on-zero:2', '2' + ]); + + const iter = sentinel.scanIterator({ MATCH: 'restart-on-zero:*' }); + // Calling .next() runs the generator body up to the first `await`, which + // means the `topology-change` listener is attached before this point. + const firstPagePromise = iter.next(); + sentinel.emit('topology-change', { + type: 'MASTER_CHANGE', + node: { host: 'synthetic', port: 0 } + }); + const firstPage = await firstPagePromise; + + assert.equal(firstPage.done, false, 'iterator must not terminate from a MASTER_CHANGE at cursor "0"'); + const collected = new Set(firstPage.value as Array); + for await (const keys of iter) { + for (const key of keys) collected.add(key); + } + assert.deepEqual(collected, new Set(['restart-on-zero:1', 'restart-on-zero:2'])); + }, GLOBAL.SENTINEL.OPEN); }); describe('scanIterator with master failover', () => { diff --git a/packages/client/lib/sentinel/index.ts b/packages/client/lib/sentinel/index.ts index 4fa08bdf328..1779bc2d2a7 100644 --- a/packages/client/lib/sentinel/index.ts +++ b/packages/client/lib/sentinel/index.ts @@ -701,7 +701,7 @@ export default class RedisSentinel< cursor = reply.cursor; yield reply.keys; - } while (cursor !== "0"); + } while (cursor !== "0" || shouldRestart); } finally { this.removeListener("topology-change", handleTopologyChange); } From 51dd022a4bf5265864e622512ab4450aa17e26df Mon Sep 17 00:00:00 2001 From: Nikolay Karadzhov Date: Wed, 3 Jun 2026 16:01:08 +0300 Subject: [PATCH 05/10] fix(sentinel): route scanIterator through _execute and fix cursor comparison - Use _execute so reserveClient:true reuses the reserved lease instead of blocking forever on an empty master pool (masterPoolSize:1). - Compare cursor.toString() to '0' so iteration terminates when a Blob String type mapping returns the cursor as Buffer. - Wrap the SCAN call so a connection-loss rejection during master failover restarts on the new master (matching the documented restart semantics) instead of propagating the transient error. - Export ScanIteratorOptions from packages/client/lib/client and drop the duplicate local interface in sentinel/index.ts. - Document that a user-supplied cursor option is honored only on the first SCAN call; failover restarts always reset to cursor 0. - Tests: reserveClient:true + masterPoolSize:1 completes, Blob String Buffer cursor terminates, user-supplied cursor is forwarded to SCAN, early break detaches the topology-change listener, synthetic regression asserts the first yield is the real scan reply, rename the "scan start" failover test to clarify it covers post-failover iteration. --- docs/sentinel.md | 2 + packages/client/lib/client/index.ts | 2 +- packages/client/lib/sentinel/index.spec.ts | 123 ++++++++++++++++++++- packages/client/lib/sentinel/index.ts | 43 +++---- 4 files changed, 148 insertions(+), 22 deletions(-) diff --git a/docs/sentinel.md b/docs/sentinel.md index e150ffe1de0..1657f0e6d3b 100644 --- a/docs/sentinel.md +++ b/docs/sentinel.md @@ -177,6 +177,8 @@ The standalone `RedisClient.scanIterator()` inherits SCAN's documented guarantee The sentinel iterator adds one extra source of duplicates: **master failover**. If the master changes mid-iteration (detected via the `topology-change` event with `type: "MASTER_CHANGE"`), the cursor is invalidated (SCAN cursors are node-local) and the iterator restarts from cursor `0` on the new master. Keys already yielded before the failover may be yielded again from the new master. +A user-supplied `cursor` option (e.g. `sentinel.scanIterator({ cursor: '1234' })`) is honored on the first SCAN call only. If a failover then resets the iterator, it restarts from cursor `0` (not the user-supplied value), because the original cursor was bound to the old master and would be meaningless on the new one. + Because Redis replication is asynchronous, the new master may also have a slightly different keyset than the old master at the moment of promotion — writes that had not yet replicated will be missing, and writes accepted on the new master after promotion will be present. If your processing must be exactly-once, deduplicate with a `Set`: diff --git a/packages/client/lib/client/index.ts b/packages/client/lib/client/index.ts index 94ccff31735..bac6efa0770 100644 --- a/packages/client/lib/client/index.ts +++ b/packages/client/lib/client/index.ts @@ -274,7 +274,7 @@ type ProxyClient = RedisClient { const firstPage = await firstPagePromise; assert.equal(firstPage.done, false, 'iterator must not terminate from a MASTER_CHANGE at cursor "0"'); + assert.ok(Array.isArray(firstPage.value), 'iterator must yield the real scan reply, not the synthetic event payload'); const collected = new Set(firstPage.value as Array); for await (const keys of iter) { for (const key of keys) collected.add(key); } assert.deepEqual(collected, new Set(['restart-on-zero:1', 'restart-on-zero:2'])); }, GLOBAL.SENTINEL.OPEN); + + testUtils.testWithClientSentinel( + 'should terminate when Blob String type mapping returns the cursor as a Buffer', + async sentinel => { + // Regression: comparing the (Buffer) cursor to the literal string '0' + // with !== never becomes false, so the iterator looped forever. + await sentinel.mSet([ + 'buffer-cursor:1', '1', + 'buffer-cursor:2', '2' + ]); + + const found = new Set(); + for await (const keyBatch of sentinel.scanIterator({ MATCH: 'buffer-cursor:*' })) { + for (const key of keyBatch) { + found.add(key.toString()); + } + } + + assert.deepEqual(found, new Set(['buffer-cursor:1', 'buffer-cursor:2'])); + }, + { + ...GLOBAL.SENTINEL.OPEN, + clientOptions: { + commandOptions: { + typeMapping: { [RESP_TYPES.BLOB_STRING]: Buffer } + } + } + } + ); + + testUtils.testWithClientSentinel('should forward a user-supplied cursor to SCAN on the first call', async sentinel => { + await sentinel.mSet([ + 'cursor-opt:1', '1', + 'cursor-opt:2', '2' + ]); + + const seenCursors: Array = []; + const sentinelAny = sentinel as unknown as { + _execute: ( + isReadonly: boolean | undefined, + fn: (client: unknown) => Promise + ) => Promise; + }; + const realExecute = sentinelAny._execute.bind(sentinel); + sentinelAny._execute = function (isReadonly, fn) { + return realExecute(isReadonly, (client: unknown) => { + const wrapped = new Proxy(client as object, { + get(target, prop, receiver) { + if (prop === 'scan') { + return (cursor: RedisArgument, opts: unknown) => { + seenCursors.push(cursor.toString()); + return (target as { scan: (c: RedisArgument, o: unknown) => unknown }).scan(cursor, opts); + }; + } + return Reflect.get(target, prop, receiver); + } + }); + return fn(wrapped); + }); + }; + + try { + const iter = sentinel.scanIterator({ MATCH: 'cursor-opt:*', cursor: '7' }); + await iter.next(); + await iter.return?.(undefined); + } finally { + sentinelAny._execute = realExecute; + } + + assert.equal(seenCursors[0], '7', 'first SCAN call must use the user-supplied cursor'); + }, GLOBAL.SENTINEL.OPEN); + + testUtils.testWithClientSentinel('should remove the topology-change listener when the consumer breaks early', async sentinel => { + await sentinel.mSet([ + 'cleanup:1', '1', + 'cleanup:2', '2' + ]); + + const before = sentinel.listenerCount('topology-change'); + for await (const _ of sentinel.scanIterator({ MATCH: 'cleanup:*' })) { + break; + } + assert.equal( + sentinel.listenerCount('topology-change'), + before, + 'scanIterator must detach its topology-change listener on early break' + ); + }, GLOBAL.SENTINEL.OPEN); + + testUtils.testWithClientSentinel( + 'should iterate without hanging when reserveClient:true and masterPoolSize:1', + async sentinel => { + // Regression: connect() consumes the only master lease into the + // reserved client, so an acquire() inside scanIterator would block + // forever on an empty pool. Routing through _execute reuses the + // reserved lease instead. + await sentinel.mSet([ + 'reserve-scan:1', '1', + 'reserve-scan:2', '2' + ]); + + const found = new Set(); + for await (const keyBatch of sentinel.scanIterator({ MATCH: 'reserve-scan:*' })) { + for (const key of keyBatch) found.add(key); + } + + assert.deepEqual(found, new Set(['reserve-scan:1', 'reserve-scan:2'])); + }, + { + ...GLOBAL.SENTINEL.OPEN, + reserveClient: true, + masterPoolSize: 1 + } + ); }); describe('scanIterator with master failover', () => { @@ -1599,7 +1714,11 @@ describe('legacy tests', () => { ); }); - it('should handle master change at scan start', async function () { + // This test waits for MASTER_CHANGE to fire before constructing the + // iterator, so the iterator never sees a failover in flight — it only + // verifies that a scan started after a completed failover talks to the + // new master. The in-flight-failover case is covered by the test above. + it('should iterate against the new master when failover completes before scan starts', async function () { this.timeout(60000); sentinel = frame.getSentinelClient({ scanInterval: 1000 }); diff --git a/packages/client/lib/sentinel/index.ts b/packages/client/lib/sentinel/index.ts index 1779bc2d2a7..3774ba86be4 100644 --- a/packages/client/lib/sentinel/index.ts +++ b/packages/client/lib/sentinel/index.ts @@ -1,6 +1,6 @@ import { EventEmitter } from 'node:events'; import { CommandArguments, RedisArgument, RedisFunctions, RedisModules, RedisScripts, ReplyUnion, RespVersions, TypeMapping, DEFAULT_RESP } from '../RESP/types'; -import RedisClient, { AnyRedisClientOptions, RedisClientOptions, RedisClientType } from '../client'; +import RedisClient, { AnyRedisClientOptions, RedisClientOptions, RedisClientType, ScanIteratorOptions } from '../client'; import { CommandOptions } from '../client/commands-queue'; import { attachConfig } from '../commander'; import { NON_STICKY_COMMANDS } from '../commands'; @@ -25,10 +25,6 @@ interface ClientInfo { id: number; } -interface ScanIteratorOptions { - cursor?: RedisArgument; -} - export class RedisSentinelClient< M extends RedisModules, F extends RedisFunctions, @@ -664,33 +660,40 @@ export default class RedisSentinel< this: RedisSentinelType, options?: ScanOptions & ScanIteratorOptions ) { - let cursor = options?.cursor ?? "0"; + let cursor: RedisArgument = options?.cursor ?? '0'; let shouldRestart = false; const handleTopologyChange = (event: RedisSentinelEvent) => { - if (event.type === "MASTER_CHANGE") { + if (event.type === 'MASTER_CHANGE') { shouldRestart = true; } }; - this.on("topology-change", handleTopologyChange); + this.on('topology-change', handleTopologyChange); try { do { if (shouldRestart) { - cursor = "0"; + cursor = '0'; shouldRestart = false; } - // Acquire the master lease only for the scan command itself, then - // release it before yielding so the consumer can run other commands - // (e.g. mGet) inside the for-await loop without exhausting the pool. - const masterClient = await this.acquire(); + // Route through _execute so reserveClient:true reuses the reserved + // lease (instead of waiting forever on an empty master pool), and the + // lease is released before yielding — consumers can issue other + // commands inside the for-await loop without exhausting the pool. let reply; try { - reply = await masterClient.scan(cursor, options); - } finally { - const release = masterClient.release(); - if (release) await release; + reply = await this._execute( + false, + client => (client as RedisClientType).scan(cursor, options) + ); + } catch (err) { + // A master failover mid-SCAN can surface as a connection error + // before the topology-change event is processed. If the iterator + // has already been flagged to restart, swallow and retry on the + // new master instead of propagating the transient error. + if (shouldRestart) continue; + throw err; } // If a topology change occurred during the scan command, the reply may @@ -701,9 +704,11 @@ export default class RedisSentinel< cursor = reply.cursor; yield reply.keys; - } while (cursor !== "0" || shouldRestart); + // Cursor may be a Buffer when a Blob String type mapping is in use; + // compare by string value so iteration actually terminates. + } while (cursor.toString() !== '0' || shouldRestart); } finally { - this.removeListener("topology-change", handleTopologyChange); + this.removeListener('topology-change', handleTopologyChange); } } } From bf8d9a61fb691aaa496706d3bb460e91d9c7d0f3 Mon Sep 17 00:00:00 2001 From: Nikolay Karadzhov Date: Thu, 4 Jun 2026 14:28:47 +0300 Subject: [PATCH 06/10] fix(sentinel): throw SentinelMasterChangeError on master failover during scanIterator Previously the iterator tried to survive failover by restarting from cursor 0 on the new master. That silently produced both duplicates (keys already yielded from the old master) and lost keys (writes that had not replicated before the failover). Throw instead and let the caller decide whether to retry, accept the partial result, or fail the surrounding operation. Co-Authored-By: Claude Opus 4.7 (1M context) --- docs/sentinel.md | 32 +++++------ packages/client/lib/errors.ts | 6 +++ packages/client/lib/sentinel/index.spec.ts | 62 ++++++---------------- packages/client/lib/sentinel/index.ts | 26 ++++----- 4 files changed, 46 insertions(+), 80 deletions(-) diff --git a/docs/sentinel.md b/docs/sentinel.md index 1657f0e6d3b..96c85cca192 100644 --- a/docs/sentinel.md +++ b/docs/sentinel.md @@ -171,31 +171,31 @@ for await (const keys of sentinel.scanIterator()) { } ``` -### Semantics and differences from standalone `scanIterator` +### Behaviour on master failover -The standalone `RedisClient.scanIterator()` inherits SCAN's documented guarantees: a full iteration returns every key present from start to end, and may return a key multiple times. See the [SCAN guarantees](https://redis.io/docs/latest/commands/scan/#scan-guarantees) page. +SCAN cursors are node-local — a cursor returned by one Redis instance is meaningless on any other instance. Because of this, the sentinel iterator cannot transparently survive a master failover: the in-flight cursor cannot be resumed on the promoted replica, and silently restarting from cursor `0` on the new master would hide both duplicate keys (already yielded from the old master) and data loss (writes that had not yet replicated before the failover). -The sentinel iterator adds one extra source of duplicates: **master failover**. If the master changes mid-iteration (detected via the `topology-change` event with `type: "MASTER_CHANGE"`), the cursor is invalidated (SCAN cursors are node-local) and the iterator restarts from cursor `0` on the new master. Keys already yielded before the failover may be yielded again from the new master. - -A user-supplied `cursor` option (e.g. `sentinel.scanIterator({ cursor: '1234' })`) is honored on the first SCAN call only. If a failover then resets the iterator, it restarts from cursor `0` (not the user-supplied value), because the original cursor was bound to the old master and would be meaningless on the new one. - -Because Redis replication is asynchronous, the new master may also have a slightly different keyset than the old master at the moment of promotion — writes that had not yet replicated will be missing, and writes accepted on the new master after promotion will be present. - -If your processing must be exactly-once, deduplicate with a `Set`: +Instead, if the master changes while an iteration is in progress, the iterator throws `SentinelMasterChangeError`. The caller decides whether to retry the iteration from scratch, accept the partial result, or fail the surrounding operation. ```javascript -const processed = new Set(); -for await (const keys of sentinel.scanIterator()) { - for (const key of keys) { - if (processed.has(key)) continue; - processed.add(key); +import { SentinelMasterChangeError } from '@redis/client'; - // process key +try { + for await (const keys of sentinel.scanIterator()) { + // ... + } +} catch (err) { + if (err instanceof SentinelMasterChangeError) { + // master failed over mid-iteration; restart from the beginning if desired + } else { + throw err; } } ``` -For very large keyspaces a `Set` may be memory-prohibitive. A Bloom filter is a lower-memory alternative but is **not** suitable for strict exactly-once processing: its false positives will cause some real keys to be skipped. Use it only when occasional skips are acceptable. +The iterator listens for the `topology-change` event with `type: "MASTER_CHANGE"`. The listener is attached when the generator body first runs (on the first `.next()` call) and is detached in a `finally` block, so an early `break` out of the `for await` loop will not leak listeners. + +The standalone `RedisClient.scanIterator()` still inherits SCAN's documented guarantees, including the possibility of returning the same key multiple times within a single iteration; see the [SCAN guarantees](https://redis.io/docs/latest/commands/scan/#scan-guarantees) page. ### Pool behaviour diff --git a/packages/client/lib/errors.ts b/packages/client/lib/errors.ts index 4fff9a9cb8d..faa6539e292 100644 --- a/packages/client/lib/errors.ts +++ b/packages/client/lib/errors.ts @@ -52,6 +52,12 @@ export class RootNodesUnavailableError extends Error { } } +export class SentinelMasterChangeError extends Error { + constructor(cause?: unknown) { + super('Sentinel master changed during operation', cause === undefined ? undefined : { cause }); + } +} + export class ReconnectStrategyError extends Error { originalError: Error; socketError: unknown; diff --git a/packages/client/lib/sentinel/index.spec.ts b/packages/client/lib/sentinel/index.spec.ts index c41fef16e6a..d74b999a99d 100644 --- a/packages/client/lib/sentinel/index.spec.ts +++ b/packages/client/lib/sentinel/index.spec.ts @@ -2,7 +2,7 @@ import { strict as assert } from 'node:assert'; import { setTimeout } from 'node:timers/promises'; import testUtils, { GLOBAL, MATH_FUNCTION } from '../test-utils'; import { RESP_TYPES } from '../RESP/decoder'; -import { WatchError } from "../errors"; +import { SentinelMasterChangeError, WatchError } from "../errors"; import { RedisSentinelConfig, SentinelFramework } from "./test-util"; import { RedisSentinelEvent, RedisSentinelType, RedisSentinelClientType, RedisNode } from "./types"; import RedisSentinel from "./index"; @@ -1448,16 +1448,13 @@ describe('legacy tests', () => { } }, GLOBAL.SENTINEL.OPEN); - testUtils.testWithClientSentinel('should restart and not silently terminate when MASTER_CHANGE fires while cursor is "0"', async sentinel => { - // Regression: a topology change during the first scan (cursor "0") used - // to `continue` past the cursor assignment, then the do/while condition - // saw the still-"0" cursor and exited, yielding nothing. + testUtils.testWithClientSentinel('should throw SentinelMasterChangeError when MASTER_CHANGE fires while cursor is "0"', async sentinel => { await sentinel.mSet([ - 'restart-on-zero:1', '1', - 'restart-on-zero:2', '2' + 'master-change-on-zero:1', '1', + 'master-change-on-zero:2', '2' ]); - const iter = sentinel.scanIterator({ MATCH: 'restart-on-zero:*' }); + const iter = sentinel.scanIterator({ MATCH: 'master-change-on-zero:*' }); // Calling .next() runs the generator body up to the first `await`, which // means the `topology-change` listener is attached before this point. const firstPagePromise = iter.next(); @@ -1465,15 +1462,8 @@ describe('legacy tests', () => { type: 'MASTER_CHANGE', node: { host: 'synthetic', port: 0 } }); - const firstPage = await firstPagePromise; - assert.equal(firstPage.done, false, 'iterator must not terminate from a MASTER_CHANGE at cursor "0"'); - assert.ok(Array.isArray(firstPage.value), 'iterator must yield the real scan reply, not the synthetic event payload'); - const collected = new Set(firstPage.value as Array); - for await (const keys of iter) { - for (const key of keys) collected.add(key); - } - assert.deepEqual(collected, new Set(['restart-on-zero:1', 'restart-on-zero:2'])); + await assert.rejects(firstPagePromise, SentinelMasterChangeError); }, GLOBAL.SENTINEL.OPEN); testUtils.testWithClientSentinel( @@ -1613,7 +1603,7 @@ describe('legacy tests', () => { await frame.cleanup(); }); - it('should restart scan from beginning when master changes during iteration', async function () { + it('should throw SentinelMasterChangeError when master changes during iteration', async function () { this.timeout(60000); sentinel = frame.getSentinelClient({ scanInterval: 1000 }); @@ -1621,25 +1611,19 @@ describe('legacy tests', () => { sentinel.on("error", () => {}); await sentinel.connect(); - // Set up test data - const testKeys = new Set(); const entries: Array = []; - for (let i = 0; i < 100; i++) { - const key = `failovertest:${i}`; - testKeys.add(key); - entries.push(key, `value${i}`); + entries.push(`failovertest:${i}`, `value${i}`); } await sentinel.mSet(entries); - // Wait for addded keys to be replicated + // Wait for added keys to be replicated await setTimeout(2000); let masterChangeDetected = false; let iterationCount = 0; const foundKeys = new Set(); - // Listen for manifest change events sentinel.on("topology-change", (event: RedisSentinelEvent) => { if (event.type === "MASTER_CHANGE") { masterChangeDetected = true; @@ -1647,17 +1631,15 @@ describe('legacy tests', () => { } }); - // Get the current master node before starting scan const originalMaster = sentinel.getMasterNode(); tracer.push(`Original master port: ${originalMaster?.port}`); - // Start scanning with a small COUNT to ensure multiple iterations const scanIterator = sentinel.scanIterator({ MATCH: "failovertest:*", COUNT: 10, }); - // Consume the scan iterator + let caught: unknown; try { for await (const keyBatch of scanIterator) { iterationCount++; @@ -1677,30 +1659,20 @@ describe('legacy tests', () => { } } } catch (error) { + caught = error; tracer.push(`Error during scan: ${error}`); - throw error; } - // Verify that master change was detected + assert.ok( + caught instanceof SentinelMasterChangeError, + `expected SentinelMasterChangeError, got ${caught}` + ); assert.equal( masterChangeDetected, true, "Master change should have been detected" ); - // Verify that we eventually got all keys despite the master change - assert.equal( - foundKeys.size, - testKeys.size, - "Should find all keys despite master failover" - ); - assert.deepEqual( - foundKeys, - testKeys, - "Found keys should match test keys" - ); - - // Verify that the master actually changed const newMaster = sentinel.getMasterNode(); tracer.push(`New master port: ${newMaster?.port}`); assert.notEqual( @@ -1708,10 +1680,6 @@ describe('legacy tests', () => { newMaster?.port, "Master should have changed" ); - - tracer.push( - `Test completed successfully with ${iterationCount} scan iterations` - ); }); // This test waits for MASTER_CHANGE to fire before constructing the diff --git a/packages/client/lib/sentinel/index.ts b/packages/client/lib/sentinel/index.ts index 3774ba86be4..7f0ef9fedf2 100644 --- a/packages/client/lib/sentinel/index.ts +++ b/packages/client/lib/sentinel/index.ts @@ -1,5 +1,6 @@ import { EventEmitter } from 'node:events'; import { CommandArguments, RedisArgument, RedisFunctions, RedisModules, RedisScripts, ReplyUnion, RespVersions, TypeMapping, DEFAULT_RESP } from '../RESP/types'; +import { SentinelMasterChangeError } from '../errors'; import RedisClient, { AnyRedisClientOptions, RedisClientOptions, RedisClientType, ScanIteratorOptions } from '../client'; import { CommandOptions } from '../client/commands-queue'; import { attachConfig } from '../commander'; @@ -661,21 +662,18 @@ export default class RedisSentinel< options?: ScanOptions & ScanIteratorOptions ) { let cursor: RedisArgument = options?.cursor ?? '0'; - let shouldRestart = false; + let masterChanged = false; const handleTopologyChange = (event: RedisSentinelEvent) => { if (event.type === 'MASTER_CHANGE') { - shouldRestart = true; + masterChanged = true; } }; this.on('topology-change', handleTopologyChange); try { do { - if (shouldRestart) { - cursor = '0'; - shouldRestart = false; - } + if (masterChanged) throw new SentinelMasterChangeError(); // Route through _execute so reserveClient:true reuses the reserved // lease (instead of waiting forever on an empty master pool), and the @@ -689,24 +687,18 @@ export default class RedisSentinel< ); } catch (err) { // A master failover mid-SCAN can surface as a connection error - // before the topology-change event is processed. If the iterator - // has already been flagged to restart, swallow and retry on the - // new master instead of propagating the transient error. - if (shouldRestart) continue; + // before the topology-change event is processed. Surface it as a + // SentinelMasterChangeError so callers can distinguish failover + // from other transient failures. + if (masterChanged) throw new SentinelMasterChangeError(err); throw err; } - // If a topology change occurred during the scan command, the reply may - // mix keys from old and new masters. Discard and restart from cursor 0. - if (shouldRestart) { - continue; - } - cursor = reply.cursor; yield reply.keys; // Cursor may be a Buffer when a Blob String type mapping is in use; // compare by string value so iteration actually terminates. - } while (cursor.toString() !== '0' || shouldRestart); + } while (cursor.toString() !== '0'); } finally { this.removeListener('topology-change', handleTopologyChange); } From 513b09eb3d033898802e76f71b0a1d75ad9fd923 Mon Sep 17 00:00:00 2001 From: Nikolay Karadzhov Date: Thu, 4 Jun 2026 15:35:34 +0300 Subject: [PATCH 07/10] refactor(sentinel): rename SentinelMasterChangeError to ScanIteratorInterruptedError and narrow wrap Reviewer flagged that the previous catch wrapped errors as SentinelMasterChangeError only when the topology-change event had already fired, leaving a race where a connection error from a dead master could surface raw before the event was processed. Considered widening the wrap to all connection-class errors, but a dropped socket is not by itself evidence of a failover (could be a transient network blip on the same master, in which case the cursor is still valid). Forcing a restart in that case is wasted work and a semantic lie. Narrow the wrap back to MASTER_CHANGE-observed only and rename the error to ScanIteratorInterruptedError to reflect what the iterator can actually detect. Connection-level errors propagate as-is; callers that want to treat both as failover-shaped can catch both. --- docs/sentinel.md | 10 +++++++--- packages/client/lib/errors.ts | 4 ++-- packages/client/lib/sentinel/index.spec.ts | 18 ++++++++++++------ packages/client/lib/sentinel/index.ts | 14 +++++++------- 4 files changed, 28 insertions(+), 18 deletions(-) diff --git a/docs/sentinel.md b/docs/sentinel.md index 96c85cca192..372070020d9 100644 --- a/docs/sentinel.md +++ b/docs/sentinel.md @@ -175,17 +175,21 @@ for await (const keys of sentinel.scanIterator()) { SCAN cursors are node-local — a cursor returned by one Redis instance is meaningless on any other instance. Because of this, the sentinel iterator cannot transparently survive a master failover: the in-flight cursor cannot be resumed on the promoted replica, and silently restarting from cursor `0` on the new master would hide both duplicate keys (already yielded from the old master) and data loss (writes that had not yet replicated before the failover). -Instead, if the master changes while an iteration is in progress, the iterator throws `SentinelMasterChangeError`. The caller decides whether to retry the iteration from scratch, accept the partial result, or fail the surrounding operation. +Instead, if the iterator observes a `MASTER_CHANGE` topology event while an iteration is in progress, it throws `ScanIteratorInterruptedError`. The caller decides whether to retry the iteration from scratch, accept the partial result, or fail the surrounding operation. + +Connection-level errors raised by the underlying client (e.g. `SocketClosedUnexpectedlyError`, `SocketTimeoutError`, `ReconnectStrategyError`) are **not** wrapped. A dropped socket is not by itself evidence of a failover — it may also be a transient network blip on the same master, in which case the cursor is still valid and a higher-level retry policy is appropriate. The original error is propagated as-is, and the caller can distinguish failover from a blip by checking for `ScanIteratorInterruptedError` versus other error types. + +In a real failover the dropped socket often precedes the Sentinel `MASTER_CHANGE` event (gated by `down-after-milliseconds`), so callers that want to treat both signals uniformly should catch both `ScanIteratorInterruptedError` **and** connection-class errors: ```javascript -import { SentinelMasterChangeError } from '@redis/client'; +import { ScanIteratorInterruptedError } from '@redis/client'; try { for await (const keys of sentinel.scanIterator()) { // ... } } catch (err) { - if (err instanceof SentinelMasterChangeError) { + if (err instanceof ScanIteratorInterruptedError) { // master failed over mid-iteration; restart from the beginning if desired } else { throw err; diff --git a/packages/client/lib/errors.ts b/packages/client/lib/errors.ts index faa6539e292..472a8ae5cc8 100644 --- a/packages/client/lib/errors.ts +++ b/packages/client/lib/errors.ts @@ -52,9 +52,9 @@ export class RootNodesUnavailableError extends Error { } } -export class SentinelMasterChangeError extends Error { +export class ScanIteratorInterruptedError extends Error { constructor(cause?: unknown) { - super('Sentinel master changed during operation', cause === undefined ? undefined : { cause }); + super('Scan iteration was interrupted by a Sentinel master change', cause === undefined ? undefined : { cause }); } } diff --git a/packages/client/lib/sentinel/index.spec.ts b/packages/client/lib/sentinel/index.spec.ts index d74b999a99d..4396d1bedb3 100644 --- a/packages/client/lib/sentinel/index.spec.ts +++ b/packages/client/lib/sentinel/index.spec.ts @@ -2,7 +2,7 @@ import { strict as assert } from 'node:assert'; import { setTimeout } from 'node:timers/promises'; import testUtils, { GLOBAL, MATH_FUNCTION } from '../test-utils'; import { RESP_TYPES } from '../RESP/decoder'; -import { SentinelMasterChangeError, WatchError } from "../errors"; +import { ScanIteratorInterruptedError, WatchError } from "../errors"; import { RedisSentinelConfig, SentinelFramework } from "./test-util"; import { RedisSentinelEvent, RedisSentinelType, RedisSentinelClientType, RedisNode } from "./types"; import RedisSentinel from "./index"; @@ -1448,7 +1448,7 @@ describe('legacy tests', () => { } }, GLOBAL.SENTINEL.OPEN); - testUtils.testWithClientSentinel('should throw SentinelMasterChangeError when MASTER_CHANGE fires while cursor is "0"', async sentinel => { + testUtils.testWithClientSentinel('should throw ScanIteratorInterruptedError when MASTER_CHANGE fires while cursor is "0"', async sentinel => { await sentinel.mSet([ 'master-change-on-zero:1', '1', 'master-change-on-zero:2', '2' @@ -1463,7 +1463,7 @@ describe('legacy tests', () => { node: { host: 'synthetic', port: 0 } }); - await assert.rejects(firstPagePromise, SentinelMasterChangeError); + await assert.rejects(firstPagePromise, ScanIteratorInterruptedError); }, GLOBAL.SENTINEL.OPEN); testUtils.testWithClientSentinel( @@ -1603,7 +1603,7 @@ describe('legacy tests', () => { await frame.cleanup(); }); - it('should throw SentinelMasterChangeError when master changes during iteration', async function () { + it('should throw when master changes during iteration', async function () { this.timeout(60000); sentinel = frame.getSentinelClient({ scanInterval: 1000 }); @@ -1663,9 +1663,15 @@ describe('legacy tests', () => { tracer.push(`Error during scan: ${error}`); } + // The error may be either ScanIteratorInterruptedError (if the + // topology-change event was observed before the next scan call) or a + // raw connection-class error (if the dropped socket raced ahead of the + // Sentinel down detection — usually gated by down-after-milliseconds). + // Both outcomes are valid; what matters is that the iteration does not + // silently continue across the failover. assert.ok( - caught instanceof SentinelMasterChangeError, - `expected SentinelMasterChangeError, got ${caught}` + caught instanceof Error, + `iteration must throw on master failover, got ${caught}` ); assert.equal( masterChangeDetected, diff --git a/packages/client/lib/sentinel/index.ts b/packages/client/lib/sentinel/index.ts index 7f0ef9fedf2..7c93b17b8cc 100644 --- a/packages/client/lib/sentinel/index.ts +++ b/packages/client/lib/sentinel/index.ts @@ -1,6 +1,6 @@ import { EventEmitter } from 'node:events'; import { CommandArguments, RedisArgument, RedisFunctions, RedisModules, RedisScripts, ReplyUnion, RespVersions, TypeMapping, DEFAULT_RESP } from '../RESP/types'; -import { SentinelMasterChangeError } from '../errors'; +import { ScanIteratorInterruptedError } from '../errors'; import RedisClient, { AnyRedisClientOptions, RedisClientOptions, RedisClientType, ScanIteratorOptions } from '../client'; import { CommandOptions } from '../client/commands-queue'; import { attachConfig } from '../commander'; @@ -673,7 +673,7 @@ export default class RedisSentinel< try { do { - if (masterChanged) throw new SentinelMasterChangeError(); + if (masterChanged) throw new ScanIteratorInterruptedError(); // Route through _execute so reserveClient:true reuses the reserved // lease (instead of waiting forever on an empty master pool), and the @@ -686,11 +686,11 @@ export default class RedisSentinel< client => (client as RedisClientType).scan(cursor, options) ); } catch (err) { - // A master failover mid-SCAN can surface as a connection error - // before the topology-change event is processed. Surface it as a - // SentinelMasterChangeError so callers can distinguish failover - // from other transient failures. - if (masterChanged) throw new SentinelMasterChangeError(err); + // Only wrap when MASTER_CHANGE has been observed; otherwise let the + // underlying error propagate. A bare socket disconnect alone is not + // sufficient evidence of a failover (could be a transient network + // blip on the same master, in which case the cursor is still valid). + if (masterChanged) throw new ScanIteratorInterruptedError(err); throw err; } From acb9ee1e27c559cc88eb60734b202de457ab0b67 Mon Sep 17 00:00:00 2001 From: Nikolay Karadzhov Date: Fri, 5 Jun 2026 10:32:49 +0300 Subject: [PATCH 08/10] docs(sentinel): tighten scanIterator failover wording and add JSDoc Clarify in docs/sentinel.md that the iterator throws only when continuing would require another SCAN call (cursor=0 on the failing call still completes cleanly). Add a JSDoc block on the sentinel scanIterator method covering the yielded shape, lease behaviour, and error contract. --- docs/sentinel.md | 4 +++- packages/client/lib/sentinel/index.ts | 18 ++++++++++++++++++ 2 files changed, 21 insertions(+), 1 deletion(-) diff --git a/docs/sentinel.md b/docs/sentinel.md index 372070020d9..1e1ab9aa34e 100644 --- a/docs/sentinel.md +++ b/docs/sentinel.md @@ -175,7 +175,9 @@ for await (const keys of sentinel.scanIterator()) { SCAN cursors are node-local — a cursor returned by one Redis instance is meaningless on any other instance. Because of this, the sentinel iterator cannot transparently survive a master failover: the in-flight cursor cannot be resumed on the promoted replica, and silently restarting from cursor `0` on the new master would hide both duplicate keys (already yielded from the old master) and data loss (writes that had not yet replicated before the failover). -Instead, if the iterator observes a `MASTER_CHANGE` topology event while an iteration is in progress, it throws `ScanIteratorInterruptedError`. The caller decides whether to retry the iteration from scratch, accept the partial result, or fail the surrounding operation. +If a `MASTER_CHANGE` topology event is observed while an iteration is in progress **and** the iterator still needs to issue another `SCAN` (i.e. the cursor has not yet returned to `0`), it throws `ScanIteratorInterruptedError` rather than send a stale, node-local cursor to a different master. The caller decides whether to retry the iteration from scratch, accept the partial result, or fail the surrounding operation. + +If the responding master returns `cursor=0` on the same call during which `MASTER_CHANGE` fires, no error is thrown — that node honored SCAN's contract ("every key present at iteration start was returned") and no further calls are needed. SCAN never claims to reflect "the current dataset" at the moment iteration ends, with or without a failover, so this case is not treated as an interruption. Connection-level errors raised by the underlying client (e.g. `SocketClosedUnexpectedlyError`, `SocketTimeoutError`, `ReconnectStrategyError`) are **not** wrapped. A dropped socket is not by itself evidence of a failover — it may also be a transient network blip on the same master, in which case the cursor is still valid and a higher-level retry policy is appropriate. The original error is propagated as-is, and the caller can distinguish failover from a blip by checking for `ScanIteratorInterruptedError` versus other error types. diff --git a/packages/client/lib/sentinel/index.ts b/packages/client/lib/sentinel/index.ts index 7c93b17b8cc..3f6b538b938 100644 --- a/packages/client/lib/sentinel/index.ts +++ b/packages/client/lib/sentinel/index.ts @@ -657,6 +657,24 @@ export default class RedisSentinel< this._self.#internal.setTracer(tracer); } + /** + * Async generator that iterates over keys on the Sentinel master by issuing + * paged `SCAN` calls. Yields one array of keys per page until the SCAN cursor + * returns to `0`. + * + * The master client lease is acquired for the duration of each `SCAN` call + * and released before yielding, so consumers can issue other commands from + * inside the `for await` loop body without deadlocking against the iterator + * — even with `masterPoolSize: 1`. + * + * Throws `ScanIteratorInterruptedError` on observed `MASTER_CHANGE`. + * Throws the underlying error on any other failure. + * + * @param options - SCAN options and an optional starting `cursor`. The + * starting cursor is honored only on the first call. + * @yields Arrays of keys returned by each `SCAN` page. Pages may be empty. + * @throws {ScanIteratorInterruptedError} On observed `MASTER_CHANGE`. + */ async *scanIterator( this: RedisSentinelType, options?: ScanOptions & ScanIteratorOptions From 3bffac4ed92408e5f043a08fd818006ad397619d Mon Sep 17 00:00:00 2001 From: Nikolay Karadzhov Date: Fri, 5 Jun 2026 10:49:04 +0300 Subject: [PATCH 09/10] test(sentinel): rewrite scanIterator synthetic failover test to observe error on the next page MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The previous test inserted two keys and expected the first page to reject when MASTER_CHANGE fired. Under the current contract — only throw when continuing would require another SCAN call — the first page returns cursor=0 and yields normally, so the test had to be reworked to fire the event between two pages. Insert enough keys with a small COUNT to force a non-terminal first page, consume it, emit the synthetic event, then assert the next page rejects with ScanIteratorInterruptedError. --- packages/client/lib/sentinel/index.spec.ts | 23 ++++++++++++---------- 1 file changed, 13 insertions(+), 10 deletions(-) diff --git a/packages/client/lib/sentinel/index.spec.ts b/packages/client/lib/sentinel/index.spec.ts index 4396d1bedb3..22fd9653f8d 100644 --- a/packages/client/lib/sentinel/index.spec.ts +++ b/packages/client/lib/sentinel/index.spec.ts @@ -1448,22 +1448,25 @@ describe('legacy tests', () => { } }, GLOBAL.SENTINEL.OPEN); - testUtils.testWithClientSentinel('should throw ScanIteratorInterruptedError when MASTER_CHANGE fires while cursor is "0"', async sentinel => { - await sentinel.mSet([ - 'master-change-on-zero:1', '1', - 'master-change-on-zero:2', '2' - ]); + testUtils.testWithClientSentinel('should throw ScanIteratorInterruptedError on the next page after MASTER_CHANGE is observed', async sentinel => { + const entries: Array = []; + for (let i = 0; i < 200; i++) { + entries.push(`master-change-mid:${i}`, String(i)); + } + await sentinel.mSet(entries); + + const iter = sentinel.scanIterator({ MATCH: 'master-change-mid:*', COUNT: 10 }); + // First page: cursor must come back non-zero so the iterator still has + // work to do after the synthetic event fires. + const firstPage = await iter.next(); + assert.equal(firstPage.done, false, 'first page must not terminate iteration'); - const iter = sentinel.scanIterator({ MATCH: 'master-change-on-zero:*' }); - // Calling .next() runs the generator body up to the first `await`, which - // means the `topology-change` listener is attached before this point. - const firstPagePromise = iter.next(); sentinel.emit('topology-change', { type: 'MASTER_CHANGE', node: { host: 'synthetic', port: 0 } }); - await assert.rejects(firstPagePromise, ScanIteratorInterruptedError); + await assert.rejects(iter.next(), ScanIteratorInterruptedError); }, GLOBAL.SENTINEL.OPEN); testUtils.testWithClientSentinel( From 0cae392d32fefdc29ca029c665c63dad297d8860 Mon Sep 17 00:00:00 2001 From: Nikolay Karadzhov Date: Fri, 5 Jun 2026 13:55:43 +0300 Subject: [PATCH 10/10] fix(sentinel): re-check masterChanged inside _execute lambda MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The pre-check at the top of the scanIterator loop only catches a MASTER_CHANGE observed before _execute is called. A failover can land while _execute is awaiting a master client lease (empty pool, common under masterPoolSize: 1 with a reserved client), in which case the lease resolves to a fresh client on the new master and SCAN would silently resume with a cursor from the old master. Re-check masterChanged inside the lambda passed to _execute so the race is closed once the lease resolves but before the SCAN call. Guard the outer catch with an instanceof check to avoid double-wrapping the synchronous throw from the lambda. Includes a regression test that emits MASTER_CHANGE synchronously inside _execute on the second page, after the pre-check has already passed. Verified red→green: with the in-lambda guard removed the test fails with "Missing expected rejection". Co-Authored-By: Claude Opus 4.7 (1M context) --- packages/client/lib/sentinel/index.spec.ts | 43 ++++++++++++++++++++++ packages/client/lib/sentinel/index.ts | 11 +++++- 2 files changed, 53 insertions(+), 1 deletion(-) diff --git a/packages/client/lib/sentinel/index.spec.ts b/packages/client/lib/sentinel/index.spec.ts index 22fd9653f8d..a54cbf5c076 100644 --- a/packages/client/lib/sentinel/index.spec.ts +++ b/packages/client/lib/sentinel/index.spec.ts @@ -1469,6 +1469,49 @@ describe('legacy tests', () => { await assert.rejects(iter.next(), ScanIteratorInterruptedError); }, GLOBAL.SENTINEL.OPEN); + testUtils.testWithClientSentinel('should throw ScanIteratorInterruptedError when MASTER_CHANGE fires during _execute (after pre-check, before scan)', async sentinel => { + const entries: Array = []; + for (let i = 0; i < 200; i++) { + entries.push(`master-change-race:${i}`, String(i)); + } + await sentinel.mSet(entries); + + let executeCalls = 0; + const sentinelAny = sentinel as unknown as { + _execute: ( + isReadonly: boolean | undefined, + fn: (client: unknown) => Promise + ) => Promise; + }; + const realExecute = sentinelAny._execute.bind(sentinel); + sentinelAny._execute = function (isReadonly, fn) { + executeCalls++; + // On the second SCAN call the iterator already ran its pre-check at the + // top of the loop (masterChanged was still false). Fire MASTER_CHANGE + // synchronously inside _execute so the listener flips masterChanged + // before the lambda passed to _execute is invoked. Without the + // in-lambda re-check, scan would proceed against the (now potentially + // stale) lease and the failure would be silent. + if (executeCalls === 2) { + sentinel.emit('topology-change', { + type: 'MASTER_CHANGE', + node: { host: 'synthetic', port: 0 } + }); + } + return realExecute(isReadonly, fn); + }; + + try { + const iter = sentinel.scanIterator({ MATCH: 'master-change-race:*', COUNT: 10 }); + const firstPage = await iter.next(); + assert.equal(firstPage.done, false, 'first page must not terminate iteration'); + + await assert.rejects(iter.next(), ScanIteratorInterruptedError); + } finally { + sentinelAny._execute = realExecute; + } + }, GLOBAL.SENTINEL.OPEN); + testUtils.testWithClientSentinel( 'should terminate when Blob String type mapping returns the cursor as a Buffer', async sentinel => { diff --git a/packages/client/lib/sentinel/index.ts b/packages/client/lib/sentinel/index.ts index 3f6b538b938..50acdd0d043 100644 --- a/packages/client/lib/sentinel/index.ts +++ b/packages/client/lib/sentinel/index.ts @@ -701,9 +701,18 @@ export default class RedisSentinel< try { reply = await this._execute( false, - client => (client as RedisClientType).scan(cursor, options) + client => { + // Re-check after the lease resolves: a failover may have landed + // while waiting on an empty master pool, in which case the lease + // now points to a fresh client on the new master and SCAN would + // resume with a cursor from the old master. + if (masterChanged) throw new ScanIteratorInterruptedError(); + return (client as RedisClientType).scan(cursor, options); + } ); } catch (err) { + // Pass through if already wrapped (from the in-lambda re-check). + if (err instanceof ScanIteratorInterruptedError) throw err; // Only wrap when MASTER_CHANGE has been observed; otherwise let the // underlying error propagate. A bare socket disconnect alone is not // sufficient evidence of a failover (could be a transient network