Skip to content

feat(client): multi-db client with automatic endpoint failover - #3435

Open
nkaradzhov wants to merge 99 commits into
redis:masterfrom
nkaradzhov:multi-db
Open

nkaradzhov wants to merge 99 commits into
redis:masterfrom
nkaradzhov:multi-db

Conversation

@nkaradzhov

@nkaradzhov nkaradzhov commented Sep 2, 2026 •

Copy link
Copy Markdown
Collaborator

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.

import { createMultiDbClient } from 'redis';

const { client, controller } = createMultiDbClient({
  databases: [
    { id: 'east', options: { url: 'redis://east:6379' }, weight: 1 },
    { id: 'west', options: { url: 'redis://west:6379' }, weight: 0.5 }
  ]
});
client.on('failover', ({ from, to, reason }) => console.log(`${from} -> ${to} (${reason})`));
await client.connect();

How it works

  • Per-member circuit breaker (CLOSED/OPEN/HALF_OPEN, clock-derived, grace period + consecutive-success recovery) plus a sliding-window failure detector fed by every command outcome and the active member's connection errors.
  • Atomic switch primitive: the repoint is one synchronous assignment; forwarders resolve the active member at call time. Pub/sub subscriptions move with the traffic on every topology that has them, and the old member is detached server-side.
  • Classified forwarding: plain commands, module/function/script namespaces (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()'s exec still gates and reports. duplicate() clones the whole multi-db client; replaceDatabase() swaps a member in one call.
  • The client is the event surface: logical lifecycle (connect/ready/end/terminated), optional error (no listener, no crash), the failover events, and member-* pass-throughs with ids in payloads. The controller is a pure admin handle (topology, weights, forced failover with pinning, runtime add/remove/replace).
  • Failure contract: commands are rejected on failover — never retried on the new member, and a dead member's unsent queue is abandoned at the switch (CommandAbandonedError) so nothing replays when it recovers. close()/destroy() are terminal; a permanently unavailable client recovers via connect().
  • Pluggable health checks (default PING; ALL/MAJORITY/ANY probe policies; a lag-aware Redis Enterprise REST check) and failure detectors; custom failover strategies are typed via the exported FailoverCandidate view.

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/client suite green.

Notes for reviewers

  • Existing public API is unchanged; the branch also fixes internal pub/sub plumbing it depends on (cluster resubscribe no longer floats promises, sentinel listener extraction snapshots live state, removeAllListeners updates isActive correctly).
  • Known accepted gap: no TLS-member integration test yet (needs certificate-configured containers).
  • Docs: docs/multi-db.md (full guide incl. behavior contracts and caveats) and examples/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-in client plus a controller for 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 with CommandAbandonedError. Pub/sub subscriptions are moved to the new active member via reworked PubSub pending-op replay and cluster _removeAllPubSubListeners / _extendAllPubSubListeners. Cluster teardown during discovery no longer leaks node clients or spurious error events; sentinel re-emits client-error on the public client for multi-db fault classification.

Adds docs/multi-db.md, examples/multi-db-failover.js, and a large @redis/client export surface (errors, health checks including LagAwareHealthCheck, strategies, types).

Reviewed by Cursor Bugbot for commit 9ffbe81. Bugbot is set up for automated code reviews on this repo. Configure here.

@cursor cursor Bot left a comment •

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Stale Bugbot comment from a previous run.

Comment thread packages/client/lib/multi-db/manager.ts
Comment thread packages/client/lib/multi-db/index.ts
Comment thread packages/redis/index.ts Outdated
Comment thread packages/client/lib/multi-db/manager.ts

@chatgpt-codex-connector chatgpt-codex-connector Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

💡 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".

Comment thread packages/client/lib/multi-db/index.ts Outdated
*/
const INTERCEPTED = new Set<PropertyKey>([
'connect', 'close', 'destroy', 'quit',
'withTypeMapping', 'withCommandOptions', 'withAbortSignal',

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1 Badge 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 👍 / 👎.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Comment on lines +184 to +186
const member = mgr.activeDatabase;
// eslint-disable-next-line @typescript-eslint/no-explicit-any -- dynamic patching
const inner = (resolve(member.client) as any).multi();

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1 Badge 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 👍 / 👎.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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).

Comment on lines +45 to +48
const INTERCEPTED = new Set<PropertyKey>([
'connect', 'close', 'destroy', 'quit',
'withTypeMapping', 'withCommandOptions', 'withAbortSignal',
'multi', 'MULTI', 'duplicate'

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge 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 👍 / 👎.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Comment thread packages/redis/index.ts Outdated
Comment on lines +124 to +126
>(options: {
databases: Array<DatabaseConfig<RedisClientOptions<M, F, S, RESP, TYPE_MAPPING>>>;
} & MultiDbConfig): MultiDbResult<RedisClientType<M, F, S, RESP, TYPE_MAPPING>> {

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge 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 👍 / 👎.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Comment thread packages/redis/index.ts Outdated
Comment on lines +128 to +133
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>>;

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge 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 👍 / 👎.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Comment on lines +224 to +227
if (!(healthCheck.timeout < healthCheck.interval)) {
throw new TypeError(
`MultiDb: healthCheck.timeout (${healthCheck.timeout}) must be less than healthCheck.interval (${healthCheck.interval})`
);

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge 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 👍 / 👎.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Comment on lines +89 to +93
(this.client as unknown as EventEmitter)
.on('error', this.#onError)
.on('client-error', this.#onClientError)
.on('ready', this.#onReady)
.on('end', this.#onEnd);

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge 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 👍 / 👎.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Comment on lines +362 to +364
const active = mgr.activeDatabase;
// eslint-disable-next-line @typescript-eslint/no-explicit-any -- dynamic
const result = (resolve(active.client) as any)[name](...args);

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1 Badge 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 👍 / 👎.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Comment thread packages/client/lib/multi-db/manager.ts Outdated
);
}

const target = this.#select(healthy);

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge 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 👍 / 👎.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Comment on lines +89 to +93
(this.client as unknown as EventEmitter)
.on('error', this.#onError)
.on('client-error', this.#onClientError)
.on('ready', this.#onReady)
.on('end', this.#onEnd);

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge 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 👍 / 👎.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

@cursor cursor Bot left a comment •

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Stale Bugbot comment from a previous run.

Comment thread packages/client/lib/multi-db/manager.ts
Comment thread packages/client/lib/multi-db/manager.ts
Comment thread packages/client/lib/multi-db/index.ts

@chatgpt-codex-connector chatgpt-codex-connector Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

💡 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".

Comment on lines +46 to +49
const INTERCEPTED = new Set<PropertyKey>([
'connect', 'close', 'destroy', 'quit',
'withTypeMapping', 'withCommandOptions', 'withAbortSignal', 'asap',
'multi', 'MULTI', 'duplicate'

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1 Badge Fan out every teardown alias

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 👍 / 👎.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Comment thread packages/client/lib/multi-db/events.ts Outdated
Comment on lines +143 to +146
export type MultiDbClientType<C extends AnyRedisClientType> =
Omit<C, 'duplicate'> &
MultiDbEventEmitter &
{

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1 Badge 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 👍 / 👎.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Comment on lines +214 to +216
if (method === 'exec' && mgr.watchDirty) {
mgr.clearWatchSession();
return Promise.reject(new WatchError('MultiDb: the active database changed after WATCH'));

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge 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 👍 / 👎.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Comment thread packages/client/lib/multi-db/manager.ts Outdated
Comment on lines +275 to +278
this.#events?.emit('failover', { from: from.id, to: target.id, reason });
}

this.#afterSwitch(from, target).catch(err => this.#emitError(err as Error));

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge 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 👍 / 👎.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Comment on lines +53 to +55
export function probeRoundBudget(options: ProbeRoundOptions): number {
return options.numProbes * options.timeout
+ (options.numProbes - 1) * options.delayBetweenProbes;

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge 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 👍 / 👎.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Comment thread packages/client/lib/multi-db/manager.ts Outdated
Comment on lines +552 to +555
const results = await Promise.all(
this.#databases.map(db => this.#establishMember(db, false))
);
const healthy = this.#databases.filter((_, index) => results[index]);

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge 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 👍 / 👎.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Comment on lines +102 to +104
(this.client as unknown as EventEmitter)
.on('error', this.#onError)
.on('client-error', this.#onClientError)

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge 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 👍 / 👎.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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);

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge 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 👍 / 👎.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Comment on lines +691 to +692
const added = await this.addDatabase(config);
await this.removeDatabase(id);

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1 Badge 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 👍 / 👎.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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).

Comment on lines +578 to +582
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])

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1 Badge 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 👍 / 👎.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

@cursor cursor Bot left a comment •

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Stale Bugbot comment from a previous run.

Comment thread packages/client/lib/multi-db/manager.ts

@chatgpt-codex-connector chatgpt-codex-connector Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

💡 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();

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge 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 👍 / 👎.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Comment on lines +537 to +539
const active = mgr.activeDatabase;
// eslint-disable-next-line @typescript-eslint/no-explicit-any -- dynamic
const result = (resolve(active.client) as any)[name](...args);

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge 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 👍 / 👎.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Comment thread packages/client/lib/multi-db/manager.ts Outdated
Comment on lines +629 to +632
if (healthy.length < required) {
// a rejected connect() must not leave live sockets or retry timers behind
this.destroy();
throw new Error(

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge 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 👍 / 👎.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Comment on lines +32 to +34
const reply = await target.sendCommand(['PING']);
// toString covers both string and Buffer type mappings
return reply?.toString() === 'PONG';

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1 Badge 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 👍 / 👎.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Comment thread packages/client/lib/multi-db/index.ts Outdated
Comment on lines +617 to +621
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

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge 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 👍 / 👎.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Comment on lines +188 to +194
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

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge 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 👍 / 👎.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Comment on lines +486 to +488
// 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 });

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge 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 👍 / 👎.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Comment on lines +157 to +161
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 ?? {

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge 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 👍 / 👎.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Comment on lines +350 to +355
function makeDerived<C extends RedisClientLike>(
mgr: MultiDbManager<C>,
resolve: ResolveClient<C>
): C {
const view = new MultiDbClientBase(mgr);
attachForwarders(view, mgr, resolve);

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge 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 👍 / 👎.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

@cursor cursor Bot left a comment •

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Stale Bugbot comment from a previous run.

Comment thread packages/client/lib/client/pub-sub.ts Outdated
Comment thread packages/client/lib/multi-db/circuit.ts
Comment thread packages/client/lib/multi-db/manager.ts Outdated
Comment thread packages/client/lib/sentinel/pub-sub-proxy.ts

@chatgpt-codex-connector chatgpt-codex-connector Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

💡 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".

Comment on lines +384 to +389
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);

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge 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 👍 / 👎.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Comment on lines +350 to +353
const DELEGATED_EMITTER_METHODS = [
'on', 'once', 'off', 'addListener', 'removeListener',
'prependListener', 'prependOnceListener'
] as const;

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge 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 👍 / 👎.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Comment on lines +157 to +160
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)

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge 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 👍 / 👎.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Comment thread packages/client/lib/client/pub-sub.ts Outdated
if (!pending.carried) removeListeners();
this.#updateIsActive();
},
reject: undefined

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge 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 👍 / 👎.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Comment on lines +696 to +704
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));
}

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge 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 👍 / 👎.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Comment thread packages/client/lib/multi-db/manager.ts Outdated
clearInterval(timer);
this.#healthTimers.delete(member);
}
this.#databases.splice(this.#databases.indexOf(member), 1);

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge 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 👍 / 👎.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Comment on lines +462 to +469
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 });

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge 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 👍 / 👎.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Comment on lines +735 to +738
// #establishMember closes the circuit once the member establishes
await this.#establishMember(member, member.skipInitialHealthCheck);
this.#startMemberChecks(member);
return member.id;

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge 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 👍 / 👎.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Confirmed. Fixed in fcf5a66: addDatabase rejects if teardown landed during establish (teardown-only recheck; a concurrent removeDatabase still resolves, and #startMemberChecks self-guards membership).

@cursor cursor Bot left a comment •

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Stale Bugbot comment from a previous run.

Comment thread packages/client/lib/multi-db/manager.ts
Comment thread packages/client/lib/cluster/cluster-slots.ts

@PavelPashov PavelPashov left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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;

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Suggested change
this.#clock = options.clock ?? Date.now;
this.#clock = options.clock ?? (() => performance.now());

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Comment thread packages/client/lib/multi-db/manager.ts Outdated
if (member.circuit.state === 'CLOSED') {
await member.client.close();
} else {
member.client.destroy();

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Suggested change
member.client.destroy();
await member.client.destroy();

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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().

Comment thread packages/client/lib/multi-db/manager.ts Outdated
Comment on lines +432 to +434
if (this.#unavailable !== 'searching') {
return; // rescued mid-delay by a forced switch
}

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Comment on lines +235 to +236
const unavailable = mgr.unavailableError;
if (unavailable) return Promise.reject(unavailable);

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Suggested change
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'));
}

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Comment thread packages/client/lib/multi-db/manager.ts Outdated
Comment on lines +418 to +420
this.#failoverInFlight = true;
this.#unavailable = 'searching';
void this.#searchLoop(reason);

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Suggested change
this.#failoverInFlight = true;
this.#unavailable = 'searching';
void this.#searchLoop(reason);
this.#failoverInFlight = true;
this.#unavailable = 'searching';
this.#detector.reset();
void this.#searchLoop(reason);

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Comment thread packages/client/lib/client/pub-sub.ts Outdated
Comment on lines +510 to +530
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();
}

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Comment on lines +106 to +107
(this.client as unknown as EventEmitter)
.on('error', this.#onError)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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)) {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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:

  1. An automatic background check on database B starts.
  2. While it is still running, setActiveDatabase('B') runs a second, separate check. It passes, so traffic moves to B and B is pinned.
  3. The first check then finishes and reports that B failed.
  4. The client acts on that result. It moves traffic away from B and clears the pin, which undoes the forced switch.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Comment thread packages/client/lib/multi-db/index.ts Outdated
// 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; };

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Suggested change
dst.removeAllListeners = (event?: string) => { eventRoot.removeAllListeners(event); return dst; };
dst.removeAllListeners = (...args: [string?]) => { eventRoot.removeAllListeners(...args); return dst; };

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

@cursor cursor Bot left a comment •

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Stale Bugbot comment from a previous run.

Comment thread packages/client/lib/multi-db/manager.ts
Comment thread packages/client/lib/multi-db/manager.ts
Comment thread packages/client/lib/multi-db/manager.ts
Comment thread packages/client/lib/multi-db/manager.ts
Comment thread packages/client/lib/client/pub-sub.ts
Comment thread packages/client/lib/multi-db/database.ts
Comment thread packages/client/lib/multi-db/index.ts Outdated
Comment thread packages/client/lib/multi-db/index.ts
@nkaradzhov

Copy link
Copy Markdown
Collaborator Author

Agreed, fixed in f55247d400. The default errorFilter is now defaultErrorFilter (exported, so you can compose it). It counts transport errors, timeouts and server-state replies (LOADING, BUSY, MASTERDOWN, CLUSTERDOWN, READONLY, NOREPLICAS, MISCONF). It ignores other error replies (WRONGTYPE, CROSSSLOT, TRYAGAIN, OOM, ...), WatchError and AbortError. Ignored results still count as traffic. A MULTI error counts if any of its replies counts.

@cursor cursor Bot left a comment •

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Stale Bugbot comment from a previous run.

Comment thread packages/client/lib/client/pub-sub.ts

@cursor cursor Bot left a comment •

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Stale Bugbot comment from a previous run.

Comment thread packages/client/lib/multi-db/index.ts
Comment thread packages/client/lib/client/pub-sub.ts
nkaradzhov and others added 11 commits October 2, 2026 15:19
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>
nkaradzhov and others added 24 commits October 2, 2026 15:19
…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.

@cursor cursor Bot left a comment •

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Stale Bugbot comment from a previous run.

// 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
};

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Fix in Cursor Fix in Web

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.

@cursor cursor Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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).

Fix All in Cursor

Reviewed by Cursor Bugbot for commit 9ffbe81. Configure here.

this.#unavailable = null;
this.#failoverInFlight = false;
this.switchTo(target, reason);
return;

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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)
Fix in Cursor Fix in Web

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
]);
}

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Fix in Cursor Fix in Web

Reviewed by Cursor Bugbot for commit 9ffbe81. Configure here.

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants