From cf895ad3feb951efc414f6c3513716e35196544c Mon Sep 17 00:00:00 2001 From: Kesavamoorthig06 Date: Sun, 11 Oct 2026 09:14:30 +0530 Subject: [PATCH] fix(pg-cursor): settle every read when reads overlap Two overlapping read() calls corrupted the cursor state machine: the row description handler shifted the queued read while the first read was still busy, overwriting the active callback so the first read never settled, and the batch delivery path never drained the queue. Guard the queue shifts in handleRowDescription/_ifNoData against the busy state and drain the queue from _sendRows once the active read settles. --- packages/pg-cursor/index.js | 10 +++++-- packages/pg-cursor/test/concurrent-reads.js | 33 +++++++++++++++++++++ 2 files changed, 41 insertions(+), 2 deletions(-) create mode 100644 packages/pg-cursor/test/concurrent-reads.js diff --git a/packages/pg-cursor/index.js b/packages/pg-cursor/index.js index 5d566882b..cdf4aa3c7 100644 --- a/packages/pg-cursor/index.js +++ b/packages/pg-cursor/index.js @@ -27,7 +27,7 @@ class Cursor extends EventEmitter { } _ifNoData() { - if (this.state !== 'done' && this.state !== 'error') { + if (this.state !== 'done' && this.state !== 'error' && this.state !== 'busy') { this.state = 'idle' this._shiftQueue() } @@ -108,7 +108,7 @@ class Cursor extends EventEmitter { handleRowDescription(msg) { this._result.addFields(msg.fields) - if (this.state !== 'done' && this.state !== 'error') { + if (this.state !== 'done' && this.state !== 'error' && this.state !== 'busy') { this.state = 'idle' this._shiftQueue() } @@ -135,6 +135,12 @@ class Cursor extends EventEmitter { cb(null, this._rows, this._result) } this._rows = [] + // the active read just settled; start any read queued behind it. + // the callback above may have started a new read itself, so only + // shift when the cursor is idle again. + if (this.state === 'idle') { + this._shiftQueue() + } }) } diff --git a/packages/pg-cursor/test/concurrent-reads.js b/packages/pg-cursor/test/concurrent-reads.js new file mode 100644 index 000000000..69d40ad61 --- /dev/null +++ b/packages/pg-cursor/test/concurrent-reads.js @@ -0,0 +1,33 @@ +const assert = require('assert') +const Cursor = require('../') +const pg = require('pg') + +const text = 'SELECT generate_series AS num FROM generate_series(1, 5)' + +describe('concurrent reads', function () { + beforeEach(function (done) { + const client = (this.client = new pg.Client()) + client.connect(done) + }) + + afterEach(function () { + this.client.end() + }) + + it('settles every read when two reads overlap', async function () { + const cursor = this.client.query(new Cursor(text)) + const first = cursor.read(1) + const second = cursor.read(1) + const [rowsA, rowsB] = await Promise.all([first, second]) + assert.strictEqual(rowsA.length, 1) + assert.strictEqual(rowsB.length, 1) + assert.strictEqual(rowsA[0].num, 1) + assert.strictEqual(rowsB[0].num, 2) + const rest = await cursor.read(10) + assert.deepStrictEqual( + rest.map((row) => row.num), + [3, 4, 5] + ) + await cursor.close() + }) +})