diff --git a/go.mod b/go.mod index 3f6dba843..4678ccc89 100644 --- a/go.mod +++ b/go.mod @@ -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 diff --git a/go.sum b/go.sum index 2754287c4..13f9ef2f8 100644 --- a/go.sum +++ b/go.sum @@ -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= diff --git a/pkg/rtc/signalling/signalcache.go b/pkg/rtc/signalling/signalcache.go index 4d39505c1..b51b41a30 100644 --- a/pkg/rtc/signalling/signalcache.go +++ b/pkg/rtc/signalling/signalcache.go @@ -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() diff --git a/pkg/rtc/signalling/signalhandlerv2.go b/pkg/rtc/signalling/signalhandlerv2.go index 5f5135297..5f2e19783 100644 --- a/pkg/rtc/signalling/signalhandlerv2.go +++ b/pkg/rtc/signalling/signalhandlerv2.go @@ -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) diff --git a/pkg/rtc/signalling/signalling.go b/pkg/rtc/signalling/signalling.go index 0251efa35..81a8e8f22 100644 --- a/pkg/rtc/signalling/signalling.go +++ b/pkg/rtc/signalling/signalling.go @@ -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, diff --git a/pkg/rtc/signalling/signallingv2.go b/pkg/rtc/signalling/signallingv2.go index 8fcdcfc93..ce127094a 100644 --- a/pkg/rtc/signalling/signallingv2.go +++ b/pkg/rtc/signalling/signallingv2.go @@ -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}, }, }, }