| 1 | 'use strict'
|
|---|
| 2 |
|
|---|
| 3 | const { EventEmitter } = require('events')
|
|---|
| 4 |
|
|---|
| 5 | const Result = require('./result')
|
|---|
| 6 | const utils = require('./utils')
|
|---|
| 7 |
|
|---|
| 8 | class Query extends EventEmitter {
|
|---|
| 9 | constructor(config, values, callback) {
|
|---|
| 10 | super()
|
|---|
| 11 |
|
|---|
| 12 | config = utils.normalizeQueryConfig(config, values, callback)
|
|---|
| 13 |
|
|---|
| 14 | this.text = config.text
|
|---|
| 15 | this.values = config.values
|
|---|
| 16 | this.rows = config.rows
|
|---|
| 17 | this.types = config.types
|
|---|
| 18 | this.name = config.name
|
|---|
| 19 | this.queryMode = config.queryMode
|
|---|
| 20 | this.binary = config.binary
|
|---|
| 21 | // use unique portal name each time
|
|---|
| 22 | this.portal = config.portal || ''
|
|---|
| 23 | this.callback = config.callback
|
|---|
| 24 | this._rowMode = config.rowMode
|
|---|
| 25 | if (process.domain && config.callback) {
|
|---|
| 26 | this.callback = process.domain.bind(config.callback)
|
|---|
| 27 | }
|
|---|
| 28 | this._result = new Result(this._rowMode, this.types)
|
|---|
| 29 |
|
|---|
| 30 | // potential for multiple results
|
|---|
| 31 | this._results = this._result
|
|---|
| 32 | this._canceledDueToError = false
|
|---|
| 33 | }
|
|---|
| 34 |
|
|---|
| 35 | requiresPreparation() {
|
|---|
| 36 | if (this.queryMode === 'extended') {
|
|---|
| 37 | return true
|
|---|
| 38 | }
|
|---|
| 39 |
|
|---|
| 40 | // named queries must always be prepared
|
|---|
| 41 | if (this.name) {
|
|---|
| 42 | return true
|
|---|
| 43 | }
|
|---|
| 44 | // always prepare if there are max number of rows expected per
|
|---|
| 45 | // portal execution
|
|---|
| 46 | if (this.rows) {
|
|---|
| 47 | return true
|
|---|
| 48 | }
|
|---|
| 49 | // don't prepare empty text queries
|
|---|
| 50 | if (!this.text) {
|
|---|
| 51 | return false
|
|---|
| 52 | }
|
|---|
| 53 | // prepare if there are values
|
|---|
| 54 | if (!this.values) {
|
|---|
| 55 | return false
|
|---|
| 56 | }
|
|---|
| 57 | return this.values.length > 0
|
|---|
| 58 | }
|
|---|
| 59 |
|
|---|
| 60 | _checkForMultirow() {
|
|---|
| 61 | // if we already have a result with a command property
|
|---|
| 62 | // then we've already executed one query in a multi-statement simple query
|
|---|
| 63 | // turn our results into an array of results
|
|---|
| 64 | if (this._result.command) {
|
|---|
| 65 | if (!Array.isArray(this._results)) {
|
|---|
| 66 | this._results = [this._result]
|
|---|
| 67 | }
|
|---|
| 68 | this._result = new Result(this._rowMode, this._result._types)
|
|---|
| 69 | this._results.push(this._result)
|
|---|
| 70 | }
|
|---|
| 71 | }
|
|---|
| 72 |
|
|---|
| 73 | // associates row metadata from the supplied
|
|---|
| 74 | // message with this query object
|
|---|
| 75 | // metadata used when parsing row results
|
|---|
| 76 | handleRowDescription(msg) {
|
|---|
| 77 | this._checkForMultirow()
|
|---|
| 78 | this._result.addFields(msg.fields)
|
|---|
| 79 | this._accumulateRows = this.callback || !this.listeners('row').length
|
|---|
| 80 | }
|
|---|
| 81 |
|
|---|
| 82 | handleDataRow(msg) {
|
|---|
| 83 | let row
|
|---|
| 84 |
|
|---|
| 85 | if (this._canceledDueToError) {
|
|---|
| 86 | return
|
|---|
| 87 | }
|
|---|
| 88 |
|
|---|
| 89 | try {
|
|---|
| 90 | row = this._result.parseRow(msg.fields)
|
|---|
| 91 | } catch (err) {
|
|---|
| 92 | this._canceledDueToError = err
|
|---|
| 93 | return
|
|---|
| 94 | }
|
|---|
| 95 |
|
|---|
| 96 | this.emit('row', row, this._result)
|
|---|
| 97 | if (this._accumulateRows) {
|
|---|
| 98 | this._result.addRow(row)
|
|---|
| 99 | }
|
|---|
| 100 | }
|
|---|
| 101 |
|
|---|
| 102 | handleCommandComplete(msg, connection) {
|
|---|
| 103 | this._checkForMultirow()
|
|---|
| 104 | this._result.addCommandComplete(msg)
|
|---|
| 105 | // need to sync after each command complete of a prepared statement
|
|---|
| 106 | // if we were using a row count which results in multiple calls to _getRows
|
|---|
| 107 | if (this.rows) {
|
|---|
| 108 | connection.sync()
|
|---|
| 109 | }
|
|---|
| 110 | }
|
|---|
| 111 |
|
|---|
| 112 | // if a named prepared statement is created with empty query text
|
|---|
| 113 | // the backend will send an emptyQuery message but *not* a command complete message
|
|---|
| 114 | // since we pipeline sync immediately after execute we don't need to do anything here
|
|---|
| 115 | // unless we have rows specified, in which case we did not pipeline the initial sync call
|
|---|
| 116 | handleEmptyQuery(connection) {
|
|---|
| 117 | if (this.rows) {
|
|---|
| 118 | connection.sync()
|
|---|
| 119 | }
|
|---|
| 120 | }
|
|---|
| 121 |
|
|---|
| 122 | handleError(err, connection) {
|
|---|
| 123 | // need to sync after error during a prepared statement
|
|---|
| 124 | if (this._canceledDueToError) {
|
|---|
| 125 | err = this._canceledDueToError
|
|---|
| 126 | this._canceledDueToError = false
|
|---|
| 127 | }
|
|---|
| 128 | // if callback supplied do not emit error event as uncaught error
|
|---|
| 129 | // events will bubble up to node process
|
|---|
| 130 | if (this.callback) {
|
|---|
| 131 | return this.callback(err)
|
|---|
| 132 | }
|
|---|
| 133 | this.emit('error', err)
|
|---|
| 134 | }
|
|---|
| 135 |
|
|---|
| 136 | handleReadyForQuery(con) {
|
|---|
| 137 | if (this._canceledDueToError) {
|
|---|
| 138 | return this.handleError(this._canceledDueToError, con)
|
|---|
| 139 | }
|
|---|
| 140 | if (this.callback) {
|
|---|
| 141 | try {
|
|---|
| 142 | this.callback(null, this._results)
|
|---|
| 143 | } catch (err) {
|
|---|
| 144 | process.nextTick(() => {
|
|---|
| 145 | throw err
|
|---|
| 146 | })
|
|---|
| 147 | }
|
|---|
| 148 | }
|
|---|
| 149 | this.emit('end', this._results)
|
|---|
| 150 | }
|
|---|
| 151 |
|
|---|
| 152 | submit(connection) {
|
|---|
| 153 | if (typeof this.text !== 'string' && typeof this.name !== 'string') {
|
|---|
| 154 | return new Error('A query must have either text or a name. Supplying neither is unsupported.')
|
|---|
| 155 | }
|
|---|
| 156 | const previous = connection.parsedStatements[this.name] || connection.submittedNamedStatements[this.name]
|
|---|
| 157 | if (this.text && previous && this.text !== previous) {
|
|---|
| 158 | return new Error(`Prepared statements must be unique - '${this.name}' was used for a different statement`)
|
|---|
| 159 | }
|
|---|
| 160 | if (this.values && !Array.isArray(this.values)) {
|
|---|
| 161 | return new Error('Query values must be an array')
|
|---|
| 162 | }
|
|---|
| 163 | if (this.requiresPreparation()) {
|
|---|
| 164 | // If we're using the extended query protocol we fire off several separate commands
|
|---|
| 165 | // to the backend. On some versions of node & some operating system versions
|
|---|
| 166 | // the network stack writes each message separately instead of buffering them together
|
|---|
| 167 | // causing the client & network to send more slowly. Corking & uncorking the stream
|
|---|
| 168 | // allows node to buffer up the messages internally before sending them all off at once.
|
|---|
| 169 | // note: we're checking for existence of cork/uncork because some versions of streams
|
|---|
| 170 | // might not have this (cloudflare?)
|
|---|
| 171 | connection.stream.cork && connection.stream.cork()
|
|---|
| 172 | try {
|
|---|
| 173 | this.prepare(connection)
|
|---|
| 174 | } finally {
|
|---|
| 175 | // while unlikely for this.prepare to throw, if it does & we don't uncork this stream
|
|---|
| 176 | // this client becomes unresponsive, so put in finally block "just in case"
|
|---|
| 177 | connection.stream.uncork && connection.stream.uncork()
|
|---|
| 178 | }
|
|---|
| 179 | } else {
|
|---|
| 180 | connection.query(this.text)
|
|---|
| 181 | }
|
|---|
| 182 | return null
|
|---|
| 183 | }
|
|---|
| 184 |
|
|---|
| 185 | hasBeenParsed(connection) {
|
|---|
| 186 | return this.name && (connection.parsedStatements[this.name] || connection.submittedNamedStatements[this.name])
|
|---|
| 187 | }
|
|---|
| 188 |
|
|---|
| 189 | handlePortalSuspended(connection) {
|
|---|
| 190 | this._getRows(connection, this.rows)
|
|---|
| 191 | }
|
|---|
| 192 |
|
|---|
| 193 | _getRows(connection, rows) {
|
|---|
| 194 | connection.execute({
|
|---|
| 195 | portal: this.portal,
|
|---|
| 196 | rows: rows,
|
|---|
| 197 | })
|
|---|
| 198 | // if we're not reading pages of rows send the sync command
|
|---|
| 199 | // to indicate the pipeline is finished
|
|---|
| 200 | if (!rows) {
|
|---|
| 201 | connection.sync()
|
|---|
| 202 | } else {
|
|---|
| 203 | // otherwise flush the call out to read more rows
|
|---|
| 204 | connection.flush()
|
|---|
| 205 | }
|
|---|
| 206 | }
|
|---|
| 207 |
|
|---|
| 208 | // http://developer.postgresql.org/pgdocs/postgres/protocol-flow.html#PROTOCOL-FLOW-EXT-QUERY
|
|---|
| 209 | prepare(connection) {
|
|---|
| 210 | // TODO refactor this poor encapsulation
|
|---|
| 211 | if (!this.hasBeenParsed(connection)) {
|
|---|
| 212 | connection.parse({
|
|---|
| 213 | text: this.text,
|
|---|
| 214 | name: this.name,
|
|---|
| 215 | types: this.types,
|
|---|
| 216 | })
|
|---|
| 217 | if (this.name) {
|
|---|
| 218 | connection.submittedNamedStatements[this.name] = this.text
|
|---|
| 219 | }
|
|---|
| 220 | }
|
|---|
| 221 |
|
|---|
| 222 | // because we're mapping user supplied values to
|
|---|
| 223 | // postgres wire protocol compatible values it could
|
|---|
| 224 | // throw an exception, so try/catch this section
|
|---|
| 225 | try {
|
|---|
| 226 | connection.bind({
|
|---|
| 227 | portal: this.portal,
|
|---|
| 228 | statement: this.name,
|
|---|
| 229 | values: this.values,
|
|---|
| 230 | binary: this.binary,
|
|---|
| 231 | valueMapper: utils.prepareValue,
|
|---|
| 232 | })
|
|---|
| 233 | } catch (err) {
|
|---|
| 234 | // we should close parse to avoid leaking connections
|
|---|
| 235 | connection.close({ type: 'S', name: this.name })
|
|---|
| 236 | connection.sync()
|
|---|
| 237 |
|
|---|
| 238 | this.handleError(err, connection)
|
|---|
| 239 | return
|
|---|
| 240 | }
|
|---|
| 241 |
|
|---|
| 242 | connection.describe({
|
|---|
| 243 | type: 'P',
|
|---|
| 244 | name: this.portal || '',
|
|---|
| 245 | })
|
|---|
| 246 |
|
|---|
| 247 | this._getRows(connection, this.rows)
|
|---|
| 248 | }
|
|---|
| 249 |
|
|---|
| 250 | handleCopyInResponse(connection) {
|
|---|
| 251 | connection.sendCopyFail('No source stream defined')
|
|---|
| 252 | }
|
|---|
| 253 |
|
|---|
| 254 | handleCopyData(msg, connection) {
|
|---|
| 255 | // noop
|
|---|
| 256 | }
|
|---|
| 257 | }
|
|---|
| 258 |
|
|---|
| 259 | module.exports = Query
|
|---|