source: node_modules/pg/lib/native/client.js

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

Project Handcraft Marketplace

  • Property mode set to 100644
File size: 10.8 KB
RevLine 
[62b2964]1const nodeUtils = require('util')
2// eslint-disable-next-line
3var Native
4// eslint-disable-next-line no-useless-catch
5try {
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}
11const TypeOverrides = require('../type-overrides')
12const EventEmitter = require('events').EventEmitter
13const util = require('util')
14const ConnectionParameters = require('../connection-parameters')
15
16const NativeQuery = require('./query')
17
18const 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
23const 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
64Client.Query = NativeQuery
65
66util.inherits(Client, EventEmitter)
67
68Client.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
88Client.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
133Client.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// }
160Client.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
249Client.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
288Client.prototype._hasActiveQuery = function () {
289 return this._activeQuery && this._activeQuery.state !== 'error' && this._activeQuery.state !== 'end'
290}
291
292Client.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
317Client.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
395Client.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
403Client.prototype.ref = function () {}
404Client.prototype.unref = function () {}
405
406Client.prototype.setTypeParser = function (oid, format, parseFn) {
407 return this._types.setTypeParser(oid, format, parseFn)
408}
409
410Client.prototype.getTypeParser = function (oid, format) {
411 return this._types.getTypeParser(oid, format)
412}
413
414Client.prototype.isConnected = function () {
415 return this._connected
416}
417
418Client.prototype.getTransactionStatus = function () {
419 return this.native.getTransactionStatus()
420}
Note: See TracBrowser for help on using the repository browser.