From 52ff6fcf1434e1476e006eff983c656bf58b3657 Mon Sep 17 00:00:00 2001 From: "Seongho.Bak" Date: Mon, 3 Aug 2026 14:03:46 +0900 Subject: [PATCH] fix: reset protocol parser after malformed messages --- .changeset/calm-parsers-recover.md | 5 ++++ packages/pg-protocol/src/parser.ts | 28 +++++++++++++------ .../pg-protocol/test/inbound-parser.test.ts | 17 +++++++++++ 3 files changed, 41 insertions(+), 9 deletions(-) create mode 100644 .changeset/calm-parsers-recover.md diff --git a/.changeset/calm-parsers-recover.md b/.changeset/calm-parsers-recover.md new file mode 100644 index 000000000..1bac0fe0d --- /dev/null +++ b/.changeset/calm-parsers-recover.md @@ -0,0 +1,5 @@ +--- +'@electric-sql/pglite': patch +--- + +Reset retained protocol parser state after a malformed backend message so later queries can recover. diff --git a/packages/pg-protocol/src/parser.ts b/packages/pg-protocol/src/parser.ts index 2582e3923..c09f89378 100644 --- a/packages/pg-protocol/src/parser.ts +++ b/packages/pg-protocol/src/parser.ts @@ -99,12 +99,18 @@ export class Parser { const length = this.#bufferView.getUint32(offset + CODE_LENGTH, false) const fullMessageLength = CODE_LENGTH + length if (fullMessageLength + offset <= bufferFullLength && length > 0) { - const message = this.#handlePacket( - offset + HEADER_LENGTH, - code, - length, - this.#bufferView.buffer, - ) + let message: BackendMessage + try { + message = this.#handlePacket( + offset + HEADER_LENGTH, + code, + length, + this.#bufferView.buffer, + ) + } catch (error) { + this.#resetBuffer() + throw error + } callback(message) offset += fullMessageLength } else { @@ -113,9 +119,7 @@ export class Parser { } if (offset === bufferFullLength) { // No more use for the buffer - this.#bufferView = new DataView(emptyBuffer) - this.#bufferRemainingLength = 0 - this.#bufferOffset = 0 + this.#resetBuffer() } else { // Adjust the cursors of remainingBuffer this.#bufferRemainingLength = bufferFullLength - offset @@ -123,6 +127,12 @@ export class Parser { } } + #resetBuffer(): void { + this.#bufferView = new DataView(emptyBuffer) + this.#bufferRemainingLength = 0 + this.#bufferOffset = 0 + } + #mergeBuffer(buffer: ArrayBuffer): void { if (this.#bufferRemainingLength > 0) { const newLength = this.#bufferRemainingLength + buffer.byteLength diff --git a/packages/pg-protocol/test/inbound-parser.test.ts b/packages/pg-protocol/test/inbound-parser.test.ts index 27f51802e..8961ea4eb 100644 --- a/packages/pg-protocol/test/inbound-parser.test.ts +++ b/packages/pg-protocol/test/inbound-parser.test.ts @@ -341,6 +341,23 @@ describe('PgPacketStream', () => { length: oneFieldBuf.byteLength - 1, }) }) + + it('recovers after a malformed data row throws', () => { + const malformedDataRow = new Uint8Array(11) + const view = new DataView(malformedDataRow.buffer) + malformedDataRow[0] = 0x44 + view.setUint32(1, 10, false) + view.setInt16(5, 2, false) + view.setInt32(7, 0x40000000, false) + + const parser = new Parser() + expect(() => parser.parse(malformedDataRow, () => {})).toThrow(RangeError) + + const messages: BackendMessage[] = [] + parser.parse(readyForQueryBuffer, (message) => messages.push(message)) + + expect(messages).toEqual([expectedReadyForQueryMessage]) + }) }) describe('notice message', () => {