Created
November 13, 2019 22:14
-
-
Save rhalff/3b94bb59983cc357e2c1f0659a165c1f to your computer and use it in GitHub Desktop.
Replay subject - clear queue
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| import 'dart:async'; | |
| import 'dart:collection'; | |
| import 'package:rxdart/rxdart.dart'; | |
| class QueueSubject<T> extends Subject<T> implements ReplayObservable<T> { | |
| final Queue<T> _queue; | |
| final int _maxSize; | |
| factory QueueSubject({ | |
| int maxSize, | |
| Queue<T> queue, | |
| void onListen(), | |
| void onCancel(), | |
| bool sync = false, | |
| }) { | |
| // ignore: close_sinks | |
| final controller = StreamController<T>.broadcast( | |
| onListen: onListen, | |
| onCancel: onCancel, | |
| sync: sync, | |
| ); | |
| final _queue = queue ?? Queue<T>(); | |
| return QueueSubject<T>._( | |
| controller, | |
| Observable<T>.defer( | |
| () => Observable<T>(controller.stream) | |
| .startWithMany(_queue.toList(growable: false)), | |
| reusable: true), | |
| _queue, | |
| maxSize, | |
| ); | |
| } | |
| QueueSubject._( | |
| StreamController<T> controller, | |
| Observable<T> observable, | |
| this._queue, | |
| this._maxSize, | |
| ) : super(controller, observable); | |
| @override | |
| void onAdd(T event) { | |
| if (_queue.length == _maxSize) { | |
| _queue.removeFirst(); | |
| } | |
| _queue.add(event); | |
| } | |
| @override | |
| List<T> get values => _queue.toList(growable: false); | |
| void clear() { | |
| _queue.clear(); | |
| } | |
| } | |
| main() async { | |
| final subject = QueueSubject<int>(); | |
| ReplayObservable<int> replayObservable = subject.stream; | |
| subject.add(1); | |
| subject.add(2); | |
| subject.add(3); | |
| replayObservable.listen(print); // prints 1, 2, 3, 4, 5 | |
| replayObservable.listen(print); // prints 1, 2, 3, 4, 5 | |
| Timer.run(() { | |
| print('clear it.'); | |
| subject.clear(); | |
| subject.add(4); | |
| subject.add(5); | |
| replayObservable.listen(print); // prints 4, 5 | |
| replayObservable.listen(print); // prints 4, 5 | |
| subject.close(); | |
| }); | |
| } |
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment