Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
22 changes: 17 additions & 5 deletions src/connection.js
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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()
Expand Down Expand Up @@ -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
}
Expand Down
114 changes: 114 additions & 0 deletions tests/index.js
Original file line number Diff line number Diff line change
Expand Up @@ -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]
)
Expand Down