fix test
rabbitprincess committed
Mar 20, 2026 at 22:50 UTC
f54402cc0d18a272882dac961a2be8ee4d662840
4 files changed
+24
-28
portal/discovery/discovery.go
+1
-1
@@ -255,7 +255,7 @@ func (s *Service) RunPollLoop(ctx context.Context, interval time.Duration, req t
255
if len(s.Bootstraps()) > 0 {
256
if _, err := s.Poll(ctx, req); err != nil {
257
if ctx.Err() != nil {
258
- return ctx.Err()
258
+ return nil
259
}
260
errText := err.Error()
261
if errText != lastPollErr {
sdk/expose.go
+17
-24
@@ -21,8 +21,8 @@ import (
21
// Exposure owns the lifecycle of one or more relay listeners and accepts
22
// traffic from all of them through one net.Listener.
23
type Exposure struct {
24
- ctx context.Context
24
cancel context.CancelFunc
25
+ done <-chan struct{}
26
27
name string
28
reverseToken string
@@ -36,7 +36,6 @@ type Exposure struct {
36
37
accepted chan net.Conn
38
datagrams chan types.DatagramFrame
39
- done chan struct{}
39
40
mu sync.RWMutex
41
relayURLs []string
@@ -94,8 +93,8 @@ func ExposeWithConfig(ctx context.Context, cfg ExposeConfig) (*Exposure, error)
93
94
exposureCtx, cancel := context.WithCancel(ctx)
95
exposure := &Exposure{
97
- ctx: exposureCtx,
96
cancel: cancel,
97
+ done: exposureCtx.Done(),
98
name: cfg.Name,
99
reverseToken: cfg.ReverseToken,
100
udpEnabled: cfg.UDPEnabled,
@@ -106,7 +105,6 @@ func ExposeWithConfig(ctx context.Context, cfg ExposeConfig) (*Exposure, error)
105
discoveryEnabled: cfg.Discovery,
106
accepted: make(chan net.Conn, max(len(relayURLs)*defaultReadyTarget*2, 1)),
107
datagrams: make(chan types.DatagramFrame, max(len(relayURLs)*32, 1)),
109
- done: make(chan struct{}),
108
listeners: make(map[string]*Listener, len(relayURLs)),
109
starting: make(map[string]struct{}, len(relayURLs)),
110
}
@@ -141,7 +139,11 @@ func ExposeWithConfig(ctx context.Context, cfg ExposeConfig) (*Exposure, error)
139
return nil, err
140
}
141
144
- go exposure.monitorStartupCounts(exposureCtx)
142
+ go exposure.monitorStartupCounts()
143
+ go func() {
144
+ <-exposure.done
145
+ _ = exposure.Close()
146
+ }()
147
if exposure.discovery != nil {
148
go func() {
149
_ = exposure.discovery.RunPollLoop(exposureCtx, 0, types.DiscoverRequest{})
@@ -359,9 +361,6 @@ func (e *Exposure) Close() error {
361
if e.cancel != nil {
362
e.cancel()
363
}
362
- if e.done != nil {
363
- close(e.done)
364
- }
364
365
listeners := e.listenersOrdered()
366
for _, listener := range listeners {
@@ -451,7 +450,7 @@ func (e *Exposure) newListener(relayURL string) (*Listener, error) {
450
RootCAPEM: append([]byte(nil), e.rootCAPEM...),
451
bootstrapService: e.discovery,
452
}
454
- return NewListener(e.ctx, relayURL, cfg)
453
+ return NewListener(context.Background(), relayURL, cfg)
454
}
455
456
func (e *Exposure) installListener(relayURL string, listener *Listener) {
@@ -462,15 +461,12 @@ func (e *Exposure) installListener(relayURL string, listener *Listener) {
461
shouldClose := false
462
e.mu.Lock()
463
delete(e.starting, relayURL)
465
- select {
466
- case <-e.done:
464
+ if e.closed() {
465
shouldClose = true
468
- default:
469
- if _, exists := e.listeners[relayURL]; exists {
470
- shouldClose = true
471
- } else {
472
- e.listeners[relayURL] = listener
473
- }
466
+ } else if _, exists := e.listeners[relayURL]; exists {
467
+ shouldClose = true
468
+ } else {
469
+ e.listeners[relayURL] = listener
470
}
471
e.mu.Unlock()
472
@@ -482,7 +478,7 @@ func (e *Exposure) installListener(relayURL string, listener *Listener) {
478
log.Info().Str("relay_url", relayURL).Msg("relay added to exposure")
479
go e.runListenerAcceptLoop(listener)
480
if e.udpEnabled {
485
- go e.attachDatagramPlane(e.ctx, listener)
481
+ go e.attachDatagramPlane(context.Background(), listener)
482
}
483
}
484
@@ -521,7 +517,7 @@ func (e *Exposure) runListenerAcceptLoop(listener *Listener) {
517
for {
518
conn, err := listener.Accept()
519
if err != nil {
524
- if listener.done() || errors.Is(err, net.ErrClosed) {
520
+ if listener.closed() || errors.Is(err, net.ErrClosed) {
521
return
522
}
523
log.Warn().Err(err).Str("relay_url", listener.relayURL).Msg("exposure listener accept failed")
@@ -708,7 +704,7 @@ func (e *Exposure) allDatagramNegotiationsResolvedWithoutDatagram() bool {
704
705
registered, enabled := listener.datagramNegotiationState()
706
if !registered {
711
- if listener.done() {
707
+ if listener.closed() {
708
resolved++
709
}
710
continue
@@ -776,7 +772,6 @@ func (e *Exposure) closed() bool {
772
if e == nil || e.done == nil {
773
return true
774
}
779
-
775
select {
776
case <-e.done:
777
return true
@@ -785,7 +780,7 @@ func (e *Exposure) closed() bool {
780
}
781
}
782
788
-func (e *Exposure) monitorStartupCounts(ctx context.Context) {
783
+func (e *Exposure) monitorStartupCounts() {
784
if e == nil {
785
return
786
}
@@ -839,8 +834,6 @@ func (e *Exposure) monitorStartupCounts(ctx context.Context) {
834
select {
835
case <-e.done:
836
return
842
- case <-ctx.Done():
843
- return
837
case <-ticker.C:
838
}
839
}
sdk/listener.go
+4
-1
@@ -600,7 +600,10 @@ type listenerAddr string
600
func (a listenerAddr) Network() string { return "portal" }
601
func (a listenerAddr) String() string { return string(a) }
602
603
-func (l *Listener) done() bool {
603
+func (l *Listener) closed() bool {
604
+ if l == nil || l.doneCh == nil {
605
+ return true
606
+ }
607
select {
608
case <-l.doneCh:
609
return true
sdk/sdk_test.go
+2
-2
@@ -291,7 +291,7 @@ func TestNewListenerClosesAfterReverseSessionRetryBudgetExhausted(t *testing.T)
291
defer listener.Close()
292
293
waitForSDKTest(t, func() bool {
294
- return listener.done()
294
+ return listener.closed()
295
})
296
if connectCount.Load() < 2 {
297
t.Fatalf("connect count = %d, want at least 2", connectCount.Load())
@@ -347,7 +347,7 @@ func TestNewListenerRetriesForeverWhenRetryCountIsNegative(t *testing.T) {
347
waitForSDKTest(t, func() bool {
348
return connectCount.Load() >= 3
349
})
350
- if listener.done() {
350
+ if listener.closed() {
351
t.Fatal("listener closed unexpectedly with negative RetryCount")
352
}
353
}