feat(client): multi-db client with automatic endpoint failover - #3435
nkaradzhov wants to merge 99 commits into
Conversation
49bd923 to
a95488d
Compare
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: f39f753780
ℹ️ About Codex in GitHub
Codex has been enabled to automatically review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
When you sign up for Codex through ChatGPT, Codex can also answer questions or update the PR, like "@codex address that feedback".
| */ | ||
| const INTERCEPTED = new Set<PropertyKey>([ | ||
| 'connect', 'close', 'destroy', 'quit', | ||
| 'withTypeMapping', 'withCommandOptions', 'withAbortSignal', |
There was a problem hiding this comment.
Keep
asap() views on the multi-db wrapper
When an application creates const view = client.asap(), asap is not intercepted like the other command-option views, so the generic forwarder returns the current member's raw proxy. After a failover, every command through view continues targeting the demoted member and bypasses the multi-db availability gate and failure detector; implement asap() through makeDerived just like withCommandOptions.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Same issue flagged by both bots, one fix. Fixed in 713539a: asap joins the intercepted derived views — resolves the active member per call, feeds the detector, rejects while all members are down, composes with the with* views. Test added.
| const member = mgr.activeDatabase; | ||
| // eslint-disable-next-line @typescript-eslint/no-explicit-any -- dynamic patching | ||
| const inner = (resolve(member.client) as any).multi(); |
There was a problem hiding this comment.
Pin WATCH and MULTI to the same member
If WATCH is sent to member A and a forced or automatic failover selects member B before multi() is created, this line builds the transaction on B. Since B has no watch state, EXEC succeeds without checking the watched key, silently defeating optimistic locking and potentially overwriting a concurrent update; the wrapper needs to remember the member used by WATCH and bind the subsequent transaction and UNWATCH to it.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Confirmed. Fixed in 527406a + d3e6651: a switch invalidates the outstanding watch session — the next EXEC rejects with WatchError (mirroring the base client's reconnect semantics and sentinel's dirty-watch on master change), so the standard retry loop re-runs the whole cycle on the new active member and nothing can commit unguarded. Documented + docker tests (switch invalidation + retry, conflict abort, unwatch release).
| const INTERCEPTED = new Set<PropertyKey>([ | ||
| 'connect', 'close', 'destroy', 'quit', | ||
| 'withTypeMapping', 'withCommandOptions', 'withAbortSignal', | ||
| 'multi', 'MULTI', 'duplicate' |
There was a problem hiding this comment.
Fan out
ref() and unref() across members
For a standalone multi-db client, ref and unref are omitted from the intercepted lifecycle surface, so the generic forwarder invokes them only on the active member. Calling client.unref() therefore leaves every passive member's socket referenced and can keep a short-lived process alive indefinitely, unlike the base-client contract; both methods need to apply to all live members.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Confirmed. Fixed in 713539a: ref()/unref() fan out over every member socket, skip the availability gate, return void; the intent is sticky so members added later come up matching it. Test added.
| >(options: { | ||
| databases: Array<DatabaseConfig<RedisClientOptions<M, F, S, RESP, TYPE_MAPPING>>>; | ||
| } & MultiDbConfig): MultiDbResult<RedisClientType<M, F, S, RESP, TYPE_MAPPING>> { |
There was a problem hiding this comment.
Preserve controller config types in Redis-package factories
This return type omits MultiDbResult's configuration generic, causing it to default to PoolDatabaseConfig<unknown>; the same omission appears in the pool, cluster, and sentinel wrappers below. Consequently, consumers importing from redis can compile invalid runtime additions such as { options: 42 }, cross-topology options, or poolOptions on a standalone controller, even though the corresponding @redis/client factories correctly preserve their member configuration types.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Confirmed. Fixed in 174b6ea: the four redis-package factories pass the CONFIG type parameter, so their controllers type-check runtime configs like the @redis/client ones.
| const result = genericCreateMultiDbClient({ | ||
| ...multiDbOptions, | ||
| databases: databases.map(db => ({ ...db, options: withStackModules(db.options) })) | ||
| }); | ||
| mergeModulesOnAdd(result.controller); | ||
| return result as unknown as MultiDbResult<RedisClientType<M, F, S, RESP, TYPE_MAPPING>>; |
There was a problem hiding this comment.
Wrap controllers returned by
duplicate() with Stack modules
mergeModulesOnAdd patches only the controller returned directly by this factory. A later client.duplicate() constructs a fresh generic controller without this wrapper, so a database added or replaced through the duplicate lacks the default json, ft, and ts modules; after that member becomes active, calls through those namespaces fail because the underlying namespace is undefined.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Same issue flagged by both bots, one fix. Fixed in 174b6ea: the redis meta-package wraps every multi-db pair recursively, so a duplicate()'s controller (at any depth) merges the Stack modules on addDatabase/replaceDatabase.
| if (!(healthCheck.timeout < healthCheck.interval)) { | ||
| throw new TypeError( | ||
| `MultiDb: healthCheck.timeout (${healthCheck.timeout}) must be less than healthCheck.interval (${healthCheck.interval})` | ||
| ); |
There was a problem hiding this comment.
Reject timer durations that Node clamps to one millisecond
The relational validation accepts Infinity and finite delays above Node's maximum timer range; for example, a healthCheck.timeout above 2^31−1 with a still-larger interval passes this check. Node clamps those timer delays to 1 ms, so probes can time out immediately while the scheduler runs in a hot loop, incorrectly opening circuits and hammering Redis. Require finite timer values within the supported delay range.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Confirmed. Fixed in 53c49b4: every duration option (setAutoFallback included) is bounded by Node's 2^31-1 timer max — Infinity and larger are rejected, so the 1ms-clamp hot loop cannot be configured. gracePeriod shares the bound.
| (this.client as unknown as EventEmitter) | ||
| .on('error', this.#onError) | ||
| .on('client-error', this.#onClientError) | ||
| .on('ready', this.#onReady) | ||
| .on('end', this.#onEnd); |
There was a problem hiding this comment.
Forward Sentinel
client-error from the internal client
This listener is attached to the public Sentinel member, but RedisSentinel emits node-level client-error events only on its internal controller and does not re-emit them on the public object; with the default passthroughClientErrorEvents: false, there is no corresponding public error either. Consequently, master, replica, and sentinel connection errors never reach the failure detector or the promised multi-db member-error event. Forward the internal event or subscribe through the Sentinel adapter.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Confirmed — the event existed only on the internal class. Fixed in 713539a + 52d3fc5: RedisSentinel re-emits client-error publicly, and the multi-db layer counts only MASTER connectivity as fault evidence (replica/sentinel-node/pub-sub-proxy errors surface as member-error only — a healthy deployment tolerates those by design). Tests added.
| const active = mgr.activeDatabase; | ||
| // eslint-disable-next-line @typescript-eslint/no-explicit-any -- dynamic | ||
| const result = (resolve(active.client) as any)[name](...args); |
There was a problem hiding this comment.
Carry SELECT state across member switches
SELECT is treated as an ordinary forwarded method, so it changes only the member active at that call. If a standalone client selects database 1 and then fails over, subsequent commands run on the new member's default database 0, silently reading and writing a different logical keyspace; the wrapper must either synchronize selected-database state across members, reapply it before a switch becomes visible, or reject runtime SELECT.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Confirmed, though narrower than P1 — only runtime SELECT on standalone members; the per-member database option always worked. Resolved in 713539a by rejecting runtime SELECT (commands and the multi builder) with guidance toward per-member database; replaying SELECT on switch would clash with the documented no-replay contract. Documented.
| ); | ||
| } | ||
|
|
||
| const target = this.#select(healthy); |
There was a problem hiding this comment.
Tear down members when initial strategy selection throws
After all member connections and health checks succeed, a custom failover strategy can throw here or return a copied candidate, which makes #select throw its explicit TypeError. Unlike the handled undefined result below, that exception escapes connect() without calling destroy(), leaving all established Redis connections and reconnect machinery alive even though the documented rejected-connect contract says every member is destroyed. Wrap initial selection in the same teardown path.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Confirmed. Fixed in 53c49b4: initial selection is guarded — a throwing strategy (or a foreign candidate) destroys every member before the rejection escapes, matching the adjacent failure branches. Test added.
| (this.client as unknown as EventEmitter) | ||
| .on('error', this.#onError) | ||
| .on('client-error', this.#onClientError) | ||
| .on('ready', this.#onReady) | ||
| .on('end', this.#onEnd); |
There was a problem hiding this comment.
Forward active-member invalidation events
When a standalone member is configured with emitInvalidate: true, Redis invalidation pushes are emitted on the hidden member client, but the manager listens only for error and lifecycle events here and the user-facing wrapper never receives invalidate. Applications using the base client's invalidation event for cache coherence therefore receive no notifications at all through the advertised drop-in client; route the active member's invalidations to the wrapper, with switch handling that avoids stale-member events.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
The mechanics are right, but the closed event surface is the documented design — and forwarding invalidate would produce FALSE coherence after a failover (the new member has no server-side tracking for keys cached via the old one). The real gap was the silent drop: 789c47c rejects emitInvalidate in member options at config resolution and documents per-member clientSideCache as the supported mode.
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 52d3fc5343
ℹ️ About Codex in GitHub
Codex has been enabled to automatically review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
When you sign up for Codex through ChatGPT, Codex can also answer questions or update the PR, like "@codex address that feedback".
| const INTERCEPTED = new Set<PropertyKey>([ | ||
| 'connect', 'close', 'destroy', 'quit', | ||
| 'withTypeMapping', 'withCommandOptions', 'withAbortSignal', 'asap', | ||
| 'multi', 'MULTI', 'duplicate' |
There was a problem hiding this comment.
When an application shuts down through client.QUIT() or the deprecated client.disconnect(), neither name is intercepted, so the generic forwarder tears down only the active member. That member's end can then trigger failover while the passive sockets and manager timers remain alive, meaning shutdown can leave the client operational and keep the process open; implement both aliases through the same all-member teardown path as quit() and destroy().
AGENTS.md reference: AGENTS.md:L75-L75
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Confirmed. Fixed in e822e3c: QUIT and disconnect are intercepted as fan-out teardowns (quit()/destroy() semantics) instead of forwarding to the active member only — shutdown no longer reads as a member failure and turns into a failover. Test added.
| export type MultiDbClientType<C extends AnyRedisClientType> = | ||
| Omit<C, 'duplicate'> & | ||
| MultiDbEventEmitter & | ||
| { |
There was a problem hiding this comment.
Keep multi-db typing on derived views
When callers use withCommandOptions(), withTypeMapping(), withAbortSignal(), or asap(), the inherited C signatures type the result as a bare member client even though makeDerived returns another multi-db wrapper. Thus await client.asap().duplicate().connect() still compiles but fails at runtime because duplicate() returns { client, controller }, which has no connect; fresh evidence beyond the earlier root duplicate() fix is that only duplicate is omitted and redeclared here, while every derived-view return type remains unchanged.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Confirmed, and fixed as part of a broader TypeScript-performance pass. In 0c15a0a MultiDbClientType is keyed by a literal member-kind name and the derived-view methods return the multi-db type with full type-mapping fidelity, the pair-shaped duplicate() and the strict event table — so a view no longer decays to the bare client type. An inference-based conditional design was measured at ~5.5M type instantiations (and collapsed cluster views to never) and rejected; a types-test pins the behavior and runs in CI. Related: 2c852aa removed the any-parameterized union constraint that was quadrupling consumer lib-check cost.
| if (method === 'exec' && mgr.watchDirty) { | ||
| mgr.clearWatchSession(); | ||
| return Promise.reject(new WatchError('MultiDb: the active database changed after WATCH')); |
There was a problem hiding this comment.
Clear WATCH on the demoted member
If member A remains connected after a WATCH and traffic switches to B, this dirty-EXEC path clears only manager metadata and rejects locally; it never sends UNWATCH or EXEC to A. When automatic fallback later selects A, its connection is still watching the old keys, so an otherwise unrelated transaction can unexpectedly abort based on changes made since the original WATCH. Fresh evidence after the earlier WATCH fix is that demoted members remain connected and eligible for fallback, unlike a reconnected base client whose old connection-scoped state disappeared.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Confirmed. Fixed in 30e7544: when a switch dirties the watch session it best-effort UNWATCHes the demoted watching member — clearing both its server-side watches and its client-side epoch — so an unrelated transaction after a later fallback to it can't abort spuriously. The one residual (a member that is down at switch time keeps its epoch until reconnect, yielding at most one fail-safe WatchError) is documented at the call site. Tests added.
| this.#events?.emit('failover', { from: from.id, to: target.id, reason }); | ||
| } | ||
|
|
||
| this.#afterSwitch(from, target).catch(err => this.#emitError(err as Error)); |
There was a problem hiding this comment.
Move subscription state before emitting the switch
Because the active member is changed and the synchronous failover event is emitted before #afterSwitch starts, an event listener that immediately calls client.unsubscribe('x') dispatches that operation to the new member while its subscription map is still empty. After the listener returns, the old member's listener is moved onto the new member, effectively resurrecting the subscription even though the unsubscribe can resolve successfully; seed the target's extractable subscription state before exposing the switch event.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Confirmed. Fixed in 73935b6: the pub/sub handover now runs before the failover/fallback event is emitted. Since all adapters seed the target's listener maps synchronously, a listener reacting to the event sees complete subscription state, so its unsubscribe is no longer resurrected by the move. Test added.
| export function probeRoundBudget(options: ProbeRoundOptions): number { | ||
| return options.numProbes * options.timeout | ||
| + (options.numProbes - 1) * options.delayBetweenProbes; |
There was a problem hiding this comment.
Bound the computed probe-round timeout
Valid individual settings can still make this computed delay exceed Node's timer range—for example, two probes with a timeout near 2^31-1—and withTimeout passes that sum directly to setTimeout, which clamps it to roughly 1 ms. Initial member connections then time out almost immediately and can make connect() destroy an otherwise healthy configuration; fresh evidence beyond the earlier per-option timer validation is this unbounded multiplication and addition, so the computed round budget also needs validation or safe bounding.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Confirmed. Fixed in 6117b06: the computed probe-round budget (numProbes * timeout + delays) is validated against the 2^31-1 timer maximum at config resolution, like the individual duration options — an overflowing product can no longer clamp the connect budget to 1ms. Spec row added.
| const results = await Promise.all( | ||
| this.#databases.map(db => this.#establishMember(db, false)) | ||
| ); | ||
| const healthy = this.#databases.filter((_, index) => results[index]); |
There was a problem hiding this comment.
Preserve member/result identity during connect
When addDatabase() or removeDatabase() mutates the member set while this Promise.all is awaiting establishment, results remains indexed to the original array but healthy is built by filtering the current array. Removing a middle member shifts later results onto the wrong databases, so initialAvailability: 'ALL' can either reject despite every remaining member being healthy or resolve while a remaining member failed its check; snapshot the member array and associate results by member identity, then reconcile removals before selection.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Confirmed. Fixed in 25ecb27: connect() now pairs the Promise.all results with a snapshot of the member set taken before the fan-out and reconciles it with the live set, so a concurrent addDatabase/removeDatabase can't shift outcomes onto the wrong members or fail the availability gate spuriously. Test added.
| (this.client as unknown as EventEmitter) | ||
| .on('error', this.#onError) | ||
| .on('client-error', this.#onClientError) |
There was a problem hiding this comment.
Do not count Sentinel passthrough errors twice
When a Sentinel member sets passthroughClientErrorEvents: true, each internal node failure is emitted publicly as both error and client-error, and these two listeners process both copies. The unconditional error path counts replica and Sentinel-node failures as data-path faults despite the client-error classifier rejecting them, while a master failure is counted twice; with a low detector threshold this can fail over a healthy Sentinel deployment after one irrelevant node error.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Confirmed. Fixed in d2bc74d: sentinel members no longer count the untyped error channel as fault evidence — with passthroughClientErrorEvents it duplicated node errors (and double-counted MASTER), and sentinel-internal housekeeping noise flowed through it too. Sentinel fault evidence is now the MASTER-typed client-error channel, end, command outcomes and health checks; every error still surfaces as member-error. Test added.
| #startFallbackTimer(intervalMs: number): void { | ||
| this.#stopFallbackTimer(); | ||
| if (!this.#schedulerRunning || intervalMs <= 0 || this.#teardown.signal.aborted) return; | ||
| this.#fallbackTimer = setInterval(() => this.#maybeFallback(), intervalMs); |
There was a problem hiding this comment.
Catch custom strategy failures in the fallback timer
When auto-fallback is enabled and a custom strategy throws or returns a copied/foreign candidate, #select throws directly from this interval callback. Unlike the initial selection in connect(), this background path has no catch boundary, so the exception can terminate the Node process instead of reaching the wrapper's guarded error event; catch strategy failures around the timer evaluation and report them through the background-error outlet.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Confirmed, and the scope was wider than the fallback path. Fixed in 7ccf76a: the auto-fallback timer, the all-down search loop and the failure handler all select through a guard that reports strategy failures on the error outlet and degrades to 'no candidate' — a throwing custom strategy no longer crashes the process or strands failover state. connect() and removeDatabase keep the raw throw (they have callers to reject to). Test added.
| const added = await this.addDatabase(config); | ||
| await this.removeDatabase(id); |
There was a problem hiding this comment.
Keep the old member when its replacement is unhealthy
When replacing a member with a different ID, addDatabase() resolves even if the new endpoint could not connect or pass its health check, leaving that new member OPEN; this code nevertheless removes the old healthy member immediately afterward. A routine endpoint rotation with a typo or temporarily unreachable replacement can therefore silently discard working redundancy and leave no usable failover target, despite the different-ID path promising that redundancy never drops; remove the old member only after the added member is verified selectable.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Partly fixed, partly by design. By design: runtime operations never throw because a member cannot connect — establishment tolerance is uniform across addDatabase and replaceDatabase, and member health is observable via circuitState and member-error rather than guaranteed. Fixed in 63d49c8: the same-id replace path now validates the new config BEFORE removing the old member (a malformed config can no longer shrink the set with no rollback), all user-supplied member configs funnel through one validator (which also closed a duplicate(overrides) validation bypass), and the docs now state what holds — the member COUNT never drops (a replacement may still be establishing; its circuit reflects its health).
| const listeners = from._getQueue().removeAllPubSubListeners(); | ||
| await Promise.all([ | ||
| to.extendPubSubListeners(PUBSUB_TYPE.CHANNELS, listeners[PUBSUB_TYPE.CHANNELS]), | ||
| to.extendPubSubListeners(PUBSUB_TYPE.PATTERNS, listeners[PUBSUB_TYPE.PATTERNS]), | ||
| to.extendPubSubListeners(PUBSUB_TYPE.SHARDED, listeners[PUBSUB_TYPE.SHARDED]) |
There was a problem hiding this comment.
Carry in-flight standalone and cluster subscriptions
If a standalone subscription command is still awaiting its server reply when a switch occurs, its listener has not yet been inserted into the PubSub maps, so this extraction returns an empty snapshot. The old command can then resolve and insert the listener back into the demoted member while the caller observes a successful subscription that receives nothing from the new active database; the cluster extractor has the same map-only behavior. Fresh evidence beyond the earlier Sentinel fix is that pending subscribe intent is tracked only in PubSubProxy, not in these standalone or cluster transfer paths.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Confirmed. Fixed in d15b39e by moving the fix into core PubSub: it now tracks in-flight subscribe intents, and removeAllPubSubListeners folds them into the returned snapshot and marks them carried — the late wire confirm no longer resurrects the subscription on the demoted client, a teardown-caused rejection resolves the caller (the subscription lives on the adopting member), and the cluster __MOVED and node-removal flows inherit the same fix. This generalizes the sentinel-proxy fix from the previous round to standalone and cluster. Unit tests added.
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 63d49c8bcd
ℹ️ About Codex in GitHub
Codex has been enabled to automatically review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
When you sign up for Codex through ChatGPT, Codex can also answer questions or update the PR, like "@codex address that feedback".
| movePubSub: async (from, to) => { | ||
| // removal (not a copy) keeps a recovering old member from re-subscribing | ||
| // server-side and double-delivering to the same listener functions | ||
| const listeners = from._getQueue().removeAllPubSubListeners(); |
There was a problem hiding this comment.
Carry pending unsubscriptions through the switch
When unsubscribe('x') is awaiting its server reply as a switch starts, this extraction still includes x and subscribes the target member to it. The old reply later removes x only from the source map, so the caller observes a successful unsubscribe while the listener remains registered and continues receiving messages from the new active member. Pending unsubscribe intent needs to be excluded from the snapshot or replayed against the target after adoption.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Confirmed — the exact mirror of the carried-subscribe fix from the last round. Fixed in 9641463: PubSub tracks in-flight UNSUBSCRIBEs and applies their removal to the extracted snapshot before returning it (per-listener unsubscribes keep co-subscribed listeners), so a channel the caller just left is no longer resurrected on the new member. Standalone and cluster share this core path.
| const active = mgr.activeDatabase; | ||
| // eslint-disable-next-line @typescript-eslint/no-explicit-any -- dynamic | ||
| const result = (resolve(active.client) as any)[name](...args); |
There was a problem hiding this comment.
Reject MONITOR on the failover wrapper
When a caller enters MONITOR mode and the active database subsequently switches, this generic forwarding leaves both the callback and the connection's monitor state on the demoted member. The callback therefore continues reporting the wrong endpoint, and if automatic fallback later selects that member, ordinary commands encounter a connection still in monitor mode until explicitly reset. Like SELECT, MONITOR must either be rejected or have its session state explicitly migrated or cleared during switches.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Confirmed. Fixed in 2ed8ffa: MONITOR is rejected through the wrapper like SELECT — it puts a connection into a permanent monitoring mode (re-applied on reconnect) that can't follow a failover — with guidance to attach it to an individual member; documented alongside the SELECT caveat.
| if (healthy.length < required) { | ||
| // a rejected connect() must not leave live sockets or retry timers behind | ||
| this.destroy(); | ||
| throw new Error( |
There was a problem hiding this comment.
Await teardown before rejecting connect
When the initial availability gate fails, connect() discards the promise returned by destroy() and rejects before asynchronous member teardown—most notably Sentinel teardown—and the logical end event have completed; the same pattern appears in the strategy-failure branches below. Fresh evidence beyond the earlier teardown fix is that MultiDbManager.destroy() now explicitly returns a promise that awaits every member, so these calls no longer satisfy the documented rejected-connect cleanup contract unless they are awaited.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Confirmed. Fixed in 3fea1df: the three connect() reject branches now await destroy() before throwing, so a rejected connect() leaves no member socket still closing or 'end' still pending by the time the caller sees the rejection.
| const reply = await target.sendCommand(['PING']); | ||
| // toString covers both string and Buffer type mappings | ||
| return reply?.toString() === 'PONG'; |
There was a problem hiding this comment.
Accept RESP2 subscribed-mode PING replies
When a standalone member uses RESP2 and has an active subscription, Redis answers a no-argument PING in subscribed mode as an array whose first element is lowercase pong, and the queue normalizes that reply to the string pong. This strict uppercase comparison therefore fails every background probe, opens the healthy active member's circuit, and can repeatedly move the subscription until every RESP2 member is considered down. Normalize the reply casing or explicitly accept the subscribed-mode response.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Confirmed, though config-gated: it bites only a RESP2 member that currently holds a subscription (not the RESP3 default). In RESP2 subscribe mode PING is allowed but its reply is array-framed with a lowercase verb, which the queue resolves to 'pong'. Fixed in ef1df0e: the default health probe compares case-insensitively.
| if (from.isReady) { | ||
| await Promise.allSettled([ | ||
| listeners[PUBSUB_TYPE.CHANNELS].size ? from.unsubscribe() : undefined, | ||
| listeners[PUBSUB_TYPE.PATTERNS].size ? from.pUnsubscribe() : undefined, | ||
| listeners[PUBSUB_TYPE.SHARDED].size ? from.sUnsubscribe() : undefined |
There was a problem hiding this comment.
Cancel stale pub/sub cleanup after a switch back
If switches occur A→B→A before the first handover finishes, the delayed cleanup from A→B can execute after the second handover has resubscribed A. It then sends UNSUBSCRIBE to the now-active A connection, leaving A's local listener map populated but its server-side subscription removed, so subsequent publications are silently lost. Before this cleanup runs, verify that from has not become active again or associate cleanup with the switch generation that created it.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Confirmed. Fixed in 6702425: isReady now reports false whenever the manager is gating commands behind unavailableError (searching or permanently failed), so it agrees with dispatch instead of reading a lingering member socket.
| Omit< | ||
| RedisClientKinds<M, F, S, RESP, TM>[K], | ||
| // the emitter methods are omitted too: the wrapper's documented event | ||
| // surface is exactly MultiDbClientEvents, and the kind's untyped | ||
| // `on(string, ...)` overload would otherwise swallow event-name typos | ||
| | 'duplicate' | 'withTypeMapping' | 'withCommandOptions' | 'withAbortSignal' | 'asap' | ||
| | keyof MultiDbEventEmitter |
There was a problem hiding this comment.
Override quit aliases to return Promise
For standalone clients, the inherited public signatures still type quit() and QUIT() as Promise<string>, while MultiDbClientBase deliberately implements both as fan-out teardown returning Promise<void>. Code such as (await client.quit()).toLowerCase() therefore compiles against MultiDbClientType and crashes at runtime because the resolved value is undefined. Omit and redeclare both aliases alongside duplicate() so their types match the wrapper contract.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Confirmed, resolved at the runtime layer instead of the type. Overriding quit/QUIT in the wrapper type tips the (already large) client-type Omit into a TS2590 'union too complex' error, and the mapped-type workaround regressed bystander import 'redis' typecheck cost ~73%. So 04bec74 instead makes quit()/QUIT() resolve the logical client's aggregate ack 'OK' once the graceful fan-out settles — the inherited Promise is now honest, with zero type surgery and zero typecheck-perf cost. QUIT is server-deprecated with no response policy, so there's no per-member reply semantics to preserve.
| // computed prop (`isOpen`, `options`) → live read from the active member | ||
| // eslint-disable-next-line @typescript-eslint/no-explicit-any -- dynamic | ||
| Object.defineProperty(dst, name, { get: () => (resolve(mgr.active) as any)[name], enumerable: false }); |
There was a problem hiding this comment.
Report an unavailable wrapper as not ready
When every circuit is open but the underlying active socket is still connected—for example, a lag-aware health check rejects otherwise reachable members—the manager gates every command with TemporarilyUnavailableError or PermanentlyUnavailableError, yet this getter continues returning the member's isReady === true. Applications that use client.isReady before dispatch will therefore send commands to a wrapper that explicitly cannot accept them. Special-case isReady so manager unavailability makes the logical client not ready.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Confirmed — the cleanup-side cousin of the earlier overlapping-handover finding. Fixed in 2669c0e with a root guard: a monotonic switch generation is stamped on every repoint, the async pub/sub handover captures it, and destructive from-cleanup runs only while it's still current — so a rapid A->B->A can't let a stale move unsubscribe a re-promoted member. setActiveDatabase additionally rejects a re-entrant forced switch.
| const subscriptions: Subscriptions = (this.#state && this.#state.connectPromise === undefined) ? { | ||
| [PUBSUB_TYPE.CHANNELS]: this.#state.client.getPubSubListeners(PUBSUB_TYPE.CHANNELS), | ||
| [PUBSUB_TYPE.PATTERNS]: this.#state.client.getPubSubListeners(PUBSUB_TYPE.PATTERNS), | ||
| [PUBSUB_TYPE.SHARDED]: this.#state.client.getPubSubListeners(PUBSUB_TYPE.SHARDED) | ||
| } : this.#subscriptions ?? { |
There was a problem hiding this comment.
Clear Sentinel's adopted snapshot after unsubscribing
After a Sentinel member adopts subscriptions during one multi-db handover, #subscriptions retains that adopted snapshot. If the user later unsubscribes from everything, #unsubscribe destroys the now-inactive pub/sub client and clears #state but leaves the snapshot intact; the next multi-db switch consequently takes this fallback and resurrects the already-unsubscribed channels on another member. Clear or refresh #subscriptions when the live pub/sub state becomes empty.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Confirmed. Fixed in 67500b1: the sentinel pub/sub proxy's inline teardowns (unsubscribe-all, and the executeCommand catch) now route through destroy(), which clears the adopted #subscriptions snapshot alongside the client and state — so a later extractListeners can't resurrect ghost channels on another member.
| function makeDerived<C extends RedisClientLike>( | ||
| mgr: MultiDbManager<C>, | ||
| resolve: ResolveClient<C> | ||
| ): C { | ||
| const view = new MultiDbClientBase(mgr); | ||
| attachForwarders(view, mgr, resolve); |
There was a problem hiding this comment.
Share failover events with derived views
After the derived-view type was changed to preserve MultiDbEventEmitter, callers can legitimately register failover, error, or lifecycle listeners on client.asap() and the other with* views. This constructor creates an independent EventEmitter, however, while the manager remains bound only to the root wrapper, so those listeners never fire. Derived views need to share or delegate to the root event emitter rather than owning an isolated event table.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Confirmed. Fixed in f4b55ce: derived views (asap(), withTypeMapping(), …) now delegate their event surface to the root wrapper the manager emits through, so a listener registered on a view actually fires. Chose delegation over dropping the emitter from the view type — the latter would add a parallel view type and risk the same TS2590/typecheck-perf hazard as the quit() case.
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: f4b55ce309
ℹ️ About Codex in GitHub
Codex has been enabled to automatically review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
When you sign up for Codex through ChatGPT, Codex can also answer questions or update the PR, like "@codex address that feedback".
| dst.withTypeMapping = (mapping: unknown) => | ||
| // eslint-disable-next-line @typescript-eslint/no-explicit-any -- member kinds type withX themselves | ||
| makeDerived(mgr, client => (resolve(client) as any).withTypeMapping(mapping), eventRoot); | ||
| dst.withCommandOptions = (options: unknown) => | ||
| // eslint-disable-next-line @typescript-eslint/no-explicit-any -- member kinds type withX themselves | ||
| makeDerived(mgr, client => (resolve(client) as any).withCommandOptions(options), eventRoot); |
There was a problem hiding this comment.
Preserve derived options when duplicating a view
When duplicate() is called on a withTypeMapping(), withCommandOptions(), withAbortSignal(), or asap() view, it still invokes the inherited MultiDbClientBase.duplicate() and creates members from the manager's original configs, discarding the view's resolve transformation. Fresh evidence beyond the earlier derived-view typing fix is that these overrides compose only further views, while the public duplicate type retains the derived mapping; for example, client.withTypeMapping({ [RESP_TYPES.BLOB_STRING]: Buffer }).duplicate().client is typed to return buffers but actually uses the default mapping and returns strings.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Confirmed. Fixed in b56e2b9: a view's duplicate() re-derives through the view's resolve over a fresh event root, so the clone keeps the view's mapping and does not cross-wire events into the original.
| const DELEGATED_EMITTER_METHODS = [ | ||
| 'on', 'once', 'off', 'addListener', 'removeListener', | ||
| 'prependListener', 'prependOnceListener' | ||
| ] as const; |
There was a problem hiding this comment.
Delegate bulk listener removal from derived views
A derived view now delegates individual listener methods to the root emitter, but its inherited removeAllListeners() still operates on the view's unused private EventEmitter table. Fresh evidence after the event-delegation fix is that view.on('failover', fn); view.removeAllListeners('failover') leaves fn registered on the root and it continues firing, even though removeAllListeners remains exposed by the underlying client type; this method and the remaining listener-inspection methods need to target the same root emitter.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Confirmed. Fixed in f9f6ad4: makeDerived delegates removeAllListeners / rawListeners / set|getMaxListeners to the root too, so clearing listeners on a view actually clears the root the manager emits through.
| const subscriptions: Subscriptions = (this.#state && this.#state.connectPromise === undefined) ? { | ||
| [PUBSUB_TYPE.CHANNELS]: this.#state.client.getPubSubListeners(PUBSUB_TYPE.CHANNELS), | ||
| [PUBSUB_TYPE.PATTERNS]: this.#state.client.getPubSubListeners(PUBSUB_TYPE.PATTERNS), | ||
| [PUBSUB_TYPE.SHARDED]: this.#state.client.getPubSubListeners(PUBSUB_TYPE.SHARDED) |
There was a problem hiding this comment.
Apply pending unsubscriptions during Sentinel handover
If a Sentinel unsubscribe() is awaiting its wire reply when a database switch occurs, these getters snapshot the still-present listener and the target Sentinel resubscribes it; the later reply only removes it from the destroyed source client, so the caller observes a successful unsubscribe while messages continue on the target. Fresh evidence after the core pending-unsubscribe fix is that Sentinel bypasses removeAllPubSubListeners()—which applies pending removals—and reads the live maps directly here.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Confirmed. Fixed in 1e2d8a1: sentinel extractListeners sources its snapshot from the inner client's removeAllPubSubListeners() (applies in-flight unsubscribes, carries in-flight subscribes) instead of the raw live maps.
| if (!pending.carried) removeListeners(); | ||
| this.#updateIsActive(); | ||
| }, | ||
| reject: undefined |
There was a problem hiding this comment.
Remove rejected unsubscriptions from pending state
When an unsubscribe command rejects before any handover, this undefined rejection hook leaves its new #pendingUnsubscribes record permanently registered. A later multi-db switch then applies that stale removal to the listener snapshot even though the caller was told the unsubscribe failed, so an existing subscription can silently disappear on failover; rejection must delete the pending record without applying its removal.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Confirmed. Fixed in f2ad81d: #unsubscribeCommand's reject hook drops the pending entry WITHOUT applying its removal, so a rejected in-flight unsubscribe no longer strands a stale removal that a later move would replay.
| if (this.#active !== target) { | ||
| // a recovery re-selection is a switch in everything but the | ||
| // announcement — the 'ready' below is its signal, not 'failover'. | ||
| // Skipping the housekeeping here once replayed a demoted member's | ||
| // unsent queue and stranded its subscriptions. | ||
| const from = this.#active; | ||
| this.#repoint(from, target); | ||
| this.#afterSwitch(from, target).catch(err => this.#emitError(err as Error)); | ||
| } |
There was a problem hiding this comment.
Restore the active role on same-member recovery
When the active member permanently ends, Database.#onEnd marks it DISCONNECTED, and a later successful connection changes it only to PASSIVE. If the documented recovery connect() selects that same member—always the case for a one-member configuration—this identity check skips #repoint, the only path that assigns ACTIVE; commands still route to the member, but getActiveDatabase() reports it as passive and the topology has no member with the active role.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Confirmed. Fixed in 606c56c: connect() asserts the selected target's role ACTIVE even when the selection did not move, so a same-member recovery no longer leaves the sole member PASSIVE.
| clearInterval(timer); | ||
| this.#healthTimers.delete(member); | ||
| } | ||
| this.#databases.splice(this.#databases.indexOf(member), 1); |
There was a problem hiding this comment.
Guard active removal against synchronous event re-entry
When removing the active member, switchTo() synchronously emits failover before this splice executes. If a failover listener also calls removeDatabase() for the from id—for example, generic cleanup attached to every failover—the nested call removes that member first; the outer call then evaluates indexOf(member) as -1 and splice(-1, 1) removes the replacement, potentially leaving an empty managed set with #active pointing outside it. Revalidate membership before splicing or mark the member as being removed before emitting the switch.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Confirmed. Fixed in fcf5a66: removeDatabase guards indexOf === -1 before the splice, so a re-entrant removeDatabase from a synchronous failover listener cannot delete the wrong member.
| if (!await runProbeRound(this.#targetFor(db), this.#healthChecks, this.#config.healthCheck, this.#teardown.signal)) { | ||
| const cause = new Error(`MultiDb: database "${db.id}" failed its health check`); | ||
| if (db === this.#active) { | ||
| this.#handleActiveFailure(cause, 'health-check'); | ||
| } else if (db.circuit.open() && this.#databases.includes(db)) { | ||
| // no announcement for a member removed while its round was in | ||
| // flight — its id may already belong to a new member | ||
| this.#events?.emit('database-unhealthy', { id: db.id, cause }); |
There was a problem hiding this comment.
Suppress health-check results after teardown
If close() or destroy() runs while a passive member's health check is awaiting a slow custom probe, teardown stops future intervals but does not cancel the in-flight check. When that probe later fails, this post-await branch still opens the circuit and emits database-unhealthy because closed members remain in #databases, so consumers can receive topology events after the logical end event; the recovery path can similarly emit database-recovered. Recheck the teardown signal after each awaited probe before mutating state or emitting.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Confirmed (recovery half; the database-unhealthy half was already covered by #handleActiveFailure's teardown guard). Fixed in fcf5a66: #recoveryProbe re-validates liveness after the probe await before emitting database-recovered.
| // #establishMember closes the circuit once the member establishes | ||
| await this.#establishMember(member, member.skipInitialHealthCheck); | ||
| this.#startMemberChecks(member); | ||
| return member.id; |
There was a problem hiding this comment.
Reject an add that completes after client teardown
If close() or destroy() occurs while addDatabase() is awaiting connection or a custom health probe, teardown includes and closes the newly pushed member, but this method never rechecks the abort signal after that await. It can therefore resolve successfully with an id after the logical client has ended, even though that member is closed and no scheduler was started; match the guarded connect() and forced-switch paths by rejecting once teardown has won the race.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Confirmed. Fixed in fcf5a66: addDatabase rejects if teardown landed during establish (teardown-only recheck; a concurrent removeDatabase still resolves, and #startMemberChecks self-guards membership).
PavelPashov
left a comment
There was a problem hiding this comment.
Should the default failure detector count every error as a failure? Errors like WRONGTYPE, CROSSSLOT, NOSCRIPT or a WATCH conflict come from a healthy database answering normally, so under load they could fail over away from it. Would it make sense to count only connection failures, timeouts and server-state errors by default?
| this.#failureRateThreshold = options.failureRateThreshold ?? MULTI_DB_DEFAULTS.failureDetector.failureRateThreshold; | ||
| this.#windowSize = options.windowSize ?? MULTI_DB_DEFAULTS.failureDetector.windowSize; | ||
| this.#errorFilter = options.errorFilter ?? (() => true); | ||
| this.#clock = options.clock ?? Date.now; |
There was a problem hiding this comment.
Could we use a monotonic clock here and in circuit.ts? System clock changes can distort both the failure window and the recovery wait. This also needs a node:perf_hooks import.
| this.#clock = options.clock ?? Date.now; | |
| this.#clock = options.clock ?? (() => performance.now()); |
There was a problem hiding this comment.
Confirmed, fixed in f55247d400: the detector and the circuit now default to performance.now(), so a wall-clock step cannot block eviction or stretch the grace period.
| if (member.circuit.state === 'CLOSED') { | ||
| await member.client.close(); | ||
| } else { | ||
| member.client.destroy(); |
There was a problem hiding this comment.
Could we await destruction before removing the member’s error listeners? Sentinel teardown is asynchronous, so removing them early can leave errors uncaught and teardown rejections unhandled.
| member.client.destroy(); | |
| await member.client.destroy(); |
There was a problem hiding this comment.
Confirmed, fixed in 744c9ca15d: removeDatabase now awaits destroy() before it disposes the member, so a late sentinel error still reaches a listener. Related: c0b7521cb2 stops cluster discovery that outlives destroy().
| if (this.#unavailable !== 'searching') { | ||
| return; // rescued mid-delay by a forced switch | ||
| } |
There was a problem hiding this comment.
Could we give each search loop an epoch and stop it if recovery or a newer search has changed that epoch? Otherwise, an old loop can wake up after recovery and a second outage, mistake the new 'searching' state for its own, and terminate the client using its nearly exhausted retry budget.
There was a problem hiding this comment.
Confirmed, fixed in 1c3ae377cb: each search has an epoch, and a loop exits when its epoch is stale. An old loop can no longer terminate the client after a rescue.
| const unavailable = mgr.unavailableError; | ||
| if (unavailable) return Promise.reject(unavailable); |
There was a problem hiding this comment.
Could we also reject here when the active database has changed since multi() was called? Right now a transaction built before a failover goes into the demoted member's offline queue and replays when that member reconnects. Using WatchError lets the usual retry loop re-run it on the new active member. Scan iterators have the same gap.
| const unavailable = mgr.unavailableError; | |
| if (unavailable) return Promise.reject(unavailable); | |
| const unavailable = mgr.unavailableError; | |
| if (unavailable) return Promise.reject(unavailable); | |
| if (member !== mgr.activeDatabase) { | |
| return Promise.reject(new WatchError('MultiDb: the active database changed after MULTI')); | |
| } |
There was a problem hiding this comment.
Confirmed, fixed in aa450b0fcc: a pinned MULTI rejects at exec/execAsPipeline/EXEC when its member is no longer the active one (identity check, so A→B→A still commits). One change from your suggestion: it throws WatchError only when a WATCH session is open; otherwise CommandAbandonedError, because a MULTI without WATCH has no optimistic-lock contract to report. Scan iterators throw CommandAbandonedError at the next batch after any switch. This replaces the old rule that exec runs on its pinned member after a forced switch; docs and tests are updated.
| this.#failoverInFlight = true; | ||
| this.#unavailable = 'searching'; | ||
| void this.#searchLoop(reason); |
There was a problem hiding this comment.
Could we reset the detector when the all-down search starts, and have onCommandResult ignore outcomes while #unavailable !== null? After a same-member recovery, switchTo returns early and skips the reset in #repoint, so the old failures and the reconnect errors that arrive during the search re-trip the detector on the first new error.
| this.#failoverInFlight = true; | |
| this.#unavailable = 'searching'; | |
| void this.#searchLoop(reason); | |
| this.#failoverInFlight = true; | |
| this.#unavailable = 'searching'; | |
| this.#detector.reset(); | |
| void this.#searchLoop(reason); |
There was a problem hiding this comment.
Confirmed, fixed in 1c3ae377cb: the detector resets when a search starts and ignores results while the client is unavailable, so a same-member rescue starts clean.
| for (const pending of this.#pendingSubscribes) { | ||
| pending.carried = true; | ||
| const typeListeners = result[pending.type]; | ||
| for (const channel of pending.channels) { | ||
| let channelListeners = typeListeners.get(channel); | ||
| if (!channelListeners) { | ||
| channelListeners = { unsubscribing: false, buffers: new Set(), strings: new Set() }; | ||
| typeListeners.set(channel, channelListeners); | ||
| } | ||
| PubSub.#listenersSet(channelListeners, pending.returnBuffers).add(pending.listener); | ||
| } | ||
| } | ||
|
|
||
| // in-flight unsubscribes will remove their channels/listeners once the | ||
| // reply lands, but those are still present in the snapshot above — apply | ||
| // the removal now (the thunk mutates the same map objects `result` | ||
| // aliases) so the move does not resurrect a channel the user just left | ||
| for (const pending of this.#pendingUnsubscribes) { | ||
| pending.carried = true; | ||
| pending.apply(); | ||
| } |
There was a problem hiding this comment.
Could we replay the pending subscribes and unsubscribes in the order they were issued, e.g. one insertion-ordered list of { kind: 'sub' | 'unsub', … } instead of two Sets? Right now every unsubscribe runs after every subscribe, so unsubscribe(); subscribe('x', fn) caught mid-flight by a failover clears x from the snapshot while the subscribe still resolves successfully, and the subscription is silently lost on the new member. PubSubProxy.extractListeners has the same ordering.
There was a problem hiding this comment.
Confirmed, fixed in bc27e7966d: in-flight pub/sub ops now sit in one ordered log and replay in issue order on a move, in both PubSub and the sentinel PubSubProxy. A carried unsubscribe now resolves instead of rejecting. A related pre-existing bug (a no-listener unsubscribe followed by a re-subscribe lost the listener) is fixed separately in d161baaaa2.
| (this.client as unknown as EventEmitter) | ||
| .on('error', this.#onError) |
There was a problem hiding this comment.
Could we also handle terminated as a member failure and remove that listener in dispose()? When reconnection is abandoned, standalone clients emit terminated without end, so traffic keeps hitting the closed member until health checks or detector thresholds trigger failover.
There was a problem hiding this comment.
Confirmed, fixed in 744c9ca15d: members now treat 'terminated' like 'end', so a client that gives up reconnecting triggers a connection-ended failover. It cannot fail over before the first 'ready'.
| await this.#recoveryProbe(db); | ||
| return; | ||
| case 'CLOSED': | ||
| if (!await runProbeRound(this.#targetFor(db), this.#healthChecks, this.#config.healthCheck, this.#teardown.signal)) { |
There was a problem hiding this comment.
Could we capture a per-member health epoch before this round, bump it on every successful verification (setActiveDatabase, connect()), and drop this round's result if the epoch changed? With pool, cluster or sentinel (masterPoolSize > 1) members, or custom/lag-aware checks, an outdated round can undo a forced switch:
- An automatic background check on database B starts.
- While it is still running,
setActiveDatabase('B')runs a second, separate check. It passes, so traffic moves to B and B is pinned. - The first check then finishes and reports that B failed.
- The client acts on that result. It moves traffic away from B and clears the pin, which undoes the forced switch.
There was a problem hiding this comment.
Confirmed, fixed in 1c3ae377cb: each member has a health epoch. A forced switch and a successful connect bump it, and a round that started before the bump drops its failure. Forced switches also pass the teardown signal, so close() no longer waits out probe delays.
| // view's own empty emitter: removeAllListeners() would clear nothing while | ||
| // the root keeps firing; rawListeners()/getMaxListeners() would read the | ||
| // wrong emitter. Chainable ones return the view, query ones the root's value. | ||
| dst.removeAllListeners = (event?: string) => { eventRoot.removeAllListeners(event); return dst; }; |
There was a problem hiding this comment.
Could we forward the arguments as given here? Node's EventEmitter checks arguments.length, so view.removeAllListeners() becomes removeAllListeners(undefined), which removes nothing and leaves every listener on the root firing.
| dst.removeAllListeners = (event?: string) => { eventRoot.removeAllListeners(event); return dst; }; | |
| dst.removeAllListeners = (...args: [string?]) => { eventRoot.removeAllListeners(...args); return dst; }; |
There was a problem hiding this comment.
Confirmed, fixed in aa450b0fcc: the view forwards its arguments as given, so removeAllListeners() with no argument clears every event. listenerCount forwards its optional listener too.
|
Agreed, fixed in |
Sketch of a multi-database client wrapper managing a homogeneous array of
underlying clients (standalone/pool/cluster/sentinel) behind one drop-in
client surface. Type mechanics only — failover/health/routing stubbed.
- @redis/client: dedicated factories (createMultiDbClient/Pool/Cluster/
Sentinel) returning { client, controller }; client typed exactly as the
base client (true drop-in), controller holds multi-db-only admin surface.
Command forwarding via prototype-walk (no runtime Proxy on hot path).
- redis meta-package: createMultiDbClient wrapper injecting default Stack
modules, mirroring createClient.
- playground.ts for manual poking.
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
…stubs Move MultiDbManager/MultiDbController out of index.ts, add typed stubs for config/circuit/database/failure-detector/health-check/failover-strategy/errors per the multi-db API contract, export the public surface from the package index, accept flat MultiDbConfig options in all factories, and drop the playground script (superseded by the quickstart smoke script). Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
…per, switch primitive - config.ts: defaults table + resolution/validation (weights in [0,1], unique ids with generated db-<n> fallback, health-check timing, >=1 db) - circuit.ts: CLOSED/OPEN/HALF_OPEN machine with clock-derived HALF_OPEN, grace-period restart on probe failure, consecutive-probe counting - database.ts: member wrapper (id/weight/circuit/role) with client lifecycle listeners feeding circuit and role; listeners are disposed only after close()/destroy() so a teardown-window 'error' emit cannot crash the process - manager.ts: atomic switch primitive with typed controller events and non-awaited old-member housekeeping; per-command result hook wired through the forwarding closures for the upcoming failure detector - controller.ts: typed event map (failover/fallback/database-unhealthy/ database-recovered/all-databases-down/error), descriptor-based getDatabases()/getActiveDatabase() - comments follow the constitution v1.0.0 principles: contract JSDoc (+@experimental) on the public surface, why-only implementation comments, cross-module constraints with file:symbol references - unit tests for circuit transitions and config resolution (no docker) Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
Remove spec-artifact references from comments (requirement IDs, research citations, data-model mentions) — the spec docs don't ship, so the references would dangle for external readers; each comment now carries its rationale inline. Task/story markers on unimplemented stubs stay as development scaffolding and are removed by the change that implements each stub. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
- health-check.ts: default PING check, chain runner with per-probe
timeout bounding the whole chain, ALL/MAJORITY/ANY early-exit round
aggregation, round budget used as the member readiness bound
- manager.ts: initial connection flow — fan-out connect with bounded
establishment, per-member initial health check, initialAvailability
gate ('all'/'majority'/'one') with full teardown on rejection,
weight-based active selection with config-order tiebreak; runtime
addDatabase (pre-opened circuit until established, skipInitialHealthCheck
honored only here), removeDatabase (active switches to replacement
first), setWeight; repeat connect() re-probes instead of tripping
circuits of already-open members
- index.ts: per-topology MemberAdapter (create + keyless sendCommand)
wired through all four factories into manager construction
- controller.ts: addDatabase/removeDatabase/setWeight admin surface
- config.ts: probe-knob validation (numProbes, timeout, delays, empty
healthChecks), shared per-member identity resolution
- tests: probe-policy unit matrix (no docker); docker integration suite —
weighted selection, availability matrix with dead members, runtime
reconfiguration incl. skip-flag semantics, repeat connect, drop-in typing
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
…g, pub/sub transfer - failure-detector.ts: sliding-window DefaultFailureDetector (count AND rate thresholds, 0 disables a condition, error filter counts filtered errors as traffic but not failures; injectable clock) - failover-strategy.ts: WeightBasedStrategy — highest-weight CLOSED member, member-order tiebreak - manager.ts: detection→failover wiring — per-command outcomes and member lifecycle errors feed the detector with source attribution (a stale in-flight rejection from the previous active must not trip the new one); failover procedure opens the failed circuit, switches to the strategy's pick, or gates traffic and retries selection up to maxFailoverAttempts before going permanently unavailable; commands fail fast with TemporarilyUnavailableError while searching and PermanentlyUnavailableError after exhaustion - pub/sub transfer on switch (standalone members): listeners are removed from the old member and re-subscribed on the new one, so a recovering old member cannot double-deliver - database.ts: per-topology signal mapping — sentinel per-node errors arrive via client-error and a sentinel-internal master change must not open the member circuit; cluster node errors aggregate into 'error' - tests: detector/strategy unit suites; failover integration (kill active → event + traffic continuity, in-flight rejection, pub/sub survival, all-down escalation); topology integration (cross-cluster failover, sentinel no-false-failover during its own master change, cross-sentinel-deployment failover) Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
- manager.ts: per-member unref'd health scheduler covering the active member too (a silently dead server never trips the organic detector); OPEN members rest through their grace period, HALF_OPEN members get recovery probes fed straight into the circuit (closing emits database-recovered), CLOSED members failing a round open with database-unhealthy — or fail over with reason 'health-check' when active; recovery probing keeps running during the all-down search so attempts can succeed; permanent unavailability stops all timers - auto-fallback loop (disabled by default): returns traffic to a strictly higher-weight healthy member, emitting 'fallback'; controller.setAutoFallback(intervalMs | false) retunes it at runtime - passive members that end or fail checks announce database-unhealthy without switching; deliberate removals stay silent - tests: recovery closes the circuit only after grace + probes (no flapping on a fast-restarting member), auto-fallback on/off/runtime toggle, passive failure reported without failover Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
…alth check - lag-aware-health-check.ts: probes the RE cluster REST API's database availability endpoint with lag verification (extend_check=lag + availability_lag_tolerance_ms, default 5s tolerance) via the global fetch — no new dependency; basic auth, per-database endpoint/uid resolvers for members on different clusters, request timeout, and a requestOptions escape hatch for custom TLS dispatchers - custom failure detectors, chained health checks and custom failover strategies were already threaded through config — now exercised end-to-end and exported (LagAwareHealthCheck added to the public surface) - tests: stub-HTTP-server suite for the lag-aware check (availability, lag/auth failures, timeout, unreachable endpoint, per-database resolvers); integration cases for an error-filtering detector, a user-supplied detector driving failover (and being reset by the switch), and an all-must-pass check chain Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
- manager.ts: setActiveDatabase(id) health-checks the target first — present reality overrides a stale OPEN circuit, so a verified target is closed (announced recovered) before the forced switch; the forced selection pins, suspending auto-fallback until releasePin(); any automatic switch away from the pin clears it, so a pin never traps traffic on a dead member; a successful force also rescues a client mid-search with every member down - controller.ts: setActiveDatabase/releasePin admin surface - tests: pin holds against auto-fallback until released; unhealthy target rejected; automatic failover off a dead pinned member clears the pin (proven by the fallback loop resuming); same-member force pins without a switch; force-rescue during the all-down search Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
- docs/multi-db.md: quick start, selection and failover lifecycle, configuration reference with defaults, controller and events reference, custom checks incl. the lag-aware example, and the behavior contracts (in-flight rejection, eventual consistency, pub/sub loss window and per-topology transfer support, per-member client-side caching, resource overhead) - examples/multi-db-failover.js: runnable two-container demo of failover, recovery and fallback - @experimental on every public multi-db export - client-side caching regression test: no stale reads across a switch, no flush needed - topology coverage: pool failover, per-factory typing assertions, forced-pin smokes on cluster and sentinel members, sentinel fallback-after-recovery; the kill tests accept either automatic failover reason — the organic detector and the background health check legitimately race on slow-to-reject topologies Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
Extend the failover pub/sub handoff (previously standalone-only) to cluster and sentinel members. On a switch, active subscriptions (channels, patterns, sharded) detach from the old member and re-subscribe on the new one, so a recovering old member cannot double-deliver. - cluster: RedisClusterSlots.removeAllPubSubListeners gathers listeners from the pub/sub node and each shard; exposed via RedisCluster and paired with the existing resubscribeAllPubSubListeners - sentinel: PubSubProxy.extractListeners/adoptListeners, surfaced through RedisSentinelInternal and the wrapped RedisSentinel - pool: no pub/sub surface, so nothing to transfer (documented) Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
…eners Mirror of the carried-subscribe fix: a channel being unsubscribed is still in the listener maps until the UNSUBSCRIBE reply lands, so a subscription move (multi-db failover, cluster __MOVED/node removal) snapshotting mid-round-trip captured it and re-subscribed it on the new member — resurrecting a channel the user just left, while their unsubscribe() resolved as success. PubSub now tracks in-flight unsubscribes and applies their removal to the extracted snapshot before returning it (per-listener unsubscribes correctly keep co-subscribed listeners); the late reply becomes a no-op against the fresh maps. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
The sentinel pub/sub proxy stores adopted subscriptions in #subscriptions (set by adoptListeners), and extractListeners falls back to that snapshot while a connect is in flight. But the unsubscribe-all teardown and the executeCommand catch tore the client down inline (client.destroy() + #state = undefined) without clearing #subscriptions — only destroy() does. So after adopt then unsubscribe-everything, a later extract returned the stale snapshot and resurrected ghost channels on the next member. Both inline teardowns now route through destroy(), which clears state, client and snapshot in one place. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
…probe PING is one of the few commands allowed while a RESP2 connection is in subscribe mode, but its reply comes back array-framed with a lowercase verb, which the queue resolves to the string 'pong'. DefaultHealthCheck compared strictly against 'PONG', so a healthy RESP2 member that happens to hold a subscription failed every probe and had its circuit opened. The compare is now case-insensitive. RESP3 is unaffected (subscribe mode does not restrict or reframe commands there), so this only ever bit RESP2 + active subscription. Ref: SUBSCRIBE docs — in RESP2 a subscribed client may only issue SUBSCRIBE/PSUBSCRIBE/SSUBSCRIBE/UNSUBSCRIBE/PUNSUBSCRIBE/SUNSUBSCRIBE/ PING/RESET/QUIT; RESP3 lifts the restriction. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
MONITOR puts a connection into a permanent monitoring mode that it re-applies on every reconnect — connection-scoped state that cannot follow a failover, the same class as SELECT. Forwarded generically it monitored only the active member and left a demoted member stuck in monitor mode. It is now rejected with guidance to attach it to an individual member, and documented alongside the SELECT caveat. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
The isReady getter read straight through to the active member's socket, so it could report ready while the manager gated every command behind unavailableError (searching, or permanently failed) — a member socket can linger connected while its circuit is open. An app that checks isReady before dispatching would send into a guaranteed rejection. isReady now means 'the logical client can serve': false whenever unavailableError is set, otherwise the active member's socket state. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
A switch's pub/sub handover runs async and unawaited, so a rapid switch-back (A→B→A) let the A→B handover's late from-cleanup run after B→A had re-subscribed A — unsubscribing the member that just became active again, silently dropping its subscriptions. switchTo is deliberately synchronous and always wins (automatic failover cannot queue behind a previous move), so the fix guards the async tail, not the switch: a monotonic #switchGeneration is bumped on every repoint, each handover captures it, and the movePubSub adapter contract now passes isCurrent() — destructive from-cleanup runs only while it holds. This is the root guard for the whole 'a switch's async continuation acts on superseded state' class. setActiveDatabase additionally rejects a re-entrant forced switch (one forced switch at a time) as a tightening on the user-facing path. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
quit()/QUIT() fanned out to Promise<void> while the inherited base-client signature said Promise<string>, so (await client.quit()).toLowerCase() compiled and crashed. Rather than override the type — which tips the wrapper's Omit over the ~700-member client type into TS2590, and the mapped-type workaround regressed bystander import-'redis' typecheck cost by ~73% — the runtime now returns the logical client's own aggregate ack 'OK' once the graceful fan-out settles (best-effort per member, like close()). That makes the inherited Promise<string> honest with zero type surgery and zero typecheck-perf cost. QUIT is server-deprecated with no response policy, so there is no per-member reply semantics to preserve; 'OK' is the conventional acknowledgement. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
The type advertises the multi-db event surface on derived views
(client.asap(), client.withTypeMapping(...)), but makeDerived built each
view as its own EventEmitter and only makeClient binds the manager to
the root wrapper — so view.on('failover', ...) compiled and never fired.
makeDerived now receives the root wrapper and forwards the view's
registration/query emitter methods to it (registrations return the view
for chaining), so a listener attached to a view observes the same
manager emissions as the root. Matches the documented single event
surface; no type change.
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
…lient connect() destroys every member on failure — correct for an initial connect (or recovery from permanent unavailability, where nothing live is lost), but connect() is documented as callable again on a live client, and a repeat connect() that misses the availability gate (e.g. default MAJORITY with 2 members while one blips down) then destroyed the healthy active member and made the instance terminal. Failure now destroys only when the client has never reached 'ready' (#everReady); a repeat/recovery connect on a client that has served rejects and leaves the live members serving. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
… stress test
Every round of review surfaced the same class: an async method awaits,
the synchronous switch / a teardown / a membership change moves shared
state underneath it, and it acts on stale state. It recurred because the
guard was hand-rolled ad-hoc at each site and forgotten at new ones.
Two named guards now express the invariant uniformly:
#memberLive(db) — teardown + membership, for per-member async work
#currentSwitch(gen) — switch supersession, for fire-and-forget tails
kept as distinct dimensions on purpose (folding teardown into the switch
generation would abort unrelated tails and obscure the reason). Every
await site was audited; the missing guards added:
- #recoveryProbe re-validates member liveness + HALF_OPEN after the
probe await before probeSucceeded/probeFailed/emit (R4-1, R4-10)
- addDatabase rejects if teardown landed during establish (R4-4)
- removeDatabase guards indexOf === -1 before splice, so a re-entrant
removeDatabase from a synchronous failover listener can't delete the
wrong member (R4-9)
- #checkMember's passive-unhealthy emit re-checks liveness (teardown)
circuit.probeFailed() is now a no-op unless HALF_OPEN, mirroring
probeSucceeded — a stale recovery probe can't reopen a closed/active
member (R4-1 at the state-machine level).
A new 're-validation after await' test block is the enforcement net:
it moves the world at each async op's await boundary (destroy mid-probe,
destroy mid-establish, re-entrant remove from a failover listener) and
asserts no effect leaks onto stale state.
Resolves round-4 items R4-1, R4-4, R4-9, R4-10.
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
When the active member ended it went DISCONNECTED, and a recovery connect() that re-selected the same member only flipped it to PASSIVE (via its 'ready' handler) — #repoint, the sole other ACTIVE setter, is skipped when the selection doesn't move. So after recovery no member held the active role and getActiveDatabase().role reported PASSIVE (guaranteed on any one-member config). connect() now asserts target.role = 'ACTIVE' unconditionally after selection. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
#unsubscribeCommand tracked a pending removal but its reject hook was undefined, so an unsubscribe that rejected (socket drop / error reply) left its entry in #pendingUnsubscribes forever. The next removeAllListeners (a subscription move) then replayed that stale removal against the snapshot and dropped a subscription that legitimately stayed. The reject hook now deletes the pending entry WITHOUT applying the removal — mirroring the subscribe path — so a failed unsubscribe leaves the channel subscribed and movable. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
extractListeners snapshotted the inner client's listener maps via getPubSubListeners, which still hold a channel that is mid-unsubscribe (removal happens on the wire reply). A failover racing an unsubscribe therefore carried the leaving channel to the new member and resurrected it. It now sources the snapshot from the inner client's removeAllPubSubListeners(), which applies the client's in-flight UNSUBSCRIBEs and carries its in-flight SUBSCRIBEs — the same reconciliation the core client fix already performs — so a leaving channel is dropped and a kept one still moves. The client is destroyed immediately after, so clearing its maps as a side effect is harmless. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
A derived view inherited the base duplicate(), which rebuilds through makeClient with the identity resolve — silently dropping the view's type mapping/command options while the type still promised them (view.withTypeMapping(...).duplicate().client.get() typed Buffer, returned string). makeDerived now overrides duplicate() to re-derive through the view's own resolve, over a FRESH event root bound to the duplicated manager so the clone's events are not cross-wired into the original. Pinned by a types-test (mapping preserved on the clone) and a runtime test (Buffer reply + event isolation). Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
…ived views
makeDerived delegated the per-listener methods to the root but left
removeAllListeners, rawListeners and set/getMaxListeners pointing at the
view's own (empty) emitter — so view.removeAllListeners('failover')
cleared nothing while the root kept firing, and rawListeners/getMaxListeners
read the wrong emitter. All of them now delegate to the root: the
chainable ones return the view, the query ones return the root's value.
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
…scue - A rescued search loop now exits once a newer all-down search exists. Before, a rescue followed by a second outage within one delay let the old loop run on and use up the new search's attempts. - The failure detector resets when a search starts and ignores command results while the gate is up, so a same-member recovery starts clean. - A background health round that started before a forced switch or a connect() verified the member drops its failed verdict instead of undoing that verification. - Forced switches and connect() pass the teardown signal to their probe round, so close() rejects a pending force at once. - Docs: recovery from all-databases-down can show as database-recovered when the active member comes back.
- A standalone client that gives up reconnecting emits 'terminated', not 'end'. Members now treat 'terminated' as down, so the active fails over at once instead of waiting for the next health round. Before the first 'ready', connect() owns the outcome and no failover runs. - removeDatabase now awaits destroy() before it detaches the member's listeners. A sentinel destroy settles asynchronously and can still emit 'error'; without the await that error had no listener. - Docs and JSDoc now say which member kinds emit 'terminated' and which rely on health-check probes.
The temporary client that fetches the slot map is not tracked by destroy(). A slots reply that lands after destroy() used to rebuild node clients that nothing would ever close, or emit 'error' when nobody was listening any more. - #discover returns early when the cluster closed while the shards were in flight, and its catch emits 'error' only while the cluster is open. - connect() cut short by destroy() during the last root node's discovery now rejects with 'Cluster closed', the same error as a destroy() between root nodes, not RootNodesUnavailableError. - A background topology refresh cut short by destroy() no longer emits 'error'.
…l emitter args A multi() or scan iterator pinned to a member kept running there after traffic moved away, so a stale transaction could commit on a demoted database and a scan could keep reading it. - multi(): exec/execAsPipeline/EXEC commit nothing once the pinned member is no longer active. Under a watch session they reject with WatchError, so the app's usual retry loop re-runs on the new member; otherwise they reject with CommandAbandonedError. The check is by identity, so a switch away and back still commits. - scan iterators: after the active member changes, the next batch rejects with CommandAbandonedError. A SCAN already in flight still yields. - derived views: removeAllListeners() and listenerCount() forward their optional arguments as-is. Node branches on arguments.length, so the old wrapper's explicit undefined made a no-arg removeAllListeners() clear nothing, and listenerCount(event, fn) ignored the filter.
A subscription move folded every in-flight subscribe into the snapshot, then applied every in-flight unsubscribe. Order was lost, so "unsubscribe x, then subscribe x again" before the replies landed moved x as unsubscribed, and a re-subscribe of the same listener was dropped. PubSub and the sentinel PubSubProxy now keep one ordered log of in-flight ops and replay it onto the snapshot in issuance order. The move consumes the log, so a late reply cannot replay a stale removal into a later move. Also: - a carried unsubscribe now resolves for its caller (PubSub exposes carried() on unsubscribe commands; the proxy tracks unsubscribes as it already did subscribes); - the proxy replays parked unsubscribes too, and skips ops a live inner client already tracks; - moved entries arrive with `unsubscribing` reset, so a flag left by a failed unsubscribe does not make the adopter re-send SUBSCRIBE.
unsubscribe(channel) and unsubscribe() without a listener did not flag the entries they remove as `unsubscribing`, unlike the per-listener path. A subscribe issued before the UNSUBSCRIBE reply landed saw the entry, joined it without sending SUBSCRIBE, and the reply then deleted the entry along with the new listener. The channel ended unsubscribed on the server and the listener silently dropped. Both paths now flag the affected entries, so the re-subscribe goes to the wire and its reply re-creates the entry after the UNSUBSCRIBE reply clears it.
…lock The default errorFilter counted every error. A burst of command errors (WRONGTYPE, CROSSSLOT, a WATCH conflict, an aborted command) could trip the detector and fail over a healthy member. The new exported defaultErrorFilter counts an error only when it is about the member: - WATCH conflicts and aborts never count; - an error reply counts only when its code reports server state: LOADING, BUSY, MASTERDOWN, CLUSTERDOWN, READONLY, NOREPLICAS, MISCONF (matched on the whole first word, so BUSYKEY does not); - a MULTI error counts when any of its replies counts; - anything else (connection errors, timeouts) counts. Ignored errors still count as traffic in the failure rate. TRYAGAIN (a healthy cluster mid-resharding) and OOM are not counted. The filter lives in its own module so config.ts and failure-detector.ts keep no import cycle. The detector window and the circuit grace period now default to performance.now() instead of Date.now(). A wall-clock step (NTP correction, VM resume) moved the window cutoff, so old failures stayed in the window or a grace period never ended. Only time differences are used, so a monotonic source is enough.
…rpose Three tests trip the detector with NOSUCHCOMMAND to check that namespace, view and exec results reach it. The new default error filter ignores that reply, so the tests now set errorFilter: () => true to keep their intent.
An unsubscribe flags the entries its reply will delete, so a re-subscribe before the reply sends SUBSCRIBE instead of joining a leaving entry. When the unsubscribe failed (socket drop, error reply), the flag stayed set on a channel that is still subscribed. Every later subscribe to that channel then sent a redundant SUBSCRIBE and attached its listener only at the reply. Each unsubscribe now records the entries it flagged. On reject, it clears the flag unless another in-flight unsubscribe flagged the same entry (that reply still deletes it, so a fast-path join would lose the listener). A carried op leaves its entries alone: they belong to the adopting client.
| // no movePubSub: a pool has no pub/sub surface (subscriptions need a | ||
| // dedicated connection, which the pool does not expose), so there is | ||
| // nothing to transfer on failover | ||
| }; |
There was a problem hiding this comment.
Pool failover can replay abandoned writes
Medium Severity
The pool adapter has no rejectQueued hook. Standalone failover rejects unsent commands with CommandAbandonedError so they cannot execute on the demoted member after it reconnects. Pooled members still hold those commands on idle clients' offline queues, so a recovered pool can apply stale writes after traffic has already moved.
Reviewed by Cursor Bugbot for commit dc7d704. Configure here.
…b/sub state The RESP2 unsubscribe wrapper checked `isActive` before the inner resolve removed the listeners and recomputed it. So the reset saw the old value and skipped. The decoder then stayed in pub/sub mode, and later replies on that connection came back as Buffers. This happens on a plain subscribe followed by unsubscribe. It also happens on a multi-db move: adopted or resubscribed channels decrement the subscribe counter without recomputing `isActive`. Resolve first, then check.
The move empties the old member's listener maps, then unsubscribes it. An argument-less UNSUBSCRIBE completes when the server count matches the pattern subscriptions still held locally, which is now 0. But the server still counts the patterns it holds, so the UNSUBSCRIBE never matched. The PUNSUBSCRIBE's last reply resolved it instead, and every later reply on that member resolved the command queued before it. Unsubscribe the moved channels, patterns and sharded channels by name. That waits for one reply per name and does not depend on the local maps.
…ands flushWaitingForReply and flushAll reset the pub/sub state before they rejected the in-flight commands. A rejected in-flight SUBSCRIBE then decremented the subscribe counter again, to -1. The pub/sub state then stayed active: on RESP2, a later full unsubscribe never reset the decoder, and isPubSubActive stayed true. Reset after the rejects.
There was a problem hiding this comment.
Cursor Bugbot has reviewed your changes using high effort and found 2 potential issues.
There are 3 total unresolved issues (including 1 from previous review).
Reviewed by Cursor Bugbot for commit 9ffbe81. Configure here.
| this.#unavailable = null; | ||
| this.#failoverInFlight = false; | ||
| this.switchTo(target, reason); | ||
| return; |
There was a problem hiding this comment.
Same-member rescue leaves PASSIVE role
Medium Severity
When the last active member ends and later comes back, search (and same-member setActiveDatabase) call switchTo which no-ops on identity, so the member stays PASSIVE after onReady. connect() already forces ACTIVE for this case, but those paths never do, so getActiveDatabase() / getDatabases() show no active member while traffic is served.
Additional Locations (2)
Reviewed by Cursor Bugbot for commit 9ffbe81. Configure here.
| listeners[PUBSUB_TYPE.PATTERNS].size ? from.pUnsubscribe([...listeners[PUBSUB_TYPE.PATTERNS].keys()]) : undefined, | ||
| listeners[PUBSUB_TYPE.SHARDED].size ? from.sUnsubscribe([...listeners[PUBSUB_TYPE.SHARDED].keys()]) : undefined | ||
| ]); | ||
| } |
There was a problem hiding this comment.
Failed pub/sub adopt skips cleanup
Medium Severity
Standalone movePubSub awaits adopt on the new member before unsubscribing the old one. If extendPubSubListeners rejects, the named unsubscribe never runs, so a still-ready RESP2 member remains in server-side subscriber mode with empty local maps. Traffic that later returns to it cannot run regular commands until that connection drops.
Reviewed by Cursor Bugbot for commit 9ffbe81. Configure here.


Adds a multi-database client to
@redis/client: one drop-in client backed by N member databases (standalone / pool / cluster / sentinel). All traffic goes to one weight-selected active member; the client fails over automatically when it degrades. Every new export is@experimental.How it works
client.json.*), and derived views (withTypeMapping& friends) all follow failover, feed the detector, and fail fast while every member is down. Deliberately pinned surfaces (multi(), scan iterators,legacy()) are documented;multi()'sexecstill gates and reports.duplicate()clones the whole multi-db client;replaceDatabase()swaps a member in one call.connect/ready/end/terminated), optionalerror(no listener, no crash), the failover events, andmember-*pass-throughs with ids in payloads. The controller is a pure admin handle (topology, weights, forced failover with pinning, runtime add/remove/replace).CommandAbandonedError) so nothing replays when it recovers.close()/destroy()are terminal; a permanently unavailable client recovers viaconnect().FailoverCandidateview.Testing
~100 multi-db tests: a fast stub-adapter unit suite for the manager's decision paths (no docker) plus docker integration suites across all four topologies — failover, fallback, forced failover with pinning, pub/sub transfer (including RESP2 and second-switch cases), permanent-unavailability escalation and recovery, half-dead clusters, sentinel-internal promotions. Full
@redis/clientsuite green.Notes for reviewers
removeAllListenersupdatesisActivecorrectly).docs/multi-db.md(full guide incl. behavior contracts and caveats) andexamples/multi-db-failover.js.🤖 Generated with Claude Code
Note
High Risk
New experimental routing/failover layer changes command, transaction, pub/sub, and session semantics across all client topologies; incorrect failover or pub/sub migration could cause lost messages, stale reads, or abandoned commands in production.
Overview
Introduces an experimental multi-database client that fronts several equivalent Redis endpoints (standalone, pool, cluster, or sentinel) with automatic failover: weight-based active selection, per-member circuit breakers, sliding-window failure detection, and background health checks. Factories (
createMultiDbClient, pool/cluster/sentinel variants) return a drop-inclientplus acontrollerfor topology, weights, forced failover, and runtime add/remove/replace.The wrapper is the only event surface (
failover,fallback,database-unhealthy,member-*, etc.); commands are not retried on switch—queued work on a demoted member can reject withCommandAbandonedError. Pub/sub subscriptions are moved to the new active member via reworkedPubSubpending-op replay and cluster_removeAllPubSubListeners/_extendAllPubSubListeners. Cluster teardown during discovery no longer leaks node clients or spuriouserrorevents; sentinel re-emitsclient-erroron the public client for multi-db fault classification.Adds
docs/multi-db.md,examples/multi-db-failover.js, and a large@redis/clientexport surface (errors, health checks includingLagAwareHealthCheck, strategies, types).Reviewed by Cursor Bugbot for commit 9ffbe81. Bugbot is set up for automated code reviews on this repo. Configure here.