| [62b2964] | 1 | "use strict";
|
|---|
| 2 | Object.defineProperty(exports, "__esModule", { value: true });
|
|---|
| 3 | exports.Parser = void 0;
|
|---|
| 4 | const messages_1 = require("./messages");
|
|---|
| 5 | const buffer_reader_1 = require("./buffer-reader");
|
|---|
| 6 | // every message is prefixed with a single byte
|
|---|
| 7 | const CODE_LENGTH = 1;
|
|---|
| 8 | // every message has an int32 length which includes itself but does
|
|---|
| 9 | // NOT include the code in the length
|
|---|
| 10 | const LEN_LENGTH = 4;
|
|---|
| 11 | const HEADER_LENGTH = CODE_LENGTH + LEN_LENGTH;
|
|---|
| 12 | // A placeholder for a `BackendMessage`’s length value that will be set after construction.
|
|---|
| 13 | const LATEINIT_LENGTH = -1;
|
|---|
| 14 | const emptyBuffer = Buffer.allocUnsafe(0);
|
|---|
| 15 | class Parser {
|
|---|
| 16 | constructor(opts) {
|
|---|
| 17 | this.buffer = emptyBuffer;
|
|---|
| 18 | this.bufferLength = 0;
|
|---|
| 19 | this.bufferOffset = 0;
|
|---|
| 20 | this.reader = new buffer_reader_1.BufferReader();
|
|---|
| 21 | if ((opts === null || opts === void 0 ? void 0 : opts.mode) === 'binary') {
|
|---|
| 22 | throw new Error('Binary mode not supported yet');
|
|---|
| 23 | }
|
|---|
| 24 | this.mode = (opts === null || opts === void 0 ? void 0 : opts.mode) || 'text';
|
|---|
| 25 | }
|
|---|
| 26 | parse(buffer, callback) {
|
|---|
| 27 | this.mergeBuffer(buffer);
|
|---|
| 28 | const bufferFullLength = this.bufferOffset + this.bufferLength;
|
|---|
| 29 | let offset = this.bufferOffset;
|
|---|
| 30 | while (offset + HEADER_LENGTH <= bufferFullLength) {
|
|---|
| 31 | // code is 1 byte long - it identifies the message type
|
|---|
| 32 | const code = this.buffer[offset];
|
|---|
| 33 | // length is 1 Uint32BE - it is the length of the message EXCLUDING the code
|
|---|
| 34 | const length = this.buffer.readUInt32BE(offset + CODE_LENGTH);
|
|---|
| 35 | const fullMessageLength = CODE_LENGTH + length;
|
|---|
| 36 | if (fullMessageLength + offset <= bufferFullLength) {
|
|---|
| 37 | const message = this.handlePacket(offset + HEADER_LENGTH, code, length, this.buffer);
|
|---|
| 38 | callback(message);
|
|---|
| 39 | offset += fullMessageLength;
|
|---|
| 40 | }
|
|---|
| 41 | else {
|
|---|
| 42 | break;
|
|---|
| 43 | }
|
|---|
| 44 | }
|
|---|
| 45 | if (offset === bufferFullLength) {
|
|---|
| 46 | // No more use for the buffer
|
|---|
| 47 | this.buffer = emptyBuffer;
|
|---|
| 48 | this.bufferLength = 0;
|
|---|
| 49 | this.bufferOffset = 0;
|
|---|
| 50 | }
|
|---|
| 51 | else {
|
|---|
| 52 | // Adjust the cursors of remainingBuffer
|
|---|
| 53 | this.bufferLength = bufferFullLength - offset;
|
|---|
| 54 | this.bufferOffset = offset;
|
|---|
| 55 | }
|
|---|
| 56 | }
|
|---|
| 57 | mergeBuffer(buffer) {
|
|---|
| 58 | if (this.bufferLength > 0) {
|
|---|
| 59 | const newLength = this.bufferLength + buffer.byteLength;
|
|---|
| 60 | const newFullLength = newLength + this.bufferOffset;
|
|---|
| 61 | if (newFullLength > this.buffer.byteLength) {
|
|---|
| 62 | // We can't concat the new buffer with the remaining one
|
|---|
| 63 | let newBuffer;
|
|---|
| 64 | if (newLength <= this.buffer.byteLength && this.bufferOffset >= this.bufferLength) {
|
|---|
| 65 | // We can move the relevant part to the beginning of the buffer instead of allocating a new buffer
|
|---|
| 66 | newBuffer = this.buffer;
|
|---|
| 67 | }
|
|---|
| 68 | else {
|
|---|
| 69 | // Allocate a new larger buffer
|
|---|
| 70 | let newBufferLength = this.buffer.byteLength * 2;
|
|---|
| 71 | while (newLength >= newBufferLength) {
|
|---|
| 72 | newBufferLength *= 2;
|
|---|
| 73 | }
|
|---|
| 74 | newBuffer = Buffer.allocUnsafe(newBufferLength);
|
|---|
| 75 | }
|
|---|
| 76 | // Move the remaining buffer to the new one
|
|---|
| 77 | this.buffer.copy(newBuffer, 0, this.bufferOffset, this.bufferOffset + this.bufferLength);
|
|---|
| 78 | this.buffer = newBuffer;
|
|---|
| 79 | this.bufferOffset = 0;
|
|---|
| 80 | }
|
|---|
| 81 | // Concat the new buffer with the remaining one
|
|---|
| 82 | buffer.copy(this.buffer, this.bufferOffset + this.bufferLength);
|
|---|
| 83 | this.bufferLength = newLength;
|
|---|
| 84 | }
|
|---|
| 85 | else {
|
|---|
| 86 | this.buffer = buffer;
|
|---|
| 87 | this.bufferOffset = 0;
|
|---|
| 88 | this.bufferLength = buffer.byteLength;
|
|---|
| 89 | }
|
|---|
| 90 | }
|
|---|
| 91 | handlePacket(offset, code, length, bytes) {
|
|---|
| 92 | const { reader } = this;
|
|---|
| 93 | // NOTE: This undesirably retains the buffer in `this.reader` if the `parse*Message` calls below throw. However, those should only throw in the case of a protocol error, which normally results in the reader being discarded.
|
|---|
| 94 | reader.setBuffer(offset, bytes);
|
|---|
| 95 | let message;
|
|---|
| 96 | switch (code) {
|
|---|
| 97 | case 50 /* MessageCodes.BindComplete */:
|
|---|
| 98 | message = messages_1.bindComplete;
|
|---|
| 99 | break;
|
|---|
| 100 | case 49 /* MessageCodes.ParseComplete */:
|
|---|
| 101 | message = messages_1.parseComplete;
|
|---|
| 102 | break;
|
|---|
| 103 | case 51 /* MessageCodes.CloseComplete */:
|
|---|
| 104 | message = messages_1.closeComplete;
|
|---|
| 105 | break;
|
|---|
| 106 | case 110 /* MessageCodes.NoData */:
|
|---|
| 107 | message = messages_1.noData;
|
|---|
| 108 | break;
|
|---|
| 109 | case 115 /* MessageCodes.PortalSuspended */:
|
|---|
| 110 | message = messages_1.portalSuspended;
|
|---|
| 111 | break;
|
|---|
| 112 | case 99 /* MessageCodes.CopyDone */:
|
|---|
| 113 | message = messages_1.copyDone;
|
|---|
| 114 | break;
|
|---|
| 115 | case 87 /* MessageCodes.ReplicationStart */:
|
|---|
| 116 | message = messages_1.replicationStart;
|
|---|
| 117 | break;
|
|---|
| 118 | case 73 /* MessageCodes.EmptyQuery */:
|
|---|
| 119 | message = messages_1.emptyQuery;
|
|---|
| 120 | break;
|
|---|
| 121 | case 68 /* MessageCodes.DataRow */:
|
|---|
| 122 | message = parseDataRowMessage(reader);
|
|---|
| 123 | break;
|
|---|
| 124 | case 67 /* MessageCodes.CommandComplete */:
|
|---|
| 125 | message = parseCommandCompleteMessage(reader);
|
|---|
| 126 | break;
|
|---|
| 127 | case 90 /* MessageCodes.ReadyForQuery */:
|
|---|
| 128 | message = parseReadyForQueryMessage(reader);
|
|---|
| 129 | break;
|
|---|
| 130 | case 65 /* MessageCodes.NotificationResponse */:
|
|---|
| 131 | message = parseNotificationMessage(reader);
|
|---|
| 132 | break;
|
|---|
| 133 | case 82 /* MessageCodes.AuthenticationResponse */:
|
|---|
| 134 | message = parseAuthenticationResponse(reader, length);
|
|---|
| 135 | break;
|
|---|
| 136 | case 83 /* MessageCodes.ParameterStatus */:
|
|---|
| 137 | message = parseParameterStatusMessage(reader);
|
|---|
| 138 | break;
|
|---|
| 139 | case 75 /* MessageCodes.BackendKeyData */:
|
|---|
| 140 | message = parseBackendKeyData(reader);
|
|---|
| 141 | break;
|
|---|
| 142 | case 69 /* MessageCodes.ErrorMessage */:
|
|---|
| 143 | message = parseErrorMessage(reader, 'error');
|
|---|
| 144 | break;
|
|---|
| 145 | case 78 /* MessageCodes.NoticeMessage */:
|
|---|
| 146 | message = parseErrorMessage(reader, 'notice');
|
|---|
| 147 | break;
|
|---|
| 148 | case 84 /* MessageCodes.RowDescriptionMessage */:
|
|---|
| 149 | message = parseRowDescriptionMessage(reader);
|
|---|
| 150 | break;
|
|---|
| 151 | case 116 /* MessageCodes.ParameterDescriptionMessage */:
|
|---|
| 152 | message = parseParameterDescriptionMessage(reader);
|
|---|
| 153 | break;
|
|---|
| 154 | case 71 /* MessageCodes.CopyIn */:
|
|---|
| 155 | message = parseCopyInMessage(reader);
|
|---|
| 156 | break;
|
|---|
| 157 | case 72 /* MessageCodes.CopyOut */:
|
|---|
| 158 | message = parseCopyOutMessage(reader);
|
|---|
| 159 | break;
|
|---|
| 160 | case 100 /* MessageCodes.CopyData */:
|
|---|
| 161 | message = parseCopyData(reader, length);
|
|---|
| 162 | break;
|
|---|
| 163 | default:
|
|---|
| 164 | return new messages_1.DatabaseError('received invalid response: ' + code.toString(16), length, 'error');
|
|---|
| 165 | }
|
|---|
| 166 | reader.setBuffer(0, emptyBuffer);
|
|---|
| 167 | message.length = length;
|
|---|
| 168 | return message;
|
|---|
| 169 | }
|
|---|
| 170 | }
|
|---|
| 171 | exports.Parser = Parser;
|
|---|
| 172 | const parseReadyForQueryMessage = (reader) => {
|
|---|
| 173 | const status = reader.string(1);
|
|---|
| 174 | return new messages_1.ReadyForQueryMessage(LATEINIT_LENGTH, status);
|
|---|
| 175 | };
|
|---|
| 176 | const parseCommandCompleteMessage = (reader) => {
|
|---|
| 177 | const text = reader.cstring();
|
|---|
| 178 | return new messages_1.CommandCompleteMessage(LATEINIT_LENGTH, text);
|
|---|
| 179 | };
|
|---|
| 180 | const parseCopyData = (reader, length) => {
|
|---|
| 181 | const chunk = reader.bytes(length - 4);
|
|---|
| 182 | return new messages_1.CopyDataMessage(LATEINIT_LENGTH, chunk);
|
|---|
| 183 | };
|
|---|
| 184 | const parseCopyInMessage = (reader) => parseCopyMessage(reader, 'copyInResponse');
|
|---|
| 185 | const parseCopyOutMessage = (reader) => parseCopyMessage(reader, 'copyOutResponse');
|
|---|
| 186 | const parseCopyMessage = (reader, messageName) => {
|
|---|
| 187 | const isBinary = reader.byte() !== 0;
|
|---|
| 188 | const columnCount = reader.int16();
|
|---|
| 189 | const message = new messages_1.CopyResponse(LATEINIT_LENGTH, messageName, isBinary, columnCount);
|
|---|
| 190 | for (let i = 0; i < columnCount; i++) {
|
|---|
| 191 | message.columnTypes[i] = reader.int16();
|
|---|
| 192 | }
|
|---|
| 193 | return message;
|
|---|
| 194 | };
|
|---|
| 195 | const parseNotificationMessage = (reader) => {
|
|---|
| 196 | const processId = reader.int32();
|
|---|
| 197 | const channel = reader.cstring();
|
|---|
| 198 | const payload = reader.cstring();
|
|---|
| 199 | return new messages_1.NotificationResponseMessage(LATEINIT_LENGTH, processId, channel, payload);
|
|---|
| 200 | };
|
|---|
| 201 | const parseRowDescriptionMessage = (reader) => {
|
|---|
| 202 | const fieldCount = reader.int16();
|
|---|
| 203 | const message = new messages_1.RowDescriptionMessage(LATEINIT_LENGTH, fieldCount);
|
|---|
| 204 | for (let i = 0; i < fieldCount; i++) {
|
|---|
| 205 | message.fields[i] = parseField(reader);
|
|---|
| 206 | }
|
|---|
| 207 | return message;
|
|---|
| 208 | };
|
|---|
| 209 | const parseField = (reader) => {
|
|---|
| 210 | const name = reader.cstring();
|
|---|
| 211 | const tableID = reader.uint32();
|
|---|
| 212 | const columnID = reader.int16();
|
|---|
| 213 | const dataTypeID = reader.uint32();
|
|---|
| 214 | const dataTypeSize = reader.int16();
|
|---|
| 215 | const dataTypeModifier = reader.int32();
|
|---|
| 216 | const mode = reader.int16() === 0 ? 'text' : 'binary';
|
|---|
| 217 | return new messages_1.Field(name, tableID, columnID, dataTypeID, dataTypeSize, dataTypeModifier, mode);
|
|---|
| 218 | };
|
|---|
| 219 | const parseParameterDescriptionMessage = (reader) => {
|
|---|
| 220 | const parameterCount = reader.int16();
|
|---|
| 221 | const message = new messages_1.ParameterDescriptionMessage(LATEINIT_LENGTH, parameterCount);
|
|---|
| 222 | for (let i = 0; i < parameterCount; i++) {
|
|---|
| 223 | // OIDs are unsigned, same as dataTypeID in parseField above
|
|---|
| 224 | message.dataTypeIDs[i] = reader.uint32();
|
|---|
| 225 | }
|
|---|
| 226 | return message;
|
|---|
| 227 | };
|
|---|
| 228 | const parseDataRowMessage = (reader) => {
|
|---|
| 229 | const fieldCount = reader.int16();
|
|---|
| 230 | const fields = new Array(fieldCount);
|
|---|
| 231 | for (let i = 0; i < fieldCount; i++) {
|
|---|
| 232 | const len = reader.int32();
|
|---|
| 233 | // a -1 for length means the value of the field is null
|
|---|
| 234 | fields[i] = len === -1 ? null : reader.string(len);
|
|---|
| 235 | }
|
|---|
| 236 | return new messages_1.DataRowMessage(LATEINIT_LENGTH, fields);
|
|---|
| 237 | };
|
|---|
| 238 | const parseParameterStatusMessage = (reader) => {
|
|---|
| 239 | const name = reader.cstring();
|
|---|
| 240 | const value = reader.cstring();
|
|---|
| 241 | return new messages_1.ParameterStatusMessage(LATEINIT_LENGTH, name, value);
|
|---|
| 242 | };
|
|---|
| 243 | const parseBackendKeyData = (reader) => {
|
|---|
| 244 | const processID = reader.int32();
|
|---|
| 245 | const secretKey = reader.int32();
|
|---|
| 246 | return new messages_1.BackendKeyDataMessage(LATEINIT_LENGTH, processID, secretKey);
|
|---|
| 247 | };
|
|---|
| 248 | const parseAuthenticationResponse = (reader, length) => {
|
|---|
| 249 | const code = reader.int32();
|
|---|
| 250 | // TODO(bmc): maybe better types here
|
|---|
| 251 | const message = {
|
|---|
| 252 | name: 'authenticationOk',
|
|---|
| 253 | length,
|
|---|
| 254 | };
|
|---|
| 255 | switch (code) {
|
|---|
| 256 | case 0: // AuthenticationOk
|
|---|
| 257 | break;
|
|---|
| 258 | case 3: // AuthenticationCleartextPassword
|
|---|
| 259 | if (message.length === 8) {
|
|---|
| 260 | message.name = 'authenticationCleartextPassword';
|
|---|
| 261 | }
|
|---|
| 262 | break;
|
|---|
| 263 | case 5: // AuthenticationMD5Password
|
|---|
| 264 | if (message.length === 12) {
|
|---|
| 265 | message.name = 'authenticationMD5Password';
|
|---|
| 266 | const salt = reader.bytes(4);
|
|---|
| 267 | return new messages_1.AuthenticationMD5Password(LATEINIT_LENGTH, salt);
|
|---|
| 268 | }
|
|---|
| 269 | break;
|
|---|
| 270 | case 10: // AuthenticationSASL
|
|---|
| 271 | {
|
|---|
| 272 | message.name = 'authenticationSASL';
|
|---|
| 273 | message.mechanisms = [];
|
|---|
| 274 | let mechanism;
|
|---|
| 275 | do {
|
|---|
| 276 | mechanism = reader.cstring();
|
|---|
| 277 | if (mechanism) {
|
|---|
| 278 | message.mechanisms.push(mechanism);
|
|---|
| 279 | }
|
|---|
| 280 | } while (mechanism);
|
|---|
| 281 | }
|
|---|
| 282 | break;
|
|---|
| 283 | case 11: // AuthenticationSASLContinue
|
|---|
| 284 | message.name = 'authenticationSASLContinue';
|
|---|
| 285 | message.data = reader.string(length - 8);
|
|---|
| 286 | break;
|
|---|
| 287 | case 12: // AuthenticationSASLFinal
|
|---|
| 288 | message.name = 'authenticationSASLFinal';
|
|---|
| 289 | message.data = reader.string(length - 8);
|
|---|
| 290 | break;
|
|---|
| 291 | default:
|
|---|
| 292 | throw new Error('Unknown authenticationOk message type ' + code);
|
|---|
| 293 | }
|
|---|
| 294 | return message;
|
|---|
| 295 | };
|
|---|
| 296 | const parseErrorMessage = (reader, name) => {
|
|---|
| 297 | const fields = {};
|
|---|
| 298 | let fieldType = reader.string(1);
|
|---|
| 299 | while (fieldType !== '\0') {
|
|---|
| 300 | fields[fieldType] = reader.cstring();
|
|---|
| 301 | fieldType = reader.string(1);
|
|---|
| 302 | }
|
|---|
| 303 | const messageValue = fields.M;
|
|---|
| 304 | const message = name === 'notice'
|
|---|
| 305 | ? new messages_1.NoticeMessage(LATEINIT_LENGTH, messageValue)
|
|---|
| 306 | : new messages_1.DatabaseError(messageValue, LATEINIT_LENGTH, name);
|
|---|
| 307 | message.severity = fields.S;
|
|---|
| 308 | message.code = fields.C;
|
|---|
| 309 | message.detail = fields.D;
|
|---|
| 310 | message.hint = fields.H;
|
|---|
| 311 | message.position = fields.P;
|
|---|
| 312 | message.internalPosition = fields.p;
|
|---|
| 313 | message.internalQuery = fields.q;
|
|---|
| 314 | message.where = fields.W;
|
|---|
| 315 | message.schema = fields.s;
|
|---|
| 316 | message.table = fields.t;
|
|---|
| 317 | message.column = fields.c;
|
|---|
| 318 | message.dataType = fields.d;
|
|---|
| 319 | message.constraint = fields.n;
|
|---|
| 320 | message.file = fields.F;
|
|---|
| 321 | message.line = fields.L;
|
|---|
| 322 | message.routine = fields.R;
|
|---|
| 323 | return message;
|
|---|
| 324 | };
|
|---|
| 325 | //# sourceMappingURL=parser.js.map |
|---|