From b93bd2dde766d750434bfa3e6aa2da549a449712 Mon Sep 17 00:00:00 2001 From: Jake Thomson Date: Wed, 30 Sep 2026 12:06:51 +0100 Subject: [PATCH 1/2] Add tests for connections closing while fetching array types When a connection closes while its first query is parked as initial, closed() reconnects without clearing the dead backend's query and error response, and without any delay: - a backend terminated during the array type fetch replays its error onto the next connection: the parked query fails with it, and the internal fetch rejects with nothing handling it (#1223) - a peer that closes cleanly during startup makes the connection reconnect with zero delay forever, and the query never settles (#1193) - a peer that fails every array type fetch, as PgBouncer does during server_login_retry, must still settle the query with its own error within connect_timeout once the replay is gone - a failed array type fetch on a live connection is an unhandled rejection All four fail on master. --- tests/index.js | 114 +++++++++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 114 insertions(+) diff --git a/tests/index.js b/tests/index.js index 845c6ce..7904739 100644 --- a/tests/index.js +++ b/tests/index.js @@ -1710,6 +1710,120 @@ t('requests works after single connect_timeout', async() => { ] }) +const backendMessage = (type, body) => { + const header = Buffer.alloc(5) + header.write(type) + header.writeInt32BE(body.length + 4, 1) + return Buffer.concat([header, body]) +} + +const errorResponse = (severity, code) => backendMessage('E', Buffer.from( + 'S' + severity + '\0V' + severity + '\0C' + code + '\0M' + code + '\0\0' +)) + +const authenticated = Buffer.concat([ + backendMessage('R', Buffer.alloc(4)), + backendMessage('Z', Buffer.from('I')) +]) + +t('Query runs on a new connection if the backend dies while fetching array types', async() => { + let killed = false + , unhandled = 0 + + const onunhandled = () => unhandled++ + process.on('unhandledRejection', onunhandled) + + const proxy = net.createServer(client => { + const server = net.connect(5432, 'localhost') + client.on('data', x => { + if (!killed && x.includes('typarray')) { + killed = true + server.destroy() + return client.end(errorResponse('FATAL', '57P01')) + } + server.write(x) + }) + server.on('data', x => client.writable && client.write(x)) + client.on('error', () => server.destroy()) + server.on('error', () => client.destroy()) + client.on('close', () => server.destroy()) + }) + + await new Promise(r => proxy.listen(0, 'localhost', r)) + const sql = postgres({ ...options, host: 'localhost', port: proxy.address().port }) + const x = await sql`select 1 as x`.then(([x]) => x.x, e => e.code) + await sql.end() + proxy.close() + process.off('unhandledRejection', onunhandled) + + return ['1,0', [x, unhandled].join(',')] +}) + +t('Clean close during startup rejects after connect_timeout instead of reconnecting forever', async() => { + let attempts = 0 + + const server = net.createServer(socket => (attempts++, socket.end())) + await new Promise(r => server.listen(0, 'localhost', r)) + const sql = postgres({ ...options, host: 'localhost', port: server.address().port }) + const code = await sql`select 1`.catch(e => e.code) + await sql.end() + server.close() + + return ['CONNECTION_CLOSED,true', [code, attempts < 10].join(',')] +}) + +t('Failing every array type fetch rejects with the peer error instead of reconnecting forever', async() => { + let attempts = 0 + , unhandled = 0 + + const onunhandled = () => unhandled++ + process.on('unhandledRejection', onunhandled) + + const server = net.createServer(socket => { + attempts++ + socket.on('error', () => socket.destroy()) + socket.on('data', x => x.includes('typarray') + ? socket.end(errorResponse('FATAL', '08P01')) + : socket.write(authenticated) + ) + }) + + await new Promise(r => server.listen(0, 'localhost', r)) + const sql = postgres({ ...options, host: 'localhost', port: server.address().port }) + const code = await sql`select 1`.catch(e => e.code) + await sql.end() + server.close() + process.off('unhandledRejection', onunhandled) + + return ['08P01,true,0', [code, attempts < 10, unhandled].join(',')] +}) + +t('Failed array type fetch on a live connection is not an unhandled rejection', async() => { + let unhandled = 0 + + const onunhandled = () => unhandled++ + process.on('unhandledRejection', onunhandled) + + const server = net.createServer(socket => { + socket.on('error', () => socket.destroy()) + socket.on('data', x => x.includes('typarray') + ? socket.write(Buffer.concat([errorResponse('ERROR', '57014'), backendMessage('Z', Buffer.from('I'))])) + : x[0] === 88 + ? socket.end() + : socket.write(authenticated) + ) + }) + + await new Promise(r => server.listen(0, 'localhost', r)) + const sql = postgres({ ...options, host: 'localhost', port: server.address().port }) + await sql`select 1`.catch(() => null) + await sql.end() + server.close() + process.off('unhandledRejection', onunhandled) + + return [0, unhandled] +}) + t('Postgres errors are of type PostgresError', async() => [true, (await sql`bad keyword`.catch(e => e)) instanceof sql.PostgresError] ) From 81af2dc979be47a8c97e4c866935c8f05feb6f28 Mon Sep 17 00:00:00 2001 From: Jake Thomson Date: Wed, 30 Sep 2026 12:08:58 +0100 Subject: [PATCH 2/2] Settle startup queries when their connection keeps closing - fixes #1223, fixes #1193 closed() reconnected a connection with a parked initial query before clearing the dead backend's query and errorResponse, so the next connection's first ReadyForQuery replayed them (#1223). It also skipped the backoff, so a peer that kept closing spun zero-delay reconnects that neither connect_timeout nor max_lifetime could bound (#1193). Clear the stale state and pace these reconnects with the shared backoff. The run of consecutive closes is bounded by connect_timeout, after which the parked query rejects with the last error the peer sent, or with CONNECTION_CLOSED. Any successful connect resets the run, including a reserve() whose initial is dropped before the array type fetch. Clearing the state alone is not enough: the replay was what settled the query after a server error followed by a close, so without the bound that case becomes the #1193 loop. fetchArrayTypes() is now caught as well. If it fails on a live connection, errored() has already settled initial, and connected() re-arms needsTypes on the next connect. --- src/connection.js | 22 +++++++++++++++++----- 1 file changed, 17 insertions(+), 5 deletions(-) diff --git a/src/connection.js b/src/connection.js index 10ab1bb..5086eba 100644 --- a/src/connection.js +++ b/src/connection.js @@ -87,6 +87,7 @@ function Connection(options, queues = {}, { onopen = noop, onend = noop, onclose , statementId = Math.random().toString(36).slice(2) , statementCount = 1 , closedTime = 0 + , closeRunStart = 0 , remaining = 0 , hostIndex = 0 , retries = 0 @@ -447,8 +448,19 @@ function Connection(options, queues = {}, { onopen = noop, onend = noop, onclose socket.removeAllListeners() socket = null - if (initial) - return reconnect() + if (initial) { + const err = errorResponse + query = errorResponse = null + closedTime = performance.now() + closeRunStart || (closeRunStart = closedTime) + options.shared.retries++ + delay = (typeof backoff === 'function' ? backoff(options.shared.retries) : backoff) * 1000 + if (closedTime + delay <= closeRunStart + (options.connect_timeout || 30) * 1000) + return reconnect() + closeRunStart = 0 + errored(err || Errors.connection('CONNECTION_CLOSED', options, socket)) + return onclose(connection, Errors.connection('CONNECTION_CLOSED', options, socket)) + } !hadError && (query || sent.length) && error(Errors.connection('CONNECTION_CLOSED', options, socket)) closedTime = performance.now() @@ -560,12 +572,12 @@ function Connection(options, queues = {}, { onopen = noop, onend = noop, onclose } if (needsTypes) { - initial.reserve && (initial = null) - return fetchArrayTypes() + initial.reserve && (closeRunStart = 0, initial = null) + return fetchArrayTypes().catch(noop) } initial && !initial.reserve && execute(initial) - options.shared.retries = retries = 0 + options.shared.retries = retries = closeRunStart = 0 initial = null return }