| 1 | 'use strict'
|
|---|
| 2 |
|
|---|
| 3 | const EventEmitter = require('events').EventEmitter
|
|---|
| 4 | const util = require('util')
|
|---|
| 5 | const utils = require('../utils')
|
|---|
| 6 |
|
|---|
| 7 | const NativeQuery = (module.exports = function (config, values, callback) {
|
|---|
| 8 | EventEmitter.call(this)
|
|---|
| 9 | config = utils.normalizeQueryConfig(config, values, callback)
|
|---|
| 10 | this.text = config.text
|
|---|
| 11 | this.values = config.values
|
|---|
| 12 | this.name = config.name
|
|---|
| 13 | this.queryMode = config.queryMode
|
|---|
| 14 | this.callback = config.callback
|
|---|
| 15 | this.state = 'new'
|
|---|
| 16 | this._arrayMode = config.rowMode === 'array'
|
|---|
| 17 |
|
|---|
| 18 | // if the 'row' event is listened for
|
|---|
| 19 | // then emit them as they come in
|
|---|
| 20 | // without setting singleRowMode to true
|
|---|
| 21 | // this has almost no meaning because libpq
|
|---|
| 22 | // reads all rows into memory before returning any
|
|---|
| 23 | this._emitRowEvents = false
|
|---|
| 24 | this.on(
|
|---|
| 25 | 'newListener',
|
|---|
| 26 | function (event) {
|
|---|
| 27 | if (event === 'row') this._emitRowEvents = true
|
|---|
| 28 | }.bind(this)
|
|---|
| 29 | )
|
|---|
| 30 | })
|
|---|
| 31 |
|
|---|
| 32 | util.inherits(NativeQuery, EventEmitter)
|
|---|
| 33 |
|
|---|
| 34 | const errorFieldMap = {
|
|---|
| 35 | sqlState: 'code',
|
|---|
| 36 | statementPosition: 'position',
|
|---|
| 37 | messagePrimary: 'message',
|
|---|
| 38 | context: 'where',
|
|---|
| 39 | schemaName: 'schema',
|
|---|
| 40 | tableName: 'table',
|
|---|
| 41 | columnName: 'column',
|
|---|
| 42 | dataTypeName: 'dataType',
|
|---|
| 43 | constraintName: 'constraint',
|
|---|
| 44 | sourceFile: 'file',
|
|---|
| 45 | sourceLine: 'line',
|
|---|
| 46 | sourceFunction: 'routine',
|
|---|
| 47 | }
|
|---|
| 48 |
|
|---|
| 49 | NativeQuery.prototype.handleError = function (err) {
|
|---|
| 50 | // copy pq error fields into the error object
|
|---|
| 51 | const fields = this.native && this.native.pq.resultErrorFields()
|
|---|
| 52 | if (fields) {
|
|---|
| 53 | for (const key in fields) {
|
|---|
| 54 | const normalizedFieldName = errorFieldMap[key] || key
|
|---|
| 55 | err[normalizedFieldName] = fields[key]
|
|---|
| 56 | }
|
|---|
| 57 | }
|
|---|
| 58 | if (this.callback) {
|
|---|
| 59 | this.callback(err)
|
|---|
| 60 | } else {
|
|---|
| 61 | this.emit('error', err)
|
|---|
| 62 | }
|
|---|
| 63 | this.state = 'error'
|
|---|
| 64 | }
|
|---|
| 65 |
|
|---|
| 66 | NativeQuery.prototype.then = function (onSuccess, onFailure) {
|
|---|
| 67 | return this._getPromise().then(onSuccess, onFailure)
|
|---|
| 68 | }
|
|---|
| 69 |
|
|---|
| 70 | NativeQuery.prototype.catch = function (callback) {
|
|---|
| 71 | return this._getPromise().catch(callback)
|
|---|
| 72 | }
|
|---|
| 73 |
|
|---|
| 74 | NativeQuery.prototype._getPromise = function () {
|
|---|
| 75 | if (this._promise) return this._promise
|
|---|
| 76 | this._promise = new Promise(
|
|---|
| 77 | function (resolve, reject) {
|
|---|
| 78 | this._once('end', resolve)
|
|---|
| 79 | this._once('error', reject)
|
|---|
| 80 | }.bind(this)
|
|---|
| 81 | )
|
|---|
| 82 | return this._promise
|
|---|
| 83 | }
|
|---|
| 84 |
|
|---|
| 85 | NativeQuery.prototype.submit = function (client) {
|
|---|
| 86 | this.state = 'running'
|
|---|
| 87 | const self = this
|
|---|
| 88 | this.native = client.native
|
|---|
| 89 | client.native.arrayMode = this._arrayMode
|
|---|
| 90 |
|
|---|
| 91 | let after = function (err, rows, results) {
|
|---|
| 92 | client.native.arrayMode = false
|
|---|
| 93 | setImmediate(function () {
|
|---|
| 94 | self.emit('_done')
|
|---|
| 95 | })
|
|---|
| 96 |
|
|---|
| 97 | // handle possible query error
|
|---|
| 98 | if (err) {
|
|---|
| 99 | return self.handleError(err)
|
|---|
| 100 | }
|
|---|
| 101 |
|
|---|
| 102 | // emit row events for each row in the result
|
|---|
| 103 | if (self._emitRowEvents) {
|
|---|
| 104 | if (results.length > 1) {
|
|---|
| 105 | rows.forEach((rowOfRows, i) => {
|
|---|
| 106 | rowOfRows.forEach((row) => {
|
|---|
| 107 | self.emit('row', row, results[i])
|
|---|
| 108 | })
|
|---|
| 109 | })
|
|---|
| 110 | } else {
|
|---|
| 111 | rows.forEach(function (row) {
|
|---|
| 112 | self.emit('row', row, results)
|
|---|
| 113 | })
|
|---|
| 114 | }
|
|---|
| 115 | }
|
|---|
| 116 |
|
|---|
| 117 | // handle successful result
|
|---|
| 118 | self.state = 'end'
|
|---|
| 119 | self.emit('end', results)
|
|---|
| 120 | if (self.callback) {
|
|---|
| 121 | self.callback(null, results)
|
|---|
| 122 | }
|
|---|
| 123 | }
|
|---|
| 124 |
|
|---|
| 125 | if (process.domain) {
|
|---|
| 126 | after = process.domain.bind(after)
|
|---|
| 127 | }
|
|---|
| 128 |
|
|---|
| 129 | // named query
|
|---|
| 130 | if (this.name) {
|
|---|
| 131 | if (this.name.length > 63) {
|
|---|
| 132 | console.error('Warning! Postgres only supports 63 characters for query names.')
|
|---|
| 133 | console.error('You supplied %s (%s)', this.name, this.name.length)
|
|---|
| 134 | console.error('This can cause conflicts and silent errors executing queries')
|
|---|
| 135 | }
|
|---|
| 136 | const values = (this.values || []).map(utils.prepareValue)
|
|---|
| 137 |
|
|---|
| 138 | // check if the client has already executed this named query
|
|---|
| 139 | // if so...just execute it again - skip the planning phase
|
|---|
| 140 | if (client.namedQueries[this.name]) {
|
|---|
| 141 | if (this.text && client.namedQueries[this.name] !== this.text) {
|
|---|
| 142 | const err = new Error(`Prepared statements must be unique - '${this.name}' was used for a different statement`)
|
|---|
| 143 | return after(err)
|
|---|
| 144 | }
|
|---|
| 145 | return client.native.execute(this.name, values, after)
|
|---|
| 146 | }
|
|---|
| 147 | // plan the named query the first time, then execute it
|
|---|
| 148 | return client.native.prepare(this.name, this.text, values.length, function (err) {
|
|---|
| 149 | if (err) return after(err)
|
|---|
| 150 | client.namedQueries[self.name] = self.text
|
|---|
| 151 | return self.native.execute(self.name, values, after)
|
|---|
| 152 | })
|
|---|
| 153 | } else if (this.values) {
|
|---|
| 154 | if (!Array.isArray(this.values)) {
|
|---|
| 155 | const err = new Error('Query values must be an array')
|
|---|
| 156 | return after(err)
|
|---|
| 157 | }
|
|---|
| 158 | const vals = this.values.map(utils.prepareValue)
|
|---|
| 159 | client.native.query(this.text, vals, after)
|
|---|
| 160 | } else if (this.queryMode === 'extended') {
|
|---|
| 161 | client.native.query(this.text, [], after)
|
|---|
| 162 | } else {
|
|---|
| 163 | client.native.query(this.text, after)
|
|---|
| 164 | }
|
|---|
| 165 | }
|
|---|