refact: simplify bootstrap logic
rabbitprincess committed
Mar 20, 2026 at 23:55 UTC
69c44afcbaabd61db0810fa990e9a15e7dabb719
14 files changed
+422
-587
cmd/demo-app/main.go
+13
-16
@@ -34,8 +34,8 @@ func main() {
34
log.Logger = log.Output(zerolog.ConsoleWriter{Out: os.Stdout, TimeFormat: time.RFC3339})
35
logger := log.With().Str("component", "demo-app").Logger()
36
37
- flag.StringVar(&flagRelayURLs, "relays", "https://localhost:4017", "additional relay API URLs (comma-separated; scheme omitted defaults to https; appended to registry.json defaults unless --default-relays=false is set) [env: RELAYS]")
38
- flag.BoolVar(&flagDefaultRelays, "default-relays", utils.ParseBoolEnv("DEFAULT_RELAYS", true), "include repository registry.json default relays [env: DEFAULT_RELAYS]")
37
+ flag.StringVar(&flagRelayURLs, "relays", "https://localhost:4017", "additional relay API URLs (comma-separated; scheme omitted defaults to https; merged with public registry relays unless --default-relays=false is set) [env: RELAYS]")
38
+ flag.BoolVar(&flagDefaultRelays, "default-relays", utils.ParseBoolEnv("DEFAULT_RELAYS", true), "include public registry relays [env: DEFAULT_RELAYS]")
39
flag.StringVar(&flagAddr, "addr", "127.0.0.1:8092", "local demo HTTP listen address (host:port or URL; disable if empty)")
40
flag.StringVar(&flagName, "name", "demo-app", "public hostname prefix (single DNS label)")
41
flag.StringVar(&flagDesc, "description", "Portal demo connectivity app", "lease description")
@@ -58,20 +58,17 @@ func runDemo() error {
58
defer stop()
59
60
relayURLs := utils.SplitCSV(flagRelayURLs)
61
- if flagDefaultRelays {
62
- relayURLs = sdk.WithDefaultRelayURLs(ctx, "", relayURLs...)
63
- }
64
- relayURLs, err := utils.NormalizeRelayURLs(relayURLs)
65
- if err != nil {
66
- return fmt.Errorf("resolve relay urls: %w", err)
67
- }
68
-
69
- exposure, err := sdk.Expose(ctx, relayURLs, flagName, false, types.LeaseMetadata{
70
- Description: flagDesc,
71
- Tags: utils.SplitCSV(flagTags),
72
- Owner: flagOwner,
73
- Thumbnail: flagThumbnail,
74
- Hide: flagHide,
61
+ exposure, err := sdk.Expose(ctx, sdk.ExposeConfig{
62
+ RelayURLs: relayURLs,
63
+ DefaultRelayEnabled: flagDefaultRelays,
64
+ Name: flagName,
65
+ Metadata: types.LeaseMetadata{
66
+ Description: flagDesc,
67
+ Tags: utils.SplitCSV(flagTags),
68
+ Owner: flagOwner,
69
+ Thumbnail: flagThumbnail,
70
+ Hide: flagHide,
71
+ },
72
})
73
if err != nil {
74
return fmt.Errorf("exposure listen error: %w", err)
cmd/demo-udp/main.go
+14
-16
@@ -36,8 +36,8 @@ func main() {
36
log.Logger = log.Output(zerolog.ConsoleWriter{Out: os.Stdout, TimeFormat: time.RFC3339})
37
logger := log.With().Str("component", "demo-udp").Logger()
38
39
- flag.StringVar(&flagRelayURLs, "relays", "https://localhost:4017", "additional relay API URLs (comma-separated; scheme omitted defaults to https; appended to registry.json defaults unless --default-relays=false is set) [env: RELAYS]")
40
- flag.BoolVar(&flagDefaultRelays, "default-relays", utils.ParseBoolEnv("DEFAULT_RELAYS", false), "include repository registry.json default relays [env: DEFAULT_RELAYS]")
39
+ flag.StringVar(&flagRelayURLs, "relays", "https://localhost:4017", "additional relay API URLs (comma-separated; scheme omitted defaults to https; merged with public registry relays unless --default-relays=false is set) [env: RELAYS]")
40
+ flag.BoolVar(&flagDefaultRelays, "default-relays", utils.ParseBoolEnv("DEFAULT_RELAYS", false), "include public registry relays [env: DEFAULT_RELAYS]")
41
flag.StringVar(&flagName, "name", "demo-udp", "public hostname prefix (single DNS label)")
42
flag.StringVar(&flagDesc, "description", "Portal demo UDP echo service", "lease description")
43
flag.StringVar(&flagTags, "tags", "demo,udp,echo", "comma-separated lease tags")
@@ -59,20 +59,18 @@ func runDemoUDP() error {
59
defer stop()
60
61
relayURLs := utils.SplitCSV(flagRelayURLs)
62
- if flagDefaultRelays {
63
- relayURLs = sdk.WithDefaultRelayURLs(ctx, "", relayURLs...)
64
- }
65
- relayURLs, err := utils.NormalizeRelayURLs(relayURLs)
66
- if err != nil {
67
- return fmt.Errorf("resolve relay urls: %w", err)
68
- }
69
-
70
- exposure, err := sdk.Expose(ctx, relayURLs, flagName, true, types.LeaseMetadata{
71
- Description: flagDesc,
72
- Tags: utils.SplitCSV(flagTags),
73
- Owner: flagOwner,
74
- Thumbnail: flagThumbnail,
75
- Hide: flagHide,
62
+ exposure, err := sdk.Expose(ctx, sdk.ExposeConfig{
63
+ RelayURLs: relayURLs,
64
+ DefaultRelayEnabled: flagDefaultRelays,
65
+ Name: flagName,
66
+ UDPEnabled: true,
67
+ Metadata: types.LeaseMetadata{
68
+ Description: flagDesc,
69
+ Tags: utils.SplitCSV(flagTags),
70
+ Owner: flagOwner,
71
+ Thumbnail: flagThumbnail,
72
+ Hide: flagHide,
73
+ },
74
})
75
if err != nil {
76
return fmt.Errorf("exposure listen error: %w", err)
cmd/portal-tunnel/README.md
+2
-2
@@ -42,7 +42,7 @@ Flags:
42
43
```text
44
--relays Portal relay API URLs (comma-separated, https only)
45
---default-relays Include repository registry.json public relays
45
+--default-relays Include public registry relays
46
--name Public hostname prefix (single DNS label); auto-generated when omitted
47
--description Service description metadata
48
--tags Service tags metadata (comma-separated)
@@ -78,7 +78,7 @@ Legacy execution compatibility has been removed:
78
- The tunnel consumes one aggregate SDK listener, so the CLI no longer manages per-relay listener loops itself.
79
- 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.
80
- 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.
81
-- The configured relay list is either `registry.json + installed/configured relay URLs` or, with `--default-relays=false`, just the explicit relay URLs. Published public URLs appear only for relays that have registered successfully.
81
+- The configured relay list is either `public registry + installed/configured relay URLs` or, with `--default-relays=false`, just the explicit relay URLs. Published public URLs appear only for relays that have registered successfully.
82
- SDK callers that do not set `ListenerConfig.RetryCount` use infinite retry semantics for each relay.
83
- Tenant TLS is provisioned automatically through the relay keyless signer. The SDK fetches the relay certificate chain and uses `/v1/sign` for remote signing.
84
- When the local service is unreachable, the tunnel returns an HTTP 503 page.
cmd/portal-tunnel/main.go
+7
-25
@@ -168,14 +168,6 @@ func runExposeCommand(args []string) error {
168
relayInputs = []string{explicitRelays}
169
}
170
171
- relayURLs, err := resolveRelayURLs(ctx, "", relayInputs, defaultRelays)
172
- if err != nil {
173
- return fmt.Errorf("resolve relay urls: %w", err)
174
- }
175
- if len(relayURLs) == 0 {
176
- return errors.New("no relay URLs configured; run the installer first or pass --relays")
177
- }
178
-
171
previousOwnerPrivateKey := cfg.OwnerPrivateKey
172
if strings.TrimSpace(privateKey) != "" {
173
cfg.OwnerPrivateKey = privateKey
@@ -185,11 +177,12 @@ func runExposeCommand(args []string) error {
177
ownerPrivateKey = &cfg.OwnerPrivateKey
178
}
179
188
- exposure, err := sdk.ExposeWithConfig(ctx, sdk.ExposeConfig{
189
- RelayURLs: relayURLs,
190
- Name: name,
191
- UDPEnabled: udp,
192
- Discovery: discoveryEnabled,
180
+ exposure, err := sdk.Expose(ctx, sdk.ExposeConfig{
181
+ RelayURLs: relayInputs,
182
+ DefaultRelayEnabled: defaultRelays,
183
+ Name: name,
184
+ UDPEnabled: udp,
185
+ Discovery: discoveryEnabled,
186
Metadata: types.LeaseMetadata{
187
Description: desc,
188
Tags: utils.SplitCSV(tags),
@@ -261,7 +254,7 @@ func runListCommand(args []string) error {
254
relayInputs = []string{explicitRelays}
255
}
256
264
- relayURLs, err := resolveRelayURLs(ctx, "", relayInputs, defaultRelays)
257
+ relayURLs, err := sdk.ResolveRelayURLs(ctx, relayInputs, defaultRelays)
258
if err != nil {
259
return fmt.Errorf("resolve relay urls: %w", err)
260
}
@@ -366,17 +359,6 @@ func runTunnel(
359
return errors.Join(waitErr, udpErr, closeErr)
360
}
361
369
-func resolveRelayURLs(ctx context.Context, registryURL string, inputs []string, includeDefaultRelays bool) ([]string, error) {
370
- if includeDefaultRelays {
371
- relayURLs := sdk.WithDefaultRelayURLs(ctx, registryURL, inputs...)
372
- if len(relayURLs) == 0 {
373
- return nil, nil
374
- }
375
- return relayURLs, nil
376
- }
377
- return utils.NormalizeRelayURLs(inputs)
378
-}
379
-
362
var exposeNameOpeners = []string{
363
"arcade", "bouncy", "bravo", "bubble", "candy", "cosmic", "dapper", "electric",
364
"fancy", "fizzy", "flashy", "fuzzy", "gentle", "glitter", "golden", "happy",
docs/architecture.md
+3
-2
@@ -115,12 +115,13 @@ That distinction matters because `/sdk/connect` stops being ordinary HTTP once h
115
116
### SDK (`sdk/`)
117
118
-- `WithDefaultRelayURLs`: fetches the default Portal relay list from the repository-root `registry.json`, appends explicit relay inputs, and normalizes the combined list
118
+- `ExposeConfig.DefaultRelayEnabled`: when true, `Expose` fetches the default Portal relay registry, merges it with explicit relay inputs, and normalizes the result
119
- Entry points can opt out of registry defaults and call `utils.NormalizeRelayURLs` directly when they need explicit relay inputs only
120
- `Listener`: validates one relay URL locally, then starts relay compatibility checks, lease registration, reverse session maintenance, and lease renewal in the background until ready
121
- `api_client.go`: internal relay client for control-plane requests, reverse session dialing, and internal QUIC tunnel setup
122
- `ListenerConfig.RetryCount <= 0` means retry forever; positive values close the listener after the retry budget is exhausted
123
-- Default app flow is `WithDefaultRelayURLs -> NewListener -> PublicURL -> http.Server.Serve(listener)` or `WithDefaultRelayURLs -> Expose -> PublicURLs -> http.Server.Serve(exposure)`, with an opt-out path for explicit relay inputs only
123
+- `NewListener` callers provide explicit normalized relay URLs
124
+- Default exposure flow is `Expose{DefaultRelayEnabled: true} -> PublicURLs -> http.Server.Serve(exposure)`, with an opt-out path for explicit relay inputs only
125
- `expose.go`: optional `RunHTTP` helper for serving one handler on both a local HTTP port and the relay listener
126
- `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
127
- `Exposure.RelayURLs()` returns the configured normalized relay URLs, while `Exposure.PublicURLs()` returns only relays that are currently registered and ready
portal/api_server.go
+5
-5
@@ -80,7 +80,7 @@ func (s *Server) apiHandler(base *http.ServeMux, keylessSignerHandler http.Handl
80
base.ServeHTTP(w, r)
81
return
82
}
83
- s.discovery.ServeHTTP(w, r)
83
+ discovery.ServeHTTP(w, r, []string{s.cfg.PortalURL}, s.discoveryBootstrapsSnapshot(), s.discover)
84
case types.PathV1Sign:
85
if keylessSignerHandler == nil {
86
http.NotFound(w, r)
@@ -517,8 +517,8 @@ func (s *Server) registerLease(req types.RegisterRequest, clientIP string) (type
517
record.Close()
518
return types.RegisterResponse{}, err
519
}
520
- if s.discovery != nil {
521
- if err := s.discovery.MergeBootstraps(bootstraps); err != nil {
520
+ if s.DiscoveryEnabled() {
521
+ if err := s.mergeDiscoveryBootstraps(bootstraps); err != nil {
522
record.Close()
523
_, _ = s.registry.Unregister(record.ID, record.ReverseToken)
524
return types.RegisterResponse{}, err
@@ -526,8 +526,8 @@ func (s *Server) registerLease(req types.RegisterRequest, clientIP string) (type
526
}
527
528
responseBootstraps := append([]string(nil), s.cfg.Bootstraps...)
529
- if s.discovery != nil {
530
- responseBootstraps = s.discovery.Bootstraps()
529
+ if s.DiscoveryEnabled() {
530
+ responseBootstraps = s.discoveryBootstrapsSnapshot()
531
} else {
532
responseBootstraps, err = utils.NormalizeRelayURLs(append(responseBootstraps, record.Bootstraps...))
533
if err != nil {
portal/discovery/discovery.go
+82
-251
@@ -3,143 +3,66 @@ package discovery
3
import (
4
"context"
5
"crypto/tls"
6
- "encoding/json"
6
"errors"
7
"fmt"
8
"net/http"
9
"net/url"
10
"strings"
12
- "sync"
11
"time"
12
15
- "github.com/rs/zerolog/log"
16
-
13
"github.com/gosuda/portal/v2/portal/keyless"
14
"github.com/gosuda/portal/v2/types"
15
"github.com/gosuda/portal/v2/utils"
16
)
17
22
-const (
23
- defaultRequestTimeout = 15 * time.Second
24
- defaultMaxPeers = 32
25
- defaultPollInterval = 30 * time.Second
26
-)
27
-
28
-type Resolver interface {
29
- Discover(context.Context, types.DiscoverRequest) (types.DiscoverResponse, error)
30
-}
31
-
32
-type ResolverFunc func(context.Context, types.DiscoverRequest) (types.DiscoverResponse, error)
33
-
34
-func (fn ResolverFunc) Discover(ctx context.Context, req types.DiscoverRequest) (types.DiscoverResponse, error) {
35
- return fn(ctx, req)
36
-}
37
-
38
-type Config struct {
39
- SelfURLs []string
40
- Bootstraps []string
41
- RootCAPEM []byte
42
- RequestTimeout time.Duration
43
- MaxPeers int
44
- OnBootstraps func([]string)
45
-}
46
-
47
-type Service struct {
48
- resolver Resolver
49
- selfURLs []string
50
- rootCAPEM []byte
51
- requestTimeout time.Duration
52
- maxPeers int
53
- onBootstraps func([]string)
18
+type Resolver func(context.Context, types.DiscoverRequest) (types.DiscoverResponse, error)
19
55
- mu sync.RWMutex
56
- bootstraps []string
57
-}
20
+const defaultRequestTimeout = 15 * time.Second
21
59
-func New(cfg Config, resolver Resolver) (*Service, error) {
60
- service := &Service{
61
- resolver: resolver,
62
- rootCAPEM: append([]byte(nil), cfg.RootCAPEM...),
63
- requestTimeout: utils.DurationOrDefault(cfg.RequestTimeout, defaultRequestTimeout),
64
- maxPeers: utils.IntOrDefault(cfg.MaxPeers, defaultMaxPeers),
65
- onBootstraps: cfg.OnBootstraps,
66
- }
67
- if err := service.SetSelfURLs(cfg.SelfURLs); err != nil {
68
- return nil, err
69
- }
70
- if err := service.MergeBootstraps(cfg.Bootstraps); err != nil {
71
- return nil, err
72
- }
73
- return service, nil
74
-}
75
-
76
-func (s *Service) Bootstraps() []string {
77
- if s == nil {
78
- return nil
79
- }
80
-
81
- s.mu.RLock()
82
- defer s.mu.RUnlock()
83
- return append([]string(nil), s.bootstraps...)
84
-}
85
-
86
-func (s *Service) MergeBootstraps(inputs []string) error {
87
- if s == nil || len(inputs) == 0 {
88
- return nil
22
+func DiscoverBootstraps(ctx context.Context, peers []string, req types.DiscoverRequest, rootCAPEM []byte) ([]string, error) {
23
+ if ctx == nil {
24
+ ctx = context.Background()
25
}
26
91
- normalized, err := utils.NormalizeRelayURLs(inputs)
27
+ peers, err := utils.NormalizeRelayURLs(peers)
28
if err != nil {
93
- return fmt.Errorf("normalize bootstraps: %w", err)
29
+ return nil, err
30
+ }
31
+ if len(peers) == 0 {
32
+ return nil, nil
33
}
34
96
- s.mu.Lock()
97
- combined, err := utils.NormalizeRelayURLs(append(append([]string(nil), s.bootstraps...), normalized...))
35
+ req, err = normalizeRequest(req)
36
if err != nil {
99
- s.mu.Unlock()
100
- return fmt.Errorf("normalize bootstraps: %w", err)
101
- }
102
- next := utils.ExcludeURLs(combined, s.selfURLs)
103
- changed := strings.Join(s.bootstraps, "\x00") != strings.Join(next, "\x00")
104
- callback := s.onBootstraps
105
- if changed {
106
- s.bootstraps = next
37
+ return nil, err
38
}
108
- bootstraps := append([]string(nil), s.bootstraps...)
109
- s.mu.Unlock()
39
111
- if changed && callback != nil {
112
- callback(bootstraps)
113
- }
114
- return nil
115
-}
40
+ bootstraps := append([]string(nil), peers...)
41
+ var discoverErr error
42
+ discovered := false
43
117
-func (s *Service) SetSelfURLs(inputs []string) error {
118
- if s == nil {
119
- return nil
120
- }
44
+ for _, peer := range peers {
45
+ resp, err := discoverPeer(ctx, peer, req, rootCAPEM)
46
+ if err != nil {
47
+ discoverErr = errors.Join(discoverErr, fmt.Errorf("discover %q: %w", peer, err))
48
+ continue
49
+ }
50
122
- selfURLs, err := utils.NormalizeRelayURLs(inputs)
123
- if err != nil {
124
- return fmt.Errorf("normalize self urls: %w", err)
51
+ bootstraps, err = utils.MergeRelayURLs(bootstraps, nil, resp.Bootstraps)
52
+ if err != nil {
53
+ discoverErr = errors.Join(discoverErr, fmt.Errorf("merge %q bootstraps: %w", peer, err))
54
+ continue
55
+ }
56
+ discovered = true
57
}
58
127
- s.mu.Lock()
128
- current := strings.Join(s.bootstraps, "\x00")
129
- s.selfURLs = selfURLs
130
- s.bootstraps = utils.ExcludeURLs(s.bootstraps, s.selfURLs)
131
- changed := current != strings.Join(s.bootstraps, "\x00")
132
- callback := s.onBootstraps
133
- bootstraps := append([]string(nil), s.bootstraps...)
134
- s.mu.Unlock()
135
-
136
- if changed && callback != nil {
137
- callback(bootstraps)
59
+ if !discovered {
60
+ return bootstraps, discoverErr
61
}
139
- return nil
62
+ return bootstraps, nil
63
}
64
142
-func (s *Service) ServeHTTP(w http.ResponseWriter, r *http.Request) {
65
+func ServeHTTP(w http.ResponseWriter, r *http.Request, selfURLs, bootstraps []string, resolver Resolver) {
66
if r.Method != http.MethodGet {
67
utils.WriteAPIError(w, http.StatusMethodNotAllowed, types.APIErrorCodeMethodNotAllowed, "method not allowed")
68
return
@@ -154,134 +77,35 @@ func (s *Service) ServeHTTP(w http.ResponseWriter, r *http.Request) {
77
return
78
}
79
157
- bootstraps, err := s.responseBootstraps(nil)
80
+ resolvedBootstraps, err := buildResponseBootstraps(selfURLs, bootstraps, nil)
81
if err != nil {
82
utils.WriteAPIError(w, http.StatusInternalServerError, types.APIErrorCodeInternal, err.Error())
83
return
84
}
85
+
86
resp := types.DiscoverResponse{
87
Found: false,
164
- Bootstraps: bootstraps,
88
+ Bootstraps: resolvedBootstraps,
89
}
166
- if s == nil || s.resolver == nil {
90
+ if resolver == nil {
91
utils.WriteAPIData(w, http.StatusOK, resp)
92
return
93
}
94
171
- localResp, err := s.resolver.Discover(r.Context(), req)
95
+ localResp, err := resolver(r.Context(), req)
96
if err != nil {
97
utils.WriteAPIError(w, http.StatusInternalServerError, types.APIErrorCodeInternal, err.Error())
98
return
99
}
100
177
- bootstraps, err = s.responseBootstraps(localResp.Bootstraps)
101
+ localResp.Bootstraps, err = buildResponseBootstraps(selfURLs, bootstraps, localResp.Bootstraps)
102
if err != nil {
103
utils.WriteAPIError(w, http.StatusInternalServerError, types.APIErrorCodeInternal, err.Error())
104
return
105
}
182
- localResp.Bootstraps = bootstraps
106
utils.WriteAPIData(w, http.StatusOK, localResp)
107
}
108
186
-func (s *Service) Poll(ctx context.Context, req types.DiscoverRequest) (types.DiscoverResponse, error) {
187
- if ctx == nil {
188
- ctx = context.Background()
189
- }
190
-
191
- req, err := normalizeRequest(req)
192
- if err != nil {
193
- return types.DiscoverResponse{}, err
194
- }
195
-
196
- discovered := s.Bootstraps()
197
- if len(discovered) == 0 {
198
- return types.DiscoverResponse{}, errors.New("at least one bootstrap is required")
199
- }
200
-
201
- queue := append([]string(nil), discovered...)
202
- seen := make(map[string]struct{}, len(queue))
203
- var lastErr error
204
- contacted := false
205
-
206
- for len(queue) > 0 && len(seen) < s.maxPeers {
207
- bootstrap := queue[0]
208
- queue = queue[1:]
209
- if _, ok := seen[bootstrap]; ok {
210
- continue
211
- }
212
- seen[bootstrap] = struct{}{}
213
-
214
- resp, err := discoverPeer(ctx, bootstrap, req, s.rootCAPEM, s.requestTimeout)
215
- if err != nil {
216
- lastErr = err
217
- continue
218
- }
219
- contacted = true
220
-
221
- if err := s.MergeBootstraps(resp.Bootstraps); err != nil {
222
- return types.DiscoverResponse{}, err
223
- }
224
- discovered = s.Bootstraps()
225
- for _, nextBootstrap := range discovered {
226
- if _, ok := seen[nextBootstrap]; !ok {
227
- queue = append(queue, nextBootstrap)
228
- }
229
- }
230
-
231
- if resp.Found {
232
- resp.Bootstraps = discovered
233
- return resp, nil
234
- }
235
- }
236
-
237
- if !contacted && lastErr != nil {
238
- return types.DiscoverResponse{}, lastErr
239
- }
240
-
241
- return types.DiscoverResponse{
242
- Found: false,
243
- Bootstraps: discovered,
244
- }, nil
245
-}
246
-
247
-func (s *Service) RunPollLoop(ctx context.Context, interval time.Duration, req types.DiscoverRequest) error {
248
- if s == nil {
249
- return nil
250
- }
251
-
252
- interval = utils.DurationOrDefault(interval, defaultPollInterval)
253
- lastPollErr := ""
254
- for {
255
- if len(s.Bootstraps()) > 0 {
256
- if _, err := s.Poll(ctx, req); err != nil {
257
- if ctx.Err() != nil {
258
- return nil
259
- }
260
- errText := err.Error()
261
- if errText != lastPollErr {
262
- log.Warn().
263
- Err(err).
264
- Int("bootstrap_count", len(s.Bootstraps())).
265
- Str("root_host", req.RootHost).
266
- Str("name", req.Name).
267
- Msg("discovery poll failed")
268
- lastPollErr = errText
269
- }
270
- } else if lastPollErr != "" {
271
- log.Info().
272
- Int("bootstrap_count", len(s.Bootstraps())).
273
- Str("root_host", req.RootHost).
274
- Str("name", req.Name).
275
- Msg("discovery poll recovered")
276
- lastPollErr = ""
277
- }
278
- }
279
- if !utils.SleepOrDone(ctx, interval) {
280
- return nil
281
- }
282
- }
283
-}
284
-
109
func normalizeRequest(req types.DiscoverRequest) (types.DiscoverRequest, error) {
110
req.RootHost = utils.NormalizeHostname(req.RootHost)
111
req.Name = strings.TrimSpace(req.Name)
@@ -299,80 +123,87 @@ func normalizeRequest(req types.DiscoverRequest) (types.DiscoverRequest, error)
123
return req, nil
124
}
125
302
-func (s *Service) responseBootstraps(extra []string) ([]string, error) {
303
- base := []string(nil)
304
- if s != nil {
305
- s.mu.RLock()
306
- base = append(base, s.selfURLs...)
307
- base = append(base, s.bootstraps...)
308
- s.mu.RUnlock()
309
- }
310
- bootstraps, err := utils.NormalizeRelayURLs(append(base, extra...))
126
+func discoverPeer(ctx context.Context, relayURL string, req types.DiscoverRequest, rootCAPEM []byte) (types.DiscoverResponse, error) {
127
+ relayURL, err := utils.NormalizeRelayURL(relayURL)
128
if err != nil {
312
- return nil, fmt.Errorf("normalize bootstraps: %w", err)
129
+ return types.DiscoverResponse{}, err
130
}
314
- return bootstraps, nil
315
-}
131
317
-func discoverPeer(ctx context.Context, bootstrap string, req types.DiscoverRequest, rootCAPEM []byte, requestTimeout time.Duration) (types.DiscoverResponse, error) {
318
- baseURL, err := url.Parse(bootstrap)
132
+ baseURL, err := url.Parse(relayURL)
133
if err != nil {
320
- return types.DiscoverResponse{}, fmt.Errorf("parse bootstrap url: %w", err)
134
+ return types.DiscoverResponse{}, fmt.Errorf("parse relay url: %w", err)
135
}
136
323
- rootCAs, err := keyless.RelayRootCAs(ctx, bootstrap, baseURL.Hostname(), rootCAPEM)
137
+ rootCAs, err := keyless.RelayRootCAs(ctx, relayURL, baseURL.Hostname(), rootCAPEM)
138
if err != nil {
139
return types.DiscoverResponse{}, err
140
}
141
142
+ ref, _ := url.Parse(types.PathDiscovery)
143
+ discoverURL := baseURL.ResolveReference(ref)
144
+ query := discoverURL.Query()
145
+ if req.RootHost != "" {
146
+ query.Set("root_host", req.RootHost)
147
+ }
148
+ if req.Name != "" {
149
+ query.Set("name", req.Name)
150
+ }
151
+ discoverURL.RawQuery = query.Encode()
152
+
153
httpClient := &http.Client{
154
Transport: &http.Transport{
155
TLSClientConfig: &tls.Config{
156
MinVersion: tls.VersionTLS12,
157
ServerName: baseURL.Hostname(),
158
RootCAs: rootCAs,
159
+ NextProtos: []string{"http/1.1"},
160
},
161
ForceAttemptHTTP2: false,
162
},
337
- Timeout: utils.DurationOrDefault(requestTimeout, defaultRequestTimeout),
338
- }
339
- defer httpClient.CloseIdleConnections()
340
-
341
- query := url.Values{}
342
- if req.RootHost != "" {
343
- query.Set("root_host", req.RootHost)
344
- }
345
- if req.Name != "" {
346
- query.Set("name", req.Name)
163
+ Timeout: defaultRequestTimeout,
164
}
165
349
- ref := &url.URL{Path: types.PathDiscovery, RawQuery: query.Encode()}
350
- httpReq, err := http.NewRequestWithContext(ctx, http.MethodGet, baseURL.ResolveReference(ref).String(), nil)
166
+ httpReq, err := http.NewRequestWithContext(ctx, http.MethodGet, discoverURL.String(), nil)
167
if err != nil {
168
return types.DiscoverResponse{}, err
169
}
170
355
- httpResp, err := httpClient.Do(httpReq)
171
+ resp, err := httpClient.Do(httpReq)
172
if err != nil {
173
return types.DiscoverResponse{}, err
174
}
359
- defer httpResp.Body.Close()
175
+ defer resp.Body.Close()
176
361
- if httpResp.StatusCode < http.StatusOK || httpResp.StatusCode >= http.StatusMultipleChoices {
362
- return types.DiscoverResponse{}, utils.DecodeAPIRequestError(httpResp)
177
+ if resp.StatusCode < http.StatusOK || resp.StatusCode >= http.StatusMultipleChoices {
178
+ return types.DiscoverResponse{}, utils.DecodeAPIRequestError(resp)
179
}
180
365
- envelope, err := utils.DecodeAPIEnvelope[json.RawMessage](httpResp.Body)
181
+ envelope, err := utils.DecodeAPIEnvelope[types.DiscoverResponse](resp.Body)
182
if err != nil {
183
return types.DiscoverResponse{}, fmt.Errorf("decode response: %w", err)
184
}
185
if !envelope.OK {
370
- return types.DiscoverResponse{}, utils.NewAPIRequestError(httpResp.StatusCode, envelope.Error)
186
+ return types.DiscoverResponse{}, utils.NewAPIRequestError(resp.StatusCode, envelope.Error)
187
}
188
+ return envelope.Data, nil
189
+}
190
373
- var resp types.DiscoverResponse
374
- if err := json.Unmarshal(envelope.Data, &resp); err != nil {
375
- return types.DiscoverResponse{}, err
191
+func buildResponseBootstraps(selfURLs, bootstraps, extra []string) ([]string, error) {
192
+ merged, err := utils.MergeRelayURLs(bootstraps, selfURLs, extra)
193
+ if err != nil {
194
+ return nil, err
195
+ }
196
+ if len(selfURLs) == 0 {
197
+ return merged, nil
198
+ }
199
+
200
+ normalizedSelf, err := utils.NormalizeRelayURLs(selfURLs)
201
+ if err != nil {
202
+ return nil, fmt.Errorf("normalize self urls: %w", err)
203
+ }
204
+ resolvedBootstraps, err := utils.NormalizeRelayURLs(append(normalizedSelf, merged...))
205
+ if err != nil {
206
+ return nil, fmt.Errorf("normalize bootstraps: %w", err)
207
}
377
- return resp, nil
208
+ return resolvedBootstraps, nil
209
}
portal/server.go
+81
-31
@@ -29,16 +29,12 @@ import (
29
const (
30
defaultLeaseTTL = 30 * time.Second
31
defaultClaimTimeout = 10 * time.Second
32
+ defaultDiscoveryInterval = 30 * time.Second
33
defaultIdleKeepalive = 15 * time.Second
34
defaultReadyQueueLimit = 8
35
defaultClientHelloWait = 2 * time.Second
36
defaultControlBodyLimit = 4 << 20
36
- defaultSessionWriteLimit = 5 * time.Second
37
- defaultQUICSNIRouteIdle = 30 * time.Second
38
- defaultQUICSNICleanup = 5 * time.Second
39
-
40
- defaultUDPPortBase = 50000
41
- defaultUDPPortCount = 0
37
+ defaultUDPPortBase = 50000
38
)
39
40
type ServerConfig struct {
@@ -61,21 +57,22 @@ type ServerConfig struct {
57
}
58
59
type Server struct {
64
- sniListener net.Listener
65
- apiListener net.Listener
66
- apiServer *http.Server
67
- apiTLSClose io.Closer
68
- acmeManager *acme.Manager
69
- quicTunnel *quic.Listener
70
- cancel context.CancelFunc
71
- group *errgroup.Group
72
- discovery *discovery.Service
73
- registry *leaseRegistry
74
- ports *transport.PortAllocator
75
- ownerIdentity discovery.Identity
76
- cfg ServerConfig
77
- rootHost string
78
- shutdownOnce sync.Once
60
+ sniListener net.Listener
61
+ apiListener net.Listener
62
+ apiServer *http.Server
63
+ apiTLSClose io.Closer
64
+ acmeManager *acme.Manager
65
+ quicTunnel *quic.Listener
66
+ cancel context.CancelFunc
67
+ group *errgroup.Group
68
+ registry *leaseRegistry
69
+ ports *transport.PortAllocator
70
+ ownerIdentity discovery.Identity
71
+ cfg ServerConfig
72
+ rootHost string
73
+ discoveryMu sync.RWMutex
74
+ discoveryBootstraps []string
75
+ shutdownOnce sync.Once
76
}
77
78
func NewServer(cfg ServerConfig) (*Server, error) {
@@ -138,14 +135,11 @@ func NewServer(cfg ServerConfig) (*Server, error) {
135
}
136
137
if cfg.DiscoveryEnabled {
141
- service, err := discovery.New(discovery.Config{
142
- SelfURLs: []string{cfg.PortalURL},
143
- Bootstraps: cfg.Bootstraps,
144
- }, discovery.ResolverFunc(s.discover))
138
+ bootstraps, err := utils.MergeRelayURLs(nil, []string{cfg.PortalURL}, cfg.Bootstraps)
139
if err != nil {
140
return nil, err
141
}
148
- s.discovery = service
142
+ s.discoveryBootstraps = bootstraps
143
}
144
145
return s, nil
@@ -198,10 +192,8 @@ func (s *Server) Start(ctx context.Context, apiMux *http.ServeMux) error {
192
group.Go(s.runAPIServer)
193
group.Go(func() error { return s.runSNIListener(groupCtx) })
194
group.Go(func() error { return s.registry.RunJanitor(groupCtx, 5*time.Second) })
201
- if s.discovery != nil {
202
- group.Go(func() error {
203
- return s.discovery.RunPollLoop(groupCtx, 0, types.DiscoverRequest{RootHost: s.rootHost})
204
- })
195
+ if s.DiscoveryEnabled() {
196
+ group.Go(func() error { return s.runDiscoveryLoop(groupCtx) })
197
}
198
group.Go(func() error { return s.watchContext(groupCtx) })
199
s.acmeManager.Start(serverCtx)
@@ -285,7 +277,7 @@ func (s *Server) QUICTunnelAddr() string {
277
}
278
279
func (s *Server) DiscoveryEnabled() bool {
288
- return s != nil && s.cfg.DiscoveryEnabled && s.discovery != nil
280
+ return s != nil && s.cfg.DiscoveryEnabled
281
}
282
283
func (s *Server) OwnerIdentity() discovery.Identity {
@@ -542,6 +534,64 @@ func (s *Server) watchContext(ctx context.Context) error {
534
return s.Shutdown(shutdownCtx)
535
}
536
537
+func (s *Server) discoveryBootstrapsSnapshot() []string {
538
+ if s == nil || !s.DiscoveryEnabled() {
539
+ return nil
540
+ }
541
+
542
+ s.discoveryMu.RLock()
543
+ defer s.discoveryMu.RUnlock()
544
+ return append([]string(nil), s.discoveryBootstraps...)
545
+}
546
+
547
+func (s *Server) mergeDiscoveryBootstraps(inputs []string) error {
548
+ if s == nil || !s.DiscoveryEnabled() || len(inputs) == 0 {
549
+ return nil
550
+ }
551
+
552
+ s.discoveryMu.Lock()
553
+ next, err := utils.MergeRelayURLs(s.discoveryBootstraps, []string{s.cfg.PortalURL}, inputs)
554
+ if err == nil {
555
+ s.discoveryBootstraps = next
556
+ }
557
+ s.discoveryMu.Unlock()
558
+ return err
559
+}
560
+
561
+func (s *Server) runDiscoveryLoop(ctx context.Context) error {
562
+ ticker := time.NewTicker(defaultDiscoveryInterval)
563
+ defer ticker.Stop()
564
+
565
+ for {
566
+ peers := s.discoveryBootstrapsSnapshot()
567
+ if len(peers) > 0 {
568
+ bootstraps, err := discovery.DiscoverBootstraps(ctx, peers, types.DiscoverRequest{}, nil)
569
+ switch {
570
+ case err == nil:
571
+ if err := s.mergeDiscoveryBootstraps(bootstraps); err != nil {
572
+ log.Warn().
573
+ Err(err).
574
+ Int("bootstrap_count", len(peers)).
575
+ Msg("merge discovered bootstraps failed")
576
+ }
577
+ case ctx.Err() != nil:
578
+ return nil
579
+ default:
580
+ log.Warn().
581
+ Err(err).
582
+ Int("bootstrap_count", len(peers)).
583
+ Msg("discover bootstraps failed")
584
+ }
585
+ }
586
+
587
+ select {
588
+ case <-ctx.Done():
589
+ return nil
590
+ case <-ticker.C:
591
+ }
592
+ }
593
+}
594
+
595
func BridgeConns(left, right net.Conn) {
596
defer left.Close()
597
defer right.Close()
sdk/api_client.go
-5
@@ -45,7 +45,6 @@ type apiClient struct {
45
rootCAPEM []byte
46
name string
47
reverseToken string
48
- discovery bool
48
metadata types.LeaseMetadata
49
ownerAddress string
50
}
@@ -88,7 +87,6 @@ func newApiClient(relayURL string, cfg ListenerConfig) (*apiClient, error) {
87
rootCAPEM: append([]byte(nil), cfg.RootCAPEM...),
88
name: name,
89
reverseToken: reverseToken,
91
- discovery: cfg.Discovery,
90
metadata: cfg.Metadata.Copy(),
91
ownerAddress: ownerAddress,
92
}, nil
@@ -105,9 +103,6 @@ func (a *apiClient) close() {
103
104
func (a *apiClient) registerLease(ctx context.Context, ttl time.Duration, udpEnabled bool, bootstraps []string) (types.RegisterResponse, error) {
105
var resp types.RegisterResponse
108
- if !a.discovery {
109
- bootstraps = nil
110
- }
106
if err := a.doJSON(ctx, http.MethodPost, types.PathSDKRegister, types.RegisterRequest{
107
Name: a.name,
108
Metadata: a.metadata.Copy(),
sdk/expose.go
+123
-80
@@ -2,6 +2,7 @@ package sdk
2
3
import (
4
"context"
5
+ "encoding/json"
6
"errors"
7
"fmt"
8
"net"
@@ -32,7 +33,6 @@ type Exposure struct {
33
ownerAddress string
34
rootCAPEM []byte
35
discoveryEnabled bool
35
- discovery *discovery.Service
36
37
accepted chan net.Conn
38
datagrams chan types.DatagramFrame
@@ -47,39 +47,32 @@ type Exposure struct {
47
}
48
49
type ExposeConfig struct {
50
- RelayURLs []string
51
- Name string
52
- ReverseToken string
53
- UDPEnabled bool
54
- Discovery bool
55
- Metadata types.LeaseMetadata
56
- OwnerAddress string
57
- OwnerPrivateKey *string
58
- RootCAPEM []byte
50
+ RelayURLs []string
51
+ DefaultRelayEnabled bool
52
+ Name string
53
+ ReverseToken string
54
+ UDPEnabled bool
55
+ Discovery bool
56
+ Metadata types.LeaseMetadata
57
+ OwnerAddress string
58
+ OwnerPrivateKey *string
59
+ RootCAPEM []byte
60
}
61
62
// Expose creates relay listeners for each normalized relay URL and exposes a
63
// dynamic listener hub for accepting traffic from all of them.
63
-func Expose(ctx context.Context, relayURLs []string, name string, udpEnabled bool, metadata types.LeaseMetadata) (*Exposure, error) {
64
- return ExposeWithConfig(ctx, ExposeConfig{
65
- RelayURLs: relayURLs,
66
- Name: name,
67
- UDPEnabled: udpEnabled,
68
- Metadata: metadata,
69
- })
70
-}
64
+func Expose(ctx context.Context, cfg ExposeConfig) (*Exposure, error) {
65
+ if ctx == nil {
66
+ ctx = context.Background()
67
+ }
68
72
-func ExposeWithConfig(ctx context.Context, cfg ExposeConfig) (*Exposure, error) {
73
- relayURLs, err := utils.NormalizeRelayURLs(cfg.RelayURLs)
69
+ relayURLs, err := ResolveRelayURLs(ctx, cfg.RelayURLs, cfg.DefaultRelayEnabled)
70
if err != nil {
71
return nil, err
72
}
73
if len(relayURLs) == 0 {
74
return nil, nil
75
}
80
- if ctx == nil {
81
- ctx = context.Background()
82
- }
76
77
ownerAddress := strings.TrimSpace(cfg.OwnerAddress)
78
identity := discovery.Identity{}
@@ -109,46 +102,19 @@ func ExposeWithConfig(ctx context.Context, cfg ExposeConfig) (*Exposure, error)
102
starting: make(map[string]struct{}, len(relayURLs)),
103
}
104
112
- if exposure.discoveryEnabled {
113
- service, err := discovery.New(discovery.Config{
114
- RootCAPEM: exposure.rootCAPEM,
115
- OnBootstraps: func(relays []string) {
116
- if err := exposure.applyRelayURLs(relays, false); err != nil {
117
- log.Warn().Err(err).Strs("relays", relays).Msg("apply discovered relay urls")
118
- }
119
- },
120
- }, nil)
121
- if err != nil {
122
- cancel()
123
- return nil, err
124
- }
125
- exposure.discovery = service
126
- }
127
-
128
- if exposure.discovery != nil {
129
- if err := exposure.discovery.MergeBootstraps(relayURLs); err != nil {
130
- _ = exposure.Close()
131
- return nil, err
132
- }
133
- if err := exposure.applyRelayURLs(exposure.discovery.Bootstraps(), true); err != nil {
134
- _ = exposure.Close()
135
- return nil, err
136
- }
137
- } else if err := exposure.applyRelayURLs(relayURLs, true); err != nil {
105
+ if err := exposure.applyRelayURLs(relayURLs, true); err != nil {
106
_ = exposure.Close()
107
return nil, err
108
}
109
110
go exposure.monitorStartupCounts()
111
+ if exposure.discoveryEnabled {
112
+ go exposure.runDiscoveryLoop(exposureCtx)
113
+ }
114
go func() {
115
<-exposure.done
116
_ = exposure.Close()
117
}()
147
- if exposure.discovery != nil {
148
- go func() {
149
- _ = exposure.discovery.RunPollLoop(exposureCtx, 0, types.DiscoverRequest{})
150
- }()
151
- }
118
119
log.Info().
120
Str("release_version", types.ReleaseVersion).
@@ -159,6 +125,50 @@ func ExposeWithConfig(ctx context.Context, cfg ExposeConfig) (*Exposure, error)
125
return exposure, nil
126
}
127
128
+const defaultDiscoveryInterval = 30 * time.Second
129
+
130
+func ResolveRelayURLs(ctx context.Context, explicit []string, includeDefaults bool) ([]string, error) {
131
+ explicit, err := utils.NormalizeRelayURLs(explicit)
132
+ if err != nil {
133
+ return nil, err
134
+ }
135
+ if !includeDefaults {
136
+ return explicit, nil
137
+ }
138
+
139
+ req, err := http.NewRequestWithContext(ctx, http.MethodGet, types.PortalRelayRegistryURL, nil)
140
+ if err != nil {
141
+ return explicit, nil
142
+ }
143
+
144
+ client := &http.Client{Timeout: defaultRequestTimeout}
145
+ resp, err := client.Do(req)
146
+ if err != nil {
147
+ return explicit, nil
148
+ }
149
+ defer resp.Body.Close()
150
+
151
+ if resp.StatusCode != http.StatusOK {
152
+ return explicit, nil
153
+ }
154
+
155
+ var registry struct {
156
+ Relays []string `json:"relays"`
157
+ }
158
+ if err := json.NewDecoder(resp.Body).Decode(®istry); err != nil {
159
+ return explicit, nil
160
+ }
161
+
162
+ defaults, err := utils.NormalizeRelayURLs(registry.Relays)
163
+ if err != nil {
164
+ return explicit, nil
165
+ }
166
+ if len(defaults) == 0 {
167
+ return explicit, nil
168
+ }
169
+ return utils.MergeRelayURLs(defaults, nil, explicit)
170
+}
171
+
172
func (e *Exposure) RelayURLs() []string {
173
if e == nil {
174
return nil
@@ -440,15 +450,19 @@ func (e *Exposure) reserveMissingRelayURLs() []string {
450
}
451
452
func (e *Exposure) newListener(relayURL string) (*Listener, error) {
453
+ bootstraps := []string(nil)
454
+ if e.discoveryEnabled {
455
+ bootstraps = e.RelayURLs()
456
+ }
457
+
458
cfg := ListenerConfig{
444
- Name: e.name,
445
- ReverseToken: e.reverseToken,
446
- UDPEnabled: e.udpEnabled,
447
- Discovery: e.discoveryEnabled,
448
- OwnerAddress: e.ownerAddress,
449
- Metadata: e.metadata.Copy(),
450
- RootCAPEM: append([]byte(nil), e.rootCAPEM...),
451
- bootstrapService: e.discovery,
459
+ Name: e.name,
460
+ ReverseToken: e.reverseToken,
461
+ UDPEnabled: e.udpEnabled,
462
+ OwnerAddress: e.ownerAddress,
463
+ RegisterBootstraps: bootstraps,
464
+ Metadata: e.metadata.Copy(),
465
+ RootCAPEM: append([]byte(nil), e.rootCAPEM...),
466
}
467
return NewListener(context.Background(), relayURL, cfg)
468
}
@@ -499,28 +513,19 @@ func (e *Exposure) listenersOrdered() []*Listener {
513
return out
514
}
515
502
-func (e *Exposure) listenerForRelayURL(relayURL string) *Listener {
503
- if e == nil {
504
- return nil
505
- }
506
-
507
- e.mu.RLock()
508
- defer e.mu.RUnlock()
509
- return e.listeners[relayURL]
510
-}
511
-
516
func (e *Exposure) runListenerAcceptLoop(listener *Listener) {
517
if e == nil || listener == nil {
518
return
519
}
520
521
+ relayURL := listener.api.baseURL.String()
522
for {
523
conn, err := listener.Accept()
524
if err != nil {
525
if listener.closed() || errors.Is(err, net.ErrClosed) {
526
return
527
}
523
- log.Warn().Err(err).Str("relay_url", listener.relayURL).Msg("exposure listener accept failed")
528
+ log.Warn().Err(err).Str("relay_url", relayURL).Msg("exposure listener accept failed")
529
return
530
}
531
@@ -596,7 +601,9 @@ func (e *Exposure) SendDatagram(frame types.DatagramFrame) error {
601
return errors.New("relay url is required")
602
}
603
599
- listener := e.listenerForRelayURL(relayURL)
604
+ e.mu.RLock()
605
+ listener := e.listeners[relayURL]
606
+ e.mu.RUnlock()
607
if listener == nil {
608
return net.ErrClosed
609
}
@@ -719,6 +726,7 @@ func (e *Exposure) allDatagramNegotiationsResolvedWithoutDatagram() bool {
726
}
727
728
func (e *Exposure) attachDatagramPlane(ctx context.Context, listener *Listener) {
729
+ relayURL := listener.api.baseURL.String()
730
err := listener.WaitDatagramReady(ctx)
731
if err != nil {
732
switch {
@@ -731,13 +739,13 @@ func (e *Exposure) attachDatagramPlane(ctx context.Context, listener *Listener)
739
default:
740
log.Warn().
741
Err(err).
734
- Str("relay_url", listener.relayURL).
742
+ Str("relay_url", relayURL).
743
Msg("attach datagram plane failed")
744
return
745
}
746
}
747
740
- e.forwardDatagrams(listener.relayURL, listener)
748
+ e.forwardDatagrams(relayURL, listener)
749
}
750
751
func (e *Exposure) forwardDatagrams(relayURL string, listener *Listener) {
@@ -799,6 +807,7 @@ func (e *Exposure) monitorStartupCounts() {
807
if listener == nil {
808
continue
809
}
810
+ relayURL := listener.api.baseURL.String()
811
812
status := listener.StartupStatus()
813
if status == listenerStatusReady {
@@ -807,14 +816,14 @@ func (e *Exposure) monitorStartupCounts() {
816
inactiveCount++
817
}
818
810
- if prev, ok := prevStatuses[listener.relayURL]; ok && prev != status {
819
+ if prev, ok := prevStatuses[relayURL]; ok && prev != status {
820
if status == listenerStatusReady {
812
- activated = append(activated, listener.relayURL)
821
+ activated = append(activated, relayURL)
822
} else {
814
- deactivated = append(deactivated, listener.relayURL)
823
+ deactivated = append(deactivated, relayURL)
824
}
825
}
817
- prevStatuses[listener.relayURL] = status
826
+ prevStatuses[relayURL] = status
827
}
828
829
if firstRun || len(activated) > 0 || len(deactivated) > 0 {
@@ -838,3 +847,37 @@ func (e *Exposure) monitorStartupCounts() {
847
}
848
}
849
}
850
+
851
+func (e *Exposure) runDiscoveryLoop(ctx context.Context) {
852
+ ticker := time.NewTicker(defaultDiscoveryInterval)
853
+ defer ticker.Stop()
854
+
855
+ for {
856
+ peers := e.RelayURLs()
857
+ if len(peers) > 0 {
858
+ relayURLs, err := discovery.DiscoverBootstraps(ctx, peers, types.DiscoverRequest{}, e.rootCAPEM)
859
+ switch {
860
+ case err == nil:
861
+ if err := e.applyRelayURLs(relayURLs, false); err != nil {
862
+ log.Warn().
863
+ Err(err).
864
+ Int("relay_count", len(peers)).
865
+ Msg("apply discovered relay urls failed")
866
+ }
867
+ case ctx.Err() != nil:
868
+ return
869
+ default:
870
+ log.Warn().
871
+ Err(err).
872
+ Int("relay_count", len(peers)).
873
+ Msg("discover relay urls failed")
874
+ }
875
+ }
876
+
877
+ select {
878
+ case <-ctx.Done():
879
+ return
880
+ case <-ticker.C:
881
+ }
882
+ }
883
+}
sdk/listener.go
+56
-72
@@ -4,6 +4,7 @@ import (
4
"context"
5
"crypto/tls"
6
"errors"
7
+ "fmt"
8
"io"
9
"net"
10
"net/url"
@@ -14,7 +15,6 @@ import (
15
"github.com/quic-go/quic-go"
16
"github.com/rs/zerolog/log"
17
17
- "github.com/gosuda/portal/v2/portal/discovery"
18
"github.com/gosuda/portal/v2/portal/keyless"
19
"github.com/gosuda/portal/v2/portal/transport"
20
"github.com/gosuda/portal/v2/types"
@@ -25,7 +25,6 @@ type ListenerConfig struct {
25
Name string
26
ReverseToken string
27
UDPEnabled bool
28
- Discovery bool
28
OwnerAddress string
29
Metadata types.LeaseMetadata
30
RootCAPEM []byte
@@ -38,7 +37,7 @@ type ListenerConfig struct {
37
RetryCount int
38
RetryWait time.Duration
39
41
- bootstrapService *discovery.Service
40
+ RegisterBootstraps []string
41
}
42
43
type listenerStatus string
@@ -49,31 +48,30 @@ const (
48
)
49
50
type Listener struct {
52
- tlsCloser io.Closer
53
- tlsConfig *tls.Config
54
- readyTarget int
55
- retryCount int
56
- retryWait time.Duration
57
- leaseTTL time.Duration
58
- renewBefore time.Duration
59
- doneCh <-chan struct{}
60
- cancel context.CancelFunc
61
- api *apiClient
62
- relayURL string
63
- bootstrapSvc *discovery.Service
51
+ api *apiClient
52
+ cancel context.CancelFunc
53
+ doneCh <-chan struct{}
54
+
55
+ retryCount int
56
+ retryWait time.Duration
57
+ leaseTTL time.Duration
58
+ renewBefore time.Duration
59
+
60
+ stream *transport.ClientStream
61
+ datagram *transport.ClientDatagram
62
+
63
+ registered chan struct{}
64
+ closeOnce sync.Once
65
+ registerOnce sync.Once
66
+
67
+ mu sync.Mutex
68
startupStatus listenerStatus
69
leaseID string
70
hostname string
71
udpAddr string
68
- udpEnabled bool
72
metadata types.LeaseMetadata
70
- stream *transport.ClientStream
71
- datagram *transport.ClientDatagram
72
-
73
- registered chan struct{}
74
- closeOnce sync.Once
75
- registerOnce sync.Once
76
- mu sync.Mutex
73
+ tlsConfig *tls.Config
74
+ tlsCloser io.Closer
75
}
76
77
// NewListener creates one relay listener and its dedicated relay transport for one relay URL.
@@ -96,20 +94,22 @@ func NewListener(ctx context.Context, relayURL string, cfg ListenerConfig) (*Lis
94
return nil, err
95
}
96
97
+ initialBootstraps, err := utils.NormalizeRelayURLs(cfg.RegisterBootstraps)
98
+ if err != nil {
99
+ cancel()
100
+ return nil, fmt.Errorf("normalize bootstraps: %w", err)
101
+ }
102
+
103
l := &Listener{
104
doneCh: listenerCtx.Done(),
105
cancel: cancel,
106
api: api,
107
registered: make(chan struct{}),
104
- relayURL: api.baseURL.String(),
105
- bootstrapSvc: cfg.bootstrapService,
108
startupStatus: listenerStatusInactive,
107
- readyTarget: readyTarget,
109
retryCount: cfg.RetryCount,
110
retryWait: retryWait,
111
leaseTTL: leaseTTL,
112
renewBefore: renewBefore,
112
- udpEnabled: cfg.UDPEnabled,
113
metadata: cfg.Metadata.Copy(),
114
}
115
l.stream = transport.NewClientStream(readyTarget, handshakeTimeout)
@@ -129,18 +129,18 @@ func NewListener(ctx context.Context, relayURL string, cfg ListenerConfig) (*Lis
129
})
130
}
131
132
- go l.runStartup(listenerCtx)
132
+ go l.runStartup(listenerCtx, initialBootstraps, readyTarget)
133
return l, nil
134
}
135
136
-func (l *Listener) runStartup(ctx context.Context) {
136
+func (l *Listener) runStartup(ctx context.Context, initialBootstraps []string, readyTarget int) {
137
var retries int
138
139
for {
140
- err := l.registerAndConfigure(ctx)
140
+ err := l.registerAndConfigure(ctx, initialBootstraps)
141
switch {
142
case err == nil:
143
- for i := 0; i < l.readyTarget; i++ {
143
+ for range readyTarget {
144
go l.stream.RunLoop(
145
ctx,
146
func(ctx context.Context) (net.Conn, error) {
@@ -162,7 +162,7 @@ func (l *Listener) runStartup(ctx context.Context) {
162
go l.runRenewLoop(ctx)
163
publicURL := l.PublicURL()
164
event := log.Info().
165
- Str("relay_url", l.relayURL).
165
+ Str("relay_url", l.api.baseURL.String()).
166
Str("lease_id", l.LeaseID())
167
if publicURL != "" {
168
event = event.Str("public_url", publicURL)
@@ -172,10 +172,14 @@ func (l *Listener) runStartup(ctx context.Context) {
172
case errors.Is(err, context.Canceled), errors.Is(err, net.ErrClosed):
173
return
174
default:
175
- if isPermanentRegistrationError(err) {
175
+ if errors.Is(err, errRelayIncompatible) ||
176
+ errors.Is(err, &types.APIRequestError{Code: types.APIErrorCodeFeatureUnavailable}) ||
177
+ errors.Is(err, &types.APIRequestError{Code: types.APIErrorCodeTransportMismatch}) ||
178
+ errors.Is(err, &types.APIRequestError{Code: types.APIErrorCodeHostnameConflict}) ||
179
+ errors.Is(err, &types.APIRequestError{Code: types.APIErrorCodeIPBanned}) {
180
log.Error().
181
Err(err).
178
- Str("relay_url", l.relayURL).
182
+ Str("relay_url", l.api.baseURL.String()).
183
Str("lease_id", l.LeaseID()).
184
Msg("lease registration failed; closing listener")
185
_ = l.Close()
@@ -266,27 +270,29 @@ func (l *Listener) Metadata() types.LeaseMetadata {
270
}
271
272
func (l *Listener) PublicURL() string {
273
+ if l == nil || l.api == nil || l.api.baseURL == nil {
274
+ return ""
275
+ }
276
+
277
l.mu.Lock()
278
hostname := l.hostname
271
- relayURL := l.relayURL
279
l.mu.Unlock()
280
281
if hostname == "" {
282
return ""
283
}
284
278
- parsed, err := url.Parse(relayURL)
279
- if err != nil || strings.TrimSpace(parsed.Scheme) == "" {
285
+ if strings.TrimSpace(l.api.baseURL.Scheme) == "" {
286
return "https://" + hostname
287
}
288
289
host := hostname
284
- if port := strings.TrimSpace(parsed.Port()); port != "" {
290
+ if port := strings.TrimSpace(l.api.baseURL.Port()); port != "" {
291
host = net.JoinHostPort(hostname, port)
292
}
293
294
return (&url.URL{
289
- Scheme: parsed.Scheme,
295
+ Scheme: l.api.baseURL.Scheme,
296
Host: host,
297
}).String()
298
}
@@ -337,7 +343,7 @@ func (l *Listener) currentDatagramState() (transport.ClientDatagramState, bool)
343
}
344
345
func (l *Listener) WaitDatagramReady(ctx context.Context) error {
340
- if l == nil || !l.UDPEnabled() {
346
+ if l == nil || l.datagram == nil {
347
return errors.New("lease does not have udp enabled")
348
}
349
if err := l.WaitRegistered(ctx); err != nil {
@@ -369,7 +375,7 @@ func (l *Listener) WaitDatagramReady(ctx context.Context) error {
375
}
376
377
func (l *Listener) activeSupportsDatagram() bool {
372
- if l == nil || !l.udpEnabled {
378
+ if l == nil || l.datagram == nil {
379
return false
380
}
381
l.mu.Lock()
@@ -443,26 +449,24 @@ func (l *Listener) renewLease(ctx context.Context) error {
449
return err
450
}
451
446
- if err := l.reregister(ctx); err != nil {
452
+ requestCtx, cancel = context.WithTimeout(ctx, 10*time.Second)
453
+ defer cancel()
454
+ if err := l.registerAndConfigure(requestCtx, nil); err != nil {
455
return err
456
}
457
return nil
458
}
459
452
-func (l *Listener) registerAndConfigure(ctx context.Context) error {
460
+func (l *Listener) registerAndConfigure(ctx context.Context, registerBootstraps []string) error {
461
if err := l.api.ensureReady(ctx); err != nil {
462
return err
463
}
464
457
- bootstraps := []string(nil)
458
- if l.bootstrapSvc != nil {
459
- bootstraps = l.bootstrapSvc.Bootstraps()
460
- }
461
- resp, err := l.api.registerLease(ctx, l.leaseTTL, l.udpEnabled, bootstraps)
465
+ resp, err := l.api.registerLease(ctx, l.leaseTTL, l.datagram != nil, registerBootstraps)
466
if err != nil {
467
return err
468
}
465
- if l.udpEnabled && !resp.UDPEnabled {
469
+ if l.datagram != nil && !resp.UDPEnabled {
470
_ = l.api.unregisterLease(context.Background(), resp.LeaseID)
471
return &types.APIRequestError{
472
Code: types.APIErrorCodeFeatureUnavailable,
@@ -509,17 +513,12 @@ func (l *Listener) registerAndConfigure(ctx context.Context) error {
513
if datagram != nil {
514
datagram.Clear("lease updated")
515
}
512
- if l.bootstrapSvc != nil && len(resp.Bootstraps) > 0 {
513
- if err := l.bootstrapSvc.MergeBootstraps(resp.Bootstraps); err != nil {
514
- log.Warn().Err(err).Strs("bootstraps", resp.Bootstraps).Msg("learn bootstraps from register response")
515
- }
516
- }
516
l.registerOnce.Do(func() { close(l.registered) })
517
return nil
518
}
519
520
func (l *Listener) SupportsDatagram() bool {
522
- return l != nil && l.udpEnabled
521
+ return l != nil && l.datagram != nil
522
}
523
524
func (l *Listener) SupportsStream() bool {
@@ -527,7 +526,7 @@ func (l *Listener) SupportsStream() bool {
526
}
527
528
func (l *Listener) UDPEnabled() bool {
530
- return l != nil && l.udpEnabled
529
+ return l != nil && l.datagram != nil
530
}
531
532
// WaitRegistered blocks until the first successful lease registration or context cancellation.
@@ -542,28 +541,13 @@ func (l *Listener) WaitRegistered(ctx context.Context) error {
541
}
542
}
543
545
-func (l *Listener) reregister(ctx context.Context) error {
546
- requestCtx, cancel := context.WithTimeout(ctx, 10*time.Second)
547
- defer cancel()
548
-
549
- return l.registerAndConfigure(requestCtx)
550
-}
551
-
552
-func isPermanentRegistrationError(err error) bool {
553
- return errors.Is(err, errRelayIncompatible) ||
554
- errors.Is(err, &types.APIRequestError{Code: types.APIErrorCodeFeatureUnavailable}) ||
555
- errors.Is(err, &types.APIRequestError{Code: types.APIErrorCodeTransportMismatch}) ||
556
- errors.Is(err, &types.APIRequestError{Code: types.APIErrorCodeHostnameConflict}) ||
557
- errors.Is(err, &types.APIRequestError{Code: types.APIErrorCodeIPBanned})
558
-}
559
-
544
func (l *Listener) retryOrClose(ctx context.Context, operation string, err error, retries int) bool {
545
if ctx.Err() != nil {
546
return false
547
}
548
549
logger := log.With().
566
- Str("relay_url", l.relayURL).
550
+ Str("relay_url", l.api.baseURL.String()).
551
Str("operation", operation).
552
Str("lease_id", l.LeaseID()).
553
Logger()
sdk/registry.go
deleted
-52
@@ -1,52 +0,0 @@
1
-package sdk
2
-
3
-import (
4
- "context"
5
- "encoding/json"
6
- "net/http"
7
-
8
- "github.com/gosuda/portal/v2/types"
9
- "github.com/gosuda/portal/v2/utils"
10
-)
11
-
12
-// WithDefaultRelayURLs fetches the default Portal relay registry and appends
13
-// any explicit relay inputs before normalization.
14
-func WithDefaultRelayURLs(ctx context.Context, registryURL string, explicit ...string) []string {
15
- if registryURL == "" {
16
- registryURL = types.PortalRelayRegistryURL
17
- }
18
-
19
- req, err := http.NewRequestWithContext(ctx, http.MethodGet, registryURL, nil)
20
- if err != nil {
21
- return explicit
22
- }
23
-
24
- client := &http.Client{Timeout: defaultRequestTimeout}
25
- resp, err := client.Do(req)
26
- if err != nil {
27
- return explicit
28
- }
29
- defer resp.Body.Close()
30
-
31
- if resp.StatusCode != http.StatusOK {
32
- return explicit
33
- }
34
-
35
- var registry struct {
36
- Relays []string `json:"relays"`
37
- }
38
- if err := json.NewDecoder(resp.Body).Decode(®istry); err != nil {
39
- return explicit
40
- }
41
-
42
- relayURLs := append(registry.Relays, explicit...)
43
- if len(relayURLs) == 0 {
44
- return nil
45
- }
46
- relayURLs, err = utils.NormalizeRelayURLs(relayURLs)
47
- if err != nil {
48
- return nil
49
- }
50
-
51
- return relayURLs
52
-}
sdk/sdk_test.go
+5
-5
@@ -353,7 +353,7 @@ func TestNewListenerRetriesForeverWhenRetryCountIsNegative(t *testing.T) {
353
}
354
355
func TestExposeNoRelayInputs(t *testing.T) {
356
- exposure, err := Expose(context.Background(), nil, "demo", false, types.LeaseMetadata{})
356
+ exposure, err := Expose(context.Background(), ExposeConfig{Name: "demo"})
357
if err != nil {
358
t.Fatalf("Expose() error = %v", err)
359
}
@@ -413,14 +413,14 @@ func TestExposeRegistersKnownRelayURLs(t *testing.T) {
413
relayB := newRelayServer()
414
defer relayB.Close()
415
416
- exposure, err := ExposeWithConfig(context.Background(), ExposeConfig{
416
+ exposure, err := Expose(context.Background(), ExposeConfig{
417
RelayURLs: []string{relayA.URL, relayB.URL},
418
Discovery: true,
419
Name: "demo",
420
OwnerAddress: "0x52908400098527886E0F7030069857D2E4169EE7",
421
})
422
if err != nil {
423
- t.Fatalf("ExposeWithConfig() error = %v", err)
423
+ t.Fatalf("Expose() error = %v", err)
424
}
425
defer exposure.Close()
426
@@ -498,13 +498,13 @@ func TestExposeResolvesOwnerPrivateKey(t *testing.T) {
498
}))
499
defer server.Close()
500
501
- exposure, err := ExposeWithConfig(context.Background(), ExposeConfig{
501
+ exposure, err := Expose(context.Background(), ExposeConfig{
502
RelayURLs: []string{server.URL},
503
Name: "demo",
504
OwnerPrivateKey: &ownerPrivateKey,
505
})
506
if err != nil {
507
- t.Fatalf("ExposeWithConfig() error = %v", err)
507
+ t.Fatalf("Expose() error = %v", err)
508
}
509
defer exposure.Close()
510
utils/utils.go
+31
-25
@@ -143,10 +143,39 @@ func NormalizeRelayURLs(inputs []string) ([]string, error) {
143
}
144
}
145
146
- return UniqueURLs(out), nil
146
+ return uniqueURLs(out), nil
147
}
148
149
-func UniqueURLs(inputs []string) []string {
149
+func MergeRelayURLs(current, excluded, inputs []string) ([]string, error) {
150
+ merged, err := NormalizeRelayURLs(append(append([]string(nil), current...), inputs...))
151
+ if err != nil {
152
+ return nil, err
153
+ }
154
+ if len(excluded) == 0 {
155
+ return merged, nil
156
+ }
157
+
158
+ excluded, err = NormalizeRelayURLs(excluded)
159
+ if err != nil {
160
+ return nil, err
161
+ }
162
+
163
+ skip := make(map[string]struct{}, len(excluded))
164
+ for _, input := range excluded {
165
+ skip[input] = struct{}{}
166
+ }
167
+
168
+ filtered := make([]string, 0, len(merged))
169
+ for _, input := range merged {
170
+ if _, ok := skip[input]; ok {
171
+ continue
172
+ }
173
+ filtered = append(filtered, input)
174
+ }
175
+ return filtered, nil
176
+}
177
+
178
+func uniqueURLs(inputs []string) []string {
179
if len(inputs) == 0 {
180
return nil
181
}
@@ -170,29 +199,6 @@ func UniqueURLs(inputs []string) []string {
199
return out
200
}
201
173
-func ExcludeURLs(inputs []string, excluded []string) []string {
174
- if len(inputs) == 0 {
175
- return nil
176
- }
177
- if len(excluded) == 0 {
178
- return UniqueURLs(inputs)
179
- }
180
-
181
- skip := make(map[string]struct{}, len(excluded))
182
- for _, input := range UniqueURLs(excluded) {
183
- skip[input] = struct{}{}
184
- }
185
-
186
- filtered := make([]string, 0, len(inputs))
187
- for _, input := range inputs {
188
- if _, ok := skip[input]; ok {
189
- continue
190
- }
191
- filtered = append(filtered, input)
192
- }
193
- return UniqueURLs(filtered)
194
-}
195
-
202
func LeaseHostname(name, rootHost string) (string, error) {
203
label, err := NormalizeDNSLabel(name)
204
if err != nil {