mirror of
https://github.com/livekit/livekit.git
synced 2026-08-27 22:34:25 +00:00
Filling out messages unlikely to change in v2. (#3806)
* Filling out messages unlikely to change in v2. * deps * remove defensive nil checks
This commit is contained in:
@@ -23,7 +23,7 @@ require (
|
||||
github.com/jxskiss/base62 v1.1.0
|
||||
github.com/livekit/mageutil v0.0.0-20250511045019-0f1ff63f7731
|
||||
github.com/livekit/mediatransportutil v0.0.0-20250519131108-fb90f5acfded
|
||||
github.com/livekit/protocol v1.39.4-0.20250720165053-c700929d2b5f
|
||||
github.com/livekit/protocol v1.39.4-0.20250721063419-93319bf9e30a
|
||||
github.com/livekit/psrpc v0.6.1-0.20250511053145-465289d72c3c
|
||||
github.com/mackerelio/go-osstat v0.2.5
|
||||
github.com/magefile/mage v1.15.0
|
||||
|
||||
@@ -167,8 +167,8 @@ github.com/livekit/mageutil v0.0.0-20250511045019-0f1ff63f7731 h1:9x+U2HGLrSw5AT
|
||||
github.com/livekit/mageutil v0.0.0-20250511045019-0f1ff63f7731/go.mod h1:Rs3MhFwutWhGwmY1VQsygw28z5bWcnEYmS1OG9OxjOQ=
|
||||
github.com/livekit/mediatransportutil v0.0.0-20250519131108-fb90f5acfded h1:ylZPdnlX1RW9Z15SD4mp87vT2D2shsk0hpLJwSPcq3g=
|
||||
github.com/livekit/mediatransportutil v0.0.0-20250519131108-fb90f5acfded/go.mod h1:mSNtYzSf6iY9xM3UX42VEI+STHvMgHmrYzEHPcdhB8A=
|
||||
github.com/livekit/protocol v1.39.4-0.20250720165053-c700929d2b5f h1:rdTSf5lLtnXbJENVntGKsasl7YNPEz+piA1T225af94=
|
||||
github.com/livekit/protocol v1.39.4-0.20250720165053-c700929d2b5f/go.mod h1:6l+zgRJZ9sY96LM7DA3EMcKQC5zsVyZVP73c+9wgvCA=
|
||||
github.com/livekit/protocol v1.39.4-0.20250721063419-93319bf9e30a h1:gJaHYMRz7ZJCT1qLIfReHEL27zD8L+pRVELMjAatdJI=
|
||||
github.com/livekit/protocol v1.39.4-0.20250721063419-93319bf9e30a/go.mod h1:6l+zgRJZ9sY96LM7DA3EMcKQC5zsVyZVP73c+9wgvCA=
|
||||
github.com/livekit/psrpc v0.6.1-0.20250511053145-465289d72c3c h1:WwEr0YBejYbKzk8LSaO9h8h0G9MnE7shyDu8yXQWmEc=
|
||||
github.com/livekit/psrpc v0.6.1-0.20250511053145-465289d72c3c/go.mod h1:kmD+AZPkWu0MaXIMv57jhNlbiSZZ/Jx4bzlxBDVmJes=
|
||||
github.com/mackerelio/go-osstat v0.2.5 h1:+MqTbZUhoIt4m8qzkVoXUJg1EuifwlAJSk4Yl2GXh+o=
|
||||
|
||||
@@ -58,12 +58,15 @@ func (s *SignalCache) SetLastProcessedRemoteMessageId(lastProcessedRemoteMessage
|
||||
s.lastProcessedRemoteMessageId = lastProcessedRemoteMessageId
|
||||
}
|
||||
|
||||
func (s *SignalCache) Add(msg *livekit.Signalv2ServerMessage) {
|
||||
func (s *SignalCache) Add(msg *livekit.Signalv2ServerMessage) *livekit.Signalv2ServerMessage {
|
||||
if msg != nil {
|
||||
s.AddBatch([]*livekit.Signalv2ServerMessage{msg})
|
||||
}
|
||||
|
||||
return msg
|
||||
}
|
||||
|
||||
// SIGNALLING-V2-TODO: may not need this API
|
||||
func (s *SignalCache) AddBatch(msgs []*livekit.Signalv2ServerMessage) {
|
||||
s.lock.Lock()
|
||||
defer s.lock.Unlock()
|
||||
|
||||
@@ -78,6 +78,13 @@ func (s *signalhandlerv2) HandleRequest(msg proto.Message) error {
|
||||
}
|
||||
|
||||
// SIGNALLING-V2-TODO: process messages
|
||||
switch payload := clientMessage.GetMessage().(type) {
|
||||
case *livekit.Signalv2ClientMessage_PublisherSdp:
|
||||
s.params.Participant.HandleOffer(FromProtoSessionDescription(payload.PublisherSdp))
|
||||
|
||||
case *livekit.Signalv2ClientMessage_SubscriberSdp:
|
||||
s.params.Participant.HandleAnswer(FromProtoSessionDescription(payload.SubscriberSdp))
|
||||
}
|
||||
|
||||
s.params.Signalling.AckMessageId(clientMessage.Sequencer.LastProcessedRemoteMessageId)
|
||||
s.params.Signalling.SetLastProcessedRemoteMessageId(clientMessage.Sequencer.MessageId)
|
||||
|
||||
@@ -40,10 +40,6 @@ func NewSignalling(params SignallingParams) ParticipantSignalling {
|
||||
}
|
||||
|
||||
func (s *signalling) SignalJoinResponse(join *livekit.JoinResponse) proto.Message {
|
||||
if join == nil {
|
||||
return nil
|
||||
}
|
||||
|
||||
return &livekit.SignalResponse{
|
||||
Message: &livekit.SignalResponse_Join{
|
||||
Join: join,
|
||||
@@ -80,10 +76,6 @@ func (s *signalling) SignalSpeakerUpdate(speakers []*livekit.SpeakerInfo) proto.
|
||||
}
|
||||
|
||||
func (s *signalling) SignalRoomUpdate(room *livekit.Room) proto.Message {
|
||||
if room == nil {
|
||||
return nil
|
||||
}
|
||||
|
||||
return &livekit.SignalResponse{
|
||||
Message: &livekit.SignalResponse_RoomUpdate{
|
||||
RoomUpdate: &livekit.RoomUpdate{
|
||||
@@ -94,10 +86,6 @@ func (s *signalling) SignalRoomUpdate(room *livekit.Room) proto.Message {
|
||||
}
|
||||
|
||||
func (s *signalling) SignalConnectionQualityUpdate(connectionQuality *livekit.ConnectionQualityUpdate) proto.Message {
|
||||
if connectionQuality == nil {
|
||||
return nil
|
||||
}
|
||||
|
||||
return &livekit.SignalResponse{
|
||||
Message: &livekit.SignalResponse_ConnectionQuality{
|
||||
ConnectionQuality: connectionQuality,
|
||||
@@ -106,10 +94,6 @@ func (s *signalling) SignalConnectionQualityUpdate(connectionQuality *livekit.Co
|
||||
}
|
||||
|
||||
func (s *signalling) SignalRefreshToken(token string) proto.Message {
|
||||
if token == "" {
|
||||
return nil
|
||||
}
|
||||
|
||||
return &livekit.SignalResponse{
|
||||
Message: &livekit.SignalResponse_RefreshToken{
|
||||
RefreshToken: token,
|
||||
@@ -118,10 +102,6 @@ func (s *signalling) SignalRefreshToken(token string) proto.Message {
|
||||
}
|
||||
|
||||
func (s *signalling) SignalRequestResponse(requestResponse *livekit.RequestResponse) proto.Message {
|
||||
if requestResponse == nil {
|
||||
return nil
|
||||
}
|
||||
|
||||
return &livekit.SignalResponse{
|
||||
Message: &livekit.SignalResponse_RequestResponse{
|
||||
RequestResponse: requestResponse,
|
||||
@@ -130,10 +110,6 @@ func (s *signalling) SignalRequestResponse(requestResponse *livekit.RequestRespo
|
||||
}
|
||||
|
||||
func (s *signalling) SignalRoomMovedResponse(roomMoved *livekit.RoomMovedResponse) proto.Message {
|
||||
if roomMoved == nil {
|
||||
return nil
|
||||
}
|
||||
|
||||
return &livekit.SignalResponse{
|
||||
Message: &livekit.SignalResponse_RoomMoved{
|
||||
RoomMoved: roomMoved,
|
||||
@@ -142,10 +118,6 @@ func (s *signalling) SignalRoomMovedResponse(roomMoved *livekit.RoomMovedRespons
|
||||
}
|
||||
|
||||
func (s *signalling) SignalReconnectResponse(reconnect *livekit.ReconnectResponse) proto.Message {
|
||||
if reconnect == nil {
|
||||
return nil
|
||||
}
|
||||
|
||||
return &livekit.SignalResponse{
|
||||
Message: &livekit.SignalResponse_Reconnect{
|
||||
Reconnect: reconnect,
|
||||
@@ -154,10 +126,6 @@ func (s *signalling) SignalReconnectResponse(reconnect *livekit.ReconnectRespons
|
||||
}
|
||||
|
||||
func (s *signalling) SignalICECandidate(trickle *livekit.TrickleRequest) proto.Message {
|
||||
if trickle == nil {
|
||||
return nil
|
||||
}
|
||||
|
||||
return &livekit.SignalResponse{
|
||||
Message: &livekit.SignalResponse_Trickle{
|
||||
Trickle: trickle,
|
||||
@@ -166,10 +134,6 @@ func (s *signalling) SignalICECandidate(trickle *livekit.TrickleRequest) proto.M
|
||||
}
|
||||
|
||||
func (s *signalling) SignalTrackMuted(mute *livekit.MuteTrackRequest) proto.Message {
|
||||
if mute == nil {
|
||||
return nil
|
||||
}
|
||||
|
||||
return &livekit.SignalResponse{
|
||||
Message: &livekit.SignalResponse_Mute{
|
||||
Mute: mute,
|
||||
@@ -178,10 +142,6 @@ func (s *signalling) SignalTrackMuted(mute *livekit.MuteTrackRequest) proto.Mess
|
||||
}
|
||||
|
||||
func (s *signalling) SignalTrackPublished(trackPublished *livekit.TrackPublishedResponse) proto.Message {
|
||||
if trackPublished == nil {
|
||||
return nil
|
||||
}
|
||||
|
||||
return &livekit.SignalResponse{
|
||||
Message: &livekit.SignalResponse_TrackPublished{
|
||||
TrackPublished: trackPublished,
|
||||
@@ -190,10 +150,6 @@ func (s *signalling) SignalTrackPublished(trackPublished *livekit.TrackPublished
|
||||
}
|
||||
|
||||
func (s *signalling) SignalTrackUnpublished(trackUnpublished *livekit.TrackUnpublishedResponse) proto.Message {
|
||||
if trackUnpublished == nil {
|
||||
return nil
|
||||
}
|
||||
|
||||
return &livekit.SignalResponse{
|
||||
Message: &livekit.SignalResponse_TrackUnpublished{
|
||||
TrackUnpublished: trackUnpublished,
|
||||
@@ -202,10 +158,6 @@ func (s *signalling) SignalTrackUnpublished(trackUnpublished *livekit.TrackUnpub
|
||||
}
|
||||
|
||||
func (s *signalling) SignalTrackSubscribed(trackSubscribed *livekit.TrackSubscribed) proto.Message {
|
||||
if trackSubscribed == nil {
|
||||
return nil
|
||||
}
|
||||
|
||||
return &livekit.SignalResponse{
|
||||
Message: &livekit.SignalResponse_TrackSubscribed{
|
||||
TrackSubscribed: trackSubscribed,
|
||||
@@ -214,10 +166,6 @@ func (s *signalling) SignalTrackSubscribed(trackSubscribed *livekit.TrackSubscri
|
||||
}
|
||||
|
||||
func (s *signalling) SignalLeaveRequest(leave *livekit.LeaveRequest) proto.Message {
|
||||
if leave == nil {
|
||||
return nil
|
||||
}
|
||||
|
||||
return &livekit.SignalResponse{
|
||||
Message: &livekit.SignalResponse_Leave{
|
||||
Leave: leave,
|
||||
@@ -226,10 +174,6 @@ func (s *signalling) SignalLeaveRequest(leave *livekit.LeaveRequest) proto.Messa
|
||||
}
|
||||
|
||||
func (s *signalling) SignalSdpAnswer(answer *livekit.SessionDescription) proto.Message {
|
||||
if answer == nil {
|
||||
return nil
|
||||
}
|
||||
|
||||
return &livekit.SignalResponse{
|
||||
Message: &livekit.SignalResponse_Answer{
|
||||
Answer: answer,
|
||||
@@ -238,10 +182,6 @@ func (s *signalling) SignalSdpAnswer(answer *livekit.SessionDescription) proto.M
|
||||
}
|
||||
|
||||
func (s *signalling) SignalSdpOffer(offer *livekit.SessionDescription) proto.Message {
|
||||
if offer == nil {
|
||||
return nil
|
||||
}
|
||||
|
||||
return &livekit.SignalResponse{
|
||||
Message: &livekit.SignalResponse_Offer{
|
||||
Offer: offer,
|
||||
@@ -250,10 +190,6 @@ func (s *signalling) SignalSdpOffer(offer *livekit.SessionDescription) proto.Mes
|
||||
}
|
||||
|
||||
func (s *signalling) SignalStreamStateUpdate(streamStateUpdate *livekit.StreamStateUpdate) proto.Message {
|
||||
if streamStateUpdate == nil {
|
||||
return nil
|
||||
}
|
||||
|
||||
return &livekit.SignalResponse{
|
||||
Message: &livekit.SignalResponse_StreamStateUpdate{
|
||||
StreamStateUpdate: streamStateUpdate,
|
||||
@@ -262,10 +198,6 @@ func (s *signalling) SignalStreamStateUpdate(streamStateUpdate *livekit.StreamSt
|
||||
}
|
||||
|
||||
func (s *signalling) SignalSubscribedQualityUpdate(subscribedQualityUpdate *livekit.SubscribedQualityUpdate) proto.Message {
|
||||
if subscribedQualityUpdate == nil {
|
||||
return nil
|
||||
}
|
||||
|
||||
return &livekit.SignalResponse{
|
||||
Message: &livekit.SignalResponse_SubscribedQualityUpdate{
|
||||
SubscribedQualityUpdate: subscribedQualityUpdate,
|
||||
@@ -274,10 +206,6 @@ func (s *signalling) SignalSubscribedQualityUpdate(subscribedQualityUpdate *live
|
||||
}
|
||||
|
||||
func (s *signalling) SignalSubscriptionResponse(subscriptionResponse *livekit.SubscriptionResponse) proto.Message {
|
||||
if subscriptionResponse == nil {
|
||||
return nil
|
||||
}
|
||||
|
||||
return &livekit.SignalResponse{
|
||||
Message: &livekit.SignalResponse_SubscriptionResponse{
|
||||
SubscriptionResponse: subscriptionResponse,
|
||||
@@ -286,10 +214,6 @@ func (s *signalling) SignalSubscriptionResponse(subscriptionResponse *livekit.Su
|
||||
}
|
||||
|
||||
func (s *signalling) SignalSubscriptionPermissionUpdate(subscriptionPermissionUpdate *livekit.SubscriptionPermissionUpdate) proto.Message {
|
||||
if subscriptionPermissionUpdate == nil {
|
||||
return nil
|
||||
}
|
||||
|
||||
return &livekit.SignalResponse{
|
||||
Message: &livekit.SignalResponse_SubscriptionPermissionUpdate{
|
||||
SubscriptionPermissionUpdate: subscriptionPermissionUpdate,
|
||||
|
||||
@@ -67,20 +67,68 @@ func (s *signallingv2) PendingMessages() proto.Message {
|
||||
}
|
||||
|
||||
func (s *signallingv2) SignalConnectResponse(connectResponse *livekit.ConnectResponse) proto.Message {
|
||||
if connectResponse == nil {
|
||||
return nil
|
||||
}
|
||||
|
||||
serverMessage := &livekit.Signalv2ServerMessage{
|
||||
Message: &livekit.Signalv2ServerMessage_ConnectResponse{
|
||||
ConnectResponse: connectResponse,
|
||||
},
|
||||
}
|
||||
s.signalCache.Add(serverMessage)
|
||||
return s.cacheAndReturnEnvelope(serverMessage)
|
||||
}
|
||||
|
||||
func (s *signallingv2) SignalSdpOffer(offer *livekit.SessionDescription) proto.Message {
|
||||
serverMessage := &livekit.Signalv2ServerMessage{
|
||||
Message: &livekit.Signalv2ServerMessage_SubscriberSdp{
|
||||
SubscriberSdp: offer,
|
||||
},
|
||||
}
|
||||
return s.cacheAndReturnEnvelope(serverMessage)
|
||||
}
|
||||
|
||||
func (s *signallingv2) SignalSdpAnswer(answer *livekit.SessionDescription) proto.Message {
|
||||
serverMessage := &livekit.Signalv2ServerMessage{
|
||||
Message: &livekit.Signalv2ServerMessage_PublisherSdp{
|
||||
PublisherSdp: answer,
|
||||
},
|
||||
}
|
||||
return s.cacheAndReturnEnvelope(serverMessage)
|
||||
}
|
||||
|
||||
func (s *signallingv2) SignalRoomUpdate(room *livekit.Room) proto.Message {
|
||||
serverMessage := &livekit.Signalv2ServerMessage{
|
||||
Message: &livekit.Signalv2ServerMessage_RoomUpdate{
|
||||
RoomUpdate: &livekit.RoomUpdate{
|
||||
Room: room,
|
||||
},
|
||||
},
|
||||
}
|
||||
return s.cacheAndReturnEnvelope(serverMessage)
|
||||
}
|
||||
|
||||
func (s *signallingv2) SignalParticipantUpdate(participants []*livekit.ParticipantInfo) proto.Message {
|
||||
if len(participants) == 0 {
|
||||
return nil
|
||||
}
|
||||
|
||||
serverMessage := &livekit.Signalv2ServerMessage{
|
||||
Message: &livekit.Signalv2ServerMessage_ParticipantUpdate{
|
||||
ParticipantUpdate: &livekit.ParticipantUpdate{
|
||||
Participants: participants,
|
||||
},
|
||||
},
|
||||
}
|
||||
return s.cacheAndReturnEnvelope(serverMessage)
|
||||
}
|
||||
|
||||
func (s *signallingv2) cacheAndReturnEnvelope(sm *livekit.Signalv2ServerMessage) proto.Message {
|
||||
sm = s.signalCache.Add(sm)
|
||||
if sm == nil {
|
||||
return nil
|
||||
}
|
||||
|
||||
return &livekit.Signalv2WireMessage{
|
||||
Message: &livekit.Signalv2WireMessage_Envelope{
|
||||
Envelope: &livekit.Envelope{
|
||||
ServerMessages: []*livekit.Signalv2ServerMessage{serverMessage},
|
||||
ServerMessages: []*livekit.Signalv2ServerMessage{sm},
|
||||
},
|
||||
},
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user