fix(webclient): remove mutex connect for sharing.

Hee Sung Son committed Nov 3, 2025 at 12:15 UTC 710cba510c656fb44ccf78f5c3ab82b50d4322ed
1 file changed +4 -222
cmd/webclient/main_js.go
+4 -222
@@ -21,9 +21,7 @@ import (
21
22 "github.com/gorilla/websocket"
23 "github.com/gosuda/portal/cmd/webclient/httpjs"
24 - "github.com/gosuda/portal/portal/core/cryptoops"
24 "github.com/gosuda/portal/sdk"
26 - "github.com/hashicorp/yamux"
25 "github.com/rs/zerolog"
26 "github.com/rs/zerolog/log"
27 "golang.org/x/net/idna"
@@ -32,229 +30,17 @@ import (
30 var (
31 bootstrapServers = []string{"ws://localhost:4017/relay", "wss://portal.gosuda.org/relay"}
32 rdClient *sdk.RDClient
35 -
36 - // Connection pool for reusing encrypted channels with yamux multiplexing
37 - muxSessions sync.Map // map[string]*muxSession
38 - muxSessionsLock sync.Mutex
39 -)
40 -
41 -const (
42 - // poolCleanupInterval is how often to clean up stale mux sessions
43 - poolCleanupInterval = 1 * time.Minute
44 - // poolIdleTimeout is how long a mux session can be idle before cleanup
45 - poolIdleTimeout = 5 * time.Minute
33 )
34
48 -// muxSession wraps a yamux session over an encrypted connection
49 -type muxSession struct {
50 - session *yamux.Session
51 - leaseID string
52 - cred *cryptoops.Credential
53 - refCount int
54 - lastUsed time.Time
55 - mu sync.Mutex
56 - closed bool
57 -}
58 -
59 -// getOrCreateMuxSession gets an existing yamux session or creates a new one with E2EE handshake
60 -func getOrCreateMuxSession(ctx context.Context, leaseID string) (*muxSession, error) {
61 - // Try to get existing session from pool
62 - if val, ok := muxSessions.Load(leaseID); ok {
63 - session := val.(*muxSession)
64 - session.mu.Lock()
65 - defer session.mu.Unlock()
66 -
67 - // Check if session is still alive
68 - if !session.closed && !session.session.IsClosed() {
69 - session.refCount++
70 - session.lastUsed = time.Now()
71 -
72 - log.Info().
73 - Str("leaseID", leaseID).
74 - Int("refCount", session.refCount).
75 - Msg("[WebClient] Reusing existing mux session (no key exchange)")
76 -
77 - return session, nil
78 - }
79 -
80 - // Session is dead, remove it
81 - log.Warn().Str("leaseID", leaseID).Msg("[WebClient] Existing mux session is closed, creating new one")
82 - muxSessions.Delete(leaseID)
83 - }
84 -
85 - // Create new session with double-checked locking
86 - muxSessionsLock.Lock()
87 - defer muxSessionsLock.Unlock()
88 -
89 - // Double-check after acquiring lock
90 - if val, ok := muxSessions.Load(leaseID); ok {
91 - session := val.(*muxSession)
92 - session.mu.Lock()
93 - defer session.mu.Unlock()
94 - if !session.closed && !session.session.IsClosed() {
95 - session.refCount++
96 - session.lastUsed = time.Now()
97 - return session, nil
98 - }
99 - muxSessions.Delete(leaseID)
100 - }
101 -
102 - log.Info().
103 - Str("leaseID", leaseID).
104 - Msg("[WebClient] Creating new encrypted channel with E2EE key exchange")
105 -
106 - // Create new credential for this connection
107 - cred := sdk.NewCredential()
108 -
109 - // Establish encrypted connection (this does the expensive key exchange)
110 - rdConn, err := rdClient.Dial(cred, leaseID, "http/1.1")
111 - if err != nil {
112 - return nil, fmt.Errorf("failed to dial: %w", err)
113 - }
114 -
115 - // Create yamux client session on top of the encrypted connection
116 - yamuxConfig := yamux.DefaultConfig()
117 - yamuxConfig.Logger = nil // Disable yamux logging
118 - yamuxSession, err := yamux.Client(rdConn, yamuxConfig)
119 - if err != nil {
120 - rdConn.Close()
121 - return nil, fmt.Errorf("failed to create yamux session: %w", err)
122 - }
123 -
124 - session := &muxSession{
125 - session: yamuxSession,
126 - leaseID: leaseID,
127 - cred: cred,
128 - refCount: 1,
129 - lastUsed: time.Now(),
130 - closed: false,
131 - }
132 -
133 - muxSessions.Store(leaseID, session)
134 -
135 - log.Info().
136 - Str("leaseID", leaseID).
137 - Msg("[WebClient] Yamux session created successfully over encrypted channel")
138 -
139 - return session, nil
140 -}
141 -
142 -// openStream opens a new stream on the mux session
143 -func (m *muxSession) openStream() (net.Conn, error) {
144 - m.mu.Lock()
145 - defer m.mu.Unlock()
146 -
147 - if m.closed || m.session.IsClosed() {
148 - return nil, fmt.Errorf("mux session is closed")
149 - }
150 -
151 - stream, err := m.session.OpenStream()
152 - if err != nil {
153 - return nil, fmt.Errorf("failed to open yamux stream: %w", err)
154 - }
155 -
156 - log.Debug().
157 - Str("leaseID", m.leaseID).
158 - Uint32("streamID", stream.StreamID()).
159 - Msg("[WebClient] Opened new stream on existing mux session")
160 -
161 - return stream, nil
162 -}
163 -
164 -// release decrements the reference count
165 -func (m *muxSession) release() {
166 - m.mu.Lock()
167 - defer m.mu.Unlock()
168 -
169 - m.refCount--
170 - m.lastUsed = time.Now()
171 -
172 - log.Debug().
173 - Str("leaseID", m.leaseID).
174 - Int("refCount", m.refCount).
175 - Msg("[WebClient] Released mux session reference")
176 -}
177 -
178 -// close closes the mux session
179 -func (m *muxSession) close() error {
180 - m.mu.Lock()
181 - defer m.mu.Unlock()
182 -
183 - if m.closed {
184 - return nil
185 - }
186 -
187 - m.closed = true
188 - return m.session.Close()
189 -}
190 -
191 -// cleanupIdleMuxSessions periodically cleans up idle mux sessions
192 -func cleanupIdleMuxSessions() {
193 - ticker := time.NewTicker(poolCleanupInterval)
194 - defer ticker.Stop()
195 -
196 - for range ticker.C {
197 - now := time.Now()
198 - var toDelete []string
199 -
200 - muxSessions.Range(func(key, value interface{}) bool {
201 - leaseID := key.(string)
202 - session := value.(*muxSession)
203 -
204 - session.mu.Lock()
205 - idle := now.Sub(session.lastUsed)
206 - shouldDelete := (session.refCount == 0 && idle > poolIdleTimeout) || session.closed || session.session.IsClosed()
207 - session.mu.Unlock()
208 -
209 - if shouldDelete {
210 - toDelete = append(toDelete, leaseID)
211 - }
212 -
213 - return true
214 - })
215 -
216 - for _, leaseID := range toDelete {
217 - if val, ok := muxSessions.LoadAndDelete(leaseID); ok {
218 - session := val.(*muxSession)
219 - session.close()
220 - log.Info().
221 - Str("leaseID", leaseID).
222 - Msg("[WebClient] Cleaned up idle mux session")
223 - }
224 - }
225 - }
226 -}
227 -
35 var rdDialer = func(ctx context.Context, network, address string) (net.Conn, error) {
36 address = strings.TrimSuffix(address, ":80")
37 address = strings.TrimSuffix(address, ":443")
231 -
232 - // Get or create mux session (does key exchange only once)
233 - session, err := getOrCreateMuxSession(ctx, address)
234 - if err != nil {
235 - return nil, fmt.Errorf("failed to get mux session: %w", err)
236 - }
237 -
238 - // Open a new stream on the mux session (no key exchange, just new yamux stream)
239 - stream, err := session.openStream()
38 + cred := sdk.NewCredential()
39 + conn, err := rdClient.Dial(cred, address, "http/1.1")
40 if err != nil {
241 - // If stream opening fails, try to create a new session
242 - log.Warn().Err(err).Str("leaseID", address).Msg("[WebClient] Failed to open stream, removing stale session")
243 - muxSessions.Delete(address)
244 -
245 - // Retry with a fresh session
246 - session, err = getOrCreateMuxSession(ctx, address)
247 - if err != nil {
248 - return nil, fmt.Errorf("failed to get mux session after retry: %w", err)
249 - }
250 -
251 - stream, err = session.openStream()
252 - if err != nil {
253 - return nil, fmt.Errorf("failed to open stream after retry: %w", err)
254 - }
41 + return nil, err
42 }
256 -
257 - return stream, nil
43 + return conn, nil
44 }
45
46 var client = &http.Client{
@@ -728,10 +514,6 @@ func main() {
514 }
515 defer rdClient.Close()
516
731 - // Start cleanup goroutine for idle mux sessions
732 - go cleanupIdleMuxSessions()
733 - log.Info().Msg("[WebClient] Started mux session cleanup goroutine")
734 -
517 // Initialize WebSocket manager
518 wsManager := NewWebSocketManager()
519 proxy := &Proxy{