| [81bc7da] | 1 | 'use strict';
|
|---|
| 2 |
|
|---|
| 3 | const Transform = require('stream').Transform;
|
|---|
| 4 |
|
|---|
| 5 | /**
|
|---|
| 6 | * MessageParser instance is a transform stream that separates message headers
|
|---|
| 7 | * from the rest of the body. Headers are emitted with the 'headers' event. Message
|
|---|
| 8 | * body is passed on as the resulting stream.
|
|---|
| 9 | */
|
|---|
| 10 | class MessageParser extends Transform {
|
|---|
| 11 | constructor(options) {
|
|---|
| 12 | super(options);
|
|---|
| 13 | this.lastBytes = Buffer.alloc(4);
|
|---|
| 14 | this.headersParsed = false;
|
|---|
| 15 | this.headerBytes = 0;
|
|---|
| 16 | this.headerChunks = [];
|
|---|
| 17 | this.rawHeaders = false;
|
|---|
| 18 | this.bodySize = 0;
|
|---|
| 19 | }
|
|---|
| 20 |
|
|---|
| 21 | /**
|
|---|
| 22 | * Keeps count of the last 4 bytes in order to detect line breaks on chunk boundaries
|
|---|
| 23 | *
|
|---|
| 24 | * @param {Buffer} data Next data chunk from the stream
|
|---|
| 25 | */
|
|---|
| 26 | updateLastBytes(data) {
|
|---|
| 27 | let lblen = this.lastBytes.length;
|
|---|
| 28 | let nblen = Math.min(data.length, lblen);
|
|---|
| 29 |
|
|---|
| 30 | // shift existing bytes
|
|---|
| 31 | for (let i = 0, len = lblen - nblen; i < len; i++) {
|
|---|
| 32 | this.lastBytes[i] = this.lastBytes[i + nblen];
|
|---|
| 33 | }
|
|---|
| 34 |
|
|---|
| 35 | // add new bytes
|
|---|
| 36 | for (let i = 1; i <= nblen; i++) {
|
|---|
| 37 | this.lastBytes[lblen - i] = data[data.length - i];
|
|---|
| 38 | }
|
|---|
| 39 | }
|
|---|
| 40 |
|
|---|
| 41 | /**
|
|---|
| 42 | * Finds and removes message headers from the remaining body. We want to keep
|
|---|
| 43 | * headers separated until final delivery to be able to modify these
|
|---|
| 44 | *
|
|---|
| 45 | * @param {Buffer} data Next chunk of data
|
|---|
| 46 | * @return {Boolean} Returns true if headers are already found or false otherwise
|
|---|
| 47 | */
|
|---|
| 48 | checkHeaders(data) {
|
|---|
| 49 | if (this.headersParsed) {
|
|---|
| 50 | return true;
|
|---|
| 51 | }
|
|---|
| 52 |
|
|---|
| 53 | let lblen = this.lastBytes.length;
|
|---|
| 54 | let headerPos = 0;
|
|---|
| 55 | this.curLinePos = 0;
|
|---|
| 56 | for (let i = 0, len = this.lastBytes.length + data.length; i < len; i++) {
|
|---|
| 57 | let chr;
|
|---|
| 58 | if (i < lblen) {
|
|---|
| 59 | chr = this.lastBytes[i];
|
|---|
| 60 | } else {
|
|---|
| 61 | chr = data[i - lblen];
|
|---|
| 62 | }
|
|---|
| 63 | if (chr === 0x0a && i) {
|
|---|
| 64 | let pr1 = i - 1 < lblen ? this.lastBytes[i - 1] : data[i - 1 - lblen];
|
|---|
| 65 | let pr2 = i > 1 ? (i - 2 < lblen ? this.lastBytes[i - 2] : data[i - 2 - lblen]) : false;
|
|---|
| 66 | if (pr1 === 0x0a) {
|
|---|
| 67 | this.headersParsed = true;
|
|---|
| 68 | headerPos = i - lblen + 1;
|
|---|
| 69 | this.headerBytes += headerPos;
|
|---|
| 70 | break;
|
|---|
| 71 | } else if (pr1 === 0x0d && pr2 === 0x0a) {
|
|---|
| 72 | this.headersParsed = true;
|
|---|
| 73 | headerPos = i - lblen + 1;
|
|---|
| 74 | this.headerBytes += headerPos;
|
|---|
| 75 | break;
|
|---|
| 76 | }
|
|---|
| 77 | }
|
|---|
| 78 | }
|
|---|
| 79 |
|
|---|
| 80 | if (this.headersParsed) {
|
|---|
| 81 | this.headerChunks.push(data.slice(0, headerPos));
|
|---|
| 82 | this.rawHeaders = Buffer.concat(this.headerChunks, this.headerBytes);
|
|---|
| 83 | this.headerChunks = null;
|
|---|
| 84 | this.emit('headers', this.parseHeaders());
|
|---|
| 85 | if (data.length - 1 > headerPos) {
|
|---|
| 86 | let chunk = data.slice(headerPos);
|
|---|
| 87 | this.bodySize += chunk.length;
|
|---|
| 88 | // this would be the first chunk of data sent downstream
|
|---|
| 89 | setImmediate(() => this.push(chunk));
|
|---|
| 90 | }
|
|---|
| 91 | return false;
|
|---|
| 92 | } else {
|
|---|
| 93 | this.headerBytes += data.length;
|
|---|
| 94 | this.headerChunks.push(data);
|
|---|
| 95 | }
|
|---|
| 96 |
|
|---|
| 97 | // store last 4 bytes to catch header break
|
|---|
| 98 | this.updateLastBytes(data);
|
|---|
| 99 |
|
|---|
| 100 | return false;
|
|---|
| 101 | }
|
|---|
| 102 |
|
|---|
| 103 | _transform(chunk, encoding, callback) {
|
|---|
| 104 | if (!chunk || !chunk.length) {
|
|---|
| 105 | return callback();
|
|---|
| 106 | }
|
|---|
| 107 |
|
|---|
| 108 | if (typeof chunk === 'string') {
|
|---|
| 109 | chunk = Buffer.from(chunk, encoding);
|
|---|
| 110 | }
|
|---|
| 111 |
|
|---|
| 112 | let headersFound;
|
|---|
| 113 |
|
|---|
| 114 | try {
|
|---|
| 115 | headersFound = this.checkHeaders(chunk);
|
|---|
| 116 | } catch (E) {
|
|---|
| 117 | return callback(E);
|
|---|
| 118 | }
|
|---|
| 119 |
|
|---|
| 120 | if (headersFound) {
|
|---|
| 121 | this.bodySize += chunk.length;
|
|---|
| 122 | this.push(chunk);
|
|---|
| 123 | }
|
|---|
| 124 |
|
|---|
| 125 | setImmediate(callback);
|
|---|
| 126 | }
|
|---|
| 127 |
|
|---|
| 128 | _flush(callback) {
|
|---|
| 129 | if (this.headerChunks) {
|
|---|
| 130 | let chunk = Buffer.concat(this.headerChunks, this.headerBytes);
|
|---|
| 131 | this.bodySize += chunk.length;
|
|---|
| 132 | this.push(chunk);
|
|---|
| 133 | this.headerChunks = null;
|
|---|
| 134 | }
|
|---|
| 135 | callback();
|
|---|
| 136 | }
|
|---|
| 137 |
|
|---|
| 138 | parseHeaders() {
|
|---|
| 139 | let lines = (this.rawHeaders || '').toString().split(/\r?\n/);
|
|---|
| 140 | for (let i = lines.length - 1; i > 0; i--) {
|
|---|
| 141 | if (/^\s/.test(lines[i])) {
|
|---|
| 142 | lines[i - 1] += '\n' + lines[i];
|
|---|
| 143 | lines.splice(i, 1);
|
|---|
| 144 | }
|
|---|
| 145 | }
|
|---|
| 146 | return lines
|
|---|
| 147 | .filter(line => line.trim())
|
|---|
| 148 | .map(line => ({
|
|---|
| 149 | key: line.substr(0, line.indexOf(':')).trim().toLowerCase(),
|
|---|
| 150 | line
|
|---|
| 151 | }));
|
|---|
| 152 | }
|
|---|
| 153 | }
|
|---|
| 154 |
|
|---|
| 155 | module.exports = MessageParser;
|
|---|