Migrate with muted track (#794)

* update pion

* migrate muted track

* refine
This commit is contained in:
cnderrauber
2022-06-30 00:14:29 +08:00
committed by GitHub
parent 2c48eafd6e
commit 93f152779e
2 changed files with 94 additions and 4 deletions
+3 -4
View File
@@ -252,11 +252,10 @@ func (t *MediaTrack) AddReceiver(receiver *webrtc.RTPReceiver, track *webrtc.Tra
func (t *MediaTrack) GetConnectionScore() float32 {
receiver := t.PrimaryReceiver()
if receiver == nil {
return 0.0
if rtcReceiver, ok := receiver.(*sfu.WebRTCReceiver); ok {
return rtcReceiver.GetConnectionScore()
}
return receiver.(*sfu.WebRTCReceiver).GetConnectionScore()
return 0.0
}
func (t *MediaTrack) SetRTT(rtt uint32) {
+91
View File
@@ -560,9 +560,31 @@ func (p *ParticipantImpl) HandleOffer(sdp webrtc.SessionDescription) (answer web
prometheus.ServiceOperationCounter.WithLabelValues("answer", "success", "").Add(1)
if p.MigrateState() == types.MigrateStateSync {
go p.handleMigrateMutedTrack()
}
return
}
func (p *ParticipantImpl) handleMigrateMutedTrack() {
// muted track won't send rtp packet, so we add mediatrack manually
var addedTrack []*MediaTrack
p.pendingTracksLock.Lock()
for cid, t := range p.pendingTracks {
if t.migrated && t.Muted && t.Type == livekit.TrackType_VIDEO {
addedTrack = append(addedTrack, p.addMigrateMutedTrack(cid, t.TrackInfo))
}
}
p.pendingTracksLock.Unlock()
for _, t := range addedTrack {
if t != nil {
p.handleTrackPublished(t)
}
}
}
// AddTrack is called when client intends to publish track.
// records track details and lets client know it's ok to proceed
func (p *ParticipantImpl) AddTrack(req *livekit.AddTrackRequest) {
@@ -1676,6 +1698,75 @@ func (p *ParticipantImpl) mediaTrackReceived(track *webrtc.TrackRemote, rtpRecei
return mt, newTrack
}
func (p *ParticipantImpl) addMigrateMutedTrack(cid string, t *livekit.TrackInfo) *MediaTrack {
p.params.Logger.Debugw("add migrate muted track", "cid", cid, "track", t.String())
var rtpReceiver *webrtc.RTPReceiver
for _, tr := range p.publisher.pc.GetTransceivers() {
if tr.Mid() == t.Mid {
rtpReceiver = tr.Receiver()
break
}
}
if rtpReceiver == nil {
p.params.Logger.Errorw("could not find receiver for migrated track", nil, "track", t.Sid)
return nil
}
mt := NewMediaTrack(MediaTrackParams{
TrackInfo: proto.Clone(t).(*livekit.TrackInfo),
SignalCid: cid,
SdpCid: cid,
ParticipantID: p.params.SID,
ParticipantIdentity: p.params.Identity,
RTCPChan: p.rtcpCh,
BufferFactory: p.params.Config.BufferFactory,
ReceiverConfig: p.params.Config.Receiver,
AudioConfig: p.params.AudioConfig,
VideoConfig: p.params.VideoConfig,
Telemetry: p.params.Telemetry,
Logger: LoggerWithTrack(p.params.Logger, livekit.TrackID(t.Sid)),
SubscriberConfig: p.params.Config.Subscriber,
PLIThrottleConfig: p.params.PLIThrottleConfig,
SimTracks: p.params.SimTracks,
})
mt.OnSubscribedMaxQualityChange(p.onSubscribedMaxQualityChange)
// add to published and clean up pending
p.UpTrackManager.AddPublishedTrack(mt)
delete(p.pendingTracks, cid)
mt.AddOnClose(func() {
// re-use track
p.lock.Lock()
p.unpublishedTracks = append(p.unpublishedTracks, t)
p.lock.Unlock()
})
potentialCodecs := make([]webrtc.RTPCodecParameters, 0, len(t.Codecs))
parameters := rtpReceiver.GetParameters()
for _, c := range t.Codecs {
for _, nc := range parameters.Codecs {
if strings.EqualFold(nc.MimeType, c.MimeType) {
potentialCodecs = append(potentialCodecs, nc)
break
}
}
}
mt.SetPotentialCodecs(potentialCodecs, parameters.HeaderExtensions)
for _, codec := range t.Codecs {
for ssrc, info := range p.params.SimTracks {
if info.Mid == codec.Mid {
mt.MediaTrackReceiver.SetLayerSsrc(codec.MimeType, info.Rid, ssrc)
}
}
}
mt.SetSimulcast(t.Simulcast)
mt.SetMuted(true)
return mt
}
func (p *ParticipantImpl) handleTrackPublished(track types.MediaTrack) {
if !p.hasPendingMigratedTrack() {
p.SetMigrateState(types.MigrateStateComplete)