From 9c573f4924da2e62843afd09645a8bd21226a4c7 Mon Sep 17 00:00:00 2001 From: Hiroshi Horie <548776+hiroshihorie@users.noreply.github.com> Date: Sat, 11 Dec 2021 14:18:35 +0700 Subject: [PATCH] Stream State Update (#53) * implemented * update docs --- lib/src/events.dart | 14 ++++++++++++ lib/src/extensions.dart | 8 +++++++ lib/src/internal/events.dart | 8 +++++++ .../publication/remote_track_publication.dart | 22 +++++++++++++++++++ lib/src/room.dart | 16 ++++++++++++++ lib/src/rtc_engine.dart | 2 ++ lib/src/signal_client.dart | 5 +++++ lib/src/types.dart | 5 +++++ 8 files changed, 80 insertions(+) diff --git a/lib/src/events.dart b/lib/src/events.dart index b52781b..acd72a4 100644 --- a/lib/src/events.dart +++ b/lib/src/events.dart @@ -186,6 +186,20 @@ class TrackUnmutedEvent with RoomEvent, ParticipantEvent { }); } +/// The [StreamState] on the [RemoteTrackPublication] has updated by the server. +/// See [RemoteTrackPublication.streamState] for more information. +/// Emitted by [Room] and [RemoteParticipant]. +class TrackStreamStateUpdatedEvent with RoomEvent, ParticipantEvent { + final RemoteParticipant participant; + final RemoteTrackPublication trackPublication; + final StreamState streamState; + const TrackStreamStateUpdatedEvent({ + required this.participant, + required this.trackPublication, + required this.streamState, + }); +} + /// Participant metadata is a simple way for app-specific state to be pushed to /// all users. When RoomService.UpdateParticipantMetadata is called to change a /// [Participant]'s state, *all* [Participant]s in the room will fire this event. diff --git a/lib/src/extensions.dart b/lib/src/extensions.dart index 8dba2ae..8e6bb1f 100644 --- a/lib/src/extensions.dart +++ b/lib/src/extensions.dart @@ -124,3 +124,11 @@ extension LKTrackSourceExt on TrackSource { }[this] ?? lk_models.TrackSource.UNKNOWN; } + +extension PBStreamStateExt on lk_rtc.StreamState { + StreamState toLKType() => + { + lk_rtc.StreamState.ACTIVE: StreamState.active, + }[this] ?? + StreamState.paused; +} diff --git a/lib/src/internal/events.dart b/lib/src/internal/events.dart index 19fc885..6fb298c 100644 --- a/lib/src/internal/events.dart +++ b/lib/src/internal/events.dart @@ -185,3 +185,11 @@ class SignalMuteTrackEvent with SignalEvent { required this.muted, }); } + +@internal +class SignalStreamStateUpdatedEvent with SignalEvent, EngineEvent { + final List updates; + const SignalStreamStateUpdatedEvent({ + required this.updates, + }); +} diff --git a/lib/src/publication/remote_track_publication.dart b/lib/src/publication/remote_track_publication.dart index a8d4c08..b626cff 100644 --- a/lib/src/publication/remote_track_publication.dart +++ b/lib/src/publication/remote_track_publication.dart @@ -13,6 +13,7 @@ import '../participant/remote_participant.dart'; import '../proto/livekit_models.pb.dart' as lk_models; import '../proto/livekit_rtc.pb.dart' as lk_rtc; import '../track/remote.dart'; +import '../types.dart'; import '../utils.dart'; import 'track_publication.dart'; @@ -40,6 +41,27 @@ class RemoteTrackPublication lk_models.VideoQuality _videoQuality = lk_models.VideoQuality.HIGH; lk_models.VideoQuality get videoQuality => _videoQuality; + StreamState _streamState = StreamState.paused; + + /// The server may pause the track when they are bandwidth limitations and resume + /// when there is more capacity. This property will be updated when the track is + /// paused / resumed by the server. See [TrackStreamStateUpdatedEvent] for the + /// relevant event. + StreamState get streamState => _streamState; + + @internal + Future updateStreamState(StreamState streamState) async { + // return if no change + if (_streamState == streamState) return; + _streamState = streamState; + [participant.events, participant.room.events] + .emit(TrackStreamStateUpdatedEvent( + participant: participant, + trackPublication: this, + streamState: streamState, + )); + } + // used to report renderer visibility to the server // and optimize final _visibilities = {}; diff --git a/lib/src/room.dart b/lib/src/room.dart index 9b30f0b..51ecc7b 100644 --- a/lib/src/room.dart +++ b/lib/src/room.dart @@ -152,6 +152,8 @@ class Room extends DisposableChangeNotifier with EventsEmittable { (event) => _onSignalSpeakersChangedEvent(event.speakers)) ..on( (event) => _onSignalConnectionQualityUpdateEvent(event.updates)) + ..on( + (event) => _onSignalStreamStateUpdateEvent(event.updates)) ..on(_onDataMessageEvent) ..on((event) async { final publication = localParticipant?.trackPublications[event.sid]; @@ -364,6 +366,20 @@ class Room extends DisposableChangeNotifier with EventsEmittable { } } + void _onSignalStreamStateUpdateEvent( + List updates) async { + for (final update in updates) { + // try to find RemoteParticipant + final participant = participants[update.participantSid]; + if (participant == null) continue; + // try to find RemoteTrackPublication + final trackPublication = participant.trackPublications[update.trackSid]; + if (trackPublication == null) continue; + // update the stream state + await trackPublication.updateStreamState(update.state.toLKType()); + } + } + void _onDataMessageEvent(EngineDataPacketReceivedEvent dataPacketEvent) { // participant may be null if data is sent from Server-API final senderSid = dataPacketEvent.packet.participantSid; diff --git a/lib/src/rtc_engine.dart b/lib/src/rtc_engine.dart index cc7a531..196db7a 100644 --- a/lib/src/rtc_engine.dart +++ b/lib/src/rtc_engine.dart @@ -563,6 +563,8 @@ class RTCEngine extends Disposable with EventsEmittable { ..on((event) => events.emit(event)) // relay ..on((event) => events.emit(event)) + // relay + ..on((event) => events.emit(event)) ..on((event) async { await close(); events.emit(const EngineDisconnectedEvent()); diff --git a/lib/src/signal_client.dart b/lib/src/signal_client.dart index 1e63422..d4b6d5f 100644 --- a/lib/src/signal_client.dart +++ b/lib/src/signal_client.dart @@ -250,6 +250,11 @@ class SignalClient extends Disposable with EventsEmittable { muted: msg.mute.muted, )); break; + case lk_rtc.SignalResponse_Message.streamStateUpdate: + events.emit(SignalStreamStateUpdatedEvent( + updates: msg.streamStateUpdate.streamStates, + )); + break; default: logger.warning('skipping unsupported signal message'); } diff --git a/lib/src/types.dart b/lib/src/types.dart index 4b59b9c..b779940 100644 --- a/lib/src/types.dart +++ b/lib/src/types.dart @@ -37,6 +37,11 @@ enum TrackSource { screenShareAudio, } +enum StreamState { + paused, + active, +} + enum CloseReason { network, // ...