diff --git a/pkg/rtc/peer.go b/pkg/rtc/peer.go index 35318808d..31c5c346c 100644 --- a/pkg/rtc/peer.go +++ b/pkg/rtc/peer.go @@ -143,6 +143,28 @@ func (p *WebRTCPeer) Close() error { return p.conn.Close() } +// Subscribes otherPeer to all of the tracks +func (p *WebRTCPeer) AddSubscriber(otherPeer *WebRTCPeer) error { + p.lock.RLock() + defer p.lock.RUnlock() + + for _, track := range p.tracks { + if err := track.AddSubscriber(otherPeer); err != nil { + return err + } + } + return nil +} + +func (p *WebRTCPeer) RemoveSubscriber(peerId string) { + p.lock.RLock() + defer p.lock.RUnlock() + + for _, track := range p.tracks { + track.RemoveSubscriber(peerId) + } +} + // when a new track is created, creates a PeerTrack and adds it to room func (p *WebRTCPeer) onTrack(track *webrtc.Track, rtpReceiver *webrtc.RTPReceiver) { @@ -151,12 +173,12 @@ func (p *WebRTCPeer) onTrack(track *webrtc.Track, rtpReceiver *webrtc.RTPReceive pt := NewPeerTrack(p.id, track, receiver) p.lock.Lock() - defer p.lock.Unlock() p.tracks = append(p.tracks, pt) + p.lock.Unlock() if p.OnPeerTrack != nil { // caller should hook up what happens when the peer track is available - p.OnPeerTrack(p, pt) + go p.OnPeerTrack(p, pt) } } diff --git a/pkg/rtc/room.go b/pkg/rtc/room.go index ffca0f85c..80f1e459a 100644 --- a/pkg/rtc/room.go +++ b/pkg/rtc/room.go @@ -8,6 +8,7 @@ import ( "github.com/pion/webrtc/v3" "github.com/pkg/errors" + "github.com/livekit/livekit-server/pkg/logger" "github.com/livekit/livekit-server/proto/livekit" ) @@ -78,8 +79,39 @@ func (r *Room) Join(peerId string, token string, sdp string) (peer *WebRTCPeer, if err != nil { return nil, errors.Wrap(err, "could not create peer") } + peer.OnPeerTrack = r.onTrackAdded r.peers[peerId] = peer + // subscribe peer to existing tracks + for _, p := range r.peers { + if err := p.AddSubscriber(peer); err != nil { + // TODO: log error? or disconnect? + logger.GetLogger().Errorw("could not subscribe to peer", + "dstPeer", peer.ID(), + "srcPeer", p.ID()) + } + } + return } + +// a peer in the room added a new track, subscribe other peers to it +func (r *Room) onTrackAdded(peer *WebRTCPeer, track *PeerTrack) { + r.lock.RLock() + defer r.lock.RUnlock() + + // subscribe all existing peers to this track + for _, p := range r.peers { + if p == peer { + // skip publishing peer + continue + } + if err := track.AddSubscriber(peer); err != nil { + logger.GetLogger().Errorw("could not subscribe to track", + "srcPeer", peer.ID(), + "track", track.id, + "dstPeer", p.ID()) + } + } +} diff --git a/pkg/rtc/track.go b/pkg/rtc/track.go index 26af6fb21..c287c618f 100644 --- a/pkg/rtc/track.go +++ b/pkg/rtc/track.go @@ -35,13 +35,14 @@ func NewPeerTrack(peerId string, track *webrtc.Track, receiver *Receiver) *PeerT // subscribes peer to current track // creates and add necessary forwarders and starts them func (t *PeerTrack) AddSubscriber(peer *WebRTCPeer) error { + return nil } // removes peer from subscription // stop all forwarders to the peer -func (t *PeerTrack) RemoveSubscriber(peerId string) error { - return nil +func (t *PeerTrack) RemoveSubscriber(peerId string) { + } // forwardWorker reads from the receiver and writes to each sender