· 10 years ago · Mar 30, 2016, 01:57 AM
1import Knex from 'knex'
2import _ from 'lodash'
3import camelize from 'camelize'
4import WaterlineSequel from 'waterline-sequel'
5
6import KnexPostgis from 'knex-postgis'
7import WaterlineError from 'waterline-errors'
8import AdapterError from './error'
9import Util from './util'
10import SpatialUtil from './spatial'
11import SQL from './sql'
12
13const Adapter = {
14
15 identity: 'waterline-postgresql',
16
17 wlSqlOptions: {
18 parameterized: true,
19 caseSensitive: false,
20 escapeCharacter: '"',
21 wlNext: {
22 caseSensitive: true
23 },
24 casting: true,
25 canReturnValues: true,
26 escapeInserts: true,
27 declareDeleteAlias: false
28 },
29
30 /**
31 * Local connections store
32 */
33 connections: new Map(),
34
35 //pkFormat: 'string',
36 syncable: true,
37
38 /**
39 * Adapter default configuration
40 */
41 defaults: {
42 schema: true,
43 debug: process.env.WL_DEBUG || false,
44
45 connection: {
46 host: 'localhost',
47 //user: 'postgres',
48 user: 'pi',
49 //password: 'postgres',
50 password: 'password',
51 //database: 'postgres',
52 database: 'homestatus',
53 port: 5432
54 },
55
56 pool: {
57 min: 1,
58 max: 16,
59 ping (knex, cb) {
60 return knex.query('SELECT 1', cb)
61 },
62 pingTimeout: 10 * 1000,
63 syncInterval: 2 * 1000,
64 idleTimeout: 30 * 1000,
65 acquireTimeout: 300 * 1000
66 }
67 },
68
69 /**
70 * This method runs when a connection is initially registered
71 * at server-start-time. This is the only required method.
72 *
73 * @param {[type]} connection [description]
74 * @param {[type]} collection [description]
75 * @param {Function} cb [description]
76 * @return {[type]} [description]
77 */
78 registerConnection (connection, collections, cb) {
79 if (!connection.identity) {
80 return cb(WaterlineError.adapter.IdentityMissing)
81 }
82 if (Adapter.connections.get(connection.identity)) {
83 return cb(WaterlineError.adapter.IdentityDuplicate)
84 }
85
86 _.defaultsDeep(connection, Adapter.defaults)
87
88 let knex = Knex({
89 client: 'pg',
90 connection: connection.url || connection.connection,
91 pool: connection.pool,
92 debug: process.env.WATERLINE_DEBUG_SQL || connection.debug
93 })
94 let cxn = {
95 identity: connection.identity,
96 schema: Adapter.buildSchema(connection, collections),
97 collections: collections,
98 config: connection,
99 knex: knex,
100 st: KnexPostgis(knex)
101 }
102
103 return Util.initializeConnection(cxn)
104 .then(() => {
105 Adapter.connections.set(connection.identity, cxn)
106 cb()
107 })
108 .catch(cb)
109 },
110
111 /**
112 * Construct the waterline schema for the given connection.
113 *
114 * @param connection
115 * @param collections[]
116 */
117 buildSchema (connection, collections) {
118 return _.chain(collections)
119 .map((model, modelName) => {
120 let definition = _.get(model, [ 'waterline', 'schema', model.identity ])
121 return _.defaultsDeep(definition, {
122 attributes: { },
123 tableName: modelName
124 })
125 })
126 .keyBy('tableName')
127 .value()
128 },
129
130 /**
131 * Return the version of the PostgreSQL server as an array
132 * e.g. for Postgres 9.3.9, return [ '9', '3', '9' ]
133 */
134 getVersion (cxn) {
135 return cxn.knex
136 .raw('select version() as version')
137 .then(({ rows: [row] }) => {
138 return row.version.split(' ')[1].split('.')
139 })
140 },
141
142 /**
143 * Describe a table. List all columns and their properties.
144 *
145 * @param connectionName
146 * @param tableName
147 */
148 describe (connectionName, tableName, cb) {
149 let cxn = Adapter.connections.get(connectionName)
150
151 return cxn.knex(tableName).columnInfo()
152 .then(columnInfo => {
153 if (_.isEmpty(columnInfo)) {
154 return cb()
155 }
156
157 return Adapter._query(cxn, SQL.indexes, [ tableName ])
158 .then(({ rows }) => {
159 _.merge(columnInfo, _.keyBy(camelize(rows), 'columnName'))
160 _.isFunction(cb) && cb(null, columnInfo)
161 })
162 })
163 .catch(AdapterError.wrap(cb))
164 },
165
166 /**
167 * Perform a direct SQL query on the database
168 *
169 * @param connectionName
170 * @param tableName
171 * @param queryString
172 * @param data
173 */
174 query (connectionName, tableName, queryString, args, cb) {
175 let cxn = Adapter.connections.get(connectionName)
176
177 return Adapter._query(cxn, queryString, args)
178 .then((result = { }) => {
179 _.isFunction(cb) && cb(null, result)
180 return result
181 })
182 .catch(AdapterError.wrap(cb))
183 },
184
185 _query (cxn, query, values) {
186 return cxn.knex.raw(Util.toKnexRawQuery(query), Util.castValues(values))
187 .then((result = { }) => result)
188 },
189
190 /**
191 * Create a new table
192 *
193 * @param connectionName
194 * @param tableName
195 * @param definition - the waterline schema definition for model
196 * @param cb
197 */
198 define (connectionName, _tableName, definition, cb) {
199 let cxn = Adapter.connections.get(connectionName)
200 let schema = cxn.collections[_tableName]
201 let tableName = _tableName.substring(0, 63)
202
203 return cxn.knex.schema
204 .hasTable(tableName)
205 .then(exists => {
206 if (exists) return
207
208 return cxn.knex.schema.createTable(tableName, table => {
209 _.each(definition, (definition, attributeName) => {
210 let newColumn = Util.toKnexColumn(table, attributeName, definition, schema, cxn.collections)
211 Util.applyColumnConstraints(newColumn, definition)
212 })
213 Util.applyTableConstraints(table, definition)
214 })
215 })
216 .then(() => {
217 //console.log('created table', tableName, schema)
218 _.isFunction(cb) && cb()
219 })
220 .catch(AdapterError.wrap(cb))
221 },
222
223 /**
224 * Drop a table
225 */
226 drop (connectionName, tableName, relations = [ ], cb = relations) {
227 let cxn = Adapter.connections.get(connectionName)
228
229 return cxn.knex.schema.dropTableIfExists(tableName)
230 .then(() => {
231 return Promise.all(_.map(relations, relation => {
232 return cxn.knex.schema.dropTableIfExists(relation)
233 }))
234 })
235 .then(() => {
236 _.isFunction(cb) && cb()
237 })
238 .catch(AdapterError.wrap(cb))
239 },
240
241 /**
242 * Add a column to a table
243 */
244 addAttribute (connectionName, tableName, attributeName, definition, cb) {
245 let cxn = Adapter.connections.get(connectionName)
246 let schema = cxn.collections[tableName]
247
248 return cxn.knex.schema
249 .table(tableName, table => {
250 let newColumn = Util.toKnexColumn(table, attributeName, definition, schema, cxn.collections)
251 Util.applyColumnConstraints(newColumn, definition)
252 })
253 .then(() => {
254 _.isFunction(cb) && cb()
255 })
256 .catch(AdapterError.wrap(cb))
257 },
258
259 /**
260 * Remove a column from a table
261 */
262 removeAttribute (connectionName, tableName, attributeName, cb) {
263 let cxn = Adapter.connections.get(connectionName)
264
265 return cxn.knex.schema
266 .table(tableName, table => {
267 table.dropColumn(attributeName)
268 })
269 .then(result => {
270 _.isFunction(cb) && cb(null, result)
271 return result
272 })
273 .catch(AdapterError.wrap(cb))
274 },
275
276 /**
277 * Create a new record
278 */
279 create (connectionName, tableName, data, cb) {
280 let cxn = Adapter.connections.get(connectionName)
281 let insertData = Util.sanitize(data, cxn.collections[tableName], cxn)
282 let schema = cxn.collections[tableName]
283 let spatialColumns = SpatialUtil.buildSpatialSelect(schema.definition, tableName, cxn)
284
285 return cxn.knex(tableName)
286 .insert(insertData)
287 .returning([ '*', ...spatialColumns ])
288 .then(rows => {
289 let casted = Util.castResultRows(rows, schema)
290 let result = _.isArray(data) ? casted : casted[0]
291
292 _.isFunction(cb) && cb(null, result)
293 return result
294 })
295 .catch(AdapterError.wrap(cb, null, data))
296 },
297
298 /**
299 * Create multiple records
300 */
301 createEach (connectionName, tableName, records, cb) {
302 // TODO use knex.batchInsert
303 return Adapter.create(connectionName, tableName, records, cb)
304 },
305
306 /**
307 * Update a record
308 */
309 update (connectionName, tableName, options, data, cb) {
310 let cxn = Adapter.connections.get(connectionName)
311 let schema = cxn.collections[tableName]
312 let wlsql = new WaterlineSequel(cxn.schema, Adapter.wlSqlOptions)
313 let spatialColumns = SpatialUtil.getSpatialColumns(schema.definition)
314 let updateData = _.omit(data, _.keys(spatialColumns))
315
316 return new Promise((resolve, reject) => {
317 if (_.isEmpty(data)) {
318 return Adapter.find(connectionName, tableName, options, cb)
319 }
320 resolve(wlsql.update(tableName, options, updateData))
321 })
322 .then(({ query, values }) => {
323 return Adapter._query(cxn, query, values)
324 })
325 .then(({ rows }) => {
326 cb && cb(null, rows)
327 })
328 .catch(AdapterError.wrap(cb, null, data))
329 },
330
331 /**
332 * Destroy a record
333 */
334 destroy (connectionName, tableName, options, cb) {
335 let cxn = Adapter.connections.get(connectionName)
336 let wlsql = new WaterlineSequel(cxn.schema, Adapter.wlSqlOptions)
337
338 return new Promise((resolve, reject) => {
339 resolve(wlsql.destroy(tableName, options))
340 })
341 .then(({ query, values }) => {
342 return Adapter._query(cxn, query, values)
343 })
344 .then(({ rows }) => {
345 cb(null, rows)
346 })
347 .catch(AdapterError.wrap(cb))
348 },
349
350 /**
351 * Populate record associations
352 */
353 join (connectionName, tableName, options, cb) {
354 let cxn = Adapter.connections.get(connectionName)
355
356 return Util.buildKnexJoinQuery (cxn, tableName, options)
357 .then(result => {
358 // return unique records only.
359 // TODO move to SQL
360 _.each(_.reject(options.joins, { select: false }), join => {
361 let alias = Util.getJoinAlias(join)
362 let pk = Adapter.getPrimaryKey(cxn, join.child)
363 let schema = cxn.collections[join.child]
364
365 _.each(result, row => {
366 row[alias] = Util.castResultRows(_.uniqBy(row[alias], pk), schema)
367 })
368 })
369
370 return result
371 })
372 .then(result => {
373 _.isFunction(cb) && cb(null, result)
374 return result
375 })
376 .catch(AdapterError.wrap(cb))
377 },
378
379 /**
380 * Get the primary key column of a table
381 */
382 getPrimaryKey ({ collections }, tableName) {
383 let definition = collections[tableName].definition
384
385
386 if (!definition._pk) {
387 let pk = _.findKey(definition, (attr, name) => {
388 return attr.primaryKey === true
389 })
390 definition._pk = pk || 'id'
391 }
392
393 return definition._pk
394 },
395
396 /**
397 * Find records
398 */
399 find (connectionName, tableName, options, cb) {
400 let cxn = Adapter.connections.get(connectionName)
401 let wlsql = new WaterlineSequel(cxn.schema, Adapter.wlSqlOptions)
402 let schema = cxn.collections[tableName]
403
404 //console.log('find', tableName, options)
405 //console.log('schema types', schema._types)
406
407 return new Promise((resolve, reject) => {
408 resolve(wlsql.find(tableName, options))
409 })
410 .then(({ query: [query], values: [values] }) => {
411 let spatialColumns = SpatialUtil.buildSpatialSelect(schema.definition, tableName, cxn)
412 let fullQuery = Util.addSelectColumns(spatialColumns, query)
413
414 //console.log('fullQuery', fullQuery)
415 //console.log('values', values)
416
417 return Adapter._query(cxn, fullQuery, values)
418 })
419 .then(({ rows }) => {
420 let result = Util.castResultRows(rows, schema)
421 _.isFunction(cb) && cb(null, result)
422 return result
423 })
424 .catch(AdapterError.wrap(cb))
425 },
426
427 /**
428 * Count the number of records
429 */
430 count (connectionName, tableName, options, cb) {
431 let cxn = Adapter.connections.get(connectionName)
432 let wlsql = new WaterlineSequel(cxn.schema, Adapter.wlSqlOptions)
433
434 return new Promise((resolve, reject) => {
435 resolve(wlsql.count(tableName, options))
436 })
437 .then(({ query: [query], values: [values] }) => {
438 return Adapter._query(cxn, query, values)
439 })
440 .then(({ rows: [row] }) => {
441 let count = Number(row.count)
442 _.isFunction(cb) && cb(null, count)
443 return count
444 })
445 .catch(AdapterError.wrap(cb))
446 },
447
448 /**
449 * Run queries inside of a transaction.
450 *
451 * Model.transaction(txn => {
452 * Model.create({ ... }, txn)
453 * .then(newModel => {
454 * return Model.update(..., txn)
455 * })
456 * })
457 * .then(txn.commit)
458 * .catch(txn.rollback)
459 */
460 transaction (connectionName, tableName, cb) {
461 let cxn = Adapter.connections.get(connectionName)
462
463 return new Promise(resolve => {
464 cxn.knex.transaction(txn => {
465 _.isFunction(cb) && cb(null, txn)
466 resolve(txn)
467 })
468 })
469 },
470
471 /**
472 * Invoke a database function, aka "stored procedure"
473 *
474 * @param connectionName
475 * @param tableName
476 * @param procedureName the name of the stored procedure to invoke
477 * @param args An array of arguments to pass to the stored procedure
478 */
479 procedure (connectionName, procedureName, args = [ ], cb = args) {
480 let cxn = Adapter.connections.get(connectionName)
481 let procedure = cxn.storedProcedures[procedureName.toLowerCase()]
482
483 if (!procedure) {
484 let error = new Error(`No stored procedure found with the name ${procedureName}`)
485 return (_.isFunction(cb) ? cb(error) : Promise.reject(error))
486 }
487
488 return procedure.invoke(args)
489 .then(result => {
490 _.isFunction(cb) && cb(null, result)
491 return result
492 })
493 .catch(AdapterError.wrap(cb))
494 },
495
496 /**
497 * Stream query results
498 *
499 * TODO not tested
500 */
501 stream (connectionName, tableName, options, outputStream) {
502 let cxn = Adapter.connections.get(connectionName)
503 let wlsql = new WaterlineSequel(cxn.schema, Adapter.wlSqlOptions)
504
505 return new Promise((resolve, reject) => {
506 resolve(wlsql.find(tableName, options))
507 })
508 .then(({ query: [query], values: [values] }) => {
509 let resultStream = cxn.knex.raw(query, values)
510 resultStream.pipe(outputStream)
511
512 return new Promise((resolve, reject) => {
513 resultStream.on('end', resolve)
514 })
515 })
516 .catch(AdapterError.wrap(cb))
517 },
518
519 /**
520 * Fired when a model is unregistered, typically when the server
521 * is killed. Useful for tearing-down remaining open connections,
522 * etc.
523 *
524 * @param {Function} cb [description]
525 * @return {[type]} [description]
526 */
527 teardown (conn, cb = conn) {
528 let connections = conn ? [ Adapter.connections.get(conn) ] : Adapter.connections.values()
529 let teardownPromises = [ ]
530
531 for (let cxn of connections) {
532 if (!cxn) continue
533
534 teardownPromises.push(cxn.knex.destroy())
535 }
536 return Promise.all(teardownPromises)
537 .then(() => {
538 // only delete connection references after all open sessions are closed
539 for (let cxn of connections) {
540 if (!cxn) continue
541 Adapter.connections.delete(cxn.identity)
542 }
543 cb()
544 })
545 .catch(cb)
546 }
547}
548export default Adapter