diff --git a/src/index.js b/src/index.js index c7fba3d..03bf739 100644 --- a/src/index.js +++ b/src/index.js @@ -210,14 +210,23 @@ 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) } @@ -225,6 +234,8 @@ function Postgres(a, b) { return sql function handler(q) { + if (closedError) + return q.reject(closedError) c.queue === full ? queue.push(q) : c.execute(q) || move(c, full) @@ -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) diff --git a/tests/index.js b/tests/index.js index 845c6ce..9afdfe5 100644 --- a/tests/index.js +++ b/tests/index.js @@ -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 = []