Refactor QUIC tunnel to QUIC backhaul
Kim committed
Apr 29, 2026 at 10:09 UTC
6dfc36a14c031652ca8ee6b9f939e4f0abf85c40
16 files changed
+366
-329
cmd/portal-tunnel/installer/update.go
+2
-2
@@ -168,13 +168,13 @@ func assetURLs(version string) (binURL, checksumURL string, ok bool) {
168
return "", "", false
169
}
170
171
- baseURL := types.OfficialReleaseBaseURL
171
+ baseURL := types.OfficialReleaseBaseURL + "/latest/download"
172
version = strings.TrimSpace(version)
173
if version != "" {
174
if !strings.HasPrefix(version, "v") {
175
version = "v" + version
176
}
177
- baseURL = types.OfficialReleaseDownloadURL + "/" + version
177
+ baseURL = types.OfficialReleaseBaseURL + "/download/" + version
178
}
179
180
binURL = baseURL + "/" + filename
cmd/relay-server/frontend.go
+1
-1
@@ -395,7 +395,7 @@ func serveInstallBinary(w http.ResponseWriter, r *http.Request) {
395
}
396
data, err := embeddedDistFS.ReadFile("dist/tunnel/" + filename)
397
if err != nil {
398
- redirectURL := types.OfficialReleaseBaseURL + "/" + filename
398
+ redirectURL := types.OfficialReleaseBaseURL + "/latest/download/" + filename
399
if checksumRequest {
400
redirectURL += ".sha256"
401
}
docs/src/routes/tcp-udp-tunneling/+page.md
+1
-1
@@ -154,5 +154,5 @@ Both ports are drawn from the same `MIN_PORT`–`MAX_PORT` range on the relay.
154
155
- **Port range capacity** — the relay can serve at most `MAX_PORT - MIN_PORT + 1` concurrent TCP+UDP leases. Plan the range accordingly.
156
- **No TLS on raw TCP** — the raw TCP path has no TLS wrapping. Traffic between the relay and connecting clients is unencrypted. Use application-level encryption (e.g. SSH, WireGuard) if the service requires confidentiality.
157
-- **UDP max packet size** — datagrams are capped at **1350 bytes** (`DefaultMaxPacketSize`). Packets larger than this are dropped. Most game protocols fit within this limit, but verify if you use a custom protocol.
157
+- **UDP max packet size** — datagrams are capped at **1350 bytes**. Packets larger than this are dropped. Most game protocols fit within this limit, but verify if you use a custom protocol.
158
- **Flow idle timeout** — UDP flows that have seen no traffic for **30 seconds** are cleaned up on the relay. Long-lived connections should send keepalive packets if they may be idle.
portal/api_server.go
+21
-48
@@ -3,7 +3,6 @@ package portal
3
import (
4
"context"
5
"crypto/tls"
6
- "encoding/json"
6
"errors"
7
"fmt"
8
"io"
@@ -47,17 +46,6 @@ var (
46
errTCPPortExhausted = &apiError{types.APIErrorCodeTCPPortExhausted, "no tcp ports available", http.StatusServiceUnavailable}
47
)
48
50
-var quicRejectTable = []struct {
51
- sentinel error
52
- code string
53
- reason string
54
-}{
55
- {errLeaseNotFound, types.APIErrorCodeLeaseNotFound, "lease not found"},
56
- {errLeaseRejected, types.APIErrorCodeLeaseRejected, "lease rejected"},
57
- {errUnauthorized, types.APIErrorCodeUnauthorized, "unauthorized"},
58
- {errTransportMismatch, types.APIErrorCodeTransportMismatch, "transport mismatch"},
59
-}
60
-
49
func writeAPIErrorResponse(w http.ResponseWriter, err error) {
50
var ae *apiError
51
if errors.As(err, &ae) {
@@ -183,7 +171,7 @@ func (s *Server) signedRelayDescriptor(now time.Time) (types.RelayDescriptor, er
171
WireGuardPublicKey: wireGuardPublicKey,
172
WireGuardPort: wireGuardPort,
173
SupportsOverlay: s.overlay != nil,
186
- SupportsUDP: s.cfg.UDPEnabled && s.quicTunnel != nil,
174
+ SupportsUDP: s.cfg.UDPEnabled && s.quicBackhaul != nil,
175
SupportsTCP: s.cfg.TCPEnabled,
176
ActiveConnections: s.proxy.activeConnectionCount(),
177
TCPBPS: s.proxy.currentTCPBPS(now),
@@ -379,7 +367,7 @@ func (s *Server) handleRegisterChallenge(w http.ResponseWriter, r *http.Request)
367
utils.WriteAPIError(w, http.StatusServiceUnavailable, types.APIErrorCodeFeatureUnavailable, errFeatureUnavailable.Error())
368
return
369
}
382
- if req.UDPEnabled && (!s.cfg.UDPEnabled || s.group != nil && s.quicTunnel == nil) {
370
+ if req.UDPEnabled && (!s.cfg.UDPEnabled || s.group != nil && s.quicBackhaul == nil) {
371
utils.WriteAPIError(w, http.StatusServiceUnavailable, types.APIErrorCodeFeatureUnavailable, errFeatureUnavailable.Error())
372
return
373
}
@@ -632,58 +620,43 @@ func (s *Server) handleConnect(w http.ResponseWriter, r *http.Request) {
620
Msg("sdk reverse connected")
621
}
622
635
-func (s *Server) handleQUICTunnelConn(conn *quic.Conn) {
636
- stream, err := conn.AcceptStream(context.Background())
623
+func (s *Server) handleQUICBackhaulConn(conn *quic.Conn) {
624
+ control, err := transport.AcceptQUICBackhaulControl(context.Background(), conn)
625
if err != nil {
638
- _ = conn.CloseWithError(1, "stream accept failed")
639
- return
640
- }
641
-
642
- _ = stream.SetReadDeadline(time.Now().Add(10 * time.Second))
643
- var msg types.QUICControlMessage
644
- if err := json.NewDecoder(io.LimitReader(stream, defaultControlBodyLimit)).Decode(&msg); err != nil {
626
_ = conn.CloseWithError(1, "control read failed")
627
return
628
}
648
- _ = stream.SetReadDeadline(time.Time{})
649
- if strings.TrimSpace(msg.AccessToken) == "" {
650
- _ = json.NewEncoder(stream).Encode(types.QUICControlResponse{OK: false, Error: "invalid_control_message"})
651
- _ = conn.CloseWithError(1, "invalid control message")
652
- return
653
- }
654
-
655
- rejectConn := func(code, reason string) {
656
- _ = json.NewEncoder(stream).Encode(types.QUICControlResponse{OK: false, Error: code})
657
- _ = conn.CloseWithError(1, reason)
658
- }
629
660
- lease, err := s.admitLeaseByToken(msg.AccessToken, true)
630
+ lease, err := s.admitLeaseByToken(control.AccessToken, true)
631
if err != nil {
632
code, reason := types.APIErrorCodeInvalidRequest, "invalid control message"
663
- for _, entry := range quicRejectTable {
664
- if errors.Is(err, entry.sentinel) {
665
- code, reason = entry.code, entry.reason
666
- break
667
- }
633
+ switch {
634
+ case errors.Is(err, errLeaseNotFound):
635
+ code, reason = types.APIErrorCodeLeaseNotFound, "lease not found"
636
+ case errors.Is(err, errLeaseRejected):
637
+ code, reason = types.APIErrorCodeLeaseRejected, "lease rejected"
638
+ case errors.Is(err, errUnauthorized):
639
+ code, reason = types.APIErrorCodeUnauthorized, "unauthorized"
640
+ case errors.Is(err, errTransportMismatch):
641
+ code, reason = types.APIErrorCodeTransportMismatch, "transport mismatch"
642
}
669
- rejectConn(code, reason)
643
+ _ = control.Reject(code, reason)
644
return
645
}
646
673
- if err := lease.datagram.Register(conn); err != nil {
674
- _ = json.NewEncoder(stream).Encode(types.QUICControlResponse{OK: false, Error: "broker_closed"})
675
- _ = conn.CloseWithError(1, "broker closed")
647
+ if err := lease.datagram.BindBackhaul(conn); err != nil {
648
+ _ = control.Reject("broker_closed", "broker closed")
649
return
650
}
651
679
- _ = json.NewEncoder(stream).Encode(types.QUICControlResponse{OK: true})
652
+ _ = control.Accept()
653
s.registry.Touch(lease.Key(), conn.RemoteAddr().String(), time.Now())
654
log.Info().
682
- Str("component", "quic-tunnel-listener").
655
+ Str("component", "quic-backhaul-listener").
656
Str("address", lease.Address).
657
Str("lease_name", lease.Name).
658
Str("remote_addr", conn.RemoteAddr().String()).
686
- Msg("quic tunnel connected")
659
+ Msg("quic backhaul connected")
660
}
661
662
func (s *Server) admitLeaseByToken(token string, requireDatagram bool) (*leaseRecord, error) {
@@ -723,7 +696,7 @@ func (s *Server) registerLease(req types.RegisterChallengeRequest, clientIP, rep
696
}
697
698
if req.UDPEnabled {
726
- if !s.cfg.UDPEnabled || s.group != nil && s.quicTunnel == nil {
699
+ if !s.cfg.UDPEnabled || s.group != nil && s.quicBackhaul == nil {
700
return types.RegisterResponse{}, errFeatureUnavailable
701
}
702
if !s.registry.policy.IsUDPEnabled() {
portal/server.go
+25
-50
@@ -112,11 +112,11 @@ type Server struct {
112
acmeManager *acme.Manager
113
proxy proxy
114
115
- apiListener net.Listener
116
- sniListener net.Listener
117
- apiServer *http.Server
118
- apiTLSClose io.Closer
119
- quicTunnel *quic.Listener
115
+ apiListener net.Listener
116
+ sniListener net.Listener
117
+ apiServer *http.Server
118
+ apiTLSClose io.Closer
119
+ quicBackhaul *quic.Listener
120
121
overlay *overlay.Overlay
122
hopMux *overlay.HopMux
@@ -179,7 +179,7 @@ func (s *Server) Start(ctx context.Context, apiMux *http.ServeMux) error {
179
var apiCloser io.Closer
180
var hopMux *overlay.HopMux
181
var ov *overlay.Overlay
182
- var quicTunnel *quic.Listener
182
+ var quicBackhaul *quic.Listener
183
defer func() {
184
if started {
185
return
@@ -233,10 +233,10 @@ func (s *Server) Start(ctx context.Context, apiMux *http.ServeMux) error {
233
}
234
}
235
if s.cfg.UDPEnabled {
236
- quicTunnel, err = s.newQUICTunnelListener(apiTLS)
236
+ quicBackhaul, err = s.newQUICBackhaulListener(apiTLS)
237
if err != nil {
238
- log.Warn().Err(err).Msg("quic tunnel listener disabled")
239
- quicTunnel = nil
238
+ log.Warn().Err(err).Msg("quic backhaul listener disabled")
239
+ quicBackhaul = nil
240
}
241
}
242
@@ -249,7 +249,7 @@ func (s *Server) Start(ctx context.Context, apiMux *http.ServeMux) error {
249
s.group = group
250
s.overlay = ov
251
s.hopMux = hopMux
252
- s.quicTunnel = quicTunnel
252
+ s.quicBackhaul = quicBackhaul
253
started = true
254
255
group.Go(s.runAPIServer)
@@ -260,8 +260,8 @@ func (s *Server) Start(ctx context.Context, apiMux *http.ServeMux) error {
260
group.Go(func() error { return s.runOverlayIngress(groupCtx) })
261
}
262
}
263
- if s.quicTunnel != nil {
264
- group.Go(s.runQUICTunnelListener)
263
+ if s.quicBackhaul != nil {
264
+ group.Go(s.runQUICBackhaulListener)
265
}
266
group.Go(func() error { return s.runRegistryJanitor(groupCtx, 5*time.Second) })
267
if s.cfg.DiscoveryEnabled {
@@ -285,10 +285,10 @@ func (s *Server) Start(ctx context.Context, apiMux *http.ServeMux) error {
285
Bool("discovery_enabled", s.cfg.DiscoveryEnabled).
286
Bool("wireguard_enabled", s.overlay != nil).
287
Bool("multihop_enabled", s.hopMux != nil).
288
- Bool("udp_enabled", s.quicTunnel != nil).
288
+ Bool("udp_enabled", s.quicBackhaul != nil).
289
Bool("tcp_enabled", s.cfg.TCPEnabled)
290
- if s.quicTunnel != nil {
291
- logEvent = logEvent.Str("internal_quic_tunnel_addr", s.quicTunnel.Addr().String())
290
+ if s.quicBackhaul != nil {
291
+ logEvent = logEvent.Str("internal_quic_backhaul_addr", s.quicBackhaul.Addr().String())
292
}
293
logEvent.Msg("relay server started")
294
@@ -352,8 +352,8 @@ func (s *Server) Shutdown(ctx context.Context) error {
352
s.cleanupRemovedRecord(ctx, lease, "delete lease remote state during shutdown")
353
}
354
355
- if s.quicTunnel != nil {
356
- _ = s.quicTunnel.Close()
355
+ if s.quicBackhaul != nil {
356
+ _ = s.quicBackhaul.Close()
357
}
358
if s.sniListener != nil {
359
if err := s.sniListener.Close(); err != nil && !errors.Is(err, net.ErrClosed) {
@@ -410,14 +410,6 @@ func (s *Server) prepareAPITLS(ctx context.Context) (keyless.TLSMaterialConfig,
410
CertPEM: certPEM,
411
KeyPEM: keyPEM,
412
}
413
- if len(apiTLS.CertPEM) == 0 {
414
- manager.Stop()
415
- return keyless.TLSMaterialConfig{}, nil, errors.New("api tls certificate is required")
416
- }
417
- if len(apiTLS.KeyPEM) == 0 && apiTLS.Keyless == nil {
418
- manager.Stop()
419
- return keyless.TLSMaterialConfig{}, nil, errors.New("api tls key or keyless signer is required")
420
- }
413
414
return apiTLS, manager, nil
415
}
@@ -585,47 +577,30 @@ func (s *Server) runRegistryJanitor(ctx context.Context, interval time.Duration)
577
}
578
}
579
588
-func (s *Server) newQUICTunnelListener(apiTLS keyless.TLSMaterialConfig) (*quic.Listener, error) {
580
+func (s *Server) newQUICBackhaulListener(apiTLS keyless.TLSMaterialConfig) (*quic.Listener, error) {
581
if len(apiTLS.KeyPEM) == 0 {
590
- return nil, fmt.Errorf("quic tunnel requires api tls key")
582
+ return nil, fmt.Errorf("quic backhaul requires api tls key")
583
}
584
tlsCert, err := tls.X509KeyPair(apiTLS.CertPEM, apiTLS.KeyPEM)
585
if err != nil {
594
- return nil, fmt.Errorf("parse quic tls keypair: %w", err)
595
- }
596
-
597
- tlsConf := &tls.Config{
598
- Certificates: []tls.Certificate{tlsCert},
599
- NextProtos: []string{"portal-tunnel"},
600
- MinVersion: tls.VersionTLS13,
601
- }
602
- quicConf := &quic.Config{
603
- EnableDatagrams: true,
604
- KeepAlivePeriod: 15 * time.Second,
605
- MaxIdleTimeout: 60 * time.Second,
606
- MaxIncomingStreams: 16,
607
- }
608
-
609
- listener, err := quic.ListenAddr(s.cfg.SNIListenAddr, tlsConf, quicConf)
610
- if err != nil {
611
- return nil, fmt.Errorf("listen quic: %w", err)
586
+ return nil, fmt.Errorf("parse quic backhaul tls keypair: %w", err)
587
}
613
- return listener, nil
588
+ return transport.ListenQUICBackhaul(s.cfg.SNIListenAddr, tlsCert)
589
}
590
616
-func (s *Server) runQUICTunnelListener() error {
617
- if s.quicTunnel == nil {
591
+func (s *Server) runQUICBackhaulListener() error {
592
+ if s.quicBackhaul == nil {
593
return nil
594
}
595
for {
621
- conn, err := s.quicTunnel.Accept(context.Background())
596
+ conn, err := s.quicBackhaul.Accept(context.Background())
597
if err != nil {
598
if errors.Is(err, quic.ErrServerClosed) || errors.Is(err, net.ErrClosed) {
599
return nil
600
}
601
return err
602
}
628
- go s.handleQUICTunnelConn(conn)
603
+ go s.handleQUICBackhaulConn(conn)
604
}
605
}
606
portal/transport/datagram_client.go
+1
-1
@@ -18,7 +18,7 @@ func NewClientDatagram(onReceiveError func(error)) *ClientDatagram {
18
}
19
}
20
21
-func (d *ClientDatagram) Bind(conn *quic.Conn) (<-chan struct{}, error) {
21
+func (d *ClientDatagram) BindBackhaul(conn *quic.Conn) (<-chan struct{}, error) {
22
if d == nil || d.session == nil {
23
if conn != nil {
24
_ = conn.CloseWithError(0, "listener closed")
portal/transport/datagram_relay.go
+15
-115
@@ -5,7 +5,6 @@ import (
5
"errors"
6
"fmt"
7
"net"
8
- "sort"
8
"sync"
9
"time"
10
@@ -16,117 +15,18 @@ import (
15
)
16
17
const (
19
- DefaultMaxPacketSize = 1350
18
+ defaultMaxPacketSize = 1350
19
defaultFlowIdleTimeout = 30 * time.Second
20
defaultFlowCleanupInterval = 30 * time.Second
21
)
22
24
-var ErrPortExhausted = errors.New("no ports available")
25
-
26
-type flowReplyFunc func([]byte) error
27
-
23
type flowState struct {
24
key string
25
lastSeen time.Time
31
- reply flowReplyFunc
32
-}
33
-
34
-type portReservation struct {
35
- port int
36
- expiresAt time.Time
37
-}
38
-
39
-// PortAllocator manages a pool of ports for dynamic per-lease allocation.
40
-type PortAllocator struct {
41
- available []int
42
- inUse map[int]string
43
- reserved map[string]portReservation
44
- grace time.Duration
45
- mu sync.Mutex
46
-}
47
-
48
-func NewPortAllocator(min, max int, grace time.Duration) *PortAllocator {
49
- if min <= 0 || max <= 0 || min > max {
50
- return &PortAllocator{
51
- available: nil,
52
- inUse: make(map[int]string),
53
- reserved: make(map[string]portReservation),
54
- grace: grace,
55
- }
56
- }
57
- available := make([]int, 0, max-min+1)
58
- for p := min; p <= max; p++ {
59
- available = append(available, p)
60
- }
61
- return &PortAllocator{
62
- available: available,
63
- inUse: make(map[int]string),
64
- reserved: make(map[string]portReservation),
65
- grace: grace,
66
- }
26
+ reply func([]byte) error
27
}
28
69
-func (a *PortAllocator) Allocate(name string) (int, error) {
70
- a.mu.Lock()
71
- defer a.mu.Unlock()
72
-
73
- a.cleanupExpiredLocked(time.Now())
74
-
75
- if res, ok := a.reserved[name]; ok {
76
- delete(a.reserved, name)
77
- a.inUse[res.port] = name
78
- return res.port, nil
79
- }
80
-
81
- if len(a.available) == 0 {
82
- return 0, ErrPortExhausted
83
- }
84
-
85
- port := a.available[0]
86
- a.available = a.available[1:]
87
- a.inUse[port] = name
88
- return port, nil
89
-}
90
-
91
-func (a *PortAllocator) Release(port int) {
92
- a.mu.Lock()
93
- defer a.mu.Unlock()
94
-
95
- name, ok := a.inUse[port]
96
- if !ok {
97
- return
98
- }
99
- delete(a.inUse, port)
100
-
101
- if prev, exists := a.reserved[name]; exists {
102
- a.sortedInsertLocked(prev.port)
103
- }
104
-
105
- a.reserved[name] = portReservation{
106
- port: port,
107
- expiresAt: time.Now().Add(a.grace),
108
- }
109
-
110
- a.cleanupExpiredLocked(time.Now())
111
-}
112
-
113
-func (a *PortAllocator) cleanupExpiredLocked(now time.Time) {
114
- for name, res := range a.reserved {
115
- if now.After(res.expiresAt) {
116
- delete(a.reserved, name)
117
- a.sortedInsertLocked(res.port)
118
- }
119
- }
120
-}
121
-
122
-func (a *PortAllocator) sortedInsertLocked(port int) {
123
- i := sort.SearchInts(a.available, port)
124
- a.available = append(a.available, 0)
125
- copy(a.available[i+1:], a.available[i:])
126
- a.available[i] = port
127
-}
128
-
129
-// Datagram owns the UDP and QUIC datagram runtime for one lease.
29
+// RelayDatagram owns UDP ingress and QUIC backhaul binding for one lease.
30
type RelayDatagram struct {
31
identityKey string
32
port int
@@ -149,9 +49,9 @@ func NewRelayDatagram(identityKey string, port int) *RelayDatagram {
49
session: newDatagramSession(256, true, func(err error) {
50
log.Warn().
51
Err(err).
152
- Str("component", "quic-flow-mux").
52
+ Str("component", "quic-backhaul").
53
Str("identity_key", identityKey).
154
- Msg("quic receive loop ended")
54
+ Msg("quic backhaul receive loop ended")
55
}),
56
flowTable: make(map[uint32]*flowState),
57
addrIndex: make(map[string]uint32),
@@ -208,27 +108,27 @@ func (d *RelayDatagram) Close() {
108
})
109
}
110
211
-func (d *RelayDatagram) Register(conn *quic.Conn) error {
111
+func (d *RelayDatagram) BindBackhaul(conn *quic.Conn) error {
112
if _, err := d.session.Bind(conn); err != nil {
113
return err
114
}
115
116
log.Info().
217
- Str("component", "quic-flow-mux").
117
+ Str("component", "quic-backhaul").
118
Str("identity_key", d.identityKey).
119
Str("remote_addr", conn.RemoteAddr().String()).
220
- Msg("quic tunnel connection registered")
120
+ Msg("quic backhaul connection registered")
121
return nil
122
}
123
224
-func (d *RelayDatagram) SendDatagram(flowID uint32, payload []byte) error {
124
+func (d *RelayDatagram) sendDatagram(flowID uint32, payload []byte) error {
125
if d == nil {
126
return net.ErrClosed
127
}
128
return d.session.Send(flowID, payload)
129
}
130
231
-func (d *RelayDatagram) TouchFlow(key string, reply func([]byte) error) uint32 {
131
+func (d *RelayDatagram) touchFlow(key string, reply func([]byte) error) uint32 {
132
now := time.Now()
133
134
d.mu.Lock()
@@ -289,7 +189,7 @@ func (d *RelayDatagram) dispatch(frame types.DatagramFrame) {
189
if err := reply(frame.Payload); err != nil {
190
log.Warn().
191
Err(err).
292
- Str("component", "quic-flow-mux").
192
+ Str("component", "udp-relay").
193
Str("identity_key", d.identityKey).
194
Uint32("flow_id", frame.FlowID).
195
Msg("flow writeback failed")
@@ -340,7 +240,7 @@ func (d *RelayDatagram) forgetFlow(flowID uint32) {
240
}
241
242
func (d *RelayDatagram) readLoop(ctx context.Context) {
343
- buf := make([]byte, DefaultMaxPacketSize)
243
+ buf := make([]byte, defaultMaxPacketSize)
244
for {
245
select {
246
case <-ctx.Done():
@@ -366,21 +266,21 @@ func (d *RelayDatagram) readLoop(ctx context.Context) {
266
return
267
}
268
369
- flowID := d.TouchFlow("udp:"+clientAddr.String(), func(payload []byte) error {
269
+ flowID := d.touchFlow("udp:"+clientAddr.String(), func(payload []byte) error {
270
_, err := d.conn.WriteToUDP(payload, clientAddr)
271
return err
272
})
273
payload := make([]byte, n)
274
copy(payload, buf[:n])
275
376
- if err := d.SendDatagram(flowID, payload); err != nil {
276
+ if err := d.sendDatagram(flowID, payload); err != nil {
277
log.Warn().
278
Str("component", "udp-relay").
279
Str("identity_key", d.identityKey).
280
Err(err).
281
Uint32("flow_id", flowID).
282
Int("bytes", n).
383
- Msg("send datagram to tunnel failed, dropping packet")
283
+ Msg("send datagram to quic backhaul failed, dropping packet")
284
continue
285
}
286
}
portal/transport/datagram_session.go
+3
-3
@@ -11,7 +11,7 @@ import (
11
"github.com/gosuda/portal-tunnel/v2/types"
12
)
13
14
-var errNoConnection = errors.New("no quic connection registered")
14
+var errNoConnection = errors.New("no quic backhaul connection registered")
15
16
// datagramSession owns one active QUIC DATAGRAM connection and exposes decoded frames.
17
type datagramSession struct {
@@ -38,11 +38,11 @@ func newDatagramSession(bufferSize int, dropIncoming bool, onReceiveError func(e
38
}
39
}
40
41
-// Bind installs a new active QUIC connection and starts the receive loop.
41
+// Bind installs a new active backhaul connection and starts the receive loop.
42
// Any previously active connection is replaced and closed.
43
func (s *datagramSession) Bind(conn *quic.Conn) (<-chan struct{}, error) {
44
if conn == nil {
45
- return nil, errors.New("quic connection is required")
45
+ return nil, errors.New("quic backhaul connection is required")
46
}
47
48
s.mu.Lock()
portal/transport/port_allocator.go
new
+104
@@ -0,0 +1,104 @@
1
+package transport
2
+
3
+import (
4
+ "errors"
5
+ "sort"
6
+ "sync"
7
+ "time"
8
+)
9
+
10
+var ErrPortExhausted = errors.New("no ports available")
11
+
12
+type portReservation struct {
13
+ port int
14
+ expiresAt time.Time
15
+}
16
+
17
+// PortAllocator manages a pool of ports for dynamic per-lease allocation.
18
+type PortAllocator struct {
19
+ available []int
20
+ inUse map[int]string
21
+ reserved map[string]portReservation
22
+ grace time.Duration
23
+ mu sync.Mutex
24
+}
25
+
26
+func NewPortAllocator(min, max int, grace time.Duration) *PortAllocator {
27
+ if min <= 0 || max <= 0 || min > max {
28
+ return &PortAllocator{
29
+ inUse: make(map[int]string),
30
+ reserved: make(map[string]portReservation),
31
+ grace: grace,
32
+ }
33
+ }
34
+
35
+ available := make([]int, 0, max-min+1)
36
+ for port := min; port <= max; port++ {
37
+ available = append(available, port)
38
+ }
39
+ return &PortAllocator{
40
+ available: available,
41
+ inUse: make(map[int]string),
42
+ reserved: make(map[string]portReservation),
43
+ grace: grace,
44
+ }
45
+}
46
+
47
+func (a *PortAllocator) Allocate(name string) (int, error) {
48
+ a.mu.Lock()
49
+ defer a.mu.Unlock()
50
+
51
+ a.cleanupExpiredLocked(time.Now())
52
+
53
+ if res, ok := a.reserved[name]; ok {
54
+ delete(a.reserved, name)
55
+ a.inUse[res.port] = name
56
+ return res.port, nil
57
+ }
58
+
59
+ if len(a.available) == 0 {
60
+ return 0, ErrPortExhausted
61
+ }
62
+
63
+ port := a.available[0]
64
+ a.available = a.available[1:]
65
+ a.inUse[port] = name
66
+ return port, nil
67
+}
68
+
69
+func (a *PortAllocator) Release(port int) {
70
+ a.mu.Lock()
71
+ defer a.mu.Unlock()
72
+
73
+ name, ok := a.inUse[port]
74
+ if !ok {
75
+ return
76
+ }
77
+ delete(a.inUse, port)
78
+
79
+ if prev, exists := a.reserved[name]; exists {
80
+ a.sortedInsertLocked(prev.port)
81
+ }
82
+
83
+ a.reserved[name] = portReservation{
84
+ port: port,
85
+ expiresAt: time.Now().Add(a.grace),
86
+ }
87
+ a.cleanupExpiredLocked(time.Now())
88
+}
89
+
90
+func (a *PortAllocator) cleanupExpiredLocked(now time.Time) {
91
+ for name, res := range a.reserved {
92
+ if now.After(res.expiresAt) {
93
+ delete(a.reserved, name)
94
+ a.sortedInsertLocked(res.port)
95
+ }
96
+ }
97
+}
98
+
99
+func (a *PortAllocator) sortedInsertLocked(port int) {
100
+ i := sort.SearchInts(a.available, port)
101
+ a.available = append(a.available, 0)
102
+ copy(a.available[i+1:], a.available[i:])
103
+ a.available[i] = port
104
+}
portal/transport/quic_backhaul.go
new
+174
@@ -0,0 +1,174 @@
1
+package transport
2
+
3
+import (
4
+ "context"
5
+ "crypto/tls"
6
+ "encoding/json"
7
+ "errors"
8
+ "fmt"
9
+ "io"
10
+ "strings"
11
+ "time"
12
+
13
+ "github.com/quic-go/quic-go"
14
+)
15
+
16
+const (
17
+ quicBackhaulALPN = "portal-tunnel"
18
+ quicBackhaulControlTimeout = 10 * time.Second
19
+ quicBackhaulControlBodyLimit = 4096
20
+)
21
+
22
+type quicBackhaulControlMessage struct {
23
+ AccessToken string `json:"access_token"`
24
+}
25
+
26
+type quicBackhaulControlResponse struct {
27
+ OK bool `json:"ok"`
28
+ Error string `json:"error,omitempty"`
29
+}
30
+
31
+type QUICBackhaulControl struct {
32
+ AccessToken string
33
+ conn *quic.Conn
34
+ stream *quic.Stream
35
+}
36
+
37
+func ListenQUICBackhaul(addr string, cert tls.Certificate) (*quic.Listener, error) {
38
+ listener, err := quic.ListenAddr(addr, quicBackhaulServerTLSConfig(cert), quicBackhaulConfig())
39
+ if err != nil {
40
+ return nil, fmt.Errorf("listen quic backhaul: %w", err)
41
+ }
42
+ return listener, nil
43
+}
44
+
45
+func DialQUICBackhaul(ctx context.Context, addr string, tlsConfig *tls.Config, accessToken string) (*quic.Conn, error) {
46
+ accessToken = strings.TrimSpace(accessToken)
47
+ if accessToken == "" {
48
+ return nil, errors.New("quic backhaul access token is required")
49
+ }
50
+
51
+ conn, err := quic.DialAddr(ctx, addr, quicBackhaulClientTLSConfig(tlsConfig), quicBackhaulConfig())
52
+ if err != nil {
53
+ return nil, fmt.Errorf("dial quic backhaul: %w", err)
54
+ }
55
+
56
+ stream, err := conn.OpenStreamSync(ctx)
57
+ if err != nil {
58
+ _ = conn.CloseWithError(1, "control stream open failed")
59
+ return nil, fmt.Errorf("open quic backhaul control stream: %w", err)
60
+ }
61
+
62
+ _ = stream.SetDeadline(time.Now().Add(quicBackhaulControlTimeout))
63
+ if err := json.NewEncoder(stream).Encode(quicBackhaulControlMessage{AccessToken: accessToken}); err != nil {
64
+ _ = conn.CloseWithError(1, "control write failed")
65
+ return nil, fmt.Errorf("write quic backhaul control message: %w", err)
66
+ }
67
+
68
+ var resp quicBackhaulControlResponse
69
+ if err := json.NewDecoder(io.LimitReader(stream, quicBackhaulControlBodyLimit)).Decode(&resp); err != nil {
70
+ _ = conn.CloseWithError(1, "control response read failed")
71
+ return nil, fmt.Errorf("read quic backhaul control response: %w", err)
72
+ }
73
+ _ = stream.SetDeadline(time.Time{})
74
+ _ = stream.Close()
75
+
76
+ if !resp.OK {
77
+ errText := strings.TrimSpace(resp.Error)
78
+ if errText == "" {
79
+ errText = "rejected"
80
+ }
81
+ _ = conn.CloseWithError(1, errText)
82
+ return nil, fmt.Errorf("quic backhaul rejected: %s", errText)
83
+ }
84
+ return conn, nil
85
+}
86
+
87
+func AcceptQUICBackhaulControl(ctx context.Context, conn *quic.Conn) (*QUICBackhaulControl, error) {
88
+ stream, err := conn.AcceptStream(ctx)
89
+ if err != nil {
90
+ return nil, fmt.Errorf("accept quic backhaul control stream: %w", err)
91
+ }
92
+
93
+ _ = stream.SetReadDeadline(time.Now().Add(quicBackhaulControlTimeout))
94
+ var msg quicBackhaulControlMessage
95
+ if err := json.NewDecoder(io.LimitReader(stream, quicBackhaulControlBodyLimit)).Decode(&msg); err != nil {
96
+ return nil, fmt.Errorf("read quic backhaul control message: %w", err)
97
+ }
98
+ _ = stream.SetReadDeadline(time.Time{})
99
+
100
+ accessToken := strings.TrimSpace(msg.AccessToken)
101
+ if accessToken == "" {
102
+ return nil, errors.New("quic backhaul access token is required")
103
+ }
104
+
105
+ return &QUICBackhaulControl{
106
+ AccessToken: accessToken,
107
+ conn: conn,
108
+ stream: stream,
109
+ }, nil
110
+}
111
+
112
+func (c *QUICBackhaulControl) Accept() error {
113
+ if c == nil || c.stream == nil {
114
+ return nil
115
+ }
116
+ err := json.NewEncoder(c.stream).Encode(quicBackhaulControlResponse{OK: true})
117
+ return errors.Join(err, c.stream.Close())
118
+}
119
+
120
+func (c *QUICBackhaulControl) Reject(code, reason string) error {
121
+ if c == nil || c.conn == nil {
122
+ return nil
123
+ }
124
+ code = strings.TrimSpace(code)
125
+ if code == "" {
126
+ code = "rejected"
127
+ }
128
+ reason = strings.TrimSpace(reason)
129
+ if reason == "" {
130
+ reason = code
131
+ }
132
+
133
+ var err error
134
+ if c.stream != nil {
135
+ err = errors.Join(
136
+ json.NewEncoder(c.stream).Encode(quicBackhaulControlResponse{OK: false, Error: code}),
137
+ c.stream.Close(),
138
+ )
139
+ }
140
+ return errors.Join(err, c.conn.CloseWithError(1, reason))
141
+}
142
+
143
+func quicBackhaulServerTLSConfig(cert tls.Certificate) *tls.Config {
144
+ return &tls.Config{
145
+ Certificates: []tls.Certificate{cert},
146
+ NextProtos: []string{quicBackhaulALPN},
147
+ MinVersion: tls.VersionTLS13,
148
+ }
149
+}
150
+
151
+func quicBackhaulClientTLSConfig(base *tls.Config) *tls.Config {
152
+ if base == nil {
153
+ return &tls.Config{
154
+ NextProtos: []string{quicBackhaulALPN},
155
+ MinVersion: tls.VersionTLS13,
156
+ }
157
+ }
158
+
159
+ cfg := base.Clone()
160
+ cfg.NextProtos = []string{quicBackhaulALPN}
161
+ if cfg.MinVersion == 0 || cfg.MinVersion < tls.VersionTLS13 {
162
+ cfg.MinVersion = tls.VersionTLS13
163
+ }
164
+ return cfg
165
+}
166
+
167
+func quicBackhaulConfig() *quic.Config {
168
+ return &quic.Config{
169
+ EnableDatagrams: true,
170
+ KeepAlivePeriod: 15 * time.Second,
171
+ MaxIdleTimeout: 60 * time.Second,
172
+ MaxIncomingStreams: 16,
173
+ }
174
+}
portal/transport/stream_client.go
-37
@@ -7,7 +7,6 @@ import (
7
"fmt"
8
"io"
9
"net"
10
- "sync"
10
"time"
11
12
"github.com/gosuda/portal-tunnel/v2/types"
@@ -15,9 +14,7 @@ import (
14
15
type ClientStream struct {
16
accepted chan net.Conn
18
- activeSessions int
17
handshakeTimeout time.Duration
20
- mu sync.Mutex
18
}
19
20
func NewClientStream(readyTarget int, handshakeTimeout time.Duration) *ClientStream {
@@ -53,16 +50,6 @@ func (s *ClientStream) RunSession(
50
return s.runSession(ctx, conn, tlsConfig)
51
}
52
56
-func (s *ClientStream) ActiveSessions() int {
57
- if s == nil {
58
- return 0
59
- }
60
-
61
- s.mu.Lock()
62
- defer s.mu.Unlock()
63
- return s.activeSessions
64
-}
65
-
53
func (s *ClientStream) Drain() {
54
if s == nil {
55
return
@@ -87,8 +74,6 @@ func (s *ClientStream) runSession(
74
if conn == nil {
75
return false, net.ErrClosed
76
}
90
- s.sessionOpened()
91
- defer s.sessionClosed()
77
78
var marker [1]byte
79
for {
@@ -151,25 +136,3 @@ func (s *ClientStream) activateRaw(ctx context.Context, conn net.Conn) error {
136
return nil
137
}
138
}
154
-
155
-func (s *ClientStream) sessionOpened() {
156
- if s == nil {
157
- return
158
- }
159
-
160
- s.mu.Lock()
161
- s.activeSessions++
162
- s.mu.Unlock()
163
-}
164
-
165
-func (s *ClientStream) sessionClosed() {
166
- if s == nil {
167
- return
168
- }
169
-
170
- s.mu.Lock()
171
- if s.activeSessions > 0 {
172
- s.activeSessions--
173
- }
174
- s.mu.Unlock()
175
-}
portal/transport/stream_relay.go
+1
-1
@@ -69,7 +69,7 @@ func (b *RelayStream) Claim(ctx context.Context) (net.Conn, error) {
69
return b.claimWithMarker(ctx, types.MarkerTLSStart)
70
}
71
72
-func (b *RelayStream) ClaimRaw(ctx context.Context) (net.Conn, error) {
72
+func (b *RelayStream) claimRaw(ctx context.Context) (net.Conn, error) {
73
return b.claimWithMarker(ctx, types.MarkerRawStart)
74
}
75
portal/transport/tcp_port_relay.go
renamed
+1
-3
@@ -13,8 +13,6 @@ import (
13
const defaultTCPPortClaimTimeout = 10 * time.Second
14
15
// RelayTCPPort owns a TCP listener on an allocated port for one lease.
16
-// Incoming connections are bridged to reverse sessions claimed from the
17
-// associated RelayStream using raw TCP (no TLS).
16
type RelayTCPPort struct {
17
identityKey string
18
port int
@@ -114,7 +112,7 @@ func (t *RelayTCPPort) handleConn(ctx context.Context, conn net.Conn) {
112
claimCtx, cancel := context.WithTimeout(ctx, defaultTCPPortClaimTimeout)
113
defer cancel()
114
117
- session, err := t.stream.ClaimRaw(claimCtx)
115
+ session, err := t.stream.claimRaw(claimCtx)
116
if err != nil {
117
_ = conn.Close()
118
log.Warn().
sdk/listener.go
+12
-52
@@ -5,7 +5,6 @@ import (
5
"bytes"
6
"context"
7
"crypto/tls"
8
- "encoding/json"
8
"errors"
9
"fmt"
10
"io"
@@ -123,9 +122,9 @@ func newListener(ctx context.Context, relayURL string, cfg listenerConfig) (*lis
122
l.datagram = transport.NewClientDatagram(func(err error) {
123
log.Info().
124
Err(err).
126
- Str("component", "sdk-datagram-plane").
125
+ Str("component", "sdk-quic-backhaul").
126
Str("address", l.identity.Address).
128
- Msg("quic datagram plane disconnected; waiting to reconnect")
127
+ Msg("quic backhaul disconnected; waiting to reconnect")
128
})
129
}
130
@@ -482,13 +481,13 @@ func (l *listener) runDatagramLoop(ctx context.Context) {
481
default:
482
}
483
485
- conn, err := l.openQUICSession(ctx)
484
+ conn, err := l.openQUICBackhaulSession(ctx)
485
if err != nil {
486
log.Info().
487
Err(err).
489
- Str("component", "sdk-datagram-plane").
488
+ Str("component", "sdk-quic-backhaul").
489
Str("address", l.identity.Address).
491
- Msg("quic datagram plane unavailable; retrying")
490
+ Msg("quic backhaul unavailable; retrying")
491
if !utils.SleepOrDone(ctx, 2*time.Second) {
492
l.datagram.Clear("lease stopped")
493
return
@@ -497,21 +496,21 @@ func (l *listener) runDatagramLoop(ctx context.Context) {
496
}
497
498
log.Info().
500
- Str("component", "sdk-datagram-plane").
499
+ Str("component", "sdk-quic-backhaul").
500
Str("address", l.identity.Address).
501
Str("remote_addr", conn.RemoteAddr().String()).
503
- Msg("quic tunnel connected")
502
+ Msg("quic backhaul connected")
503
505
- recvDone, err := l.datagram.Bind(conn)
504
+ recvDone, err := l.datagram.BindBackhaul(conn)
505
if err != nil {
506
if ctx.Err() != nil {
507
return
508
}
509
log.Info().
510
Err(err).
512
- Str("component", "sdk-datagram-plane").
511
+ Str("component", "sdk-quic-backhaul").
512
Str("address", l.identity.Address).
514
- Msg("quic datagram plane did not bind cleanly; retrying")
513
+ Msg("quic backhaul did not bind cleanly; retrying")
514
if !utils.SleepOrDone(ctx, time.Second) {
515
return
516
}
@@ -581,7 +580,7 @@ func (l *listener) openReverseSession(ctx context.Context) (net.Conn, error) {
580
return wrapBufferedConn(conn, reader), nil
581
}
582
584
-func (l *listener) openQUICSession(ctx context.Context) (*quic.Conn, error) {
583
+func (l *listener) openQUICBackhaulSession(ctx context.Context) (*quic.Conn, error) {
584
lease, ok := l.leaseSnapshot()
585
if !ok || lease.accessToken == "" {
586
return nil, errors.New("access token is not available")
@@ -592,51 +591,12 @@ func (l *listener) openQUICSession(ctx context.Context) (*quic.Conn, error) {
591
if l.tlsConfig == nil {
592
return nil, errors.New("relay tls config is unavailable")
593
}
595
- tlsConf := l.tlsConfig.Clone()
596
- tlsConf.NextProtos = []string{"portal-tunnel"}
597
-
598
- quicConf := &quic.Config{
599
- EnableDatagrams: true,
600
- KeepAlivePeriod: 15 * time.Second,
601
- MaxIdleTimeout: 60 * time.Second,
602
- }
603
-
594
host := strings.TrimSpace(l.relayURL.Hostname())
595
if host == "" {
596
host = strings.TrimSpace(l.relayURL.Host)
597
}
598
dialAddr := net.JoinHostPort(host, fmt.Sprintf("%d", lease.sniPort))
609
- conn, err := quic.DialAddr(ctx, dialAddr, tlsConf, quicConf)
610
- if err != nil {
611
- return nil, fmt.Errorf("quic dial: %w", err)
612
- }
613
-
614
- stream, err := conn.OpenStreamSync(ctx)
615
- if err != nil {
616
- _ = conn.CloseWithError(1, "stream open failed")
617
- return nil, fmt.Errorf("open control stream: %w", err)
618
- }
619
-
620
- controlMsg := types.QUICControlMessage{
621
- AccessToken: lease.accessToken,
622
- }
623
- if err := json.NewEncoder(stream).Encode(controlMsg); err != nil {
624
- _ = conn.CloseWithError(1, "control write failed")
625
- return nil, fmt.Errorf("write control: %w", err)
626
- }
627
-
628
- _ = stream.SetReadDeadline(time.Now().Add(10 * time.Second))
629
- var resp types.QUICControlResponse
630
- if err := json.NewDecoder(io.LimitReader(stream, 4096)).Decode(&resp); err != nil {
631
- _ = conn.CloseWithError(1, "control read failed")
632
- return nil, fmt.Errorf("read control response: %w", err)
633
- }
634
- if !resp.OK {
635
- _ = conn.CloseWithError(1, resp.Error)
636
- return nil, fmt.Errorf("quic connect rejected: %s", resp.Error)
637
- }
638
-
639
- return conn, nil
599
+ return transport.DialQUICBackhaul(ctx, dialAddr, l.tlsConfig, lease.accessToken)
600
}
601
602
func (l *listener) runRenewLoop(ctx context.Context) error {
types/api.go
-9
@@ -107,15 +107,6 @@ type DiscoveryAnnounceResponse struct {
107
Accepted bool `json:"accepted"`
108
}
109
110
-type QUICControlMessage struct {
111
- AccessToken string `json:"access_token"`
112
-}
113
-
114
-type QUICControlResponse struct {
115
- OK bool `json:"ok"`
116
- Error string `json:"error,omitempty"`
117
-}
118
-
110
type RenewRequest struct {
111
AccessToken string `json:"access_token"`
112
TTL int `json:"ttl,omitempty"`
types/types.go
+5
-6
@@ -1,12 +1,11 @@
1
package types
2
3
const (
4
- ReleaseVersion = "v2.1.8"
5
- SDKVersion = "6"
6
- DiscoveryVersion = "7"
7
- PortalRelayRegistryURL = "https://raw.githubusercontent.com/gosuda/portal-tunnel/main/registry.json"
8
- OfficialReleaseBaseURL = "https://github.com/gosuda/portal-tunnel/releases/latest/download"
9
- OfficialReleaseDownloadURL = "https://github.com/gosuda/portal-tunnel/releases/download"
4
+ ReleaseVersion = "v2.1.8"
5
+ SDKVersion = "6"
6
+ DiscoveryVersion = "7"
7
+ PortalRelayRegistryURL = "https://raw.githubusercontent.com/gosuda/portal-tunnel/main/registry.json"
8
+ OfficialReleaseBaseURL = "https://github.com/gosuda/portal-tunnel/releases"
9
10
HeaderAccessToken = "X-Portal-Access-Token"
11
MarkerKeepalive = byte(0x00)