handle combined participant update. (#130)
* handle combined participant update (WIP). * always emit events for new publications. * update info for exist participant. * fix tests. * update. * Move emitWhenConnected to Room layer.
This commit is contained in:
@@ -406,6 +406,27 @@ class Engine extends Disposable with EventsEmittable<EngineEvent> {
|
|||||||
logger.fine('[WebRTC] stream.onRemoveTrack');
|
logger.fine('[WebRTC] stream.onRemoveTrack');
|
||||||
};
|
};
|
||||||
|
|
||||||
|
if (connectionState == ConnectionState.reconnecting ||
|
||||||
|
connectionState == ConnectionState.connecting) {
|
||||||
|
final track = event.track;
|
||||||
|
final receiver = event.receiver;
|
||||||
|
events.on<EngineConnectionStateUpdatedEvent>((event) async {
|
||||||
|
Timer(const Duration(milliseconds: 10), () {
|
||||||
|
events.emit(EngineTrackAddedEvent(
|
||||||
|
track: track,
|
||||||
|
stream: stream,
|
||||||
|
receiver: receiver,
|
||||||
|
));
|
||||||
|
});
|
||||||
|
});
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
if (connectionState == ConnectionState.disconnected) {
|
||||||
|
logger.warning('skipping incoming track after Room disconnected');
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
events.emit(EngineTrackAddedEvent(
|
events.emit(EngineTrackAddedEvent(
|
||||||
track: event.track,
|
track: event.track,
|
||||||
stream: stream,
|
stream: stream,
|
||||||
|
|||||||
+28
-5
@@ -1,4 +1,5 @@
|
|||||||
import 'package:collection/collection.dart';
|
import 'package:collection/collection.dart';
|
||||||
|
import 'package:meta/meta.dart';
|
||||||
|
|
||||||
import '../core/signal_client.dart';
|
import '../core/signal_client.dart';
|
||||||
import '../events.dart';
|
import '../events.dart';
|
||||||
@@ -190,10 +191,16 @@ class Room extends DisposableChangeNotifier with EventsEmittable<RoomEvent> {
|
|||||||
}
|
}
|
||||||
//
|
//
|
||||||
await publication.updateSubscriptionAllowed(event.allowed);
|
await publication.updateSubscriptionAllowed(event.allowed);
|
||||||
|
emitWhenConnected(TrackSubscriptionPermissionChangedEvent(
|
||||||
|
participant: participant,
|
||||||
|
publication: publication,
|
||||||
|
state: publication.subscriptionState,
|
||||||
|
));
|
||||||
})
|
})
|
||||||
..on<SignalRoomUpdateEvent>((event) async {
|
..on<SignalRoomUpdateEvent>((event) async {
|
||||||
_metadata = event.room.metadata;
|
_metadata = event.room.metadata;
|
||||||
events.emit(RoomMetadataChangedEvent(metadata: event.room.metadata));
|
emitWhenConnected(
|
||||||
|
RoomMetadataChangedEvent(metadata: event.room.metadata));
|
||||||
})
|
})
|
||||||
..on<SignalConnectionStateUpdatedEvent>((event) {
|
..on<SignalConnectionStateUpdatedEvent>((event) {
|
||||||
// during reconnection, need to send sync state upon signal connection.
|
// during reconnection, need to send sync state upon signal connection.
|
||||||
@@ -279,6 +286,9 @@ class Room extends DisposableChangeNotifier with EventsEmittable<RoomEvent> {
|
|||||||
String sid, lk_models.ParticipantInfo? info) {
|
String sid, lk_models.ParticipantInfo? info) {
|
||||||
RemoteParticipant? participant = _participants[sid];
|
RemoteParticipant? participant = _participants[sid];
|
||||||
if (participant != null) {
|
if (participant != null) {
|
||||||
|
if (info != null) {
|
||||||
|
participant.updateFromInfo(info);
|
||||||
|
}
|
||||||
return participant;
|
return participant;
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -323,7 +333,8 @@ class Room extends DisposableChangeNotifier with EventsEmittable<RoomEvent> {
|
|||||||
|
|
||||||
if (isNew) {
|
if (isNew) {
|
||||||
hasChanged = true;
|
hasChanged = true;
|
||||||
events.emit(ParticipantConnectedEvent(participant: participant));
|
// fire connected event
|
||||||
|
emitWhenConnected(ParticipantConnectedEvent(participant: participant));
|
||||||
} else {
|
} else {
|
||||||
await participant.updateFromInfo(info);
|
await participant.updateFromInfo(info);
|
||||||
}
|
}
|
||||||
@@ -357,7 +368,7 @@ class Room extends DisposableChangeNotifier with EventsEmittable<RoomEvent> {
|
|||||||
final activeSpeakers = lastSpeakers.values.toList();
|
final activeSpeakers = lastSpeakers.values.toList();
|
||||||
activeSpeakers.sort((a, b) => b.audioLevel.compareTo(a.audioLevel));
|
activeSpeakers.sort((a, b) => b.audioLevel.compareTo(a.audioLevel));
|
||||||
_activeSpeakers = activeSpeakers;
|
_activeSpeakers = activeSpeakers;
|
||||||
events.emit(ActiveSpeakersChangedEvent(speakers: activeSpeakers));
|
emitWhenConnected(ActiveSpeakersChangedEvent(speakers: activeSpeakers));
|
||||||
}
|
}
|
||||||
|
|
||||||
// from data channel
|
// from data channel
|
||||||
@@ -391,7 +402,7 @@ class Room extends DisposableChangeNotifier with EventsEmittable<RoomEvent> {
|
|||||||
}
|
}
|
||||||
|
|
||||||
_activeSpeakers = activeSpeakers;
|
_activeSpeakers = activeSpeakers;
|
||||||
events.emit(ActiveSpeakersChangedEvent(speakers: activeSpeakers));
|
emitWhenConnected(ActiveSpeakersChangedEvent(speakers: activeSpeakers));
|
||||||
}
|
}
|
||||||
|
|
||||||
void _onSignalConnectionQualityUpdateEvent(
|
void _onSignalConnectionQualityUpdateEvent(
|
||||||
@@ -422,6 +433,11 @@ class Room extends DisposableChangeNotifier with EventsEmittable<RoomEvent> {
|
|||||||
if (trackPublication == null) continue;
|
if (trackPublication == null) continue;
|
||||||
// update the stream state
|
// update the stream state
|
||||||
await trackPublication.updateStreamState(update.state.toLKType());
|
await trackPublication.updateStreamState(update.state.toLKType());
|
||||||
|
emitWhenConnected(TrackStreamStateUpdatedEvent(
|
||||||
|
participant: participant,
|
||||||
|
publication: trackPublication,
|
||||||
|
streamState: update.state.toLKType(),
|
||||||
|
));
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -452,7 +468,7 @@ class Room extends DisposableChangeNotifier with EventsEmittable<RoomEvent> {
|
|||||||
|
|
||||||
await participant.unpublishAllTracks(notify: true);
|
await participant.unpublishAllTracks(notify: true);
|
||||||
|
|
||||||
events.emit(ParticipantDisconnectedEvent(participant: participant));
|
emitWhenConnected(ParticipantDisconnectedEvent(participant: participant));
|
||||||
}
|
}
|
||||||
|
|
||||||
Future<void> _sendSyncState() async {
|
Future<void> _sendSyncState() async {
|
||||||
@@ -512,6 +528,13 @@ extension RoomPrivateMethods on Room {
|
|||||||
_serverVersion = null;
|
_serverVersion = null;
|
||||||
_serverRegion = null;
|
_serverRegion = null;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@internal
|
||||||
|
void emitWhenConnected(RoomEvent event) {
|
||||||
|
if (connectionState == ConnectionState.connected) {
|
||||||
|
events.emit(event);
|
||||||
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
extension RoomDebugMethods on Room {
|
extension RoomDebugMethods on Room {
|
||||||
|
|||||||
@@ -138,7 +138,6 @@ class RemoteParticipant extends Participant<RemoteTrackPublication> {
|
|||||||
@internal
|
@internal
|
||||||
Future<void> updateFromInfo(lk_models.ParticipantInfo info) async {
|
Future<void> updateFromInfo(lk_models.ParticipantInfo info) async {
|
||||||
logger.fine('RemoteParticipant.updateFromInfo(info: $info)');
|
logger.fine('RemoteParticipant.updateFromInfo(info: $info)');
|
||||||
final hadInfo = hasInfo;
|
|
||||||
super.updateFromInfo(info);
|
super.updateFromInfo(info);
|
||||||
|
|
||||||
// figuring out deltas between tracks
|
// figuring out deltas between tracks
|
||||||
@@ -168,13 +167,13 @@ class RemoteParticipant extends Participant<RemoteTrackPublication> {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// notify listeners when it's not a new participant
|
// always emit events for new publications, Room will not forward them unless it's ready
|
||||||
if (hadInfo) {
|
|
||||||
for (final pub in newPubs) {
|
for (final pub in newPubs) {
|
||||||
final event = TrackPublishedEvent(
|
final event = TrackPublishedEvent(
|
||||||
participant: this,
|
participant: this,
|
||||||
publication: pub,
|
publication: pub,
|
||||||
);
|
);
|
||||||
|
if (room.connectionState == ConnectionState.connected) {
|
||||||
[events, room.events].emit(event);
|
[events, room.events].emit(event);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -69,7 +69,6 @@ class RemoteTrackPublication<T extends RemoteTrack>
|
|||||||
_streamState = streamState;
|
_streamState = streamState;
|
||||||
[
|
[
|
||||||
participant.events,
|
participant.events,
|
||||||
participant.room.events,
|
|
||||||
].emit(TrackStreamStateUpdatedEvent(
|
].emit(TrackStreamStateUpdatedEvent(
|
||||||
participant: participant,
|
participant: participant,
|
||||||
publication: this,
|
publication: this,
|
||||||
@@ -304,7 +303,6 @@ class RemoteTrackPublication<T extends RemoteTrack>
|
|||||||
// emit events
|
// emit events
|
||||||
[
|
[
|
||||||
participant.events,
|
participant.events,
|
||||||
participant.room.events,
|
|
||||||
].emit(TrackSubscriptionPermissionChangedEvent(
|
].emit(TrackSubscriptionPermissionChangedEvent(
|
||||||
participant: participant,
|
participant: participant,
|
||||||
publication: this,
|
publication: this,
|
||||||
|
|||||||
@@ -36,10 +36,16 @@ void main() {
|
|||||||
test('participant join', () async {
|
test('participant join', () async {
|
||||||
expect(
|
expect(
|
||||||
room.events.streamCtrl.stream,
|
room.events.streamCtrl.stream,
|
||||||
emits(predicate<ParticipantConnectedEvent>(
|
emitsInOrder(<Matcher>[
|
||||||
|
predicate<TrackPublishedEvent>(
|
||||||
(event) => event.participant.sid == remoteParticipantData.sid,
|
(event) => event.participant.sid == remoteParticipantData.sid,
|
||||||
)),
|
),
|
||||||
|
predicate<ParticipantConnectedEvent>(
|
||||||
|
(event) => event.participant.sid == remoteParticipantData.sid,
|
||||||
|
)
|
||||||
|
]),
|
||||||
);
|
);
|
||||||
|
|
||||||
ws.onData(participantJoinResponse.writeToBuffer());
|
ws.onData(participantJoinResponse.writeToBuffer());
|
||||||
|
|
||||||
await room.events.waitFor<ParticipantConnectedEvent>(
|
await room.events.waitFor<ParticipantConnectedEvent>(
|
||||||
|
|||||||
Reference in New Issue
Block a user