diff --git a/cjs/src/connection.js b/cjs/src/connection.js index 07f6716702ac888215c86ad6e0071aaf55ae519c..dc24963ce2dd0a219f366f80ba9847548ba23a52 100644 --- a/cjs/src/connection.js +++ b/cjs/src/connection.js @@ -438,6 +438,7 @@ function Connection(options, queues = {}, { onopen = noop, onend = noop, onclose remaining = 0 incomings = null clearImmediate(nextWriteTimer) + chunk = nextWriteTimer = null socket.removeListener('data', data) socket.removeListener('connect', connected) idleTimer.cancel() @@ -451,6 +452,9 @@ function Connection(options, queues = {}, { onopen = noop, onend = noop, onclose return reconnect() !hadError && (query || sent.length) && error(Errors.connection('CONNECTION_CLOSED', options, socket)) + query = results = errorResponse = null + result = new Result() + rows = 0 closedTime = performance.now() hadError && options.shared.retries++ delay = (typeof backoff === 'function' ? backoff(options.shared.retries) : backoff) * 1000 diff --git a/cjs/src/index.js b/cjs/src/index.js index f09c61c7e998f411560b95c9a8e1ed6a97972c11..04dded5c7e2cdc34ac01b29a87cc93248ec8c9c6 100644 --- a/cjs/src/index.js +++ b/cjs/src/index.js @@ -216,8 +216,19 @@ function Postgres(a, b) { : move(c, reserved) c.reserved.release = true + let closedError + c.onclose = error => { + closedError = error + while (queue.length) + queue.shift().reject(error) + } + const sql = Sql(handler) sql.release = () => { + if (closedError) + return + closedError = Errors.connection('CONNECTION_RELEASED', options) + c.onclose = null c.reserved = null onopen(c) } @@ -225,6 +236,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) @@ -236,13 +249,19 @@ function Postgres(a, b) { const queries = Queue() let savepoints = 0 , connection + , closedError , prepare = null try { await sql.unsafe('begin ' + options.replace(/[^a-z ]/ig, ''), [], { onexecute }).execute() return await Promise.race([ scope(connection, fn), - new Promise((_, reject) => connection.onclose = reject) + new Promise((_, reject) => connection.onclose = error => { + closedError = error + while (queries.length) + queries.shift().reject(error) + reject(error) + }) ]) } catch (error) { throw error @@ -290,6 +309,8 @@ function Postgres(a, b) { function handler(q) { q.catch(e => uncaughtError || (uncaughtError = e)) + if (closedError) + return q.reject(closedError) c.queue === full ? queries.push(q) : c.execute(q) || move(c, full) diff --git a/src/connection.js b/src/connection.js index 1b1cccde43b4570d5d071d6ffaab2669b3c2065a..c95c326a2747c680c2c22b04b1c5f2e27a4e8c90 100644 --- a/src/connection.js +++ b/src/connection.js @@ -438,6 +438,7 @@ function Connection(options, queues = {}, { onopen = noop, onend = noop, onclose remaining = 0 incomings = null clearImmediate(nextWriteTimer) + chunk = nextWriteTimer = null socket.removeListener('data', data) socket.removeListener('connect', connected) idleTimer.cancel() @@ -451,6 +452,9 @@ function Connection(options, queues = {}, { onopen = noop, onend = noop, onclose return reconnect() !hadError && (query || sent.length) && error(Errors.connection('CONNECTION_CLOSED', options, socket)) + query = results = errorResponse = null + result = new Result() + rows = 0 closedTime = performance.now() hadError && options.shared.retries++ delay = (typeof backoff === 'function' ? backoff(options.shared.retries) : backoff) * 1000 diff --git a/src/index.js b/src/index.js index c7fba3dace20173f4cb13f3f138bb6d213b43683..f0888986add8ba02c8853f792ae42dd1cec88c82 100644 --- a/src/index.js +++ b/src/index.js @@ -216,8 +216,19 @@ function Postgres(a, b) { : move(c, reserved) c.reserved.release = true + let closedError + c.onclose = error => { + closedError = error + while (queue.length) + queue.shift().reject(error) + } + const sql = Sql(handler) sql.release = () => { + if (closedError) + return + closedError = Errors.connection('CONNECTION_RELEASED', options) + c.onclose = null c.reserved = null onopen(c) } @@ -225,6 +236,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) @@ -236,13 +249,19 @@ function Postgres(a, b) { const queries = Queue() let savepoints = 0 , connection + , closedError , prepare = null try { await sql.unsafe('begin ' + options.replace(/[^a-z ]/ig, ''), [], { onexecute }).execute() return await Promise.race([ scope(connection, fn), - new Promise((_, reject) => connection.onclose = reject) + new Promise((_, reject) => connection.onclose = error => { + closedError = error + while (queries.length) + queries.shift().reject(error) + reject(error) + }) ]) } catch (error) { throw error @@ -290,6 +309,8 @@ function Postgres(a, b) { function handler(q) { q.catch(e => uncaughtError || (uncaughtError = e)) + if (closedError) + return q.reject(closedError) c.queue === full ? queries.push(q) : c.execute(q) || move(c, full)