All files / src/utils ContractsEventsSubscription.js

100% Statements 34/34
100% Branches 6/6
100% Functions 7/7
100% Lines 31/31

Press n or j to go to the next uncovered block, b, p or k for the previous block.

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 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86            1x 1x   1x   4x 4x   4x   4x                 4x     1x               5x 5x 5x 5x 5x 5x   5x       5x 5x       5x 5x       5x   5x 1x     4x   4x               4x   4x 7x     4x   4x 4x      
/**
 * Copyright (c) 2018-present, Leap DAO (leapdao.org)
 *
 * This source code is licensed under the Mozilla Public License Version 2.0
 * found in the LICENSE file in the root directory of this source tree.
 */
const EventEmitter = require('events');
const getBlockAverageTime = require('../utils/getBlockAverageTime');
 
const BATCH_SIZE = 5000;
async function getPastEvents(contract, eventName, fromBlock, toBlock) {
  const batchCount = Math.ceil((toBlock - fromBlock) / BATCH_SIZE);
  const events = [];
 
  for (let i = 0; i < batchCount; i += 1) {
    /* eslint-disable no-await-in-loop */
    events.push(
      await contract.getPastEvents(eventName, {
        fromBlock: i * BATCH_SIZE + fromBlock,
        toBlock: Math.min(toBlock, i * BATCH_SIZE + fromBlock + BATCH_SIZE),
      })
    );
    /* eslint-enable */
  }
 
  return events.reduce((result, ev) => result.concat(ev), []);
}
 
module.exports = class ContractsEventsSubscription extends EventEmitter {
  constructor(
    web3,
    contracts,
    eventsBuffer,
    fromBlock = null,
    eventName = 'allEvents'
  ) {
    super();
    this.fromBlock = fromBlock;
    this.web3 = web3;
    this.contracts = contracts;
    this.eventName = eventName;
    this.fetchEvents = this.fetchEvents.bind(this);
 
    this.eventsBuffer = eventsBuffer;
  }
 
  async init() {
    this.initialEvents = await this.fetchEvents();
    const eventsInterval = Math.max(
      1,
      (await getBlockAverageTime(this.web3)) * 0.7
    );
    setInterval(this.fetchEvents, eventsInterval * 1000);
    return this.initialEvents;
  }
 
  async fetchEvents() {
    const blockNumber = await this.web3.eth.getBlockNumber();
 
    if (this.fromBlock === blockNumber) {
      return null;
    }
 
    const eventsList = await Promise.all(
      this.contracts.map(contract => {
        return getPastEvents(
          contract,
          this.eventName,
          this.fromBlock || 0,
          blockNumber
        );
      })
    );
    const events = eventsList.reduce((acc, evnts) => acc.concat(evnts), []);
 
    for (const event of events) {
      this.eventsBuffer.push(event);
    }
 
    this.fromBlock = blockNumber + 1;
 
    this.emit('newEvents', events);
    return events;
  }
};