diff --git a/lib/client.js b/lib/client.js index 46f9eda19..43d20a92c 100644 --- a/lib/client.js +++ b/lib/client.js @@ -108,6 +108,14 @@ p.connect = function(callback) { con.on('copyInResponse', function(msg) { self.activeQuery.streamData(self.connection); }); + con.on('copyOutResponse', function(msg) { + if (self.activeQuery.stream === undefined) { + self.activeQuery._canceledDueToError = new Error('No destination stream defined'); + //canceling query requires creation of new connection + //look for postgres frontend/backend protocol + (new self.constructor({port: self.port, host: self.host})).cancel(self, self.activeQuery); + } + }); con.on('copyData', function (msg) { self.activeQuery.handleCopyFromChunk(msg.chunk); }); diff --git a/lib/native/index.js b/lib/native/index.js index fd405f29f..eb9307d93 100644 --- a/lib/native/index.js +++ b/lib/native/index.js @@ -72,8 +72,8 @@ p.copyTo = function (text) { p.sendCopyFromChunk = function (chunk) { this._sendCopyFromChunk(chunk); }; -p.endCopyFrom = function () { - this._endCopyFrom(); +p.endCopyFrom = function (msg) { + this._endCopyFrom(msg); }; p.query = function(config, values, callback) { var query = (config instanceof NativeQuery) ? config : new NativeQuery(config, values, callback); @@ -134,7 +134,9 @@ p.resumeDrain = function() { }; this._drainPaused = 0; }; - +p.sendCopyFail = function(msg) { + this.endCopyFrom(msg); +}; var clientBuilder = function(config) { config = config || {}; var connection = new Connection(); @@ -198,6 +200,12 @@ var clientBuilder = function(config) { //start to send data from stream connection._activeQuery.streamData(connection); }); + connection.on('copyOutResponse', function(msg) { + if (connection._activeQuery.stream === undefined) { + connection._activeQuery._canceledDueToError = new Error('No destination stream defined'); + (new clientBuilder({port: connection.port, host: connection.host})).cancel(connection, connection._activeQuery); + } + }); connection.on('copyData', function (chunk) { //recieve chunk from connection //move it to stream diff --git a/lib/native/query.js b/lib/native/query.js index 9b334b1a8..c6dc0f9f3 100644 --- a/lib/native/query.js +++ b/lib/native/query.js @@ -26,6 +26,7 @@ var NativeQuery = function(config, values, callback) { this.values[i] = utils.prepareValue(this.values[i]); } } + this._canceledDueToError = false; }; util.inherits(NativeQuery, EventEmitter); @@ -50,6 +51,10 @@ p.handleRow = function(rowData) { }; p.handleError = function(error) { + if (this._canceledDueToError) { + error = this._canceledDueToError; + this._canceledDueToError = false; + } if(this.callback) { this.callback(error); this.callback = null; @@ -59,6 +64,9 @@ p.handleError = function(error) { } p.handleReadyForQuery = function(meta) { + if (this._canceledDueToError) { + return this.handleError(this._canceledDueToError); + } if(meta) { this._result.addCommandComplete(meta); } @@ -68,9 +76,15 @@ p.handleReadyForQuery = function(meta) { this.emit('end', this._result); }; p.streamData = function (connection) { - this.stream.startStreamingToConnection(connection); + if ( this.stream ) this.stream.startStreamingToConnection(connection); + else connection.sendCopyFail('No source stream defined'); }; p.handleCopyFromChunk = function (chunk) { - this.stream.handleChunk(chunk); + if ( this.stream ) { + this.stream.handleChunk(chunk); + } + //if there are no stream (for example when copy to query was sent by + //query method instead of copyTo) error will be handled + //on copyOutResponse event, so silently ignore this error here } module.exports = NativeQuery; diff --git a/lib/query.js b/lib/query.js index 5f22a66b7..7934cd4ec 100644 --- a/lib/query.js +++ b/lib/query.js @@ -25,6 +25,7 @@ var Query = function(config, values, callback) { this._fieldConverters = []; this._result = new Result(); this.isPreparedStatement = false; + this._canceledDueToError = false; EventEmitter.call(this); }; @@ -92,6 +93,9 @@ p.handleCommandComplete = function(msg) { }; p.handleReadyForQuery = function() { + if (this._canceledDueToError) { + return this.handleError(this._canceledDueToError); + } if(this.callback) { this.callback(null, this._result); } @@ -99,6 +103,10 @@ p.handleReadyForQuery = function() { }; p.handleError = function(err) { + if (this._canceledDueToError) { + err = this._canceledDueToError; + this._canceledDueToError = false; + } //if callback supplied do not emit error event as uncaught error //events will bubble up to node process if(this.callback) { @@ -174,10 +182,11 @@ p.streamData = function (connection) { else connection.sendCopyFail('No source stream defined'); }; p.handleCopyFromChunk = function (chunk) { - if ( this.stream ) this.stream.handleChunk(chunk); - else { - // TODO: signal the problem somehow - //this.handleError(new Error('error', 'No destination stream defined')); - } + if ( this.stream ) { + this.stream.handleChunk(chunk); + } + //if there are no stream (for example when copy to query was sent by + //query method instead of copyTo) error will be handled + //on copyOutResponse event, so silently ignore this error here } module.exports = Query; diff --git a/src/binding.cc b/src/binding.cc index 982aa9696..c26adc1cf 100644 --- a/src/binding.cc +++ b/src/binding.cc @@ -248,12 +248,14 @@ class Connection : public ObjectWrap { bool connecting_; bool ioInitialized_; bool copyOutMode_; + bool copyInMode_; Connection () : ObjectWrap () { connection_ = NULL; connecting_ = false; ioInitialized_ = false; copyOutMode_ = false; + copyInMode_ = false; TRACE("Initializing ev watchers"); read_watcher_.data = this; write_watcher_.data = this; @@ -278,8 +280,13 @@ class Connection : public ObjectWrap { EndCopyFrom(const Arguments& args) { HandleScope scope; Connection *self = ObjectWrap::Unwrap(args.This()); + char * error_msg = NULL; + if (args[0]->IsString()) { + error_msg = MallocCString(args[0]); + } //TODO handle errors in some way - self->EndCopyFrom(); + self->EndCopyFrom(error_msg); + free(error_msg); return Undefined(); } @@ -433,23 +440,19 @@ class Connection : public ObjectWrap { if (this->copyOutMode_) { this->HandleCopyOut(); } - if (PQisBusy(connection_) == 0) { + if (!this->copyInMode_ && !this->copyOutMode_ && PQisBusy(connection_) == 0) { PGresult *result; bool didHandleResult = false; while ((result = PQgetResult(connection_))) { - if (PGRES_COPY_IN == PQresultStatus(result)) { - didHandleResult = false; - Emit("copyInResponse"); - PQclear(result); + didHandleResult = HandleResult(result); + PQclear(result); + if(!didHandleResult) { + //this means that we are in copy in or copy out mode + //in this situation PQgetResult will return same + //result untill all data will be read (copy out) or + //until data end notification (copy in) + //and because of this, we need to break cycle break; - } else if (PGRES_COPY_OUT == PQresultStatus(result)) { - PQclear(result); - this->copyOutMode_ = true; - didHandleResult = this->HandleCopyOut(); - } else { - HandleResult(result); - didHandleResult = true; - PQclear(result); } } //might have fired from notification @@ -479,37 +482,29 @@ class Connection : public ObjectWrap { } bool HandleCopyOut () { char * buffer = NULL; - int copied = PQgetCopyData(connection_, &buffer, 1); - if (copied > 0) { - Buffer * chunk = Buffer::New(buffer, copied); + int copied; + Buffer * chunk; + copied = PQgetCopyData(connection_, &buffer, 1); + while (copied > 0) { + chunk = Buffer::New(buffer, copied); Handle node_chunk = chunk->handle_; Emit("copyData", &node_chunk); PQfreemem(buffer); - //result was not handled copmpletely - return false; - } else if (copied == 0) { + copied = PQgetCopyData(connection_, &buffer, 1); + } + if (copied == 0) { //wait for next read ready //result was not handled copmpletely return false; } else if (copied == -1) { - PGresult *result; - //result is handled completely this->copyOutMode_ = false; - if (PQisBusy(connection_) == 0 && (result = PQgetResult(connection_))) { - HandleResult(result); - PQclear(result); - return true; - } else { - return false; - } + return true; } else if (copied == -2) { - //TODO error handling - //result is handled with error - HandleErrorResult(NULL); + this->copyOutMode_ = false; return true; } } - void HandleResult(PGresult* result) + bool HandleResult(PGresult* result) { ExecStatusType status = PQresultStatus(result); switch(status) { @@ -517,14 +512,35 @@ class Connection : public ObjectWrap { { HandleTuplesResult(result); EmitCommandMetaData(result); + return true; } break; case PGRES_FATAL_ERROR: - HandleErrorResult(result); + { + HandleErrorResult(result); + return true; + } break; case PGRES_COMMAND_OK: case PGRES_EMPTY_QUERY: - EmitCommandMetaData(result); + { + EmitCommandMetaData(result); + return true; + } + break; + case PGRES_COPY_IN: + { + this->copyInMode_ = true; + Emit("copyInResponse"); + return false; + } + break; + case PGRES_COPY_OUT: + { + this->copyOutMode_ = true; + Emit("copyOutResponse"); + return this->HandleCopyOut(); + } break; default: printf("YOU SHOULD NEVER SEE THIS! PLEASE OPEN AN ISSUE ON GITHUB! Unrecogized query status: %s\n", PQresStatus(status)); @@ -772,8 +788,9 @@ class Connection : public ObjectWrap { void SendCopyFromChunk(Handle chunk) { PQputCopyData(connection_, Buffer::Data(chunk), Buffer::Length(chunk)); } - void EndCopyFrom() { - PQputCopyEnd(connection_, NULL); + void EndCopyFrom(char * error_msg) { + PQputCopyEnd(connection_, error_msg); + this->copyInMode_ = false; } }; diff --git a/test/integration/client/copy-tests.js b/test/integration/client/copy-tests.js index c0e533506..577ee96a4 100644 --- a/test/integration/client/copy-tests.js +++ b/test/integration/client/copy-tests.js @@ -95,4 +95,66 @@ test('COPY TO, queue queries', function () { }); }); }); +test("COPY TO incorrect usage with large data", function () { + //when many data is loaded from database (and it takes a lot of time) + //there are chance, that query will be canceled before it ends + //but if there are not so much data, cancel message may be + //send after copy query ends + //so we need to test both situations + pg.connect(helper.config, function (error, client) { + assert.equal(error, null, "Failed to connect: " + helper.sys.inspect(error)); + //intentionally incorrect usage of copy. + //this has to report error in standart way, instead of just throwing exception + client.query( + "COPY (SELECT GENERATE_SERIES(1, 10000000)) TO STDOUT WITH CSV", + assert.calls(function (error) { + assert.ok(error, "error should be reported when sending copy to query with query method"); + client.query("SELECT 1", assert.calls(function (error, result) { + assert.isNull(error, "incorrect copy usage should not break connection"); + assert.ok(result, "incorrect copy usage should not break connection"); + pg.end(helper.config); + })); + }) + ); + }); +}); +test("COPY TO incorrect usage with small data", function () { + pg.connect(helper.config, function (error, client) { + assert.equal(error, null, "Failed to connect: " + helper.sys.inspect(error)); + //intentionally incorrect usage of copy. + //this has to report error in standart way, instead of just throwing exception + client.query( + "COPY (SELECT GENERATE_SERIES(1, 1)) TO STDOUT WITH CSV", + assert.calls(function (error) { + assert.ok(error, "error should be reported when sending copy to query with query method"); + client.query("SELECT 1", assert.calls(function (error, result) { + assert.isNull(error, "incorrect copy usage should not break connection"); + assert.ok(result, "incorrect copy usage should not break connection"); + pg.end(helper.config); + })); + }) + ); + }); +}); + +test("COPY FROM incorrect usage", function () { + pg.connect(helper.config, function (error, client) { + assert.equal(error, null, "Failed to connect: " + helper.sys.inspect(error)); + prepareTable(client, function () { + //intentionally incorrect usage of copy. + //this has to report error in standart way, instead of just throwing exception + client.query( + "COPY copy_test from STDIN WITH CSV", + assert.calls(function (error) { + assert.ok(error, "error should be reported when sending copy to query with query method"); + client.query("SELECT 1", assert.calls(function (error, result) { + assert.isNull(error, "incorrect copy usage should not break connection"); + assert.ok(result, "incorrect copy usage should not break connection"); + pg.end(helper.config); + })); + }) + ); + }); + }); +}); diff --git a/test/native/copy-events-tests.js b/test/native/copy-events-tests.js index 0633b5661..76f7e2921 100644 --- a/test/native/copy-events-tests.js +++ b/test/native/copy-events-tests.js @@ -20,6 +20,10 @@ test('COPY FROM events check', function () { test('COPY TO events check', function () { var con = new Client(helper.config), stdoutStream = con.copyTo('COPY person TO STDOUT'); + assert.emits(con, 'copyOutResponse', + function () {}, + "backend should emit copyOutResponse on copyOutResponse message from server" + ); assert.emits(con, 'copyData', function () { }, diff --git a/test/native/copyto-largedata-tests.js b/test/native/copyto-largedata-tests.js new file mode 100644 index 000000000..518514e54 --- /dev/null +++ b/test/native/copyto-largedata-tests.js @@ -0,0 +1,23 @@ +var helper = require(__dirname+"/../test-helper"); +var Client = require(__dirname + "/../../lib/native"); +test("COPY TO large amount of data from postgres", function () { + //there were a bug in native implementation of COPY TO: + //if there were too much data (if we face situation + //when data is not ready while calling PQgetCopyData); + //while loop in Connection::HandleIOEvent becomes infinite + //in such way hanging node, consumes 100% cpu, and making connection unusable + var con = new Client(helper.config), + rowCount = 100000, + stdoutStream = con.copyTo('COPY (select generate_series(1, ' + rowCount + ')) TO STDOUT'); + con.connect(); + stdoutStream.on('data', function () { + rowCount --; + }); + stdoutStream.on('end', function () { + assert.equal(rowCount, 1, "copy to should load exactly requested number of rows" + rowCount); + con.query("SELECT 1", assert.calls(function (error, result) { + assert.ok(!error && result, "loading large amount of data by copy to should not break connection"); + con.end(); + })); + }); +});