diff --git a/pkg/rtc/participant.go b/pkg/rtc/participant.go index 48a26cc20..e1a4aab7c 100644 --- a/pkg/rtc/participant.go +++ b/pkg/rtc/participant.go @@ -30,8 +30,8 @@ type Participant struct { once sync.Once // callbacks & handlers - // OnPeerTrack - remote peer added a mediaTrack - OnPeerTrack func(*Participant, *Track) + // OnParticipantTrack - remote peer added a mediaTrack + OnParticipantTrack func(*Participant, *Track) // OnOffer - offer is ready for remote peer OnOffer func(webrtc.SessionDescription) // OnIceCandidate - ice candidate discovered for local peer @@ -122,6 +122,10 @@ func (p *Participant) ToProto() *livekit.ParticipantInfo { Name: p.name, State: p.state, } + + for _, t := range p.tracks { + info.Tracks = append(info.Tracks, t.ToProto()) + } return info } @@ -302,9 +306,9 @@ func (p *Participant) onTrack(track *webrtc.Track, rtpReceiver *webrtc.RTPReceiv p.lock.Unlock() pt.Start() - if p.OnPeerTrack != nil { + if p.OnParticipantTrack != nil { // caller should hook up what happens when the peer mediaTrack is available - go p.OnPeerTrack(p, pt) + go p.OnParticipantTrack(p, pt) } } diff --git a/pkg/rtc/room.go b/pkg/rtc/room.go index b69572643..5e96d330b 100644 --- a/pkg/rtc/room.go +++ b/pkg/rtc/room.go @@ -63,7 +63,7 @@ func (r *Room) Join(participant *Participant) error { log := logger.GetLogger() // it's important to set this before connection, we don't want to miss out on any tracks - participant.OnPeerTrack = r.onTrackAdded + participant.OnParticipantTrack = r.onTrackAdded participant.OnStateChange = func(p *Participant, oldState livekit.ParticipantInfo_State) { log.Debugw("participant state changed", "state", p.state, "participant", p.id) r.broadcastParticipantState(p) @@ -114,19 +114,23 @@ func (r *Room) RemoveParticipant(id string) { } // a peer in the room added a new mediaTrack, subscribe other participants to it -func (r *Room) onTrackAdded(peer *Participant, track *Track) { +func (r *Room) onTrackAdded(participant *Participant, track *Track) { + // publish participant update, since track state is changed + r.broadcastParticipantState(participant) + r.lock.RLock() defer r.lock.RUnlock() // subscribe all existing participants to this mediaTrack + // this is the default behavior. in the future this could be more selective for _, existingParticipant := range r.participants { - if existingParticipant == peer { + if existingParticipant == participant { // skip publishing peer continue } if err := track.AddSubscriber(existingParticipant); err != nil { logger.GetLogger().Errorw("could not subscribe to mediaTrack", - "srcParticipant", peer.ID(), + "srcParticipant", participant.ID(), "mediaTrack", track.id, "dstParticipant", existingParticipant.ID()) } diff --git a/pkg/rtc/track.go b/pkg/rtc/track.go index cb0c8f234..8c67f6a01 100644 --- a/pkg/rtc/track.go +++ b/pkg/rtc/track.go @@ -9,6 +9,7 @@ import ( "github.com/livekit/livekit-server/pkg/logger" "github.com/livekit/livekit-server/pkg/utils" + "github.com/livekit/livekit-server/proto/livekit" ) var ( @@ -57,6 +58,10 @@ func (t *Track) Start() { go t.forwardWorker() } +func (t *Track) Kind() webrtc.RTPCodecType { + return t.mediaTrack.Kind() +} + // subscribes participant to current mediaTrack // creates and add necessary forwarders and starts them func (t *Track) AddSubscriber(participant *Participant) error { @@ -125,6 +130,22 @@ func (t *Track) RemoveAllSubscribers() { } } +func (t *Track) ToProto() *livekit.TrackInfo { + var kind livekit.TrackInfo_Type + switch t.Kind() { + case webrtc.RTPCodecTypeAudio: + kind = livekit.TrackInfo_AUDIO + case webrtc.RTPCodecTypeVideo: + kind = livekit.TrackInfo_VIDEO + } + + return &livekit.TrackInfo{ + Sid: t.id, + Type: kind, + Name: t.mediaTrack.Label(), + } +} + // forwardWorker reads from the receiver and writes to each sender func (t *Track) forwardWorker() { for pkt := range t.receiver.RTPChan() { diff --git a/proto/livekit/model.pb.go b/proto/livekit/model.pb.go index ad86b15bb..5b3c0e60b 100644 --- a/proto/livekit/model.pb.go +++ b/proto/livekit/model.pb.go @@ -412,9 +412,10 @@ type ParticipantInfo struct { sizeCache protoimpl.SizeCache unknownFields protoimpl.UnknownFields - Sid string `protobuf:"bytes,1,opt,name=sid,proto3" json:"sid,omitempty"` - Name string `protobuf:"bytes,2,opt,name=name,proto3" json:"name,omitempty"` - State ParticipantInfo_State `protobuf:"varint,3,opt,name=state,proto3,enum=livekit.ParticipantInfo_State" json:"state,omitempty"` + Sid string `protobuf:"bytes,1,opt,name=sid,proto3" json:"sid,omitempty"` + Name string `protobuf:"bytes,2,opt,name=name,proto3" json:"name,omitempty"` + State ParticipantInfo_State `protobuf:"varint,3,opt,name=state,proto3,enum=livekit.ParticipantInfo_State" json:"state,omitempty"` + Tracks []*TrackInfo `protobuf:"bytes,4,rep,name=tracks,proto3" json:"tracks,omitempty"` } func (x *ParticipantInfo) Reset() { @@ -470,6 +471,13 @@ func (x *ParticipantInfo) GetState() ParticipantInfo_State { return ParticipantInfo_JOINING } +func (x *ParticipantInfo) GetTracks() []*TrackInfo { + if x != nil { + return x.Tracks + } + return nil +} + // describing type TrackInfo struct { state protoimpl.MessageState @@ -622,35 +630,37 @@ var file_model_proto_rawDesc = []byte{ 0x72, 0x65, 0x61, 0x74, 0x69, 0x6f, 0x6e, 0x5f, 0x74, 0x69, 0x6d, 0x65, 0x18, 0x03, 0x20, 0x01, 0x28, 0x03, 0x52, 0x0c, 0x63, 0x72, 0x65, 0x61, 0x74, 0x69, 0x6f, 0x6e, 0x54, 0x69, 0x6d, 0x65, 0x12, 0x14, 0x0a, 0x05, 0x74, 0x6f, 0x6b, 0x65, 0x6e, 0x18, 0x04, 0x20, 0x01, 0x28, 0x09, 0x52, - 0x05, 0x74, 0x6f, 0x6b, 0x65, 0x6e, 0x22, 0xad, 0x01, 0x0a, 0x0f, 0x50, 0x61, 0x72, 0x74, 0x69, + 0x05, 0x74, 0x6f, 0x6b, 0x65, 0x6e, 0x22, 0xd9, 0x01, 0x0a, 0x0f, 0x50, 0x61, 0x72, 0x74, 0x69, 0x63, 0x69, 0x70, 0x61, 0x6e, 0x74, 0x49, 0x6e, 0x66, 0x6f, 0x12, 0x10, 0x0a, 0x03, 0x73, 0x69, 0x64, 0x18, 0x01, 0x20, 0x01, 0x28, 0x09, 0x52, 0x03, 0x73, 0x69, 0x64, 0x12, 0x12, 0x0a, 0x04, 0x6e, 0x61, 0x6d, 0x65, 0x18, 0x02, 0x20, 0x01, 0x28, 0x09, 0x52, 0x04, 0x6e, 0x61, 0x6d, 0x65, 0x12, 0x34, 0x0a, 0x05, 0x73, 0x74, 0x61, 0x74, 0x65, 0x18, 0x03, 0x20, 0x01, 0x28, 0x0e, 0x32, 0x1e, 0x2e, 0x6c, 0x69, 0x76, 0x65, 0x6b, 0x69, 0x74, 0x2e, 0x50, 0x61, 0x72, 0x74, 0x69, 0x63, 0x69, 0x70, 0x61, 0x6e, 0x74, 0x49, 0x6e, 0x66, 0x6f, 0x2e, 0x53, 0x74, 0x61, 0x74, 0x65, 0x52, - 0x05, 0x73, 0x74, 0x61, 0x74, 0x65, 0x22, 0x3e, 0x0a, 0x05, 0x53, 0x74, 0x61, 0x74, 0x65, 0x12, - 0x0b, 0x0a, 0x07, 0x4a, 0x4f, 0x49, 0x4e, 0x49, 0x4e, 0x47, 0x10, 0x00, 0x12, 0x0a, 0x0a, 0x06, - 0x4a, 0x4f, 0x49, 0x4e, 0x45, 0x44, 0x10, 0x01, 0x12, 0x0a, 0x0a, 0x06, 0x41, 0x43, 0x54, 0x49, - 0x56, 0x45, 0x10, 0x02, 0x12, 0x10, 0x0a, 0x0c, 0x44, 0x49, 0x53, 0x43, 0x4f, 0x4e, 0x4e, 0x45, - 0x43, 0x54, 0x45, 0x44, 0x10, 0x03, 0x22, 0x86, 0x01, 0x0a, 0x09, 0x54, 0x72, 0x61, 0x63, 0x6b, - 0x49, 0x6e, 0x66, 0x6f, 0x12, 0x10, 0x0a, 0x03, 0x73, 0x69, 0x64, 0x18, 0x01, 0x20, 0x01, 0x28, - 0x09, 0x52, 0x03, 0x73, 0x69, 0x64, 0x12, 0x2b, 0x0a, 0x04, 0x74, 0x79, 0x70, 0x65, 0x18, 0x02, - 0x20, 0x01, 0x28, 0x0e, 0x32, 0x17, 0x2e, 0x6c, 0x69, 0x76, 0x65, 0x6b, 0x69, 0x74, 0x2e, 0x54, - 0x72, 0x61, 0x63, 0x6b, 0x49, 0x6e, 0x66, 0x6f, 0x2e, 0x54, 0x79, 0x70, 0x65, 0x52, 0x04, 0x74, - 0x79, 0x70, 0x65, 0x12, 0x12, 0x0a, 0x04, 0x6e, 0x61, 0x6d, 0x65, 0x18, 0x03, 0x20, 0x01, 0x28, - 0x09, 0x52, 0x04, 0x6e, 0x61, 0x6d, 0x65, 0x22, 0x26, 0x0a, 0x04, 0x54, 0x79, 0x70, 0x65, 0x12, - 0x09, 0x0a, 0x05, 0x41, 0x55, 0x44, 0x49, 0x4f, 0x10, 0x00, 0x12, 0x09, 0x0a, 0x05, 0x56, 0x49, - 0x44, 0x45, 0x4f, 0x10, 0x01, 0x12, 0x08, 0x0a, 0x04, 0x44, 0x41, 0x54, 0x41, 0x10, 0x02, 0x22, - 0x46, 0x0a, 0x0b, 0x44, 0x61, 0x74, 0x61, 0x43, 0x68, 0x61, 0x6e, 0x6e, 0x65, 0x6c, 0x12, 0x1d, - 0x0a, 0x0a, 0x73, 0x65, 0x73, 0x73, 0x69, 0x6f, 0x6e, 0x5f, 0x69, 0x64, 0x18, 0x01, 0x20, 0x01, - 0x28, 0x09, 0x52, 0x09, 0x73, 0x65, 0x73, 0x73, 0x69, 0x6f, 0x6e, 0x49, 0x64, 0x12, 0x18, 0x0a, - 0x07, 0x70, 0x61, 0x79, 0x6c, 0x6f, 0x61, 0x64, 0x18, 0x02, 0x20, 0x01, 0x28, 0x0c, 0x52, 0x07, - 0x70, 0x61, 0x79, 0x6c, 0x6f, 0x61, 0x64, 0x42, 0x31, 0x5a, 0x2f, 0x67, 0x69, 0x74, 0x68, 0x75, - 0x62, 0x2e, 0x63, 0x6f, 0x6d, 0x2f, 0x6c, 0x69, 0x76, 0x65, 0x6b, 0x69, 0x74, 0x2f, 0x6c, 0x69, - 0x76, 0x65, 0x6b, 0x69, 0x74, 0x2d, 0x73, 0x65, 0x72, 0x76, 0x65, 0x72, 0x2f, 0x70, 0x72, 0x6f, - 0x74, 0x6f, 0x2f, 0x6c, 0x69, 0x76, 0x65, 0x6b, 0x69, 0x74, 0x62, 0x06, 0x70, 0x72, 0x6f, 0x74, - 0x6f, 0x33, + 0x05, 0x73, 0x74, 0x61, 0x74, 0x65, 0x12, 0x2a, 0x0a, 0x06, 0x74, 0x72, 0x61, 0x63, 0x6b, 0x73, + 0x18, 0x04, 0x20, 0x03, 0x28, 0x0b, 0x32, 0x12, 0x2e, 0x6c, 0x69, 0x76, 0x65, 0x6b, 0x69, 0x74, + 0x2e, 0x54, 0x72, 0x61, 0x63, 0x6b, 0x49, 0x6e, 0x66, 0x6f, 0x52, 0x06, 0x74, 0x72, 0x61, 0x63, + 0x6b, 0x73, 0x22, 0x3e, 0x0a, 0x05, 0x53, 0x74, 0x61, 0x74, 0x65, 0x12, 0x0b, 0x0a, 0x07, 0x4a, + 0x4f, 0x49, 0x4e, 0x49, 0x4e, 0x47, 0x10, 0x00, 0x12, 0x0a, 0x0a, 0x06, 0x4a, 0x4f, 0x49, 0x4e, + 0x45, 0x44, 0x10, 0x01, 0x12, 0x0a, 0x0a, 0x06, 0x41, 0x43, 0x54, 0x49, 0x56, 0x45, 0x10, 0x02, + 0x12, 0x10, 0x0a, 0x0c, 0x44, 0x49, 0x53, 0x43, 0x4f, 0x4e, 0x4e, 0x45, 0x43, 0x54, 0x45, 0x44, + 0x10, 0x03, 0x22, 0x86, 0x01, 0x0a, 0x09, 0x54, 0x72, 0x61, 0x63, 0x6b, 0x49, 0x6e, 0x66, 0x6f, + 0x12, 0x10, 0x0a, 0x03, 0x73, 0x69, 0x64, 0x18, 0x01, 0x20, 0x01, 0x28, 0x09, 0x52, 0x03, 0x73, + 0x69, 0x64, 0x12, 0x2b, 0x0a, 0x04, 0x74, 0x79, 0x70, 0x65, 0x18, 0x02, 0x20, 0x01, 0x28, 0x0e, + 0x32, 0x17, 0x2e, 0x6c, 0x69, 0x76, 0x65, 0x6b, 0x69, 0x74, 0x2e, 0x54, 0x72, 0x61, 0x63, 0x6b, + 0x49, 0x6e, 0x66, 0x6f, 0x2e, 0x54, 0x79, 0x70, 0x65, 0x52, 0x04, 0x74, 0x79, 0x70, 0x65, 0x12, + 0x12, 0x0a, 0x04, 0x6e, 0x61, 0x6d, 0x65, 0x18, 0x03, 0x20, 0x01, 0x28, 0x09, 0x52, 0x04, 0x6e, + 0x61, 0x6d, 0x65, 0x22, 0x26, 0x0a, 0x04, 0x54, 0x79, 0x70, 0x65, 0x12, 0x09, 0x0a, 0x05, 0x41, + 0x55, 0x44, 0x49, 0x4f, 0x10, 0x00, 0x12, 0x09, 0x0a, 0x05, 0x56, 0x49, 0x44, 0x45, 0x4f, 0x10, + 0x01, 0x12, 0x08, 0x0a, 0x04, 0x44, 0x41, 0x54, 0x41, 0x10, 0x02, 0x22, 0x46, 0x0a, 0x0b, 0x44, + 0x61, 0x74, 0x61, 0x43, 0x68, 0x61, 0x6e, 0x6e, 0x65, 0x6c, 0x12, 0x1d, 0x0a, 0x0a, 0x73, 0x65, + 0x73, 0x73, 0x69, 0x6f, 0x6e, 0x5f, 0x69, 0x64, 0x18, 0x01, 0x20, 0x01, 0x28, 0x09, 0x52, 0x09, + 0x73, 0x65, 0x73, 0x73, 0x69, 0x6f, 0x6e, 0x49, 0x64, 0x12, 0x18, 0x0a, 0x07, 0x70, 0x61, 0x79, + 0x6c, 0x6f, 0x61, 0x64, 0x18, 0x02, 0x20, 0x01, 0x28, 0x0c, 0x52, 0x07, 0x70, 0x61, 0x79, 0x6c, + 0x6f, 0x61, 0x64, 0x42, 0x31, 0x5a, 0x2f, 0x67, 0x69, 0x74, 0x68, 0x75, 0x62, 0x2e, 0x63, 0x6f, + 0x6d, 0x2f, 0x6c, 0x69, 0x76, 0x65, 0x6b, 0x69, 0x74, 0x2f, 0x6c, 0x69, 0x76, 0x65, 0x6b, 0x69, + 0x74, 0x2d, 0x73, 0x65, 0x72, 0x76, 0x65, 0x72, 0x2f, 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x2f, 0x6c, + 0x69, 0x76, 0x65, 0x6b, 0x69, 0x74, 0x62, 0x06, 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x33, } var ( @@ -681,12 +691,13 @@ var file_model_proto_goTypes = []interface{}{ var file_model_proto_depIdxs = []int32{ 3, // 0: livekit.Node.stats:type_name -> livekit.NodeStats 0, // 1: livekit.ParticipantInfo.state:type_name -> livekit.ParticipantInfo.State - 1, // 2: livekit.TrackInfo.type:type_name -> livekit.TrackInfo.Type - 3, // [3:3] is the sub-list for method output_type - 3, // [3:3] is the sub-list for method input_type - 3, // [3:3] is the sub-list for extension type_name - 3, // [3:3] is the sub-list for extension extendee - 0, // [0:3] is the sub-list for field type_name + 7, // 2: livekit.ParticipantInfo.tracks:type_name -> livekit.TrackInfo + 1, // 3: livekit.TrackInfo.type:type_name -> livekit.TrackInfo.Type + 4, // [4:4] is the sub-list for method output_type + 4, // [4:4] is the sub-list for method input_type + 4, // [4:4] is the sub-list for extension type_name + 4, // [4:4] is the sub-list for extension extendee + 0, // [0:4] is the sub-list for field type_name } func init() { file_model_proto_init() } diff --git a/proto/model.proto b/proto/model.proto index a15e33f45..235cb5d16 100644 --- a/proto/model.proto +++ b/proto/model.proto @@ -45,6 +45,7 @@ message ParticipantInfo { string sid = 1; string name = 2; State state = 3; + repeated TrackInfo tracks = 4; } // describing