-
Notifications
You must be signed in to change notification settings - Fork 72
Expand file tree
/
Copy pathstream.js
More file actions
62 lines (60 loc) · 1.65 KB
/
Copy pathstream.js
File metadata and controls
62 lines (60 loc) · 1.65 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
import { Transform } from 'stream';
import { Packr } from './pack.js';
import { Unpackr } from './unpack.js';
var DEFAULT_OPTIONS = {objectMode: true};
export class PackrStream extends Transform {
constructor(options) {
if (!options)
options = {};
options.writableObjectMode = true;
super(options);
options.sequential = true;
this.packr = options.packr || new Packr(options);
}
_transform(value, encoding, callback) {
this.push(this.packr.pack(value));
callback();
}
}
export class UnpackrStream extends Transform {
constructor(options) {
if (!options)
options = {};
options.objectMode = true;
super(options);
options.structures = [];
this.maxIncompleteBufferSize = options.maxIncompleteBufferSize !== undefined ? options.maxIncompleteBufferSize : 0x4000000;
this.unpackr = options.unpackr || new Unpackr(options);
}
_transform(chunk, encoding, callback) {
if (this.incompleteBuffer) {
chunk = Buffer.concat([this.incompleteBuffer, chunk]);
this.incompleteBuffer = null;
}
let values;
try {
values = this.unpackr.unpackMultiple(chunk);
} catch(error) {
if (error.incomplete) {
let incompleteBuffer = chunk.slice(error.lastPosition);
if (incompleteBuffer.length > this.maxIncompleteBufferSize) {
this.incompleteBuffer = null;
return callback(new Error('Maximum incomplete buffer size exceeded'));
}
this.incompleteBuffer = incompleteBuffer;
values = error.values;
} else {
return callback(error);
}
}
for (let value of values || []) {
if (value === null)
value = this.getNullValue();
this.push(value);
}
callback();
}
getNullValue() {
return Symbol.for(null);
}
}