fix: resolve tunnel death hang, goroutine leak, and re-registration failures

Hee Sung Son committed Mar 9, 2026 at 18:56 UTC af7d046d94c81d0c9e02dafb23bacd43a31dadab
4 files changed +151 -15
portal/broker.go
+16 -7
@@ -17,6 +17,7 @@ var (
17 errLeaseDropped = errors.New("lease dropped")
18 errLeaseStopped = errors.New("lease stopped")
19 errBrokerFull = errors.New("broker ready queue full")
20 + errNoSessions = errors.New("no reverse sessions available")
21 )
22
23 type brokerState int
@@ -28,13 +29,14 @@ const (
29 )
30
31 type leaseBroker struct {
31 - notify chan struct{}
32 - leaseID string
33 - ready []*reverseSession
34 - idleInterval time.Duration
35 - readyLimit int
36 - state brokerState
37 - mu sync.Mutex
32 + notify chan struct{}
33 + leaseID string
34 + ready []*reverseSession
35 + idleInterval time.Duration
36 + readyLimit int
37 + totalSessions int
38 + state brokerState
39 + mu sync.Mutex
40 }
41
42 func newLeaseBroker(leaseID string, idleInterval time.Duration, readyLimit int) *leaseBroker {
@@ -67,6 +69,7 @@ func (b *leaseBroker) Offer(session *reverseSession) error {
69
70 session.StartIdle()
71 b.ready = append(b.ready, session)
72 + b.totalSessions++
73 b.signalLocked()
74 go b.watchSession(session)
75 return nil
@@ -98,6 +101,10 @@ func (b *leaseBroker) Claim(ctx context.Context) (*reverseSession, error) {
101 }
102 return session, nil
103 }
104 + if b.totalSessions == 0 {
105 + b.mu.Unlock()
106 + return nil, errNoSessions
107 + }
108 b.mu.Unlock()
109
110 select {
@@ -154,11 +161,13 @@ func (b *leaseBroker) watchSession(session *reverseSession) {
161 break
162 }
163 }
164 + b.totalSessions--
165 log.Info().
166 Str("component", "relay-server").
167 Str("lease_id", b.leaseID).
168 Str("remote_addr", session.RemoteAddr()).
169 Int("ready", len(b.ready)).
170 + Int("total_sessions", b.totalSessions).
171 Msg("sdk reverse disconnected")
172 b.signalLocked()
173 }
portal/server.go
+5
@@ -586,7 +586,11 @@ func (s *Server) authorizeLeaseToken(record *leaseRecord, token string) error {
586
587 func (s *Server) findLeaseByHostnameLocked(host string) *leaseRecord {
588 host = normalizeHostname(host)
589 + now := time.Now()
590 for _, lease := range s.leases {
591 + if now.After(lease.ExpiresAt) {
592 + continue
593 + }
594 for _, candidate := range lease.Hostnames {
595 if normalizeHostname(candidate) == host {
596 return lease
@@ -663,6 +667,7 @@ func (s *Server) handleSNIConn(conn net.Conn) {
667 }
668
669 bridgeConns(wrappedConn, session.Conn())
670 + _ = session.Close()
671 }
672
673 func (s *Server) bridgeToFallback(conn net.Conn) {
sdk/client.go
+9
@@ -234,6 +234,7 @@ func (c *Client) Listen(ctx context.Context, req ListenRequest) (*Listener, erro
234 baseContext: func() context.Context { return listenerCtx },
235 ctxDone: listenerCtx.Done(),
236 cancel: cancel,
237 + name: strings.TrimSpace(req.Name),
238 leaseID: registerResp.LeaseID,
239 hostnames: registerResp.Hostnames,
240 metadata: registerResp.Metadata,
@@ -298,6 +299,14 @@ func (c *Client) doJSON(ctx context.Context, method, path string, payload any, o
299 return json.Unmarshal(envelope.Data, out)
300 }
301
302 +func (c *Client) registerLease(ctx context.Context, req types.RegisterRequest) (types.RegisterResponse, error) {
303 + var resp types.RegisterResponse
304 + if err := c.doJSON(ctx, http.MethodPost, types.PathSDKRegister, req, &resp); err != nil {
305 + return types.RegisterResponse{}, err
306 + }
307 + return resp, nil
308 +}
309 +
310 func (c *Client) renewLease(ctx context.Context, leaseID, reverseToken string, ttl time.Duration) error {
311 return c.doJSON(ctx, http.MethodPost, types.PathSDKRenew, types.RenewRequest{
312 LeaseID: leaseID,
sdk/listener.go
+121 -8
@@ -10,6 +10,9 @@ import (
10 "sync"
11 "time"
12
13 + "github.com/rs/zerolog/log"
14 +
15 + "github.com/gosuda/portal/v2/portal/keyless"
16 "github.com/gosuda/portal/v2/types"
17 )
18
@@ -31,6 +34,7 @@ type Listener struct {
34 client *Client
35 signal chan struct{}
36 accepted chan net.Conn
37 + name string
38 leaseID string
39 reverseToken string
40 hostnames []string
@@ -60,13 +64,18 @@ func (l *Listener) Close() error {
64 l.closeOnce.Do(func() {
65 l.cancel()
66
67 + l.mu.Lock()
68 + leaseID := l.leaseID
69 + tlsCloser := l.tlsCloser
70 + l.mu.Unlock()
71 +
72 ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
73 defer cancel()
65 - if err := l.client.unregisterLease(ctx, l.leaseID, l.reverseToken); err != nil {
74 + if err := l.client.unregisterLease(ctx, leaseID, l.reverseToken); err != nil {
75 closeErr = err
76 }
68 - if l.tlsCloser != nil {
69 - closeErr = errors.Join(closeErr, l.tlsCloser.Close())
77 + if tlsCloser != nil {
78 + closeErr = errors.Join(closeErr, tlsCloser.Close())
79 }
80 })
81 return closeErr
@@ -77,20 +86,30 @@ func (l *Listener) Addr() net.Addr {
86 }
87
88 func (l *Listener) LeaseID() string {
89 + l.mu.Lock()
90 + defer l.mu.Unlock()
91 return l.leaseID
92 }
93
94 func (l *Listener) Hostnames() []string {
95 + l.mu.Lock()
96 + defer l.mu.Unlock()
97 return l.hostnames
98 }
99
100 func (l *Listener) Metadata() types.LeaseMetadata {
101 + l.mu.Lock()
102 + defer l.mu.Unlock()
103 return l.metadata
104 }
105
106 func (l *Listener) PublicURLs() []string {
92 - urls := make([]string, 0, len(l.hostnames))
93 - for _, host := range l.hostnames {
107 + l.mu.Lock()
108 + hostnames := l.hostnames
109 + l.mu.Unlock()
110 +
111 + urls := make([]string, 0, len(hostnames))
112 + for _, host := range hostnames {
113 urls = append(urls, "https://"+host)
114 }
115 return urls
@@ -125,14 +144,55 @@ func (l *Listener) runRenewLoop() {
144 ticker := time.NewTicker(interval)
145 defer ticker.Stop()
146
147 + var consecutiveFailures int
148 +
149 for {
150 select {
151 case <-l.ctxDone:
152 return
153 case <-ticker.C:
154 + l.mu.Lock()
155 + leaseID := l.leaseID
156 + l.mu.Unlock()
157 +
158 ctx, cancel := context.WithTimeout(l.context(), 10*time.Second)
134 - _ = l.client.renewLease(ctx, l.leaseID, l.reverseToken, l.leaseTTL)
159 + err := l.client.renewLease(ctx, leaseID, l.reverseToken, l.leaseTTL)
160 cancel()
161 +
162 + if err != nil {
163 + if isLeaseNotFound(err) {
164 + log.Warn().
165 + Str("component", "sdk-listener").
166 + Str("lease_id", leaseID).
167 + Msg("lease not found on relay, attempting re-registration")
168 + if reregErr := l.reregister(); reregErr != nil {
169 + log.Error().Err(reregErr).
170 + Str("component", "sdk-listener").
171 + Msg("lease re-registration failed")
172 + } else {
173 + consecutiveFailures = 0
174 + log.Info().
175 + Str("component", "sdk-listener").
176 + Str("lease_id", l.LeaseID()).
177 + Strs("hostnames", l.Hostnames()).
178 + Msg("lease re-registered successfully")
179 + continue
180 + }
181 + }
182 +
183 + consecutiveFailures++
184 + event := log.Warn()
185 + if consecutiveFailures >= 3 {
186 + event = log.Error()
187 + }
188 + event.Err(err).
189 + Str("component", "sdk-listener").
190 + Str("lease_id", l.LeaseID()).
191 + Int("consecutive_failures", consecutiveFailures).
192 + Msg("lease renewal failed")
193 + } else {
194 + consecutiveFailures = 0
195 + }
196 }
197 }
198 }
@@ -141,7 +201,10 @@ func (l *Listener) runSession() {
201 defer l.releaseSessionSlot()
202
203 sessionCtx := l.context()
144 - conn, err := l.client.openReverseSession(sessionCtx, l.leaseID, l.reverseToken)
204 + l.mu.Lock()
205 + leaseID := l.leaseID
206 + l.mu.Unlock()
207 + conn, err := l.client.openReverseSession(sessionCtx, leaseID, l.reverseToken)
208 if err != nil {
209 sleepOrDone(sessionCtx, time.Second)
210 return
@@ -176,7 +239,10 @@ func (l *Listener) awaitActivation(conn net.Conn) error {
239 }
240
241 func (l *Listener) activate(conn net.Conn) error {
179 - tlsConn := tls.Server(conn, l.tlsConfig.Clone())
242 + l.mu.Lock()
243 + tlsCfg := l.tlsConfig
244 + l.mu.Unlock()
245 + tlsConn := tls.Server(conn, tlsCfg.Clone())
246 handshakeCtx, cancel := context.WithTimeout(l.context(), l.client.handshakeTimeout)
247 defer cancel()
248 if err := tlsConn.HandshakeContext(handshakeCtx); err != nil {
@@ -192,6 +258,53 @@ func (l *Listener) activate(conn net.Conn) error {
258 }
259 }
260
261 +func (l *Listener) reregister() error {
262 + ctx, cancel := context.WithTimeout(l.context(), 10*time.Second)
263 + defer cancel()
264 +
265 + l.mu.Lock()
266 + hostnames := l.hostnames
267 + l.mu.Unlock()
268 +
269 + resp, err := l.client.registerLease(ctx, types.RegisterRequest{
270 + Name: l.name,
271 + Hostnames: hostnames,
272 + Metadata: l.metadata,
273 + ReverseToken: l.reverseToken,
274 + TLS: true,
275 + TTLSeconds: int(l.leaseTTL / time.Second),
276 + })
277 + if err != nil {
278 + return err
279 + }
280 +
281 + tlsConf, tlsCloser, err := keyless.BuildClientTLSConfig(l.client.baseURL.String(), resp.Hostnames)
282 + if err != nil {
283 + _ = l.client.unregisterLease(ctx, resp.LeaseID, l.reverseToken)
284 + return err
285 + }
286 +
287 + l.mu.Lock()
288 + oldCloser := l.tlsCloser
289 + l.leaseID = resp.LeaseID
290 + l.hostnames = resp.Hostnames
291 + l.metadata = resp.Metadata
292 + l.tlsConfig = tlsConf
293 + l.tlsCloser = tlsCloser
294 + l.mu.Unlock()
295 +
296 + if oldCloser != nil {
297 + _ = oldCloser.Close()
298 + }
299 +
300 + l.notify()
301 + return nil
302 +}
303 +
304 +func isLeaseNotFound(err error) bool {
305 + return errors.Is(err, &types.APIRequestError{Code: types.APIErrorCodeLeaseNotFound})
306 +}
307 +
308 func (l *Listener) reserveSessionSlot() bool {
309 l.mu.Lock()
310 defer l.mu.Unlock()