From 81ac1071d6bc082a1d3a852d858f3ea826b1b646 Mon Sep 17 00:00:00 2001 From: zfaustk <4340287+zfaustk@users.noreply.github.com> Date: Thu, 13 Aug 2026 00:56:28 +0800 Subject: [PATCH 1/3] fix: close Cloudflare sockets after failed connections Emit the close event when the explicit close operation resolves and support the write(data, callback) overload used by Connection.end(). Fixes #3152 --- packages/pg-cloudflare/src/index.ts | 22 +++++++++--- packages/pg-esm-test/pg-cloudflare.test.js | 41 ++++++++++++++++++++++ 2 files changed, 58 insertions(+), 5 deletions(-) diff --git a/packages/pg-cloudflare/src/index.ts b/packages/pg-cloudflare/src/index.ts index 357131911..b0beca71d 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) @@ -107,7 +111,10 @@ 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() + .then(() => this._emitClose()) + .catch((e) => this.emit('error', e)) if (callback) callback(err) }) return this @@ -138,15 +145,20 @@ export class CloudflareSocket extends EventEmitter { _addClosedHandler() { this._cfSocket!.closed.then(() => { if (!this._upgrading) { - log('CF socket closed') - this._cfSocket = null - this.emit('close') + this._emitClose() } else { this._upgrading = false this._upgraded = true } }).catch((e) => this.emit('error', e)) } + + private _emitClose() { + if (!this._cfSocket) return + log('CF socket closed') + this._cfSocket = null + this.emit('close') + } } const debug = false diff --git a/packages/pg-esm-test/pg-cloudflare.test.js b/packages/pg-esm-test/pg-cloudflare.test.js index c140f0bb1..880d45b7d 100644 --- a/packages/pg-esm-test/pg-cloudflare.test.js +++ b/packages/pg-esm-test/pg-cloudflare.test.js @@ -18,4 +18,45 @@ describe('pg-cloudflare', () => { assert.doesNotThrow(() => socket.end()) }) + + it('should emit close when the underlying close call resolves', async () => { + const socket = new CloudflareSocket() + const underlyingSocket = { + closed: new Promise(() => {}), + close: () => Promise.resolve(), + } + socket._cfSocket = underlyingSocket + socket._addClosedHandler() + + const close = new Promise((resolve) => socket.once('close', resolve)) + socket.end() + + await close + assert.equal(socket._cfSocket, null) + }) + + it('should support the write(data, callback) overload', async () => { + const socket = new CloudflareSocket() + socket._cfWriter = { write: () => Promise.resolve() } + + await new Promise((resolve) => socket.write(Buffer.from('x'), resolve)) + }) + + it('should emit close only once when both close signals resolve', async () => { + const socket = new CloudflareSocket() + const underlyingSocket = { + closed: Promise.resolve(), + close: () => Promise.resolve(), + } + socket._cfSocket = underlyingSocket + socket._addClosedHandler() + + let closeCount = 0 + socket.on('close', () => closeCount++) + socket.end() + + await underlyingSocket.closed + await new Promise((resolve) => setImmediate(resolve)) + assert.equal(closeCount, 1) + }) }) From 4b8e0bda4838adeee4120eaec34ec5e127b57d70 Mon Sep 17 00:00:00 2001 From: zfaustk <4340287+zfaustk@users.noreply.github.com> Date: Thu, 13 Aug 2026 18:10:32 +0800 Subject: [PATCH 2/3] fix: support Cloudflare write callbacks --- packages/pg-cloudflare/src/index.ts | 30 ++++++---------- packages/pg-esm-test/pg-cloudflare.test.js | 41 ++++------------------ 2 files changed, 18 insertions(+), 53 deletions(-) diff --git a/packages/pg-cloudflare/src/index.ts b/packages/pg-cloudflare/src/index.ts index b0beca71d..d625beee0 100644 --- a/packages/pg-cloudflare/src/index.ts +++ b/packages/pg-cloudflare/src/index.ts @@ -82,15 +82,15 @@ export class CloudflareSocket extends EventEmitter { this.emit('data', Buffer.from(value)) } + write(data: Uint8Array | string, callback?: (error?: unknown) => void): true | void + write(data: Uint8Array | string, encoding?: BufferEncoding, callback?: (error?: unknown) => void): true | void write( data: Uint8Array | string, - encoding: BufferEncoding | ((...args: unknown[]) => void) = 'utf8', - callback: (...args: unknown[]) => void = () => {} - ) { - if (typeof encoding === 'function') { - callback = encoding - encoding = 'utf8' - } + encodingOrCallback: BufferEncoding | ((error?: unknown) => void) = 'utf8', + callback: (error?: unknown) => void = () => {} + ): true | void { + const encoding = typeof encodingOrCallback === 'function' ? 'utf8' : encodingOrCallback + if (typeof encodingOrCallback === 'function') callback = encodingOrCallback if (data.length === 0) return callback() if (typeof data === 'string') data = Buffer.from(data, encoding) @@ -111,10 +111,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() - .then(() => this._emitClose()) - .catch((e) => this.emit('error', e)) + this._cfSocket?.close() if (callback) callback(err) }) return this @@ -145,20 +142,15 @@ export class CloudflareSocket extends EventEmitter { _addClosedHandler() { this._cfSocket!.closed.then(() => { if (!this._upgrading) { - this._emitClose() + log('CF socket closed') + this._cfSocket = null + this.emit('close') } else { this._upgrading = false this._upgraded = true } }).catch((e) => this.emit('error', e)) } - - private _emitClose() { - if (!this._cfSocket) return - log('CF socket closed') - this._cfSocket = null - this.emit('close') - } } const debug = false diff --git a/packages/pg-esm-test/pg-cloudflare.test.js b/packages/pg-esm-test/pg-cloudflare.test.js index 880d45b7d..f62572c7e 100644 --- a/packages/pg-esm-test/pg-cloudflare.test.js +++ b/packages/pg-esm-test/pg-cloudflare.test.js @@ -19,44 +19,17 @@ describe('pg-cloudflare', () => { assert.doesNotThrow(() => socket.end()) }) - it('should emit close when the underlying close call resolves', async () => { - const socket = new CloudflareSocket() - const underlyingSocket = { - closed: new Promise(() => {}), - close: () => Promise.resolve(), - } - socket._cfSocket = underlyingSocket - socket._addClosedHandler() - - const close = new Promise((resolve) => socket.once('close', resolve)) - socket.end() - - await close - assert.equal(socket._cfSocket, null) - }) - - it('should support the write(data, callback) overload', async () => { + it('should call the write(data, callback) callback exactly once', async () => { const socket = new CloudflareSocket() socket._cfWriter = { write: () => Promise.resolve() } - await new Promise((resolve) => socket.write(Buffer.from('x'), resolve)) - }) + let callbackCount = 0 + socket.write(Buffer.from('x'), (error) => { + assert.ifError(error) + callbackCount++ + }) - it('should emit close only once when both close signals resolve', async () => { - const socket = new CloudflareSocket() - const underlyingSocket = { - closed: Promise.resolve(), - close: () => Promise.resolve(), - } - socket._cfSocket = underlyingSocket - socket._addClosedHandler() - - let closeCount = 0 - socket.on('close', () => closeCount++) - socket.end() - - await underlyingSocket.closed await new Promise((resolve) => setImmediate(resolve)) - assert.equal(closeCount, 1) + assert.equal(callbackCount, 1) }) }) From 952b113b5cd57554ea2d634a9ed2900650fbd827 Mon Sep 17 00:00:00 2001 From: zfaustk <4340287+zfaustk@users.noreply.github.com> Date: Sat, 15 Aug 2026 03:02:57 +0800 Subject: [PATCH 3/3] test(pg-cloudflare): await write callback deterministically ## Context - Principle: Callback tests must observe callback completion directly across the supported Node matrix - Why: The write overload regression should fail on a missing or duplicate callback without relying on event-loop timing ## Key Deltas - pg-cloudflare callback test: waited one event-loop tick and counted callbacks afterward -> awaits the callback and rejects a duplicate invocation inside it; why: the assertion now follows the callback contract while remaining compatible with Node 16. Key refs: packages/pg-esm-test/pg-cloudflare.test.js:21 ## Verification - Result: passed --- packages/pg-esm-test/pg-cloudflare.test.js | 13 +++++++++---- 1 file changed, 9 insertions(+), 4 deletions(-) diff --git a/packages/pg-esm-test/pg-cloudflare.test.js b/packages/pg-esm-test/pg-cloudflare.test.js index f62572c7e..75d1f9957 100644 --- a/packages/pg-esm-test/pg-cloudflare.test.js +++ b/packages/pg-esm-test/pg-cloudflare.test.js @@ -23,13 +23,18 @@ describe('pg-cloudflare', () => { const socket = new CloudflareSocket() socket._cfWriter = { write: () => Promise.resolve() } - let callbackCount = 0 + let resolve + const promise = new Promise((resolvePromise) => { + resolve = resolvePromise + }) + let called = false socket.write(Buffer.from('x'), (error) => { assert.ifError(error) - callbackCount++ + assert(!called) + called = true + resolve() }) - await new Promise((resolve) => setImmediate(resolve)) - assert.equal(callbackCount, 1) + await promise }) })