| [62b2964] | 1 | const nodeUtils = require('util')
|
|---|
| 2 | // eslint-disable-next-line
|
|---|
| 3 | var Native
|
|---|
| 4 | // eslint-disable-next-line no-useless-catch
|
|---|
| 5 | try {
|
|---|
| 6 | // Wrap this `require()` in a try-catch to avoid upstream bundlers from complaining that this might not be available since it is an optional import
|
|---|
| 7 | Native = require('pg-native')
|
|---|
| 8 | } catch (e) {
|
|---|
| 9 | throw e
|
|---|
| 10 | }
|
|---|
| 11 | const TypeOverrides = require('../type-overrides')
|
|---|
| 12 | const EventEmitter = require('events').EventEmitter
|
|---|
| 13 | const util = require('util')
|
|---|
| 14 | const ConnectionParameters = require('../connection-parameters')
|
|---|
| 15 |
|
|---|
| 16 | const NativeQuery = require('./query')
|
|---|
| 17 |
|
|---|
| 18 | const queryQueueLengthDeprecationNotice = nodeUtils.deprecate(
|
|---|
| 19 | () => {},
|
|---|
| 20 | 'Calling client.query() when the client is already executing a query is deprecated and will be removed in pg@9.0. Use async/await or an external async flow control mechanism instead.'
|
|---|
| 21 | )
|
|---|
| 22 |
|
|---|
| 23 | const Client = (module.exports = function (config) {
|
|---|
| 24 | EventEmitter.call(this)
|
|---|
| 25 | config = config || {}
|
|---|
| 26 |
|
|---|
| 27 | this._Promise = config.Promise || global.Promise
|
|---|
| 28 | this._types = new TypeOverrides(config.types)
|
|---|
| 29 |
|
|---|
| 30 | this.native = new Native({
|
|---|
| 31 | types: this._types,
|
|---|
| 32 | })
|
|---|
| 33 |
|
|---|
| 34 | this._queryQueue = []
|
|---|
| 35 | this._ending = false
|
|---|
| 36 | this._connecting = false
|
|---|
| 37 | this._connected = false
|
|---|
| 38 | this._queryable = true
|
|---|
| 39 | this.pipeline = Boolean(config.pipeline)
|
|---|
| 40 | this._pipelineInFlight = false
|
|---|
| 41 |
|
|---|
| 42 | // keep these on the object for legacy reasons
|
|---|
| 43 | // for the time being. TODO: deprecate all this jazz
|
|---|
| 44 | const cp = (this.connectionParameters = new ConnectionParameters(config))
|
|---|
| 45 | if (config.nativeConnectionString) cp.nativeConnectionString = config.nativeConnectionString
|
|---|
| 46 | this.user = cp.user
|
|---|
| 47 |
|
|---|
| 48 | // "hiding" the password so it doesn't show up in stack traces
|
|---|
| 49 | // or if the client is console.logged
|
|---|
| 50 | Object.defineProperty(this, 'password', {
|
|---|
| 51 | configurable: true,
|
|---|
| 52 | enumerable: false,
|
|---|
| 53 | writable: true,
|
|---|
| 54 | value: cp.password,
|
|---|
| 55 | })
|
|---|
| 56 | this.database = cp.database
|
|---|
| 57 | this.host = cp.host
|
|---|
| 58 | this.port = cp.port
|
|---|
| 59 |
|
|---|
| 60 | // a hash to hold named queries
|
|---|
| 61 | this.namedQueries = {}
|
|---|
| 62 | })
|
|---|
| 63 |
|
|---|
| 64 | Client.Query = NativeQuery
|
|---|
| 65 |
|
|---|
| 66 | util.inherits(Client, EventEmitter)
|
|---|
| 67 |
|
|---|
| 68 | Client.prototype._errorAllQueries = function (err) {
|
|---|
| 69 | const enqueueError = (query) => {
|
|---|
| 70 | process.nextTick(() => {
|
|---|
| 71 | query.native = this.native
|
|---|
| 72 | query.handleError(err)
|
|---|
| 73 | })
|
|---|
| 74 | }
|
|---|
| 75 |
|
|---|
| 76 | if (this._hasActiveQuery()) {
|
|---|
| 77 | enqueueError(this._activeQuery)
|
|---|
| 78 | this._activeQuery = null
|
|---|
| 79 | }
|
|---|
| 80 |
|
|---|
| 81 | this._queryQueue.forEach(enqueueError)
|
|---|
| 82 | this._queryQueue.length = 0
|
|---|
| 83 | }
|
|---|
| 84 |
|
|---|
| 85 | // connect to the backend
|
|---|
| 86 | // pass an optional callback to be called once connected
|
|---|
| 87 | // or with an error if there was a connection error
|
|---|
| 88 | Client.prototype._connect = function (cb) {
|
|---|
| 89 | const self = this
|
|---|
| 90 |
|
|---|
| 91 | if (this._connecting) {
|
|---|
| 92 | process.nextTick(() => cb(new Error('Client has already been connected. You cannot reuse a client.')))
|
|---|
| 93 | return
|
|---|
| 94 | }
|
|---|
| 95 |
|
|---|
| 96 | this._connecting = true
|
|---|
| 97 |
|
|---|
| 98 | this.connectionParameters.getLibpqConnectionString(function (err, conString) {
|
|---|
| 99 | if (self.connectionParameters.nativeConnectionString) conString = self.connectionParameters.nativeConnectionString
|
|---|
| 100 | if (err) return cb(err)
|
|---|
| 101 | self.native.connect(conString, function (err) {
|
|---|
| 102 | if (err) {
|
|---|
| 103 | self.native.end()
|
|---|
| 104 | return cb(err)
|
|---|
| 105 | }
|
|---|
| 106 |
|
|---|
| 107 | // set internal states to connected
|
|---|
| 108 | self._connected = true
|
|---|
| 109 |
|
|---|
| 110 | // handle connection errors from the native layer
|
|---|
| 111 | self.native.on('error', function (err) {
|
|---|
| 112 | self._queryable = false
|
|---|
| 113 | self._errorAllQueries(err)
|
|---|
| 114 | self.emit('error', err)
|
|---|
| 115 | })
|
|---|
| 116 |
|
|---|
| 117 | self.native.on('notification', function (msg) {
|
|---|
| 118 | self.emit('notification', {
|
|---|
| 119 | channel: msg.relname,
|
|---|
| 120 | payload: msg.extra,
|
|---|
| 121 | })
|
|---|
| 122 | })
|
|---|
| 123 |
|
|---|
| 124 | // signal we are connected now
|
|---|
| 125 | self.emit('connect')
|
|---|
| 126 | self._pulseQueryQueue(true)
|
|---|
| 127 |
|
|---|
| 128 | cb(null, this)
|
|---|
| 129 | })
|
|---|
| 130 | })
|
|---|
| 131 | }
|
|---|
| 132 |
|
|---|
| 133 | Client.prototype.connect = function (callback) {
|
|---|
| 134 | if (callback) {
|
|---|
| 135 | this._connect(callback)
|
|---|
| 136 | return
|
|---|
| 137 | }
|
|---|
| 138 |
|
|---|
| 139 | return new this._Promise((resolve, reject) => {
|
|---|
| 140 | this._connect((error) => {
|
|---|
| 141 | if (error) {
|
|---|
| 142 | reject(error)
|
|---|
| 143 | } else {
|
|---|
| 144 | resolve(this)
|
|---|
| 145 | }
|
|---|
| 146 | })
|
|---|
| 147 | })
|
|---|
| 148 | }
|
|---|
| 149 |
|
|---|
| 150 | // send a query to the server
|
|---|
| 151 | // this method is highly overloaded to take
|
|---|
| 152 | // 1) string query, optional array of parameters, optional function callback
|
|---|
| 153 | // 2) object query with {
|
|---|
| 154 | // string query
|
|---|
| 155 | // optional array values,
|
|---|
| 156 | // optional function callback instead of as a separate parameter
|
|---|
| 157 | // optional string name to name & cache the query plan
|
|---|
| 158 | // optional string rowMode = 'array' for an array of results
|
|---|
| 159 | // }
|
|---|
| 160 | Client.prototype.query = function (config, values, callback) {
|
|---|
| 161 | let query
|
|---|
| 162 | let result
|
|---|
| 163 | let readTimeout
|
|---|
| 164 | let readTimeoutTimer
|
|---|
| 165 | let queryCallback
|
|---|
| 166 |
|
|---|
| 167 | if (config === null || config === undefined) {
|
|---|
| 168 | throw new TypeError('Client was passed a null or undefined query')
|
|---|
| 169 | } else if (typeof config.submit === 'function') {
|
|---|
| 170 | readTimeout = config.query_timeout || this.connectionParameters.query_timeout
|
|---|
| 171 | result = query = config
|
|---|
| 172 | // accept query(new Query(...), (err, res) => { }) style
|
|---|
| 173 | if (typeof values === 'function') {
|
|---|
| 174 | config.callback = values
|
|---|
| 175 | }
|
|---|
| 176 | } else {
|
|---|
| 177 | readTimeout = config.query_timeout || this.connectionParameters.query_timeout
|
|---|
| 178 | query = new NativeQuery(config, values, callback)
|
|---|
| 179 | if (!query.callback) {
|
|---|
| 180 | let resolveOut, rejectOut
|
|---|
| 181 | result = new this._Promise((resolve, reject) => {
|
|---|
| 182 | resolveOut = resolve
|
|---|
| 183 | rejectOut = reject
|
|---|
| 184 | }).catch((err) => {
|
|---|
| 185 | Error.captureStackTrace(err)
|
|---|
| 186 | throw err
|
|---|
| 187 | })
|
|---|
| 188 | query.callback = (err, res) => (err ? rejectOut(err) : resolveOut(res))
|
|---|
| 189 | }
|
|---|
| 190 | }
|
|---|
| 191 |
|
|---|
| 192 | if (readTimeout) {
|
|---|
| 193 | queryCallback = query.callback || (() => {})
|
|---|
| 194 |
|
|---|
| 195 | readTimeoutTimer = setTimeout(() => {
|
|---|
| 196 | const error = new Error('Query read timeout')
|
|---|
| 197 |
|
|---|
| 198 | process.nextTick(() => {
|
|---|
| 199 | query.handleError(error, this.connection)
|
|---|
| 200 | })
|
|---|
| 201 |
|
|---|
| 202 | queryCallback(error)
|
|---|
| 203 |
|
|---|
| 204 | // we already returned an error,
|
|---|
| 205 | // just do nothing if query completes
|
|---|
| 206 | query.callback = () => {}
|
|---|
| 207 |
|
|---|
| 208 | // Remove from queue
|
|---|
| 209 | const index = this._queryQueue.indexOf(query)
|
|---|
| 210 | if (index > -1) {
|
|---|
| 211 | this._queryQueue.splice(index, 1)
|
|---|
| 212 | }
|
|---|
| 213 |
|
|---|
| 214 | this._pulseQueryQueue()
|
|---|
| 215 | }, readTimeout)
|
|---|
| 216 |
|
|---|
| 217 | query.callback = (err, res) => {
|
|---|
| 218 | clearTimeout(readTimeoutTimer)
|
|---|
| 219 | queryCallback(err, res)
|
|---|
| 220 | }
|
|---|
| 221 | }
|
|---|
| 222 |
|
|---|
| 223 | if (!this._queryable) {
|
|---|
| 224 | query.native = this.native
|
|---|
| 225 | process.nextTick(() => {
|
|---|
| 226 | query.handleError(new Error('Client has encountered a connection error and is not queryable'))
|
|---|
| 227 | })
|
|---|
| 228 | return result
|
|---|
| 229 | }
|
|---|
| 230 |
|
|---|
| 231 | if (this._ending) {
|
|---|
| 232 | query.native = this.native
|
|---|
| 233 | process.nextTick(() => {
|
|---|
| 234 | query.handleError(new Error('Client was closed and is not queryable'))
|
|---|
| 235 | })
|
|---|
| 236 | return result
|
|---|
| 237 | }
|
|---|
| 238 |
|
|---|
| 239 | if (this._queryQueue.length > 0 && !this.pipeline) {
|
|---|
| 240 | queryQueueLengthDeprecationNotice()
|
|---|
| 241 | }
|
|---|
| 242 |
|
|---|
| 243 | this._queryQueue.push(query)
|
|---|
| 244 | this._pulseQueryQueue()
|
|---|
| 245 | return result
|
|---|
| 246 | }
|
|---|
| 247 |
|
|---|
| 248 | // disconnect from the backend server
|
|---|
| 249 | Client.prototype.end = function (cb) {
|
|---|
| 250 | const self = this
|
|---|
| 251 |
|
|---|
| 252 | this._ending = true
|
|---|
| 253 |
|
|---|
| 254 | if (this._connecting && !this._connected) {
|
|---|
| 255 | this.once('connect', () => {
|
|---|
| 256 | this.end(() => {})
|
|---|
| 257 | })
|
|---|
| 258 | }
|
|---|
| 259 | let result
|
|---|
| 260 | if (!cb) {
|
|---|
| 261 | result = new this._Promise(function (resolve, reject) {
|
|---|
| 262 | cb = (err) => (err ? reject(err) : resolve())
|
|---|
| 263 | })
|
|---|
| 264 | }
|
|---|
| 265 |
|
|---|
| 266 | const doEnd = function () {
|
|---|
| 267 | self.native.end(function () {
|
|---|
| 268 | self._connected = false
|
|---|
| 269 |
|
|---|
| 270 | self._errorAllQueries(new Error('Connection terminated'))
|
|---|
| 271 |
|
|---|
| 272 | process.nextTick(() => {
|
|---|
| 273 | self.emit('end')
|
|---|
| 274 | if (cb) cb()
|
|---|
| 275 | })
|
|---|
| 276 | })
|
|---|
| 277 | }
|
|---|
| 278 |
|
|---|
| 279 | // If pipeline has in-flight or queued queries, wait for them to drain before closing
|
|---|
| 280 | if (this.pipeline && (this._pipelineInFlight || this._queryQueue.length > 0)) {
|
|---|
| 281 | this.once('drain', doEnd)
|
|---|
| 282 | } else {
|
|---|
| 283 | doEnd()
|
|---|
| 284 | }
|
|---|
| 285 | return result
|
|---|
| 286 | }
|
|---|
| 287 |
|
|---|
| 288 | Client.prototype._hasActiveQuery = function () {
|
|---|
| 289 | return this._activeQuery && this._activeQuery.state !== 'error' && this._activeQuery.state !== 'end'
|
|---|
| 290 | }
|
|---|
| 291 |
|
|---|
| 292 | Client.prototype._pulseQueryQueue = function (initialConnection) {
|
|---|
| 293 | if (!this._connected) {
|
|---|
| 294 | return
|
|---|
| 295 | }
|
|---|
| 296 | if (this.pipeline && !initialConnection) {
|
|---|
| 297 | return this._pulsePipelinedQueryQueue()
|
|---|
| 298 | }
|
|---|
| 299 | if (this._hasActiveQuery()) {
|
|---|
| 300 | return
|
|---|
| 301 | }
|
|---|
| 302 | const query = this._queryQueue.shift()
|
|---|
| 303 | if (!query) {
|
|---|
| 304 | if (!initialConnection) {
|
|---|
| 305 | this.emit('drain')
|
|---|
| 306 | }
|
|---|
| 307 | return
|
|---|
| 308 | }
|
|---|
| 309 | this._activeQuery = query
|
|---|
| 310 | query.submit(this)
|
|---|
| 311 | const self = this
|
|---|
| 312 | query.once('_done', function () {
|
|---|
| 313 | self._pulseQueryQueue()
|
|---|
| 314 | })
|
|---|
| 315 | }
|
|---|
| 316 |
|
|---|
| 317 | Client.prototype._pulsePipelinedQueryQueue = function () {
|
|---|
| 318 | if (!this._connected || this._pipelineInFlight) {
|
|---|
| 319 | return
|
|---|
| 320 | }
|
|---|
| 321 | if (this._queryQueue.length === 0) {
|
|---|
| 322 | if (this.hasExecuted) {
|
|---|
| 323 | this.emit('drain')
|
|---|
| 324 | }
|
|---|
| 325 | return
|
|---|
| 326 | }
|
|---|
| 327 |
|
|---|
| 328 | this._pipelineInFlight = true
|
|---|
| 329 | const self = this
|
|---|
| 330 | const queries = []
|
|---|
| 331 | const nativeQueries = []
|
|---|
| 332 | const utils = require('../utils')
|
|---|
| 333 |
|
|---|
| 334 | while (this._queryQueue.length > 0) {
|
|---|
| 335 | const query = this._queryQueue.shift()
|
|---|
| 336 | this.hasExecuted = true
|
|---|
| 337 | nativeQueries.push(query)
|
|---|
| 338 |
|
|---|
| 339 | const values = query.values ? query.values.map(utils.prepareValue) : null
|
|---|
| 340 | const pipelineEntry = { text: query.text, name: query.name }
|
|---|
| 341 | if (values) {
|
|---|
| 342 | pipelineEntry.values = values
|
|---|
| 343 | }
|
|---|
| 344 | if (query.name && this.namedQueries[query.name]) {
|
|---|
| 345 | pipelineEntry._alreadyPrepared = true
|
|---|
| 346 | }
|
|---|
| 347 | queries.push(pipelineEntry)
|
|---|
| 348 | }
|
|---|
| 349 |
|
|---|
| 350 | this.native.pipeline(queries, function (err, results) {
|
|---|
| 351 | self._pipelineInFlight = false
|
|---|
| 352 |
|
|---|
| 353 | if (err) {
|
|---|
| 354 | // Total pipeline failure — error all queries
|
|---|
| 355 | for (let i = 0; i < nativeQueries.length; i++) {
|
|---|
| 356 | const q = nativeQueries[i]
|
|---|
| 357 | q.native = self.native
|
|---|
| 358 | q.handleError(err)
|
|---|
| 359 | }
|
|---|
| 360 | self._pulsePipelinedQueryQueue()
|
|---|
| 361 | return
|
|---|
| 362 | }
|
|---|
| 363 |
|
|---|
| 364 | // Deliver results to each query
|
|---|
| 365 | for (let i = 0; i < nativeQueries.length; i++) {
|
|---|
| 366 | const q = nativeQueries[i]
|
|---|
| 367 | const r = results[i]
|
|---|
| 368 | q.native = self.native
|
|---|
| 369 |
|
|---|
| 370 | if (r.err) {
|
|---|
| 371 | q.handleError(r.err)
|
|---|
| 372 | } else {
|
|---|
| 373 | // Track named queries on success
|
|---|
| 374 | if (q.name) {
|
|---|
| 375 | self.namedQueries[q.name] = q.text
|
|---|
| 376 | }
|
|---|
| 377 | q.state = 'end'
|
|---|
| 378 | q.emit('end', r.result)
|
|---|
| 379 | if (q.callback) {
|
|---|
| 380 | q.callback(null, r.result)
|
|---|
| 381 | }
|
|---|
| 382 | }
|
|---|
| 383 |
|
|---|
| 384 | setImmediate(function () {
|
|---|
| 385 | q.emit('_done')
|
|---|
| 386 | })
|
|---|
| 387 | }
|
|---|
| 388 |
|
|---|
| 389 | // Process any queries that arrived while we were reading
|
|---|
| 390 | self._pulsePipelinedQueryQueue()
|
|---|
| 391 | })
|
|---|
| 392 | }
|
|---|
| 393 |
|
|---|
| 394 | // attempt to cancel an in-progress query
|
|---|
| 395 | Client.prototype.cancel = function (query) {
|
|---|
| 396 | if (this._activeQuery === query) {
|
|---|
| 397 | this.native.cancel(function () {})
|
|---|
| 398 | } else if (this._queryQueue.indexOf(query) !== -1) {
|
|---|
| 399 | this._queryQueue.splice(this._queryQueue.indexOf(query), 1)
|
|---|
| 400 | }
|
|---|
| 401 | }
|
|---|
| 402 |
|
|---|
| 403 | Client.prototype.ref = function () {}
|
|---|
| 404 | Client.prototype.unref = function () {}
|
|---|
| 405 |
|
|---|
| 406 | Client.prototype.setTypeParser = function (oid, format, parseFn) {
|
|---|
| 407 | return this._types.setTypeParser(oid, format, parseFn)
|
|---|
| 408 | }
|
|---|
| 409 |
|
|---|
| 410 | Client.prototype.getTypeParser = function (oid, format) {
|
|---|
| 411 | return this._types.getTypeParser(oid, format)
|
|---|
| 412 | }
|
|---|
| 413 |
|
|---|
| 414 | Client.prototype.isConnected = function () {
|
|---|
| 415 | return this._connected
|
|---|
| 416 | }
|
|---|
| 417 |
|
|---|
| 418 | Client.prototype.getTransactionStatus = function () {
|
|---|
| 419 | return this.native.getTransactionStatus()
|
|---|
| 420 | }
|
|---|