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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 3 additions & 1 deletion lib/Redis.ts
Original file line number Diff line number Diff line change
Expand Up @@ -381,7 +381,7 @@
* and may lose some pending replies that haven't written to client.
* If you want to wait for the pending replies, use Redis#quit instead.
*/
disconnect(reconnect = false) {

Check warning on line 384 in lib/Redis.ts

View workflow job for this annotation

GitHub Actions / test / test (24.x, 8.8.0)

Member disconnect should be declared before all private instance method definitions
if (!reconnect) {
this.manuallyClosing = true;
}
Expand All @@ -401,7 +401,7 @@
*
* @deprecated
*/
end() {

Check warning on line 404 in lib/Redis.ts

View workflow job for this annotation

GitHub Actions / test / test (24.x, 8.8.0)

Member end should be declared before all private instance method definitions
this.disconnect();
}

Expand All @@ -414,7 +414,7 @@
* var anotherRedis = redis.duplicate();
* ```
*/
duplicate<Override extends Partial<RedisOptions> | undefined = undefined>(

Check warning on line 417 in lib/Redis.ts

View workflow job for this annotation

GitHub Actions / test / test (24.x, 8.8.0)

Member duplicate should be declared before all private instance method definitions
override?: Override
): Redis<ReplyMappingFromOptions<ReplyMapping, Override>> {
return new Redis({
Expand Down Expand Up @@ -463,7 +463,7 @@
* });
* ```
*/
monitor(callback?: Callback<Redis>): Promise<Redis> {

Check warning on line 466 in lib/Redis.ts

View workflow job for this annotation

GitHub Actions / test / test (24.x, 8.8.0)

Member monitor should be declared before all private instance method definitions
const monitorInstance = this.duplicate({
monitor: true,
lazyConnect: false,
Expand Down Expand Up @@ -498,7 +498,7 @@
*
* @ignore
*/
sendCommand(command: Command, stream?: WriteableStream): unknown {

Check warning on line 501 in lib/Redis.ts

View workflow job for this annotation

GitHub Actions / test / test (24.x, 8.8.0)

Member sendCommand should be declared before all private instance method definitions
command.setReplyContext(this.condition ?? this.options);

if (this.status === "wait") {
Expand Down Expand Up @@ -741,23 +741,23 @@
});
}

scanStream(options?: ScanStreamOptions) {

Check warning on line 744 in lib/Redis.ts

View workflow job for this annotation

GitHub Actions / test / test (24.x, 8.8.0)

Member scanStream should be declared before all private instance method definitions
return this.createScanStream("scan", { options });
}

scanBufferStream(options?: ScanStreamOptions) {

Check warning on line 748 in lib/Redis.ts

View workflow job for this annotation

GitHub Actions / test / test (24.x, 8.8.0)

Member scanBufferStream should be declared before all private instance method definitions
return this.createScanStream("scanBuffer", { options });
}

sscanStream(key: string, options?: ScanStreamOptions) {

Check warning on line 752 in lib/Redis.ts

View workflow job for this annotation

GitHub Actions / test / test (24.x, 8.8.0)

Member sscanStream should be declared before all private instance method definitions
return this.createScanStream("sscan", { key, options });
}

sscanBufferStream(key: string, options?: ScanStreamOptions) {

Check warning on line 756 in lib/Redis.ts

View workflow job for this annotation

GitHub Actions / test / test (24.x, 8.8.0)

Member sscanBufferStream should be declared before all private instance method definitions
return this.createScanStream("sscanBuffer", { key, options });
}

hscanStream(key: string, options?: ScanStreamOptions) {

Check warning on line 760 in lib/Redis.ts

View workflow job for this annotation

GitHub Actions / test / test (24.x, 8.8.0)

Member hscanStream should be declared before all private instance method definitions
return this.createScanStream("hscan", { key, options });
}

Expand Down Expand Up @@ -853,7 +853,9 @@
this.condition?.select !== item.select &&
item.command.name !== "select"
) {
this.select(item.select);
this.select(item.select).catch((err) =>
this.silentEmit("error", err)
);
}
// TODO
// @ts-expect-error
Expand Down
83 changes: 83 additions & 0 deletions test/unit/reconnectOnError.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,83 @@
import { expect } from "chai";
import MockServer from "../helpers/mock_server";
import Redis from "../../lib/Redis";

const READONLY_ERROR = "READONLY You can't write against a read only replica.";
const INVALID_DB_INDEX = "ERR DB index is out of range";

// The restoring `SELECT` is protocol-independent, but the connection setup
// around it is not: under RESP3 the client sends HELLO before anything else.
const PROTOCOLS = [3, 2] as const;

PROTOCOLS.forEach((protocol, index) => {
describe(`reconnectOnError db restoration (RESP${protocol})`, () => {
const port = 17930 + index;
let savedListeners: any[] = [];
let unhandled: string[] = [];

beforeEach(() => {
unhandled = [];
// test/helpers/global.ts installs an unhandledRejection listener that
// throws; record rejections here instead and restore it afterwards.
savedListeners = process.listeners("unhandledRejection");
process.removeAllListeners("unhandledRejection");
process.on("unhandledRejection", (reason) => {
unhandled.push(String(reason));
});
});

afterEach(() => {
process.removeAllListeners("unhandledRejection");
for (const listener of savedListeners) {
process.on("unhandledRejection", listener);
}
});

it("surfaces a failing db-restoring SELECT as an error event", async () => {
// `get` fails on the first connection only, so reconnectOnError fires
// once; `select` fails from the reconnect onwards.
let connections = 0;
const server = new MockServer(port, (argv) => {
const name = String(argv[0]).toLowerCase();
if (name === "info") {
return "# Server\r\nredis_version:7.0.0\r\n";
}
if (name === "get" && connections < 2) {
return new Error(READONLY_ERROR);
}
if (name === "select" && connections >= 2) {
return new Error(INVALID_DB_INDEX);
}
return "OK";
});
server.on("connect", () => connections++);

const redis = new Redis({
port,
protocol,
lazyConnect: true,
retryStrategy: () => 40,
// 2 = reconnect and resend, the branch that restores the command's db.
reconnectOnError: (err: Error) =>
err.message.startsWith("READONLY") ? 2 : false,
});
const errors: string[] = [];
redis.on("error", (err: Error) => errors.push(err.message));
await redis.connect();

// `select(2)` is sent before the `get` reply arrives, so the failed
// `get` (issued against db 0) needs its db restored before the resend.
await Promise.all([
redis.get("foo").catch(() => {}),
redis.select(2).catch(() => {}),
]);
await new Promise((resolve) => setTimeout(resolve, 500));

redis.disconnect();
await server.disconnectPromise();

expect(unhandled).to.eql([]);
expect(errors).to.include(INVALID_DB_INDEX);
});
});
});
Loading