Skip to content

Instantly share code, notes, and snippets.

@rhalff
Created November 13, 2019 22:14
Show Gist options
  • Select an option

  • Save rhalff/3b94bb59983cc357e2c1f0659a165c1f to your computer and use it in GitHub Desktop.

Select an option

Save rhalff/3b94bb59983cc357e2c1f0659a165c1f to your computer and use it in GitHub Desktop.
Replay subject - clear queue
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