sdk: retry all relays

Kim committed Mar 11, 2026 at 17:43 UTC 46e13a39a35086ae171dd50ec77b39020ec4e637
7 files changed +303 -166
cmd/portal-tunnel/README.md
+4 -1
@@ -32,6 +32,9 @@ Portal-tunnel connects a local service to a Portal relay with the legacy CLI sha
32 - Multiple relay URLs are registered independently. Each relay gets its own lease ID and public URLs.
33 - Relay publishes each service at `<name>.<portal root host>`.
34 - Portal-tunnel now consumes one aggregate SDK listener, so the CLI no longer manages per-relay listener loops itself.
35 -- Startup is fail-fast: if any configured relay cannot register, the tunnel exits instead of partially publishing.
35 +- Relay startup and reconnect failures are retried independently in the background. A relay that is down does not stop healthy relays from continuing to serve traffic.
36 +- The tunnel starts once relay URLs pass local validation. Remote compatibility checks, lease registration, and reconnects continue in the background until each relay becomes ready.
37 +- The configured relay list stays fixed, but published public URLs appear only for relays that have registered successfully.
38 +- SDK callers that do not set `ListenerConfig.RetryCount` use infinite retry semantics for each relay.
39 - Tenant TLS is provisioned automatically through the relay keyless signer. The SDK fetches the relay certificate chain and uses `/v1/sign` for remote signing.
40 - When the local service is unreachable, the tunnel returns an HTTP 503 page.
docs/architecture.md
+5 -2
@@ -54,10 +54,13 @@ That distinction matters because `/sdk/connect` stops being ordinary HTTP once h
54
55 ### SDK (`sdk/`)
56
57 -- `Listener`: validates one relay URL, registers one lease per relay, maintains per-entry `readyTarget` reverse sessions, renews lease TTLs, and yields accepted tenant TLS connections
57 +- `Listener`: validates one relay URL locally, then starts relay compatibility checks, lease registration, reverse session maintenance, and lease renewal in the background until ready
58 - `relayclient.go`: internal relay transport helper for control-plane requests and reverse session dialing
59 -- Default app flow is `RelayURL -> NewListener -> PublicURLs -> http.Server.Serve(listener)` or `RelayURLs -> Expose -> PublicURLs -> http.Server.Serve(exposure)`
59 +- `ListenerConfig.RetryCount <= 0` means retry forever; positive values close the listener after the retry budget is exhausted
60 +- Default app flow is `RelayURL -> NewListener -> PublicURL -> http.Server.Serve(listener)` or `RelayURLs -> Expose -> PublicURLs -> http.Server.Serve(exposure)`
61 - `expose.go`: optional `RunHTTP` helper for serving one handler on both a local HTTP port and the relay listener
62 +- `Expose` keeps one listener per configured relay URL. Relay startup and reconnect failures are retried independently per relay, and successful relays remain available while failed relays keep retrying in the background
63 +- `Exposure.RelayURLs()` returns the configured normalized relay URLs, while `Exposure.PublicURLs()` returns only relays that are currently registered and ready
64 - Relay-aware entry inspection is reserved for advanced callers such as `portal-tunnel`
65 - Tenant TLS is created automatically through the relay keyless signer; callers do not provide a local self-signed fallback path
66
sdk/api_client.go
+79 -56
@@ -6,6 +6,7 @@ import (
6 "context"
7 "crypto/tls"
8 "encoding/json"
9 + "errors"
10 "fmt"
11 "io"
12 "net"
@@ -31,16 +32,18 @@ const (
32 )
33
34 type apiClient struct {
34 - baseURL *url.URL
35 - httpClient *http.Client
36 - rawTLSConfig *tls.Config
37 - dialTimeout time.Duration
38 - name string
39 - reverseToken string
40 - metadata types.LeaseMetadata
35 + baseURL *url.URL
36 + httpClient *http.Client
37 + rawTLSConfig *tls.Config
38 + dialTimeout time.Duration
39 + requestTimeout time.Duration
40 + rootCAPEM []byte
41 + name string
42 + reverseToken string
43 + metadata types.LeaseMetadata
44 }
45
43 -func newApiClient(ctx context.Context, relayURL string, cfg ListenerConfig) (*apiClient, error) {
46 +func newApiClient(relayURL string, cfg ListenerConfig) (*apiClient, error) {
47 name, err := utils.NormalizeDNSLabel(cfg.Name)
48 if err != nil {
49 return nil, err
@@ -61,55 +64,17 @@ func newApiClient(ctx context.Context, relayURL string, cfg ListenerConfig) (*ap
64 return nil, fmt.Errorf("parse relay url: %w", err)
65 }
66
64 - rootCAPEM := append([]byte(nil), cfg.RootCAPEM...)
65 - if len(rootCAPEM) == 0 && utils.IsLocalRelayHost(baseURL.Hostname()) {
66 - bootstrapParent := ctx
67 - if bootstrapParent == nil {
68 - bootstrapParent = context.Background()
69 - }
70 - bootstrapCtx, cancel := context.WithTimeout(bootstrapParent, defaultDialTimeout+defaultHandshakeTimeout)
71 - defer cancel()
72 -
73 - _, resolvedCAPEM, bootstrapErr := keyless.ResolveMaterials(bootstrapCtx, baseURL.String(), baseURL.Hostname())
74 - if bootstrapErr != nil {
75 - return nil, fmt.Errorf("bootstrap localhost relay trust: %w", bootstrapErr)
76 - }
77 - rootCAPEM = resolvedCAPEM
78 - }
79 -
80 - rootCAs, err := utils.CertPoolFromPEM(rootCAPEM)
81 - if err != nil {
82 - return nil, err
83 - }
84 -
67 dialTimeout := utils.DurationOrDefault(cfg.DialTimeout, defaultDialTimeout)
68 requestTimeout := utils.DurationOrDefault(cfg.RequestTimeout, defaultRequestTimeout)
69
88 - baseTLS := &tls.Config{
89 - MinVersion: tls.VersionTLS12,
90 - ServerName: baseURL.Hostname(),
91 - RootCAs: rootCAs,
92 - NextProtos: []string{"http/1.1"},
93 - }
94 -
95 - transport := &http.Transport{
96 - TLSClientConfig: baseTLS.Clone(),
97 - ForceAttemptHTTP2: false,
98 - }
99 -
70 api := &apiClient{
101 - baseURL: baseURL,
102 - httpClient: &http.Client{Transport: transport, Timeout: requestTimeout},
103 - rawTLSConfig: baseTLS,
104 - dialTimeout: dialTimeout,
105 - name: name,
106 - reverseToken: reverseToken,
107 - metadata: cfg.Metadata.Copy(),
108 - }
109 -
110 - if err := api.ensureCompatible(ctx); err != nil {
111 - api.close()
112 - return nil, err
71 + baseURL: baseURL,
72 + dialTimeout: dialTimeout,
73 + requestTimeout: requestTimeout,
74 + rootCAPEM: append([]byte(nil), cfg.RootCAPEM...),
75 + name: name,
76 + reverseToken: reverseToken,
77 + metadata: cfg.Metadata.Copy(),
78 }
79 return api, nil
80 }
@@ -136,9 +101,59 @@ func (a *apiClient) registerLease(ctx context.Context, ttl time.Duration) (types
101 return resp, nil
102 }
103
139 -func (a *apiClient) ensureCompatible(ctx context.Context) error {
104 +func (a *apiClient) ensureReady(ctx context.Context) error {
105 + if a.httpClient != nil && a.rawTLSConfig != nil {
106 + return nil
107 + }
108 +
109 + rootCAPEM := append([]byte(nil), a.rootCAPEM...)
110 + if len(rootCAPEM) == 0 && utils.IsLocalRelayHost(a.baseURL.Hostname()) {
111 + bootstrapParent := ctx
112 + if bootstrapParent == nil {
113 + bootstrapParent = context.Background()
114 + }
115 + bootstrapCtx, cancel := context.WithTimeout(bootstrapParent, defaultDialTimeout+defaultHandshakeTimeout)
116 + defer cancel()
117 +
118 + _, resolvedCAPEM, err := keyless.ResolveMaterials(bootstrapCtx, a.baseURL.String(), a.baseURL.Hostname())
119 + if err != nil {
120 + return fmt.Errorf("bootstrap localhost relay trust: %w", err)
121 + }
122 + rootCAPEM = resolvedCAPEM
123 + }
124 +
125 + rootCAs, err := utils.CertPoolFromPEM(rootCAPEM)
126 + if err != nil {
127 + return err
128 + }
129 +
130 + rawTLSConfig := &tls.Config{
131 + MinVersion: tls.VersionTLS12,
132 + ServerName: a.baseURL.Hostname(),
133 + RootCAs: rootCAs,
134 + NextProtos: []string{"http/1.1"},
135 + }
136 + httpClient := &http.Client{
137 + Transport: &http.Transport{
138 + TLSClientConfig: rawTLSConfig.Clone(),
139 + ForceAttemptHTTP2: false,
140 + },
141 + Timeout: a.requestTimeout,
142 + }
143 + if err := a.ensureCompatible(ctx, httpClient); err != nil {
144 + a.close()
145 + return err
146 + }
147 +
148 + a.close()
149 + a.httpClient = httpClient
150 + a.rawTLSConfig = rawTLSConfig
151 + return nil
152 +}
153 +
154 +func (a *apiClient) ensureCompatible(ctx context.Context, httpClient *http.Client) error {
155 var resp types.DomainResponse
141 - if err := a.doJSON(ctx, http.MethodGet, types.PathSDKDomain, nil, &resp); err != nil {
156 + if err := a.doJSONWithClient(ctx, httpClient, http.MethodGet, types.PathSDKDomain, nil, &resp); err != nil {
157 return fmt.Errorf("check relay compatibility: %w", err)
158 }
159 if strings.TrimSpace(resp.Version) != types.SDKProtocolVersion {
@@ -211,6 +226,14 @@ func (a *apiClient) openReverseSession(ctx context.Context, leaseID string) (net
226 }
227
228 func (a *apiClient) doJSON(ctx context.Context, method, path string, payload any, out any) error {
229 + return a.doJSONWithClient(ctx, a.httpClient, method, path, payload, out)
230 +}
231 +
232 +func (a *apiClient) doJSONWithClient(ctx context.Context, httpClient *http.Client, method, path string, payload any, out any) error {
233 + if httpClient == nil {
234 + return errors.New("api client is not ready")
235 + }
236 +
237 var body io.Reader
238 if payload != nil {
239 buf, err := json.Marshal(payload)
@@ -227,7 +250,7 @@ func (a *apiClient) doJSON(ctx context.Context, method, path string, payload any
250 }
251 req.Header.Set("Content-Type", "application/json")
252
230 - resp, err := a.httpClient.Do(req)
253 + resp, err := httpClient.Do(req)
254 if err != nil {
255 return err
256 }
sdk/expose.go
+1 -7
@@ -83,7 +83,7 @@ func Expose(ctx context.Context, relayUrls []string, name string, metadata types
83 Int("relay_count", len(exposure.relays)).
84 Strs("relays", exposure.RelayURLs()).
85 Strs("public_urls", exposure.PublicURLs()).
86 - Msg("exposure ready")
86 + Msg("exposure started")
87
88 return exposure, nil
89 }
@@ -373,7 +373,6 @@ func mergeListeners(listeners ...net.Listener) (net.Listener, error) {
373 merged.listeners = append(merged.listeners, listener)
374 }
375
376 - merged.addr = merged.buildAddr()
376 merged.active = len(merged.listeners)
377 for _, listener := range merged.listeners {
378 source := listener
@@ -386,7 +385,6 @@ type mergedListener struct {
385 listeners []net.Listener
386 accepted chan net.Conn
387 closed chan struct{}
389 - addr net.Addr
388
389 closeOnce sync.Once
390 mu sync.Mutex
@@ -426,10 +424,6 @@ func (l *mergedListener) Close() error {
424 }
425
426 func (l *mergedListener) Addr() net.Addr {
429 - return l.addr
430 -}
431 -
432 -func (l *mergedListener) buildAddr() net.Addr {
427 if len(l.listeners) == 1 {
428 return l.listeners[0].Addr()
429 }
sdk/listener.go
+62 -61
@@ -54,6 +54,7 @@ type Listener struct {
54 }
55
56 // NewListener creates one relay listener and its dedicated relay transport for one relay URL.
57 +// Only local config validation fails immediately; relay startup runs in the background until ready.
58 func NewListener(ctx context.Context, relayURL string, cfg ListenerConfig) (*Listener, error) {
59 if ctx == nil {
60 ctx = context.Background()
@@ -66,7 +67,7 @@ func NewListener(ctx context.Context, relayURL string, cfg ListenerConfig) (*Lis
67 renewBefore := utils.DurationOrDefault(cfg.RenewBefore, defaultRenewBefore)
68 retryWait := utils.DurationOrDefault(cfg.RetryWait, defaultRetryWait)
69
69 - api, err := newApiClient(listenerCtx, relayURL, cfg)
70 + api, err := newApiClient(relayURL, cfg)
71 if err != nil {
72 cancel()
73 return nil, err
@@ -85,42 +86,36 @@ func NewListener(ctx context.Context, relayURL string, cfg ListenerConfig) (*Lis
86 handshakeTimeout: handshakeTimeout,
87 }
88
88 - resp, err := api.registerLease(listenerCtx, leaseTTL)
89 - if err != nil {
90 - api.close()
91 - cancel()
92 - return nil, err
93 - }
94 -
95 - tlsConf, tlsCloser, err := keyless.BuildClientTLSConfig(api.baseURL.String(), []string{resp.Hostname})
96 - if err != nil {
97 - _ = api.unregisterLease(context.Background(), resp.LeaseID)
98 - api.close()
99 - cancel()
100 - return nil, err
101 - }
102 -
103 - if listenerCtx.Err() != nil {
104 - _ = api.unregisterLease(context.Background(), resp.LeaseID)
105 - _ = tlsCloser.Close()
106 - api.close()
107 - cancel()
108 - return nil, listenerCtx.Err()
109 - }
89 + go l.runStartup(listenerCtx)
90 + return l, nil
91 +}
92
111 - l.mu.Lock()
112 - l.leaseID = resp.LeaseID
113 - l.hostname = resp.Hostname
114 - l.metadata = resp.Metadata.Copy()
115 - l.tlsConfig = tlsConf
116 - l.tlsCloser = tlsCloser
117 - l.mu.Unlock()
93 +func (l *Listener) runStartup(ctx context.Context) {
94 + var retries int
95
119 - for i := 0; i < l.readyTarget; i++ {
120 - go l.runSessionLoop(listenerCtx)
96 + for {
97 + err := l.registerAndConfigure(ctx)
98 + switch {
99 + case err == nil:
100 + for i := 0; i < l.readyTarget; i++ {
101 + go l.runSessionLoop(ctx)
102 + }
103 + go l.runRenewLoop(ctx)
104 + log.Info().
105 + Str("component", "sdk-listener").
106 + Str("lease_id", l.LeaseID()).
107 + Str("hostname", l.Hostname()).
108 + Msg("listener ready")
109 + return
110 + case errors.Is(err, context.Canceled), errors.Is(err, net.ErrClosed):
111 + return
112 + default:
113 + retries++
114 + if !l.retryOrClose(ctx, "lease registration", err, retries) {
115 + return
116 + }
117 + }
118 }
122 - go l.runRenewLoop(listenerCtx)
123 - return l, nil
119 }
120
121 func (l *Listener) Accept() (net.Conn, error) {
@@ -333,36 +328,25 @@ func (l *Listener) renewLease(ctx context.Context) error {
328 return err
329 }
330
336 - log.Warn().
337 - Err(err).
338 - Str("component", "sdk-listener").
339 - Str("lease_id", leaseID).
340 - Msg("lease not found on relay, attempting re-registration")
341 -
331 if err := l.reregister(ctx); err != nil {
332 return err
333 }
345 -
346 - log.Info().
347 - Str("component", "sdk-listener").
348 - Str("lease_id", l.LeaseID()).
349 - Str("hostname", l.Hostname()).
350 - Msg("lease re-registered successfully")
334 return nil
335 }
336
354 -func (l *Listener) reregister(ctx context.Context) error {
355 - requestCtx, cancel := context.WithTimeout(ctx, 10*time.Second)
356 - defer cancel()
337 +func (l *Listener) registerAndConfigure(ctx context.Context) error {
338 + if err := l.api.ensureReady(ctx); err != nil {
339 + return err
340 + }
341
358 - resp, err := l.api.registerLease(requestCtx, l.leaseTTL)
342 + resp, err := l.api.registerLease(ctx, l.leaseTTL)
343 if err != nil {
344 return err
345 }
346
347 tlsConf, tlsCloser, err := keyless.BuildClientTLSConfig(l.api.baseURL.String(), []string{resp.Hostname})
348 if err != nil {
365 - _ = l.api.unregisterLease(requestCtx, resp.LeaseID)
349 + _ = l.api.unregisterLease(context.Background(), resp.LeaseID)
350 return err
351 }
352
@@ -373,6 +357,12 @@ func (l *Listener) reregister(ctx context.Context) error {
357 }
358
359 l.mu.Lock()
360 + if ctx.Err() != nil {
361 + l.mu.Unlock()
362 + _ = l.api.unregisterLease(context.Background(), resp.LeaseID)
363 + _ = tlsCloser.Close()
364 + return ctx.Err()
365 + }
366 oldCloser := l.tlsCloser
367 l.leaseID = resp.LeaseID
368 l.hostname = resp.Hostname
@@ -387,6 +377,13 @@ func (l *Listener) reregister(ctx context.Context) error {
377 return nil
378 }
379
380 +func (l *Listener) reregister(ctx context.Context) error {
381 + requestCtx, cancel := context.WithTimeout(ctx, 10*time.Second)
382 + defer cancel()
383 +
384 + return l.registerAndConfigure(requestCtx)
385 +}
386 +
387 func isLeaseNotFound(err error) bool {
388 return errors.Is(err, &types.APIRequestError{Code: types.APIErrorCodeLeaseNotFound})
389 }
@@ -403,20 +400,24 @@ func (l *Listener) retryOrClose(ctx context.Context, operation string, err error
400 Logger()
401
402 if l.retryCount > 0 && retries > l.retryCount {
406 - logger.Error().
407 - Err(err).
408 - Int("retry_count", l.retryCount).
409 - Msg("retry budget exhausted; closing listener")
403 + if operation != "lease renewal" {
404 + logger.Error().
405 + Err(err).
406 + Int("retry_count", l.retryCount).
407 + Msg("retry budget exhausted; closing listener")
408 + }
409 _ = l.Close()
410 return false
411 }
412
414 - logger.Warn().
415 - Err(err).
416 - Int("retry_attempt", retries).
417 - Int("retry_count", l.retryCount).
418 - Dur("retry_wait", l.retryWait).
419 - Msg("operation failed; retrying")
413 + if operation != "lease renewal" {
414 + logger.Warn().
415 + Err(err).
416 + Int("retry_attempt", retries).
417 + Int("retry_count", l.retryCount).
418 + Dur("retry_wait", l.retryWait).
419 + Msg("operation failed; retrying")
420 + }
421
422 return utils.SleepOrDone(ctx, l.retryWait)
423 }
sdk/sdk_test.go
+152 -39
@@ -5,6 +5,7 @@ import (
5 "encoding/json"
6 "net/http"
7 "net/http/httptest"
8 + "strings"
9 "sync/atomic"
10 "testing"
11 "time"
@@ -12,12 +13,16 @@ import (
13 "github.com/gosuda/portal/v2/types"
14 )
15
15 -func TestNewListenerFailsFastOnRegisterError(t *testing.T) {
16 - t.Parallel()
17 -
16 +func TestNewListenerRetriesInitialStartupUntilReady(t *testing.T) {
17 + var domainCount atomic.Int32
18 + var registerCount atomic.Int32
19 server := httptest.NewTLSServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
20 switch r.URL.Path {
21 case types.PathSDKDomain:
22 + if domainCount.Add(1) == 1 {
23 + http.Error(w, "temporarily unavailable", http.StatusBadGateway)
24 + return
25 + }
26 writeSDKTestEnvelope(w, http.StatusOK, types.APIEnvelope[types.DomainResponse]{
27 OK: true,
28 Data: types.DomainResponse{
@@ -25,31 +30,50 @@ func TestNewListenerFailsFastOnRegisterError(t *testing.T) {
30 },
31 })
32 case types.PathSDKRegister:
28 - writeSDKTestEnvelope(w, http.StatusConflict, types.APIEnvelope[any]{
33 + if registerCount.Add(1) == 1 {
34 + http.Error(w, "temporarily unavailable", http.StatusBadGateway)
35 + return
36 + }
37 + writeSDKTestEnvelope(w, http.StatusCreated, types.APIEnvelope[types.RegisterResponse]{
38 + OK: true,
39 + Data: types.RegisterResponse{
40 + LeaseID: "lease-1",
41 + Hostname: "127.0.0.1",
42 + },
43 + })
44 + case types.PathSDKConnect:
45 + writeSDKTestEnvelope(w, http.StatusForbidden, types.APIEnvelope[any]{
46 OK: false,
30 - Error: &types.APIError{Code: types.APIErrorCodeHostnameConflict, Message: "hostname already registered"},
47 + Error: &types.APIError{Code: types.APIErrorCodeUnauthorized, Message: "not used in test"},
48 })
49 + case types.PathSDKRenew:
50 + writeSDKTestEnvelope(w, http.StatusOK, types.APIEnvelope[types.RenewResponse]{
51 + OK: true,
52 + Data: types.RenewResponse{LeaseID: "lease-1"},
53 + })
54 + case types.PathSDKUnregister:
55 + writeSDKTestEnvelope(w, http.StatusOK, types.APIEnvelope[any]{OK: true})
56 default:
57 http.NotFound(w, r)
34 - return
58 }
59 }))
60 defer server.Close()
61
62 listener, err := NewListener(context.Background(), server.URL, ListenerConfig{
40 - Name: "demo",
63 + Name: "demo",
64 + RetryWait: 10 * time.Millisecond,
65 })
42 - if err == nil {
43 - t.Fatal("NewListener() error = nil, want register failure")
44 - }
45 - if listener != nil {
46 - t.Fatalf("NewListener() listener = %#v, want nil", listener)
66 + if err != nil {
67 + t.Fatalf("NewListener() error = %v", err)
68 }
69 + defer listener.Close()
70 +
71 + waitForSDKTest(t, func() bool {
72 + return domainCount.Load() >= 2 && registerCount.Load() >= 2 && listener.LeaseID() == "lease-1"
73 + })
74 }
75
76 func TestNewListenerRejectsInvalidName(t *testing.T) {
51 - t.Parallel()
52 -
77 listener, err := NewListener(context.Background(), "https://relay.example.com", ListenerConfig{Name: "demo app"})
78 if err == nil {
79 t.Fatal("NewListener() error = nil, want invalid name error")
@@ -60,9 +84,7 @@ func TestNewListenerRejectsInvalidName(t *testing.T) {
84 }
85
86 func TestNewListenerRegistersLeaseWithMainContract(t *testing.T) {
63 - t.Parallel()
64 -
65 - var registerReq types.RegisterRequest
87 + registerReqCh := make(chan types.RegisterRequest, 1)
88 server := httptest.NewTLSServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
89 switch r.URL.Path {
90 case types.PathSDKDomain:
@@ -73,9 +95,14 @@ func TestNewListenerRegistersLeaseWithMainContract(t *testing.T) {
95 },
96 })
97 case types.PathSDKRegister:
98 + var registerReq types.RegisterRequest
99 if err := json.NewDecoder(r.Body).Decode(&registerReq); err != nil {
100 t.Fatalf("decode register request: %v", err)
101 }
102 + select {
103 + case registerReqCh <- registerReq:
104 + default:
105 + }
106 writeSDKTestEnvelope(w, http.StatusCreated, types.APIEnvelope[types.RegisterResponse]{
107 OK: true,
108 Data: types.RegisterResponse{
@@ -112,6 +139,19 @@ func TestNewListenerRegistersLeaseWithMainContract(t *testing.T) {
139 }
140 defer listener.Close()
141
142 + var registerReq types.RegisterRequest
143 + waitForSDKTest(t, func() bool {
144 + select {
145 + case registerReq = <-registerReqCh:
146 + return true
147 + default:
148 + return false
149 + }
150 + })
151 + waitForSDKTest(t, func() bool {
152 + return listener.LeaseID() == "lease-1"
153 + })
154 +
155 if registerReq.TTL != 42 {
156 t.Fatalf("register request TTL = %d, want 42", registerReq.TTL)
157 }
@@ -133,8 +173,6 @@ func TestNewListenerRegistersLeaseWithMainContract(t *testing.T) {
173 }
174
175 func TestNewListenerReregistersOnLeaseNotFound(t *testing.T) {
136 - t.Parallel()
137 -
176 var registerCount atomic.Int32
177 server := httptest.NewTLSServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
178 switch r.URL.Path {
@@ -203,8 +241,6 @@ func TestNewListenerReregistersOnLeaseNotFound(t *testing.T) {
241 }
242
243 func TestNewListenerClosesAfterReverseSessionRetryBudgetExhausted(t *testing.T) {
206 - t.Parallel()
207 -
244 var connectCount atomic.Int32
245 var unregisterCount atomic.Int32
246 server := httptest.NewTLSServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
@@ -260,10 +296,58 @@ func TestNewListenerClosesAfterReverseSessionRetryBudgetExhausted(t *testing.T)
296 }
297 }
298
263 -func TestExposeFailsFastWhenAnyRelayCannotRegister(t *testing.T) {
264 - t.Parallel()
299 +func TestNewListenerRetriesForeverWhenRetryCountIsNegative(t *testing.T) {
300 + var connectCount atomic.Int32
301 + server := httptest.NewTLSServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
302 + switch r.URL.Path {
303 + case types.PathSDKDomain:
304 + writeSDKTestEnvelope(w, http.StatusOK, types.APIEnvelope[types.DomainResponse]{
305 + OK: true,
306 + Data: types.DomainResponse{
307 + Version: types.SDKProtocolVersion,
308 + },
309 + })
310 + case types.PathSDKRegister:
311 + writeSDKTestEnvelope(w, http.StatusCreated, types.APIEnvelope[types.RegisterResponse]{
312 + OK: true,
313 + Data: types.RegisterResponse{
314 + LeaseID: "lease-1",
315 + Hostname: "127.0.0.1",
316 + },
317 + })
318 + case types.PathSDKConnect:
319 + connectCount.Add(1)
320 + writeSDKTestEnvelope(w, http.StatusForbidden, types.APIEnvelope[any]{
321 + OK: false,
322 + Error: &types.APIError{Code: types.APIErrorCodeUnauthorized, Message: "reverse session denied"},
323 + })
324 + case types.PathSDKUnregister:
325 + writeSDKTestEnvelope(w, http.StatusOK, types.APIEnvelope[any]{OK: true})
326 + default:
327 + http.NotFound(w, r)
328 + }
329 + }))
330 + defer server.Close()
331
266 - var unregisterCount atomic.Int32
332 + listener, err := NewListener(context.Background(), server.URL, ListenerConfig{
333 + Name: "demo",
334 + RetryCount: -1,
335 + RetryWait: 10 * time.Millisecond,
336 + })
337 + if err != nil {
338 + t.Fatalf("NewListener() error = %v", err)
339 + }
340 + defer listener.Close()
341 +
342 + waitForSDKTest(t, func() bool {
343 + return connectCount.Load() >= 3
344 + })
345 + if listener.done() {
346 + t.Fatal("listener closed unexpectedly with negative RetryCount")
347 + }
348 +}
349 +
350 +func TestExposeAddsRecoveredRelayWithoutDroppingHealthyRelay(t *testing.T) {
351 goodServer := httptest.NewTLSServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
352 switch r.URL.Path {
353 case types.PathSDKDomain:
@@ -292,7 +376,6 @@ func TestExposeFailsFastWhenAnyRelayCannotRegister(t *testing.T) {
376 Data: types.RenewResponse{LeaseID: "lease-good"},
377 })
378 case types.PathSDKUnregister:
295 - unregisterCount.Add(1)
379 writeSDKTestEnvelope(w, http.StatusOK, types.APIEnvelope[any]{OK: true})
380 default:
381 http.NotFound(w, r)
@@ -300,9 +383,14 @@ func TestExposeFailsFastWhenAnyRelayCannotRegister(t *testing.T) {
383 }))
384 defer goodServer.Close()
385
303 - badServer := httptest.NewTLSServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
386 + var delayedDomainCount atomic.Int32
387 + delayedServer := httptest.NewTLSServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
388 switch r.URL.Path {
389 case types.PathSDKDomain:
390 + if delayedDomainCount.Add(1) == 1 {
391 + http.Error(w, "temporarily unavailable", http.StatusBadGateway)
392 + return
393 + }
394 writeSDKTestEnvelope(w, http.StatusOK, types.APIEnvelope[types.DomainResponse]{
395 OK: true,
396 Data: types.DomainResponse{
@@ -310,33 +398,52 @@ func TestExposeFailsFastWhenAnyRelayCannotRegister(t *testing.T) {
398 },
399 })
400 case types.PathSDKRegister:
313 - writeSDKTestEnvelope(w, http.StatusConflict, types.APIEnvelope[any]{
401 + writeSDKTestEnvelope(w, http.StatusCreated, types.APIEnvelope[types.RegisterResponse]{
402 + OK: true,
403 + Data: types.RegisterResponse{
404 + LeaseID: "lease-delayed",
405 + Hostname: "127.0.0.2",
406 + },
407 + })
408 + case types.PathSDKConnect:
409 + writeSDKTestEnvelope(w, http.StatusForbidden, types.APIEnvelope[any]{
410 OK: false,
315 - Error: &types.APIError{Code: types.APIErrorCodeHostnameConflict, Message: "hostname already registered"},
411 + Error: &types.APIError{Code: types.APIErrorCodeUnauthorized, Message: "not used in test"},
412 + })
413 + case types.PathSDKRenew:
414 + writeSDKTestEnvelope(w, http.StatusOK, types.APIEnvelope[types.RenewResponse]{
415 + OK: true,
416 + Data: types.RenewResponse{LeaseID: "lease-delayed"},
417 })
418 + case types.PathSDKUnregister:
419 + writeSDKTestEnvelope(w, http.StatusOK, types.APIEnvelope[any]{OK: true})
420 default:
421 http.NotFound(w, r)
319 - return
422 }
423 }))
322 - defer badServer.Close()
424 + defer delayedServer.Close()
425
324 - exposure, err := Expose(context.Background(), []string{goodServer.URL, badServer.URL}, "demo", types.LeaseMetadata{})
325 - if err == nil {
326 - t.Fatal("Expose() error = nil, want register failure")
426 + exposure, err := Expose(context.Background(), []string{goodServer.URL, delayedServer.URL}, "demo", types.LeaseMetadata{})
427 + if err != nil {
428 + t.Fatalf("Expose() error = %v", err)
429 }
328 - if exposure != nil {
329 - t.Fatalf("Expose() exposure = %#v, want nil", exposure)
430 + if exposure == nil {
431 + t.Fatal("Expose() exposure = nil, want non-nil")
432 }
433 + defer exposure.Close()
434
435 waitForSDKTest(t, func() bool {
333 - return unregisterCount.Load() > 0
436 + return len(exposure.PublicURLs()) >= 1
437 + })
438 + waitForSDKTestWithTimeout(t, 15*time.Second, func() bool {
439 + return delayedDomainCount.Load() >= 2 &&
440 + len(exposure.PublicURLs()) == 2 &&
441 + strings.Contains(exposure.Addr().String(), "lease-good") &&
442 + strings.Contains(exposure.Addr().String(), "lease-delayed")
443 })
444 }
445
446 func TestExposeNoRelayInputs(t *testing.T) {
338 - t.Parallel()
339 -
447 exposure, err := Expose(context.Background(), nil, "demo", types.LeaseMetadata{})
448 if err != nil {
449 t.Fatalf("Expose() error = %v", err)
@@ -355,7 +462,13 @@ func writeSDKTestEnvelope[T any](w http.ResponseWriter, status int, envelope typ
462 func waitForSDKTest(t *testing.T, fn func() bool) {
463 t.Helper()
464
358 - deadline := time.Now().Add(5 * time.Second)
465 + waitForSDKTestWithTimeout(t, 5*time.Second, fn)
466 +}
467 +
468 +func waitForSDKTestWithTimeout(t *testing.T, timeout time.Duration, fn func() bool) {
469 + t.Helper()
470 +
471 + deadline := time.Now().Add(timeout)
472 for time.Now().Before(deadline) {
473 if fn() {
474 return
types/error.go renamed