LCOV - code coverage report
Current view: top level - lib - disruptor.js (source / functions) Hit Total Coverage
Test: lcov_final.info Lines: 335 335 100.0 %
Date: 2023-02-17 22:37:01 Functions: 23 23 100.0 %
Branches: 111 111 100.0 %

           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;

Generated by: LCOV version 1.16