-
-
Notifications
You must be signed in to change notification settings - Fork 0
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
- Loading branch information
Showing
7 changed files
with
90 additions
and
174 deletions.
There are no files selected for viewing
This file contains 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
This file contains 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
This file was deleted.
Oops, something went wrong.
This file contains 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
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,39 @@ | ||
import 'dart:async'; | ||
|
||
import 'package:centrifuge_dart/src/model/channel_presence.dart'; | ||
import 'package:centrifuge_dart/src/model/event.dart'; | ||
import 'package:centrifuge_dart/src/model/publication.dart'; | ||
|
||
/// Stream of received events. | ||
/// {@category Entity} | ||
/// {@subCategory Event} | ||
/// {@subCategory Channel} | ||
final class CentrifugeEventStream extends StreamView<CentrifugeEvent> { | ||
/// Stream of received events. | ||
CentrifugeEventStream(super.stream); | ||
|
||
/// Publications stream. | ||
late final Stream<CentrifugePublication> publications = | ||
whereType<CentrifugePublication>(); | ||
|
||
/// Stream of presence (join & leave) events. | ||
late final Stream<CentrifugeChannelPresenceEvent> presenceEvents = | ||
whereType<CentrifugeChannelPresenceEvent>(); | ||
|
||
/// Join events | ||
late final Stream<CentrifugeJoinEvent> joinEvents = | ||
whereType<CentrifugeJoinEvent>(); | ||
|
||
/// Leave events | ||
late final Stream<CentrifugeLeaveEvent> leaveEvents = | ||
whereType<CentrifugeLeaveEvent>(); | ||
|
||
/// Filtered stream of data of [CentrifugeEvent]. | ||
Stream<T> whereType<T extends CentrifugeEvent>() => | ||
transform<T>(StreamTransformer<CentrifugeEvent, T>.fromHandlers( | ||
handleData: (data, sink) => switch (data) { | ||
T valid => sink.add(valid), | ||
_ => null, | ||
}, | ||
)).asBroadcastStream(); | ||
} |
This file contains 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
Oops, something went wrong.