node-postgres/lib/connection.js

437 lines
9.6 KiB
JavaScript
Raw Normal View History

var sys = require('sys');
var net = require('net');
var crypto = require('crypto');
var EventEmitter = require('events').EventEmitter;
var utils = require(__dirname + '/utils');
var Writer = require(__dirname + '/writer');
var Connection = function(config) {
EventEmitter.call(this);
config = config || {};
this.stream = config.stream || new net.Stream();
this.lastBuffer = false;
this.lastOffset = 0;
this.buffer = null;
this.offset = null;
this.encoding = 'utf8';
this.parsedStatements = {};
this.writer = new Writer();
};
sys.inherits(Connection, EventEmitter);
var p = Connection.prototype;
2010-10-24 06:36:04 +08:00
p.connect = function(port, host) {
2010-10-24 09:26:24 +08:00
if(this.stream.readyState === 'closed'){
2010-10-24 06:36:04 +08:00
this.stream.connect(port, host);
2010-10-24 09:26:24 +08:00
}
else if(this.stream.readyState == 'open') {
2010-10-24 06:36:04 +08:00
this.emit('connect');
}
2010-10-24 09:26:24 +08:00
var self = this;
2010-10-24 06:36:04 +08:00
this.stream.on('connect', function() {
2010-10-24 06:36:04 +08:00
self.emit('connect');
});
2010-10-24 06:36:04 +08:00
this.stream.on('data', function(buffer) {
self.setBuffer(buffer);
var msg;
while(msg = self.parseMessage()) {
self.emit('message', msg);
self.emit(msg.name, msg);
}
});
2010-10-31 09:23:54 +08:00
this.stream.on('error', function(error) {
self.emit('error', error);
});
2010-10-24 06:36:04 +08:00
};
p.startup = function(config) {
var bodyBuffer = this.writer
2010-10-24 06:36:04 +08:00
.addInt16(3)
.addInt16(0)
.addCString('user')
.addCString(config.user)
.addCString('database')
.addCString(config.database)
.addCString('').flush();
2010-11-01 07:36:35 +08:00
//this message is sent without a code
var length = bodyBuffer.length + 4;
2010-11-01 07:36:35 +08:00
var buffer = new Writer()
.addInt32(length)
.add(bodyBuffer)
.join();
this.stream.write(buffer);
};
p.password = function(password) {
2010-11-01 07:36:35 +08:00
//0x70 = 'p'
this._send(0x70, this.writer.addCString(password));
2010-10-24 08:02:13 +08:00
};
p._send = function(code, writer) {
return this.stream.write(writer.flush(code));
}
2010-11-01 07:36:35 +08:00
var termBuffer = new Buffer([0x58, 0, 0, 0, 4]);
p.end = function() {
2010-11-01 07:36:35 +08:00
var wrote = this.stream.write(termBuffer);
};
p.query = function(text) {
2010-11-01 07:36:35 +08:00
//0x51 = Q
this.stream.write(this.writer.addCString(text).flush(0x51));
};
p.parse = function(query) {
//expect something like this:
// { name: 'queryName',
// text: 'select * from blah',
// types: ['int8', 'bool'] }
//normalize missing query names to allow for null
query.name = query.name || '';
//normalize null type array
query.types = query.types || [];
2010-10-24 13:18:48 +08:00
var len = query.types.length;
var buffer = this.writer
.addCString(query.name) //name of query
.addCString(query.text) //actual query text
2010-10-24 13:18:48 +08:00
.addInt16(len);
for(var i = 0; i < len; i++) {
buffer.addInt32(query.types[i]);
}
2010-10-28 13:27:08 +08:00
2010-11-01 07:36:35 +08:00
//0x50 = 'P'
this._send(0x50, buffer);
return this;
};
p.bind = function(config) {
//normalize config
config = config || {};
2010-10-25 02:46:50 +08:00
config.portal = config.portal || '';
config.statement = config.statement || '';
var values = config.values || [];
var len = values.length;
var buffer = this.writer
2010-10-25 02:46:50 +08:00
.addCString(config.portal)
.addCString(config.statement)
.addInt16(0) //always use default text format
2010-10-25 02:46:50 +08:00
.addInt16(len); //number of parameters
for(var i = 0; i < len; i++) {
var val = values[i];
if(val === null) {
buffer.addInt32(-1);
} else {
val = val.toString();
buffer.addInt32(Buffer.byteLength(val));
buffer.addString(val);
2010-10-25 02:46:50 +08:00
}
}
buffer.addInt16(0); //no format codes, use text
2010-11-01 07:36:35 +08:00
//0x42 = 'B'
this._send(0x42, buffer);
};
2010-10-25 02:46:50 +08:00
p.execute = function(config) {
config = config || {};
config.portal = config.portal || '';
config.rows = config.rows || '';
var buffer = this.writer
2010-10-25 02:46:50 +08:00
.addCString(config.portal)
.addInt32(config.rows);
2010-11-01 07:36:35 +08:00
//0x45 = 'E'
this._send(0x45, buffer);
};
var emptyBuffer = Buffer(0);
p.flush = function() {
2010-11-01 07:36:35 +08:00
//0x48 = 'H'
this._send(0x48,this.writer.add(emptyBuffer));
}
p.sync = function() {
2010-11-01 07:36:35 +08:00
//0x53 = 'S'
this._send(0x53, this.writer.add(emptyBuffer));
};
2010-10-24 08:28:57 +08:00
p.end = function() {
2010-11-01 07:36:35 +08:00
//0x58 = 'X'
this._send(0x58, this.writer.add(emptyBuffer));
2010-10-24 08:28:57 +08:00
};
2010-10-28 13:27:08 +08:00
p.describe = function(msg) {
this._send(0x44, this.writer.addCString(msg.type + (msg.name || '')));
2010-10-28 13:27:08 +08:00
};
//parsing methods
p.setBuffer = function(buffer) {
if(this.lastBuffer) { //we have unfinished biznaz
//need to combine last two buffers
var remaining = this.lastBuffer.length - this.lastOffset;
var combinedBuffer = new Buffer(buffer.length + remaining);
this.lastBuffer.copy(combinedBuffer, 0, this.lastOffset);
buffer.copy(combinedBuffer, remaining, 0);
buffer = combinedBuffer;
}
this.buffer = buffer;
this.offset = 0;
};
p.parseMessage = function() {
var remaining = this.buffer.length - (this.offset);
if(remaining < 5) {
//cannot read id + length without at least 5 bytes
//just abort the read now
this.lastBuffer = this.buffer;
this.lastOffset = this.offset;
return false;
}
2010-11-01 06:58:32 +08:00
//read message id code
var id = this.buffer[this.offset++];
//read message length
var length = this.parseInt32();
2010-11-01 06:58:32 +08:00
if(remaining <= length) {
this.lastBuffer = this.buffer;
//rewind the last 5 bytes we read
this.lastOffset = this.offset-5;
return false;
}
2010-11-01 06:58:32 +08:00
var msg = {
length: length
};
switch(id)
{
case 0x52: //R
2010-11-01 06:58:32 +08:00
msg.name = 'authenticationOk';
return this.parseR(msg);
case 0x53: //S
2010-11-01 06:58:32 +08:00
msg.name = 'parameterStatus';
return this.parseS(msg);
case 0x4b: //K
2010-11-01 06:58:32 +08:00
msg.name = 'backendKeyData';
return this.parseK(msg);
case 0x43: //C
2010-11-01 06:58:32 +08:00
msg.name = 'commandComplete';
return this.parseC(msg);
case 0x5a: //Z
2010-11-01 06:58:32 +08:00
msg.name = 'readyForQuery';
return this.parseZ(msg);
case 0x54: //T
2010-11-01 06:58:32 +08:00
msg.name = 'rowDescription';
return this.parseT(msg);
case 0x44: //D
2010-11-01 06:58:32 +08:00
msg.name = 'dataRow';
return this.parseD(msg);
case 0x45: //E
2010-11-01 06:58:32 +08:00
msg.name = 'error';
return this.parseE(msg);
case 0x4e: //N
2010-11-01 06:58:32 +08:00
msg.name = 'notice';
return this.parseN(msg);
case 0x31: //1
2010-11-01 06:58:32 +08:00
msg.name = 'parseComplete';
return msg;
case 0x32: //2
2010-11-01 06:58:32 +08:00
msg.name = 'bindComplete';
return msg;
case 0x41: //A
2010-11-01 06:58:32 +08:00
msg.name = 'notification';
return this.parseA(msg);
case 0x6e: //n
2010-11-01 06:58:32 +08:00
msg.name = 'noData';
return msg;
case 0x49: //I
2010-11-01 06:58:32 +08:00
msg.name = 'emptyQuery';
return msg;
2010-11-15 07:44:36 +08:00
case 0x73: //s
msg.name = 'portalSuspended';
return msg;
default:
2010-11-01 06:58:32 +08:00
throw new Error("Unrecognized message code " + id);
}
};
p.parseR = function(msg) {
var code = 0;
if(msg.length === 8) {
code = this.parseInt32();
if(code === 3) {
msg.name = 'authenticationCleartextPassword';
}
return msg;
}
if(msg.length === 12) {
code = this.parseInt32();
if(code === 5) { //md5 required
msg.name = 'authenticationMD5Password';
msg.salt = new Buffer(4);
this.buffer.copy(msg.salt, 0, this.offset, this.offset + 4);
this.offset += 4;
return msg;
}
}
throw new Error("Unknown authenticatinOk message type" + sys.inspect(msg));
};
p.parseS = function(msg) {
msg.parameterName = this.parseCString();
msg.parameterValue = this.parseCString();
return msg;
};
p.parseK = function(msg) {
msg.processID = this.parseInt32();
msg.secretKey = this.parseInt32();
return msg;
};
p.parseC = function(msg) {
msg.text = this.parseCString();
return msg;
};
p.parseZ = function(msg) {
msg.status = this.readChar();
return msg;
};
p.parseT = function(msg) {
msg.fieldCount = this.parseInt16();
var fields = [];
for(var i = 0; i < msg.fieldCount; i++){
fields[i] = this.parseField();
}
msg.fields = fields;
return msg;
};
p.parseField = function() {
var field = {
name: this.parseCString(),
tableID: this.parseInt32(),
columnID: this.parseInt16(),
dataTypeID: this.parseInt32(),
dataTypeSize: this.parseInt16(),
dataTypeModifier: this.parseInt32(),
format: this.parseInt16() === 0 ? 'text' : 'binary'
};
return field;
};
p.parseD = function(msg) {
var fieldCount = this.parseInt16();
var fields = [];
for(var i = 0; i < fieldCount; i++) {
var length = this.parseInt32();
fields[i] = (length === -1 ? null : this.readString(length))
};
msg.fieldCount = fieldCount;
msg.fields = fields;
return msg;
};
//parses error
p.parseE = function(msg) {
var fields = {};
var fieldType = this.readString(1);
while(fieldType != '\0') {
fields[fieldType] = this.parseCString();
fieldType = this.readString(1);
}
msg.severity = fields.S;
msg.code = fields.C;
msg.message = fields.M;
msg.detail = fields.D;
msg.hint = fields.H;
msg.position = fields.P;
msg.internalPosition = fields.p;
msg.internalQuery = fields.q;
msg.where = fields.W;
msg.file = fields.F;
msg.line = fields.L;
msg.routine = fields.R;
return msg;
};
2010-10-25 03:43:25 +08:00
//same thing, different name
p.parseN = p.parseE;
2010-10-24 11:31:43 +08:00
p.parseA = function(msg) {
msg.processId = this.parseInt32();
msg.channel = this.parseCString();
2010-10-24 11:45:03 +08:00
msg.payload = this.parseCString();
2010-10-24 11:31:43 +08:00
return msg;
};
p.readChar = function() {
return Buffer([this.buffer[this.offset++]]).toString(this.encoding);
};
p.parseInt32 = function() {
var value = this.peekInt32();
this.offset += 4;
return value;
};
p.peekInt32 = function(offset) {
offset = offset || this.offset;
var buffer = this.buffer;
return ((buffer[offset++] << 24) +
(buffer[offset++] << 16) +
(buffer[offset++] << 8) +
buffer[offset++]);
};
p.parseInt16 = function() {
return ((this.buffer[this.offset++] << 8) +
(this.buffer[this.offset++] << 0));
};
p.readString = function(length) {
return this.buffer.toString(this.encoding, this.offset, (this.offset += length));
};
p.parseCString = function() {
var start = this.offset;
while(this.buffer[this.offset++]) { };
return this.buffer.toString(this.encoding, start, this.offset - 1);
};
//end parsing methods
2010-10-24 06:36:04 +08:00
module.exports = Connection;