fixed server initiated negotiation

This commit is contained in:
David Zhao
2020-11-11 00:02:37 -08:00
parent 40be24bf60
commit 256328c4ff
5 changed files with 65 additions and 28 deletions
+31 -4
View File
@@ -102,7 +102,7 @@ func NewRTCClient(conn *websocket.Conn) (*RTCClient, error) {
})
peerConn.OnTrack(func(track *webrtc.Track, r *webrtc.RTPReceiver) {
c.AppendLog("track received", "track", track.Label(), "ssrc", track.SSRC())
c.AppendLog("track received", "label", track.Label(), "id", track.ID())
// TODO: set up track consumer to read
})
@@ -214,9 +214,10 @@ func (c *RTCClient) Run() error {
c.pendingCandidates = nil
c.lock.Unlock()
case *livekit.SignalResponse_Negotiate:
c.AppendLog("received negotate answer")
answer := service.FromProtoSessionDescription(msg.Negotiate)
if err := c.PeerConn.SetRemoteDescription(answer); err != nil {
c.AppendLog("received negotate answer",
"type", msg.Negotiate.Type)
desc := service.FromProtoSessionDescription(msg.Negotiate)
if err := c.handleNegotiate(desc); err != nil {
return err
}
case *livekit.SignalResponse_Trickle:
@@ -321,6 +322,32 @@ func (c *RTCClient) Negotiate() error {
return nil
}
func (c *RTCClient) handleNegotiate(desc webrtc.SessionDescription) error {
// always set remote description
if err := c.PeerConn.SetRemoteDescription(desc); err != nil {
return err
}
if desc.Type == webrtc.SDPTypeOffer {
answer, err := c.PeerConn.CreateAnswer(nil)
if err != nil {
return err
}
if err := c.PeerConn.SetLocalDescription(answer); err != nil {
return err
}
// send remote an answer
return c.SendRequest(&livekit.SignalRequest{
Message: &livekit.SignalRequest_Negotiate{
Negotiate: service.ToProtoSessionDescription(answer),
},
})
}
return nil
}
func (c *RTCClient) AddTrack(path string, codecType webrtc.RTPCodecType, id string, label string) error {
// determine file type
format, ok := extFormatMapping[filepath.Ext(path)]
+2
View File
@@ -76,10 +76,12 @@ func (w *TrackWriter) writeOgg() {
if err == io.EOF {
logger.GetLogger().Infow("all audio samples parsed and sent")
w.onWriteComplete()
return
}
if err != nil {
logger.GetLogger().Errorw("could not parse ogg page", "err", err)
return
}
// The amount of samples is the difference between the last and current timestamp
+23 -17
View File
@@ -56,23 +56,6 @@ func NewWebRTCPeer(id string, me *MediaEngine, conf WebRTCConfig) (*WebRTCPeer,
pc.OnTrack(peer.onTrack)
pc.OnNegotiationNeeded(func() {
offer, err := pc.CreateOffer(nil)
if err != nil {
// TODO: log
return
}
err = pc.SetLocalDescription(offer)
if err != nil {
// TODO: log
return
}
if peer.OnOffer != nil {
peer.OnOffer(offer)
}
})
pc.OnICECandidate(func(c *webrtc.ICECandidate) {
if c == nil {
return
@@ -113,6 +96,25 @@ func (p *WebRTCPeer) Answer(sdp webrtc.SessionDescription) (answer webrtc.Sessio
return
}
// only set after answered
p.conn.OnNegotiationNeeded(func() {
logger.GetLogger().Debugw("negotiation needed")
offer, err := p.conn.CreateOffer(nil)
if err != nil {
// TODO: log
return
}
err = p.conn.SetLocalDescription(offer)
if err != nil {
// TODO: log
return
}
if p.OnOffer != nil {
p.OnOffer(offer)
}
})
return
}
@@ -155,6 +157,10 @@ func (p *WebRTCPeer) AddSubscriber(otherPeer *WebRTCPeer) error {
defer p.lock.RUnlock()
for _, track := range p.tracks {
logger.GetLogger().Debugw("subscribing to track",
"srcPeer", p.ID(),
"dstPeer", otherPeer.ID(),
"track", track.id)
if err := track.AddSubscriber(otherPeer); err != nil {
return err
}
+4 -2
View File
@@ -67,6 +67,8 @@ func (r *Room) Join(peerId string, sdp string) (peer *WebRTCPeer, err error) {
return nil, ErrPeerExists
}
logger.GetLogger().Infow("new peer joined", "peerId", peerId,
"roomId", r.RoomId)
offer := webrtc.SessionDescription{
Type: webrtc.SDPTypeOffer,
SDP: sdp,
@@ -83,8 +85,6 @@ func (r *Room) Join(peerId string, sdp string) (peer *WebRTCPeer, err error) {
}
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 {
@@ -95,6 +95,8 @@ func (r *Room) Join(peerId string, sdp string) (peer *WebRTCPeer, err error) {
}
}
r.peers[peerId] = peer
return
}
+5 -5
View File
@@ -16,7 +16,7 @@ var (
// Peer track represents a track that needs to be forwarded
type PeerTrack struct {
id uint32
id string
ctx context.Context
peerId string
// source track
@@ -32,7 +32,7 @@ type PeerTrack struct {
func NewPeerTrack(ctx context.Context, peerId string, rtcpWriter RTCPWriter, track *webrtc.Track, receiver *Receiver) *PeerTrack {
return &PeerTrack{
id: track.SSRC(),
id: track.ID(),
ctx: ctx,
peerId: peerId,
track: track,
@@ -111,9 +111,9 @@ func (t *PeerTrack) forwardWorker() {
}
now := time.Now()
logger.GetLogger().Debugw("read packet from track",
"peerId", t.peerId,
"track", t.track.ID())
//logger.GetLogger().Debugw("read packet from track",
// "peerId", t.peerId,
// "track", t.track.ID())
t.lock.RLock()
for dstPeerId, forwarder := range t.forwarders {
// There exists a bug in chrome where setLocalDescription