diff --git a/pkg/rtc/forwarder.go b/pkg/rtc/forwarder.go index e1b302619..91f2c06b4 100644 --- a/pkg/rtc/forwarder.go +++ b/pkg/rtc/forwarder.go @@ -89,6 +89,9 @@ func (f *SimpleForwarder) ChannelType() ChannelType { func (f *SimpleForwarder) Start() { f.once.Do(func() { + defer func() { + recover() + }() go f.rtcpWorker() }) } diff --git a/pkg/rtc/participant.go b/pkg/rtc/participant.go index 7937f722a..2ff0d5ebf 100644 --- a/pkg/rtc/participant.go +++ b/pkg/rtc/participant.go @@ -469,6 +469,9 @@ func (p *Participant) downTracksRTCPWorker() { func (p *Participant) rtcpSendWorker() { // read from rtcpChan for pkts := range p.rtcpCh { + for _, pkt := range pkts { + logger.GetLogger().Debugw("writing RTCP", "packet", pkt) + } if err := p.peerConn.WriteRTCP(pkts); err != nil { logger.GetLogger().Errorw("could not write RTCP to participant", "participant", p.id, diff --git a/pkg/rtc/room.go b/pkg/rtc/room.go index 5a0e20cee..3aad14a3e 100644 --- a/pkg/rtc/room.go +++ b/pkg/rtc/room.go @@ -82,6 +82,8 @@ func (r *Room) Join(participant *Participant) error { "srcParticipant", p.ID()) } } + // start the workers once connectivity is established + participant.Start() } } diff --git a/pkg/rtc/track.go b/pkg/rtc/track.go index 5fa537c6a..be53fdb6b 100644 --- a/pkg/rtc/track.go +++ b/pkg/rtc/track.go @@ -17,8 +17,9 @@ import ( ) var ( - creationDelay = 500 * time.Millisecond - feedbackTypes = []webrtc.RTCPFeedback{{"goog-remb", ""}, {"nack", ""}, {"nack", "pli"}} + creationDelay = 500 * time.Millisecond + maxPLIFrequency = 1 * time.Second + feedbackTypes = []webrtc.RTCPFeedback{{"goog-remb", ""}, {"nack", ""}, {"nack", "pli"}} ) // Track represents a remoteTrack that needs to be forwarded @@ -35,6 +36,7 @@ type Track struct { forwarders map[string]Forwarder receiver *Receiver lastNack int64 + lastPLI time.Time } func NewTrack(pId string, rtcpCh chan []rtcp.Packet, track *webrtc.TrackRemote, receiver *Receiver) *Track { @@ -203,13 +205,19 @@ func (t *Track) forwardRTPWorker() { } if err == sfu.ErrRequiresKeyFrame { + delta := time.Now().Sub(t.lastPLI) + if delta < maxPLIFrequency { + continue + } logger.GetLogger().Infow("keyframe required, sending PLI") + rtcpPkts := []rtcp.Packet{ + &rtcp.PictureLossIndication{SenderSSRC: uint32(t.remoteTrack.SSRC()), MediaSSRC: pkt.SSRC}, + } // queue up a PLI, but don't block channel go func() { - t.rtcpCh <- []rtcp.Packet{ - &rtcp.PictureLossIndication{SenderSSRC: forwarder.Track().SSRC(), MediaSSRC: pkt.SSRC}, - } + t.rtcpCh <- rtcpPkts }() + t.lastPLI = time.Now() } else if err != nil { logger.GetLogger().Warnw("could not forward packet to participant", "src", t.participantId, diff --git a/pkg/sfu/downtrack.go b/pkg/sfu/downtrack.go index 8ed86476d..8c2a54407 100644 --- a/pkg/sfu/downtrack.go +++ b/pkg/sfu/downtrack.go @@ -229,8 +229,9 @@ func (d *DownTrack) writeSimpleRTP(pkt rtp.Packet) error { } } if !relay { - // TODO: how do we sent PLI to the source? - //return ErrRequiresKeyFrame + // when we are writing to a new client and there isn't a keyframe, it makes it impossible + // for clients to render a frame. we'll send an error for the track writer. + return ErrRequiresKeyFrame } } d.snOffset = pkt.SequenceNumber - d.lastSN - 1