mirror of
https://github.com/livekit/livekit.git
synced 2026-08-28 20:18:18 +00:00
write signal messages from media without blocking (#1580)
This commit is contained in:
+55
-13
@@ -2,6 +2,8 @@ package service
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"sync"
|
||||
|
||||
"github.com/pkg/errors"
|
||||
"google.golang.org/protobuf/proto"
|
||||
@@ -106,6 +108,12 @@ func (r *signalService) RelaySignal(stream psrpc.ServerStream[*rpc.RelaySignalRe
|
||||
return errors.Wrap(err, "failed to read participant from session")
|
||||
}
|
||||
|
||||
l := logger.GetLogger().WithValues(
|
||||
"room", ss.RoomName,
|
||||
"participant", ss.Identity,
|
||||
"connectionID", ss.ConnectionId,
|
||||
)
|
||||
|
||||
reqChan := routing.NewDefaultMessageChannel()
|
||||
defer reqChan.Close()
|
||||
|
||||
@@ -115,14 +123,13 @@ func (r *signalService) RelaySignal(stream psrpc.ServerStream[*rpc.RelaySignalRe
|
||||
*pi,
|
||||
livekit.ConnectionID(ss.ConnectionId),
|
||||
reqChan,
|
||||
&relaySignalResponseSink{stream},
|
||||
&relaySignalResponseSink{
|
||||
ServerStream: stream,
|
||||
logger: l,
|
||||
},
|
||||
)
|
||||
if err != nil {
|
||||
logger.Errorw("could not handle new participant", err,
|
||||
"room", ss.RoomName,
|
||||
"participant", ss.Identity,
|
||||
"connectionID", ss.ConnectionId,
|
||||
)
|
||||
l.Errorw("could not handle new participant", err)
|
||||
}
|
||||
|
||||
for msg := range stream.Channel() {
|
||||
@@ -131,16 +138,17 @@ func (r *signalService) RelaySignal(stream psrpc.ServerStream[*rpc.RelaySignalRe
|
||||
}
|
||||
}
|
||||
|
||||
logger.Debugw("participant signal stream closed",
|
||||
"room", ss.RoomName,
|
||||
"participant", ss.Identity,
|
||||
"connectionID", ss.ConnectionId,
|
||||
)
|
||||
l.Debugw("participant signal stream closed")
|
||||
return
|
||||
}
|
||||
|
||||
type relaySignalResponseSink struct {
|
||||
psrpc.ServerStream[*rpc.RelaySignalResponse, *rpc.RelaySignalRequest]
|
||||
logger logger.Logger
|
||||
|
||||
mu sync.Mutex
|
||||
queue []*livekit.SignalResponse
|
||||
writing bool
|
||||
}
|
||||
|
||||
func (s *relaySignalResponseSink) Close() {
|
||||
@@ -151,6 +159,40 @@ func (s *relaySignalResponseSink) IsClosed() bool {
|
||||
return s.Context().Err() != nil
|
||||
}
|
||||
|
||||
func (s *relaySignalResponseSink) WriteMessage(msg proto.Message) error {
|
||||
return s.Send(&rpc.RelaySignalResponse{Response: msg.(*livekit.SignalResponse)})
|
||||
func (s *relaySignalResponseSink) write() {
|
||||
for {
|
||||
s.mu.Lock()
|
||||
var msg *livekit.SignalResponse
|
||||
if len(s.queue) != 0 && !s.IsClosed() {
|
||||
msg = s.queue[0]
|
||||
s.queue = s.queue[1:]
|
||||
} else {
|
||||
s.writing = false
|
||||
s.mu.Unlock()
|
||||
return
|
||||
}
|
||||
s.mu.Unlock()
|
||||
|
||||
if err := s.Send(&rpc.RelaySignalResponse{Response: msg}); err != nil {
|
||||
s.logger.Warnw(
|
||||
"could not send message to participant", err,
|
||||
"messageType", fmt.Sprintf("%T", msg.Message),
|
||||
)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func (s *relaySignalResponseSink) WriteMessage(msg proto.Message) error {
|
||||
if err := s.Context().Err(); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
s.mu.Lock()
|
||||
s.queue = append(s.queue, msg.(*livekit.SignalResponse))
|
||||
if !s.writing {
|
||||
s.writing = true
|
||||
go s.write()
|
||||
}
|
||||
s.mu.Unlock()
|
||||
return nil
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user