From 61ac44e5f7bd6a039d3f5d4f61992c7ab3f3da7d Mon Sep 17 00:00:00 2001 From: cnderrauber Date: Tue, 15 Mar 2022 19:30:10 +0800 Subject: [PATCH] Revert data track change (#513) * Revert data track change * clean code --- pkg/rtc/datatrack.go | 198 --------- pkg/rtc/participant.go | 132 +++--- pkg/rtc/participant_internal_test.go | 6 +- pkg/rtc/room.go | 10 +- pkg/rtc/room_test.go | 12 +- pkg/rtc/types/interfaces.go | 16 +- pkg/rtc/types/typesfakes/fake_data_track.go | 377 ------------------ .../typesfakes/fake_local_participant.go | 138 ++----- pkg/rtc/types/typesfakes/fake_participant.go | 86 +--- pkg/rtc/utils.go | 2 +- pkg/service/roommanager.go | 8 +- 11 files changed, 101 insertions(+), 884 deletions(-) delete mode 100644 pkg/rtc/datatrack.go delete mode 100644 pkg/rtc/types/typesfakes/fake_data_track.go diff --git a/pkg/rtc/datatrack.go b/pkg/rtc/datatrack.go deleted file mode 100644 index 66484bffa..000000000 --- a/pkg/rtc/datatrack.go +++ /dev/null @@ -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{}{} -} diff --git a/pkg/rtc/participant.go b/pkg/rtc/participant.go index 7ce9f05ba..522021a33 100644 --- a/pkg/rtc/participant.go +++ b/pkg/rtc/participant.go @@ -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 { diff --git a/pkg/rtc/participant_internal_test.go b/pkg/rtc/participant_internal_test.go index 9d60e87a7..c92f83833 100644 --- a/pkg/rtc/participant_internal_test.go +++ b/pkg/rtc/participant_internal_test.go @@ -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) } diff --git a/pkg/rtc/room.go b/pkg/rtc/room.go index 7b6d8d4c4..182ff0609 100644 --- a/pkg/rtc/room.go +++ b/pkg/rtc/room.go @@ -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) diff --git a/pkg/rtc/room_test.go b/pkg/rtc/room_test.go index e3f2c782c..babf748e6 100644 --- a/pkg/rtc/room_test.go +++ b/pkg/rtc/room_test.go @@ -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 diff --git a/pkg/rtc/types/interfaces.go b/pkg/rtc/types/interfaces.go index 27ee4cc26..7a635a0e9 100644 --- a/pkg/rtc/types/interfaces.go +++ b/pkg/rtc/types/interfaces.go @@ -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 -} diff --git a/pkg/rtc/types/typesfakes/fake_data_track.go b/pkg/rtc/types/typesfakes/fake_data_track.go deleted file mode 100644 index 624bec193..000000000 --- a/pkg/rtc/types/typesfakes/fake_data_track.go +++ /dev/null @@ -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) diff --git a/pkg/rtc/types/typesfakes/fake_local_participant.go b/pkg/rtc/types/typesfakes/fake_local_participant.go index 1d468cbb9..c8782c0a1 100644 --- a/pkg/rtc/types/typesfakes/fake_local_participant.go +++ b/pkg/rtc/types/typesfakes/fake_local_participant.go @@ -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() diff --git a/pkg/rtc/types/typesfakes/fake_participant.go b/pkg/rtc/types/typesfakes/fake_participant.go index 4cd71f9ae..dbba2c0a6 100644 --- a/pkg/rtc/types/typesfakes/fake_participant.go +++ b/pkg/rtc/types/typesfakes/fake_participant.go @@ -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() diff --git a/pkg/rtc/utils.go b/pkg/rtc/utils.go index 9616dfb04..628e0c076 100644 --- a/pkg/rtc/utils.go +++ b/pkg/rtc/utils.go @@ -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 } diff --git a/pkg/service/roommanager.go b/pkg/service/roommanager.go index 7edb044e1..170fc35fc 100644 --- a/pkg/service/roommanager.go +++ b/pkg/service/roommanager.go @@ -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) } }