mirror of
https://github.com/livekit/livekit.git
synced 2026-08-28 00:44:12 +00:00
Revert data track change (#513)
* Revert data track change * clean code
This commit is contained in:
@@ -1,198 +0,0 @@
|
||||
package rtc
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"sync"
|
||||
|
||||
"github.com/livekit/livekit-server/pkg/sfu"
|
||||
"github.com/livekit/livekit-server/pkg/sfu/buffer"
|
||||
"github.com/livekit/protocol/livekit"
|
||||
"github.com/livekit/protocol/logger"
|
||||
"github.com/pion/webrtc/v3"
|
||||
"google.golang.org/protobuf/proto"
|
||||
)
|
||||
|
||||
type DataTrackSender interface {
|
||||
sfu.TrackSender
|
||||
Write(label string, data []byte)
|
||||
}
|
||||
|
||||
type DataTrack struct {
|
||||
trackID livekit.TrackID
|
||||
participantID livekit.ParticipantID
|
||||
logger logger.Logger
|
||||
lock sync.RWMutex
|
||||
downTracks []DataTrackSender
|
||||
onDataPacket func(*livekit.DataPacket)
|
||||
onClose []func()
|
||||
}
|
||||
|
||||
func NewDataTrack(trackID livekit.TrackID, participantID livekit.ParticipantID, logger logger.Logger) *DataTrack {
|
||||
t := &DataTrack{
|
||||
trackID: trackID,
|
||||
participantID: participantID,
|
||||
logger: logger,
|
||||
}
|
||||
return t
|
||||
}
|
||||
|
||||
func (t *DataTrack) onData(label string, data []byte) {
|
||||
t.lock.RLock()
|
||||
f := t.onDataPacket
|
||||
dts := t.downTracks
|
||||
t.lock.RUnlock()
|
||||
|
||||
for _, dt := range dts {
|
||||
dt.Write(label, data)
|
||||
}
|
||||
|
||||
if f != nil {
|
||||
dp, err := DataPacketFromBytes(label, data)
|
||||
if err != nil {
|
||||
t.logger.Warnw("invalid data", err, "label", label)
|
||||
return
|
||||
}
|
||||
// only forward on user payloads
|
||||
switch payload := dp.Value.(type) {
|
||||
case *livekit.DataPacket_User:
|
||||
payload.User.ParticipantSid = string(t.participantID)
|
||||
f(dp)
|
||||
default:
|
||||
t.logger.Warnw("received unsupported data packet", nil, "payload", payload)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func (t *DataTrack) OnDataPacket(f func(*livekit.DataPacket)) {
|
||||
t.lock.Lock()
|
||||
t.onDataPacket = f
|
||||
t.lock.Unlock()
|
||||
}
|
||||
|
||||
func (t *DataTrack) ID() livekit.TrackID {
|
||||
return t.trackID
|
||||
}
|
||||
|
||||
func (t *DataTrack) TrackID() livekit.TrackID {
|
||||
return t.trackID
|
||||
}
|
||||
|
||||
func (t *DataTrack) Write(label string, data []byte) {
|
||||
t.onData(label, data)
|
||||
}
|
||||
|
||||
func (t *DataTrack) AddDownTrack(dt sfu.TrackSender) error {
|
||||
dataDt, ok := dt.(DataTrackSender)
|
||||
if !ok {
|
||||
return errors.New("invalid DownTrack type, expect DataTrackSender")
|
||||
}
|
||||
t.lock.Lock()
|
||||
defer t.lock.Unlock()
|
||||
t.downTracks = append(t.downTracks, dataDt)
|
||||
return nil
|
||||
}
|
||||
|
||||
func (t *DataTrack) DeleteDownTrack(peerID livekit.ParticipantID) {
|
||||
t.lock.Lock()
|
||||
defer t.lock.Unlock()
|
||||
for k, v := range t.downTracks {
|
||||
if v.PeerID() == peerID {
|
||||
t.downTracks[k] = t.downTracks[len(t.downTracks)-1]
|
||||
t.downTracks = t.downTracks[:len(t.downTracks)-1]
|
||||
break
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func (t *DataTrack) AddOnClose(f func()) {
|
||||
if f == nil {
|
||||
return
|
||||
}
|
||||
t.lock.Lock()
|
||||
t.onClose = append(t.onClose, f)
|
||||
t.lock.Unlock()
|
||||
}
|
||||
|
||||
func (t *DataTrack) Close() {
|
||||
t.lock.Lock()
|
||||
fs := t.onClose
|
||||
t.lock.Unlock()
|
||||
|
||||
for _, f := range fs {
|
||||
f()
|
||||
}
|
||||
}
|
||||
|
||||
func (t *DataTrack) Receiver() sfu.TrackReceiver {
|
||||
return t
|
||||
}
|
||||
|
||||
func (t *DataTrack) ToProto() *livekit.TrackInfo {
|
||||
return &livekit.TrackInfo{
|
||||
Sid: string(t.trackID),
|
||||
Type: livekit.TrackType_DATA,
|
||||
}
|
||||
}
|
||||
|
||||
func DataPacketFromBytes(label string, data []byte) (*livekit.DataPacket, error) {
|
||||
dp := livekit.DataPacket{}
|
||||
if err := proto.Unmarshal(data, &dp); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
switch label {
|
||||
case reliableDataChannel:
|
||||
dp.Kind = livekit.DataPacket_RELIABLE
|
||||
case lossyDataChannel:
|
||||
dp.Kind = livekit.DataPacket_LOSSY
|
||||
default:
|
||||
return nil, errors.New("unsupported datachannel added")
|
||||
}
|
||||
|
||||
return &dp, nil
|
||||
}
|
||||
|
||||
func (t *DataTrack) Kind() livekit.TrackType {
|
||||
return livekit.TrackType_DATA
|
||||
}
|
||||
|
||||
//---------------------------------------------
|
||||
// no op methods for sfu.TrackReceiver
|
||||
func (t *DataTrack) StreamID() string {
|
||||
return ""
|
||||
}
|
||||
|
||||
func (t *DataTrack) Codec() webrtc.RTPCodecCapability {
|
||||
return webrtc.RTPCodecCapability{}
|
||||
}
|
||||
|
||||
func (t *DataTrack) ReadRTP(buf []byte, layer uint8, sn uint16) (int, error) {
|
||||
return 0, nil
|
||||
}
|
||||
|
||||
func (t *DataTrack) GetSenderReportTime(layer int32) (rtpTS uint32, ntpTS buffer.NtpTime) {
|
||||
return
|
||||
}
|
||||
|
||||
func (t *DataTrack) GetBitrateTemporalCumulative() sfu.Bitrates {
|
||||
return sfu.Bitrates{}
|
||||
}
|
||||
|
||||
func (t *DataTrack) SendPLI(layer int32) {
|
||||
}
|
||||
|
||||
func (t *DataTrack) LastPLI() int64 {
|
||||
return 0
|
||||
}
|
||||
|
||||
func (t *DataTrack) SetUpTrackPaused(paused bool) {
|
||||
|
||||
}
|
||||
|
||||
func (t *DataTrack) SetMaxExpectedSpatialLayer(layer int32) {
|
||||
|
||||
}
|
||||
|
||||
func (t *DataTrack) DebugInfo() map[string]interface{} {
|
||||
return map[string]interface{}{}
|
||||
}
|
||||
+47
-85
@@ -31,10 +31,11 @@ import (
|
||||
)
|
||||
|
||||
const (
|
||||
lossyDataChannel = "_lossy"
|
||||
reliableDataChannel = "_reliable"
|
||||
sdBatchSize = 20
|
||||
rttUpdateInterval = 5 * time.Second
|
||||
LossyDataChannel = "_lossy"
|
||||
ReliableDataChannel = "_reliable"
|
||||
|
||||
sdBatchSize = 20
|
||||
rttUpdateInterval = 5 * time.Second
|
||||
)
|
||||
|
||||
type pendingTrackInfo struct {
|
||||
@@ -110,21 +111,18 @@ type ParticipantImpl struct {
|
||||
updateLock sync.Mutex
|
||||
version atomic.Uint32
|
||||
|
||||
dataTrack *DataTrack
|
||||
|
||||
// callbacks & handlers
|
||||
onTrackPublished func(types.LocalParticipant, types.MediaTrack)
|
||||
onTrackUpdated func(types.LocalParticipant, types.MediaTrack)
|
||||
onStateChange func(p types.LocalParticipant, oldState livekit.ParticipantInfo_State)
|
||||
onMetadataUpdate func(types.LocalParticipant)
|
||||
onDataPacket func(types.LocalParticipant, *livekit.DataPacket)
|
||||
|
||||
migrateState atomic.Value // types.MigrateState
|
||||
pendingOffer *webrtc.SessionDescription
|
||||
pendingDataChannels []*livekit.DataChannelInfo
|
||||
onClose func(types.LocalParticipant, map[livekit.TrackID]livekit.ParticipantID)
|
||||
onClaimsChanged func(participant types.LocalParticipant)
|
||||
|
||||
onDataTrackPublished func(types.LocalParticipant, types.DataTrack)
|
||||
}
|
||||
|
||||
func NewParticipant(params ParticipantParams, perms *livekit.ParticipantPermission) (*ParticipantImpl, error) {
|
||||
@@ -200,14 +198,14 @@ func NewParticipant(params ParticipantParams, perms *livekit.ParticipantPermissi
|
||||
primaryPC = p.subscriber.pc
|
||||
ordered := true
|
||||
// also create data channels for subs
|
||||
p.reliableDCSub, err = primaryPC.CreateDataChannel(reliableDataChannel, &webrtc.DataChannelInit{
|
||||
p.reliableDCSub, err = primaryPC.CreateDataChannel(ReliableDataChannel, &webrtc.DataChannelInit{
|
||||
Ordered: &ordered,
|
||||
})
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
retransmits := uint16(0)
|
||||
p.lossyDCSub, err = primaryPC.CreateDataChannel(lossyDataChannel, &webrtc.DataChannelInit{
|
||||
p.lossyDCSub, err = primaryPC.CreateDataChannel(LossyDataChannel, &webrtc.DataChannelInit{
|
||||
Ordered: &ordered,
|
||||
MaxRetransmits: &retransmits,
|
||||
})
|
||||
@@ -299,7 +297,7 @@ func (p *ParticipantImpl) SetPermission(permission *livekit.ParticipantPermissio
|
||||
}
|
||||
}
|
||||
|
||||
func (p *ParticipantImpl) ToProto(mediaTrackOnly bool) *livekit.ParticipantInfo {
|
||||
func (p *ParticipantImpl) ToProto() *livekit.ParticipantInfo {
|
||||
info := &livekit.ParticipantInfo{
|
||||
Sid: string(p.params.SID),
|
||||
Identity: string(p.params.Identity),
|
||||
@@ -315,12 +313,6 @@ func (p *ParticipantImpl) ToProto(mediaTrackOnly bool) *livekit.ParticipantInfo
|
||||
info.Metadata = p.params.Grants.Metadata
|
||||
}
|
||||
|
||||
p.lock.RLock()
|
||||
if !mediaTrackOnly && p.dataTrack != nil {
|
||||
info.Tracks = append(info.Tracks, p.dataTrack.ToProto())
|
||||
}
|
||||
p.lock.RUnlock()
|
||||
|
||||
return info
|
||||
}
|
||||
|
||||
@@ -354,6 +346,10 @@ func (p *ParticipantImpl) OnMetadataUpdate(callback func(types.LocalParticipant)
|
||||
p.onMetadataUpdate = callback
|
||||
}
|
||||
|
||||
func (p *ParticipantImpl) OnDataPacket(callback func(types.LocalParticipant, *livekit.DataPacket)) {
|
||||
p.onDataPacket = callback
|
||||
}
|
||||
|
||||
func (p *ParticipantImpl) OnClose(callback func(types.LocalParticipant, map[livekit.TrackID]livekit.ParticipantID)) {
|
||||
p.onClose = callback
|
||||
}
|
||||
@@ -501,16 +497,6 @@ func (p *ParticipantImpl) Close(sendLeave bool) error {
|
||||
})
|
||||
}
|
||||
|
||||
var dt *DataTrack
|
||||
p.lock.Lock()
|
||||
dt = p.dataTrack
|
||||
p.dataTrack = nil
|
||||
p.lock.Unlock()
|
||||
|
||||
if dt != nil {
|
||||
dt.Close()
|
||||
}
|
||||
|
||||
p.UpTrackManager.Close()
|
||||
|
||||
p.pendingTracksLock.Lock()
|
||||
@@ -616,7 +602,7 @@ func (p *ParticipantImpl) SendJoinResponse(
|
||||
Message: &livekit.SignalResponse_Join{
|
||||
Join: &livekit.JoinResponse{
|
||||
Room: roomInfo,
|
||||
Participant: p.ToProto(true),
|
||||
Participant: p.ToProto(),
|
||||
OtherParticipants: otherParticipants,
|
||||
ServerVersion: version.Version,
|
||||
ServerRegion: region,
|
||||
@@ -1089,53 +1075,50 @@ func (p *ParticipantImpl) onMediaTrack(track *webrtc.TrackRemote, rtpReceiver *w
|
||||
}
|
||||
}
|
||||
|
||||
func (p *ParticipantImpl) OnDataTrackPublished(f func(types.LocalParticipant, types.DataTrack)) {
|
||||
p.onDataTrackPublished = f
|
||||
}
|
||||
|
||||
func (p *ParticipantImpl) onDataChannel(dc *webrtc.DataChannel) {
|
||||
p.lock.Lock()
|
||||
created := p.onDataChannelLocked(dc)
|
||||
dt := p.dataTrack
|
||||
p.lock.Unlock()
|
||||
|
||||
if created && p.onDataTrackPublished != nil && dt != nil {
|
||||
p.onDataTrackPublished(p, dt)
|
||||
}
|
||||
}
|
||||
|
||||
func (p *ParticipantImpl) onDataChannelLocked(dc *webrtc.DataChannel) (created bool) {
|
||||
if p.State() == livekit.ParticipantInfo_DISCONNECTED {
|
||||
return
|
||||
}
|
||||
if p.dataTrack == nil {
|
||||
p.dataTrack = NewDataTrack(livekit.TrackID("DT_"+p.params.SID), p.params.SID, p.params.Logger)
|
||||
created = true
|
||||
}
|
||||
|
||||
label := dc.Label()
|
||||
switch label {
|
||||
case reliableDataChannel:
|
||||
switch dc.Label() {
|
||||
case ReliableDataChannel:
|
||||
p.reliableDC = dc
|
||||
dc.OnMessage(func(msg webrtc.DataChannelMessage) {
|
||||
if !p.CanPublishData() {
|
||||
return
|
||||
if p.CanPublishData() {
|
||||
p.handleDataMessage(livekit.DataPacket_RELIABLE, msg.Data)
|
||||
}
|
||||
p.dataTrack.Write(label, msg.Data)
|
||||
})
|
||||
case lossyDataChannel:
|
||||
case LossyDataChannel:
|
||||
p.lossyDC = dc
|
||||
dc.OnMessage(func(msg webrtc.DataChannelMessage) {
|
||||
if !p.CanPublishData() {
|
||||
return
|
||||
if p.CanPublishData() {
|
||||
p.handleDataMessage(livekit.DataPacket_LOSSY, msg.Data)
|
||||
}
|
||||
p.dataTrack.Write(label, msg.Data)
|
||||
})
|
||||
default:
|
||||
p.params.Logger.Warnw("unsupported datachannel added", nil, "label", dc.Label())
|
||||
}
|
||||
}
|
||||
|
||||
return
|
||||
func (p *ParticipantImpl) handleDataMessage(kind livekit.DataPacket_Kind, data []byte) {
|
||||
dp := livekit.DataPacket{}
|
||||
if err := proto.Unmarshal(data, &dp); err != nil {
|
||||
p.params.Logger.Warnw("could not parse data packet", err)
|
||||
return
|
||||
}
|
||||
|
||||
// trust the channel that it came in as the source of truth
|
||||
dp.Kind = kind
|
||||
|
||||
// only forward on user payloads
|
||||
switch payload := dp.Value.(type) {
|
||||
case *livekit.DataPacket_User:
|
||||
if p.onDataPacket != nil {
|
||||
payload.User.ParticipantSid = string(p.params.SID)
|
||||
p.onDataPacket(p, &dp)
|
||||
}
|
||||
default:
|
||||
p.params.Logger.Warnw("received unsupported data packet", nil, "payload", payload)
|
||||
}
|
||||
}
|
||||
|
||||
func (p *ParticipantImpl) handlePrimaryStateChange(state webrtc.PeerConnectionState) {
|
||||
@@ -1638,20 +1621,7 @@ func (p *ParticipantImpl) DebugInfo() map[string]interface{} {
|
||||
return info
|
||||
}
|
||||
|
||||
func (p *ParticipantImpl) GetDataTrack() types.DataTrack {
|
||||
p.lock.RLock()
|
||||
defer p.lock.RUnlock()
|
||||
|
||||
if dt := p.dataTrack; dt != nil {
|
||||
return dt
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (p *ParticipantImpl) handlePendingDataChannels() {
|
||||
p.lock.Lock()
|
||||
created := false
|
||||
ordered := true
|
||||
negotiated := true
|
||||
for _, ci := range p.pendingDataChannels {
|
||||
@@ -1659,18 +1629,18 @@ func (p *ParticipantImpl) handlePendingDataChannels() {
|
||||
dc *webrtc.DataChannel
|
||||
err error
|
||||
)
|
||||
if ci.Label == lossyDataChannel && p.lossyDC == nil {
|
||||
if ci.Label == LossyDataChannel && p.lossyDC == nil {
|
||||
retransmits := uint16(0)
|
||||
id := uint16(ci.GetId())
|
||||
dc, err = p.publisher.pc.CreateDataChannel(lossyDataChannel, &webrtc.DataChannelInit{
|
||||
dc, err = p.publisher.pc.CreateDataChannel(LossyDataChannel, &webrtc.DataChannelInit{
|
||||
Ordered: &ordered,
|
||||
MaxRetransmits: &retransmits,
|
||||
Negotiated: &negotiated,
|
||||
ID: &id,
|
||||
})
|
||||
} else if ci.Label == reliableDataChannel && p.reliableDC == nil {
|
||||
} else if ci.Label == ReliableDataChannel && p.reliableDC == nil {
|
||||
id := uint16(ci.GetId())
|
||||
dc, err = p.publisher.pc.CreateDataChannel(reliableDataChannel, &webrtc.DataChannelInit{
|
||||
dc, err = p.publisher.pc.CreateDataChannel(ReliableDataChannel, &webrtc.DataChannelInit{
|
||||
Ordered: &ordered,
|
||||
Negotiated: &negotiated,
|
||||
ID: &id,
|
||||
@@ -1679,18 +1649,10 @@ func (p *ParticipantImpl) handlePendingDataChannels() {
|
||||
if err != nil {
|
||||
p.params.Logger.Errorw("create migrated data channel failed", err, "label", ci.Label)
|
||||
} else if dc != nil {
|
||||
creating := p.onDataChannelLocked(dc)
|
||||
created = created || creating
|
||||
p.onDataChannel(dc)
|
||||
}
|
||||
}
|
||||
p.pendingDataChannels = nil
|
||||
|
||||
dt := p.dataTrack
|
||||
p.lock.Unlock()
|
||||
|
||||
if created && p.onDataTrackPublished != nil && dt != nil {
|
||||
p.onDataTrackPublished(p, dt)
|
||||
}
|
||||
}
|
||||
|
||||
func (p *ParticipantImpl) GetSubscribedTracks() []types.SubscribedTrack {
|
||||
|
||||
@@ -170,9 +170,9 @@ func TestOutOfOrderUpdates(t *testing.T) {
|
||||
p := newParticipantForTest("test")
|
||||
p.SetMetadata("initial metadata")
|
||||
sink := p.GetResponseSink().(*routingfakes.FakeMessageSink)
|
||||
pi1 := p.ToProto(true)
|
||||
pi1 := p.ToProto()
|
||||
p.SetMetadata("second update")
|
||||
pi2 := p.ToProto(true)
|
||||
pi2 := p.ToProto()
|
||||
|
||||
require.Greater(t, pi2.Version, pi1.Version)
|
||||
|
||||
@@ -208,7 +208,7 @@ func TestDisconnectTiming(t *testing.T) {
|
||||
|
||||
func TestCorrectJoinedAt(t *testing.T) {
|
||||
p := newParticipantForTest("test")
|
||||
info := p.ToProto(true)
|
||||
info := p.ToProto()
|
||||
require.NotZero(t, info.JoinedAt)
|
||||
require.True(t, time.Now().Unix()-info.JoinedAt <= 1)
|
||||
}
|
||||
|
||||
+3
-7
@@ -221,11 +221,7 @@ func (r *Room) Join(participant types.LocalParticipant, opts *ParticipantOptions
|
||||
})
|
||||
participant.OnTrackUpdated(r.onTrackUpdated)
|
||||
participant.OnMetadataUpdate(r.onParticipantMetadataUpdate)
|
||||
participant.OnDataTrackPublished(func(p types.LocalParticipant, dt types.DataTrack) {
|
||||
dt.OnDataPacket(func(dp *livekit.DataPacket) {
|
||||
r.onDataPacket(p, dp)
|
||||
})
|
||||
})
|
||||
participant.OnDataPacket(r.onDataPacket)
|
||||
r.Logger.Infow("new participant joined",
|
||||
"pID", participant.ID(),
|
||||
"participant", participant.Identity(),
|
||||
@@ -244,7 +240,7 @@ func (r *Room) Join(participant types.LocalParticipant, opts *ParticipantOptions
|
||||
otherParticipants := make([]*livekit.ParticipantInfo, 0, len(r.participants))
|
||||
for _, p := range r.participants {
|
||||
if p.ID() != participant.ID() && !p.Hidden() {
|
||||
otherParticipants = append(otherParticipants, p.ToProto(true))
|
||||
otherParticipants = append(otherParticipants, p.ToProto())
|
||||
}
|
||||
}
|
||||
|
||||
@@ -332,7 +328,7 @@ func (r *Room) RemoveParticipant(identity livekit.ParticipantIdentity) {
|
||||
p.OnTrackPublished(nil)
|
||||
p.OnStateChange(nil)
|
||||
p.OnMetadataUpdate(nil)
|
||||
p.OnDataTrackPublished(nil)
|
||||
p.OnDataPacket(nil)
|
||||
|
||||
// close participant as well
|
||||
_ = p.Close(true)
|
||||
|
||||
@@ -420,9 +420,7 @@ func TestDataChannel(t *testing.T) {
|
||||
},
|
||||
},
|
||||
}
|
||||
dataTrack := &typesfakes.FakeDataTrack{}
|
||||
p.OnDataTrackPublishedArgsForCall(0)(p, dataTrack)
|
||||
dataTrack.OnDataPacketArgsForCall(0)(&packet)
|
||||
p.OnDataPacketArgsForCall(0)(p, &packet)
|
||||
|
||||
// ensure everyone has received the packet
|
||||
for _, op := range participants {
|
||||
@@ -453,9 +451,7 @@ func TestDataChannel(t *testing.T) {
|
||||
},
|
||||
},
|
||||
}
|
||||
dataTrack := &typesfakes.FakeDataTrack{}
|
||||
p.OnDataTrackPublishedArgsForCall(0)(p, dataTrack)
|
||||
dataTrack.OnDataPacketArgsForCall(0)(&packet)
|
||||
p.OnDataPacketArgsForCall(0)(p, &packet)
|
||||
|
||||
// only p1 should receive the data
|
||||
for _, op := range participants {
|
||||
@@ -483,10 +479,8 @@ func TestDataChannel(t *testing.T) {
|
||||
},
|
||||
},
|
||||
}
|
||||
dataTrack := &typesfakes.FakeDataTrack{}
|
||||
p.OnDataTrackPublishedArgsForCall(0)(p, dataTrack)
|
||||
if p.CanPublishData() {
|
||||
dataTrack.OnDataPacketArgsForCall(0)(&packet)
|
||||
p.OnDataPacketArgsForCall(0)(p, &packet)
|
||||
}
|
||||
|
||||
// no one should've been sent packet
|
||||
|
||||
@@ -53,13 +53,12 @@ type Participant interface {
|
||||
ID() livekit.ParticipantID
|
||||
Identity() livekit.ParticipantIdentity
|
||||
|
||||
ToProto(mediaTrackOnly bool) *livekit.ParticipantInfo
|
||||
ToProto() *livekit.ParticipantInfo
|
||||
|
||||
SetMetadata(metadata string)
|
||||
|
||||
GetPublishedTrack(sid livekit.TrackID) MediaTrack
|
||||
GetPublishedTracks() []MediaTrack
|
||||
GetDataTrack() DataTrack
|
||||
|
||||
AddSubscriber(op LocalParticipant, params AddSubscriberParams) (int, error)
|
||||
RemoveSubscriber(op LocalParticipant, trackID livekit.TrackID, resume bool)
|
||||
@@ -146,9 +145,9 @@ type LocalParticipant interface {
|
||||
// OnTrackUpdated - one of its publishedTracks changed in status
|
||||
OnTrackUpdated(callback func(LocalParticipant, MediaTrack))
|
||||
OnMetadataUpdate(callback func(LocalParticipant))
|
||||
OnDataPacket(callback func(LocalParticipant, *livekit.DataPacket))
|
||||
OnClose(_callback func(LocalParticipant, map[livekit.TrackID]livekit.ParticipantID))
|
||||
OnClaimsChanged(_callback func(LocalParticipant))
|
||||
OnDataTrackPublished(callback func(LocalParticipant, DataTrack))
|
||||
|
||||
// session migration
|
||||
SetMigrateState(s MigrateState)
|
||||
@@ -240,14 +239,3 @@ type SubscribedTrack interface {
|
||||
// selects appropriate video layer according to subscriber preferences
|
||||
UpdateVideoLayer()
|
||||
}
|
||||
|
||||
// DataTrack is the interface representing a data track published to the room
|
||||
//counterfeiter:generate . DataTrack
|
||||
type DataTrack interface {
|
||||
ID() livekit.TrackID
|
||||
TrackID() livekit.TrackID
|
||||
Receiver() sfu.TrackReceiver
|
||||
AddOnClose(func())
|
||||
OnDataPacket(callback func(*livekit.DataPacket))
|
||||
Kind() livekit.TrackType
|
||||
}
|
||||
|
||||
@@ -1,377 +0,0 @@
|
||||
// Code generated by counterfeiter. DO NOT EDIT.
|
||||
package typesfakes
|
||||
|
||||
import (
|
||||
"sync"
|
||||
|
||||
"github.com/livekit/livekit-server/pkg/rtc/types"
|
||||
"github.com/livekit/livekit-server/pkg/sfu"
|
||||
"github.com/livekit/protocol/livekit"
|
||||
)
|
||||
|
||||
type FakeDataTrack struct {
|
||||
AddOnCloseStub func(func())
|
||||
addOnCloseMutex sync.RWMutex
|
||||
addOnCloseArgsForCall []struct {
|
||||
arg1 func()
|
||||
}
|
||||
IDStub func() livekit.TrackID
|
||||
iDMutex sync.RWMutex
|
||||
iDArgsForCall []struct {
|
||||
}
|
||||
iDReturns struct {
|
||||
result1 livekit.TrackID
|
||||
}
|
||||
iDReturnsOnCall map[int]struct {
|
||||
result1 livekit.TrackID
|
||||
}
|
||||
KindStub func() livekit.TrackType
|
||||
kindMutex sync.RWMutex
|
||||
kindArgsForCall []struct {
|
||||
}
|
||||
kindReturns struct {
|
||||
result1 livekit.TrackType
|
||||
}
|
||||
kindReturnsOnCall map[int]struct {
|
||||
result1 livekit.TrackType
|
||||
}
|
||||
OnDataPacketStub func(func(*livekit.DataPacket))
|
||||
onDataPacketMutex sync.RWMutex
|
||||
onDataPacketArgsForCall []struct {
|
||||
arg1 func(*livekit.DataPacket)
|
||||
}
|
||||
ReceiverStub func() sfu.TrackReceiver
|
||||
receiverMutex sync.RWMutex
|
||||
receiverArgsForCall []struct {
|
||||
}
|
||||
receiverReturns struct {
|
||||
result1 sfu.TrackReceiver
|
||||
}
|
||||
receiverReturnsOnCall map[int]struct {
|
||||
result1 sfu.TrackReceiver
|
||||
}
|
||||
TrackIDStub func() livekit.TrackID
|
||||
trackIDMutex sync.RWMutex
|
||||
trackIDArgsForCall []struct {
|
||||
}
|
||||
trackIDReturns struct {
|
||||
result1 livekit.TrackID
|
||||
}
|
||||
trackIDReturnsOnCall map[int]struct {
|
||||
result1 livekit.TrackID
|
||||
}
|
||||
invocations map[string][][]interface{}
|
||||
invocationsMutex sync.RWMutex
|
||||
}
|
||||
|
||||
func (fake *FakeDataTrack) AddOnClose(arg1 func()) {
|
||||
fake.addOnCloseMutex.Lock()
|
||||
fake.addOnCloseArgsForCall = append(fake.addOnCloseArgsForCall, struct {
|
||||
arg1 func()
|
||||
}{arg1})
|
||||
stub := fake.AddOnCloseStub
|
||||
fake.recordInvocation("AddOnClose", []interface{}{arg1})
|
||||
fake.addOnCloseMutex.Unlock()
|
||||
if stub != nil {
|
||||
fake.AddOnCloseStub(arg1)
|
||||
}
|
||||
}
|
||||
|
||||
func (fake *FakeDataTrack) AddOnCloseCallCount() int {
|
||||
fake.addOnCloseMutex.RLock()
|
||||
defer fake.addOnCloseMutex.RUnlock()
|
||||
return len(fake.addOnCloseArgsForCall)
|
||||
}
|
||||
|
||||
func (fake *FakeDataTrack) AddOnCloseCalls(stub func(func())) {
|
||||
fake.addOnCloseMutex.Lock()
|
||||
defer fake.addOnCloseMutex.Unlock()
|
||||
fake.AddOnCloseStub = stub
|
||||
}
|
||||
|
||||
func (fake *FakeDataTrack) AddOnCloseArgsForCall(i int) func() {
|
||||
fake.addOnCloseMutex.RLock()
|
||||
defer fake.addOnCloseMutex.RUnlock()
|
||||
argsForCall := fake.addOnCloseArgsForCall[i]
|
||||
return argsForCall.arg1
|
||||
}
|
||||
|
||||
func (fake *FakeDataTrack) ID() livekit.TrackID {
|
||||
fake.iDMutex.Lock()
|
||||
ret, specificReturn := fake.iDReturnsOnCall[len(fake.iDArgsForCall)]
|
||||
fake.iDArgsForCall = append(fake.iDArgsForCall, struct {
|
||||
}{})
|
||||
stub := fake.IDStub
|
||||
fakeReturns := fake.iDReturns
|
||||
fake.recordInvocation("ID", []interface{}{})
|
||||
fake.iDMutex.Unlock()
|
||||
if stub != nil {
|
||||
return stub()
|
||||
}
|
||||
if specificReturn {
|
||||
return ret.result1
|
||||
}
|
||||
return fakeReturns.result1
|
||||
}
|
||||
|
||||
func (fake *FakeDataTrack) IDCallCount() int {
|
||||
fake.iDMutex.RLock()
|
||||
defer fake.iDMutex.RUnlock()
|
||||
return len(fake.iDArgsForCall)
|
||||
}
|
||||
|
||||
func (fake *FakeDataTrack) IDCalls(stub func() livekit.TrackID) {
|
||||
fake.iDMutex.Lock()
|
||||
defer fake.iDMutex.Unlock()
|
||||
fake.IDStub = stub
|
||||
}
|
||||
|
||||
func (fake *FakeDataTrack) IDReturns(result1 livekit.TrackID) {
|
||||
fake.iDMutex.Lock()
|
||||
defer fake.iDMutex.Unlock()
|
||||
fake.IDStub = nil
|
||||
fake.iDReturns = struct {
|
||||
result1 livekit.TrackID
|
||||
}{result1}
|
||||
}
|
||||
|
||||
func (fake *FakeDataTrack) IDReturnsOnCall(i int, result1 livekit.TrackID) {
|
||||
fake.iDMutex.Lock()
|
||||
defer fake.iDMutex.Unlock()
|
||||
fake.IDStub = nil
|
||||
if fake.iDReturnsOnCall == nil {
|
||||
fake.iDReturnsOnCall = make(map[int]struct {
|
||||
result1 livekit.TrackID
|
||||
})
|
||||
}
|
||||
fake.iDReturnsOnCall[i] = struct {
|
||||
result1 livekit.TrackID
|
||||
}{result1}
|
||||
}
|
||||
|
||||
func (fake *FakeDataTrack) Kind() livekit.TrackType {
|
||||
fake.kindMutex.Lock()
|
||||
ret, specificReturn := fake.kindReturnsOnCall[len(fake.kindArgsForCall)]
|
||||
fake.kindArgsForCall = append(fake.kindArgsForCall, struct {
|
||||
}{})
|
||||
stub := fake.KindStub
|
||||
fakeReturns := fake.kindReturns
|
||||
fake.recordInvocation("Kind", []interface{}{})
|
||||
fake.kindMutex.Unlock()
|
||||
if stub != nil {
|
||||
return stub()
|
||||
}
|
||||
if specificReturn {
|
||||
return ret.result1
|
||||
}
|
||||
return fakeReturns.result1
|
||||
}
|
||||
|
||||
func (fake *FakeDataTrack) KindCallCount() int {
|
||||
fake.kindMutex.RLock()
|
||||
defer fake.kindMutex.RUnlock()
|
||||
return len(fake.kindArgsForCall)
|
||||
}
|
||||
|
||||
func (fake *FakeDataTrack) KindCalls(stub func() livekit.TrackType) {
|
||||
fake.kindMutex.Lock()
|
||||
defer fake.kindMutex.Unlock()
|
||||
fake.KindStub = stub
|
||||
}
|
||||
|
||||
func (fake *FakeDataTrack) KindReturns(result1 livekit.TrackType) {
|
||||
fake.kindMutex.Lock()
|
||||
defer fake.kindMutex.Unlock()
|
||||
fake.KindStub = nil
|
||||
fake.kindReturns = struct {
|
||||
result1 livekit.TrackType
|
||||
}{result1}
|
||||
}
|
||||
|
||||
func (fake *FakeDataTrack) KindReturnsOnCall(i int, result1 livekit.TrackType) {
|
||||
fake.kindMutex.Lock()
|
||||
defer fake.kindMutex.Unlock()
|
||||
fake.KindStub = nil
|
||||
if fake.kindReturnsOnCall == nil {
|
||||
fake.kindReturnsOnCall = make(map[int]struct {
|
||||
result1 livekit.TrackType
|
||||
})
|
||||
}
|
||||
fake.kindReturnsOnCall[i] = struct {
|
||||
result1 livekit.TrackType
|
||||
}{result1}
|
||||
}
|
||||
|
||||
func (fake *FakeDataTrack) OnDataPacket(arg1 func(*livekit.DataPacket)) {
|
||||
fake.onDataPacketMutex.Lock()
|
||||
fake.onDataPacketArgsForCall = append(fake.onDataPacketArgsForCall, struct {
|
||||
arg1 func(*livekit.DataPacket)
|
||||
}{arg1})
|
||||
stub := fake.OnDataPacketStub
|
||||
fake.recordInvocation("OnDataPacket", []interface{}{arg1})
|
||||
fake.onDataPacketMutex.Unlock()
|
||||
if stub != nil {
|
||||
fake.OnDataPacketStub(arg1)
|
||||
}
|
||||
}
|
||||
|
||||
func (fake *FakeDataTrack) OnDataPacketCallCount() int {
|
||||
fake.onDataPacketMutex.RLock()
|
||||
defer fake.onDataPacketMutex.RUnlock()
|
||||
return len(fake.onDataPacketArgsForCall)
|
||||
}
|
||||
|
||||
func (fake *FakeDataTrack) OnDataPacketCalls(stub func(func(*livekit.DataPacket))) {
|
||||
fake.onDataPacketMutex.Lock()
|
||||
defer fake.onDataPacketMutex.Unlock()
|
||||
fake.OnDataPacketStub = stub
|
||||
}
|
||||
|
||||
func (fake *FakeDataTrack) OnDataPacketArgsForCall(i int) func(*livekit.DataPacket) {
|
||||
fake.onDataPacketMutex.RLock()
|
||||
defer fake.onDataPacketMutex.RUnlock()
|
||||
argsForCall := fake.onDataPacketArgsForCall[i]
|
||||
return argsForCall.arg1
|
||||
}
|
||||
|
||||
func (fake *FakeDataTrack) Receiver() sfu.TrackReceiver {
|
||||
fake.receiverMutex.Lock()
|
||||
ret, specificReturn := fake.receiverReturnsOnCall[len(fake.receiverArgsForCall)]
|
||||
fake.receiverArgsForCall = append(fake.receiverArgsForCall, struct {
|
||||
}{})
|
||||
stub := fake.ReceiverStub
|
||||
fakeReturns := fake.receiverReturns
|
||||
fake.recordInvocation("Receiver", []interface{}{})
|
||||
fake.receiverMutex.Unlock()
|
||||
if stub != nil {
|
||||
return stub()
|
||||
}
|
||||
if specificReturn {
|
||||
return ret.result1
|
||||
}
|
||||
return fakeReturns.result1
|
||||
}
|
||||
|
||||
func (fake *FakeDataTrack) ReceiverCallCount() int {
|
||||
fake.receiverMutex.RLock()
|
||||
defer fake.receiverMutex.RUnlock()
|
||||
return len(fake.receiverArgsForCall)
|
||||
}
|
||||
|
||||
func (fake *FakeDataTrack) ReceiverCalls(stub func() sfu.TrackReceiver) {
|
||||
fake.receiverMutex.Lock()
|
||||
defer fake.receiverMutex.Unlock()
|
||||
fake.ReceiverStub = stub
|
||||
}
|
||||
|
||||
func (fake *FakeDataTrack) ReceiverReturns(result1 sfu.TrackReceiver) {
|
||||
fake.receiverMutex.Lock()
|
||||
defer fake.receiverMutex.Unlock()
|
||||
fake.ReceiverStub = nil
|
||||
fake.receiverReturns = struct {
|
||||
result1 sfu.TrackReceiver
|
||||
}{result1}
|
||||
}
|
||||
|
||||
func (fake *FakeDataTrack) ReceiverReturnsOnCall(i int, result1 sfu.TrackReceiver) {
|
||||
fake.receiverMutex.Lock()
|
||||
defer fake.receiverMutex.Unlock()
|
||||
fake.ReceiverStub = nil
|
||||
if fake.receiverReturnsOnCall == nil {
|
||||
fake.receiverReturnsOnCall = make(map[int]struct {
|
||||
result1 sfu.TrackReceiver
|
||||
})
|
||||
}
|
||||
fake.receiverReturnsOnCall[i] = struct {
|
||||
result1 sfu.TrackReceiver
|
||||
}{result1}
|
||||
}
|
||||
|
||||
func (fake *FakeDataTrack) TrackID() livekit.TrackID {
|
||||
fake.trackIDMutex.Lock()
|
||||
ret, specificReturn := fake.trackIDReturnsOnCall[len(fake.trackIDArgsForCall)]
|
||||
fake.trackIDArgsForCall = append(fake.trackIDArgsForCall, struct {
|
||||
}{})
|
||||
stub := fake.TrackIDStub
|
||||
fakeReturns := fake.trackIDReturns
|
||||
fake.recordInvocation("TrackID", []interface{}{})
|
||||
fake.trackIDMutex.Unlock()
|
||||
if stub != nil {
|
||||
return stub()
|
||||
}
|
||||
if specificReturn {
|
||||
return ret.result1
|
||||
}
|
||||
return fakeReturns.result1
|
||||
}
|
||||
|
||||
func (fake *FakeDataTrack) TrackIDCallCount() int {
|
||||
fake.trackIDMutex.RLock()
|
||||
defer fake.trackIDMutex.RUnlock()
|
||||
return len(fake.trackIDArgsForCall)
|
||||
}
|
||||
|
||||
func (fake *FakeDataTrack) TrackIDCalls(stub func() livekit.TrackID) {
|
||||
fake.trackIDMutex.Lock()
|
||||
defer fake.trackIDMutex.Unlock()
|
||||
fake.TrackIDStub = stub
|
||||
}
|
||||
|
||||
func (fake *FakeDataTrack) TrackIDReturns(result1 livekit.TrackID) {
|
||||
fake.trackIDMutex.Lock()
|
||||
defer fake.trackIDMutex.Unlock()
|
||||
fake.TrackIDStub = nil
|
||||
fake.trackIDReturns = struct {
|
||||
result1 livekit.TrackID
|
||||
}{result1}
|
||||
}
|
||||
|
||||
func (fake *FakeDataTrack) TrackIDReturnsOnCall(i int, result1 livekit.TrackID) {
|
||||
fake.trackIDMutex.Lock()
|
||||
defer fake.trackIDMutex.Unlock()
|
||||
fake.TrackIDStub = nil
|
||||
if fake.trackIDReturnsOnCall == nil {
|
||||
fake.trackIDReturnsOnCall = make(map[int]struct {
|
||||
result1 livekit.TrackID
|
||||
})
|
||||
}
|
||||
fake.trackIDReturnsOnCall[i] = struct {
|
||||
result1 livekit.TrackID
|
||||
}{result1}
|
||||
}
|
||||
|
||||
func (fake *FakeDataTrack) Invocations() map[string][][]interface{} {
|
||||
fake.invocationsMutex.RLock()
|
||||
defer fake.invocationsMutex.RUnlock()
|
||||
fake.addOnCloseMutex.RLock()
|
||||
defer fake.addOnCloseMutex.RUnlock()
|
||||
fake.iDMutex.RLock()
|
||||
defer fake.iDMutex.RUnlock()
|
||||
fake.kindMutex.RLock()
|
||||
defer fake.kindMutex.RUnlock()
|
||||
fake.onDataPacketMutex.RLock()
|
||||
defer fake.onDataPacketMutex.RUnlock()
|
||||
fake.receiverMutex.RLock()
|
||||
defer fake.receiverMutex.RUnlock()
|
||||
fake.trackIDMutex.RLock()
|
||||
defer fake.trackIDMutex.RUnlock()
|
||||
copiedInvocations := map[string][][]interface{}{}
|
||||
for key, value := range fake.invocations {
|
||||
copiedInvocations[key] = value
|
||||
}
|
||||
return copiedInvocations
|
||||
}
|
||||
|
||||
func (fake *FakeDataTrack) recordInvocation(key string, args []interface{}) {
|
||||
fake.invocationsMutex.Lock()
|
||||
defer fake.invocationsMutex.Unlock()
|
||||
if fake.invocations == nil {
|
||||
fake.invocations = map[string][][]interface{}{}
|
||||
}
|
||||
if fake.invocations[key] == nil {
|
||||
fake.invocations[key] = [][]interface{}{}
|
||||
}
|
||||
fake.invocations[key] = append(fake.invocations[key], args)
|
||||
}
|
||||
|
||||
var _ types.DataTrack = new(FakeDataTrack)
|
||||
@@ -143,16 +143,6 @@ type FakeLocalParticipant struct {
|
||||
getConnectionQualityReturnsOnCall map[int]struct {
|
||||
result1 *livekit.ConnectionQualityInfo
|
||||
}
|
||||
GetDataTrackStub func() types.DataTrack
|
||||
getDataTrackMutex sync.RWMutex
|
||||
getDataTrackArgsForCall []struct {
|
||||
}
|
||||
getDataTrackReturns struct {
|
||||
result1 types.DataTrack
|
||||
}
|
||||
getDataTrackReturnsOnCall map[int]struct {
|
||||
result1 types.DataTrack
|
||||
}
|
||||
GetLoggerStub func() logger.Logger
|
||||
getLoggerMutex sync.RWMutex
|
||||
getLoggerArgsForCall []struct {
|
||||
@@ -322,10 +312,10 @@ type FakeLocalParticipant struct {
|
||||
onCloseArgsForCall []struct {
|
||||
arg1 func(types.LocalParticipant, map[livekit.TrackID]livekit.ParticipantID)
|
||||
}
|
||||
OnDataTrackPublishedStub func(func(types.LocalParticipant, types.DataTrack))
|
||||
onDataTrackPublishedMutex sync.RWMutex
|
||||
onDataTrackPublishedArgsForCall []struct {
|
||||
arg1 func(types.LocalParticipant, types.DataTrack)
|
||||
OnDataPacketStub func(func(types.LocalParticipant, *livekit.DataPacket))
|
||||
onDataPacketMutex sync.RWMutex
|
||||
onDataPacketArgsForCall []struct {
|
||||
arg1 func(types.LocalParticipant, *livekit.DataPacket)
|
||||
}
|
||||
OnMetadataUpdateStub func(func(types.LocalParticipant))
|
||||
onMetadataUpdateMutex sync.RWMutex
|
||||
@@ -548,10 +538,9 @@ type FakeLocalParticipant struct {
|
||||
arg2 livekit.TrackID
|
||||
arg3 bool
|
||||
}
|
||||
ToProtoStub func(bool) *livekit.ParticipantInfo
|
||||
ToProtoStub func() *livekit.ParticipantInfo
|
||||
toProtoMutex sync.RWMutex
|
||||
toProtoArgsForCall []struct {
|
||||
arg1 bool
|
||||
}
|
||||
toProtoReturns struct {
|
||||
result1 *livekit.ParticipantInfo
|
||||
@@ -1308,59 +1297,6 @@ func (fake *FakeLocalParticipant) GetConnectionQualityReturnsOnCall(i int, resul
|
||||
}{result1}
|
||||
}
|
||||
|
||||
func (fake *FakeLocalParticipant) GetDataTrack() types.DataTrack {
|
||||
fake.getDataTrackMutex.Lock()
|
||||
ret, specificReturn := fake.getDataTrackReturnsOnCall[len(fake.getDataTrackArgsForCall)]
|
||||
fake.getDataTrackArgsForCall = append(fake.getDataTrackArgsForCall, struct {
|
||||
}{})
|
||||
stub := fake.GetDataTrackStub
|
||||
fakeReturns := fake.getDataTrackReturns
|
||||
fake.recordInvocation("GetDataTrack", []interface{}{})
|
||||
fake.getDataTrackMutex.Unlock()
|
||||
if stub != nil {
|
||||
return stub()
|
||||
}
|
||||
if specificReturn {
|
||||
return ret.result1
|
||||
}
|
||||
return fakeReturns.result1
|
||||
}
|
||||
|
||||
func (fake *FakeLocalParticipant) GetDataTrackCallCount() int {
|
||||
fake.getDataTrackMutex.RLock()
|
||||
defer fake.getDataTrackMutex.RUnlock()
|
||||
return len(fake.getDataTrackArgsForCall)
|
||||
}
|
||||
|
||||
func (fake *FakeLocalParticipant) GetDataTrackCalls(stub func() types.DataTrack) {
|
||||
fake.getDataTrackMutex.Lock()
|
||||
defer fake.getDataTrackMutex.Unlock()
|
||||
fake.GetDataTrackStub = stub
|
||||
}
|
||||
|
||||
func (fake *FakeLocalParticipant) GetDataTrackReturns(result1 types.DataTrack) {
|
||||
fake.getDataTrackMutex.Lock()
|
||||
defer fake.getDataTrackMutex.Unlock()
|
||||
fake.GetDataTrackStub = nil
|
||||
fake.getDataTrackReturns = struct {
|
||||
result1 types.DataTrack
|
||||
}{result1}
|
||||
}
|
||||
|
||||
func (fake *FakeLocalParticipant) GetDataTrackReturnsOnCall(i int, result1 types.DataTrack) {
|
||||
fake.getDataTrackMutex.Lock()
|
||||
defer fake.getDataTrackMutex.Unlock()
|
||||
fake.GetDataTrackStub = nil
|
||||
if fake.getDataTrackReturnsOnCall == nil {
|
||||
fake.getDataTrackReturnsOnCall = make(map[int]struct {
|
||||
result1 types.DataTrack
|
||||
})
|
||||
}
|
||||
fake.getDataTrackReturnsOnCall[i] = struct {
|
||||
result1 types.DataTrack
|
||||
}{result1}
|
||||
}
|
||||
|
||||
func (fake *FakeLocalParticipant) GetLogger() logger.Logger {
|
||||
fake.getLoggerMutex.Lock()
|
||||
ret, specificReturn := fake.getLoggerReturnsOnCall[len(fake.getLoggerArgsForCall)]
|
||||
@@ -2271,35 +2207,35 @@ func (fake *FakeLocalParticipant) OnCloseArgsForCall(i int) func(types.LocalPart
|
||||
return argsForCall.arg1
|
||||
}
|
||||
|
||||
func (fake *FakeLocalParticipant) OnDataTrackPublished(arg1 func(types.LocalParticipant, types.DataTrack)) {
|
||||
fake.onDataTrackPublishedMutex.Lock()
|
||||
fake.onDataTrackPublishedArgsForCall = append(fake.onDataTrackPublishedArgsForCall, struct {
|
||||
arg1 func(types.LocalParticipant, types.DataTrack)
|
||||
func (fake *FakeLocalParticipant) OnDataPacket(arg1 func(types.LocalParticipant, *livekit.DataPacket)) {
|
||||
fake.onDataPacketMutex.Lock()
|
||||
fake.onDataPacketArgsForCall = append(fake.onDataPacketArgsForCall, struct {
|
||||
arg1 func(types.LocalParticipant, *livekit.DataPacket)
|
||||
}{arg1})
|
||||
stub := fake.OnDataTrackPublishedStub
|
||||
fake.recordInvocation("OnDataTrackPublished", []interface{}{arg1})
|
||||
fake.onDataTrackPublishedMutex.Unlock()
|
||||
stub := fake.OnDataPacketStub
|
||||
fake.recordInvocation("OnDataPacket", []interface{}{arg1})
|
||||
fake.onDataPacketMutex.Unlock()
|
||||
if stub != nil {
|
||||
fake.OnDataTrackPublishedStub(arg1)
|
||||
fake.OnDataPacketStub(arg1)
|
||||
}
|
||||
}
|
||||
|
||||
func (fake *FakeLocalParticipant) OnDataTrackPublishedCallCount() int {
|
||||
fake.onDataTrackPublishedMutex.RLock()
|
||||
defer fake.onDataTrackPublishedMutex.RUnlock()
|
||||
return len(fake.onDataTrackPublishedArgsForCall)
|
||||
func (fake *FakeLocalParticipant) OnDataPacketCallCount() int {
|
||||
fake.onDataPacketMutex.RLock()
|
||||
defer fake.onDataPacketMutex.RUnlock()
|
||||
return len(fake.onDataPacketArgsForCall)
|
||||
}
|
||||
|
||||
func (fake *FakeLocalParticipant) OnDataTrackPublishedCalls(stub func(func(types.LocalParticipant, types.DataTrack))) {
|
||||
fake.onDataTrackPublishedMutex.Lock()
|
||||
defer fake.onDataTrackPublishedMutex.Unlock()
|
||||
fake.OnDataTrackPublishedStub = stub
|
||||
func (fake *FakeLocalParticipant) OnDataPacketCalls(stub func(func(types.LocalParticipant, *livekit.DataPacket))) {
|
||||
fake.onDataPacketMutex.Lock()
|
||||
defer fake.onDataPacketMutex.Unlock()
|
||||
fake.OnDataPacketStub = stub
|
||||
}
|
||||
|
||||
func (fake *FakeLocalParticipant) OnDataTrackPublishedArgsForCall(i int) func(types.LocalParticipant, types.DataTrack) {
|
||||
fake.onDataTrackPublishedMutex.RLock()
|
||||
defer fake.onDataTrackPublishedMutex.RUnlock()
|
||||
argsForCall := fake.onDataTrackPublishedArgsForCall[i]
|
||||
func (fake *FakeLocalParticipant) OnDataPacketArgsForCall(i int) func(types.LocalParticipant, *livekit.DataPacket) {
|
||||
fake.onDataPacketMutex.RLock()
|
||||
defer fake.onDataPacketMutex.RUnlock()
|
||||
argsForCall := fake.onDataPacketArgsForCall[i]
|
||||
return argsForCall.arg1
|
||||
}
|
||||
|
||||
@@ -3560,18 +3496,17 @@ func (fake *FakeLocalParticipant) SubscriptionPermissionUpdateArgsForCall(i int)
|
||||
return argsForCall.arg1, argsForCall.arg2, argsForCall.arg3
|
||||
}
|
||||
|
||||
func (fake *FakeLocalParticipant) ToProto(arg1 bool) *livekit.ParticipantInfo {
|
||||
func (fake *FakeLocalParticipant) ToProto() *livekit.ParticipantInfo {
|
||||
fake.toProtoMutex.Lock()
|
||||
ret, specificReturn := fake.toProtoReturnsOnCall[len(fake.toProtoArgsForCall)]
|
||||
fake.toProtoArgsForCall = append(fake.toProtoArgsForCall, struct {
|
||||
arg1 bool
|
||||
}{arg1})
|
||||
}{})
|
||||
stub := fake.ToProtoStub
|
||||
fakeReturns := fake.toProtoReturns
|
||||
fake.recordInvocation("ToProto", []interface{}{arg1})
|
||||
fake.recordInvocation("ToProto", []interface{}{})
|
||||
fake.toProtoMutex.Unlock()
|
||||
if stub != nil {
|
||||
return stub(arg1)
|
||||
return stub()
|
||||
}
|
||||
if specificReturn {
|
||||
return ret.result1
|
||||
@@ -3585,19 +3520,12 @@ func (fake *FakeLocalParticipant) ToProtoCallCount() int {
|
||||
return len(fake.toProtoArgsForCall)
|
||||
}
|
||||
|
||||
func (fake *FakeLocalParticipant) ToProtoCalls(stub func(bool) *livekit.ParticipantInfo) {
|
||||
func (fake *FakeLocalParticipant) ToProtoCalls(stub func() *livekit.ParticipantInfo) {
|
||||
fake.toProtoMutex.Lock()
|
||||
defer fake.toProtoMutex.Unlock()
|
||||
fake.ToProtoStub = stub
|
||||
}
|
||||
|
||||
func (fake *FakeLocalParticipant) ToProtoArgsForCall(i int) bool {
|
||||
fake.toProtoMutex.RLock()
|
||||
defer fake.toProtoMutex.RUnlock()
|
||||
argsForCall := fake.toProtoArgsForCall[i]
|
||||
return argsForCall.arg1
|
||||
}
|
||||
|
||||
func (fake *FakeLocalParticipant) ToProtoReturns(result1 *livekit.ParticipantInfo) {
|
||||
fake.toProtoMutex.Lock()
|
||||
defer fake.toProtoMutex.Unlock()
|
||||
@@ -3993,8 +3921,6 @@ func (fake *FakeLocalParticipant) Invocations() map[string][][]interface{} {
|
||||
defer fake.getAudioLevelMutex.RUnlock()
|
||||
fake.getConnectionQualityMutex.RLock()
|
||||
defer fake.getConnectionQualityMutex.RUnlock()
|
||||
fake.getDataTrackMutex.RLock()
|
||||
defer fake.getDataTrackMutex.RUnlock()
|
||||
fake.getLoggerMutex.RLock()
|
||||
defer fake.getLoggerMutex.RUnlock()
|
||||
fake.getPublishedTrackMutex.RLock()
|
||||
@@ -4031,8 +3957,8 @@ func (fake *FakeLocalParticipant) Invocations() map[string][][]interface{} {
|
||||
defer fake.onClaimsChangedMutex.RUnlock()
|
||||
fake.onCloseMutex.RLock()
|
||||
defer fake.onCloseMutex.RUnlock()
|
||||
fake.onDataTrackPublishedMutex.RLock()
|
||||
defer fake.onDataTrackPublishedMutex.RUnlock()
|
||||
fake.onDataPacketMutex.RLock()
|
||||
defer fake.onDataPacketMutex.RUnlock()
|
||||
fake.onMetadataUpdateMutex.RLock()
|
||||
defer fake.onMetadataUpdateMutex.RUnlock()
|
||||
fake.onStateChangeMutex.RLock()
|
||||
|
||||
@@ -44,16 +44,6 @@ type FakeParticipant struct {
|
||||
debugInfoReturnsOnCall map[int]struct {
|
||||
result1 map[string]interface{}
|
||||
}
|
||||
GetDataTrackStub func() types.DataTrack
|
||||
getDataTrackMutex sync.RWMutex
|
||||
getDataTrackArgsForCall []struct {
|
||||
}
|
||||
getDataTrackReturns struct {
|
||||
result1 types.DataTrack
|
||||
}
|
||||
getDataTrackReturnsOnCall map[int]struct {
|
||||
result1 types.DataTrack
|
||||
}
|
||||
GetPublishedTrackStub func(livekit.TrackID) types.MediaTrack
|
||||
getPublishedTrackMutex sync.RWMutex
|
||||
getPublishedTrackArgsForCall []struct {
|
||||
@@ -141,10 +131,9 @@ type FakeParticipant struct {
|
||||
subscriptionPermissionReturnsOnCall map[int]struct {
|
||||
result1 *livekit.SubscriptionPermission
|
||||
}
|
||||
ToProtoStub func(bool) *livekit.ParticipantInfo
|
||||
ToProtoStub func() *livekit.ParticipantInfo
|
||||
toProtoMutex sync.RWMutex
|
||||
toProtoArgsForCall []struct {
|
||||
arg1 bool
|
||||
}
|
||||
toProtoReturns struct {
|
||||
result1 *livekit.ParticipantInfo
|
||||
@@ -384,59 +373,6 @@ func (fake *FakeParticipant) DebugInfoReturnsOnCall(i int, result1 map[string]in
|
||||
}{result1}
|
||||
}
|
||||
|
||||
func (fake *FakeParticipant) GetDataTrack() types.DataTrack {
|
||||
fake.getDataTrackMutex.Lock()
|
||||
ret, specificReturn := fake.getDataTrackReturnsOnCall[len(fake.getDataTrackArgsForCall)]
|
||||
fake.getDataTrackArgsForCall = append(fake.getDataTrackArgsForCall, struct {
|
||||
}{})
|
||||
stub := fake.GetDataTrackStub
|
||||
fakeReturns := fake.getDataTrackReturns
|
||||
fake.recordInvocation("GetDataTrack", []interface{}{})
|
||||
fake.getDataTrackMutex.Unlock()
|
||||
if stub != nil {
|
||||
return stub()
|
||||
}
|
||||
if specificReturn {
|
||||
return ret.result1
|
||||
}
|
||||
return fakeReturns.result1
|
||||
}
|
||||
|
||||
func (fake *FakeParticipant) GetDataTrackCallCount() int {
|
||||
fake.getDataTrackMutex.RLock()
|
||||
defer fake.getDataTrackMutex.RUnlock()
|
||||
return len(fake.getDataTrackArgsForCall)
|
||||
}
|
||||
|
||||
func (fake *FakeParticipant) GetDataTrackCalls(stub func() types.DataTrack) {
|
||||
fake.getDataTrackMutex.Lock()
|
||||
defer fake.getDataTrackMutex.Unlock()
|
||||
fake.GetDataTrackStub = stub
|
||||
}
|
||||
|
||||
func (fake *FakeParticipant) GetDataTrackReturns(result1 types.DataTrack) {
|
||||
fake.getDataTrackMutex.Lock()
|
||||
defer fake.getDataTrackMutex.Unlock()
|
||||
fake.GetDataTrackStub = nil
|
||||
fake.getDataTrackReturns = struct {
|
||||
result1 types.DataTrack
|
||||
}{result1}
|
||||
}
|
||||
|
||||
func (fake *FakeParticipant) GetDataTrackReturnsOnCall(i int, result1 types.DataTrack) {
|
||||
fake.getDataTrackMutex.Lock()
|
||||
defer fake.getDataTrackMutex.Unlock()
|
||||
fake.GetDataTrackStub = nil
|
||||
if fake.getDataTrackReturnsOnCall == nil {
|
||||
fake.getDataTrackReturnsOnCall = make(map[int]struct {
|
||||
result1 types.DataTrack
|
||||
})
|
||||
}
|
||||
fake.getDataTrackReturnsOnCall[i] = struct {
|
||||
result1 types.DataTrack
|
||||
}{result1}
|
||||
}
|
||||
|
||||
func (fake *FakeParticipant) GetPublishedTrack(arg1 livekit.TrackID) types.MediaTrack {
|
||||
fake.getPublishedTrackMutex.Lock()
|
||||
ret, specificReturn := fake.getPublishedTrackReturnsOnCall[len(fake.getPublishedTrackArgsForCall)]
|
||||
@@ -906,18 +842,17 @@ func (fake *FakeParticipant) SubscriptionPermissionReturnsOnCall(i int, result1
|
||||
}{result1}
|
||||
}
|
||||
|
||||
func (fake *FakeParticipant) ToProto(arg1 bool) *livekit.ParticipantInfo {
|
||||
func (fake *FakeParticipant) ToProto() *livekit.ParticipantInfo {
|
||||
fake.toProtoMutex.Lock()
|
||||
ret, specificReturn := fake.toProtoReturnsOnCall[len(fake.toProtoArgsForCall)]
|
||||
fake.toProtoArgsForCall = append(fake.toProtoArgsForCall, struct {
|
||||
arg1 bool
|
||||
}{arg1})
|
||||
}{})
|
||||
stub := fake.ToProtoStub
|
||||
fakeReturns := fake.toProtoReturns
|
||||
fake.recordInvocation("ToProto", []interface{}{arg1})
|
||||
fake.recordInvocation("ToProto", []interface{}{})
|
||||
fake.toProtoMutex.Unlock()
|
||||
if stub != nil {
|
||||
return stub(arg1)
|
||||
return stub()
|
||||
}
|
||||
if specificReturn {
|
||||
return ret.result1
|
||||
@@ -931,19 +866,12 @@ func (fake *FakeParticipant) ToProtoCallCount() int {
|
||||
return len(fake.toProtoArgsForCall)
|
||||
}
|
||||
|
||||
func (fake *FakeParticipant) ToProtoCalls(stub func(bool) *livekit.ParticipantInfo) {
|
||||
func (fake *FakeParticipant) ToProtoCalls(stub func() *livekit.ParticipantInfo) {
|
||||
fake.toProtoMutex.Lock()
|
||||
defer fake.toProtoMutex.Unlock()
|
||||
fake.ToProtoStub = stub
|
||||
}
|
||||
|
||||
func (fake *FakeParticipant) ToProtoArgsForCall(i int) bool {
|
||||
fake.toProtoMutex.RLock()
|
||||
defer fake.toProtoMutex.RUnlock()
|
||||
argsForCall := fake.toProtoArgsForCall[i]
|
||||
return argsForCall.arg1
|
||||
}
|
||||
|
||||
func (fake *FakeParticipant) ToProtoReturns(result1 *livekit.ParticipantInfo) {
|
||||
fake.toProtoMutex.Lock()
|
||||
defer fake.toProtoMutex.Unlock()
|
||||
@@ -1225,8 +1153,6 @@ func (fake *FakeParticipant) Invocations() map[string][][]interface{} {
|
||||
defer fake.closeMutex.RUnlock()
|
||||
fake.debugInfoMutex.RLock()
|
||||
defer fake.debugInfoMutex.RUnlock()
|
||||
fake.getDataTrackMutex.RLock()
|
||||
defer fake.getDataTrackMutex.RUnlock()
|
||||
fake.getPublishedTrackMutex.RLock()
|
||||
defer fake.getPublishedTrackMutex.RUnlock()
|
||||
fake.getPublishedTracksMutex.RLock()
|
||||
|
||||
+1
-1
@@ -48,7 +48,7 @@ func UnpackDataTrackLabel(packed string) (peerID livekit.ParticipantID, trackID
|
||||
func ToProtoParticipants(participants []types.LocalParticipant) []*livekit.ParticipantInfo {
|
||||
infos := make([]*livekit.ParticipantInfo, 0, len(participants))
|
||||
for _, op := range participants {
|
||||
infos = append(infos, op.ToProto(true))
|
||||
infos = append(infos, op.ToProto())
|
||||
}
|
||||
return infos
|
||||
}
|
||||
|
||||
@@ -270,7 +270,7 @@ func (r *RoomManager) StartSession(ctx context.Context, roomName livekit.RoomNam
|
||||
_ = participant.Close(true)
|
||||
return
|
||||
}
|
||||
if err = r.roomStore.StoreParticipant(ctx, roomName, participant.ToProto(true)); err != nil {
|
||||
if err = r.roomStore.StoreParticipant(ctx, roomName, participant.ToProto()); err != nil {
|
||||
pLogger.Errorw("could not store participant", err)
|
||||
}
|
||||
|
||||
@@ -287,7 +287,7 @@ func (r *RoomManager) StartSession(ctx context.Context, roomName livekit.RoomNam
|
||||
updateParticipantCount()
|
||||
|
||||
clientMeta := &livekit.AnalyticsClientMeta{Region: r.currentNode.Region, Node: r.currentNode.Id}
|
||||
r.telemetry.ParticipantJoined(ctx, room.Room, participant.ToProto(true), pi.Client, clientMeta)
|
||||
r.telemetry.ParticipantJoined(ctx, room.Room, participant.ToProto(), pi.Client, clientMeta)
|
||||
participant.OnClose(func(p types.LocalParticipant, disallowedSubscriptions map[livekit.TrackID]livekit.ParticipantID) {
|
||||
if err := r.roomStore.DeleteParticipant(ctx, roomName, p.Identity()); err != nil {
|
||||
pLogger.Errorw("could not delete participant", err)
|
||||
@@ -295,7 +295,7 @@ func (r *RoomManager) StartSession(ctx context.Context, roomName livekit.RoomNam
|
||||
|
||||
// update room store with new numParticipants
|
||||
updateParticipantCount()
|
||||
r.telemetry.ParticipantLeft(ctx, room.Room, p.ToProto(true))
|
||||
r.telemetry.ParticipantLeft(ctx, room.Room, p.ToProto())
|
||||
|
||||
room.RemoveDisallowedSubscriptions(p, disallowedSubscriptions)
|
||||
})
|
||||
@@ -359,7 +359,7 @@ func (r *RoomManager) getOrCreateRoom(ctx context.Context, roomName livekit.Room
|
||||
|
||||
newRoom.OnParticipantChanged(func(p types.LocalParticipant) {
|
||||
if p.State() != livekit.ParticipantInfo_DISCONNECTED {
|
||||
if err := r.roomStore.StoreParticipant(ctx, roomName, p.ToProto(true)); err != nil {
|
||||
if err := r.roomStore.StoreParticipant(ctx, roomName, p.ToProto()); err != nil {
|
||||
logger.Errorw("could not handle participant change", err)
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user