Sitelet https://github.com/redis/node-redis/pull/3486/files
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
27 changes: 27 additions & 0 deletions packages/client/lib/RESP/decoder.spec.ts
Original file line number Diff line number Diff line change
Expand Up @@ -502,4 +502,31 @@ describe('RESP Decoder', () => {
pushReplies: [[0, 1, 2, 3, 4, 5, 6, 7, 8, 9]]
});
});

it('keeps decoding the chunk when a callback throws', () => {
// Regression: the error stopped write(), and the rest of the chunk was lost
const pushError = new Error('push listener error'),
nullError = new Error('reply listener error'),
onPush = spy(() => {
throw pushError;
}),
onReply = spy((reply: unknown) => {
if (reply === null) throw nullError;
}),
decoder = new Decoder({
getTypeMapping: () => ({}),
onReply,
onErrorReply: spy(),
onPush
});

assert.equal(decoder.write(Buffer.from('>2\r\n+mess')), undefined);
assert.deepEqual(
decoder.write(Buffer.from('age\r\n+1\r\n:1\r\n_\r\n>2\r\n+message\r\n+2\r\n$2\r\nv')),
[pushError, nullError, pushError]
);
assert.equal(decoder.write(Buffer.from('1\r\n')), undefined);
assert.deepEqual(onPush.args, [[['message', '1']], [['message', '2']]]);
assert.deepEqual(onReply.args, [[1], [null], ['v1']]);
});
});
19 changes: 16 additions & 3 deletions packages/client/lib/RESP/decoder.ts
Original file line number Diff line number Diff line change
Expand Up @@ -83,6 +83,7 @@ export class Decoder {
getTypeMapping;
#cursor = 0;
#next;
#callbackErrors?: Array<unknown>;

constructor(config: DecoderOptions) {
this.onReply = config.onReply;
Expand All @@ -97,6 +98,12 @@ export class Decoder {
}

write(chunk) {
this.#callbackErrors = undefined;
this.#write(chunk);
return this.#callbackErrors;
}

#write(chunk) {
if (this.#cursor >= chunk.length) {
this.#cursor -= chunk.length;
return;
Expand Down Expand Up @@ -131,8 +138,10 @@ export class Decoder {
#decodeTypeValue(type, chunk) {
switch (type) {
case RESP_TYPES.NULL:
this.onReply(this.#decodeNull());
return false;
return this.#handleDecodedValue(
this.onReply,
this.#decodeNull()
);

case RESP_TYPES.BOOLEAN:
return this.#handleDecodedValue(
Expand Down Expand Up @@ -241,7 +250,11 @@ export class Decoder {
return true;
}

cb(value);
try {
cb(value);
} catch (err) {
(this.#callbackErrors ??= []).push(err);
}
return false;
}

Expand Down
54 changes: 54 additions & 0 deletions packages/client/lib/client/index.spec.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1321,6 +1321,60 @@ describe('Client', () => {
}
}, GLOBAL.SERVERS.OPEN);

for (const resp of [2, 3] as const) {
testUtils.testWithClient(`should keep delivering messages after a listener throws (RESP${resp})`, async publisher => {
const subscriber = await publisher.duplicate().connect(),
listenerError = new Error('listener error'),
errors: Array<unknown> = [],
listener = spy((message: string) => {
if (message === '1') throw listenerError;
});
subscriber.on('error', err => errors.push(err));

try {
await subscriber.subscribe('channel', listener);
// published in one pipeline, the first three messages reach the subscriber in one chunk
await Promise.all(['1', '2', '3'].map(message => publisher.publish('channel', message)));
await publisher.publish('channel', '4');
await subscriber.ping();

assert.deepEqual(listener.args.map(([message]) => message), ['1', '2', '3', '4']);
assert.deepEqual(errors, [listenerError]);
} finally {
subscriber.destroy();
}
}, {
...GLOBAL.SERVERS.OPEN,
clientOptions: { ...GLOBAL.SERVERS.OPEN.clientOptions, RESP: resp }
});
}

testUtils.testWithClient('should keep replies matched to their commands when a listener throws', async client => {
await client.mSet({ k1: 'v1', k2: 'v2', k3: 'v3', k4: 'v4' });
const duplicate = await client.duplicate().connect(),
listenerError = new Error('listener error');

try {
await duplicate.subscribe('channel', () => {
throw listenerError;
});
// the message and the two replies after it arrive in one chunk
const errorEvent = once(duplicate, 'error'),
replies = Promise.all([
duplicate.publish('channel', 'message'),
duplicate.get('k1'),
duplicate.get('k2')
]);
assert.deepEqual(await errorEvent, [listenerError]);

const laterReplies = Promise.all([duplicate.get('k3'), duplicate.get('k4')]);
assert.deepEqual(await replies, [1, 'v1', 'v2']);
assert.deepEqual(await laterReplies, ['v3', 'v4']);
} finally {
duplicate.destroy();
}
}, GLOBAL.SERVERS.OPEN);

testUtils.testWithClient('should be able to quit in PubSub mode', async client => {
await client.subscribe('channel', () => {
// noop
Expand Down
4 changes: 3 additions & 1 deletion packages/client/lib/client/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1048,12 +1048,14 @@ export default class RedisClient<

#attachListeners(socket: RedisSocket) {
socket.on('data', chunk => {
let errors: Array<unknown> | undefined;
try {
this.#queue.decoder.write(chunk);
errors = this.#queue.decoder.write(chunk);
} catch (err) {
this.#queue.resetDecoder();
this.emit('error', err);
}
errors?.forEach(err => this.emit('error', err));
})
.on('error', err => {
this.emit('error', err);
Expand Down
Loading