feat: add health monitoring and reconnection to relay clients
- Add Ping() method to RelayClient for connectivity checks - Extend RDClientConfig with health check interval, max retries, and reconnect interval - Add option functions for configuring new RDClient settings - Implement healthCheckWorker to periodically verify relay health and trigger reconnections - Introduce reconnectRelay function to handle connection recovery with configurable retries - Update RDClient and rdRelay structs to support the new features, ensuring reliable client-server connections via automatic failover and health monitoring
lemon-mint committed
Oct 29, 2025 at 15:30 UTC
6ede4eb8ee50218e1c902963ab992b3283f4d1c3
2 files changed
+226
-5
relaydns/client.go
+4
@@ -89,6 +89,10 @@ func NewRelayClient(conn io.ReadWriteCloser) *RelayClient {
89
return g
90
}
91
92
+func (g *RelayClient) Ping() (time.Duration, error) {
93
+ return g.sess.Ping()
94
+}
95
+
96
// Close는 서버와의 연결을 종료합니다.
97
func (g *RelayClient) Close() error {
98
log.Debug().Msg("[RelayClient] Closing relay client")
sdk/sdk.go
+222
-5
@@ -54,16 +54,51 @@ func webSocketDialer() func(context.Context, string) (io.ReadWriteCloser, error)
54
}
55
56
type RDClientConfig struct {
57
- BootstrapServers []string
58
- Dialer func(context.Context, string) (io.ReadWriteCloser, error)
57
+ BootstrapServers []string
58
+ Dialer func(context.Context, string) (io.ReadWriteCloser, error)
59
+ HealthCheckInterval time.Duration // Interval for health checks (default: 10 seconds)
60
+ ReconnectMaxRetries int // Maximum reconnection attempts (default: 3, 0 = infinite)
61
+ ReconnectInterval time.Duration // Interval between reconnection attempts (default: 5 seconds)
62
}
63
64
type Option func(*RDClientConfig)
65
66
+func WithBootstrapServers(servers []string) Option {
67
+ return func(c *RDClientConfig) {
68
+ c.BootstrapServers = servers
69
+ }
70
+}
71
+
72
+func WithDialer(dialer func(context.Context, string) (io.ReadWriteCloser, error)) Option {
73
+ return func(c *RDClientConfig) {
74
+ c.Dialer = dialer
75
+ }
76
+}
77
+
78
+func WithHealthCheckInterval(interval time.Duration) Option {
79
+ return func(c *RDClientConfig) {
80
+ c.HealthCheckInterval = interval
81
+ }
82
+}
83
+
84
+func WithReconnectMaxRetries(retries int) Option {
85
+ return func(c *RDClientConfig) {
86
+ c.ReconnectMaxRetries = retries
87
+ }
88
+}
89
+
90
+func WithReconnectInterval(interval time.Duration) Option {
91
+ return func(c *RDClientConfig) {
92
+ c.ReconnectInterval = interval
93
+ }
94
+}
95
+
96
type rdRelay struct {
97
addr string
98
client *relaydns.RelayClient
99
+ dialer func(context.Context, string) (io.ReadWriteCloser, error)
100
stop chan struct{}
101
+ mu sync.Mutex
102
}
103
104
var _ net.Conn = (*RDConnection)(nil)
@@ -136,6 +171,7 @@ type RDClient struct {
171
172
relays map[string]*rdRelay
173
listeners map[string]*RDListener
174
+ config *RDClientConfig
175
176
stopch chan struct{}
177
waitGroup sync.WaitGroup // Track all listener workers
@@ -155,7 +191,10 @@ func NewClient(opt ...Option) (*RDClient, error) {
191
log.Debug().Msg("[SDK] Creating new RDClient")
192
193
config := &RDClientConfig{
158
- Dialer: webSocketDialer(),
194
+ Dialer: webSocketDialer(),
195
+ HealthCheckInterval: 10 * time.Second,
196
+ ReconnectMaxRetries: 9,
197
+ ReconnectInterval: 5 * time.Second,
198
}
199
200
for _, o := range opt {
@@ -165,6 +204,7 @@ func NewClient(opt ...Option) (*RDClient, error) {
204
client := &RDClient{
205
relays: make(map[string]*rdRelay),
206
listeners: make(map[string]*RDListener),
207
+ config: config,
208
stopch: make(chan struct{}),
209
}
210
@@ -188,11 +228,17 @@ func NewClient(opt ...Option) (*RDClient, error) {
228
}
229
230
log.Debug().Str("server", server).Msg("[SDK] Successfully connected to bootstrap server")
191
- client.relays[server] = &rdRelay{
231
+ relay := &rdRelay{
232
addr: server,
233
client: relayClient,
234
+ dialer: config.Dialer,
235
stop: make(chan struct{}),
236
}
237
+ client.relays[server] = relay
238
+
239
+ // Start health monitoring for this relay
240
+ client.waitGroup.Add(1)
241
+ go client.healthCheckWorker(relay, config)
242
}
243
244
// If no relays were successfully connected, return an error
@@ -448,6 +494,171 @@ func (g *RDClient) Close() error {
494
return nil
495
}
496
497
+// healthCheckWorker periodically checks relay health and reconnects if needed
498
+func (g *RDClient) healthCheckWorker(relay *rdRelay, config *RDClientConfig) {
499
+ defer g.waitGroup.Done()
500
+
501
+ ticker := time.NewTicker(config.HealthCheckInterval)
502
+ defer ticker.Stop()
503
+
504
+ log.Debug().Str("relay", relay.addr).Msg("[SDK] Health check worker started")
505
+
506
+ for {
507
+ select {
508
+ case <-g.stopch:
509
+ log.Debug().Str("relay", relay.addr).Msg("[SDK] Health check worker stopped")
510
+ return
511
+ case <-relay.stop:
512
+ log.Debug().Str("relay", relay.addr).Msg("[SDK] Relay stopped, health check worker exiting")
513
+ return
514
+ case <-ticker.C:
515
+ relay.mu.Lock()
516
+ client := relay.client
517
+ relay.mu.Unlock()
518
+
519
+ if client == nil {
520
+ log.Warn().Str("relay", relay.addr).Msg("[SDK] Relay client is nil, attempting reconnection")
521
+ g.reconnectRelay(relay, config)
522
+ continue
523
+ }
524
+
525
+ // Perform health check using Ping
526
+ _, err := client.Ping()
527
+ if err != nil {
528
+ log.Warn().
529
+ Err(err).
530
+ Str("relay", relay.addr).
531
+ Msg("[SDK] Health check failed, attempting reconnection")
532
+ g.reconnectRelay(relay, config)
533
+ } else {
534
+ log.Debug().Str("relay", relay.addr).Msg("[SDK] Health check passed")
535
+ }
536
+ }
537
+ }
538
+}
539
+
540
+// reconnectRelay attempts to reconnect to a relay server
541
+func (g *RDClient) reconnectRelay(relay *rdRelay, config *RDClientConfig) {
542
+ relay.mu.Lock()
543
+
544
+ // Close old client if exists
545
+ if relay.client != nil {
546
+ log.Debug().Str("relay", relay.addr).Msg("[SDK] Closing old relay client")
547
+ relay.client.Close()
548
+ relay.client = nil
549
+ }
550
+ relay.mu.Unlock()
551
+
552
+ maxRetries := config.ReconnectMaxRetries
553
+ if maxRetries == 0 {
554
+ maxRetries = -1 // Infinite retries
555
+ }
556
+
557
+ attempt := 0
558
+ for {
559
+ // Check if we should stop
560
+ select {
561
+ case <-g.stopch:
562
+ log.Debug().Str("relay", relay.addr).Msg("[SDK] Client stopped, abandoning reconnection")
563
+ return
564
+ case <-relay.stop:
565
+ log.Debug().Str("relay", relay.addr).Msg("[SDK] Relay stopped, abandoning reconnection")
566
+ return
567
+ default:
568
+ }
569
+
570
+ attempt++
571
+ if maxRetries > 0 && attempt > maxRetries {
572
+ log.Error().
573
+ Str("relay", relay.addr).
574
+ Int("attempts", attempt-1).
575
+ Msg("[SDK] Max reconnection attempts reached, giving up")
576
+ return
577
+ }
578
+
579
+ log.Debug().
580
+ Str("relay", relay.addr).
581
+ Int("attempt", attempt).
582
+ Msg("[SDK] Attempting to reconnect")
583
+
584
+ // Attempt to connect
585
+ conn, err := relay.dialer(context.Background(), relay.addr)
586
+ if err != nil {
587
+ log.Warn().
588
+ Err(err).
589
+ Str("relay", relay.addr).
590
+ Int("attempt", attempt).
591
+ Msg("[SDK] Reconnection attempt failed")
592
+
593
+ // Wait before next retry
594
+ select {
595
+ case <-g.stopch:
596
+ return
597
+ case <-relay.stop:
598
+ return
599
+ case <-time.After(config.ReconnectInterval):
600
+ continue
601
+ }
602
+ }
603
+
604
+ // Create new relay client
605
+ relayClient := relaydns.NewRelayClient(conn)
606
+ if relayClient == nil {
607
+ log.Error().Str("relay", relay.addr).Msg("[SDK] Failed to create relay client after reconnection")
608
+ conn.Close()
609
+
610
+ // Wait before next retry
611
+ select {
612
+ case <-g.stopch:
613
+ return
614
+ case <-relay.stop:
615
+ return
616
+ case <-time.After(config.ReconnectInterval):
617
+ continue
618
+ }
619
+ }
620
+
621
+ relay.mu.Lock()
622
+ relay.client = relayClient
623
+ relay.mu.Unlock()
624
+
625
+ log.Info().
626
+ Str("relay", relay.addr).
627
+ Int("attempt", attempt).
628
+ Msg("[SDK] Successfully reconnected to relay")
629
+
630
+ // Re-register all leases with the reconnected relay
631
+ g.mu.Lock()
632
+ for _, listener := range g.listeners {
633
+ go func(l *RDListener) {
634
+ l.mu.Lock()
635
+ cred := l.cred
636
+ lease := l.lease
637
+ l.mu.Unlock()
638
+
639
+ if lease != nil {
640
+ err := relayClient.RegisterLease(cred, lease.Name, lease.Alpn)
641
+ if err != nil {
642
+ log.Error().
643
+ Err(err).
644
+ Str("relay", relay.addr).
645
+ Str("lease_id", cred.ID()).
646
+ Msg("[SDK] Failed to re-register lease after reconnection")
647
+ } else {
648
+ log.Debug().
649
+ Str("relay", relay.addr).
650
+ Str("lease_id", cred.ID()).
651
+ Msg("[SDK] Lease re-registered after reconnection")
652
+ }
653
+ }
654
+ }(listener)
655
+ }
656
+ g.mu.Unlock()
657
+
658
+ return
659
+ }
660
+}
661
+
662
// Implement net.Listener interface for RDListener
663
func (l *RDListener) Accept() (net.Conn, error) {
664
conn, ok := <-l.connCh
@@ -513,11 +724,17 @@ func (g *RDClient) AddRelay(addr string, dialer func(context.Context, string) (i
724
}
725
726
// Add relay
516
- g.relays[addr] = &rdRelay{
727
+ relay := &rdRelay{
728
addr: addr,
729
client: relayClient,
730
+ dialer: dialer,
731
stop: make(chan struct{}),
732
}
733
+ g.relays[addr] = relay
734
+
735
+ // Start health monitoring for this relay
736
+ g.waitGroup.Add(1)
737
+ go g.healthCheckWorker(relay, g.config)
738
739
return nil
740
}