diff --git a/packages/client/lib/RESP/decoder.spec.ts b/packages/client/lib/RESP/decoder.spec.ts index 19377f90e78..e7b8cfa39d9 100644 --- a/packages/client/lib/RESP/decoder.spec.ts +++ b/packages/client/lib/RESP/decoder.spec.ts @@ -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']]); + }); }); diff --git a/packages/client/lib/RESP/decoder.ts b/packages/client/lib/RESP/decoder.ts index a834b4b3ed4..c15b85baad7 100644 --- a/packages/client/lib/RESP/decoder.ts +++ b/packages/client/lib/RESP/decoder.ts @@ -83,6 +83,7 @@ export class Decoder { getTypeMapping; #cursor = 0; #next; + #callbackErrors?: Array; constructor(config: DecoderOptions) { this.onReply = config.onReply; @@ -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; @@ -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( @@ -241,7 +250,11 @@ export class Decoder { return true; } - cb(value); + try { + cb(value); + } catch (err) { + (this.#callbackErrors ??= []).push(err); + } return false; } diff --git a/packages/client/lib/client/index.spec.ts b/packages/client/lib/client/index.spec.ts index 3a7b7b43341..ad34f26f779 100644 --- a/packages/client/lib/client/index.spec.ts +++ b/packages/client/lib/client/index.spec.ts @@ -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 = [], + 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 diff --git a/packages/client/lib/client/index.ts b/packages/client/lib/client/index.ts index 78d1ce0302d..46f4a0bc440 100644 --- a/packages/client/lib/client/index.ts +++ b/packages/client/lib/client/index.ts @@ -1048,12 +1048,14 @@ export default class RedisClient< #attachListeners(socket: RedisSocket) { socket.on('data', chunk => { + let errors: Array | 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);