enhance proxy connection handling, and restructure SDK client

gosunuts committed Feb 24, 2026 at 18:50 UTC 52e700d15c5cd1e4cb4976e0fa6734fdf59191b7
8 files changed +190 -86
cmd/demo-app/main.go
+20
@@ -15,6 +15,7 @@ import (
15 "time"
16
17 "github.com/rs/zerolog/log"
18 + "golang.org/x/net/websocket"
19
20 "gosuda.org/portal/sdk"
21 )
@@ -98,6 +99,25 @@ func runDemo() error {
99 }
100 })
101
102 + // WebSocket echo endpoint
103 + mux.Handle("/ws", websocket.Handler(func(conn *websocket.Conn) {
104 + defer conn.Close()
105 + for {
106 + var msg string
107 + if err := websocket.Message.Receive(conn, &msg); err != nil {
108 + if err.Error() != "EOF" {
109 + log.Error().Err(err).Msg("websocket read error")
110 + }
111 + break
112 + }
113 + log.Debug().Str("msg", msg).Msg("websocket received")
114 + if err := websocket.Message.Send(conn, "echo: "+msg); err != nil {
115 + log.Error().Err(err).Msg("websocket write error")
116 + break
117 + }
118 + }
119 + }))
120 +
121 // Test endpoint for multiple Set-Cookie headers
122 // Note: HttpOnly cookies cannot be set via Service Worker (browser security limitation)
123 mux.HandleFunc("/api/test-cookies", func(w http.ResponseWriter, r *http.Request) {
cmd/portal-tunnel/main.go
+53 -37
@@ -225,51 +225,30 @@ func splitCSV(raw string) []string {
225 return out
226 }
227
228 -// proxyConnection proxies data between relay and local service.
229 -// If local service is not available, it retries with backoff instead of failing immediately.
228 +// proxyConnection proxies data between relay and local service using raw TCP.
229 +// It ensures complete data transfer before closing connections.
230 func proxyConnection(ctx context.Context, localAddr string, relayConn net.Conn) error {
231 defer relayConn.Close()
232
233 - // Try to connect to local service with retry
234 - var localConn net.Conn
235 - var err error
236 -
237 - maxRetries := 30
238 - retryDelay := 500 * time.Millisecond
239 -
240 - for i := 0; i < maxRetries; i++ {
241 - select {
242 - case <-ctx.Done():
243 - return ctx.Err()
244 - default:
245 - }
246 -
247 - dialer := &net.Dialer{Timeout: 5 * time.Second}
248 - localConn, err = dialer.DialContext(ctx, "tcp", localAddr)
249 - if err == nil {
250 - break
251 - }
252 -
253 - if i == 0 {
254 - log.Warn().
255 - Str("local_addr", localAddr).
256 - Err(err).
257 - Msg("Local service not ready, retrying...")
258 - }
259 -
260 - time.Sleep(retryDelay)
261 - }
262 -
233 + // Try to connect to local service (no retry)
234 + dialer := &net.Dialer{Timeout: 5 * time.Second}
235 + localConn, err := dialer.DialContext(ctx, "tcp", localAddr)
236 if err != nil {
264 - return fmt.Errorf("local service unavailable: %w", err)
237 + log.Debug().
238 + Str("local_addr", localAddr).
239 + Err(err).
240 + Msg("Local service unavailable, returning service unavailable page")
241 + return writeEmptyHTTPResponse(relayConn)
242 }
266 -
243 defer localConn.Close()
244
245 log.Info().Str("local_addr", localAddr).Msg("Connected to local service")
246
247 + // Use bidirectional copy with proper error handling
248 errCh := make(chan error, 2)
249 stopCh := make(chan struct{})
250 +
251 + // Context cancellation handler
252 go func() {
253 select {
254 case <-ctx.Done():
@@ -279,24 +258,61 @@ func proxyConnection(ctx context.Context, localAddr string, relayConn net.Conn)
258 }
259 }()
260
261 + // Relay -> Local
262 go func() {
263 buf := *bufferPool.Get().(*[]byte)
264 defer bufferPool.Put(&buf)
265 _, err := io.CopyBuffer(localConn, relayConn, buf)
266 + if err != nil {
267 + log.Debug().Err(err).Msg("relay->local copy ended")
268 + }
269 + // Shut down localConn write side to signal EOF to local service
270 + if tcpConn, ok := localConn.(*net.TCPConn); ok {
271 + tcpConn.CloseWrite()
272 + }
273 errCh <- err
274 }()
275
276 + // Local -> Relay
277 go func() {
278 buf := *bufferPool.Get().(*[]byte)
279 defer bufferPool.Put(&buf)
280 _, err := io.CopyBuffer(relayConn, localConn, buf)
281 + if err != nil {
282 + log.Debug().Err(err).Msg("local->relay copy ended")
283 + }
284 errCh <- err
285 }()
286
296 - err = <-errCh
287 + // Wait for both directions to finish
288 + var firstErr error
289 + for i := 0; i < 2; i++ {
290 + if err := <-errCh; err != nil && firstErr == nil {
291 + firstErr = err
292 + }
293 + }
294 +
295 close(stopCh)
298 - relayConn.Close()
299 - <-errCh
296 + return firstErr
297 +}
298
299 +// writeEmptyHTTPResponse writes an HTML response indicating the service is unavailable.
300 +// Used when the local service is unavailable to avoid showing browser error pages.
301 +func writeEmptyHTTPResponse(conn net.Conn) error {
302 + htmlBody := `<!DOCTYPE html>
303 +<html>
304 +<head><title>Service Unavailable</title></head>
305 +<body style="font-family:sans-serif;text-align:center;padding:50px;">
306 +<h1>🔌 Service Unavailable</h1>
307 +<p>The local service is not currently running.</p>
308 +<p>Please start your local application and refresh this page.</p>
309 +</body>
310 +</html>`
311 + response := fmt.Sprintf("HTTP/1.1 503 Service Unavailable\r\n"+
312 + "Content-Type: text/html; charset=utf-8\r\n"+
313 + "Content-Length: %d\r\n"+
314 + "Connection: close\r\n"+
315 + "\r\n%s", len(htmlBody), htmlBody)
316 + _, err := conn.Write([]byte(response))
317 return err
318 }
cmd/relay-server/frontend.go
+8 -1
@@ -256,6 +256,13 @@ func convertLeaseEntriesToRows(serv *portal.RelayServer, admin *Admin) []leaseRo
256 bps = bpsMgr.GetBPSLimit(identityID)
257 }
258
259 + metadataStr := ""
260 + if b, err := json.Marshal(metadata); err == nil {
261 + metadataStr = string(b)
262 + } else {
263 + log.Warn().Err(err).Str("lease_id", identityID).Msg("[Frontend] Failed to marshal lease metadata")
264 + }
265 +
266 row := leaseRow{
267 Peer: identityID,
268 Name: name,
@@ -269,7 +276,7 @@ func convertLeaseEntriesToRows(serv *portal.RelayServer, admin *Admin) []leaseRo
276 Link: link,
277 StaleRed: !connected && since >= 15*time.Second,
278 Hide: leaseEntry.ParsedMetadata != nil && leaseEntry.ParsedMetadata.Hide,
272 - Metadata: "",
279 + Metadata: metadataStr,
280 BPS: bps,
281 }
282
cmd/relay-server/registry.go
+11 -28
@@ -13,6 +13,7 @@ import (
13 "gosuda.org/portal/cmd/relay-server/manager"
14 "gosuda.org/portal/portal"
15 "gosuda.org/portal/portal/utils/sni"
16 + "gosuda.org/portal/sdk"
17 "gosuda.org/portal/utils"
18 )
19
@@ -33,24 +34,6 @@ func NewSDKRegistry(server *portal.RelayServer, sniRouter *sni.Router) *SDKRegis
34 }
35 }
36
36 -// RegisterRequest represents an SDK lease registration request
37 -type RegisterRequest struct {
38 - LeaseID string `json:"lease_id"`
39 - Name string `json:"name"`
40 - Address string `json:"address"` // Backend address for TCP connection
41 - Metadata portal.Metadata `json:"metadata"`
42 - TLSEnabled bool `json:"tls_enabled"` // Whether the backend handles TLS termination
43 - ReverseToken string `json:"reverse_token"`
44 -}
45 -
46 -// RegisterResponse represents an SDK lease registration response
47 -type RegisterResponse struct {
48 - Success bool `json:"success"`
49 - Message string `json:"message,omitempty"`
50 - LeaseID string `json:"lease_id,omitempty"`
51 - PublicURL string `json:"public_url,omitempty"`
52 -}
53 -
37 // HandleRegister handles SDK lease registration requests
38 func (r *SDKRegistry) HandleRegister(w http.ResponseWriter, req *http.Request) {
39 if req.Method != http.MethodPost {
@@ -59,10 +42,10 @@ func (r *SDKRegistry) HandleRegister(w http.ResponseWriter, req *http.Request) {
42 return
43 }
44
62 - var registerReq RegisterRequest
45 + var registerReq sdk.RegisterRequest
46 if err := json.NewDecoder(req.Body).Decode(&registerReq); err != nil {
47 log.Error().Err(err).Msg("[Registry] Failed to decode registration request")
65 - writeJSON(w, RegisterResponse{
48 + writeJSON(w, sdk.RegisterResponse{
49 Success: false,
50 Message: "invalid request body",
51 })
@@ -71,7 +54,7 @@ func (r *SDKRegistry) HandleRegister(w http.ResponseWriter, req *http.Request) {
54
55 // Validate request
56 if registerReq.LeaseID == "" {
74 - writeJSON(w, RegisterResponse{
57 + writeJSON(w, sdk.RegisterResponse{
58 Success: false,
59 Message: "lease_id is required",
60 })
@@ -79,7 +62,7 @@ func (r *SDKRegistry) HandleRegister(w http.ResponseWriter, req *http.Request) {
62 }
63
64 if registerReq.Name == "" {
82 - writeJSON(w, RegisterResponse{
65 + writeJSON(w, sdk.RegisterResponse{
66 Success: false,
67 Message: "name is required",
68 })
@@ -87,14 +70,14 @@ func (r *SDKRegistry) HandleRegister(w http.ResponseWriter, req *http.Request) {
70 }
71
72 if registerReq.Address == "" {
90 - writeJSON(w, RegisterResponse{
73 + writeJSON(w, sdk.RegisterResponse{
74 Success: false,
75 Message: "address is required",
76 })
77 return
78 }
79 if strings.TrimSpace(registerReq.ReverseToken) == "" {
97 - writeJSON(w, RegisterResponse{
80 + writeJSON(w, sdk.RegisterResponse{
81 Success: false,
82 Message: "reverse_token is required",
83 })
@@ -103,7 +86,7 @@ func (r *SDKRegistry) HandleRegister(w http.ResponseWriter, req *http.Request) {
86
87 resolvedAddr, err := resolveLeaseAddress(req, registerReq.Address)
88 if err != nil {
106 - writeJSON(w, RegisterResponse{
89 + writeJSON(w, sdk.RegisterResponse{
90 Success: false,
91 Message: err.Error(),
92 })
@@ -123,7 +106,7 @@ func (r *SDKRegistry) HandleRegister(w http.ResponseWriter, req *http.Request) {
106
107 // Register with lease manager
108 if !r.server.GetLeaseManager().UpdateLease(lease) {
126 - writeJSON(w, RegisterResponse{
109 + writeJSON(w, sdk.RegisterResponse{
110 Success: false,
111 Message: "failed to register lease (name conflict or policy violation)",
112 })
@@ -133,7 +116,7 @@ func (r *SDKRegistry) HandleRegister(w http.ResponseWriter, req *http.Request) {
116 if err := r.registerSNIRoute(registerReq.LeaseID, registerReq.Name, resolvedAddr); err != nil {
117 // Keep lease and route state consistent on partial failure.
118 r.server.GetLeaseManager().DeleteLease(registerReq.LeaseID)
136 - writeJSON(w, RegisterResponse{
119 + writeJSON(w, sdk.RegisterResponse{
120 Success: false,
121 Message: fmt.Sprintf("failed to register SNI route: %v", err),
122 })
@@ -151,7 +134,7 @@ func (r *SDKRegistry) HandleRegister(w http.ResponseWriter, req *http.Request) {
134 // Build public URL
135 publicURL := utils.ServicePublicURL(flagPortalURL, registerReq.Name)
136
154 - writeJSON(w, RegisterResponse{
137 + writeJSON(w, sdk.RegisterResponse{
138 Success: true,
139 LeaseID: registerReq.LeaseID,
140 PublicURL: publicURL,
portal/reverse_hub.go
+8 -5
@@ -12,11 +12,12 @@ import (
12 )
13
14 const (
15 - ReverseStartMarker = byte(0x01)
16 - ReverseQueueSize = 64
17 - ReverseAcquireWait = 2 * time.Second
18 - ReverseHTTPWait = 1500 * time.Millisecond
19 - ReverseSNIAcquireWait = 2 * time.Second
15 + ReverseStartMarker = byte(0x01)
16 + ReverseQueueSize = 64
17 + ReverseAcquireWait = 2 * time.Second
18 + ReverseHTTPWait = 1500 * time.Millisecond
19 + ReverseSNIAcquireWait = 2 * time.Second
20 + ReverseHandleConnectDelay = 2 * time.Second
21 )
22
23 type ReverseConn struct {
@@ -193,11 +194,13 @@ func (h *ReverseHub) HandleConnect(ws *websocket.Conn) {
194 }
195 if leaseID == "" {
196 log.Warn().Msg("[ReverseHub] Missing lease_id on reverse connect")
197 + time.Sleep(ReverseHandleConnectDelay)
198 ws.Close()
199 return
200 }
201 if !h.isAuthorized(leaseID, token) {
202 log.Warn().Str("lease_id", leaseID).Msg("[ReverseHub] Unauthorized reverse connect")
203 + time.Sleep(ReverseHandleConnectDelay)
204 ws.Close()
205 return
206 }
sdk/client.go renamed
sdk/listener.go
+56 -15
@@ -253,6 +253,20 @@ func (l *Listener) keepaliveLoop() {
253 return
254 case <-ticker.C:
255 if err := l.sendKeepalive(); err != nil {
256 + if isLeaseNotFoundError(err) {
257 + if rerr := l.reRegisterLease(); rerr != nil {
258 + log.Warn().
259 + Err(rerr).
260 + Str("lease_id", l.lease.ID).
261 + Msg("[SDK] Relay keepalive failed and re-register failed")
262 + } else {
263 + log.Info().
264 + Str("lease_id", l.lease.ID).
265 + Str("name", l.lease.Name).
266 + Msg("[SDK] Lease re-registered after relay reset")
267 + }
268 + continue
269 + }
270 log.Warn().Err(err).Str("lease_id", l.lease.ID).Msg("[SDK] Relay keepalive failed")
271 }
272 }
@@ -381,14 +395,7 @@ func (l *Listener) waitForReverseStart(conn net.Conn) error {
395 }
396
397 func (l *Listener) registerWithRelay(tunnelAddr string) error {
384 - reqBody := struct {
385 - LeaseID string `json:"lease_id"`
386 - Name string `json:"name"`
387 - Address string `json:"address"`
388 - Metadata portal.Metadata `json:"metadata"`
389 - TLSEnabled bool `json:"tls_enabled"`
390 - ReverseToken string `json:"reverse_token"`
391 - }{
398 + reqBody := RegisterRequest{
399 LeaseID: l.lease.ID,
400 Name: l.lease.Name,
401 Address: tunnelAddr,
@@ -401,25 +408,30 @@ func (l *Listener) registerWithRelay(tunnelAddr string) error {
408 }
409
410 func (l *Listener) unregisterFromRelay() error {
404 - reqBody := struct {
405 - LeaseID string `json:"lease_id"`
406 - }{
411 + reqBody := UnregisterRequest{
412 LeaseID: l.lease.ID,
413 }
414 return l.postJSON("/api/unregister", reqBody)
415 }
416
417 func (l *Listener) sendKeepalive() error {
413 - reqBody := struct {
414 - LeaseID string `json:"lease_id"`
415 - ReverseToken string `json:"reverse_token"`
416 - }{
418 + reqBody := RenewRequest{
419 LeaseID: l.lease.ID,
420 ReverseToken: l.lease.ReverseToken,
421 }
422 return l.postJSON("/api/renew", reqBody)
423 }
424
425 +func (l *Listener) reRegisterLease() error {
426 + l.mu.RLock()
427 + addr := strings.TrimSpace(l.lease.Address)
428 + l.mu.RUnlock()
429 + if addr == "" {
430 + return fmt.Errorf("lease address is empty; cannot re-register")
431 + }
432 + return l.registerWithRelay(addr)
433 +}
434 +
435 func (l *Listener) postJSON(path string, body any) error {
436 payload, err := json.Marshal(body)
437 if err != nil {
@@ -437,9 +449,38 @@ func (l *Listener) postJSON(path string, body any) error {
449 data, _ := io.ReadAll(resp.Body)
450 return fmt.Errorf("POST %s failed: status=%d body=%s", path, resp.StatusCode, strings.TrimSpace(string(data)))
451 }
452 +
453 + data, _ := io.ReadAll(resp.Body)
454 + if len(data) == 0 {
455 + return nil
456 + }
457 +
458 + var apiResp struct {
459 + Success *bool `json:"success"`
460 + Message string `json:"message"`
461 + }
462 + if err := json.Unmarshal(data, &apiResp); err != nil {
463 + // Non-JSON success payloads are treated as successful.
464 + return nil
465 + }
466 + if apiResp.Success != nil && !*apiResp.Success {
467 + msg := strings.TrimSpace(apiResp.Message)
468 + if msg == "" {
469 + msg = strings.TrimSpace(string(data))
470 + }
471 + return fmt.Errorf("POST %s rejected: %s", path, msg)
472 + }
473 +
474 return nil
475 }
476
477 +func isLeaseNotFoundError(err error) bool {
478 + if err == nil {
479 + return false
480 + }
481 + return strings.Contains(strings.ToLower(err.Error()), "lease not found")
482 +}
483 +
484 func normalizeRelayAPIURL(raw string) (string, error) {
485 raw = strings.TrimSpace(raw)
486 if raw == "" {
sdk/types.go
+34
@@ -5,6 +5,8 @@ import (
5 "errors"
6 "io"
7 "time"
8 +
9 + "gosuda.org/portal/portal"
10 )
11
12 var (
@@ -149,3 +151,35 @@ func WithHide(hide bool) MetadataOption {
151 m.Hide = hide
152 }
153 }
154 +
155 +// API Types for /api/ endpoints
156 +// These types are shared between SDK and relay server
157 +type RegisterRequest struct {
158 + LeaseID string `json:"lease_id"`
159 + Name string `json:"name"`
160 + Address string `json:"address"` // Backend address for TCP connection
161 + Metadata portal.Metadata `json:"metadata"`
162 + TLSEnabled bool `json:"tls_enabled"` // Whether the backend handles TLS termination
163 + ReverseToken string `json:"reverse_token"`
164 +}
165 +
166 +type RegisterResponse struct {
167 + Success bool `json:"success"`
168 + Message string `json:"message,omitempty"`
169 + LeaseID string `json:"lease_id,omitempty"`
170 + PublicURL string `json:"public_url,omitempty"`
171 +}
172 +
173 +type UnregisterRequest struct {
174 + LeaseID string `json:"lease_id"`
175 +}
176 +
177 +type RenewRequest struct {
178 + LeaseID string `json:"lease_id"`
179 + ReverseToken string `json:"reverse_token"`
180 +}
181 +
182 +type APIResponse struct {
183 + Success bool `json:"success"`
184 + Message string `json:"message,omitempty"`
185 +}