diff --git a/cmd/cli/client/client.go b/cmd/cli/client/client.go index df8df9ff2..5fe6b5b51 100644 --- a/cmd/cli/client/client.go +++ b/cmd/cli/client/client.go @@ -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)] diff --git a/cmd/cli/client/trackwriter.go b/cmd/cli/client/trackwriter.go index 68cbdd85d..223798587 100644 --- a/cmd/cli/client/trackwriter.go +++ b/cmd/cli/client/trackwriter.go @@ -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 diff --git a/pkg/rtc/peer.go b/pkg/rtc/peer.go index f494ab770..9ecb5d4c7 100644 --- a/pkg/rtc/peer.go +++ b/pkg/rtc/peer.go @@ -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 } diff --git a/pkg/rtc/room.go b/pkg/rtc/room.go index 5564da55c..113db39dc 100644 --- a/pkg/rtc/room.go +++ b/pkg/rtc/room.go @@ -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 } diff --git a/pkg/rtc/track.go b/pkg/rtc/track.go index da32eb7a3..ae95279b5 100644 --- a/pkg/rtc/track.go +++ b/pkg/rtc/track.go @@ -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