Don't clone stream in packet.parse()

This commit is contained in:
Daniel Huigens 2018-06-12 18:55:00 +02:00
parent ddda6a0b16
commit 589b666ac7
2 changed files with 75 additions and 71 deletions

View File

@ -138,8 +138,8 @@ export default {
* @returns {Boolean} Returns false if the stream was empty and parsing is done, and true otherwise.
*/
read: async function(input, callback) {
let reader = stream.getReader(input);
let controller;
const reader = stream.getReader(input);
let writer;
try {
const peekedBytes = await reader.peekBytes(2);
// some sanity checks
@ -168,6 +168,16 @@ export default {
const streaming = this.supportsStreaming(tag);
let packet = null;
let callbackReturned;
if (streaming) {
const transform = new TransformStream();
writer = stream.getWriter(transform.writable);
packet = transform.readable;
callbackReturned = callback({ tag, packet });
}
let wasPartialLength;
do {
if (!format) {
// 4.2.1. Old Format Packet Lengths
switch (packet_length_type) {
@ -202,8 +212,6 @@ export default {
break;
}
} else { // 4.2.2. New Format Packet Lengths
let wasPartialLength;
do {
// 4.2.2.1. One-Octet Lengths
const lengthByte = await reader.readByte();
wasPartialLength = false;
@ -218,51 +226,47 @@ export default {
wasPartialLength = true;
if (!streaming) {
throw new TypeError('This packet type does not support partial lengths.');
} else if (!packet) {
packet = new ReadableStream({
// eslint-disable-next-line no-loop-func
async start(_controller) {
controller = _controller;
},
cancel: stream.cancel.bind(input)
});
callback({ tag, packet });
}
// 4.2.2.3. Five-Octet Lengths
} else {
packet_length = (await reader.readByte() << 24) | (await reader.readByte() << 16) | (await reader.readByte() <<
8) | await reader.readByte();
}
if (controller) {
controller.enqueue(await reader.readBytes(packet_length));
}
if (writer) {
let bytesRead = 0;
while (true) {
await writer.ready;
const { done, value } = await reader.read();
if (done) {
if (packet_length === Infinity) break;
throw new Error('Unexpected end of packet');
}
await writer.write(value.slice(0, packet_length - bytesRead));
bytesRead += value.length;
if (bytesRead >= packet_length) {
reader.unshift(value.slice(packet_length - bytesRead + value.length));
break;
}
}
}
} while(wasPartialLength);
}
if (!packet) {
if (streaming) {
// Send the remainder of the packet to the callback as a stream
reader.releaseLock();
packet = stream.slice(stream.clone(input), 0, packet_length);
await callback({ tag, packet });
// Read the entire packet before parsing the next one
reader = stream.getReader(input);
await reader.readBytes(packet_length);
} else {
if (!streaming) {
packet = await reader.readBytes(packet_length);
await callback({ tag, packet });
}
}
const { done, value } = await reader.read();
if (!done) reader.unshift(value);
if (controller) {
controller.close();
if (writer) {
await writer.ready;
await writer.close();
}
if (streaming) await callbackReturned;
return done || !value || !value.length;
} catch(e) {
if (controller) {
controller.error(e);
if (writer) {
writer.abort(e);
return true;
} else {
throw e;

View File

@ -47,7 +47,7 @@ List.prototype.read = async function (bytes) {
packet.packets = new List();
packet.fromStream = util.isStream(parsed.packet);
await packet.read(parsed.packet);
writer.write(packet);
await writer.write(packet);
} catch (e) {
if (!config.tolerant ||
parsed.tag === enums.packet.symmetricallyEncrypted ||