From 57c2ea91e5f1b6777d8c83415724062d3a0535d8 Mon Sep 17 00:00:00 2001 From: David Zhao Date: Fri, 23 Jul 2021 16:11:38 -0700 Subject: [PATCH] Room implementation, library exports --- lib/livekit_client.dart | 21 +- lib/src/participant/remote_participant.dart | 9 +- lib/src/room.dart | 257 +++++++++++++++++++- lib/src/rtc_engine.dart | 7 +- pubspec.lock | 14 ++ pubspec.yaml | 1 + 6 files changed, 289 insertions(+), 20 deletions(-) diff --git a/lib/livekit_client.dart b/lib/livekit_client.dart index d6e6cea..3e77ad8 100644 --- a/lib/livekit_client.dart +++ b/lib/livekit_client.dart @@ -1,7 +1,18 @@ library livekit_client; -/// A Calculator. -class Calculator { - /// Returns [value] plus 1. - int addOne(int value) => value + 1; -} +export 'src/errors.dart'; +export 'src/room.dart'; +export 'src/rtc_engine.dart'; +export 'src/participant/participant.dart'; +export 'src/participant/local_participant.dart'; +export 'src/participant/remote_participant.dart'; +export 'src/proto/livekit_models.pbenum.dart'; +export 'src/proto/livekit_rtc.pbenum.dart'; +export 'src/participant/local_participant.dart'; +export 'src/track/track.dart'; +export 'src/track/video_track.dart'; +export 'src/track/local_audio_track.dart'; +export 'src/track/local_video_track.dart'; +export 'src/track/track_publication.dart'; +export 'src/track/local_track_publication.dart'; +export 'src/track/remote_track_publication.dart'; diff --git a/lib/src/participant/remote_participant.dart b/lib/src/participant/remote_participant.dart index 0341322..87ea531 100644 --- a/lib/src/participant/remote_participant.dart +++ b/lib/src/participant/remote_participant.dart @@ -26,7 +26,14 @@ class RemoteParticipant extends Participant { } } - addSubscribedMediaTrack(MediaStreamTrack mediaTrack, String sid) async { + addSubscribedMediaTrack(MediaStreamTrack mediaTrack, String? sid) async { + if (sid == null) { + var msg = 'addSubscribedMediaTrack received null sid'; + delegate?.onTrackSubscriptionFailed(this, '', msg); + roomDelegate?.onTrackSubscriptionFailed(this, '', msg); + return; + } + var pub = getTrackPublication(sid); if (pub == null) { // we may have received the track prior to metadata. wait up to 3s diff --git a/lib/src/room.dart b/lib/src/room.dart index 4482626..a5b3dd1 100644 --- a/lib/src/room.dart +++ b/lib/src/room.dart @@ -1,6 +1,8 @@ import 'dart:async'; import 'package:flutter_webrtc/flutter_webrtc.dart'; +import 'package:tuple/tuple.dart'; +import 'errors.dart'; import 'extensions.dart'; import 'logger.dart'; import 'participant/local_participant.dart'; @@ -10,6 +12,9 @@ import 'proto/livekit_models.pb.dart'; import 'proto/livekit_rtc.pb.dart'; import 'rtc_engine.dart'; import 'signal_client.dart'; +import 'track/remote_track_publication.dart'; +import 'track/track.dart'; +import 'track/track_publication.dart'; enum RoomState { Disconnected, @@ -17,7 +22,33 @@ enum RoomState { Reconnecting, } -class Room { +mixin RoomDelegate { + // room level callbacks + void onReconnecting() {} + void onReconnected() {} + void onDisconnected() {} + void onParticipantConnected(Participant participant) {} + void onParticipantDisconnected(Participant participant) {} + void onActiveSpeakersChanged(List participants) {} + + // callbacks about participant events + void onMetadataChanged(Participant participant) {} + void onTrackMuted(Participant participant, TrackPublication publication) {} + void onTrackUnmuted(Participant participant, TrackPublication publication) {} + void onTrackPublished( + RemoteParticipant participant, RemoteTrackPublication publication) {} + void onTrackUnpublished( + RemoteParticipant participant, RemoteTrackPublication publication) {} + void onTrackSubscribed(RemoteParticipant participant, Track track, + RemoteTrackPublication publication) {} + void onTrackUnsubscribed(RemoteParticipant participant, Track track, + RemoteTrackPublication publication) {} + void onDataReceived(RemoteParticipant participant, List data) {} + void onTrackSubscriptionFailed( + RemoteParticipant participant, String sid, String? message) {} +} + +class Room with ParticipantDelegate { RoomState state = RoomState.Disconnected; /// map of SID to RemoteParticipant @@ -35,6 +66,9 @@ class Room { /// a list of participants that are actively speaking, including local participant. List activeSpeakers = []; + /// delegate for room events + RoomDelegate? delegate; + RTCEngine _engine; Completer? _connectCompleter; @@ -42,6 +76,7 @@ class Room { Room(SignalClient client, RTCConfiguration? rtcConfig) : _engine = new RTCEngine(client, rtcConfig) { _engine.onTrack = _onTrackAdded; + _engine.onICEConnected = _handleICEConnected; _engine.onDisconnected = _handleDisconnect; _engine.onParticipantUpdateCallback = _handleParticipantUpdate; _engine.onActiveSpeakerchangedCallback = _handleSpeakerUpdate; @@ -50,7 +85,7 @@ class Room { // TODO: handle reconnecting & reconnected events } - Future _connect(String url, String token, JoinOptions? opts) async { + Future connect(String url, String token, JoinOptions? opts) async { var completer = new Completer(); _connectCompleter = completer; @@ -59,20 +94,224 @@ class Room { 'connected to LiveKit server, version: ${joinResponse.serverVersion}'); state = RoomState.Connected; - var pi = joinResponse.participant; - localParticipant = new LocalParticipant(pi.sid, pi.identity, _engine); + localParticipant = new LocalParticipant( + engine: _engine, + info: joinResponse.participant, + ); + localParticipant.roomDelegate = this; + + sid = joinResponse.room.sid; + name = joinResponse.room.name; + + for (var info in joinResponse.otherParticipants) { + _getOrCreateRemoteParticipant(info.sid, info); + } + + // room is not ready until ICE is connected. so we would return a completer for now + // if it times out, we'll fail the completer + Timer(Duration(seconds: 5), () { + _connectCompleter?.completeError(ConnectError()); + _connectCompleter = null; + }); return completer.future; } - _handleDisconnect() {} + disconnect() { + _engine.client.sendLeave(); + _handleDisconnect(); + } - _handleParticipantUpdate(List participants) {} + RemoteParticipant _getOrCreateRemoteParticipant( + String sid, ParticipantInfo? info) { + var participant = participants[sid]; + if (participant != null) { + return participant; + } - _handleSpeakerUpdate(List speakers) {} + if (info == null) { + participant = RemoteParticipant(_engine.client, sid, ''); + } else { + participant = RemoteParticipant.fromInfo(_engine.client, info); + } + participant.roomDelegate = this; + participants[sid] = participant; - _handleDataPacket(UserPacket packet, DataPacket_Kind kind) {} + return participant; + } + + _handleICEConnected() { + _connectCompleter?.complete(this); + _connectCompleter = null; + } + + _handleDisconnect() { + if (state == RoomState.Disconnected) { + return; + } + + for (var p in participants.values) { + for (var pub in p.tracks.values) { + p.unpublishTrack(pub.sid); + } + } + for (var pub in localParticipant.tracks.values) { + pub.track?.stop(); + } + + _engine.close(); + participants.clear(); + activeSpeakers.clear(); + state = RoomState.Disconnected; + delegate?.onDisconnected(); + } + + _handleParticipantUpdate(List updates) { + for (var info in updates) { + if (localParticipant.sid == info.sid) { + localParticipant.updateFromInfo(info); + continue; + } + + if (info.state == ParticipantInfo_State.DISCONNECTED) { + _handleParticipantDisconnect(info.sid); + continue; + } + + var isNew = !participants.containsKey(info.sid); + var participant = _getOrCreateRemoteParticipant(info.sid, info); + + if (isNew) { + delegate?.onParticipantConnected(participant); + } else { + participant.updateFromInfo(info); + } + } + } + + _handleSpeakerUpdate(List speakers) { + var seenSids = Set(); + List newSpeakers = []; + for (var info in speakers) { + seenSids.add(info.sid); + + if (info.sid == localParticipant.sid) { + localParticipant.audioLevel = info.level; + localParticipant.isSpeaking = true; + newSpeakers.add(localParticipant); + continue; + } + + var participant = participants[info.sid]; + if (participant != null) { + participant.audioLevel = info.level; + participant.isSpeaking = true; + newSpeakers.add(participant); + } + } + + // clear previous speakers + if (seenSids.contains(localParticipant.sid)) { + localParticipant.audioLevel = 0; + localParticipant.isSpeaking = false; + } + for (var participant in participants.values) { + if (!seenSids.contains(participant.sid)) { + participant.audioLevel = 0; + participant.isSpeaking = false; + } + } + activeSpeakers = newSpeakers; + delegate?.onActiveSpeakersChanged(newSpeakers); + } + + _handleDataPacket(UserPacket packet, DataPacket_Kind kind) { + var participant = participants[packet.participantSid]; + if (participant == null) { + return; + } + + participant.delegate?.onDataReceived(participant, packet.payload); + delegate?.onDataReceived(participant, packet.payload); + } _onTrackAdded( - MediaStreamTrack track, MediaStream? stream, RTCRtpReceiver? receiver) {} + MediaStreamTrack track, MediaStream? stream, RTCRtpReceiver? receiver) { + if (stream == null) { + // we need the stream to get the track's id + logger.severe('received track without mediastream'); + return; + } + + var parsed = unpackStreamId(stream.id); + var trackSid = parsed.item2; + if (trackSid == null) { + trackSid = track.id; + } + + var participant = _getOrCreateRemoteParticipant(parsed.item1, null); + participant.addSubscribedMediaTrack(track, trackSid); + } + + _handleParticipantDisconnect(String sid) { + var participant = participants.remove(sid); + if (participant == null) { + return; + } + + for (var track in participant.tracks.values) { + participant.unpublishTrack(track.sid, true); + } + delegate?.onParticipantDisconnected(participant); + } + + //----------------- forward participant delegate calls ---------------------// + + void onMetadataChanged(Participant participant) { + delegate?.onMetadataChanged(participant); + } + + void onTrackMuted(Participant participant, TrackPublication publication) { + delegate?.onTrackMuted(participant, publication); + } + + void onTrackUnmuted(Participant participant, TrackPublication publication) { + delegate?.onTrackUnmuted(participant, publication); + } + + void onTrackPublished( + RemoteParticipant participant, RemoteTrackPublication publication) { + delegate?.onTrackPublished(participant, publication); + } + + void onTrackUnpublished( + RemoteParticipant participant, RemoteTrackPublication publication) { + delegate?.onTrackUnpublished(participant, publication); + } + + void onTrackSubscribed(RemoteParticipant participant, Track track, + RemoteTrackPublication publication) { + delegate?.onTrackSubscribed(participant, track, publication); + } + + void onTrackUnsubscribed(RemoteParticipant participant, Track track, + RemoteTrackPublication publication) { + delegate?.onTrackUnsubscribed(participant, track, publication); + } + + // omitted because data dispatching is handled in _handleDataPacket + void onDataReceived(RemoteParticipant participant, List data) {} + + void onTrackSubscriptionFailed( + RemoteParticipant participant, String sid, String? message) { + delegate?.onTrackSubscriptionFailed(participant, sid, message); + } +} + +Tuple2 unpackStreamId(String streamId) { + var parts = streamId.split('|'); + if (parts.length != 2) { + return Tuple2(parts[0], null); + } + return Tuple2(parts[0], parts[1]); } diff --git a/lib/src/rtc_engine.dart b/lib/src/rtc_engine.dart index ee38db3..95b60ae 100644 --- a/lib/src/rtc_engine.dart +++ b/lib/src/rtc_engine.dart @@ -286,11 +286,8 @@ class RTCEngine with SignalClientDelegate { } void onLocalTrackPublished(TrackPublishedResponse response) { - var completer = pendingTrackResolvers[response.cid]; - if (completer != null) { - completer.complete(Future.value(response.track)); - pendingTrackResolvers.remove(response.cid); - } + var completer = pendingTrackResolvers.remove(response.cid); + completer?.complete(Future.value(response.track)); } void onActiveSpeakersChanged(List speakers) { diff --git a/pubspec.lock b/pubspec.lock index 807fd1f..cec208d 100644 --- a/pubspec.lock +++ b/pubspec.lock @@ -186,6 +186,13 @@ packages: url: "https://pub.dartlang.org" source: hosted version: "2.0.0" + quiver: + dependency: transitive + description: + name: quiver + url: "https://pub.dartlang.org" + source: hosted + version: "3.0.1" sky_engine: dependency: transitive description: flutter @@ -233,6 +240,13 @@ packages: url: "https://pub.dartlang.org" source: hosted version: "0.3.0" + tuple: + dependency: "direct main" + description: + name: tuple + url: "https://pub.dartlang.org" + source: hosted + version: "2.0.0" typed_data: dependency: transitive description: diff --git a/pubspec.yaml b/pubspec.yaml index b67fc47..247bef1 100644 --- a/pubspec.yaml +++ b/pubspec.yaml @@ -13,6 +13,7 @@ dependencies: flutter_webrtc: ^0.6.4 logging: ^1.0.1 protobuf: ^2.0.0 + tuple: ^2.0.0 web_socket_channel: ^2.1.0 uuid: ^3.0.4