2013-10-29 06:31:11 +08:00
|
|
|
var assert = require('assert')
|
|
|
|
var gonna = require('gonna')
|
|
|
|
|
|
|
|
var _ = require('lodash')
|
2014-03-28 16:17:05 +08:00
|
|
|
var async = require('async')
|
|
|
|
var concat = require('concat-stream')
|
2016-05-04 02:20:04 +08:00
|
|
|
var pg = require('pg')
|
2013-10-29 06:31:11 +08:00
|
|
|
|
2014-03-28 16:17:05 +08:00
|
|
|
var copy = require('../').to
|
2016-07-29 06:49:09 +08:00
|
|
|
var code = require('../message-formats')
|
2013-10-29 06:31:11 +08:00
|
|
|
|
2014-03-28 16:17:05 +08:00
|
|
|
var client = function() {
|
|
|
|
var client = new pg.Client()
|
|
|
|
client.connect()
|
|
|
|
return client
|
|
|
|
}
|
2013-10-29 06:31:11 +08:00
|
|
|
|
2014-09-16 03:01:39 +08:00
|
|
|
var testConstruction = function() {
|
|
|
|
var txt = 'COPY (SELECT * FROM generate_series(0, 10)) TO STDOUT'
|
|
|
|
var stream = copy(txt, {highWaterMark: 10})
|
|
|
|
assert.equal(stream._readableState.highWaterMark, 10, 'Client should have been set with a correct highWaterMark.')
|
|
|
|
}
|
|
|
|
|
|
|
|
testConstruction()
|
|
|
|
|
2016-07-29 06:49:09 +08:00
|
|
|
var testComparators = function() {
|
|
|
|
var copy1 = copy();
|
|
|
|
copy1.pipe(concat(function(buf) {
|
|
|
|
assert(copy1._gotCopyOutResponse, 'should have received CopyOutResponse')
|
|
|
|
assert(!copy1._remainder, 'Message with no additional data (len=Int4Len+0) should not leave a remainder')
|
|
|
|
}))
|
|
|
|
copy1.end(new Buffer([code.CopyOutResponse, 0x00, 0x00, 0x00, 0x04]));
|
|
|
|
|
|
|
|
|
|
|
|
}
|
|
|
|
testComparators();
|
|
|
|
|
2014-03-28 16:17:05 +08:00
|
|
|
var testRange = function(top) {
|
|
|
|
var fromClient = client()
|
2013-10-29 10:27:15 +08:00
|
|
|
var txt = 'COPY (SELECT * from generate_series(0, ' + (top - 1) + ')) TO STDOUT'
|
2016-05-04 02:20:04 +08:00
|
|
|
var res;
|
|
|
|
|
2013-10-29 06:31:11 +08:00
|
|
|
|
|
|
|
var stream = fromClient.query(copy(txt))
|
|
|
|
var done = gonna('finish piping out', 1000, function() {
|
|
|
|
fromClient.end()
|
|
|
|
})
|
|
|
|
|
|
|
|
stream.pipe(concat(function(buf) {
|
2016-05-04 02:20:04 +08:00
|
|
|
res = buf.toString('utf8')
|
|
|
|
}))
|
|
|
|
|
|
|
|
stream.on('end', function() {
|
2013-10-29 10:27:15 +08:00
|
|
|
var expected = _.range(0, top).join('\n') + '\n'
|
2013-10-29 06:31:11 +08:00
|
|
|
assert.equal(res, expected)
|
2013-10-29 10:50:45 +08:00
|
|
|
assert.equal(stream.rowCount, top, 'should have rowCount ' + top + ' but got ' + stream.rowCount)
|
2013-10-29 06:31:11 +08:00
|
|
|
done()
|
2016-05-04 02:20:04 +08:00
|
|
|
});
|
2013-10-29 06:31:11 +08:00
|
|
|
}
|
|
|
|
|
|
|
|
testRange(10000)
|
2014-03-28 16:17:05 +08:00
|
|
|
|
2014-04-07 13:02:24 +08:00
|
|
|
var testInternalPostgresError = function() {
|
2014-09-16 12:24:14 +08:00
|
|
|
var cancelClient = client()
|
|
|
|
var queryClient = client()
|
2014-04-07 13:02:24 +08:00
|
|
|
|
|
|
|
var runStream = function(callback) {
|
2014-09-16 12:24:14 +08:00
|
|
|
var txt = "COPY (SELECT pg_sleep(10)) TO STDOUT"
|
|
|
|
var stream = queryClient.query(copy(txt))
|
2014-04-07 13:02:24 +08:00
|
|
|
stream.on('data', function(data) {
|
|
|
|
// Just throw away the data.
|
|
|
|
})
|
|
|
|
stream.on('error', callback)
|
2014-09-16 12:24:14 +08:00
|
|
|
|
|
|
|
setTimeout(function() {
|
|
|
|
var cancelQuery = "SELECT pg_cancel_backend(pid) FROM pg_stat_activity WHERE query ~ 'pg_sleep' AND NOT query ~ 'pg_cancel_backend'"
|
|
|
|
cancelClient.query(cancelQuery)
|
|
|
|
}, 50)
|
2014-04-07 13:02:24 +08:00
|
|
|
}
|
|
|
|
|
|
|
|
runStream(function(err) {
|
|
|
|
assert.notEqual(err, null)
|
2014-09-16 12:24:14 +08:00
|
|
|
var expectedMessage = 'canceling statement due to user request'
|
|
|
|
assert.notEqual(err.toString().indexOf(expectedMessage), -1, 'Error message should mention reason for query failure.')
|
|
|
|
cancelClient.end()
|
|
|
|
queryClient.end()
|
2014-04-07 13:02:24 +08:00
|
|
|
})
|
|
|
|
}
|
|
|
|
testInternalPostgresError()
|
2016-07-29 02:15:59 +08:00
|
|
|
|
|
|
|
var testNoticeResponse = function() {
|
|
|
|
// we use a special trick to generate a warning
|
|
|
|
// on the copy stream.
|
|
|
|
var queryClient = client()
|
|
|
|
var set = '';
|
|
|
|
set += 'SET SESSION client_min_messages = WARNING;'
|
|
|
|
set += 'SET SESSION standard_conforming_strings = off;'
|
|
|
|
set += 'SET SESSION escape_string_warning = on;'
|
|
|
|
queryClient.query(set, function(err, res) {
|
|
|
|
assert.equal(err, null, 'testNoticeResponse - could not SET parameters')
|
|
|
|
var runStream = function(callback) {
|
|
|
|
var txt = "COPY (SELECT '\\\n') TO STDOUT"
|
|
|
|
var stream = queryClient.query(copy(txt))
|
|
|
|
stream.on('data', function(data) {
|
|
|
|
})
|
|
|
|
stream.on('error', callback)
|
|
|
|
|
|
|
|
// make sure stream is pulled from
|
|
|
|
stream.pipe(concat(callback.bind(null,null)))
|
|
|
|
}
|
|
|
|
|
|
|
|
runStream(function(err) {
|
|
|
|
assert.equal(err, null, err)
|
|
|
|
queryClient.end()
|
|
|
|
})
|
|
|
|
|
|
|
|
})
|
|
|
|
}
|
|
|
|
|
|
|
|
testNoticeResponse();
|
|
|
|
|
|
|
|
|