update method
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;
}