Branch data Line data Source code
1 [ + ]: 506 : const { promisify } = require('util');
2 : 506 : const { Readable, Writable } = require('stream');
3 : 506 : const Disruptor = require('bindings')('disruptor.node').Disruptor;
4 : 506 :
5 : 506 : const status_eof = 1;
6 : 506 : const status_error = 2;
7 : 506 :
8 [ + ]: 8200277 : function check(cb, r, arg1, arg2, arg3)
9 : 8200277 : {
10 : 8200277 : if (r !== undefined)
11 [ + ]: 8200277 : {
12 : 7713704 : process.nextTick(cb, null, r, arg1, arg2, arg3);
13 : 7713704 : }
14 : 8200277 : }
15 : 506 :
16 : 506 : class Disruptor2 extends Disruptor
17 : 506 : {
18 [ + ]: 506 : constructor(...args)
19 : 5771 : {
20 : 5771 : super(...args);
21 : 5771 :
22 [ + ]: 5771 : this._consumeNewAsync = promisify(cb => {
23 [ + ]: 6529800 : this._consumeNew((err, bufs, start) => {
24 : 6529800 : cb(err, { bufs, start });
25 : 6529800 : });
26 : 5771 : });
27 : 5771 :
28 [ + ]: 5771 : this._produceClaimAsync = promisify(cb => {
29 [ + ]: 467557 : this._produceClaim((err, buf, claimStart, claimEnd, allConsumersIgnoring) => {
30 : 467557 : cb(err, { buf, claimStart, claimEnd, allConsumersIgnoring });
31 : 467557 : });
32 : 5771 : });
33 : 5771 :
34 [ + ]: 5771 : this._produceClaimManyAsync = promisify((n, cb) => {
35 [ + ]: 27 : this._produceClaimMany(n, (err, bufs, claimStart, claimEnd, allConsumersIgnoring) => {
36 : 27 : cb(err, { bufs, claimStart, claimEnd, allConsumersIgnoring });
37 : 27 : });
38 : 5771 : });
39 : 5771 :
40 [ + ]: 5771 : this._produceClaimAvailAsync = promisify((max, cb) => {
41 [ + ]: 84308 : this._produceClaimAvail(max, (err, bufs, claimStart, claimEnd, allConsumersIgnoring) => {
42 : 84308 : cb(err, { bufs, claimStart, claimEnd, allConsumersIgnoring });
43 : 84308 : });
44 : 5771 : });
45 : 5771 :
46 [ + ]: 5771 : this._produceCommitAsync = promisify(function (claimStart, claimEnd, cb) {
47 [ + ]: 469725 : if (arguments.length >= 2) {
48 : 2163 : return this._produceCommit(claimStart, claimEnd, cb);
49 : 2163 : }
50 [ + ]: 467562 : this._produceCommit(claimStart);
51 : 5771 : });
52 : 5771 : }
53 : 506 :
54 [ + ]: 506 : _consumeNew(cb)
55 : 7074862 : {
56 : 7074862 : check(cb,
57 : 7074862 : super.consumeNew(cb),
58 : 7074862 : super.prevConsumeStart);
59 : 7074862 : }
60 : 506 :
61 [ + ]: 506 : consumeNew(cb)
62 : 7074862 : {
63 : 7074862 : if (cb)
64 [ + ]: 7074862 : {
65 : 545062 : return this._consumeNew(cb);
66 : 545062 : }
67 [ + ]: 6529800 :
68 : 6529800 : return this._consumeNewAsync();
69 : 7074862 : }
70 : 506 :
71 [ + ]: 506 : _produceClaim(cb)
72 : 519445 : {
73 : 519445 : check(cb,
74 : 519445 : super.produceClaim(cb),
75 : 519445 : super.prevClaimStart,
76 : 519445 : super.prevClaimEnd,
77 : 519445 : super.allConsumersIgnoring);
78 : 519445 : }
79 : 506 :
80 [ + ]: 506 : produceClaim(cb)
81 : 519445 : {
82 : 519445 : if (cb)
83 [ + ]: 519445 : {
84 : 51888 : return this._produceClaim(cb);
85 : 51888 : }
86 [ + ]: 467557 :
87 : 467557 : return this._produceClaimAsync();
88 : 519445 : }
89 : 506 :
90 [ + ]: 506 : _produceClaimMany(n, cb)
91 : 40 : {
92 : 40 : check(cb,
93 : 40 : super.produceClaimMany(n, cb),
94 : 40 : super.prevClaimStart,
95 : 40 : super.prevClaimEnd,
96 : 40 : super.allConsumersIgnoring);
97 : 40 : }
98 : 506 :
99 [ + ]: 506 : produceClaimMany(n, cb)
100 : 40 : {
101 : 40 : if (cb)
102 [ + ]: 40 : {
103 : 13 : return this._produceClaimMany(n, cb);
104 : 13 : }
105 [ + ]: 27 :
106 : 27 : return this._produceClaimManyAsync(n);
107 : 40 : }
108 : 506 :
109 [ + ]: 506 : _produceClaimAvail(max, cb)
110 : 84312 : {
111 : 84312 : check(cb,
112 : 84312 : super.produceClaimAvail(max, cb),
113 : 84312 : super.prevClaimStart,
114 : 84312 : super.prevClaimEnd,
115 : 84312 : super.allConsumersIgnoring);
116 : 84312 : }
117 : 506 :
118 [ + ]: 506 : produceClaimAvail(max, cb)
119 : 84312 : {
120 : 84312 : if (cb)
121 [ + ]: 84312 : {
122 : 4 : return this._produceClaimAvail(max, cb);
123 : 4 : }
124 [ + ]: 84308 :
125 : 84308 : return this._produceClaimAvailAsync(max);
126 : 84312 : }
127 : 506 :
128 [ + ]: 506 : _produceCommit(claimStart, claimEnd, cb)
129 : 521618 : {
130 : 521618 : if (arguments.length >= 2)
131 [ + ]: 521618 : {
132 : 2171 : return check(cb, super.produceCommit(claimStart, claimEnd, cb));
133 : 2171 : }
134 [ + ]: 519447 :
135 : 519447 : check(claimStart, super.produceCommit(claimStart));
136 : 521618 : }
137 : 506 :
138 [ + ]: 506 : produceCommit(claimStart, claimEnd, cb)
139 : 521618 : {
140 : 521618 : if (arguments.length >= 2)
141 [ + ]: 521618 : {
142 : 2171 : if (cb)
143 [ + ]: 2171 : {
144 : 8 : return this._produceCommit(claimStart, claimEnd, cb);
145 : 8 : }
146 [ + ]: 2163 :
147 : 2163 : return this._produceCommitAsync(claimStart, claimEnd);
148 : 2163 : }
149 [ + ][ + ]: 519447 :
150 : 519447 : if (claimStart)
151 [ + ]: 520618 : {
152 : 51885 : return this._produceCommit(claimStart);
153 : 51885 : }
154 [ + ]: 467562 :
155 : 467562 : return this._produceCommitAsync();
156 : 521618 : }
157 : 506 : }
158 : 506 :
159 : 506 : class DisruptorReadStream extends Readable {
160 [ + ]: 506 : constructor(disruptor, options) {
161 : 2012 : super(options);
162 [ + ]: 2012 : if (disruptor.elementSize !== 1) {
163 : 1 : throw new Error('element size must be 1');
164 : 1 : }
165 [ + ][ + ]: 2012 : if (disruptor.spin) {
166 : 1 : throw new Error('spin must be false');
167 : 1 : }
168 [ + ]: 2010 : this.disruptor = disruptor;
169 : 2010 : this._reading = false;
170 : 2012 : }
171 : 506 :
172 [ + ]: 506 : async __read() {
173 [ + ]: 6128592 : if (this._reading) {
174 : 1 : return;
175 : 1 : }
176 [ + ]: 6128591 : let buf;
177 : 6128591 : this._reading = true;
178 : 6128591 : try {
179 [ + ]: 6128592 : buf = Buffer.concat((await this.disruptor.consumeNew()).bufs);
180 : 6128590 : this.disruptor.consumeCommit();
181 : 6128590 : }
182 [ + ]: 6128592 : catch (ex) {
183 : 1 : this._reading = false; // must be done in case _destroy below gets called
184 : 1 : this.emit('error', ex);
185 : 1 : }
186 [ + ]: 6128591 : this._reading = false;
187 [ + ]: 6128592 : if (this._destroy_info) {
188 : 1001 : const { err, cb } = this._destroy_info;
189 : 1001 : cb(err);
190 : 1001 : return;
191 : 1001 : }
192 [ + ]: 6127590 : return buf;
193 : 6128592 : }
194 : 506 :
195 [ + ]: 506 : async _read() {
196 [ + ]: 6235515 : if (this.destroyed) {
197 : 107927 : return;
198 : 107927 : }
199 [ + ]: 6127588 : const buf = await this.__read();
200 [ + ]: 6235515 : if (buf) {
201 [ + ]: 6126585 : if (buf.length > 0) {
202 [ + ]: 107926 : if (this.push(buf)) {
203 : 107924 : this.read_again();
204 : 107924 : }
205 [ + ][ + ]: 6126585 : } else if (this.disruptor.status === status_eof) {
206 : 1004 : // Check for data the writer added between us reading zero bytes
207 : 1004 : // and it setting status_eof
208 : 1004 : const buf2 = await this.__read();
209 : 1004 : if (buf2) {
210 : 1004 : this.push(buf2);
211 : 1004 : this.push(null);
212 : 1004 : }
213 [ + ][ + ]: 6018659 : } else if (this.disruptor.status === status_error) {
214 : 1 : this.emit('error', new Error('writer errored'));
215 [ + ]: 6017655 : } else {
216 : 6017654 : this.read_again();
217 : 6017654 : }
218 : 6126585 : }
219 : 6235515 : }
220 : 506 :
221 [ + ]: 506 : read_again() {
222 [ + ]: 6125578 : setImmediate(() => this._read());
223 : 6125578 : }
224 : 506 :
225 [ + ]: 506 : _destroy(err, cb) {
226 [ + ]: 2010 : if (this._reading) {
227 : 1001 : this._destroy_info = { err, cb };
228 [ + ]: 2010 : } else {
229 : 1009 : cb(err);
230 : 1009 : }
231 : 2010 : }
232 : 506 : }
233 : 506 :
234 : 506 : class DisruptorWriteStream extends Writable {
235 [ + ]: 506 : constructor(disruptor, options) {
236 : 16 : super(options);
237 [ + ]: 16 : if (disruptor.elementSize !== 1) {
238 : 1 : throw new Error('element size must be 1');
239 : 1 : }
240 [ + ][ + ]: 16 : if (disruptor.spin) {
241 : 1 : throw new Error('spin must be false');
242 : 1 : }
243 [ + ]: 14 : this.disruptor = disruptor;
244 : 14 : this._writing = false;
245 : 16 : }
246 : 506 :
247 [ + ]: 506 : async _write(chunk, encoding, cb) {
248 [ + ]: 84304 : if (this.destroyed) {
249 : 1 : return;
250 : 1 : }
251 [ + ]: 84303 : let bufs, claimStart, claimEnd, allConsumersIgnoring;
252 : 84303 : this._writing = true;
253 : 84303 : try {
254 : 84303 : ({ bufs,
255 : 84303 : claimStart,
256 : 84303 : claimEnd,
257 : 84303 : allConsumersIgnoring
258 [ + ]: 84304 : } = await this.disruptor.produceClaimAvail(chunk.length));
259 [ + ]: 84304 : } catch (ex) {
260 : 1 : this._writing = false; // must be done before onwrite called _destroy above
261 : 1 : return cb(ex);
262 : 1 : }
263 [ + ]: 84302 : this._writing = false;
264 [ + ]: 84304 : if (this._destroy_info) {
265 : 1 : const { err, cb } = this._destroy_info;
266 : 1 : return cb(err);
267 : 1 : }
268 [ + ]: 84301 : let i = 0;
269 [ + ]: 84304 : for (const buf of bufs) {
270 : 2167 : chunk.copy(buf, 0, i, i + buf.length);
271 : 2167 : i += buf.length;
272 : 2167 : }
273 [ + ][ + ]: 84304 : if (i === 0) {
274 [ + ]: 82138 : if (allConsumersIgnoring) {
275 : 1 : return cb(new Error('no consumers'));
276 : 1 : }
277 [ + ]: 82137 : return this.write_again(chunk, encoding, cb);
278 : 82137 : }
279 [ + ]: 2163 : await this._commit(claimStart, claimEnd, chunk.slice(i), encoding, cb);
280 : 84304 : }
281 : 506 :
282 [ + ]: 506 : _final(cb) {
283 : 5 : this.disruptor.status = status_eof;
284 : 5 : cb();
285 : 5 : }
286 : 506 :
287 [ + ]: 506 : _destroy(err, cb) {
288 [ + ]: 14 : if (err) {
289 : 6 : this.disruptor.status = status_error;
290 : 6 : }
291 [ + ]: 14 : if (this._writing) {
292 : 2 : this._destroy_info = { err, cb };
293 [ + ]: 14 : } else {
294 : 12 : cb(err);
295 : 12 : }
296 : 14 : }
297 : 506 :
298 [ + ]: 506 : async _commit(claimStart, claimEnd, chunk, encoding, cb) {
299 [ + ]: 2164 : if (this.destroyed) {
300 : 1 : return;
301 : 1 : }
302 [ + ]: 2163 : let committed;
303 : 2163 : this._writing = true;
304 : 2163 : try {
305 : 2163 : committed = await this.disruptor.produceCommit(claimStart, claimEnd);
306 [ + ][ + ]: 2164 : } catch (ex) {
307 : 1 : this._writing = false; // must be done before onwrite called _destroy above
308 : 1 : return cb(ex);
309 : 1 : }
310 [ + ]: 2162 : this._writing = false;
311 [ + ]: 2164 : if (this._destroy_info) {
312 : 1 : const { err, cb } = this._destroy_info;
313 : 1 : return cb(err);
314 : 1 : }
315 [ + ][ + ]: 2164 : if (!committed) {
316 : 1 : return this.commit_again(claimStart, claimEnd, chunk, encoding, cb);
317 : 1 : }
318 [ + ][ + ]: 2164 : if (chunk.length > 0) {
319 : 1153 : return this.write_again(chunk, encoding, cb);
320 : 1153 : }
321 [ + ]: 1007 : cb();
322 : 2164 : }
323 : 506 :
324 [ + ]: 506 : commit_again(claimStart, claimEnd, chunk, encoding, cb) {
325 [ + ]: 1 : setImmediate(() => this._commit(claimStart, claimEnd, chunk, encoding, cb));
326 : 1 : }
327 : 506 :
328 [ + ]: 506 : write_again(chunk, encoding, cb) {
329 [ + ]: 83290 : setImmediate(() => this._write(chunk, encoding, cb));
330 : 83290 : }
331 : 506 : }
332 : 506 :
333 : 506 : exports.Disruptor = Disruptor2;
334 : 506 : exports.DisruptorReadStream = DisruptorReadStream;
335 : 506 : exports.DisruptorWriteStream = DisruptorWriteStream;
|