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
12 changes: 12 additions & 0 deletions src/index.js
Original file line number Diff line number Diff line change
Expand Up @@ -210,21 +210,32 @@ function Postgres(a, b) {
closed.length && connect(closed.shift(), query)
})

let closedError = null

move(c, reserved)
c.reserved = () => queue.length
? c.execute(queue.shift())
: move(c, reserved)
c.reserved.release = true
c.reserved.closed = error => {
closedError = error
while (queue.length)
queue.shift().reject(error)
}

const sql = Sql(handler)
sql.release = () => {
if (closedError)
return
c.reserved = null
onopen(c)
}

return sql

function handler(q) {
if (closedError)
return q.reject(closedError)
c.queue === full
? queue.push(q)
: c.execute(q) || move(c, full)
Expand Down Expand Up @@ -420,6 +431,7 @@ function Postgres(a, b) {

function onclose(c, e) {
move(c, closed)
c.reserved && c.reserved.closed && c.reserved.closed(e)
c.reserved = null
c.onclose && (c.onclose(e), c.onclose = null)
options.onclose && options.onclose(c.id)
Expand Down
30 changes: 30 additions & 0 deletions tests/index.js
Original file line number Diff line number Diff line change
Expand Up @@ -889,6 +889,36 @@ t('listen reconnects', { timeout: 2 }, async() => {
return [connects, 2]
})

t('Reserved connection rejects queries after its backend is terminated', async() => {
const sql = postgres({ ...options, max: 2 })
const reserved = await sql.reserve()
const [{ pid }] = await reserved`select pg_backend_pid() as pid`
await sql`select pg_terminate_backend(${ pid })`
await delay(100)
const first = await reserved`select 1`.catch(e => e.code)
const second = await reserved`select 2`.catch(e => e.code)
reserved.release()
await sql.end()
return ['CONNECTION_CLOSED,CONNECTION_CLOSED', first + ',' + second]
})

t('Releasing a terminated reserved connection does not return it to the pool', async() => {
const sql = postgres({ ...options, max: 1 })
const killer = postgres(options)
const reserved = await sql.reserve()
const [{ pid }] = await reserved`select pg_backend_pid() as pid`
await killer`select pg_terminate_backend(${ pid })`
await delay(100)
await reserved`select 1`.catch(() => { /* rejects, connection closed */ })
reserved.release()
const [{ x }] = await Promise.race([
sql`select 1 as x`,
delay(2000).then(() => [{ x: 'hung' }])
])
await Promise.all([sql.end(), killer.end()])
return [1, x]
})

t('listen result reports correct connection state after reconnection', async() => {
const sql = postgres(options)
, xs = []
Expand Down