From 5564bc531f22d89cf4885f889a80641b901c4453 Mon Sep 17 00:00:00 2001 From: Paul Wells Date: Wed, 5 Apr 2023 03:42:59 -0700 Subject: [PATCH] write signal messages from media without blocking (#1580) --- pkg/service/signal.go | 68 ++++++++++++++++++++++++++++++++++--------- 1 file changed, 55 insertions(+), 13 deletions(-) diff --git a/pkg/service/signal.go b/pkg/service/signal.go index 835788b34..b4838945e 100644 --- a/pkg/service/signal.go +++ b/pkg/service/signal.go @@ -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 }