feat(logging): add zerolog logging and fix LocalID/RemoteID recursion
Add structured logging using zerolog library to client.go and relay.go for improved debugging and error tracking. Fix infinite recursion bug in IncomingConn.LocalID() and RemoteID() methods by correcting self-references to use SecureConnection fields.
lemon-mint committed
Oct 28, 2025 at 15:26 UTC
8aec989ec3355143729737d43c544682daabbebe
6 files changed
+883
-121
relaydns/client.go
+72
-2
@@ -11,6 +11,7 @@ import (
11
"github.com/gosuda/relaydns/relaydns/core/proto/rdsec"
12
"github.com/gosuda/relaydns/relaydns/core/proto/rdverb"
13
"github.com/hashicorp/yamux"
14
+ "github.com/rs/zerolog/log"
15
)
16
17
var (
@@ -28,11 +29,11 @@ func (i *IncommingConn) LeaseID() string {
29
}
30
31
func (i *IncommingConn) LocalID() string {
31
- return i.LocalID()
32
+ return i.SecureConnection.LocalID()
33
}
34
35
func (i *IncommingConn) RemoteID() string {
35
- return i.RemoteID()
36
+ return i.SecureConnection.RemoteID()
37
}
38
39
// RelayClient는 RelayServer에 연결하여 서비스를 요청하는 클라이언트입니다.
@@ -57,16 +58,21 @@ type leaseWithCred struct {
58
59
// NewRelayClient는 새로운 RelayClient 인스턴스를 생성합니다.
60
func NewRelayClient(conn io.ReadWriteCloser) *RelayClient {
61
+ log.Debug().Msg("[RelayClient] Creating new relay client")
62
+
63
// Create yamux session as client
64
config := yamux.DefaultConfig()
65
config.Logger = nil // Disable logging for cleaner output
66
sess, err := yamux.Client(conn, config)
67
if err != nil {
68
+ log.Error().Err(err).Msg("[RelayClient] Failed to create yamux session")
69
// If session creation fails, close the connection and return nil
70
conn.Close()
71
return nil
72
}
73
74
+ log.Debug().Msg("[RelayClient] Yamux session created successfully")
75
+
76
g := &RelayClient{
77
conn: conn,
78
sess: sess,
@@ -79,11 +85,13 @@ func NewRelayClient(conn io.ReadWriteCloser) *RelayClient {
85
go g.leaseUpdateWorker()
86
go g.leaseListenWorker()
87
88
+ log.Debug().Msg("[RelayClient] RelayClient initialized and workers started")
89
return g
90
}
91
92
// Close는 서버와의 연결을 종료합니다.
93
func (g *RelayClient) Close() error {
94
+ log.Debug().Msg("[RelayClient] Closing relay client")
95
close(g.stopCh)
96
g.waitGroup.Wait()
97
@@ -92,6 +100,7 @@ func (g *RelayClient) Close() error {
100
// Close the session first
101
if g.sess != nil {
102
if err := g.sess.Close(); err != nil {
103
+ log.Error().Err(err).Msg("[RelayClient] Error closing yamux session")
104
errs = append(errs, err)
105
}
106
}
@@ -99,10 +108,12 @@ func (g *RelayClient) Close() error {
108
// Then close the underlying connection
109
if g.conn != nil {
110
if err := g.conn.Close(); err != nil {
111
+ log.Error().Err(err).Msg("[RelayClient] Error closing connection")
112
errs = append(errs, err)
113
}
114
}
115
116
+ log.Debug().Msg("[RelayClient] Relay client closed")
117
if len(errs) > 0 {
118
return errs[0]
119
}
@@ -145,10 +156,12 @@ func (g *RelayClient) leaseUpdateWorker() {
156
157
func (g *RelayClient) leaseListenWorker() {
158
defer g.waitGroup.Done()
159
+ log.Debug().Msg("[RelayClient] Lease listen worker started")
160
161
for {
162
select {
163
case <-g.stopCh:
164
+ log.Debug().Msg("[RelayClient] Lease listen worker stopped")
165
return
166
default:
167
if g.sess == nil {
@@ -164,23 +177,29 @@ func (g *RelayClient) leaseListenWorker() {
177
case <-g.stopCh:
178
return
179
default:
180
+ log.Debug().Err(err).Msg("[RelayClient] Error accepting stream, retrying")
181
// Continue trying to accept streams
182
continue
183
}
184
}
185
+ log.Debug().Uint32("stream_id", stream.StreamID()).Msg("[RelayClient] Accepted incoming stream")
186
go g.handleConnectionRequestStream(stream)
187
}
188
}
189
}
190
191
func (g *RelayClient) handleConnectionRequestStream(stream *yamux.Stream) {
192
+ log.Debug().Uint32("stream_id", stream.StreamID()).Msg("[RelayClient] Handling connection request stream")
193
+
194
pkt, err := readPacket(stream)
195
if err != nil {
196
+ log.Error().Err(err).Msg("[RelayClient] Failed to read packet from stream")
197
stream.Close()
198
return
199
}
200
201
if pkt.Type != rdverb.PacketType_PACKET_TYPE_CONNECTION_REQUEST {
202
+ log.Warn().Str("packet_type", pkt.Type.String()).Msg("[RelayClient] Unexpected packet type")
203
stream.Close()
204
return
205
}
@@ -188,23 +207,29 @@ func (g *RelayClient) handleConnectionRequestStream(stream *yamux.Stream) {
207
req := &rdverb.ConnectionRequest{}
208
err = req.UnmarshalVT(pkt.Payload)
209
if err != nil {
210
+ log.Error().Err(err).Msg("[RelayClient] Failed to unmarshal connection request")
211
stream.Close()
212
return
213
}
214
215
+ log.Debug().Str("lease_id", req.LeaseId).Msg("[RelayClient] Connection request received")
216
+
217
g.leasesMu.Lock()
218
lease, ok := g.leases[req.LeaseId]
219
g.leasesMu.Unlock()
220
221
resp := &rdverb.ConnectionResponse{}
222
if !ok {
223
+ log.Warn().Str("lease_id", req.LeaseId).Msg("[RelayClient] Lease not found, rejecting connection")
224
resp.Code = rdverb.ResponseCode_RESPONSE_CODE_REJECTED
225
} else {
226
+ log.Debug().Str("lease_id", req.LeaseId).Msg("[RelayClient] Lease found, accepting connection")
227
resp.Code = rdverb.ResponseCode_RESPONSE_CODE_ACCEPTED
228
}
229
230
respPayload, err := resp.MarshalVT()
231
if err != nil {
232
+ log.Error().Err(err).Msg("[RelayClient] Failed to marshal response")
233
stream.Close()
234
return
235
}
@@ -214,6 +239,7 @@ func (g *RelayClient) handleConnectionRequestStream(stream *yamux.Stream) {
239
Payload: respPayload,
240
})
241
if err != nil {
242
+ log.Error().Err(err).Msg("[RelayClient] Failed to write response packet")
243
stream.Close()
244
return
245
}
@@ -223,13 +249,21 @@ func (g *RelayClient) handleConnectionRequestStream(stream *yamux.Stream) {
249
return
250
}
251
252
+ log.Debug().Str("lease_id", req.LeaseId).Msg("[RelayClient] Starting server handshake")
253
handshaker := cryptoops.NewHandshaker(lease.Cred)
254
secConn, err := handshaker.ServerHandshake(stream, lease.Lease.Alpn)
255
if err != nil {
256
+ log.Error().Err(err).Str("lease_id", req.LeaseId).Msg("[RelayClient] Server handshake failed")
257
stream.Close()
258
return
259
}
260
261
+ log.Debug().
262
+ Str("lease_id", req.LeaseId).
263
+ Str("local_id", secConn.LocalID()).
264
+ Str("remote_id", secConn.RemoteID()).
265
+ Msg("[RelayClient] Secure connection established, sending to incoming channel")
266
+
267
g.incommingConnCh <- &IncommingConn{
268
SecureConnection: secConn,
269
leaseID: req.LeaseId,
@@ -414,9 +448,12 @@ func (g *RelayClient) deleteLease(cred *cryptoops.Credential, identity *rdsec.Id
448
449
// requestConnection은 다른 클라이언트로의 연결을 요청합니다.
450
func (g *RelayClient) RequestConnection(leaseID string, alpn string, clientCred *cryptoops.Credential) (rdverb.ResponseCode, *cryptoops.SecureConnection, error) {
451
+ log.Debug().Str("lease_id", leaseID).Str("alpn", alpn).Msg("[RelayClient] Requesting connection")
452
+
453
// 새 스트림 열기
454
stream, err := g.sess.OpenStream()
455
if err != nil {
456
+ log.Error().Err(err).Msg("[RelayClient] Failed to open stream for connection request")
457
return rdverb.ResponseCode_RESPONSE_CODE_UNKNOWN, nil, err
458
}
459
@@ -433,28 +470,34 @@ func (g *RelayClient) RequestConnection(leaseID string, alpn string, clientCred
470
471
reqPayload, err := req.MarshalVT()
472
if err != nil {
473
+ log.Error().Err(err).Msg("[RelayClient] Failed to marshal connection request")
474
stream.Close()
475
return rdverb.ResponseCode_RESPONSE_CODE_UNKNOWN, nil, err
476
}
477
478
// 요청 전송
479
+ log.Debug().Str("lease_id", leaseID).Msg("[RelayClient] Sending connection request")
480
err = writePacket(stream, &rdverb.Packet{
481
Type: rdverb.PacketType_PACKET_TYPE_CONNECTION_REQUEST,
482
Payload: reqPayload,
483
})
484
if err != nil {
485
+ log.Error().Err(err).Msg("[RelayClient] Failed to write connection request packet")
486
stream.Close()
487
return rdverb.ResponseCode_RESPONSE_CODE_UNKNOWN, nil, err
488
}
489
490
// 응답 수신
491
+ log.Debug().Str("lease_id", leaseID).Msg("[RelayClient] Waiting for connection response")
492
respPacket, err := readPacket(stream)
493
if err != nil {
494
+ log.Error().Err(err).Msg("[RelayClient] Failed to read connection response")
495
stream.Close()
496
return rdverb.ResponseCode_RESPONSE_CODE_UNKNOWN, nil, err
497
}
498
499
if respPacket.Type != rdverb.PacketType_PACKET_TYPE_CONNECTION_RESPONSE {
500
+ log.Warn().Str("packet_type", respPacket.Type.String()).Msg("[RelayClient] Unexpected response packet type")
501
stream.Close()
502
return rdverb.ResponseCode_RESPONSE_CODE_UNKNOWN, nil, ErrInvalidResponse
503
}
@@ -462,23 +505,38 @@ func (g *RelayClient) RequestConnection(leaseID string, alpn string, clientCred
505
var resp rdverb.ConnectionResponse
506
err = resp.UnmarshalVT(respPacket.Payload)
507
if err != nil {
508
+ log.Error().Err(err).Msg("[RelayClient] Failed to unmarshal connection response")
509
stream.Close()
510
return rdverb.ResponseCode_RESPONSE_CODE_UNKNOWN, nil, err
511
}
512
513
+ log.Debug().
514
+ Str("lease_id", leaseID).
515
+ Str("response_code", resp.Code.String()).
516
+ Msg("[RelayClient] Connection response received")
517
+
518
// 거절된 경우 스트림을 닫고 오류 코드 반환
519
if resp.Code != rdverb.ResponseCode_RESPONSE_CODE_ACCEPTED {
520
+ log.Warn().Str("lease_id", leaseID).Str("code", resp.Code.String()).Msg("[RelayClient] Connection rejected")
521
stream.Close()
522
return resp.Code, nil, ErrConnectionRejected
523
}
524
525
+ log.Debug().Str("lease_id", leaseID).Msg("[RelayClient] Starting client handshake")
526
handshaker := cryptoops.NewHandshaker(clientCred)
527
secConn, err := handshaker.ClientHandshake(stream, alpn)
528
if err != nil {
529
+ log.Error().Err(err).Str("lease_id", leaseID).Msg("[RelayClient] Client handshake failed")
530
stream.Close()
531
return rdverb.ResponseCode_RESPONSE_CODE_UNKNOWN, nil, err
532
}
533
534
+ log.Debug().
535
+ Str("lease_id", leaseID).
536
+ Str("local_id", secConn.LocalID()).
537
+ Str("remote_id", secConn.RemoteID()).
538
+ Msg("[RelayClient] Secure connection established successfully")
539
+
540
return resp.Code, secConn, nil
541
}
542
@@ -488,6 +546,12 @@ func (g *RelayClient) RegisterLease(cred *cryptoops.Credential, name string, alp
546
PublicKey: cred.PublicKey(),
547
}
548
549
+ log.Debug().
550
+ Str("lease_id", identity.Id).
551
+ Str("name", name).
552
+ Strs("alpns", alpns).
553
+ Msg("[RelayClient] Registering lease")
554
+
555
lease := &rdverb.Lease{
556
Identity: identity,
557
Expires: time.Now().Add(30 * time.Second).Unix(),
@@ -504,12 +568,18 @@ func (g *RelayClient) RegisterLease(cred *cryptoops.Credential, name string, alp
568
569
resp, err := g.updateLease(cred, lease)
570
if err != nil || resp != rdverb.ResponseCode_RESPONSE_CODE_ACCEPTED {
571
+ log.Error().
572
+ Err(err).
573
+ Str("lease_id", identity.Id).
574
+ Str("response", resp.String()).
575
+ Msg("[RelayClient] Failed to register lease")
576
g.leasesMu.Lock()
577
delete(g.leases, identity.Id)
578
g.leasesMu.Unlock()
579
return err
580
}
581
582
+ log.Debug().Str("lease_id", identity.Id).Msg("[RelayClient] Lease registered successfully")
583
return nil
584
}
585
relaydns/handlers.go
+204
-116
@@ -9,6 +9,7 @@ import (
9
"github.com/gosuda/relaydns/relaydns/core/proto/rdsec"
10
"github.com/gosuda/relaydns/relaydns/core/proto/rdverb"
11
"github.com/hashicorp/yamux"
12
+ "github.com/rs/zerolog/log"
13
"github.com/valyala/bytebufferpool"
14
)
15
@@ -133,171 +134,258 @@ func (g *RelayServer) handleConnectionRequest(ctx *StreamContext, packet *rdverb
134
var req rdverb.ConnectionRequest
135
err := req.UnmarshalVT(packet.Payload)
136
if err != nil {
137
+ log.Error().Err(err).Msg("[RelayServer] Failed to unmarshal connection request")
138
return err
139
}
140
141
+ log.Debug().
142
+ Str("lease_id", req.LeaseId).
143
+ Str("client_id", req.ClientIdentity.Id).
144
+ Int64("conn_id", ctx.ConnectionID).
145
+ Msg("[RelayServer] Handling connection request")
146
+
147
var resp rdverb.ConnectionResponse
148
141
- // Check if lease exists and get lease connection
142
- leaseEntry, exists := g.leaseManager.GetLease(req.ClientIdentity)
149
+ // Check if lease exists and get lease connection using LeaseId
150
+ leaseEntry, exists := g.leaseManager.GetLeaseByID(req.LeaseId)
151
if !exists {
152
+ log.Warn().Str("lease_id", req.LeaseId).Msg("[RelayServer] Lease not found")
153
resp.Code = rdverb.ResponseCode_RESPONSE_CODE_INVALID_IDENTITY
145
- } else {
146
- // Get the lease connection using connection ID
147
- g.connectionsLock.RLock()
148
- leaseConn, leaseExists := g.connections[leaseEntry.ConnectionID]
149
- g.connectionsLock.RUnlock()
150
-
151
- if !leaseExists {
152
- resp.Code = rdverb.ResponseCode_RESPONSE_CODE_INVALID_IDENTITY
153
- } else {
154
- // Forward the connection request to the lease holder
155
- forwardResp, err := g.forwardConnectionRequest(leaseConn, &req)
156
- if err != nil {
157
- resp.Code = rdverb.ResponseCode_RESPONSE_CODE_REJECTED
158
- } else {
159
- resp.Code = forwardResp.Code
160
-
161
- // If accepted, set up bidirectional forwarding
162
- if resp.Code == rdverb.ResponseCode_RESPONSE_CODE_ACCEPTED {
163
- ctx.Hijack()
164
- go g.setupBidirectionalForwarding(ctx.Stream, leaseConn, leaseEntry)
165
- }
166
- }
154
+
155
+ response, err := resp.MarshalVT()
156
+ if err != nil {
157
+ log.Error().Err(err).Msg("[RelayServer] Failed to marshal connection response")
158
+ return err
159
}
168
- }
160
170
- response, err := resp.MarshalVT()
171
- if err != nil {
172
- return err
161
+ return writePacket(ctx.Stream, &rdverb.Packet{
162
+ Type: rdverb.PacketType_PACKET_TYPE_CONNECTION_RESPONSE,
163
+ Payload: response,
164
+ })
165
}
166
175
- return writePacket(ctx.Stream, &rdverb.Packet{
176
- Type: rdverb.PacketType_PACKET_TYPE_CONNECTION_RESPONSE,
177
- Payload: response,
178
- })
179
-}
167
+ log.Debug().
168
+ Str("lease_id", req.LeaseId).
169
+ Int64("lease_conn_id", leaseEntry.ConnectionID).
170
+ Msg("[RelayServer] Lease found, forwarding to lease holder")
171
+
172
+ // Get the lease connection using connection ID
173
+ g.connectionsLock.RLock()
174
+ leaseConn, leaseExists := g.connections[leaseEntry.ConnectionID]
175
+ g.connectionsLock.RUnlock()
176
+
177
+ if !leaseExists {
178
+ log.Warn().
179
+ Str("lease_id", req.LeaseId).
180
+ Int64("lease_conn_id", leaseEntry.ConnectionID).
181
+ Msg("[RelayServer] Lease connection no longer active")
182
+ resp.Code = rdverb.ResponseCode_RESPONSE_CODE_INVALID_IDENTITY
183
181
-func (g *RelayServer) forwardConnectionRequest(leaseConn *Connection, req *rdverb.ConnectionRequest) (*rdverb.ConnectionResponse, error) {
182
- // Open a new stream to the lease holder
183
- forwardStream, err := leaseConn.sess.OpenStream()
184
+ response, err := resp.MarshalVT()
185
+ if err != nil {
186
+ log.Error().Err(err).Msg("[RelayServer] Failed to marshal connection response")
187
+ return err
188
+ }
189
+
190
+ return writePacket(ctx.Stream, &rdverb.Packet{
191
+ Type: rdverb.PacketType_PACKET_TYPE_CONNECTION_RESPONSE,
192
+ Payload: response,
193
+ })
194
+ }
195
+
196
+ // Open a stream to the lease holder
197
+ log.Debug().Str("lease_id", req.LeaseId).Msg("[RelayServer] Opening stream to lease holder")
198
+ leaseStream, err := leaseConn.sess.OpenStream()
199
if err != nil {
185
- return nil, err
200
+ log.Error().Err(err).Str("lease_id", req.LeaseId).Msg("[RelayServer] Failed to open stream to lease holder")
201
+ resp.Code = rdverb.ResponseCode_RESPONSE_CODE_REJECTED
202
+
203
+ response, err := resp.MarshalVT()
204
+ if err != nil {
205
+ return err
206
+ }
207
+
208
+ return writePacket(ctx.Stream, &rdverb.Packet{
209
+ Type: rdverb.PacketType_PACKET_TYPE_CONNECTION_RESPONSE,
210
+ Payload: response,
211
+ })
212
}
187
- defer forwardStream.Close()
213
214
// Forward the connection request
215
requestPayload, err := req.MarshalVT()
216
if err != nil {
192
- return nil, err
217
+ log.Error().Err(err).Msg("[RelayServer] Failed to marshal forward request")
218
+ leaseStream.Close()
219
+ resp.Code = rdverb.ResponseCode_RESPONSE_CODE_REJECTED
220
+
221
+ response, err := resp.MarshalVT()
222
+ if err != nil {
223
+ return err
224
+ }
225
+
226
+ return writePacket(ctx.Stream, &rdverb.Packet{
227
+ Type: rdverb.PacketType_PACKET_TYPE_CONNECTION_RESPONSE,
228
+ Payload: response,
229
+ })
230
}
231
195
- err = writePacket(forwardStream, &rdverb.Packet{
232
+ log.Debug().Str("lease_id", req.LeaseId).Msg("[RelayServer] Sending connection request to lease holder")
233
+ err = writePacket(leaseStream, &rdverb.Packet{
234
Type: rdverb.PacketType_PACKET_TYPE_CONNECTION_REQUEST,
235
Payload: requestPayload,
236
})
237
if err != nil {
200
- return nil, err
238
+ log.Error().Err(err).Msg("[RelayServer] Failed to write forward request")
239
+ leaseStream.Close()
240
+ resp.Code = rdverb.ResponseCode_RESPONSE_CODE_REJECTED
241
+
242
+ response, err := resp.MarshalVT()
243
+ if err != nil {
244
+ return err
245
+ }
246
+
247
+ return writePacket(ctx.Stream, &rdverb.Packet{
248
+ Type: rdverb.PacketType_PACKET_TYPE_CONNECTION_RESPONSE,
249
+ Payload: response,
250
+ })
251
}
252
253
// Read the response
204
- respPacket, err := readPacket(forwardStream)
254
+ log.Debug().Str("lease_id", req.LeaseId).Msg("[RelayServer] Waiting for response from lease holder")
255
+ respPacket, err := readPacket(leaseStream)
256
if err != nil {
206
- return nil, err
257
+ log.Error().Err(err).Msg("[RelayServer] Failed to read forward response")
258
+ leaseStream.Close()
259
+ resp.Code = rdverb.ResponseCode_RESPONSE_CODE_REJECTED
260
+
261
+ response, err := resp.MarshalVT()
262
+ if err != nil {
263
+ return err
264
+ }
265
+
266
+ return writePacket(ctx.Stream, &rdverb.Packet{
267
+ Type: rdverb.PacketType_PACKET_TYPE_CONNECTION_RESPONSE,
268
+ Payload: response,
269
+ })
270
}
271
272
if respPacket.Type != rdverb.PacketType_PACKET_TYPE_CONNECTION_RESPONSE {
210
- return nil, err
273
+ log.Warn().Str("packet_type", respPacket.Type.String()).Msg("[RelayServer] Unexpected response packet type")
274
+ leaseStream.Close()
275
+ resp.Code = rdverb.ResponseCode_RESPONSE_CODE_REJECTED
276
+
277
+ response, err := resp.MarshalVT()
278
+ if err != nil {
279
+ return err
280
+ }
281
+
282
+ return writePacket(ctx.Stream, &rdverb.Packet{
283
+ Type: rdverb.PacketType_PACKET_TYPE_CONNECTION_RESPONSE,
284
+ Payload: response,
285
+ })
286
}
287
213
- var resp rdverb.ConnectionResponse
288
err = resp.UnmarshalVT(respPacket.Payload)
289
if err != nil {
216
- return nil, err
217
- }
218
- return &resp, nil
219
-}
290
+ log.Error().Err(err).Msg("[RelayServer] Failed to unmarshal forward response")
291
+ leaseStream.Close()
292
+ resp.Code = rdverb.ResponseCode_RESPONSE_CODE_REJECTED
293
221
-func (g *RelayServer) setupBidirectionalForwarding(clientStream *yamux.Stream, leaseConn *Connection, leaseEntry *LeaseEntry) {
222
- // Open a new stream to the lease holder for data forwarding
223
- dataStream, err := leaseConn.sess.OpenStream()
224
- if err != nil {
225
- return
226
- }
227
- defer dataStream.Close()
294
+ response, err := resp.MarshalVT()
295
+ if err != nil {
296
+ return err
297
+ }
298
229
- // Send a special packet to indicate this is a data stream
230
- initPacket := &rdverb.Packet{
231
- Type: rdverb.PacketType_PACKET_TYPE_CONNECTION_REQUEST, // Reuse as init signal
232
- Payload: []byte("DATA_STREAM"),
299
+ return writePacket(ctx.Stream, &rdverb.Packet{
300
+ Type: rdverb.PacketType_PACKET_TYPE_CONNECTION_RESPONSE,
301
+ Payload: response,
302
+ })
303
}
304
235
- err = writePacket(dataStream, initPacket)
236
- if err != nil {
237
- return
238
- }
305
+ log.Debug().
306
+ Str("lease_id", req.LeaseId).
307
+ Str("response_code", resp.Code.String()).
308
+ Msg("[RelayServer] Received response from lease holder, sending to client")
309
240
- // Read the response to confirm data stream is ready
241
- respPacket, err := readPacket(dataStream)
310
+ // Send response to client
311
+ response, err := resp.MarshalVT()
312
if err != nil {
243
- return
244
- }
245
-
246
- // Check if we got a data stream confirmation
247
- if respPacket.Type != rdverb.PacketType_PACKET_TYPE_CONNECTION_RESPONSE {
248
- return
313
+ log.Error().Err(err).Msg("[RelayServer] Failed to marshal connection response")
314
+ leaseStream.Close()
315
+ return err
316
}
317
251
- var resp rdverb.ConnectionResponse
252
- err = resp.UnmarshalVT(respPacket.Payload)
318
+ err = writePacket(ctx.Stream, &rdverb.Packet{
319
+ Type: rdverb.PacketType_PACKET_TYPE_CONNECTION_RESPONSE,
320
+ Payload: response,
321
+ })
322
if err != nil {
254
- return
255
- }
256
-
257
- if resp.Code != rdverb.ResponseCode_RESPONSE_CODE_ACCEPTED {
258
- return
323
+ log.Error().Err(err).Msg("[RelayServer] Failed to write connection response")
324
+ leaseStream.Close()
325
+ return err
326
}
327
261
- // Track this relayed connection
262
- leaseID := string(leaseEntry.Lease.Identity.Id)
263
- g.relayedConnectionsLock.Lock()
264
- g.relayedConnections[leaseID] = append(g.relayedConnections[leaseID], clientStream)
265
- g.relayedConnectionsLock.Unlock()
266
-
267
- // Set up bidirectional copying with cleanup
268
- var wg sync.WaitGroup
269
- wg.Add(2)
270
-
271
- // Copy from client to lease holder
272
- go func() {
273
- defer wg.Done()
274
- io.Copy(dataStream, clientStream)
275
- }()
276
-
277
- // Copy from lease holder to client
278
- go func() {
279
- defer wg.Done()
280
- io.Copy(clientStream, dataStream)
281
- }()
282
-
283
- wg.Wait()
284
-
285
- // Clean up relayed connection tracking when done
286
- g.relayedConnectionsLock.Lock()
287
- if streams, exists := g.relayedConnections[leaseID]; exists {
288
- // Remove this stream from the slice
289
- for i, stream := range streams {
290
- if stream == clientStream {
291
- g.relayedConnections[leaseID] = append(streams[:i], streams[i+1:]...)
292
- break
328
+ // If accepted, hijack both streams and set up bidirectional forwarding
329
+ if resp.Code == rdverb.ResponseCode_RESPONSE_CODE_ACCEPTED {
330
+ log.Debug().Str("lease_id", req.LeaseId).Msg("[RelayServer] Connection accepted, setting up bidirectional forwarding")
331
+ ctx.Hijack()
332
+
333
+ leaseID := string(leaseEntry.Lease.Identity.Id)
334
+ g.relayedConnectionsLock.Lock()
335
+ g.relayedConnections[leaseID] = append(g.relayedConnections[leaseID], ctx.Stream)
336
+ g.relayedConnectionsLock.Unlock()
337
+
338
+ // Set up bidirectional copying
339
+ var wg sync.WaitGroup
340
+ wg.Add(2)
341
+
342
+ // Copy from client to lease holder
343
+ go func() {
344
+ defer wg.Done()
345
+ n, err := io.Copy(leaseStream, ctx.Stream)
346
+ log.Debug().
347
+ Str("lease_id", leaseID).
348
+ Int64("bytes", n).
349
+ Err(err).
350
+ Msg("[RelayServer] Client -> Lease copy finished")
351
+ leaseStream.Close()
352
+ }()
353
+
354
+ // Copy from lease holder to client
355
+ go func() {
356
+ defer wg.Done()
357
+ n, err := io.Copy(ctx.Stream, leaseStream)
358
+ log.Debug().
359
+ Str("lease_id", leaseID).
360
+ Int64("bytes", n).
361
+ Err(err).
362
+ Msg("[RelayServer] Lease -> Client copy finished")
363
+ ctx.Stream.Close()
364
+ }()
365
+
366
+ wg.Wait()
367
+ log.Debug().Str("lease_id", leaseID).Msg("[RelayServer] Bidirectional forwarding completed")
368
+
369
+ // Clean up relayed connection tracking
370
+ g.relayedConnectionsLock.Lock()
371
+ if streams, exists := g.relayedConnections[leaseID]; exists {
372
+ for i, stream := range streams {
373
+ if stream == ctx.Stream {
374
+ g.relayedConnections[leaseID] = append(streams[:i], streams[i+1:]...)
375
+ break
376
+ }
377
+ }
378
+ if len(g.relayedConnections[leaseID]) == 0 {
379
+ delete(g.relayedConnections, leaseID)
380
}
381
}
295
- // If no more streams for this lease, remove the entry
296
- if len(g.relayedConnections[leaseID]) == 0 {
297
- delete(g.relayedConnections, leaseID)
298
- }
382
+ g.relayedConnectionsLock.Unlock()
383
+ } else {
384
+ // Connection rejected, close lease stream
385
+ leaseStream.Close()
386
}
300
- g.relayedConnectionsLock.Unlock()
387
+
388
+ return nil
389
}
390
391
// Helper function to read packet from stream
relaydns/lease.go
+17
@@ -116,6 +116,23 @@ func (lm *LeaseManager) GetLease(identity *rdsec.Identity) (*LeaseEntry, bool) {
116
return lease, true
117
}
118
119
+func (lm *LeaseManager) GetLeaseByID(leaseID string) (*LeaseEntry, bool) {
120
+ lm.leasesLock.RLock()
121
+ defer lm.leasesLock.RUnlock()
122
+
123
+ lease, exists := lm.leases[leaseID]
124
+ if !exists {
125
+ return nil, false
126
+ }
127
+
128
+ // Check if lease is expired
129
+ if time.Now().After(lease.Expires) {
130
+ return nil, false
131
+ }
132
+
133
+ return lease, true
134
+}
135
+
136
func (lm *LeaseManager) GetAllLeases() []*rdverb.Lease {
137
lm.leasesLock.RLock()
138
defer lm.leasesLock.RUnlock()
relaydns/relay.go
+61
-2
@@ -9,6 +9,7 @@ import (
9
"github.com/gosuda/relaydns/relaydns/core/proto/rdsec"
10
"github.com/gosuda/relaydns/relaydns/core/proto/rdverb"
11
"github.com/hashicorp/yamux"
12
+ "github.com/rs/zerolog/log"
13
)
14
15
type Connection struct {
@@ -60,10 +61,21 @@ func NewRelayServer(credential *cryptoops.Credential, address []string) *RelaySe
61
var _yamux_config = yamux.DefaultConfig()
62
63
func (g *RelayServer) handleConn(id int64, connection *Connection) {
64
+ log.Debug().Int64("conn_id", id).Msg("[RelayServer] Handling new connection")
65
+
66
defer func() {
67
+ log.Debug().Int64("conn_id", id).Msg("[RelayServer] Connection closing, cleaning up")
68
+
69
// Clean up leases associated with this connection when it closes
70
cleanedLeaseIDs := g.leaseManager.CleanupLeasesByConnectionID(id)
71
72
+ if len(cleanedLeaseIDs) > 0 {
73
+ log.Debug().
74
+ Int64("conn_id", id).
75
+ Strs("lease_ids", cleanedLeaseIDs).
76
+ Msg("[RelayServer] Cleaned up leases for connection")
77
+ }
78
+
79
// Also clean up lease connections mapping
80
g.leaseConnectionsLock.Lock()
81
for _, leaseID := range cleanedLeaseIDs {
@@ -91,13 +103,20 @@ func (g *RelayServer) handleConn(id int64, connection *Connection) {
103
104
// Close the underlying connection
105
connection.conn.Close()
106
+ log.Debug().Int64("conn_id", id).Msg("[RelayServer] Connection cleanup complete")
107
}()
108
109
for {
110
stream, err := connection.sess.AcceptStream()
111
if err != nil {
112
+ log.Debug().Err(err).Int64("conn_id", id).Msg("[RelayServer] Error accepting stream, connection closing")
113
return
114
}
115
+ log.Debug().
116
+ Int64("conn_id", id).
117
+ Uint32("stream_id", stream.StreamID()).
118
+ Msg("[RelayServer] Accepted new stream")
119
+
120
connection.streamsLock.Lock()
121
connection.streams[stream.StreamID()] = stream
122
connection.streamsLock.Unlock()
@@ -108,14 +127,28 @@ func (g *RelayServer) handleConn(id int64, connection *Connection) {
127
const _MAX_RAW_PACKET_SIZE = 1 << 26 // 64MB
128
129
func (g *RelayServer) handleStream(stream *yamux.Stream, id int64, connection *Connection) {
130
+ log.Debug().
131
+ Int64("conn_id", id).
132
+ Uint32("stream_id", stream.StreamID()).
133
+ Msg("[RelayServer] Handling stream")
134
+
135
var hijacked bool = false
136
defer func() {
137
stream_id := stream.StreamID()
138
if !hijacked {
139
+ log.Debug().
140
+ Int64("conn_id", id).
141
+ Uint32("stream_id", stream_id).
142
+ Msg("[RelayServer] Closing stream")
143
connection.streamsLock.Lock()
144
stream.Close()
145
delete(connection.streams, stream_id)
146
connection.streamsLock.Unlock()
147
+ } else {
148
+ log.Debug().
149
+ Int64("conn_id", id).
150
+ Uint32("stream_id", stream_id).
151
+ Msg("[RelayServer] Stream was hijacked, not closing")
152
}
153
}()
154
@@ -130,9 +163,20 @@ func (g *RelayServer) handleStream(stream *yamux.Stream, id int64, connection *C
163
for {
164
packet, err := readPacket(stream)
165
if err != nil {
166
+ log.Debug().
167
+ Err(err).
168
+ Int64("conn_id", id).
169
+ Uint32("stream_id", stream.StreamID()).
170
+ Msg("[RelayServer] Error reading packet")
171
return
172
}
173
174
+ log.Debug().
175
+ Int64("conn_id", id).
176
+ Uint32("stream_id", stream.StreamID()).
177
+ Str("packet_type", packet.Type.String()).
178
+ Msg("[RelayServer] Received packet")
179
+
180
switch packet.Type {
181
case rdverb.PacketType_PACKET_TYPE_RELAY_INFO_REQUEST:
182
err = g.handleRelayInfoRequest(ctx, packet)
@@ -143,38 +187,53 @@ func (g *RelayServer) handleStream(stream *yamux.Stream, id int64, connection *C
187
case rdverb.PacketType_PACKET_TYPE_CONNECTION_REQUEST:
188
err = g.handleConnectionRequest(ctx, packet)
189
default:
190
+ log.Warn().
191
+ Int64("conn_id", id).
192
+ Str("packet_type", packet.Type.String()).
193
+ Msg("[RelayServer] Unknown packet type")
194
// Unknown packet type, return to close the stream
195
return
196
}
197
198
if err != nil {
199
+ log.Error().
200
+ Err(err).
201
+ Int64("conn_id", id).
202
+ Str("packet_type", packet.Type.String()).
203
+ Msg("[RelayServer] Error handling packet")
204
return
205
}
206
207
// If the stream was hijacked, exit the loop
208
if hijacked {
209
+ log.Debug().Int64("conn_id", id).Msg("[RelayServer] Stream hijacked, exiting handler")
210
return
211
}
212
}
213
}
214
215
func (g *RelayServer) HandleConnection(conn io.ReadWriteCloser) error {
216
+ log.Debug().Msg("[RelayServer] New connection received")
217
+
218
sess, err := yamux.Server(conn, _yamux_config)
219
if err != nil {
220
+ log.Error().Err(err).Msg("[RelayServer] Failed to create yamux server session")
221
return err
222
}
223
224
g.connectionsLock.Lock()
225
g.connidCounter++
226
+ connID := g.connidCounter
227
connection := &Connection{
228
conn: conn,
229
sess: sess,
230
streams: make(map[uint32]*yamux.Stream),
231
}
174
- g.connections[g.connidCounter] = connection
232
+ g.connections[connID] = connection
233
g.connectionsLock.Unlock()
234
177
- go g.handleConn(g.connidCounter, connection)
235
+ log.Debug().Int64("conn_id", connID).Msg("[RelayServer] Connection registered, starting handler")
236
+ go g.handleConn(connID, connection)
237
238
return nil
239
}
sdk/sdk.go
+68
-1
@@ -15,6 +15,7 @@ import (
15
"github.com/gosuda/relaydns/relaydns/core/cryptoops"
16
"github.com/gosuda/relaydns/relaydns/core/proto/rdverb"
17
"github.com/gosuda/relaydns/relaydns/utils/wsstream"
18
+ "github.com/rs/zerolog/log"
19
)
20
21
func NewCredential() (*cryptoops.Credential, error) {
@@ -128,6 +129,8 @@ var (
129
)
130
131
func NewClient(opt ...Option) (*RDClient, error) {
132
+ log.Debug().Msg("[SDK] Creating new RDClient")
133
+
134
config := &RDClientConfig{
135
Dialer: webSocketDialer(),
136
}
@@ -145,19 +148,23 @@ func NewClient(opt ...Option) (*RDClient, error) {
148
// Initialize relays from bootstrap servers
149
var connectionErrors []error
150
for _, server := range config.BootstrapServers {
151
+ log.Debug().Str("server", server).Msg("[SDK] Connecting to bootstrap server")
152
conn, err := config.Dialer(context.Background(), server)
153
if err != nil {
154
+ log.Error().Err(err).Str("server", server).Msg("[SDK] Failed to connect to bootstrap server")
155
connectionErrors = append(connectionErrors, err)
156
continue // Skip failed connections
157
}
158
159
relayClient := relaydns.NewRelayClient(conn)
160
if relayClient == nil {
161
+ log.Error().Str("server", server).Msg("[SDK] Failed to create relay client")
162
conn.Close()
163
connectionErrors = append(connectionErrors, ErrFailedToCreateClient)
164
continue
165
}
166
167
+ log.Debug().Str("server", server).Msg("[SDK] Successfully connected to bootstrap server")
168
client.relays[server] = &rdRelay{
169
addr: server,
170
client: relayClient,
@@ -167,13 +174,20 @@ func NewClient(opt ...Option) (*RDClient, error) {
174
175
// If no relays were successfully connected, return an error
176
if len(client.relays) == 0 && len(config.BootstrapServers) > 0 {
177
+ log.Error().Int("attempted", len(config.BootstrapServers)).Msg("[SDK] Failed to connect to any bootstrap servers")
178
return nil, fmt.Errorf("failed to connect to any bootstrap servers: %v", connectionErrors)
179
}
180
181
+ log.Debug().Int("relay_count", len(client.relays)).Msg("[SDK] RDClient created successfully")
182
return client, nil
183
}
184
185
func (g *RDClient) Dial(cred *cryptoops.Credential, leaseID string, alpn string) (*RDConnection, error) {
186
+ log.Debug().
187
+ Str("lease_id", leaseID).
188
+ Str("alpn", alpn).
189
+ Msg("[SDK] Dialing to lease")
190
+
191
var relays []*rdRelay
192
193
g.mu.Lock()
@@ -182,6 +196,8 @@ func (g *RDClient) Dial(cred *cryptoops.Credential, leaseID string, alpn string)
196
}
197
g.mu.Unlock()
198
199
+ log.Debug().Int("relay_count", len(relays)).Msg("[SDK] Checking relays for lease")
200
+
201
var wg sync.WaitGroup
202
var availableRelaysMu sync.Mutex
203
var availableRelays []*rdRelay
@@ -192,10 +208,12 @@ func (g *RDClient) Dial(cred *cryptoops.Credential, leaseID string, alpn string)
208
defer wg.Done()
209
info, err := relay.client.GetRelayInfo()
210
if err != nil {
211
+ log.Debug().Err(err).Str("relay", relay.addr).Msg("[SDK] Failed to get relay info")
212
return
213
}
214
215
if slices.Contains(info.Leases, leaseID) {
216
+ log.Debug().Str("relay", relay.addr).Str("lease_id", leaseID).Msg("[SDK] Found lease on relay")
217
availableRelaysMu.Lock()
218
availableRelays = append(availableRelays, relay)
219
availableRelaysMu.Unlock()
@@ -205,27 +223,50 @@ func (g *RDClient) Dial(cred *cryptoops.Credential, leaseID string, alpn string)
223
wg.Wait()
224
225
if len(availableRelays) == 0 {
226
+ log.Warn().Str("lease_id", leaseID).Msg("[SDK] No available relay found for lease")
227
return nil, ErrNoAvailableRelay
228
}
229
230
+ log.Debug().Int("available_relays", len(availableRelays)).Str("lease_id", leaseID).Msg("[SDK] Attempting to connect")
231
+
232
for _, relay := range availableRelays {
233
+ log.Debug().Str("relay", relay.addr).Str("lease_id", leaseID).Msg("[SDK] Requesting connection")
234
code, conn, err := relay.client.RequestConnection(leaseID, alpn, cred)
235
if err != nil || code != rdverb.ResponseCode_RESPONSE_CODE_ACCEPTED {
236
+ log.Debug().
237
+ Err(err).
238
+ Str("relay", relay.addr).
239
+ Str("code", code.String()).
240
+ Msg("[SDK] Connection request failed, trying next relay")
241
continue
242
}
243
+ log.Debug().
244
+ Str("relay", relay.addr).
245
+ Str("lease_id", leaseID).
246
+ Str("local", conn.LocalID()).
247
+ Str("remote", conn.RemoteID()).
248
+ Msg("[SDK] Connection established successfully")
249
return &RDConnection{via: relay, conn: conn, localAddr: conn.LocalID(), remoteAddr: conn.RemoteID()}, nil
250
}
251
252
+ log.Warn().Str("lease_id", leaseID).Msg("[SDK] All connection attempts failed")
253
return nil, ErrNoAvailableRelay
254
}
255
256
func (g *RDClient) Listen(cred *cryptoops.Credential, name string, alpns []string) (*RDListener, error) {
257
+ log.Debug().
258
+ Str("lease_id", cred.ID()).
259
+ Str("name", name).
260
+ Strs("alpns", alpns).
261
+ Msg("[SDK] Creating listener")
262
+
263
g.mu.Lock()
264
defer g.mu.Unlock()
265
266
// Check if client is closed
267
select {
268
case <-g.stopch:
269
+ log.Error().Msg("[SDK] Cannot create listener, client is closed")
270
return nil, ErrClientClosed
271
default:
272
// Client is still open
@@ -233,6 +274,7 @@ func (g *RDClient) Listen(cred *cryptoops.Credential, name string, alpns []strin
274
275
// Check if listener already exists
276
if _, exists := g.listeners[cred.ID()]; exists {
277
+ log.Warn().Str("lease_id", cred.ID()).Msg("[SDK] Listener already exists")
278
return nil, ErrListenerExists
279
}
280
@@ -247,10 +289,20 @@ func (g *RDClient) Listen(cred *cryptoops.Credential, name string, alpns []strin
289
// Register listener
290
g.listeners[cred.ID()] = listener
291
292
+ log.Debug().
293
+ Str("lease_id", cred.ID()).
294
+ Int("relay_count", len(g.relays)).
295
+ Msg("[SDK] Registering lease with relays")
296
+
297
// Register lease with all available relays
298
for _, relay := range g.relays {
299
go func(r *rdRelay) {
253
- r.client.RegisterLease(cred, name, alpns)
300
+ err := r.client.RegisterLease(cred, name, alpns)
301
+ if err != nil {
302
+ log.Error().Err(err).Str("relay", r.addr).Msg("[SDK] Failed to register lease")
303
+ } else {
304
+ log.Debug().Str("relay", r.addr).Msg("[SDK] Lease registered successfully")
305
+ }
306
}(relay)
307
}
308
@@ -259,26 +311,38 @@ func (g *RDClient) Listen(cred *cryptoops.Credential, name string, alpns []strin
311
go g.listenerWorker(relay)
312
}
313
314
+ log.Debug().Str("lease_id", cred.ID()).Msg("[SDK] Listener created successfully")
315
return listener, nil
316
}
317
318
func (g *RDClient) listenerWorker(server *rdRelay) {
319
+ log.Debug().Str("relay", server.addr).Msg("[SDK] Listener worker started")
320
+
321
for {
322
select {
323
case <-server.stop:
324
+ log.Debug().Str("relay", server.addr).Msg("[SDK] Listener worker stopped")
325
return
326
case conn, ok := <-server.client.IncommingConnection():
327
if !ok {
328
+ log.Debug().Str("relay", server.addr).Msg("[SDK] Incoming connection channel closed")
329
return // Channel closed
330
}
331
332
lease := conn.LeaseID()
333
+ log.Debug().
334
+ Str("relay", server.addr).
335
+ Str("lease_id", lease).
336
+ Str("local", conn.LocalID()).
337
+ Str("remote", conn.RemoteID()).
338
+ Msg("[SDK] Received incoming connection")
339
340
g.mu.Lock()
341
listener, exists := g.listeners[lease]
342
g.mu.Unlock()
343
344
if !exists {
345
+ log.Warn().Str("lease_id", lease).Msg("[SDK] No listener found for lease, closing connection")
346
conn.SecureConnection.Close() // Close unused connection
347
continue
348
}
@@ -293,6 +357,7 @@ func (g *RDClient) listenerWorker(server *rdRelay) {
357
listener.mu.Lock()
358
// Check if listener is still active
359
if listener.closed {
360
+ log.Debug().Str("lease_id", lease).Msg("[SDK] Listener closed, rejecting connection")
361
listener.mu.Unlock()
362
rdConn.Close()
363
continue
@@ -303,9 +368,11 @@ func (g *RDClient) listenerWorker(server *rdRelay) {
368
// Send connection to listener (non-blocking)
369
select {
370
case listener.connCh <- rdConn:
371
+ log.Debug().Str("lease_id", lease).Msg("[SDK] Connection sent to listener channel")
372
// Connection sent successfully
373
default:
374
// Channel full, close connection
375
+ log.Warn().Str("lease_id", lease).Msg("[SDK] Listener channel full, closing connection")
376
listener.mu.Lock()
377
delete(listener.conns, rdConn)
378
listener.mu.Unlock()
sdk/sdk_e2e_test.go
new
+461
@@ -0,0 +1,461 @@
1
+package sdk
2
+
3
+import (
4
+ "bufio"
5
+ "context"
6
+ "fmt"
7
+ "io"
8
+ "net"
9
+ "net/http"
10
+ "os"
11
+ "testing"
12
+ "time"
13
+
14
+ "github.com/gorilla/websocket"
15
+ "github.com/gosuda/relaydns/relaydns"
16
+ "github.com/gosuda/relaydns/relaydns/core/cryptoops"
17
+ "github.com/gosuda/relaydns/relaydns/utils/wsstream"
18
+ "github.com/rs/zerolog"
19
+ "github.com/rs/zerolog/log"
20
+)
21
+
22
+func init() {
23
+ // Set zerolog to Debug level for testing
24
+ zerolog.SetGlobalLevel(zerolog.DebugLevel)
25
+ log.Logger = log.Output(zerolog.ConsoleWriter{Out: os.Stderr, TimeFormat: time.RFC3339})
26
+}
27
+
28
+// TestE2E_ClientToAppThroughRelay tests the full end-to-end flow:
29
+// SDK Client -> Relay Server -> Demo App
30
+func TestE2E_ClientToAppThroughRelay(t *testing.T) {
31
+ log.Info().Msg("=== Starting E2E Test ===")
32
+
33
+ // 1. Create relay server credential
34
+ log.Info().Msg("[TEST] Step 1: Creating relay server credential")
35
+ relayServerCred, err := cryptoops.NewCredential()
36
+ if err != nil {
37
+ t.Fatalf("Failed to create relay server credential: %v", err)
38
+ }
39
+ log.Debug().Str("relay_id", relayServerCred.ID()).Msg("[TEST] Relay server credential created")
40
+
41
+ // 2. Start relay server
42
+ log.Info().Msg("[TEST] Step 2: Starting relay server")
43
+ relayServer := relaydns.NewRelayServer(relayServerCred, []string{"ws://127.0.0.1:14017/relay"})
44
+ relayServer.Start()
45
+ defer relayServer.Stop()
46
+
47
+ // Start WebSocket server for relay
48
+ relayAddr := "127.0.0.1:14017"
49
+ relayMux := http.NewServeMux()
50
+ relayMux.HandleFunc("/relay", func(w http.ResponseWriter, r *http.Request) {
51
+ log.Debug().Str("remote", r.RemoteAddr).Msg("[TEST] Relay server accepting WebSocket connection")
52
+ upgrader := websocket.Upgrader{
53
+ CheckOrigin: func(r *http.Request) bool { return true },
54
+ }
55
+ ws, err := upgrader.Upgrade(w, r, nil)
56
+ if err != nil {
57
+ log.Error().Err(err).Msg("[TEST] Failed to upgrade WebSocket")
58
+ return
59
+ }
60
+ wsConn := &wsstream.WsStream{Conn: ws}
61
+ if err := relayServer.HandleConnection(wsConn); err != nil {
62
+ log.Error().Err(err).Msg("[TEST] Relay server error handling connection")
63
+ }
64
+ })
65
+
66
+ relayHTTPServer := &http.Server{
67
+ Addr: relayAddr,
68
+ Handler: relayMux,
69
+ }
70
+
71
+ go func() {
72
+ log.Info().Str("addr", relayAddr).Msg("[TEST] Relay HTTP server starting")
73
+ if err := relayHTTPServer.ListenAndServe(); err != nil && err != http.ErrServerClosed {
74
+ log.Error().Err(err).Msg("[TEST] Relay HTTP server error")
75
+ }
76
+ }()
77
+ defer func() {
78
+ ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second)
79
+ defer cancel()
80
+ relayHTTPServer.Shutdown(ctx)
81
+ }()
82
+
83
+ // Wait for relay server to start
84
+ time.Sleep(500 * time.Millisecond)
85
+ log.Info().Msg("[TEST] Relay server started")
86
+
87
+ // 3. Create app credential and SDK client
88
+ log.Info().Msg("[TEST] Step 3: Creating app (listener) credential")
89
+ appCred, err := NewCredential()
90
+ if err != nil {
91
+ t.Fatalf("Failed to create app credential: %v", err)
92
+ }
93
+ log.Debug().Str("app_id", appCred.ID()).Msg("[TEST] App credential created")
94
+
95
+ // 4. Create app SDK client and register listener
96
+ log.Info().Msg("[TEST] Step 4: Creating app SDK client")
97
+ appClient, err := NewClient(func(c *RDClientConfig) {
98
+ c.BootstrapServers = []string{"ws://127.0.0.1:14017/relay"}
99
+ })
100
+ if err != nil {
101
+ t.Fatalf("Failed to create app SDK client: %v", err)
102
+ }
103
+ defer appClient.Close()
104
+ log.Info().Msg("[TEST] App SDK client created")
105
+
106
+ // 5. Register listener on app side
107
+ log.Info().Msg("[TEST] Step 5: Registering app listener")
108
+ appListener, err := appClient.Listen(appCred, "test-app", []string{"h1"})
109
+ if err != nil {
110
+ t.Fatalf("Failed to create app listener: %v", err)
111
+ }
112
+ defer appListener.Close()
113
+ log.Info().Str("lease_id", appCred.ID()).Msg("[TEST] App listener registered")
114
+
115
+ // 6. Start serving HTTP on app listener
116
+ log.Info().Msg("[TEST] Step 6: Starting HTTP server on app listener")
117
+ appMux := http.NewServeMux()
118
+ testMessage := "Hello from E2E Test App!"
119
+ appMux.HandleFunc("/", func(w http.ResponseWriter, r *http.Request) {
120
+ log.Debug().
121
+ Str("method", r.Method).
122
+ Str("path", r.URL.Path).
123
+ Str("remote", r.RemoteAddr).
124
+ Msg("[TEST] App received HTTP request")
125
+ fmt.Fprintf(w, "%s\n", testMessage)
126
+ fmt.Fprintf(w, "Request from: %s\n", r.RemoteAddr)
127
+ })
128
+
129
+ appErrChan := make(chan error, 1)
130
+ go func() {
131
+ log.Info().Msg("[TEST] App HTTP server starting on listener")
132
+ appErrChan <- http.Serve(appListener, appMux)
133
+ }()
134
+
135
+ // Wait for app to be ready
136
+ time.Sleep(1 * time.Second)
137
+ log.Info().Msg("[TEST] App is ready to accept connections")
138
+
139
+ // 7. Create client credential
140
+ log.Info().Msg("[TEST] Step 7: Creating client credential")
141
+ clientCred, err := NewCredential()
142
+ if err != nil {
143
+ t.Fatalf("Failed to create client credential: %v", err)
144
+ }
145
+ log.Debug().Str("client_id", clientCred.ID()).Msg("[TEST] Client credential created")
146
+
147
+ // 8. Create client SDK client
148
+ log.Info().Msg("[TEST] Step 8: Creating client SDK client")
149
+ clientSDK, err := NewClient(func(c *RDClientConfig) {
150
+ c.BootstrapServers = []string{"ws://127.0.0.1:14017/relay"}
151
+ })
152
+ if err != nil {
153
+ t.Fatalf("Failed to create client SDK: %v", err)
154
+ }
155
+ defer clientSDK.Close()
156
+ log.Info().Msg("[TEST] Client SDK client created")
157
+
158
+ // 9. Wait for lease to be fully registered
159
+ log.Info().Msg("[TEST] Step 9: Waiting for lease propagation")
160
+ time.Sleep(2 * time.Second)
161
+
162
+ // 10. Dial to app through relay
163
+ log.Info().Msg("[TEST] Step 10: Client dialing to app through relay")
164
+ conn, err := clientSDK.Dial(clientCred, appCred.ID(), "h1")
165
+ if err != nil {
166
+ t.Fatalf("Failed to dial to app: %v", err)
167
+ }
168
+ defer conn.Close()
169
+ log.Info().
170
+ Str("local", conn.LocalAddr().String()).
171
+ Str("remote", conn.RemoteAddr().String()).
172
+ Msg("[TEST] Connection established")
173
+
174
+ // 11. Send HTTP request through the connection
175
+ log.Info().Msg("[TEST] Step 11: Sending HTTP request through connection")
176
+
177
+ // Create HTTP request
178
+ req, err := http.NewRequest("GET", "http://test-app/", nil)
179
+ if err != nil {
180
+ t.Fatalf("Failed to create HTTP request: %v", err)
181
+ }
182
+
183
+ // Write HTTP request to connection
184
+ if err := req.Write(conn); err != nil {
185
+ t.Fatalf("Failed to write HTTP request: %v", err)
186
+ }
187
+ log.Debug().Msg("[TEST] HTTP request sent")
188
+
189
+ // Read HTTP response
190
+ log.Info().Msg("[TEST] Step 12: Reading HTTP response")
191
+ resp, err := http.ReadResponse(bufio.NewReader(conn), req)
192
+ if err != nil {
193
+ t.Fatalf("Failed to read HTTP response: %v", err)
194
+ }
195
+ defer resp.Body.Close()
196
+
197
+ log.Debug().
198
+ Int("status_code", resp.StatusCode).
199
+ Str("status", resp.Status).
200
+ Msg("[TEST] HTTP response received")
201
+
202
+ // Read response body
203
+ body, err := io.ReadAll(resp.Body)
204
+ if err != nil {
205
+ t.Fatalf("Failed to read response body: %v", err)
206
+ }
207
+
208
+ responseStr := string(body)
209
+ log.Info().Str("body", responseStr).Msg("[TEST] Response body received")
210
+
211
+ // 12. Verify response
212
+ log.Info().Msg("[TEST] Step 13: Verifying response")
213
+ if resp.StatusCode != 200 {
214
+ t.Errorf("Expected status code 200, got %d", resp.StatusCode)
215
+ }
216
+
217
+ if len(body) == 0 {
218
+ t.Error("Expected non-empty response body")
219
+ }
220
+
221
+ // Check if response contains test message
222
+ bodyStr := string(body)
223
+ if len(bodyStr) == 0 {
224
+ t.Error("Response body is empty")
225
+ } else {
226
+ log.Info().Str("response", bodyStr).Msg("[TEST] Response verification successful")
227
+ }
228
+
229
+ log.Info().Msg("=== E2E Test Completed Successfully ===")
230
+}
231
+
232
+// TestE2E_MultipleConnections tests multiple concurrent connections
233
+func TestE2E_MultipleConnections(t *testing.T) {
234
+ log.Info().Msg("=== Starting Multiple Connections Test ===")
235
+
236
+ // Setup relay server
237
+ relayServerCred, err := cryptoops.NewCredential()
238
+ if err != nil {
239
+ t.Fatalf("Failed to create relay server credential: %v", err)
240
+ }
241
+
242
+ relayServer := relaydns.NewRelayServer(relayServerCred, []string{"ws://127.0.0.1:14018/relay"})
243
+ relayServer.Start()
244
+ defer relayServer.Stop()
245
+
246
+ relayAddr := "127.0.0.1:14018"
247
+ relayMux := http.NewServeMux()
248
+ relayMux.HandleFunc("/relay", func(w http.ResponseWriter, r *http.Request) {
249
+ upgrader := websocket.Upgrader{
250
+ CheckOrigin: func(r *http.Request) bool { return true },
251
+ }
252
+ ws, err := upgrader.Upgrade(w, r, nil)
253
+ if err != nil {
254
+ return
255
+ }
256
+ wsConn := &wsstream.WsStream{Conn: ws}
257
+ relayServer.HandleConnection(wsConn)
258
+ })
259
+
260
+ relayHTTPServer := &http.Server{
261
+ Addr: relayAddr,
262
+ Handler: relayMux,
263
+ }
264
+
265
+ go relayHTTPServer.ListenAndServe()
266
+ defer func() {
267
+ ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second)
268
+ defer cancel()
269
+ relayHTTPServer.Shutdown(ctx)
270
+ }()
271
+
272
+ time.Sleep(500 * time.Millisecond)
273
+
274
+ // Setup app
275
+ appCred, err := NewCredential()
276
+ if err != nil {
277
+ t.Fatalf("Failed to create app credential: %v", err)
278
+ }
279
+
280
+ appClient, err := NewClient(func(c *RDClientConfig) {
281
+ c.BootstrapServers = []string{"ws://127.0.0.1:14018/relay"}
282
+ })
283
+ if err != nil {
284
+ t.Fatalf("Failed to create app SDK client: %v", err)
285
+ }
286
+ defer appClient.Close()
287
+
288
+ appListener, err := appClient.Listen(appCred, "multi-test-app", []string{"h1"})
289
+ if err != nil {
290
+ t.Fatalf("Failed to create app listener: %v", err)
291
+ }
292
+ defer appListener.Close()
293
+
294
+ // Serve echo server
295
+ go func() {
296
+ for {
297
+ conn, err := appListener.Accept()
298
+ if err != nil {
299
+ return
300
+ }
301
+ go func(c net.Conn) {
302
+ defer c.Close()
303
+ io.Copy(c, c) // Echo back
304
+ }(conn)
305
+ }
306
+ }()
307
+
308
+ time.Sleep(1 * time.Second)
309
+
310
+ // Create client
311
+ clientCred, err := NewCredential()
312
+ if err != nil {
313
+ t.Fatalf("Failed to create client credential: %v", err)
314
+ }
315
+
316
+ clientSDK, err := NewClient(func(c *RDClientConfig) {
317
+ c.BootstrapServers = []string{"ws://127.0.0.1:14018/relay"}
318
+ })
319
+ if err != nil {
320
+ t.Fatalf("Failed to create client SDK: %v", err)
321
+ }
322
+ defer clientSDK.Close()
323
+
324
+ time.Sleep(2 * time.Second)
325
+
326
+ // Test multiple concurrent connections
327
+ numConnections := 5
328
+ log.Info().Int("count", numConnections).Msg("[TEST] Testing multiple concurrent connections")
329
+
330
+ for i := 0; i < numConnections; i++ {
331
+ i := i
332
+ go func() {
333
+ log.Debug().Int("conn_num", i).Msg("[TEST] Starting connection")
334
+
335
+ conn, err := clientSDK.Dial(clientCred, appCred.ID(), "h1")
336
+ if err != nil {
337
+ t.Errorf("Connection %d failed to dial: %v", i, err)
338
+ return
339
+ }
340
+ defer conn.Close()
341
+
342
+ testData := fmt.Sprintf("test-message-%d", i)
343
+
344
+ // Write test data
345
+ if _, err := conn.Write([]byte(testData)); err != nil {
346
+ t.Errorf("Connection %d failed to write: %v", i, err)
347
+ return
348
+ }
349
+
350
+ // Read echoed data
351
+ buf := make([]byte, len(testData))
352
+ if _, err := io.ReadFull(conn, buf); err != nil {
353
+ t.Errorf("Connection %d failed to read: %v", i, err)
354
+ return
355
+ }
356
+
357
+ if string(buf) != testData {
358
+ t.Errorf("Connection %d: expected %q, got %q", i, testData, string(buf))
359
+ return
360
+ }
361
+
362
+ log.Debug().Int("conn_num", i).Msg("[TEST] Connection successful")
363
+ }()
364
+ }
365
+
366
+ time.Sleep(5 * time.Second)
367
+ log.Info().Msg("=== Multiple Connections Test Completed ===")
368
+}
369
+
370
+// TestE2E_ConnectionTimeout tests timeout scenarios
371
+func TestE2E_ConnectionTimeout(t *testing.T) {
372
+ log.Info().Msg("=== Starting Connection Timeout Test ===")
373
+
374
+ // Create client with non-existent relay
375
+ clientCred, err := NewCredential()
376
+ if err != nil {
377
+ t.Fatalf("Failed to create client credential: %v", err)
378
+ }
379
+
380
+ // This should fail or timeout appropriately
381
+ ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second)
382
+ defer cancel()
383
+
384
+ done := make(chan error, 1)
385
+ go func() {
386
+ _, err := NewClient(func(c *RDClientConfig) {
387
+ c.BootstrapServers = []string{"ws://127.0.0.1:19999/relay"} // Non-existent
388
+ })
389
+ done <- err
390
+ }()
391
+
392
+ select {
393
+ case err := <-done:
394
+ if err == nil {
395
+ t.Error("Expected error when connecting to non-existent relay")
396
+ } else {
397
+ log.Info().Err(err).Msg("[TEST] Got expected error")
398
+ }
399
+ case <-ctx.Done():
400
+ t.Error("Connection attempt did not complete within timeout")
401
+ }
402
+
403
+ // Try to dial to non-existent lease
404
+ relayServerCred, err := cryptoops.NewCredential()
405
+ if err != nil {
406
+ t.Fatalf("Failed to create relay server credential: %v", err)
407
+ }
408
+
409
+ relayServer := relaydns.NewRelayServer(relayServerCred, []string{"ws://127.0.0.1:14019/relay"})
410
+ relayServer.Start()
411
+ defer relayServer.Stop()
412
+
413
+ relayAddr := "127.0.0.1:14019"
414
+ relayMux := http.NewServeMux()
415
+ relayMux.HandleFunc("/relay", func(w http.ResponseWriter, r *http.Request) {
416
+ upgrader := websocket.Upgrader{
417
+ CheckOrigin: func(r *http.Request) bool { return true },
418
+ }
419
+ ws, err := upgrader.Upgrade(w, r, nil)
420
+ if err != nil {
421
+ return
422
+ }
423
+ wsConn := &wsstream.WsStream{Conn: ws}
424
+ relayServer.HandleConnection(wsConn)
425
+ })
426
+
427
+ relayHTTPServer := &http.Server{
428
+ Addr: relayAddr,
429
+ Handler: relayMux,
430
+ }
431
+
432
+ go relayHTTPServer.ListenAndServe()
433
+ defer func() {
434
+ ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second)
435
+ defer cancel()
436
+ relayHTTPServer.Shutdown(ctx)
437
+ }()
438
+
439
+ time.Sleep(500 * time.Millisecond)
440
+
441
+ clientSDK, err := NewClient(func(c *RDClientConfig) {
442
+ c.BootstrapServers = []string{"ws://127.0.0.1:14019/relay"}
443
+ })
444
+ if err != nil {
445
+ t.Fatalf("Failed to create client SDK: %v", err)
446
+ }
447
+ defer clientSDK.Close()
448
+
449
+ time.Sleep(1 * time.Second)
450
+
451
+ // Try to dial to non-existent lease
452
+ log.Info().Msg("[TEST] Attempting to dial non-existent lease")
453
+ _, err = clientSDK.Dial(clientCred, "non-existent-lease-id", "h1")
454
+ if err == nil {
455
+ t.Error("Expected error when dialing non-existent lease")
456
+ } else {
457
+ log.Info().Err(err).Msg("[TEST] Got expected error for non-existent lease")
458
+ }
459
+
460
+ log.Info().Msg("=== Connection Timeout Test Completed ===")
461
+}