diff --git a/lib/src/managers/event.dart b/lib/src/managers/event.dart index a4d58f2..80c299d 100644 --- a/lib/src/managers/event.dart +++ b/lib/src/managers/event.dart @@ -1,4 +1,5 @@ import 'dart:async'; +import 'dart:collection'; import 'package:meta/meta.dart'; import 'package:synchronized/synchronized.dart' as sync; @@ -22,9 +23,13 @@ class EventsEmitter extends EventsListenable { // suppport for multiple event listeners final streamCtrl = StreamController.broadcast(sync: false); + bool _queueMode = false; + final _queue = Queue(); + EventsEmitter({ bool listenSynchronized = false, }) : super(synchronized: listenSynchronized) { + // clean up onDispose(() async => await streamCtrl.close()); } @@ -33,14 +38,40 @@ class EventsEmitter extends EventsListenable { @internal void emit(T event) { - // do nothing if already closed - if (streamCtrl.isClosed) { + // check if already disposed + if (isDisposed) { logger.warning('failed to emit event ${event} on a disposed emitter'); return; } + // queue mode + if (_queueMode) { + _queue.add(event); + return; + } // emit the event streamCtrl.add(event); } + + @internal + void updateQueueMode(bool newValue, {bool shouldEmitQueued = true}) { + // check if already disposed + if (isDisposed) { + logger.warning('failed to update queueMode on a disposed emitter'); + return; + } + if (_queueMode == newValue) return; + _queueMode = newValue; + if (!_queueMode && shouldEmitQueued) emitQueued(); + } + + @internal + void emitQueued() { + while (_queue.isNotEmpty) { + final event = _queue.removeFirst(); + // emit the event + streamCtrl.add(event); + } + } } // for listening only