update method

  1. @override
int update()

Removes consumed and expired events.

Returns the largest number of events missed by one reader.

Implementation

@override
int update() {
  _readerLagged = false;
  if (_readers.isEmpty) {
    // Apply retention before the first reader exists.
    final maxPasses = retainedUpdates;
    if (maxPasses == null) {
      _base = _end;
      _events.clear();
      _retainedEnds.clear();
      return 0;
    }
    var floor = _base;
    final window = maxPasses - 1;
    if (window == 0) {
      floor = _end;
    } else {
      if (_retainedEnds.length == window) {
        floor = _retainedEnds.removeAt(0);
      }
      _retainedEnds.add(_end);
    }
    final drop = floor - _base;
    if (drop > 0) {
      _events.removeRange(0, drop);
      _base = floor;
    }
    return 0;
  }
  final maxPasses = retainedUpdates;
  var floor = _base;
  if (maxPasses != null) {
    // Events recorded [maxPasses - 1] passes ago have now been observable
    // for maxPasses frame windows; expire them. With maxPasses == 1 that is
    // everything sent before this pass.
    final window = maxPasses - 1;
    if (window == 0) {
      floor = _end;
    } else {
      if (_retainedEnds.length == window) {
        floor = _retainedEnds.removeAt(0);
      }
      _retainedEnds.add(_end);
    }
  }
  var minCursor = _end;
  var maxSkipped = 0;
  for (final reader in _readers) {
    // Advance readers past expired events.
    final lag = floor - reader._cursor;
    if (lag > 0) {
      reader._cursor = floor;
      if (lag > maxSkipped) maxSkipped = lag;
    }
    if (reader._cursor < minCursor) minCursor = reader._cursor;
  }
  final drop = minCursor - _base;
  if (drop > 0) {
    _events.removeRange(0, drop);
    _base += drop;
  }
  _readerLagged = maxSkipped > 0;
  return maxSkipped;
}