source: node_modules/pg-pool/index.js@ 2d1ec46

main
Last change on this file since 2d1ec46 was 62b2964, checked in by Klimentina Efremova <klimentina08642@…>, 2 weeks ago

Project Handcraft Marketplace

  • Property mode set to 100644
File size: 14.1 KB
Line 
1'use strict'
2const EventEmitter = require('events').EventEmitter
3
4const NOOP = function () {}
5
6const removeWhere = (list, predicate) => {
7 const i = list.findIndex(predicate)
8
9 return i === -1 ? undefined : list.splice(i, 1)[0]
10}
11
12class IdleItem {
13 constructor(client, idleListener, timeoutId) {
14 this.client = client
15 this.idleListener = idleListener
16 this.timeoutId = timeoutId
17 }
18}
19
20class PendingItem {
21 constructor(callback) {
22 this.callback = callback
23 }
24}
25
26function throwOnDoubleRelease() {
27 throw new Error('Release called on client which has already been released to the pool.')
28}
29
30function promisify(Promise, callback) {
31 if (callback) {
32 return { callback: callback, result: undefined }
33 }
34 let rej
35 let res
36 const cb = function (err, client) {
37 err ? rej(err) : res(client)
38 }
39 const result = new Promise(function (resolve, reject) {
40 res = resolve
41 rej = reject
42 }).catch((err) => {
43 // replace the stack trace that leads to `TCP.onStreamRead` with one that leads back to the
44 // application that created the query
45 Error.captureStackTrace(err)
46 throw err
47 })
48 return { callback: cb, result: result }
49}
50
51function makeIdleListener(pool, client) {
52 return function idleListener(err) {
53 err.client = client
54
55 client.removeListener('error', idleListener)
56 client.on('error', () => {
57 pool.log('additional client error after disconnection due to error', err)
58 })
59 pool._remove(client)
60 // TODO - document that once the pool emits an error
61 // the client has already been closed & purged and is unusable
62 pool.emit('error', err, client)
63 }
64}
65
66class Pool extends EventEmitter {
67 constructor(options, Client) {
68 super()
69 this.options = Object.assign({}, options)
70
71 if (options != null && 'password' in options) {
72 // "hiding" the password so it doesn't show up in stack traces
73 // or if the client is console.logged
74 Object.defineProperty(this.options, 'password', {
75 configurable: true,
76 enumerable: false,
77 writable: true,
78 value: options.password,
79 })
80 }
81 if (options != null && options.ssl && options.ssl.key) {
82 // "hiding" the ssl->key so it doesn't show up in stack traces
83 // or if the client is console.logged
84 Object.defineProperty(this.options.ssl, 'key', {
85 enumerable: false,
86 })
87 }
88
89 this.options.max = this.options.max || this.options.poolSize || 10
90 this.options.min = this.options.min || 0
91 this.options.maxUses = this.options.maxUses || Infinity
92 this.options.allowExitOnIdle = this.options.allowExitOnIdle || false
93 this.options.maxLifetimeSeconds = this.options.maxLifetimeSeconds || 0
94 this.log = this.options.log || function () {}
95 this.Client = this.options.Client || Client || require('pg').Client
96 this.Promise = this.options.Promise || global.Promise
97
98 if (typeof this.options.idleTimeoutMillis === 'undefined') {
99 this.options.idleTimeoutMillis = 10000
100 }
101
102 this._clients = []
103 this._idle = []
104 this._expired = new WeakSet()
105 this._pendingQueue = []
106 this._endCallback = undefined
107 this.ending = false
108 this.ended = false
109 }
110
111 _promiseTry(f) {
112 const Promise = this.Promise
113 if (typeof Promise.try === 'function') {
114 return Promise.try(f)
115 }
116 return new Promise((resolve) => resolve(f()))
117 }
118
119 _isFull() {
120 return this._clients.length >= this.options.max
121 }
122
123 _isAboveMin() {
124 return this._clients.length > this.options.min
125 }
126
127 _pulseQueue() {
128 this.log('pulse queue')
129 if (this.ended) {
130 this.log('pulse queue ended')
131 return
132 }
133 if (this.ending) {
134 this.log('pulse queue on ending')
135 if (this._idle.length) {
136 this._idle.slice().map((item) => {
137 this._remove(item.client)
138 })
139 }
140 if (!this._clients.length) {
141 this.ended = true
142 this._endCallback()
143 }
144 return
145 }
146
147 // if we don't have any waiting, do nothing
148 if (!this._pendingQueue.length) {
149 this.log('no queued requests')
150 return
151 }
152 // if we don't have any idle clients and we have no more room do nothing
153 if (!this._idle.length && this._isFull()) {
154 return
155 }
156 const pendingItem = this._pendingQueue.shift()
157 if (this._idle.length) {
158 const idleItem = this._idle.pop()
159 clearTimeout(idleItem.timeoutId)
160 const client = idleItem.client
161 client.ref && client.ref()
162 const idleListener = idleItem.idleListener
163
164 return this._acquireClient(client, pendingItem, idleListener, false)
165 }
166 if (!this._isFull()) {
167 return this.newClient(pendingItem)
168 }
169 throw new Error('unexpected condition')
170 }
171
172 _remove(client, callback) {
173 const removed = removeWhere(this._idle, (item) => item.client === client)
174
175 if (removed !== undefined) {
176 clearTimeout(removed.timeoutId)
177 }
178
179 this._clients = this._clients.filter((c) => c !== client)
180 const context = this
181 client.end(() => {
182 context.emit('remove', client)
183
184 if (typeof callback === 'function') {
185 callback()
186 }
187 })
188 }
189
190 connect(cb) {
191 if (this.ending) {
192 const err = new Error('Cannot use a pool after calling end on the pool')
193 return cb ? cb(err) : this.Promise.reject(err)
194 }
195
196 const response = promisify(this.Promise, cb)
197 const result = response.result
198
199 // if we don't have to connect a new client, don't do so
200 if (this._isFull() || this._idle.length) {
201 // if we have idle clients schedule a pulse immediately
202 if (this._idle.length) {
203 process.nextTick(() => this._pulseQueue())
204 }
205
206 if (!this.options.connectionTimeoutMillis) {
207 this._pendingQueue.push(new PendingItem(response.callback))
208 return result
209 }
210
211 const queueCallback = (err, res, done) => {
212 clearTimeout(tid)
213 response.callback(err, res, done)
214 }
215
216 const pendingItem = new PendingItem(queueCallback)
217
218 // set connection timeout on checking out an existing client
219 const tid = setTimeout(() => {
220 // remove the callback from pending waiters because
221 // we're going to call it with a timeout error
222 removeWhere(this._pendingQueue, (i) => i.callback === queueCallback)
223 pendingItem.timedOut = true
224 response.callback(new Error('timeout exceeded when trying to connect'))
225 }, this.options.connectionTimeoutMillis)
226
227 if (tid.unref) {
228 tid.unref()
229 }
230
231 this._pendingQueue.push(pendingItem)
232 return result
233 }
234
235 this.newClient(new PendingItem(response.callback))
236
237 return result
238 }
239
240 newClient(pendingItem) {
241 const client = new this.Client(this.options)
242 this._clients.push(client)
243 const idleListener = makeIdleListener(this, client)
244
245 this.log('checking client timeout')
246
247 // connection timeout logic
248 let tid
249 let timeoutHit = false
250 if (this.options.connectionTimeoutMillis) {
251 tid = setTimeout(() => {
252 if (client.connection) {
253 this.log('ending client due to timeout')
254 timeoutHit = true
255 client.connection.stream.destroy()
256 } else if (!client.isConnected()) {
257 this.log('ending client due to timeout')
258 timeoutHit = true
259 // force kill the node driver, and let libpq do its teardown
260 client.end()
261 }
262 }, this.options.connectionTimeoutMillis)
263 }
264
265 this.log('connecting new client')
266 client.connect((err) => {
267 if (tid) {
268 clearTimeout(tid)
269 }
270 client.on('error', idleListener)
271 if (err) {
272 this.log('client failed to connect', err)
273 // remove the dead client from our list of clients
274 this._clients = this._clients.filter((c) => c !== client)
275 if (timeoutHit) {
276 err = new Error('Connection terminated due to connection timeout', { cause: err })
277 }
278
279 // this client won’t be released, so move on immediately
280 this._pulseQueue()
281
282 if (!pendingItem.timedOut) {
283 pendingItem.callback(err, undefined, NOOP)
284 }
285 } else {
286 this.log('new client connected')
287
288 if (this.options.onConnect) {
289 this._promiseTry(() => this.options.onConnect(client)).then(
290 () => {
291 this._afterConnect(client, pendingItem, idleListener)
292 },
293 (hookErr) => {
294 this._clients = this._clients.filter((c) => c !== client)
295 client.end(() => {
296 this._pulseQueue()
297 if (!pendingItem.timedOut) {
298 pendingItem.callback(hookErr, undefined, NOOP)
299 }
300 })
301 }
302 )
303 return
304 }
305
306 return this._afterConnect(client, pendingItem, idleListener)
307 }
308 })
309 }
310
311 _afterConnect(client, pendingItem, idleListener) {
312 if (this.options.maxLifetimeSeconds !== 0) {
313 const maxLifetimeTimeout = setTimeout(() => {
314 this.log('ending client due to expired lifetime')
315 this._expired.add(client)
316 const idleIndex = this._idle.findIndex((idleItem) => idleItem.client === client)
317 if (idleIndex !== -1) {
318 this._acquireClient(
319 client,
320 new PendingItem((err, client, clientRelease) => clientRelease()),
321 idleListener,
322 false
323 )
324 }
325 }, this.options.maxLifetimeSeconds * 1000)
326
327 maxLifetimeTimeout.unref()
328 client.once('end', () => clearTimeout(maxLifetimeTimeout))
329 }
330
331 return this._acquireClient(client, pendingItem, idleListener, true)
332 }
333
334 // acquire a client for a pending work item
335 _acquireClient(client, pendingItem, idleListener, isNew) {
336 if (isNew) {
337 this.emit('connect', client)
338 }
339
340 this.emit('acquire', client)
341
342 client.release = this._releaseOnce(client, idleListener)
343
344 client.removeListener('error', idleListener)
345
346 if (!pendingItem.timedOut) {
347 if (isNew && this.options.verify) {
348 this.options.verify(client, (err) => {
349 if (err) {
350 client.release(err)
351 return pendingItem.callback(err, undefined, NOOP)
352 }
353
354 pendingItem.callback(undefined, client, client.release)
355 })
356 } else {
357 pendingItem.callback(undefined, client, client.release)
358 }
359 } else {
360 if (isNew && this.options.verify) {
361 this.options.verify(client, client.release)
362 } else {
363 client.release()
364 }
365 }
366 }
367
368 // returns a function that wraps _release and throws if called more than once
369 _releaseOnce(client, idleListener) {
370 let released = false
371
372 return (err) => {
373 if (released) {
374 throwOnDoubleRelease()
375 }
376
377 released = true
378 this._release(client, idleListener, err)
379 }
380 }
381
382 // release a client back to the poll, include an error
383 // to remove it from the pool
384 _release(client, idleListener, err) {
385 client.on('error', idleListener)
386
387 client._poolUseCount = (client._poolUseCount || 0) + 1
388
389 this.emit('release', err, client)
390
391 // TODO(bmc): expose a proper, public interface _queryable and _ending
392 if (err || this.ending || !client._queryable || client._ending || client._poolUseCount >= this.options.maxUses) {
393 if (client._poolUseCount >= this.options.maxUses) {
394 this.log('remove expended client')
395 }
396
397 return this._remove(client, this._pulseQueue.bind(this))
398 }
399
400 const isExpired = this._expired.has(client)
401 if (isExpired) {
402 this.log('remove expired client')
403 this._expired.delete(client)
404 return this._remove(client, this._pulseQueue.bind(this))
405 }
406
407 // idle timeout
408 let tid
409 if (this.options.idleTimeoutMillis && this._isAboveMin()) {
410 tid = setTimeout(() => {
411 if (this._isAboveMin()) {
412 this.log('remove idle client')
413 this._remove(client, this._pulseQueue.bind(this))
414 }
415 }, this.options.idleTimeoutMillis)
416
417 if (this.options.allowExitOnIdle) {
418 // allow Node to exit if this is all that's left
419 tid.unref()
420 }
421 }
422
423 if (this.options.allowExitOnIdle) {
424 client.unref()
425 }
426
427 this._idle.push(new IdleItem(client, idleListener, tid))
428 this._pulseQueue()
429 }
430
431 query(text, values, cb) {
432 // guard clause against passing a function as the first parameter
433 if (typeof text === 'function') {
434 const response = promisify(this.Promise, text)
435 setImmediate(function () {
436 return response.callback(new Error('Passing a function as the first parameter to pool.query is not supported'))
437 })
438 return response.result
439 }
440
441 // allow plain text query without values, but callback
442 if (typeof values === 'function') {
443 cb = values
444 values = undefined
445 }
446 const response = promisify(this.Promise, cb)
447 cb = response.callback
448
449 this.connect((err, client) => {
450 if (err) {
451 return cb(err)
452 }
453
454 let clientReleased = false
455 const onError = (err) => {
456 if (clientReleased) {
457 return
458 }
459 clientReleased = true
460 client.release(err)
461 cb(err)
462 }
463
464 client.once('error', onError)
465 this.log('dispatching query')
466 try {
467 client.query(text, values, (err, res) => {
468 this.log('query dispatched')
469 client.removeListener('error', onError)
470 if (clientReleased) {
471 return
472 }
473 clientReleased = true
474 client.release(err)
475 if (err) {
476 return cb(err)
477 }
478 return cb(undefined, res)
479 })
480 } catch (err) {
481 client.release(err)
482 return cb(err)
483 }
484 })
485 return response.result
486 }
487
488 end(cb) {
489 this.log('ending')
490 if (this.ending) {
491 const err = new Error('Called end on pool more than once')
492 return cb ? cb(err) : this.Promise.reject(err)
493 }
494 this.ending = true
495 const promised = promisify(this.Promise, cb)
496 this._endCallback = promised.callback
497 this._pulseQueue()
498 return promised.result
499 }
500
501 get waitingCount() {
502 return this._pendingQueue.length
503 }
504
505 get idleCount() {
506 return this._idle.length
507 }
508
509 get expiredCount() {
510 return this._clients.reduce((acc, client) => acc + (this._expired.has(client) ? 1 : 0), 0)
511 }
512
513 get totalCount() {
514 return this._clients.length
515 }
516}
517module.exports = Pool
Note: See TracBrowser for help on using the repository browser.