From 33d2f24fcf7f7bfffc5375c446d2a23b1bf3d842 Mon Sep 17 00:00:00 2001 From: Pratik Dulal Date: Sun, 9 Aug 2026 11:48:35 +0545 Subject: [PATCH 1/3] fix(pg-cloudflare): safely end closed sockets --- packages/pg-cloudflare/src/index.ts | 2 +- packages/pg-esm-test/pg-cloudflare.test.js | 12 ++++++++++++ 2 files changed, 13 insertions(+), 1 deletion(-) diff --git a/packages/pg-cloudflare/src/index.ts b/packages/pg-cloudflare/src/index.ts index 9b1e517ba..357131911 100644 --- a/packages/pg-cloudflare/src/index.ts +++ b/packages/pg-cloudflare/src/index.ts @@ -107,7 +107,7 @@ export class CloudflareSocket extends EventEmitter { end(data = Buffer.alloc(0), encoding: BufferEncoding = 'utf8', callback: (...args: unknown[]) => void = () => {}) { log('ending CF socket') this.write(data, encoding, (err) => { - this._cfSocket!.close() + this._cfSocket?.close() if (callback) callback(err) }) return this diff --git a/packages/pg-esm-test/pg-cloudflare.test.js b/packages/pg-esm-test/pg-cloudflare.test.js index a42620253..c140f0bb1 100644 --- a/packages/pg-esm-test/pg-cloudflare.test.js +++ b/packages/pg-esm-test/pg-cloudflare.test.js @@ -6,4 +6,16 @@ describe('pg-cloudflare', () => { it('should export CloudflareSocket constructor', () => { assert.ok(new CloudflareSocket()) }) + + it('should safely end after the underlying socket has closed', async () => { + const socket = new CloudflareSocket() + const underlyingSocket = { closed: Promise.resolve() } + socket._cfSocket = underlyingSocket + socket._addClosedHandler() + + await underlyingSocket.closed + assert.equal(socket._cfSocket, null) + + assert.doesNotThrow(() => socket.end()) + }) }) From e747cc4c122c0c1d7914ef105ad99527763ab2bb Mon Sep 17 00:00:00 2001 From: Pratik Dulal Date: Mon, 17 Aug 2026 14:08:17 +0545 Subject: [PATCH 2/3] fix(pg-cloudflare): complete socket shutdown handling --- packages/pg-cloudflare/src/index.ts | 52 ++++++++++++++++------ packages/pg-esm-test/pg-cloudflare.test.js | 32 +++++++++++++ 2 files changed, 71 insertions(+), 13 deletions(-) diff --git a/packages/pg-cloudflare/src/index.ts b/packages/pg-cloudflare/src/index.ts index 357131911..98e97139d 100644 --- a/packages/pg-cloudflare/src/index.ts +++ b/packages/pg-cloudflare/src/index.ts @@ -84,9 +84,13 @@ export class CloudflareSocket extends EventEmitter { write( data: Uint8Array | string, - encoding: BufferEncoding = 'utf8', + encoding: BufferEncoding | ((...args: unknown[]) => void) = 'utf8', callback: (...args: unknown[]) => void = () => {} ) { + if (typeof encoding === 'function') { + callback = encoding + encoding = 'utf8' + } if (data.length === 0) return callback() if (typeof data === 'string') data = Buffer.from(data, encoding) @@ -104,10 +108,27 @@ export class CloudflareSocket extends EventEmitter { return true } - end(data = Buffer.alloc(0), encoding: BufferEncoding = 'utf8', callback: (...args: unknown[]) => void = () => {}) { + end( + data = Buffer.alloc(0), + encoding: BufferEncoding | ((...args: unknown[]) => void) = 'utf8', + callback: (...args: unknown[]) => void = () => {} + ) { + if (typeof encoding === 'function') { + callback = encoding + encoding = 'utf8' + } log('ending CF socket') this.write(data, encoding, (err) => { - this._cfSocket?.close() + const socket = this._cfSocket + const closePromise = socket?.close() + closePromise + ?.then(() => { + if (this._cfSocket === socket) { + this._cfSocket = null + this.emit('close') + } + }) + .catch((e) => this.emit('error', e)) if (callback) callback(err) }) return this @@ -136,16 +157,21 @@ export class CloudflareSocket extends EventEmitter { } _addClosedHandler() { - this._cfSocket!.closed.then(() => { - if (!this._upgrading) { - log('CF socket closed') - this._cfSocket = null - this.emit('close') - } else { - this._upgrading = false - this._upgraded = true - } - }).catch((e) => this.emit('error', e)) + const socket = this._cfSocket! + socket.closed + .then(() => { + if (!this._upgrading) { + log('CF socket closed') + if (this._cfSocket === socket) { + this._cfSocket = null + this.emit('close') + } + } else { + this._upgrading = false + this._upgraded = true + } + }) + .catch((e) => this.emit('error', e)) } } diff --git a/packages/pg-esm-test/pg-cloudflare.test.js b/packages/pg-esm-test/pg-cloudflare.test.js index c140f0bb1..48960de01 100644 --- a/packages/pg-esm-test/pg-cloudflare.test.js +++ b/packages/pg-esm-test/pg-cloudflare.test.js @@ -18,4 +18,36 @@ describe('pg-cloudflare', () => { assert.doesNotThrow(() => socket.end()) }) + + it('should invoke callbacks passed as the write encoding argument', async () => { + const socket = new CloudflareSocket() + socket._cfSocket = { close() {} } + socket._cfWriter = { write: async () => {} } + + await new Promise((resolve, reject) => { + const timer = setTimeout(() => reject(new Error('write callback was not called')), 100) + socket.write(Buffer.from('terminate'), () => { + clearTimeout(timer) + resolve() + }) + }) + }) + + it('should emit close when ending a socket whose closed promise never settles', async () => { + const socket = new CloudflareSocket() + socket._cfSocket = { + closed: new Promise(() => {}), + close: async () => {}, + } + socket._addClosedHandler() + + await new Promise((resolve, reject) => { + const timer = setTimeout(() => reject(new Error('close event was not emitted')), 100) + socket.once('close', () => { + clearTimeout(timer) + resolve() + }) + socket.end() + }) + }) }) From 9c7a2406594cb6cd8e8316b5291bcd3a325c4470 Mon Sep 17 00:00:00 2001 From: Pratik Dulal Date: Mon, 17 Aug 2026 14:12:28 +0545 Subject: [PATCH 3/3] test(pg-cloudflare): retain write callback coverage --- packages/pg-esm-test/pg-cloudflare.test.js | 19 +++++++++++++++++++ 1 file changed, 19 insertions(+) diff --git a/packages/pg-esm-test/pg-cloudflare.test.js b/packages/pg-esm-test/pg-cloudflare.test.js index c99323956..683229158 100644 --- a/packages/pg-esm-test/pg-cloudflare.test.js +++ b/packages/pg-esm-test/pg-cloudflare.test.js @@ -19,6 +19,25 @@ describe('pg-cloudflare', () => { assert.doesNotThrow(() => socket.end()) }) + it('should call the write(data, callback) callback exactly once', async () => { + const socket = new CloudflareSocket() + socket._cfWriter = { write: () => Promise.resolve() } + + let resolve + const promise = new Promise((resolvePromise) => { + resolve = resolvePromise + }) + let called = false + socket.write(Buffer.from('x'), (error) => { + assert.ifError(error) + assert(!called) + called = true + resolve() + }) + + await promise + }) + it('should emit close when ending a socket whose closed promise never settles', async () => { const socket = new CloudflareSocket() socket._cfSocket = {