mirror of
https://github.com/livekit/livekit.git
synced 2026-07-29 01:09:25 +00:00
handle peer subscription
This commit is contained in:
+24
-2
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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())
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
+3
-2
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user