feat: add overlay relay descriptor functionality and improve hop stream handling
Kim committed
Apr 17, 2026 at 14:30 UTC
29a421371589d283873c64dbf231c99c1e5713c9
5 files changed
+96
-32
portal/discovery/relayset.go
+17
@@ -247,6 +247,23 @@ func (s *RelaySet) OverlayPeerStates() []RelayState {
247
return out
248
}
249
250
+func (s *RelaySet) OverlayRelayDescriptor(relayURL string, now time.Time) (types.RelayDescriptor, bool) {
251
+ if now.IsZero() {
252
+ now = time.Now().UTC()
253
+ } else {
254
+ now = now.UTC()
255
+ }
256
+ relayURL = strings.TrimSpace(relayURL)
257
+
258
+ s.mu.RLock()
259
+ state := s.relays[relayURL]
260
+ s.mu.RUnlock()
261
+ if state.Banned || !state.hasObservedDescriptor() || !state.Descriptor.ExpiresAt.After(now) || !state.Descriptor.HasOverlayPeer() {
262
+ return types.RelayDescriptor{}, false
263
+ }
264
+ return state.Descriptor, true
265
+}
266
+
267
// BootstrapRelayURLs returns configured bootstrap discovery endpoints that
268
// can receive this relay's periodic self-announce.
269
func (s *RelaySet) BootstrapRelayURLs() []string {
portal/overlay/hop_mux.go
+20
-1
@@ -121,7 +121,7 @@ func (m *HopMux) OpenStream(ctx context.Context, overlayIPv4, token string) (net
121
if err != nil {
122
return nil, err
123
}
124
- stream, err := session.OpenStream()
124
+ stream, err := openYamuxStream(ctx, session)
125
if err != nil {
126
m.mu.Lock()
127
if m.outbound[overlayIPv4] == session {
@@ -150,6 +150,25 @@ func (m *HopMux) OpenStream(ctx context.Context, overlayIPv4, token string) (net
150
return stream, nil
151
}
152
153
+func openYamuxStream(ctx context.Context, session *yamux.Session) (*yamux.Stream, error) {
154
+ type openResult struct {
155
+ stream *yamux.Stream
156
+ err error
157
+ }
158
+ result := make(chan openResult, 1)
159
+ go func() {
160
+ stream, err := session.OpenStream()
161
+ result <- openResult{stream: stream, err: err}
162
+ }()
163
+ select {
164
+ case res := <-result:
165
+ return res.stream, res.err
166
+ case <-ctx.Done():
167
+ _ = session.Close()
168
+ return nil, ctx.Err()
169
+ }
170
+}
171
+
172
func (m *HopMux) session(ctx context.Context, overlayIPv4 string) (*yamux.Session, error) {
173
m.mu.Lock()
174
select {
portal/server.go
+44
-16
@@ -34,6 +34,7 @@ const (
34
defaultReadyQueueLimit = 8
35
defaultClientHelloWait = 2 * time.Second
36
defaultControlBodyLimit = 4 << 20
37
+ defaultHopOpenRetryWait = 250 * time.Millisecond
38
)
39
40
type ServerConfig struct {
@@ -484,7 +485,12 @@ func (s *Server) runSNIListener(ctx context.Context) error {
485
}
486
487
record, ok := s.registry.Lookup(serverName)
487
- if !ok || !s.bridgeLeaseConn(ctx, wrappedConn, record) {
488
+ if !ok {
489
+ _ = wrappedConn.Close()
490
+ return
491
+ }
492
+ if err := s.bridgeLeaseConn(ctx, wrappedConn, record); err != nil {
493
+ log.Warn().Err(err).Str("server_name", serverName).Msg("bridge lease connection")
494
_ = wrappedConn.Close()
495
return
496
}
@@ -525,8 +531,8 @@ func (s *Server) runHopMux(ctx context.Context) error {
531
return
532
}
533
log.Info().Str("remote_addr", stream.RemoteAddr).Bool("forward", record.isHopForward()).Msg("hop stream received")
528
- if !s.bridgeLeaseConn(groupCtx, stream.Conn, record) {
529
- log.Warn().Str("remote_addr", stream.RemoteAddr).Msg("hop stream bridge failed")
534
+ if err := s.bridgeLeaseConn(groupCtx, stream.Conn, record); err != nil {
535
+ log.Warn().Err(err).Str("remote_addr", stream.RemoteAddr).Msg("hop stream bridge failed")
536
_ = stream.Conn.Close()
537
}
538
}(stream)
@@ -535,35 +541,57 @@ func (s *Server) runHopMux(ctx context.Context) error {
541
return group.Wait()
542
}
543
538
-func (s *Server) bridgeLeaseConn(ctx context.Context, conn net.Conn, record *leaseRecord) bool {
544
+func (s *Server) bridgeLeaseConn(ctx context.Context, conn net.Conn, record *leaseRecord) error {
545
if s == nil || s.registry == nil || record == nil || time.Now().After(record.ExpiresAt) {
540
- return false
546
+ return errLeaseNotFound
547
}
548
if record.isHopForward() {
549
overlayIPv4 := strings.TrimSpace(record.hopNextOverlayIPv4)
550
forwardToken := strings.TrimSpace(record.hopNextToken)
545
- if s.hopMux == nil || overlayIPv4 == "" || forwardToken == "" {
546
- return false
551
+ switch {
552
+ case s.hopMux == nil:
553
+ return errFeatureUnavailable
554
+ case overlayIPv4 == "":
555
+ return errors.New("next hop overlay ipv4 is required")
556
+ case forwardToken == "":
557
+ return errors.New("next hop token is required")
558
}
548
- next, err := s.hopMux.OpenStream(ctx, overlayIPv4, forwardToken)
549
- if err != nil {
550
- log.Warn().Err(err).Str("next_overlay_ipv4", overlayIPv4).Msg("open next hop stream")
551
- return false
559
+
560
+ openCtx, cancel := context.WithTimeout(ctx, defaultClaimTimeout)
561
+ defer cancel()
562
+ var next net.Conn
563
+ var lastErr error
564
+ for {
565
+ var err error
566
+ next, err = s.hopMux.OpenStream(openCtx, overlayIPv4, forwardToken)
567
+ if err == nil {
568
+ break
569
+ }
570
+ lastErr = err
571
+ if errors.Is(err, net.ErrClosed) {
572
+ return fmt.Errorf("open next hop stream: %w", err)
573
+ }
574
+ if !utils.SleepOrDone(openCtx, defaultHopOpenRetryWait) {
575
+ return fmt.Errorf("open next hop stream within %s: %w", defaultClaimTimeout, errors.Join(lastErr, openCtx.Err()))
576
+ }
577
}
578
s.proxy.bridge(conn, next, "", nil)
554
- return true
579
+ return nil
580
+ }
581
+ if record.stream == nil {
582
+ return errors.New("lease stream is not ready")
583
}
556
- if record.stream == nil || !s.registry.policy.IsIdentityRoutable(record.Key()) {
557
- return false
584
+ if !s.registry.policy.IsIdentityRoutable(record.Key()) {
585
+ return errLeaseRejected
586
}
587
claimCtx, cancel := context.WithTimeout(ctx, defaultClaimTimeout)
588
session, err := record.stream.Claim(claimCtx)
589
cancel()
590
if err != nil {
563
- return false
591
+ return fmt.Errorf("claim lease stream: %w", err)
592
}
593
s.proxy.bridge(conn, session, record.Key(), s.registry.policy.BPSManager())
566
- return true
594
+ return nil
595
}
596
597
func (s *Server) runLeaseJanitor(ctx context.Context, interval time.Duration) error {
sdk/api_client.go
+14
-14
@@ -98,23 +98,12 @@ func (l *listener) registerLease(ctx context.Context, ttl time.Duration, udpEnab
98
return types.RegisterResponse{}, nil, errors.New("multi-hop relay set is unavailable")
99
}
100
101
- states := l.relaySet.AggregateRelays()
102
- descriptors := make(map[string]types.RelayDescriptor, len(states))
103
- for _, state := range states {
104
- desc := state.Descriptor
105
- if desc.APIHTTPSAddr != "" {
106
- descriptors[desc.APIHTTPSAddr] = desc
107
- }
108
- }
109
-
101
+ now := time.Now().UTC()
102
hopPath := make([]types.RelayDescriptor, 0, len(l.multiHop))
103
for i, relayURL := range l.multiHop {
112
- desc, ok := descriptors[relayURL]
104
+ desc, ok := l.relaySet.OverlayRelayDescriptor(relayURL, now)
105
if !ok {
114
- return types.RegisterResponse{}, nil, fmt.Errorf("multi-hop relay %d descriptor was not discovered", i)
115
- }
116
- if !desc.SupportsOverlay {
117
- return types.RegisterResponse{}, nil, fmt.Errorf("multi-hop relay %d does not support overlay", i)
106
+ return types.RegisterResponse{}, nil, fmt.Errorf("multi-hop relay %d descriptor is unavailable", i)
107
}
108
hopPath = append(hopPath, desc)
109
}
@@ -228,7 +217,18 @@ func (l *listener) syncHopRoutes(ctx context.Context, method string, expiresAt t
217
218
orderedRoutes := routes
219
if method == http.MethodPost {
220
+ if l.relaySet == nil {
221
+ return errors.New("multi-hop relay set is unavailable")
222
+ }
223
orderedRoutes = append([]types.HopRoute(nil), routes...)
224
+ now := time.Now().UTC()
225
+ for i := range orderedRoutes {
226
+ desc, ok := l.relaySet.OverlayRelayDescriptor(orderedRoutes[i].ForwardRelay.APIHTTPSAddr, now)
227
+ if !ok {
228
+ return fmt.Errorf("multi-hop forward relay %d descriptor is unavailable", i)
229
+ }
230
+ orderedRoutes[i].ForwardRelay = desc
231
+ }
232
slices.Reverse(orderedRoutes)
233
}
234
sdk/expose.go
+1
-1
@@ -165,7 +165,7 @@ func Expose(ctx context.Context, cfg ExposeConfig) (*Exposure, error) {
165
}
166
}
167
168
- if cfg.Discovery {
168
+ if cfg.Discovery || len(multiHop) > 0 {
169
go exposure.runDiscoveryLoop(exposureCtx)
170
}
171