| [62b2964] | 1 | const EventEmitter = require('events').EventEmitter
|
|---|
| 2 | const utils = require('./utils')
|
|---|
| 3 | const nodeUtils = require('util')
|
|---|
| 4 | const sasl = require('./crypto/sasl')
|
|---|
| 5 | const TypeOverrides = require('./type-overrides')
|
|---|
| 6 |
|
|---|
| 7 | const ConnectionParameters = require('./connection-parameters')
|
|---|
| 8 | const Query = require('./query')
|
|---|
| 9 | const defaults = require('./defaults')
|
|---|
| 10 | const Connection = require('./connection')
|
|---|
| 11 | const crypto = require('./crypto/utils')
|
|---|
| 12 |
|
|---|
| 13 | const activeQueryDeprecationNotice = nodeUtils.deprecate(
|
|---|
| 14 | () => {},
|
|---|
| 15 | 'Client.activeQuery is deprecated and will be removed in pg@9.0'
|
|---|
| 16 | )
|
|---|
| 17 |
|
|---|
| 18 | const queryQueueDeprecationNotice = nodeUtils.deprecate(
|
|---|
| 19 | () => {},
|
|---|
| 20 | 'Client.queryQueue is deprecated and will be removed in pg@9.0.'
|
|---|
| 21 | )
|
|---|
| 22 |
|
|---|
| 23 | const pgPassDeprecationNotice = nodeUtils.deprecate(
|
|---|
| 24 | () => {},
|
|---|
| 25 | 'pgpass support is deprecated and will be removed in pg@9.0. ' +
|
|---|
| 26 | 'You can provide an async function as the password property to the Client/Pool constructor that returns a password instead. Within this function you can call the pgpass module in your own code.'
|
|---|
| 27 | )
|
|---|
| 28 |
|
|---|
| 29 | const byoPromiseDeprecationNotice = nodeUtils.deprecate(
|
|---|
| 30 | () => {},
|
|---|
| 31 | 'Passing a custom Promise implementation to the Client/Pool constructor is deprecated and will be removed in pg@9.0.'
|
|---|
| 32 | )
|
|---|
| 33 |
|
|---|
| 34 | const queryQueueLengthDeprecationNotice = nodeUtils.deprecate(
|
|---|
| 35 | () => {},
|
|---|
| 36 | '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.'
|
|---|
| 37 | )
|
|---|
| 38 |
|
|---|
| 39 | function coerceNumberOrDefault(value, defaultValue) {
|
|---|
| 40 | if (typeof value === 'number') {
|
|---|
| 41 | return Number.isFinite(value) ? value : defaultValue
|
|---|
| 42 | }
|
|---|
| 43 | if (typeof value === 'string' && value.trim() !== '') {
|
|---|
| 44 | const n = Number(value)
|
|---|
| 45 | return Number.isFinite(n) ? n : defaultValue
|
|---|
| 46 | }
|
|---|
| 47 | return defaultValue
|
|---|
| 48 | }
|
|---|
| 49 |
|
|---|
| 50 | class Client extends EventEmitter {
|
|---|
| 51 | constructor(config) {
|
|---|
| 52 | super()
|
|---|
| 53 |
|
|---|
| 54 | this.connectionParameters = new ConnectionParameters(config)
|
|---|
| 55 | this.user = this.connectionParameters.user
|
|---|
| 56 | this.database = this.connectionParameters.database
|
|---|
| 57 | this.port = this.connectionParameters.port
|
|---|
| 58 | this.host = this.connectionParameters.host
|
|---|
| 59 |
|
|---|
| 60 | // "hiding" the password so it doesn't show up in stack traces
|
|---|
| 61 | // or if the client is console.logged
|
|---|
| 62 | Object.defineProperty(this, 'password', {
|
|---|
| 63 | configurable: true,
|
|---|
| 64 | enumerable: false,
|
|---|
| 65 | writable: true,
|
|---|
| 66 | value: this.connectionParameters.password,
|
|---|
| 67 | })
|
|---|
| 68 |
|
|---|
| 69 | this.replication = this.connectionParameters.replication
|
|---|
| 70 |
|
|---|
| 71 | const c = config || {}
|
|---|
| 72 |
|
|---|
| 73 | if (c.Promise) {
|
|---|
| 74 | byoPromiseDeprecationNotice()
|
|---|
| 75 | }
|
|---|
| 76 | this._Promise = c.Promise || global.Promise
|
|---|
| 77 | this._types = new TypeOverrides(c.types)
|
|---|
| 78 | this._ending = false
|
|---|
| 79 | this._ended = false
|
|---|
| 80 | this._connecting = false
|
|---|
| 81 | this._connected = false
|
|---|
| 82 | this._connectionError = false
|
|---|
| 83 | this._queryable = true
|
|---|
| 84 | this._activeQuery = null
|
|---|
| 85 | this._txStatus = null
|
|---|
| 86 |
|
|---|
| 87 | this.enableChannelBinding = Boolean(c.enableChannelBinding) // set true to use SCRAM-SHA-256-PLUS when offered
|
|---|
| 88 | this.scramMaxIterations = coerceNumberOrDefault(c.scramMaxIterations, sasl.DEFAULT_MAX_SCRAM_ITERATIONS)
|
|---|
| 89 | this.connection =
|
|---|
| 90 | c.connection ||
|
|---|
| 91 | new Connection({
|
|---|
| 92 | stream: c.stream,
|
|---|
| 93 | ssl: this.connectionParameters.ssl,
|
|---|
| 94 | sslNegotiation: this.connectionParameters.sslnegotiation,
|
|---|
| 95 | keepAlive: c.keepAlive || false,
|
|---|
| 96 | keepAliveInitialDelayMillis: c.keepAliveInitialDelayMillis || 0,
|
|---|
| 97 | encoding: this.connectionParameters.client_encoding || 'utf8',
|
|---|
| 98 | })
|
|---|
| 99 | this._queryQueue = []
|
|---|
| 100 | this._sentQueryQueue = []
|
|---|
| 101 | this.pipeline = Boolean(c.pipeline)
|
|---|
| 102 | this.binary = c.binary || defaults.binary
|
|---|
| 103 | this.processID = null
|
|---|
| 104 | this.secretKey = null
|
|---|
| 105 | this.ssl = this.connectionParameters.ssl || false
|
|---|
| 106 | this.sslNegotiation = this.connectionParameters.sslnegotiation || 'postgres'
|
|---|
| 107 | // As with Password, make SSL->Key (the private key) non-enumerable.
|
|---|
| 108 | // It won't show up in stack traces
|
|---|
| 109 | // or if the client is console.logged
|
|---|
| 110 | if (this.ssl && this.ssl.key) {
|
|---|
| 111 | Object.defineProperty(this.ssl, 'key', {
|
|---|
| 112 | enumerable: false,
|
|---|
| 113 | })
|
|---|
| 114 | }
|
|---|
| 115 |
|
|---|
| 116 | this._connectionTimeoutMillis = c.connectionTimeoutMillis || 0
|
|---|
| 117 | }
|
|---|
| 118 |
|
|---|
| 119 | get activeQuery() {
|
|---|
| 120 | activeQueryDeprecationNotice()
|
|---|
| 121 | return this._activeQuery
|
|---|
| 122 | }
|
|---|
| 123 |
|
|---|
| 124 | set activeQuery(val) {
|
|---|
| 125 | activeQueryDeprecationNotice()
|
|---|
| 126 | this._activeQuery = val
|
|---|
| 127 | }
|
|---|
| 128 |
|
|---|
| 129 | _getActiveQuery() {
|
|---|
| 130 | return this._activeQuery
|
|---|
| 131 | }
|
|---|
| 132 |
|
|---|
| 133 | _errorAllQueries(err) {
|
|---|
| 134 | const enqueueError = (query) => {
|
|---|
| 135 | process.nextTick(() => {
|
|---|
| 136 | query.handleError(err, this.connection)
|
|---|
| 137 | })
|
|---|
| 138 | }
|
|---|
| 139 |
|
|---|
| 140 | const activeQuery = this._getActiveQuery()
|
|---|
| 141 | if (activeQuery) {
|
|---|
| 142 | enqueueError(activeQuery)
|
|---|
| 143 | this._activeQuery = null
|
|---|
| 144 | }
|
|---|
| 145 |
|
|---|
| 146 | this._sentQueryQueue.forEach(enqueueError)
|
|---|
| 147 | this._sentQueryQueue.length = 0
|
|---|
| 148 |
|
|---|
| 149 | this._queryQueue.forEach(enqueueError)
|
|---|
| 150 | this._queryQueue.length = 0
|
|---|
| 151 | }
|
|---|
| 152 |
|
|---|
| 153 | _connect(callback) {
|
|---|
| 154 | const self = this
|
|---|
| 155 | const con = this.connection
|
|---|
| 156 | this._connectionCallback = callback
|
|---|
| 157 |
|
|---|
| 158 | if (this._connecting || this._connected) {
|
|---|
| 159 | const err = new Error('Client has already been connected. You cannot reuse a client.')
|
|---|
| 160 | process.nextTick(() => {
|
|---|
| 161 | callback(err)
|
|---|
| 162 | })
|
|---|
| 163 | return
|
|---|
| 164 | }
|
|---|
| 165 | this._connecting = true
|
|---|
| 166 |
|
|---|
| 167 | if (this._connectionTimeoutMillis > 0) {
|
|---|
| 168 | this.connectionTimeoutHandle = setTimeout(() => {
|
|---|
| 169 | con._ending = true
|
|---|
| 170 | con.stream.destroy(new Error('timeout expired'))
|
|---|
| 171 | }, this._connectionTimeoutMillis)
|
|---|
| 172 |
|
|---|
| 173 | if (this.connectionTimeoutHandle.unref) {
|
|---|
| 174 | this.connectionTimeoutHandle.unref()
|
|---|
| 175 | }
|
|---|
| 176 | }
|
|---|
| 177 |
|
|---|
| 178 | if (this.host && this.host.indexOf('/') === 0) {
|
|---|
| 179 | con.connect(this.host + '/.s.PGSQL.' + this.port)
|
|---|
| 180 | } else {
|
|---|
| 181 | con.connect(this.port, this.host)
|
|---|
| 182 | }
|
|---|
| 183 |
|
|---|
| 184 | // once connection is established send startup message
|
|---|
| 185 | con.on('connect', function () {
|
|---|
| 186 | if (self.ssl) {
|
|---|
| 187 | // With direct SSL negotiation the connection upgrades to TLS without an
|
|---|
| 188 | // SSLRequest packet, so the startup message is sent after 'sslconnect'.
|
|---|
| 189 | if (self.sslNegotiation !== 'direct') {
|
|---|
| 190 | con.requestSsl()
|
|---|
| 191 | }
|
|---|
| 192 | } else {
|
|---|
| 193 | con.startup(self.getStartupConf())
|
|---|
| 194 | }
|
|---|
| 195 | })
|
|---|
| 196 |
|
|---|
| 197 | con.on('sslconnect', function () {
|
|---|
| 198 | con.startup(self.getStartupConf())
|
|---|
| 199 | })
|
|---|
| 200 |
|
|---|
| 201 | this._attachListeners(con)
|
|---|
| 202 |
|
|---|
| 203 | con.once('end', () => {
|
|---|
| 204 | const error = this._ending ? new Error('Connection terminated') : new Error('Connection terminated unexpectedly')
|
|---|
| 205 |
|
|---|
| 206 | clearTimeout(this.connectionTimeoutHandle)
|
|---|
| 207 | this._errorAllQueries(error)
|
|---|
| 208 | this._ended = true
|
|---|
| 209 |
|
|---|
| 210 | if (!this._ending) {
|
|---|
| 211 | // if the connection is ended without us calling .end()
|
|---|
| 212 | // on this client then we have an unexpected disconnection
|
|---|
| 213 | // treat this as an error unless we've already emitted an error
|
|---|
| 214 | // during connection.
|
|---|
| 215 | if (this._connecting && !this._connectionError) {
|
|---|
| 216 | if (this._connectionCallback) {
|
|---|
| 217 | this._connectionCallback(error)
|
|---|
| 218 | } else {
|
|---|
| 219 | this._handleErrorEvent(error)
|
|---|
| 220 | }
|
|---|
| 221 | } else if (!this._connectionError) {
|
|---|
| 222 | this._handleErrorEvent(error)
|
|---|
| 223 | }
|
|---|
| 224 | }
|
|---|
| 225 |
|
|---|
| 226 | process.nextTick(() => {
|
|---|
| 227 | this.emit('end')
|
|---|
| 228 | })
|
|---|
| 229 | })
|
|---|
| 230 | }
|
|---|
| 231 |
|
|---|
| 232 | connect(callback) {
|
|---|
| 233 | if (callback) {
|
|---|
| 234 | this._connect(callback)
|
|---|
| 235 | return
|
|---|
| 236 | }
|
|---|
| 237 |
|
|---|
| 238 | return new this._Promise((resolve, reject) => {
|
|---|
| 239 | this._connect((error) => {
|
|---|
| 240 | if (error) {
|
|---|
| 241 | reject(error)
|
|---|
| 242 | } else {
|
|---|
| 243 | resolve(this)
|
|---|
| 244 | }
|
|---|
| 245 | })
|
|---|
| 246 | })
|
|---|
| 247 | }
|
|---|
| 248 |
|
|---|
| 249 | _attachListeners(con) {
|
|---|
| 250 | // password request handling
|
|---|
| 251 | con.on('authenticationCleartextPassword', this._handleAuthCleartextPassword.bind(this))
|
|---|
| 252 | // password request handling
|
|---|
| 253 | con.on('authenticationMD5Password', this._handleAuthMD5Password.bind(this))
|
|---|
| 254 | // password request handling (SASL)
|
|---|
| 255 | con.on('authenticationSASL', this._handleAuthSASL.bind(this))
|
|---|
| 256 | con.on('authenticationSASLContinue', this._handleAuthSASLContinue.bind(this))
|
|---|
| 257 | con.on('authenticationSASLFinal', this._handleAuthSASLFinal.bind(this))
|
|---|
| 258 | con.on('backendKeyData', this._handleBackendKeyData.bind(this))
|
|---|
| 259 | con.on('error', this._handleErrorEvent.bind(this))
|
|---|
| 260 | con.on('errorMessage', this._handleErrorMessage.bind(this))
|
|---|
| 261 | con.on('readyForQuery', this._handleReadyForQuery.bind(this))
|
|---|
| 262 | con.on('notice', this._handleNotice.bind(this))
|
|---|
| 263 | con.on('rowDescription', this._handleRowDescription.bind(this))
|
|---|
| 264 | con.on('dataRow', this._handleDataRow.bind(this))
|
|---|
| 265 | con.on('portalSuspended', this._handlePortalSuspended.bind(this))
|
|---|
| 266 | con.on('emptyQuery', this._handleEmptyQuery.bind(this))
|
|---|
| 267 | con.on('commandComplete', this._handleCommandComplete.bind(this))
|
|---|
| 268 | con.on('parseComplete', this._handleParseComplete.bind(this))
|
|---|
| 269 | con.on('copyInResponse', this._handleCopyInResponse.bind(this))
|
|---|
| 270 | con.on('copyData', this._handleCopyData.bind(this))
|
|---|
| 271 | con.on('notification', this._handleNotification.bind(this))
|
|---|
| 272 | }
|
|---|
| 273 |
|
|---|
| 274 | _getPassword(cb) {
|
|---|
| 275 | const con = this.connection
|
|---|
| 276 | if (typeof this.password === 'function') {
|
|---|
| 277 | this._Promise
|
|---|
| 278 | .resolve()
|
|---|
| 279 | .then(() => this.password(this.connectionParameters))
|
|---|
| 280 | .then((pass) => {
|
|---|
| 281 | if (pass !== undefined) {
|
|---|
| 282 | if (typeof pass !== 'string') {
|
|---|
| 283 | con.emit('error', new TypeError('Password must be a string'))
|
|---|
| 284 | return
|
|---|
| 285 | }
|
|---|
| 286 | this.connectionParameters.password = this.password = pass
|
|---|
| 287 | } else {
|
|---|
| 288 | this.connectionParameters.password = this.password = null
|
|---|
| 289 | }
|
|---|
| 290 | cb()
|
|---|
| 291 | })
|
|---|
| 292 | .catch((err) => {
|
|---|
| 293 | con.emit('error', err)
|
|---|
| 294 | })
|
|---|
| 295 | } else if (this.password !== null) {
|
|---|
| 296 | cb()
|
|---|
| 297 | } else {
|
|---|
| 298 | try {
|
|---|
| 299 | const pgPass = require('pgpass')
|
|---|
| 300 | pgPass(this.connectionParameters, (pass) => {
|
|---|
| 301 | if (undefined !== pass) {
|
|---|
| 302 | pgPassDeprecationNotice()
|
|---|
| 303 | this.connectionParameters.password = this.password = pass
|
|---|
| 304 | }
|
|---|
| 305 | cb()
|
|---|
| 306 | })
|
|---|
| 307 | } catch (e) {
|
|---|
| 308 | this.emit('error', e)
|
|---|
| 309 | }
|
|---|
| 310 | }
|
|---|
| 311 | }
|
|---|
| 312 |
|
|---|
| 313 | _handleAuthCleartextPassword(msg) {
|
|---|
| 314 | this._getPassword(() => {
|
|---|
| 315 | this.connection.password(this.password)
|
|---|
| 316 | })
|
|---|
| 317 | }
|
|---|
| 318 |
|
|---|
| 319 | _handleAuthMD5Password(msg) {
|
|---|
| 320 | this._getPassword(async () => {
|
|---|
| 321 | try {
|
|---|
| 322 | const hashedPassword = await crypto.postgresMd5PasswordHash(this.user, this.password, msg.salt)
|
|---|
| 323 | this.connection.password(hashedPassword)
|
|---|
| 324 | } catch (e) {
|
|---|
| 325 | this.emit('error', e)
|
|---|
| 326 | }
|
|---|
| 327 | })
|
|---|
| 328 | }
|
|---|
| 329 |
|
|---|
| 330 | _handleAuthSASL(msg) {
|
|---|
| 331 | this._getPassword(() => {
|
|---|
| 332 | try {
|
|---|
| 333 | this.saslSession = sasl.startSession(
|
|---|
| 334 | msg.mechanisms,
|
|---|
| 335 | this.enableChannelBinding && this.connection.stream,
|
|---|
| 336 | this.scramMaxIterations
|
|---|
| 337 | )
|
|---|
| 338 | this.connection.sendSASLInitialResponseMessage(this.saslSession.mechanism, this.saslSession.response)
|
|---|
| 339 | } catch (err) {
|
|---|
| 340 | this.connection.emit('error', err)
|
|---|
| 341 | }
|
|---|
| 342 | })
|
|---|
| 343 | }
|
|---|
| 344 |
|
|---|
| 345 | async _handleAuthSASLContinue(msg) {
|
|---|
| 346 | try {
|
|---|
| 347 | await sasl.continueSession(
|
|---|
| 348 | this.saslSession,
|
|---|
| 349 | this.password,
|
|---|
| 350 | msg.data,
|
|---|
| 351 | this.enableChannelBinding && this.connection.stream
|
|---|
| 352 | )
|
|---|
| 353 | this.connection.sendSCRAMClientFinalMessage(this.saslSession.response)
|
|---|
| 354 | } catch (err) {
|
|---|
| 355 | this.connection.emit('error', err)
|
|---|
| 356 | }
|
|---|
| 357 | }
|
|---|
| 358 |
|
|---|
| 359 | _handleAuthSASLFinal(msg) {
|
|---|
| 360 | try {
|
|---|
| 361 | sasl.finalizeSession(this.saslSession, msg.data)
|
|---|
| 362 | this.saslSession = null
|
|---|
| 363 | } catch (err) {
|
|---|
| 364 | this.connection.emit('error', err)
|
|---|
| 365 | }
|
|---|
| 366 | }
|
|---|
| 367 |
|
|---|
| 368 | _handleBackendKeyData(msg) {
|
|---|
| 369 | this.processID = msg.processID
|
|---|
| 370 | this.secretKey = msg.secretKey
|
|---|
| 371 | }
|
|---|
| 372 |
|
|---|
| 373 | _handleReadyForQuery(msg) {
|
|---|
| 374 | if (this._connecting) {
|
|---|
| 375 | this._connecting = false
|
|---|
| 376 | this._connected = true
|
|---|
| 377 | clearTimeout(this.connectionTimeoutHandle)
|
|---|
| 378 |
|
|---|
| 379 | // process possible callback argument to Client#connect
|
|---|
| 380 | if (this._connectionCallback) {
|
|---|
| 381 | this._connectionCallback(null, this)
|
|---|
| 382 | // remove callback for proper error handling
|
|---|
| 383 | // after the connect event
|
|---|
| 384 | this._connectionCallback = null
|
|---|
| 385 | }
|
|---|
| 386 | this.emit('connect')
|
|---|
| 387 | }
|
|---|
| 388 | const activeQuery = this._getActiveQuery()
|
|---|
| 389 | this._activeQuery = null
|
|---|
| 390 | this._txStatus = msg?.status ?? null
|
|---|
| 391 | this.readyForQuery = true
|
|---|
| 392 | if (activeQuery) {
|
|---|
| 393 | activeQuery.handleReadyForQuery(this.connection)
|
|---|
| 394 | }
|
|---|
| 395 | this._pulseQueryQueue()
|
|---|
| 396 | }
|
|---|
| 397 |
|
|---|
| 398 | // if we receive an error event or error message
|
|---|
| 399 | // during the connection process we handle it here
|
|---|
| 400 | _handleErrorWhileConnecting(err) {
|
|---|
| 401 | if (this._connectionError) {
|
|---|
| 402 | // TODO(bmc): this is swallowing errors - we shouldn't do this
|
|---|
| 403 | return
|
|---|
| 404 | }
|
|---|
| 405 | this._connectionError = true
|
|---|
| 406 | clearTimeout(this.connectionTimeoutHandle)
|
|---|
| 407 | if (this._connectionCallback) {
|
|---|
| 408 | return this._connectionCallback(err)
|
|---|
| 409 | }
|
|---|
| 410 | this.emit('error', err)
|
|---|
| 411 | }
|
|---|
| 412 |
|
|---|
| 413 | // if we're connected and we receive an error event from the connection
|
|---|
| 414 | // this means the socket is dead - do a hard abort of all queries and emit
|
|---|
| 415 | // the socket error on the client as well
|
|---|
| 416 | _handleErrorEvent(err) {
|
|---|
| 417 | if (this._connecting) {
|
|---|
| 418 | return this._handleErrorWhileConnecting(err)
|
|---|
| 419 | }
|
|---|
| 420 | this._queryable = false
|
|---|
| 421 | this._errorAllQueries(err)
|
|---|
| 422 | this.emit('error', err)
|
|---|
| 423 | }
|
|---|
| 424 |
|
|---|
| 425 | // handle error messages from the postgres backend
|
|---|
| 426 | _handleErrorMessage(msg) {
|
|---|
| 427 | if (this._connecting) {
|
|---|
| 428 | return this._handleErrorWhileConnecting(msg)
|
|---|
| 429 | }
|
|---|
| 430 | const activeQuery = this._getActiveQuery()
|
|---|
| 431 |
|
|---|
| 432 | if (!activeQuery) {
|
|---|
| 433 | this._handleErrorEvent(msg)
|
|---|
| 434 | return
|
|---|
| 435 | }
|
|---|
| 436 |
|
|---|
| 437 | this._activeQuery = null
|
|---|
| 438 | if (activeQuery.name) {
|
|---|
| 439 | delete this.connection.submittedNamedStatements[activeQuery.name]
|
|---|
| 440 | }
|
|---|
| 441 | activeQuery.handleError(msg, this.connection)
|
|---|
| 442 | }
|
|---|
| 443 |
|
|---|
| 444 | _handleRowDescription(msg) {
|
|---|
| 445 | const activeQuery = this._getActiveQuery()
|
|---|
| 446 | if (activeQuery == null) {
|
|---|
| 447 | const error = new Error('Received unexpected rowDescription message from backend.')
|
|---|
| 448 | this._handleErrorEvent(error)
|
|---|
| 449 | return
|
|---|
| 450 | }
|
|---|
| 451 | // delegate rowDescription to active query
|
|---|
| 452 | activeQuery.handleRowDescription(msg)
|
|---|
| 453 | }
|
|---|
| 454 |
|
|---|
| 455 | _handleDataRow(msg) {
|
|---|
| 456 | const activeQuery = this._getActiveQuery()
|
|---|
| 457 | if (activeQuery == null) {
|
|---|
| 458 | const error = new Error('Received unexpected dataRow message from backend.')
|
|---|
| 459 | this._handleErrorEvent(error)
|
|---|
| 460 | return
|
|---|
| 461 | }
|
|---|
| 462 | // delegate dataRow to active query
|
|---|
| 463 | activeQuery.handleDataRow(msg)
|
|---|
| 464 | }
|
|---|
| 465 |
|
|---|
| 466 | _handlePortalSuspended(msg) {
|
|---|
| 467 | const activeQuery = this._getActiveQuery()
|
|---|
| 468 | if (activeQuery == null) {
|
|---|
| 469 | const error = new Error('Received unexpected portalSuspended message from backend.')
|
|---|
| 470 | this._handleErrorEvent(error)
|
|---|
| 471 | return
|
|---|
| 472 | }
|
|---|
| 473 | // delegate portalSuspended to active query
|
|---|
| 474 | activeQuery.handlePortalSuspended(this.connection)
|
|---|
| 475 | }
|
|---|
| 476 |
|
|---|
| 477 | _handleEmptyQuery(msg) {
|
|---|
| 478 | const activeQuery = this._getActiveQuery()
|
|---|
| 479 | if (activeQuery == null) {
|
|---|
| 480 | const error = new Error('Received unexpected emptyQuery message from backend.')
|
|---|
| 481 | this._handleErrorEvent(error)
|
|---|
| 482 | return
|
|---|
| 483 | }
|
|---|
| 484 | // delegate emptyQuery to active query
|
|---|
| 485 | activeQuery.handleEmptyQuery(this.connection)
|
|---|
| 486 | }
|
|---|
| 487 |
|
|---|
| 488 | _handleCommandComplete(msg) {
|
|---|
| 489 | const activeQuery = this._getActiveQuery()
|
|---|
| 490 | if (activeQuery == null) {
|
|---|
| 491 | const error = new Error('Received unexpected commandComplete message from backend.')
|
|---|
| 492 | this._handleErrorEvent(error)
|
|---|
| 493 | return
|
|---|
| 494 | }
|
|---|
| 495 | // delegate commandComplete to active query
|
|---|
| 496 | activeQuery.handleCommandComplete(msg, this.connection)
|
|---|
| 497 | }
|
|---|
| 498 |
|
|---|
| 499 | _handleParseComplete() {
|
|---|
| 500 | const activeQuery = this._getActiveQuery()
|
|---|
| 501 | if (activeQuery == null) {
|
|---|
| 502 | const error = new Error('Received unexpected parseComplete message from backend.')
|
|---|
| 503 | this._handleErrorEvent(error)
|
|---|
| 504 | return
|
|---|
| 505 | }
|
|---|
| 506 | // if a prepared statement has a name and properly parses
|
|---|
| 507 | // we track that its already been executed so we don't parse
|
|---|
| 508 | // it again on the same client
|
|---|
| 509 | if (activeQuery.name) {
|
|---|
| 510 | this.connection.parsedStatements[activeQuery.name] = activeQuery.text
|
|---|
| 511 | delete this.connection.submittedNamedStatements[activeQuery.name]
|
|---|
| 512 | }
|
|---|
| 513 | }
|
|---|
| 514 |
|
|---|
| 515 | _handleCopyInResponse(msg) {
|
|---|
| 516 | const activeQuery = this._getActiveQuery()
|
|---|
| 517 | if (activeQuery == null) {
|
|---|
| 518 | const error = new Error('Received unexpected copyInResponse message from backend.')
|
|---|
| 519 | this._handleErrorEvent(error)
|
|---|
| 520 | return
|
|---|
| 521 | }
|
|---|
| 522 | activeQuery.handleCopyInResponse(this.connection)
|
|---|
| 523 | }
|
|---|
| 524 |
|
|---|
| 525 | _handleCopyData(msg) {
|
|---|
| 526 | const activeQuery = this._getActiveQuery()
|
|---|
| 527 | if (activeQuery == null) {
|
|---|
| 528 | const error = new Error('Received unexpected copyData message from backend.')
|
|---|
| 529 | this._handleErrorEvent(error)
|
|---|
| 530 | return
|
|---|
| 531 | }
|
|---|
| 532 | activeQuery.handleCopyData(msg, this.connection)
|
|---|
| 533 | }
|
|---|
| 534 |
|
|---|
| 535 | _handleNotification(msg) {
|
|---|
| 536 | this.emit('notification', msg)
|
|---|
| 537 | }
|
|---|
| 538 |
|
|---|
| 539 | _handleNotice(msg) {
|
|---|
| 540 | this.emit('notice', msg)
|
|---|
| 541 | }
|
|---|
| 542 |
|
|---|
| 543 | getStartupConf() {
|
|---|
| 544 | const params = this.connectionParameters
|
|---|
| 545 |
|
|---|
| 546 | const data = {
|
|---|
| 547 | user: params.user,
|
|---|
| 548 | database: params.database,
|
|---|
| 549 | }
|
|---|
| 550 |
|
|---|
| 551 | const appName = params.application_name || params.fallback_application_name
|
|---|
| 552 | if (appName) {
|
|---|
| 553 | data.application_name = appName
|
|---|
| 554 | }
|
|---|
| 555 | if (params.replication) {
|
|---|
| 556 | data.replication = '' + params.replication
|
|---|
| 557 | }
|
|---|
| 558 | if (params.statement_timeout) {
|
|---|
| 559 | data.statement_timeout = String(parseInt(params.statement_timeout, 10))
|
|---|
| 560 | }
|
|---|
| 561 | if (params.lock_timeout) {
|
|---|
| 562 | data.lock_timeout = String(parseInt(params.lock_timeout, 10))
|
|---|
| 563 | }
|
|---|
| 564 | if (params.idle_in_transaction_session_timeout) {
|
|---|
| 565 | data.idle_in_transaction_session_timeout = String(parseInt(params.idle_in_transaction_session_timeout, 10))
|
|---|
| 566 | }
|
|---|
| 567 | if (params.options) {
|
|---|
| 568 | data.options = params.options
|
|---|
| 569 | }
|
|---|
| 570 |
|
|---|
| 571 | return data
|
|---|
| 572 | }
|
|---|
| 573 |
|
|---|
| 574 | cancel(client, query) {
|
|---|
| 575 | if (client.activeQuery === query) {
|
|---|
| 576 | const con = this.connection
|
|---|
| 577 |
|
|---|
| 578 | if (this.host && this.host.indexOf('/') === 0) {
|
|---|
| 579 | con.connect(this.host + '/.s.PGSQL.' + this.port)
|
|---|
| 580 | } else {
|
|---|
| 581 | con.connect(this.port, this.host)
|
|---|
| 582 | }
|
|---|
| 583 |
|
|---|
| 584 | // once connection is established send cancel message
|
|---|
| 585 | con.on('connect', function () {
|
|---|
| 586 | con.cancel(client.processID, client.secretKey)
|
|---|
| 587 | })
|
|---|
| 588 | } else if (client._queryQueue.indexOf(query) !== -1) {
|
|---|
| 589 | client._queryQueue.splice(client._queryQueue.indexOf(query), 1)
|
|---|
| 590 | } else if (client._sentQueryQueue.indexOf(query) !== -1) {
|
|---|
| 591 | // Query already sent on wire — can't remove it without corrupting the
|
|---|
| 592 | // pipeline. No-op the callback so the result is silently discarded.
|
|---|
| 593 | query.callback = () => {}
|
|---|
| 594 | }
|
|---|
| 595 | }
|
|---|
| 596 |
|
|---|
| 597 | setTypeParser(oid, format, parseFn) {
|
|---|
| 598 | return this._types.setTypeParser(oid, format, parseFn)
|
|---|
| 599 | }
|
|---|
| 600 |
|
|---|
| 601 | getTypeParser(oid, format) {
|
|---|
| 602 | return this._types.getTypeParser(oid, format)
|
|---|
| 603 | }
|
|---|
| 604 |
|
|---|
| 605 | // escapeIdentifier and escapeLiteral moved to utility functions & exported
|
|---|
| 606 | // on PG
|
|---|
| 607 | // re-exported here for backwards compatibility
|
|---|
| 608 | escapeIdentifier(str) {
|
|---|
| 609 | return utils.escapeIdentifier(str)
|
|---|
| 610 | }
|
|---|
| 611 |
|
|---|
| 612 | escapeLiteral(str) {
|
|---|
| 613 | return utils.escapeLiteral(str)
|
|---|
| 614 | }
|
|---|
| 615 |
|
|---|
| 616 | _pulseQueryQueue() {
|
|---|
| 617 | if (this.pipeline) {
|
|---|
| 618 | this._pulsePipelinedQueryQueue()
|
|---|
| 619 | return
|
|---|
| 620 | }
|
|---|
| 621 | if (this.readyForQuery === true) {
|
|---|
| 622 | this._activeQuery = this._queryQueue.shift()
|
|---|
| 623 | const activeQuery = this._getActiveQuery()
|
|---|
| 624 | if (activeQuery) {
|
|---|
| 625 | this.readyForQuery = false
|
|---|
| 626 | this.hasExecuted = true
|
|---|
| 627 |
|
|---|
| 628 | const queryError = activeQuery.submit(this.connection)
|
|---|
| 629 | if (queryError) {
|
|---|
| 630 | process.nextTick(() => {
|
|---|
| 631 | activeQuery.handleError(queryError, this.connection)
|
|---|
| 632 | this.readyForQuery = true
|
|---|
| 633 | this._pulseQueryQueue()
|
|---|
| 634 | })
|
|---|
| 635 | }
|
|---|
| 636 | } else if (this.hasExecuted) {
|
|---|
| 637 | this._activeQuery = null
|
|---|
| 638 | this.emit('drain')
|
|---|
| 639 | }
|
|---|
| 640 | }
|
|---|
| 641 | }
|
|---|
| 642 |
|
|---|
| 643 | _pulsePipelinedQueryQueue() {
|
|---|
| 644 | if (!this._connected || !this._queryable) {
|
|---|
| 645 | return
|
|---|
| 646 | }
|
|---|
| 647 | while (this._queryQueue.length > 0) {
|
|---|
| 648 | const query = this._queryQueue.shift()
|
|---|
| 649 | this.hasExecuted = true
|
|---|
| 650 | const queryError = query.submit(this.connection)
|
|---|
| 651 | if (queryError) {
|
|---|
| 652 | process.nextTick(() => {
|
|---|
| 653 | query.handleError(queryError, this.connection)
|
|---|
| 654 | })
|
|---|
| 655 | continue
|
|---|
| 656 | }
|
|---|
| 657 | this._sentQueryQueue.push(query)
|
|---|
| 658 | }
|
|---|
| 659 | if (this.readyForQuery && !this._activeQuery && this._sentQueryQueue.length > 0) {
|
|---|
| 660 | this._activeQuery = this._sentQueryQueue.shift()
|
|---|
| 661 | this.readyForQuery = false
|
|---|
| 662 | }
|
|---|
| 663 | if (!this._activeQuery && this._sentQueryQueue.length === 0 && this._queryQueue.length === 0 && this.hasExecuted) {
|
|---|
| 664 | this.emit('drain')
|
|---|
| 665 | }
|
|---|
| 666 | }
|
|---|
| 667 |
|
|---|
| 668 | query(config, values, callback) {
|
|---|
| 669 | // can take in strings, config object or query object
|
|---|
| 670 | let query
|
|---|
| 671 | let result
|
|---|
| 672 |
|
|---|
| 673 | if (config == null) {
|
|---|
| 674 | throw new TypeError('Client was passed a null or undefined query')
|
|---|
| 675 | }
|
|---|
| 676 |
|
|---|
| 677 | if (typeof config.submit === 'function') {
|
|---|
| 678 | result = query = config
|
|---|
| 679 | if (!query.callback) {
|
|---|
| 680 | if (typeof values === 'function') {
|
|---|
| 681 | query.callback = values
|
|---|
| 682 | } else if (callback) {
|
|---|
| 683 | query.callback = callback
|
|---|
| 684 | }
|
|---|
| 685 | }
|
|---|
| 686 | } else {
|
|---|
| 687 | query = new Query(config, values, callback)
|
|---|
| 688 | if (!query.callback) {
|
|---|
| 689 | result = new this._Promise((resolve, reject) => {
|
|---|
| 690 | query.callback = (err, res) => (err ? reject(err) : resolve(res))
|
|---|
| 691 | }).catch((err) => {
|
|---|
| 692 | // replace the stack trace that leads to `TCP.onStreamRead` with one that leads back to the
|
|---|
| 693 | // application that created the query
|
|---|
| 694 | Error.captureStackTrace(err)
|
|---|
| 695 | throw err
|
|---|
| 696 | })
|
|---|
| 697 | } else if (typeof query.callback !== 'function') {
|
|---|
| 698 | throw new TypeError('callback is not a function')
|
|---|
| 699 | }
|
|---|
| 700 | }
|
|---|
| 701 |
|
|---|
| 702 | const readTimeout = config.query_timeout || this.connectionParameters.query_timeout
|
|---|
| 703 | if (readTimeout) {
|
|---|
| 704 | const queryCallback = query.callback || (() => {})
|
|---|
| 705 |
|
|---|
| 706 | const readTimeoutTimer = setTimeout(() => {
|
|---|
| 707 | const error = new Error('Query read timeout')
|
|---|
| 708 |
|
|---|
| 709 | process.nextTick(() => {
|
|---|
| 710 | query.handleError(error, this.connection)
|
|---|
| 711 | })
|
|---|
| 712 |
|
|---|
| 713 | queryCallback(error)
|
|---|
| 714 |
|
|---|
| 715 | // we already returned an error,
|
|---|
| 716 | // just do nothing if query completes
|
|---|
| 717 | query.callback = () => {}
|
|---|
| 718 |
|
|---|
| 719 | // Remove from queue (only safe if not yet sent)
|
|---|
| 720 | const index = this._queryQueue.indexOf(query)
|
|---|
| 721 | if (index > -1) {
|
|---|
| 722 | this._queryQueue.splice(index, 1)
|
|---|
| 723 | } else if (this.pipeline) {
|
|---|
| 724 | // Query already sent — the pipeline is blocked until it completes.
|
|---|
| 725 | // Destroy the connection to unblock all remaining pipelined queries.
|
|---|
| 726 | this.connection.stream.destroy()
|
|---|
| 727 | return
|
|---|
| 728 | }
|
|---|
| 729 |
|
|---|
| 730 | this._pulseQueryQueue()
|
|---|
| 731 | }, readTimeout)
|
|---|
| 732 |
|
|---|
| 733 | query.callback = (err, res) => {
|
|---|
| 734 | clearTimeout(readTimeoutTimer)
|
|---|
| 735 | queryCallback(err, res)
|
|---|
| 736 | }
|
|---|
| 737 | }
|
|---|
| 738 |
|
|---|
| 739 | if (this.binary && !query.binary) {
|
|---|
| 740 | query.binary = true
|
|---|
| 741 | }
|
|---|
| 742 |
|
|---|
| 743 | if (query._result && !query._result._types) {
|
|---|
| 744 | query._result._types = this._types
|
|---|
| 745 | }
|
|---|
| 746 |
|
|---|
| 747 | if (!this._queryable) {
|
|---|
| 748 | process.nextTick(() => {
|
|---|
| 749 | query.handleError(new Error('Client has encountered a connection error and is not queryable'), this.connection)
|
|---|
| 750 | })
|
|---|
| 751 | return result
|
|---|
| 752 | }
|
|---|
| 753 |
|
|---|
| 754 | if (this._ending) {
|
|---|
| 755 | process.nextTick(() => {
|
|---|
| 756 | query.handleError(new Error('Client was closed and is not queryable'), this.connection)
|
|---|
| 757 | })
|
|---|
| 758 | return result
|
|---|
| 759 | }
|
|---|
| 760 |
|
|---|
| 761 | if (this._queryQueue.length > 0 && !this.pipeline) {
|
|---|
| 762 | queryQueueLengthDeprecationNotice()
|
|---|
| 763 | }
|
|---|
| 764 | this._queryQueue.push(query)
|
|---|
| 765 | this._pulseQueryQueue()
|
|---|
| 766 | return result
|
|---|
| 767 | }
|
|---|
| 768 |
|
|---|
| 769 | ref() {
|
|---|
| 770 | this.connection.ref()
|
|---|
| 771 | }
|
|---|
| 772 |
|
|---|
| 773 | unref() {
|
|---|
| 774 | this.connection.unref()
|
|---|
| 775 | }
|
|---|
| 776 |
|
|---|
| 777 | getTransactionStatus() {
|
|---|
| 778 | return this._txStatus
|
|---|
| 779 | }
|
|---|
| 780 |
|
|---|
| 781 | end(cb) {
|
|---|
| 782 | this._ending = true
|
|---|
| 783 |
|
|---|
| 784 | // if we have never connected, then end is a noop, callback immediately
|
|---|
| 785 | if (!this.connection._connecting || this._ended) {
|
|---|
| 786 | if (cb) {
|
|---|
| 787 | cb()
|
|---|
| 788 | return
|
|---|
| 789 | } else {
|
|---|
| 790 | return this._Promise.resolve()
|
|---|
| 791 | }
|
|---|
| 792 | }
|
|---|
| 793 |
|
|---|
| 794 | if (!this._queryable) {
|
|---|
| 795 | // socket is dead — force close
|
|---|
| 796 | this.connection.stream.destroy()
|
|---|
| 797 | } else if (
|
|---|
| 798 | this.pipeline &&
|
|---|
| 799 | (this._getActiveQuery() || this._sentQueryQueue.length > 0 || this._queryQueue.length > 0)
|
|---|
| 800 | ) {
|
|---|
| 801 | // pipelined queries are already on the wire (or queued to send) and will
|
|---|
| 802 | // complete normally; wait for drain then do a graceful goodbye
|
|---|
| 803 | this.once('drain', () => this.connection.end())
|
|---|
| 804 | } else if (this._getActiveQuery()) {
|
|---|
| 805 | // non-pipeline: a hung query could block end forever — force disconnect
|
|---|
| 806 | this.connection.stream.destroy()
|
|---|
| 807 | } else {
|
|---|
| 808 | this.connection.end()
|
|---|
| 809 | }
|
|---|
| 810 |
|
|---|
| 811 | if (cb) {
|
|---|
| 812 | this.connection.once('end', cb)
|
|---|
| 813 | } else {
|
|---|
| 814 | return new this._Promise((resolve) => {
|
|---|
| 815 | this.connection.once('end', resolve)
|
|---|
| 816 | })
|
|---|
| 817 | }
|
|---|
| 818 | }
|
|---|
| 819 | get queryQueue() {
|
|---|
| 820 | queryQueueDeprecationNotice()
|
|---|
| 821 | return this._queryQueue
|
|---|
| 822 | }
|
|---|
| 823 | }
|
|---|
| 824 |
|
|---|
| 825 | // expose a Query constructor
|
|---|
| 826 | Client.Query = Query
|
|---|
| 827 |
|
|---|
| 828 | module.exports = Client
|
|---|