refactor(sdk): move stream/datagram reconnect ownership into listener

rabbitprincess committed Apr 12, 2026 at 19:48 UTC c0906bd43660c87bb61b63083dd51fcebb20d2b7
4 files changed +245 -196
portal/transport/datagram_client.go
+6 -76
@@ -1,22 +1,13 @@
1 package transport
2
3 import (
4 - "context"
4 "net"
6 - "time"
5
6 "github.com/quic-go/quic-go"
9 - "github.com/rs/zerolog/log"
7
8 "github.com/gosuda/portal-tunnel/v2/types"
12 - "github.com/gosuda/portal-tunnel/v2/utils"
9 )
10
15 -type ClientDatagramState struct {
16 - Identity types.Identity
17 - AccessToken string
18 -}
19 -
11 type ClientDatagram struct {
12 session *datagramSession
13 }
@@ -27,75 +18,14 @@ func NewClientDatagram(onReceiveError func(error)) *ClientDatagram {
18 }
19 }
20
30 -func (d *ClientDatagram) RunLoop(
31 - ctx context.Context,
32 - currentState func() (ClientDatagramState, bool),
33 - open func(context.Context, ClientDatagramState) (*quic.Conn, error),
34 -) {
35 - for {
36 - select {
37 - case <-ctx.Done():
38 - d.session.Stop("listener context closed")
39 - return
40 - default:
41 - }
42 -
43 - state, ok := currentState()
44 - if !ok {
45 - if !utils.SleepOrDone(ctx, time.Second) {
46 - d.session.Stop("listener context closed")
47 - return
48 - }
49 - continue
50 - }
51 -
52 - conn, err := open(ctx, state)
53 - if err != nil {
54 - log.Info().
55 - Err(err).
56 - Str("component", "sdk-datagram-plane").
57 - Str("address", state.Identity.Address).
58 - Msg("quic datagram plane unavailable; retrying")
59 - if !utils.SleepOrDone(ctx, 2*time.Second) {
60 - d.session.Stop("listener context closed")
61 - return
62 - }
63 - continue
64 - }
65 -
66 - log.Info().
67 - Str("component", "sdk-datagram-plane").
68 - Str("address", state.Identity.Address).
69 - Str("remote_addr", conn.RemoteAddr().String()).
70 - Msg("quic tunnel connected")
71 -
72 - recvDone, err := d.session.Bind(conn)
73 - if err != nil {
74 - if ctx.Err() != nil {
75 - return
76 - }
77 - log.Info().
78 - Err(err).
79 - Str("component", "sdk-datagram-plane").
80 - Str("address", state.Identity.Address).
81 - Msg("quic datagram plane did not bind cleanly; retrying")
82 - if !utils.SleepOrDone(ctx, time.Second) {
83 - return
84 - }
85 - continue
86 - }
87 -
88 - select {
89 - case <-ctx.Done():
90 - d.session.Stop("listener context closed")
91 - return
92 - case <-recvDone:
93 - }
94 -
95 - if !utils.SleepOrDone(ctx, time.Second) {
96 - return
21 +func (d *ClientDatagram) Bind(conn *quic.Conn) (<-chan struct{}, error) {
22 + if d == nil || d.session == nil {
23 + if conn != nil {
24 + _ = conn.CloseWithError(0, "listener closed")
25 }
26 + return nil, net.ErrClosed
27 }
28 + return d.session.Bind(conn)
29 }
30
31 func (d *ClientDatagram) Accept(done <-chan struct{}) (types.DatagramFrame, error) {
portal/transport/stream_client.go
+5 -29
@@ -42,39 +42,15 @@ func (s *ClientStream) Accept(done <-chan struct{}) (net.Conn, error) {
42 }
43 }
44
45 -func (s *ClientStream) RunLoop(
45 +func (s *ClientStream) RunSession(
46 ctx context.Context,
47 open func(context.Context) (net.Conn, error),
48 currentTLSConfig func() *tls.Config,
49 - retry func(context.Context, string, error, int) bool,
50 - resetRetries <-chan struct{},
51 -) {
52 - var retries int
53 -
54 - for {
55 - // Drain any pending reset signals (e.g. from system wake) before
56 - // evaluating retry budget. This is non-blocking so it never stalls.
57 - select {
58 - case <-resetRetries:
59 - retries = 0
60 - default:
61 - }
62 -
63 - claimed, err := s.runSession(ctx, open, currentTLSConfig)
64 - switch {
65 - case err == nil:
66 - retries = 0
67 - case errors.Is(err, context.Canceled), errors.Is(err, net.ErrClosed):
68 - return
69 - case claimed:
70 - retries = 0
71 - default:
72 - retries++
73 - if retry == nil || !retry(ctx, "reverse session connect", err, retries) {
74 - return
75 - }
76 - }
49 +) (bool, error) {
50 + if s == nil {
51 + return false, net.ErrClosed
52 }
53 + return s.runSession(ctx, open, currentTLSConfig)
54 }
55
56 func (s *ClientStream) ActiveSessions() int {
sdk/api_client.go
+68 -22
@@ -76,15 +76,25 @@ func newApiClient(relayURL string, cfg ListenerConfig) (*apiClient, error) {
76 }, nil
77 }
78
79 -func (a *apiClient) close() {
80 - if a == nil || a.httpClient == nil {
79 +func closeIdleHTTPClient(httpClient *http.Client) {
80 + if httpClient == nil {
81 return
82 }
83 - if transport, ok := a.httpClient.Transport.(*http.Transport); ok {
83 + if transport, ok := httpClient.Transport.(*http.Transport); ok {
84 transport.CloseIdleConnections()
85 }
86 }
87
88 +func (a *apiClient) close() {
89 + if a == nil {
90 + return
91 + }
92 + a.mu.RLock()
93 + httpClient := a.httpClient
94 + a.mu.RUnlock()
95 + closeIdleHTTPClient(httpClient)
96 +}
97 +
98 // resetTransport tears down the cached HTTP client and TLS config so the next
99 // API call creates fresh TCP connections. Call this after detecting a system
100 // sleep/wake cycle where pooled connections are almost certainly dead.
@@ -92,15 +102,24 @@ func (a *apiClient) resetTransport() {
102 if a == nil {
103 return
104 }
95 - a.close()
105 + a.mu.Lock()
106 + httpClient := a.httpClient
107 a.httpClient = nil
108 a.rawTLSConfig = nil
109 + a.mu.Unlock()
110 + closeIdleHTTPClient(httpClient)
111 }
112
113 func (a *apiClient) registerLease(ctx context.Context, ttl time.Duration, udpEnabled, tcpEnabled bool) (types.RegisterResponse, error) {
114 if err := a.ensureHTTPClient(ctx); err != nil {
115 return types.RegisterResponse{}, err
116 }
117 + a.mu.RLock()
118 + httpClient := a.httpClient
119 + a.mu.RUnlock()
120 + if httpClient == nil {
121 + return types.RegisterResponse{}, errors.New("relay http client is unavailable")
122 + }
123
124 var challenge types.RegisterChallengeResponse
125 challengeReq := types.RegisterChallengeRequest{
@@ -110,7 +129,7 @@ func (a *apiClient) registerLease(ctx context.Context, ttl time.Duration, udpEna
129 UDPEnabled: udpEnabled,
130 TCPEnabled: tcpEnabled,
131 }
113 - if err := utils.HTTPDoAPIPath(ctx, a.httpClient, a.baseURL, http.MethodPost, types.PathSDKRegisterChallenge, challengeReq, nil, &challenge); err != nil {
132 + if err := utils.HTTPDoAPIPath(ctx, httpClient, a.baseURL, http.MethodPost, types.PathSDKRegisterChallenge, challengeReq, nil, &challenge); err != nil {
133 return types.RegisterResponse{}, err
134 }
135
@@ -120,7 +139,7 @@ func (a *apiClient) registerLease(ctx context.Context, ttl time.Duration, udpEna
139 }
140
141 var resp types.RegisterResponse
123 - if err := utils.HTTPDoAPIPath(ctx, a.httpClient, a.baseURL, http.MethodPost, types.PathSDKRegister, types.RegisterRequest{
142 + if err := utils.HTTPDoAPIPath(ctx, httpClient, a.baseURL, http.MethodPost, types.PathSDKRegister, types.RegisterRequest{
143 ChallengeID: challenge.ChallengeID,
144 SIWEMessage: challenge.SIWEMessage,
145 SIWESignature: signature,
@@ -158,9 +177,12 @@ func (a *apiClient) registerLease(ctx context.Context, ttl time.Duration, udpEna
177 }
178
179 func (a *apiClient) ensureHTTPClient(ctx context.Context) error {
180 + a.mu.RLock()
181 if a.httpClient != nil && a.rawTLSConfig != nil {
182 + a.mu.RUnlock()
183 return nil
184 }
185 + a.mu.RUnlock()
186
187 bootstrapCtx, cancel := context.WithTimeout(ctx, defaultDialTimeout+defaultHandshakeTimeout)
188 defer cancel()
@@ -170,16 +192,21 @@ func (a *apiClient) ensureHTTPClient(ctx context.Context) error {
192 return err
193 }
194 if err := a.ensureCompatible(ctx, httpClient); err != nil {
173 - if transport, ok := httpClient.Transport.(*http.Transport); ok {
174 - transport.CloseIdleConnections()
175 - }
176 - a.close()
195 + closeIdleHTTPClient(httpClient)
196 return err
197 }
198
180 - a.close()
199 + a.mu.Lock()
200 + if a.httpClient != nil && a.rawTLSConfig != nil {
201 + a.mu.Unlock()
202 + closeIdleHTTPClient(httpClient)
203 + return nil
204 + }
205 + oldHTTPClient := a.httpClient
206 a.httpClient = httpClient
207 a.rawTLSConfig = rawTLSConfig
208 + a.mu.Unlock()
209 + closeIdleHTTPClient(oldHTTPClient)
210
211 return nil
212 }
@@ -218,14 +245,18 @@ func (a *apiClient) renewLease(ctx context.Context, ttl time.Duration) error {
245 }
246
247 a.mu.RLock()
248 + httpClient := a.httpClient
249 accessToken := a.accessToken
250 a.mu.RUnlock()
251 + if httpClient == nil {
252 + return errors.New("relay http client is unavailable")
253 + }
254 if strings.TrimSpace(accessToken) == "" {
255 return errors.New("access token is not available")
256 }
257
258 var resp types.RenewResponse
228 - if err := utils.HTTPDoAPIPath(ctx, a.httpClient, a.baseURL, http.MethodPost, types.PathSDKRenew, types.RenewRequest{
259 + if err := utils.HTTPDoAPIPath(ctx, httpClient, a.baseURL, http.MethodPost, types.PathSDKRenew, types.RenewRequest{
260 AccessToken: accessToken,
261 TTL: int(ttl / time.Second),
262 ReportedIP: a.reportedIP(ctx),
@@ -247,10 +278,17 @@ func (a *apiClient) renewLease(ctx context.Context, ttl time.Duration) error {
278 }
279
280 func (a *apiClient) unregisterLease(ctx context.Context) error {
281 + if err := a.ensureHTTPClient(ctx); err != nil {
282 + return err
283 + }
284 a.mu.RLock()
285 + httpClient := a.httpClient
286 accessToken := a.accessToken
287 a.mu.RUnlock()
253 - return utils.HTTPDoAPIPath(ctx, a.httpClient, a.baseURL, http.MethodPost, types.PathSDKUnregister, types.UnregisterRequest{
288 + if httpClient == nil {
289 + return errors.New("relay http client is unavailable")
290 + }
291 + return utils.HTTPDoAPIPath(ctx, httpClient, a.baseURL, http.MethodPost, types.PathSDKUnregister, types.UnregisterRequest{
292 AccessToken: accessToken,
293 }, nil, nil)
294 }
@@ -259,10 +297,17 @@ func (a *apiClient) openReverseSession(ctx context.Context) (net.Conn, error) {
297 if err := a.ensureHTTPClient(ctx); err != nil {
298 return nil, err
299 }
300 + a.mu.RLock()
301 + rawTLSConfig := a.rawTLSConfig
302 + accessToken := a.accessToken
303 + a.mu.RUnlock()
304 + if rawTLSConfig == nil {
305 + return nil, errors.New("relay tls config is unavailable")
306 + }
307
308 dialer := &tls.Dialer{
309 NetDialer: &net.Dialer{Timeout: a.dialTimeout},
265 - Config: a.rawTLSConfig.Clone(),
310 + Config: rawTLSConfig.Clone(),
311 }
312
313 conn, err := dialer.DialContext(ctx, "tcp", utils.EnsurePort(a.baseURL.Host))
@@ -276,9 +321,6 @@ func (a *apiClient) openReverseSession(ctx context.Context) (net.Conn, error) {
321 Host: a.baseURL.Host,
322 Header: make(http.Header),
323 }
279 - a.mu.RLock()
280 - accessToken := a.accessToken
281 - a.mu.RUnlock()
324 req.Header.Set(types.HeaderAccessToken, accessToken)
325 req.Header.Set("Connection", "keep-alive")
326
@@ -333,7 +375,15 @@ func (a *apiClient) openQUICSession(ctx context.Context, accessToken string) (*q
375 return nil, err
376 }
377
336 - tlsConf := a.rawTLSConfig.Clone()
378 + a.mu.RLock()
379 + rawTLSConfig := a.rawTLSConfig
380 + sniPort := a.sniPort
381 + a.mu.RUnlock()
382 + if rawTLSConfig == nil {
383 + return nil, errors.New("relay tls config is unavailable")
384 + }
385 +
386 + tlsConf := rawTLSConfig.Clone()
387 tlsConf.NextProtos = []string{"portal-tunnel"}
388
389 quicConf := &quic.Config{
@@ -342,10 +392,6 @@ func (a *apiClient) openQUICSession(ctx context.Context, accessToken string) (*q
392 MaxIdleTimeout: 60 * time.Second,
393 }
394
345 - a.mu.RLock()
346 - sniPort := a.sniPort
347 - a.mu.RUnlock()
348 -
395 if sniPort <= 0 {
396 return nil, errors.New("sni port is not available")
397 }
sdk/listener.go
+166 -69
@@ -11,7 +11,6 @@ import (
11 "sync"
12 "time"
13
14 - "github.com/quic-go/quic-go"
14 "github.com/rs/zerolog/log"
15
16 "github.com/gosuda/portal-tunnel/v2/portal/discovery"
@@ -44,6 +43,7 @@ type Listener struct {
43 cancel context.CancelFunc
44 doneCh <-chan struct{}
45
46 + readyTarget int
47 retryCount int
48 retryWait time.Duration
49 leaseTTL time.Duration
@@ -56,13 +56,7 @@ type Listener struct {
56 registered chan struct{}
57 closeOnce sync.Once
58 registerOnce sync.Once
59 -
60 - // wakeBroadcast is closed and replaced each time the renew loop detects
61 - // a system sleep/wake cycle. Stream RunLoops select on this channel to
62 - // reset their retry counters so they don't exhaust budget on stale-conn
63 - // errors that are really caused by the OS suspending the process.
64 - wakeBroadcast chan struct{}
65 - wakeMu sync.Mutex
59 + streamCancel context.CancelFunc
60
61 banMITM bool
62 tcpEnabled bool
@@ -93,20 +87,20 @@ func NewListener(ctx context.Context, relayURL string, cfg ListenerConfig) (*Lis
87 }
88
89 l := &Listener{
96 - doneCh: listenerCtx.Done(),
97 - cancel: cancel,
98 - api: api,
99 - registered: make(chan struct{}),
100 - wakeBroadcast: make(chan struct{}),
101 - retryCount: cfg.RetryCount,
102 - retryWait: retryWait,
103 - leaseTTL: leaseTTL,
104 - renewBefore: renewBefore,
105 - identity: api.identity.Copy(),
106 - metadata: cfg.Metadata.Copy(),
107 - banMITM: cfg.BanMITM,
108 - tcpEnabled: cfg.TCPEnabled,
109 - relaySet: cfg.relaySet,
90 + doneCh: listenerCtx.Done(),
91 + cancel: cancel,
92 + api: api,
93 + registered: make(chan struct{}),
94 + readyTarget: readyTarget,
95 + retryCount: cfg.RetryCount,
96 + retryWait: retryWait,
97 + leaseTTL: leaseTTL,
98 + renewBefore: renewBefore,
99 + identity: api.identity.Copy(),
100 + metadata: cfg.Metadata.Copy(),
101 + banMITM: cfg.BanMITM,
102 + tcpEnabled: cfg.TCPEnabled,
103 + relaySet: cfg.relaySet,
104 }
105 l.mitmManager = newMITMManager(listenerCtx, l)
106 l.stream = transport.NewClientStream(readyTarget, handshakeTimeout)
@@ -118,36 +112,22 @@ func NewListener(ctx context.Context, relayURL string, cfg ListenerConfig) (*Lis
112 Str("address", l.Address()).
113 Msg("quic datagram plane disconnected; waiting to reconnect")
114 })
121 - go l.datagram.RunLoop(listenerCtx, l.currentDatagramState, func(ctx context.Context, state transport.ClientDatagramState) (*quic.Conn, error) {
122 - return l.api.openQUICSession(ctx, state.AccessToken)
123 - })
115 }
116
126 - go l.runStartup(listenerCtx, readyTarget)
117 + go l.runStartup(listenerCtx)
118 return l, nil
119 }
120
130 -func (l *Listener) runStartup(ctx context.Context, readyTarget int) {
121 +func (l *Listener) runStartup(ctx context.Context) {
122 var retries int
123
124 for {
125 err := l.registerAndConfigure(ctx)
126 switch {
127 case err == nil:
137 - for range readyTarget {
138 - go l.stream.RunLoop(
139 - ctx,
140 - func(ctx context.Context) (net.Conn, error) {
141 - return l.api.openReverseSession(ctx)
142 - },
143 - func() *tls.Config {
144 - l.mu.Lock()
145 - defer l.mu.Unlock()
146 - return l.tlsConfig
147 - },
148 - l.retryOrClose,
149 - l.wakeChannel(),
150 - )
128 + l.startStreamLoops(ctx)
129 + if l.datagram != nil {
130 + go l.runDatagramLoop(ctx)
131 }
132 go l.runRenewLoop(ctx)
133 publicURL := l.PublicURL()
@@ -194,6 +174,7 @@ func (l *Listener) Close() error {
174 identity := l.identity.Copy()
175 registered := l.hostname != ""
176 tlsCloser := l.tlsCloser
177 + streamCancel := l.streamCancel
178 stream := l.stream
179 datagram := l.datagram
180 api := l.api
@@ -201,11 +182,15 @@ func (l *Listener) Close() error {
182 l.udpAddr = ""
183 l.tlsConfig = nil
184 l.tlsCloser = nil
185 + l.streamCancel = nil
186 l.mu.Unlock()
187
188 if l.mitmManager != nil {
189 l.mitmManager.reset()
190 }
191 + if streamCancel != nil {
192 + streamCancel()
193 + }
194
195 if stream != nil {
196 stream.Drain()
@@ -375,43 +360,154 @@ func (l *Listener) PublicURL() string {
360 }).String()
361 }
362
378 -func (l *Listener) currentDatagramState() (transport.ClientDatagramState, bool) {
379 - if l.datagram == nil {
380 - return transport.ClientDatagramState{}, false
363 +func (l *Listener) startStreamLoops(parentCtx context.Context) {
364 + if l.stream == nil || l.readyTarget <= 0 {
365 + return
366 }
367
368 + streamCtx, cancel := context.WithCancel(parentCtx)
369 +
370 l.mu.Lock()
384 - defer l.mu.Unlock()
371 + if l.streamCancel != nil {
372 + l.streamCancel()
373 + }
374 + l.streamCancel = cancel
375 + readyTarget := l.readyTarget
376 + l.mu.Unlock()
377
386 - if l.api == nil || l.identity.Key() == "" || l.udpAddr == "" {
387 - return transport.ClientDatagramState{}, false
378 + for range readyTarget {
379 + go l.runReverseSessionLoop(streamCtx)
380 }
389 - l.api.mu.RLock()
390 - accessToken := l.api.accessToken
391 - l.api.mu.RUnlock()
381 +}
382 +
383 +func (l *Listener) stopStreamLoops() {
384 + l.mu.Lock()
385 + cancel := l.streamCancel
386 + l.streamCancel = nil
387 + l.mu.Unlock()
388
393 - return transport.ClientDatagramState{
394 - Identity: l.identity.Copy(),
395 - AccessToken: accessToken,
396 - }, true
389 + if cancel != nil {
390 + cancel()
391 + }
392 }
393
399 -// notifyWake closes the current wakeBroadcast channel (waking all stream
400 -// RunLoops that select on it) and replaces it with a fresh channel.
401 -func (l *Listener) notifyWake() {
402 - l.wakeMu.Lock()
403 - ch := l.wakeBroadcast
404 - l.wakeBroadcast = make(chan struct{})
405 - l.wakeMu.Unlock()
406 - close(ch)
394 +func (l *Listener) runReverseSessionLoop(ctx context.Context) {
395 + if l.stream == nil {
396 + return
397 + }
398 +
399 + var retries int
400 + for {
401 + claimed, err := l.stream.RunSession(
402 + ctx,
403 + func(ctx context.Context) (net.Conn, error) {
404 + return l.api.openReverseSession(ctx)
405 + },
406 + func() *tls.Config {
407 + l.mu.Lock()
408 + defer l.mu.Unlock()
409 + return l.tlsConfig
410 + },
411 + )
412 + switch {
413 + case err == nil:
414 + retries = 0
415 + case errors.Is(err, context.Canceled), errors.Is(err, net.ErrClosed):
416 + return
417 + case claimed:
418 + retries = 0
419 + default:
420 + retries++
421 + if !l.retryOrClose(ctx, "reverse session connect", err, retries) {
422 + return
423 + }
424 + }
425 + }
426 }
427
409 -// wakeChannel returns the current wake broadcast channel. Stream RunLoops
410 -// select on this; when it is closed they reset their retry counters.
411 -func (l *Listener) wakeChannel() <-chan struct{} {
412 - l.wakeMu.Lock()
413 - defer l.wakeMu.Unlock()
414 - return l.wakeBroadcast
428 +func (l *Listener) runDatagramLoop(ctx context.Context) {
429 + if l.datagram == nil {
430 + return
431 + }
432 +
433 + for {
434 + select {
435 + case <-ctx.Done():
436 + l.datagram.Close()
437 + return
438 + default:
439 + }
440 +
441 + l.mu.Lock()
442 + api := l.api
443 + identity := l.identity.Copy()
444 + udpAddr := l.udpAddr
445 + l.mu.Unlock()
446 + if api == nil || identity.Key() == "" || udpAddr == "" {
447 + if !utils.SleepOrDone(ctx, time.Second) {
448 + l.datagram.Close()
449 + return
450 + }
451 + continue
452 + }
453 + api.mu.RLock()
454 + accessToken := api.accessToken
455 + api.mu.RUnlock()
456 + if strings.TrimSpace(accessToken) == "" {
457 + if !utils.SleepOrDone(ctx, time.Second) {
458 + l.datagram.Close()
459 + return
460 + }
461 + continue
462 + }
463 +
464 + conn, err := api.openQUICSession(ctx, accessToken)
465 + if err != nil {
466 + log.Info().
467 + Err(err).
468 + Str("component", "sdk-datagram-plane").
469 + Str("address", identity.Address).
470 + Msg("quic datagram plane unavailable; retrying")
471 + if !utils.SleepOrDone(ctx, 2*time.Second) {
472 + l.datagram.Close()
473 + return
474 + }
475 + continue
476 + }
477 +
478 + log.Info().
479 + Str("component", "sdk-datagram-plane").
480 + Str("address", identity.Address).
481 + Str("remote_addr", conn.RemoteAddr().String()).
482 + Msg("quic tunnel connected")
483 +
484 + recvDone, err := l.datagram.Bind(conn)
485 + if err != nil {
486 + if ctx.Err() != nil {
487 + return
488 + }
489 + log.Info().
490 + Err(err).
491 + Str("component", "sdk-datagram-plane").
492 + Str("address", identity.Address).
493 + Msg("quic datagram plane did not bind cleanly; retrying")
494 + if !utils.SleepOrDone(ctx, time.Second) {
495 + return
496 + }
497 + continue
498 + }
499 +
500 + select {
501 + case <-ctx.Done():
502 + l.datagram.Close()
503 + return
504 + case <-recvDone:
505 + }
506 +
507 + if !utils.SleepOrDone(ctx, time.Second) {
508 + return
509 + }
510 + }
511 }
512
513 func (l *Listener) runRenewLoop(ctx context.Context) {
@@ -450,12 +546,13 @@ func (l *Listener) runRenewLoop(ctx context.Context) {
546 Str("address", l.Address()).
547 Msg("system sleep/wake detected; resetting transport and re-registering")
548
549 + l.stopStreamLoops()
550 l.api.resetTransport()
454 - l.notifyWake()
551
552 var retries int
553 for {
554 if err := l.registerAndConfigure(ctx); err == nil {
555 + l.startStreamLoops(ctx)
556 break
557 } else if errors.Is(err, context.Canceled) || errors.Is(err, net.ErrClosed) {
558 return