handle new UpdateParticipant API, enable permission updates

This commit is contained in:
David Zhao
2021-03-20 22:27:47 -07:00
parent 6aee50b397
commit 537edda4c3
14 changed files with 387 additions and 423 deletions
+1 -2
View File
@@ -13,7 +13,7 @@ func newMockParticipant(identity string) *typesfakes.FakeParticipant {
p.IdentityReturns(identity)
p.StateReturns(livekit.ParticipantInfo_JOINED)
p.SetMetadataStub = func(m map[string]interface{}) error {
p.SetMetadataStub = func(m string) {
var f func(participant types.Participant)
if p.OnMetadataUpdateCallCount() > 0 {
f = p.OnMetadataUpdateArgsForCall(p.OnMetadataUpdateCallCount() - 1)
@@ -21,7 +21,6 @@ func newMockParticipant(identity string) *typesfakes.FakeParticipant {
if f != nil {
f(p)
}
return nil
}
updateTrack := func() {
var f func(participant types.Participant, track types.PublishedTrack)
+2 -12
View File
@@ -1,7 +1,6 @@
package rtc
import (
"encoding/json"
"fmt"
"io"
"sync"
@@ -139,21 +138,12 @@ func (p *ParticipantImpl) IsReady() bool {
}
// attach metadata to the participant
func (p *ParticipantImpl) SetMetadata(metadata map[string]interface{}) error {
if metadata == nil {
p.metadata = ""
} else {
if data, err := json.Marshal(metadata); err != nil {
return err
} else {
p.metadata = string(data)
}
}
func (p *ParticipantImpl) SetMetadata(metadata string) {
p.metadata = metadata
if p.onMetadataUpdate != nil {
p.onMetadataUpdate(p)
}
return nil
}
func (p *ParticipantImpl) SetPermission(permission *livekit.ParticipantPermission) {
+8 -8
View File
@@ -152,7 +152,7 @@ func TestParticipantUpdate(t *testing.T) {
"track metadata updates are sent to everyone",
true,
func(p types.Participant) {
p.SetMetadata(map[string]interface{}{})
p.SetMetadata("")
},
},
{
@@ -285,7 +285,7 @@ func TestActiveSpeakers(t *testing.T) {
assert.Equal(t, p2.ID(), speakers[1].Sid)
})
t.Run("participants are getting updates when active", func(t *testing.T) {
t.Run("participants are getting audio updates", func(t *testing.T) {
rm := newRoomWithParticipants(t, 2)
participants := rm.GetParticipants()
p := participants[0].(*typesfakes.FakeParticipant)
@@ -293,29 +293,29 @@ func TestActiveSpeakers(t *testing.T) {
p.GetAudioLevelReturns(30, true)
speakers := rm.GetActiveSpeakers()
assert.NotEmpty(t, speakers)
assert.Equal(t, p.ID(), speakers[0].Sid)
require.NotEmpty(t, speakers)
require.Equal(t, p.ID(), speakers[0].Sid)
time.Sleep(audioUpdateDuration)
// everyone should've received updates
for _, op := range participants {
op := op.(*typesfakes.FakeParticipant)
assert.Equal(t, 1, op.SendActiveSpeakersCallCount())
require.Equal(t, 1, op.SendActiveSpeakersCallCount())
}
// after another cycle, we are not getting any new updates since unchanged
time.Sleep(audioUpdateDuration)
for _, op := range participants {
op := op.(*typesfakes.FakeParticipant)
assert.Equal(t, 1, op.SendActiveSpeakersCallCount())
require.Equal(t, 1, op.SendActiveSpeakersCallCount())
}
// no longer speaking, send update with empty items
p.GetAudioLevelReturns(127, false)
time.Sleep(audioUpdateDuration)
assert.Equal(t, 2, p.SendActiveSpeakersCallCount())
assert.Empty(t, p.SendActiveSpeakersArgsForCall(1))
require.Equal(t, 2, p.SendActiveSpeakersCallCount())
require.Empty(t, p.SendActiveSpeakersArgsForCall(1))
})
}
+1 -1
View File
@@ -28,7 +28,7 @@ type Participant interface {
IsReady() bool
ToProto() *livekit.ParticipantInfo
RTCPChan() chan []rtcp.Packet
SetMetadata(metadata map[string]interface{}) error
SetMetadata(metadata string)
SetPermission(permission *livekit.ParticipantPermission)
GetResponseSink() routing.MessageSink
SetResponseSink(sink routing.MessageSink)
+7 -42
View File
@@ -255,16 +255,10 @@ type FakeParticipant struct {
sendParticipantUpdateReturnsOnCall map[int]struct {
result1 error
}
SetMetadataStub func(map[string]interface{}) error
SetMetadataStub func(string)
setMetadataMutex sync.RWMutex
setMetadataArgsForCall []struct {
arg1 map[string]interface{}
}
setMetadataReturns struct {
result1 error
}
setMetadataReturnsOnCall map[int]struct {
result1 error
arg1 string
}
SetPermissionStub func(*livekit.ParticipantPermission)
setPermissionMutex sync.RWMutex
@@ -1661,23 +1655,17 @@ func (fake *FakeParticipant) SendParticipantUpdateReturnsOnCall(i int, result1 e
}{result1}
}
func (fake *FakeParticipant) SetMetadata(arg1 map[string]interface{}) error {
func (fake *FakeParticipant) SetMetadata(arg1 string) {
fake.setMetadataMutex.Lock()
ret, specificReturn := fake.setMetadataReturnsOnCall[len(fake.setMetadataArgsForCall)]
fake.setMetadataArgsForCall = append(fake.setMetadataArgsForCall, struct {
arg1 map[string]interface{}
arg1 string
}{arg1})
stub := fake.SetMetadataStub
fakeReturns := fake.setMetadataReturns
fake.recordInvocation("SetMetadata", []interface{}{arg1})
fake.setMetadataMutex.Unlock()
if stub != nil {
return stub(arg1)
fake.SetMetadataStub(arg1)
}
if specificReturn {
return ret.result1
}
return fakeReturns.result1
}
func (fake *FakeParticipant) SetMetadataCallCount() int {
@@ -1686,42 +1674,19 @@ func (fake *FakeParticipant) SetMetadataCallCount() int {
return len(fake.setMetadataArgsForCall)
}
func (fake *FakeParticipant) SetMetadataCalls(stub func(map[string]interface{}) error) {
func (fake *FakeParticipant) SetMetadataCalls(stub func(string)) {
fake.setMetadataMutex.Lock()
defer fake.setMetadataMutex.Unlock()
fake.SetMetadataStub = stub
}
func (fake *FakeParticipant) SetMetadataArgsForCall(i int) map[string]interface{} {
func (fake *FakeParticipant) SetMetadataArgsForCall(i int) string {
fake.setMetadataMutex.RLock()
defer fake.setMetadataMutex.RUnlock()
argsForCall := fake.setMetadataArgsForCall[i]
return argsForCall.arg1
}
func (fake *FakeParticipant) SetMetadataReturns(result1 error) {
fake.setMetadataMutex.Lock()
defer fake.setMetadataMutex.Unlock()
fake.SetMetadataStub = nil
fake.setMetadataReturns = struct {
result1 error
}{result1}
}
func (fake *FakeParticipant) SetMetadataReturnsOnCall(i int, result1 error) {
fake.setMetadataMutex.Lock()
defer fake.setMetadataMutex.Unlock()
fake.SetMetadataStub = nil
if fake.setMetadataReturnsOnCall == nil {
fake.setMetadataReturnsOnCall = make(map[int]struct {
result1 error
})
}
fake.setMetadataReturnsOnCall[i] = struct {
result1 error
}{result1}
}
func (fake *FakeParticipant) SetPermission(arg1 *livekit.ParticipantPermission) {
fake.setPermissionMutex.Lock()
fake.setPermissionArgsForCall = append(fake.setPermissionArgsForCall, struct {
+8 -14
View File
@@ -1,7 +1,6 @@
package service
import (
"encoding/json"
"fmt"
"sync"
"time"
@@ -216,10 +215,7 @@ func (r *RoomManager) StartSession(roomName string, pi routing.ParticipantInit,
return
}
if pi.Metadata != "" {
var md map[string]interface{}
if err := json.Unmarshal([]byte(pi.Metadata), &md); err == nil {
participant.SetMetadata(md)
}
participant.SetMetadata(pi.Metadata)
}
if pi.Permission != nil {
@@ -383,16 +379,14 @@ func (r *RoomManager) handleRTCMessage(roomName, identity string, msg *livekit.R
logger.Debugw("setting track muted", "room", roomName, "participant", identity,
"track", rm.MuteTrack.TrackSid, "muted", rm.MuteTrack.Muted)
participant.SetTrackMuted(rm.MuteTrack.TrackSid, rm.MuteTrack.Muted)
case *livekit.RTCNodeMessage_UpdateMetadata:
logger.Debugw("updating metadata", "room", roomName, "participant", identity)
var md map[string]interface{}
if rm.UpdateMetadata.Metadata != "" {
if err := json.Unmarshal([]byte(rm.UpdateMetadata.Metadata), &md); err != nil {
logger.Errorw("could not update metadata", "error", err)
return
}
case *livekit.RTCNodeMessage_UpdateParticipant:
logger.Debugw("updating participant", "room", roomName, "participant", identity)
if rm.UpdateParticipant.Metadata != "" {
participant.SetMetadata(rm.UpdateParticipant.Metadata)
}
if rm.UpdateParticipant.Permission != nil {
participant.SetPermission(rm.UpdateParticipant.Permission)
}
participant.SetMetadata(md)
}
}
+5 -5
View File
@@ -150,21 +150,21 @@ func (s *RoomService) MutePublishedTrack(ctx context.Context, req *livekit.MuteR
return
}
func (s *RoomService) UpdateParticipantMetadata(ctx context.Context, req *livekit.UpdateParticipantMetadataRequest) (*livekit.ParticipantInfo, error) {
rtcSink, err := s.createRTCSink(ctx, req.Target.Room, req.Target.Identity)
func (s *RoomService) UpdateParticipant(ctx context.Context, req *livekit.UpdateParticipantRequest) (*livekit.ParticipantInfo, error) {
rtcSink, err := s.createRTCSink(ctx, req.Room, req.Identity)
if err != nil {
return nil, err
}
defer rtcSink.Close()
participant, err := s.roomManager.roomStore.GetParticipant(req.Target.Room, req.Target.Identity)
participant, err := s.roomManager.roomStore.GetParticipant(req.Room, req.Identity)
if err != nil {
return nil, err
}
err = rtcSink.WriteMessage(&livekit.RTCNodeMessage{
Message: &livekit.RTCNodeMessage_UpdateMetadata{
UpdateMetadata: req,
Message: &livekit.RTCNodeMessage_UpdateParticipant{
UpdateParticipant: req,
},
})
+2 -7
View File
@@ -1,7 +1,6 @@
package service
import (
"encoding/json"
"fmt"
"io"
"net/http"
@@ -91,12 +90,8 @@ func (s *RTCService) ServeHTTP(w http.ResponseWriter, r *http.Request) {
return
}
if claims.Metadata != nil {
if data, err := json.Marshal(claims.Metadata); err != nil {
logger.Warnw("unable to encode metadata", "error", err)
} else {
pi.Metadata = string(data)
}
if claims.Metadata != "" {
pi.Metadata = claims.Metadata
}
// this needs to be started first *before* using router functions on this node