mirror of
https://github.com/livekit/livekit.git
synced 2026-08-29 07:39:09 +00:00
retry signal stream start (#2410)
This commit is contained in:
+19
-2
@@ -20,6 +20,7 @@ import (
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"github.com/avast/retry-go/v4"
|
||||
"go.uber.org/atomic"
|
||||
"google.golang.org/protobuf/proto"
|
||||
|
||||
@@ -99,13 +100,19 @@ func (r *signalClient) StartParticipantSignal(
|
||||
|
||||
l.Debugw("starting signal connection")
|
||||
|
||||
stream, err := r.client.RelaySignal(ctx, nodeID)
|
||||
var stream psrpc.ClientStream[*rpc.RelaySignalRequest, *rpc.RelaySignalResponse]
|
||||
err = r.retry(ctx, func() (err error) {
|
||||
stream, err = r.client.RelaySignal(ctx, nodeID)
|
||||
return
|
||||
})
|
||||
if err != nil {
|
||||
prometheus.MessageCounter.WithLabelValues("signal", "failure").Add(1)
|
||||
return
|
||||
}
|
||||
|
||||
err = stream.Send(&rpc.RelaySignalRequest{StartSession: ss})
|
||||
err = r.retry(ctx, func() error {
|
||||
return stream.Send(&rpc.RelaySignalRequest{StartSession: ss})
|
||||
})
|
||||
if err != nil {
|
||||
stream.Close(err)
|
||||
prometheus.MessageCounter.WithLabelValues("signal", "failure").Add(1)
|
||||
@@ -141,6 +148,16 @@ func (r *signalClient) StartParticipantSignal(
|
||||
return connectionID, sink, resChan, nil
|
||||
}
|
||||
|
||||
func (r *signalClient) retry(ctx context.Context, fn retry.RetryableFunc) error {
|
||||
return retry.Do(
|
||||
fn,
|
||||
retry.Context(ctx),
|
||||
retry.Delay(r.config.MinRetryInterval),
|
||||
retry.MaxDelay(r.config.MaxRetryInterval),
|
||||
retry.DelayType(retry.BackOffDelay),
|
||||
)
|
||||
}
|
||||
|
||||
type signalRequestMessageWriter struct{}
|
||||
|
||||
func (e signalRequestMessageWriter) Write(seq uint64, close bool, msgs []proto.Message) *rpc.RelaySignalRequest {
|
||||
|
||||
Reference in New Issue
Block a user