diff --git a/index.js b/index.js index d6f0399..ccac224 100644 --- a/index.js +++ b/index.js @@ -43,7 +43,7 @@ CopyStreamQuery.prototype._flush = function(cb) { var Int32Len = 4; var finBuffer = Buffer([code.CopyDone, 0, 0, 0, Int32Len]) this.push(finBuffer) - cb() + this.cb_flush = cb } CopyStreamQuery.prototype.handleError = function(e) { @@ -62,8 +62,13 @@ CopyStreamQuery.prototype.handleCommandComplete = function(msg) { this.rowCount = parseInt(match[1], 10) } - this.unpipe() - this.emit('end') + // we delay the _flush cb so that the 'end' event is + // triggered after CommandComplete + this.cb_flush() + + // unpipe from connection + this.unpipe(this.connection) + this.connection = null } CopyStreamQuery.prototype.handleReadyForQuery = function() { diff --git a/test/copy-from.js b/test/copy-from.js index b8e6e0a..1765ca5 100644 --- a/test/copy-from.js +++ b/test/copy-from.js @@ -39,6 +39,7 @@ var testRange = function(top) { fromClient.query('SELECT COUNT(*) FROM numbers', function(err, res) { assert.ifError(err) assert.equal(res.rows[0].count, top, 'expected ' + top + ' rows but got ' + res.rows[0].count) + assert.equal(stream.rowCount, top, 'expected ' + top + ' rows but db count is ' + stream.rowCount) //console.log('found ', res.rows.length, 'rows') countDone() var firstRowDone = gonna('have correct result') @@ -54,3 +55,19 @@ var testRange = function(top) { } testRange(1000) + +var testSingleEnd = function() { + var fromClient = client() + fromClient.query('CREATE TEMP TABLE numbers(num int)') + var txt = 'COPY numbers FROM STDIN'; + var stream = fromClient.query(copy(txt)) + var count = 0; + stream.on('end', function() { + count++; + assert(count==1, '`end` Event was triggered ' + count + ' times'); + if (count == 1) fromClient.end(); + }) + stream.end(Buffer('1\n')) + +} +testSingleEnd()