gofmt all files
Kim committed
Oct 21, 2025 at 12:00 UTC
34d7b2225b2ae48d3caeebc6d85b69baaf754d1b
5 files changed
+721
-721
cmd/example_client/main.go
+142
-142
@@ -1,142 +1,142 @@
1
-package main
2
-
3
-import (
4
- "context"
5
- "fmt"
6
- "net/http"
7
- "os"
8
- "os/signal"
9
- "syscall"
10
- "text/template"
11
- "time"
12
-
13
- "github.com/gosuda/relaydns/relaydns"
14
- "github.com/rs/zerolog/log"
15
- "github.com/spf13/cobra"
16
-)
17
-
18
-var rootCmd = &cobra.Command{
19
- Use: "relaydns-client",
20
- Short: "RelayDNS demo client (local HTTP backend + libp2p advertiser)",
21
- RunE: runClient,
22
-}
23
-
24
-var (
25
- flagServerURL string
26
- flagBootstraps []string
27
- flagBackendHTTP string
28
-)
29
-
30
-func init() {
31
- flags := rootCmd.PersistentFlags()
32
- flags.StringVar(&flagServerURL, "server-url", "http://localhost:8080", "relayserver admin base URL to auto-fetch multiaddrs from /health")
33
- flags.StringSliceVar(&flagBootstraps, "bootstrap", nil, "multiaddrs with /p2p/ (supports /dnsaddr/ that resolves to /p2p/)")
34
- flags.StringVar(&flagBackendHTTP, "backend-http", ":8081", "local backend HTTP listen address")
35
-}
36
-
37
-func main() {
38
- if err := rootCmd.Execute(); err != nil {
39
- log.Fatal().Err(err).Msg("execute root command")
40
- }
41
-}
42
-
43
-func runClient(cmd *cobra.Command, args []string) error {
44
- ctx, cancel := context.WithCancel(context.Background())
45
- defer cancel()
46
-
47
- // 1) HTTP backend
48
- go func() {
49
- mux := http.NewServeMux()
50
- mux.HandleFunc("/", func(w http.ResponseWriter, r *http.Request) {
51
- data := struct {
52
- Now string
53
- Host string
54
- Addr string
55
- }{
56
- Now: time.Now().Format(time.RFC1123),
57
- Host: r.Host,
58
- Addr: flagBackendHTTP,
59
- }
60
- _ = pageTmpl.Execute(w, data)
61
- })
62
- mux.HandleFunc("/healthz", func(w http.ResponseWriter, r *http.Request) {
63
- w.WriteHeader(http.StatusOK)
64
- _, _ = w.Write([]byte("ok"))
65
- })
66
-
67
- log.Info().Msgf("[client] local backend http %s", flagBackendHTTP)
68
- if err := http.ListenAndServe(flagBackendHTTP, mux); err != nil {
69
- log.Error().Err(err).Msg("[client] http backend error")
70
- cancel()
71
- }
72
- }()
73
-
74
- // 2) libp2p host
75
- h, err := relaydns.MakeHost(ctx, 0, true)
76
- if err != nil {
77
- return fmt.Errorf("make host: %w", err)
78
- }
79
-
80
- client, err := relaydns.NewClient(ctx, h, relaydns.ClientConfig{
81
- Protocol: "/relaydns/http/1.0",
82
- Topic: "relaydns.backends",
83
- AdvertiseEvery: 3 * time.Second,
84
- TargetTCP: addrToTarget(flagBackendHTTP),
85
-
86
- ServerURL: flagServerURL,
87
- Bootstraps: flagBootstraps,
88
- HTTPTimeout: 3 * time.Second,
89
- PreferQUIC: true,
90
- PreferLocal: true,
91
- })
92
- if err != nil {
93
- return fmt.Errorf("new client: %w", err)
94
- }
95
- defer client.Close()
96
-
97
- if addrs := h.Addrs(); len(addrs) > 0 {
98
- for _, a := range addrs {
99
- log.Info().Msgf("[client] host addr: %s/p2p/%s", a.String(), h.ID().String())
100
- }
101
- } else {
102
- log.Info().Msgf("[client] host peer: %s (no listen addrs yet)", h.ID().String())
103
- }
104
-
105
- // wait for termination
106
- sig := make(chan os.Signal, 1)
107
- signal.Notify(sig, syscall.SIGINT, syscall.SIGTERM)
108
- <-sig
109
- log.Info().Msg("[client] shutting down")
110
- time.Sleep(200 * time.Millisecond)
111
- return nil
112
-}
113
-
114
-func addrToTarget(listen string) string {
115
- if len(listen) > 0 && listen[0] == ':' {
116
- return "127.0.0.1" + listen
117
- }
118
- return listen
119
-}
120
-
121
-var pageTmpl = template.Must(template.New("index").Parse(`<!DOCTYPE html>
122
-<html lang="en">
123
-<head>
124
- <meta charset="UTF-8">
125
- <title>RelayDNS Backend</title>
126
- <style>
127
- body { font-family: sans-serif; background: #f9f9f9; padding: 40px; }
128
- h1 { color: #333; }
129
- footer { margin-top: 40px; color: #666; font-size: 0.9em; }
130
- .card { background: white; border-radius: 12px; padding: 24px; box-shadow: 0 2px 6px rgba(0,0,0,0.1); }
131
- </style>
132
-</head>
133
-<body>
134
- <div class="card">
135
- <h1>🚀 RelayDNS Backend</h1>
136
- <p>This page is served from the backend node.</p>
137
- <p>Current time: <b>{{.Now}}</b></p>
138
- <p>Hostname: <b>{{.Host}}</b></p>
139
- </div>
140
- <footer>relaydns demo client — served locally at {{.Addr}}</footer>
141
-</body>
142
-</html>`))
1
+package main
2
+
3
+import (
4
+ "context"
5
+ "fmt"
6
+ "net/http"
7
+ "os"
8
+ "os/signal"
9
+ "syscall"
10
+ "text/template"
11
+ "time"
12
+
13
+ "github.com/gosuda/relaydns/relaydns"
14
+ "github.com/rs/zerolog/log"
15
+ "github.com/spf13/cobra"
16
+)
17
+
18
+var rootCmd = &cobra.Command{
19
+ Use: "relaydns-client",
20
+ Short: "RelayDNS demo client (local HTTP backend + libp2p advertiser)",
21
+ RunE: runClient,
22
+}
23
+
24
+var (
25
+ flagServerURL string
26
+ flagBootstraps []string
27
+ flagBackendHTTP string
28
+)
29
+
30
+func init() {
31
+ flags := rootCmd.PersistentFlags()
32
+ flags.StringVar(&flagServerURL, "server-url", "http://localhost:8080", "relayserver admin base URL to auto-fetch multiaddrs from /health")
33
+ flags.StringSliceVar(&flagBootstraps, "bootstrap", nil, "multiaddrs with /p2p/ (supports /dnsaddr/ that resolves to /p2p/)")
34
+ flags.StringVar(&flagBackendHTTP, "backend-http", ":8081", "local backend HTTP listen address")
35
+}
36
+
37
+func main() {
38
+ if err := rootCmd.Execute(); err != nil {
39
+ log.Fatal().Err(err).Msg("execute root command")
40
+ }
41
+}
42
+
43
+func runClient(cmd *cobra.Command, args []string) error {
44
+ ctx, cancel := context.WithCancel(context.Background())
45
+ defer cancel()
46
+
47
+ // 1) HTTP backend
48
+ go func() {
49
+ mux := http.NewServeMux()
50
+ mux.HandleFunc("/", func(w http.ResponseWriter, r *http.Request) {
51
+ data := struct {
52
+ Now string
53
+ Host string
54
+ Addr string
55
+ }{
56
+ Now: time.Now().Format(time.RFC1123),
57
+ Host: r.Host,
58
+ Addr: flagBackendHTTP,
59
+ }
60
+ _ = pageTmpl.Execute(w, data)
61
+ })
62
+ mux.HandleFunc("/healthz", func(w http.ResponseWriter, r *http.Request) {
63
+ w.WriteHeader(http.StatusOK)
64
+ _, _ = w.Write([]byte("ok"))
65
+ })
66
+
67
+ log.Info().Msgf("[client] local backend http %s", flagBackendHTTP)
68
+ if err := http.ListenAndServe(flagBackendHTTP, mux); err != nil {
69
+ log.Error().Err(err).Msg("[client] http backend error")
70
+ cancel()
71
+ }
72
+ }()
73
+
74
+ // 2) libp2p host
75
+ h, err := relaydns.MakeHost(ctx, 0, true)
76
+ if err != nil {
77
+ return fmt.Errorf("make host: %w", err)
78
+ }
79
+
80
+ client, err := relaydns.NewClient(ctx, h, relaydns.ClientConfig{
81
+ Protocol: "/relaydns/http/1.0",
82
+ Topic: "relaydns.backends",
83
+ AdvertiseEvery: 3 * time.Second,
84
+ TargetTCP: addrToTarget(flagBackendHTTP),
85
+
86
+ ServerURL: flagServerURL,
87
+ Bootstraps: flagBootstraps,
88
+ HTTPTimeout: 3 * time.Second,
89
+ PreferQUIC: true,
90
+ PreferLocal: true,
91
+ })
92
+ if err != nil {
93
+ return fmt.Errorf("new client: %w", err)
94
+ }
95
+ defer client.Close()
96
+
97
+ if addrs := h.Addrs(); len(addrs) > 0 {
98
+ for _, a := range addrs {
99
+ log.Info().Msgf("[client] host addr: %s/p2p/%s", a.String(), h.ID().String())
100
+ }
101
+ } else {
102
+ log.Info().Msgf("[client] host peer: %s (no listen addrs yet)", h.ID().String())
103
+ }
104
+
105
+ // wait for termination
106
+ sig := make(chan os.Signal, 1)
107
+ signal.Notify(sig, syscall.SIGINT, syscall.SIGTERM)
108
+ <-sig
109
+ log.Info().Msg("[client] shutting down")
110
+ time.Sleep(200 * time.Millisecond)
111
+ return nil
112
+}
113
+
114
+func addrToTarget(listen string) string {
115
+ if len(listen) > 0 && listen[0] == ':' {
116
+ return "127.0.0.1" + listen
117
+ }
118
+ return listen
119
+}
120
+
121
+var pageTmpl = template.Must(template.New("index").Parse(`<!DOCTYPE html>
122
+<html lang="en">
123
+<head>
124
+ <meta charset="UTF-8">
125
+ <title>RelayDNS Backend</title>
126
+ <style>
127
+ body { font-family: sans-serif; background: #f9f9f9; padding: 40px; }
128
+ h1 { color: #333; }
129
+ footer { margin-top: 40px; color: #666; font-size: 0.9em; }
130
+ .card { background: white; border-radius: 12px; padding: 24px; box-shadow: 0 2px 6px rgba(0,0,0,0.1); }
131
+ </style>
132
+</head>
133
+<body>
134
+ <div class="card">
135
+ <h1>🚀 RelayDNS Backend</h1>
136
+ <p>This page is served from the backend node.</p>
137
+ <p>Current time: <b>{{.Now}}</b></p>
138
+ <p>Hostname: <b>{{.Host}}</b></p>
139
+ </div>
140
+ <footer>relaydns demo client — served locally at {{.Addr}}</footer>
141
+</body>
142
+</html>`))
relaydns/client.go
+237
-237
@@ -1,237 +1,237 @@
1
-package relaydns
2
-
3
-import (
4
- "context"
5
- "encoding/json"
6
- "errors"
7
- "fmt"
8
- "io"
9
- "net"
10
- "net/http"
11
- "net/url"
12
- "sort"
13
- "strings"
14
- "sync"
15
- "time"
16
-
17
- pubsub "github.com/libp2p/go-libp2p-pubsub"
18
- "github.com/libp2p/go-libp2p/core/host"
19
- "github.com/libp2p/go-libp2p/core/network"
20
- "github.com/libp2p/go-libp2p/core/protocol"
21
- "github.com/rs/zerolog/log"
22
-)
23
-
24
-type ClientConfig struct {
25
- ServerURL string
26
- Bootstraps []string
27
- HTTPTimeout time.Duration
28
- PreferQUIC bool
29
- PreferLocal bool
30
-
31
- // libp2p stream protocol id (e.g. "/relaydns/ssh/1.0")
32
- Protocol string
33
- // pubsub topic for backend adverts (e.g. "relaydns.backends")
34
- Topic string
35
- // advertise interval
36
- AdvertiseEvery time.Duration
37
- // optional metadata
38
- Name string
39
- DNS string
40
-
41
- // One of the following:
42
- // 1) Provide a custom stream handler
43
- Handler func(s network.Stream)
44
- // 2) Or just set TargetTCP to auto-pipe bytes to a local TCP service (e.g. "127.0.0.1:22")
45
- TargetTCP string
46
-}
47
-
48
-type RelayClient struct {
49
- h host.Host
50
- cfg ClientConfig
51
- protoID protocol.ID
52
-
53
- ps *pubsub.PubSub
54
- t *pubsub.Topic
55
- wg sync.WaitGroup
56
- stop context.CancelFunc
57
-}
58
-
59
-// NewClient wires a reusable backend node that other apps can embed.
60
-// It registers a stream handler and starts an advertiser loop.
61
-// Call Close() to stop.
62
-func NewClient(ctx context.Context, h host.Host, cfg ClientConfig) (*RelayClient, error) {
63
- if cfg.AdvertiseEvery <= 0 {
64
- cfg.AdvertiseEvery = 5 * time.Second
65
- }
66
- if cfg.HTTPTimeout <= 0 {
67
- cfg.HTTPTimeout = 3 * time.Second
68
- }
69
- b := &RelayClient{
70
- h: h,
71
- cfg: cfg,
72
- protoID: protocol.ID(cfg.Protocol),
73
- }
74
-
75
- boot := make([]string, 0, len(cfg.Bootstraps)+4)
76
- if len(cfg.Bootstraps) > 0 {
77
- boot = append(boot, cfg.Bootstraps...)
78
- }
79
- if cfg.ServerURL != "" {
80
- if addrs, err := fetchMultiaddrsFromHealth(cfg.ServerURL, cfg.HTTPTimeout); err != nil {
81
- log.Warn().Err(err).Msgf("relaydns: fetch /health from %s failed", cfg.ServerURL)
82
- } else {
83
- sortMultiaddrs(addrs, cfg.PreferQUIC, cfg.PreferLocal)
84
- boot = append(boot, addrs...)
85
- }
86
- }
87
- boot = uniq(boot)
88
- if len(boot) > 0 {
89
- ConnectBootstraps(ctx, h, boot)
90
- } else {
91
- log.Warn().Msg("relaydns: no bootstrap sources provided (Bootstraps/ServerURL); discovery may fail")
92
- }
93
-
94
- // 1) stream handler
95
- switch {
96
- case cfg.Handler != nil:
97
- h.SetStreamHandler(b.protoID, cfg.Handler)
98
- case cfg.TargetTCP != "":
99
- h.SetStreamHandler(b.protoID, func(s network.Stream) {
100
- defer s.Close()
101
- c, err := net.Dial("tcp", cfg.TargetTCP)
102
- if err != nil {
103
- log.Error().Err(err).Msgf("relaydns: dial %s", cfg.TargetTCP)
104
- return
105
- }
106
- defer c.Close()
107
- // raw byte pipe
108
- go io.Copy(c, s)
109
- io.Copy(s, c)
110
- })
111
- default:
112
- return nil, fmt.Errorf("relaydns: either Handler or TargetTCP must be set")
113
- }
114
-
115
- // 2) pubsub join
116
- ps, err := pubsub.NewGossipSub(ctx, h, pubsub.WithMessageSigning(true))
117
- if err != nil {
118
- return nil, err
119
- }
120
- t, err := ps.Join(cfg.Topic)
121
- if err != nil {
122
- return nil, err
123
- }
124
- b.ps, b.t = ps, t
125
-
126
- // 3) advertiser loop
127
- advCtx, cancel := context.WithCancel(ctx)
128
- b.stop = cancel
129
- b.wg.Add(1)
130
- go func() {
131
- defer b.wg.Done()
132
- ticker := time.NewTicker(cfg.AdvertiseEvery)
133
- defer ticker.Stop()
134
- for {
135
- select {
136
- case <-advCtx.Done():
137
- return
138
- case <-ticker.C:
139
- addrs := b.h.Addrs()
140
- enc := make([]string, 0, len(addrs))
141
- for _, a := range addrs {
142
- enc = append(enc, fmt.Sprintf("%s/p2p/%s", a.String(), b.h.ID().String()))
143
- }
144
- ad := Advertise{
145
- Peer: b.h.ID().String(),
146
- Name: b.cfg.Name,
147
- DNS: b.cfg.DNS,
148
- Addrs: enc,
149
- Ready: true,
150
- Load: 0.0,
151
- TS: time.Now().UTC(),
152
- }
153
- payload, _ := json.Marshal(ad)
154
- _ = b.t.Publish(advCtx, payload)
155
- }
156
- }
157
- }()
158
-
159
- return b, nil
160
-}
161
-
162
-func (b *RelayClient) Close() error {
163
- if b.stop != nil {
164
- b.stop()
165
- }
166
- b.wg.Wait()
167
- // leaving topic is optional; libp2p will clean up on host close
168
- return nil
169
-}
170
-func fetchMultiaddrsFromHealth(base string, timeout time.Duration) ([]string, error) {
171
- u, err := url.Parse(base)
172
- if err != nil {
173
- return nil, fmt.Errorf("parse server-url: %w", err)
174
- }
175
- // ensure path ends with /health
176
- if !strings.HasSuffix(u.Path, "/health") {
177
- if u.Path == "" || u.Path == "/" {
178
- u.Path = "/health"
179
- } else {
180
- u.Path = strings.TrimSuffix(u.Path, "/") + "/health"
181
- }
182
- }
183
- client := &http.Client{Timeout: timeout}
184
- req, _ := http.NewRequest(http.MethodGet, u.String(), nil)
185
- resp, err := client.Do(req)
186
- if err != nil {
187
- return nil, err
188
- }
189
- defer resp.Body.Close()
190
-
191
- var payload struct {
192
- Status string `json:"status"`
193
- PeerID string `json:"peerId"`
194
- Multiaddrs []string `json:"multiaddrs"`
195
- }
196
- if err := json.NewDecoder(resp.Body).Decode(&payload); err != nil {
197
- return nil, err
198
- }
199
- if payload.Status != "ok" {
200
- return nil, errors.New("health not ok")
201
- }
202
- addrs := make([]string, 0, len(payload.Multiaddrs))
203
- for _, s := range payload.Multiaddrs {
204
- // 아주 기본적인 sanity check
205
- if strings.Contains(s, "/p2p/") && (strings.Contains(s, "/ip4/") || strings.Contains(s, "/ip6/")) {
206
- addrs = append(addrs, s)
207
- }
208
- }
209
- return addrs, nil
210
-}
211
-
212
-func sortMultiaddrs(addrs []string, preferQUIC, preferLocal bool) {
213
- score := func(a string) int {
214
- sc := 0
215
- if preferQUIC && strings.Contains(a, "/quic-v1") {
216
- sc += 2
217
- }
218
- if preferLocal && (strings.Contains(a, "/ip4/127.0.0.1/") || strings.Contains(a, "/ip6/::1/")) {
219
- sc += 1
220
- }
221
- return sc
222
- }
223
- sort.SliceStable(addrs, func(i, j int) bool { return score(addrs[i]) > score(addrs[j]) })
224
-}
225
-
226
-func uniq(ss []string) []string {
227
- seen := map[string]struct{}{}
228
- out := make([]string, 0, len(ss))
229
- for _, s := range ss {
230
- if _, ok := seen[s]; ok {
231
- continue
232
- }
233
- seen[s] = struct{}{}
234
- out = append(out, s)
235
- }
236
- return out
237
-}
1
+package relaydns
2
+
3
+import (
4
+ "context"
5
+ "encoding/json"
6
+ "errors"
7
+ "fmt"
8
+ "io"
9
+ "net"
10
+ "net/http"
11
+ "net/url"
12
+ "sort"
13
+ "strings"
14
+ "sync"
15
+ "time"
16
+
17
+ pubsub "github.com/libp2p/go-libp2p-pubsub"
18
+ "github.com/libp2p/go-libp2p/core/host"
19
+ "github.com/libp2p/go-libp2p/core/network"
20
+ "github.com/libp2p/go-libp2p/core/protocol"
21
+ "github.com/rs/zerolog/log"
22
+)
23
+
24
+type ClientConfig struct {
25
+ ServerURL string
26
+ Bootstraps []string
27
+ HTTPTimeout time.Duration
28
+ PreferQUIC bool
29
+ PreferLocal bool
30
+
31
+ // libp2p stream protocol id (e.g. "/relaydns/ssh/1.0")
32
+ Protocol string
33
+ // pubsub topic for backend adverts (e.g. "relaydns.backends")
34
+ Topic string
35
+ // advertise interval
36
+ AdvertiseEvery time.Duration
37
+ // optional metadata
38
+ Name string
39
+ DNS string
40
+
41
+ // One of the following:
42
+ // 1) Provide a custom stream handler
43
+ Handler func(s network.Stream)
44
+ // 2) Or just set TargetTCP to auto-pipe bytes to a local TCP service (e.g. "127.0.0.1:22")
45
+ TargetTCP string
46
+}
47
+
48
+type RelayClient struct {
49
+ h host.Host
50
+ cfg ClientConfig
51
+ protoID protocol.ID
52
+
53
+ ps *pubsub.PubSub
54
+ t *pubsub.Topic
55
+ wg sync.WaitGroup
56
+ stop context.CancelFunc
57
+}
58
+
59
+// NewClient wires a reusable backend node that other apps can embed.
60
+// It registers a stream handler and starts an advertiser loop.
61
+// Call Close() to stop.
62
+func NewClient(ctx context.Context, h host.Host, cfg ClientConfig) (*RelayClient, error) {
63
+ if cfg.AdvertiseEvery <= 0 {
64
+ cfg.AdvertiseEvery = 5 * time.Second
65
+ }
66
+ if cfg.HTTPTimeout <= 0 {
67
+ cfg.HTTPTimeout = 3 * time.Second
68
+ }
69
+ b := &RelayClient{
70
+ h: h,
71
+ cfg: cfg,
72
+ protoID: protocol.ID(cfg.Protocol),
73
+ }
74
+
75
+ boot := make([]string, 0, len(cfg.Bootstraps)+4)
76
+ if len(cfg.Bootstraps) > 0 {
77
+ boot = append(boot, cfg.Bootstraps...)
78
+ }
79
+ if cfg.ServerURL != "" {
80
+ if addrs, err := fetchMultiaddrsFromHealth(cfg.ServerURL, cfg.HTTPTimeout); err != nil {
81
+ log.Warn().Err(err).Msgf("relaydns: fetch /health from %s failed", cfg.ServerURL)
82
+ } else {
83
+ sortMultiaddrs(addrs, cfg.PreferQUIC, cfg.PreferLocal)
84
+ boot = append(boot, addrs...)
85
+ }
86
+ }
87
+ boot = uniq(boot)
88
+ if len(boot) > 0 {
89
+ ConnectBootstraps(ctx, h, boot)
90
+ } else {
91
+ log.Warn().Msg("relaydns: no bootstrap sources provided (Bootstraps/ServerURL); discovery may fail")
92
+ }
93
+
94
+ // 1) stream handler
95
+ switch {
96
+ case cfg.Handler != nil:
97
+ h.SetStreamHandler(b.protoID, cfg.Handler)
98
+ case cfg.TargetTCP != "":
99
+ h.SetStreamHandler(b.protoID, func(s network.Stream) {
100
+ defer s.Close()
101
+ c, err := net.Dial("tcp", cfg.TargetTCP)
102
+ if err != nil {
103
+ log.Error().Err(err).Msgf("relaydns: dial %s", cfg.TargetTCP)
104
+ return
105
+ }
106
+ defer c.Close()
107
+ // raw byte pipe
108
+ go io.Copy(c, s)
109
+ io.Copy(s, c)
110
+ })
111
+ default:
112
+ return nil, fmt.Errorf("relaydns: either Handler or TargetTCP must be set")
113
+ }
114
+
115
+ // 2) pubsub join
116
+ ps, err := pubsub.NewGossipSub(ctx, h, pubsub.WithMessageSigning(true))
117
+ if err != nil {
118
+ return nil, err
119
+ }
120
+ t, err := ps.Join(cfg.Topic)
121
+ if err != nil {
122
+ return nil, err
123
+ }
124
+ b.ps, b.t = ps, t
125
+
126
+ // 3) advertiser loop
127
+ advCtx, cancel := context.WithCancel(ctx)
128
+ b.stop = cancel
129
+ b.wg.Add(1)
130
+ go func() {
131
+ defer b.wg.Done()
132
+ ticker := time.NewTicker(cfg.AdvertiseEvery)
133
+ defer ticker.Stop()
134
+ for {
135
+ select {
136
+ case <-advCtx.Done():
137
+ return
138
+ case <-ticker.C:
139
+ addrs := b.h.Addrs()
140
+ enc := make([]string, 0, len(addrs))
141
+ for _, a := range addrs {
142
+ enc = append(enc, fmt.Sprintf("%s/p2p/%s", a.String(), b.h.ID().String()))
143
+ }
144
+ ad := Advertise{
145
+ Peer: b.h.ID().String(),
146
+ Name: b.cfg.Name,
147
+ DNS: b.cfg.DNS,
148
+ Addrs: enc,
149
+ Ready: true,
150
+ Load: 0.0,
151
+ TS: time.Now().UTC(),
152
+ }
153
+ payload, _ := json.Marshal(ad)
154
+ _ = b.t.Publish(advCtx, payload)
155
+ }
156
+ }
157
+ }()
158
+
159
+ return b, nil
160
+}
161
+
162
+func (b *RelayClient) Close() error {
163
+ if b.stop != nil {
164
+ b.stop()
165
+ }
166
+ b.wg.Wait()
167
+ // leaving topic is optional; libp2p will clean up on host close
168
+ return nil
169
+}
170
+func fetchMultiaddrsFromHealth(base string, timeout time.Duration) ([]string, error) {
171
+ u, err := url.Parse(base)
172
+ if err != nil {
173
+ return nil, fmt.Errorf("parse server-url: %w", err)
174
+ }
175
+ // ensure path ends with /health
176
+ if !strings.HasSuffix(u.Path, "/health") {
177
+ if u.Path == "" || u.Path == "/" {
178
+ u.Path = "/health"
179
+ } else {
180
+ u.Path = strings.TrimSuffix(u.Path, "/") + "/health"
181
+ }
182
+ }
183
+ client := &http.Client{Timeout: timeout}
184
+ req, _ := http.NewRequest(http.MethodGet, u.String(), nil)
185
+ resp, err := client.Do(req)
186
+ if err != nil {
187
+ return nil, err
188
+ }
189
+ defer resp.Body.Close()
190
+
191
+ var payload struct {
192
+ Status string `json:"status"`
193
+ PeerID string `json:"peerId"`
194
+ Multiaddrs []string `json:"multiaddrs"`
195
+ }
196
+ if err := json.NewDecoder(resp.Body).Decode(&payload); err != nil {
197
+ return nil, err
198
+ }
199
+ if payload.Status != "ok" {
200
+ return nil, errors.New("health not ok")
201
+ }
202
+ addrs := make([]string, 0, len(payload.Multiaddrs))
203
+ for _, s := range payload.Multiaddrs {
204
+ // 아주 기본적인 sanity check
205
+ if strings.Contains(s, "/p2p/") && (strings.Contains(s, "/ip4/") || strings.Contains(s, "/ip6/")) {
206
+ addrs = append(addrs, s)
207
+ }
208
+ }
209
+ return addrs, nil
210
+}
211
+
212
+func sortMultiaddrs(addrs []string, preferQUIC, preferLocal bool) {
213
+ score := func(a string) int {
214
+ sc := 0
215
+ if preferQUIC && strings.Contains(a, "/quic-v1") {
216
+ sc += 2
217
+ }
218
+ if preferLocal && (strings.Contains(a, "/ip4/127.0.0.1/") || strings.Contains(a, "/ip6/::1/")) {
219
+ sc += 1
220
+ }
221
+ return sc
222
+ }
223
+ sort.SliceStable(addrs, func(i, j int) bool { return score(addrs[i]) > score(addrs[j]) })
224
+}
225
+
226
+func uniq(ss []string) []string {
227
+ seen := map[string]struct{}{}
228
+ out := make([]string, 0, len(ss))
229
+ for _, s := range ss {
230
+ if _, ok := seen[s]; ok {
231
+ continue
232
+ }
233
+ seen[s] = struct{}{}
234
+ out = append(out, s)
235
+ }
236
+ return out
237
+}
relaydns/director.go
+215
-215
@@ -1,215 +1,215 @@
1
-package relaydns
2
-
3
-import (
4
- "context"
5
- "encoding/json"
6
- "fmt"
7
- "io"
8
- "net"
9
- "net/http"
10
- "sync"
11
- "time"
12
-
13
- pubsub "github.com/libp2p/go-libp2p-pubsub"
14
- "github.com/libp2p/go-libp2p/core/host"
15
- "github.com/libp2p/go-libp2p/core/peer"
16
- ma "github.com/multiformats/go-multiaddr"
17
- "github.com/rs/zerolog/log"
18
-)
19
-
20
-type Director struct {
21
- ctx context.Context
22
- h host.Host
23
- protocol string
24
- topicName string
25
- sub *pubsub.Subscription
26
-
27
- storeMu sync.Mutex
28
- store map[string]HostEntry
29
- ttl time.Duration
30
- pick *Picker
31
-}
32
-
33
-func NewDirector(ctx context.Context, h host.Host, protocol, topic string) (*Director, error) {
34
- ps, err := pubsub.NewGossipSub(ctx, h)
35
- if err != nil {
36
- return nil, err
37
- }
38
- t, err := ps.Join(topic)
39
- if err != nil {
40
- return nil, err
41
- }
42
- sub, err := t.Subscribe()
43
- if err != nil {
44
- return nil, err
45
- }
46
- d := &Director{
47
- ctx: ctx,
48
- h: h,
49
- protocol: protocol,
50
- topicName: topic,
51
- sub: sub,
52
- store: map[string]HostEntry{},
53
- ttl: 45 * time.Second,
54
- pick: &Picker{},
55
- }
56
- go d.collect()
57
- go d.gc()
58
- return d, nil
59
-}
60
-
61
-func (d *Director) Close() error {
62
- d.sub.Cancel()
63
- return nil
64
-}
65
-
66
-func (d *Director) collect() {
67
- for {
68
- msg, err := d.sub.Next(d.ctx)
69
- if err != nil {
70
- return
71
- }
72
- var ad Advertise
73
- if err := json.Unmarshal(msg.Data, &ad); err != nil {
74
- continue
75
- }
76
- var ai *peer.AddrInfo
77
- // pick any addr that includes /p2p/ (dnsaddr resolved entries will)
78
- for _, s := range ad.Addrs {
79
- m, err := ma.NewMultiaddr(s)
80
- if err != nil {
81
- continue
82
- }
83
- if a, err := peer.AddrInfoFromP2pAddr(m); err == nil {
84
- ai = a
85
- break
86
- }
87
- }
88
- if ai == nil {
89
- continue
90
- }
91
- d.storeMu.Lock()
92
- d.store[ad.Peer] = HostEntry{Info: ad, AddrInfo: ai, LastSeen: time.Now()}
93
- // refresh picker snapshot
94
- snap := make([]HostEntry, 0, len(d.store))
95
- for _, v := range d.store {
96
- snap = append(snap, v)
97
- }
98
- d.storeMu.Unlock()
99
- d.pick.update(snap)
100
- }
101
-}
102
-
103
-func (d *Director) gc() {
104
- t := time.NewTicker(5 * time.Second)
105
- defer t.Stop()
106
- for {
107
- select {
108
- case <-d.ctx.Done():
109
- return
110
- case <-t.C:
111
- now := time.Now()
112
- d.storeMu.Lock()
113
- for k, v := range d.store {
114
- if now.Sub(v.LastSeen) > d.ttl {
115
- delete(d.store, k)
116
- }
117
- }
118
- snap := make([]HostEntry, 0, len(d.store))
119
- for _, v := range d.store {
120
- snap = append(snap, v)
121
- }
122
- d.storeMu.Unlock()
123
- d.pick.update(snap)
124
- }
125
- }
126
-}
127
-
128
-func (d *Director) ServeTCP(addr string) error {
129
- ln, err := net.Listen("tcp", addr)
130
- if err != nil {
131
- return err
132
- }
133
-
134
- log.Info().Msgf("director TCP listening on %s (protocol %s)", addr, d.protocol)
135
- for {
136
- c, err := ln.Accept()
137
- if err != nil {
138
- continue
139
- }
140
- go d.handleConn(c)
141
- }
142
-}
143
-
144
-func (d *Director) handleConn(c net.Conn) {
145
- defer c.Close()
146
- entry, ok := d.pick.choose()
147
- if !ok {
148
- log.Warn().Msg("no backend peers available")
149
- return
150
- }
151
- // ensure we're connected (AddrInfo contains addrs)
152
- if err := d.h.Connect(d.ctx, *entry.AddrInfo); err != nil {
153
- log.Error().Err(err).Msgf("connect %s failed", entry.AddrInfo.ID)
154
- return
155
- }
156
- s, err := d.h.NewStream(d.ctx, entry.AddrInfo.ID, protocolID(d.protocol))
157
- if err != nil {
158
- log.Error().Err(err).Msg("new stream")
159
- return
160
- }
161
- defer s.Close()
162
- // raw byte tunnel
163
- go io.Copy(s, c)
164
- io.Copy(c, s)
165
-}
166
-
167
-func (d *Director) ServeHTTP(addr string) error {
168
- mux := http.NewServeMux()
169
- mux.HandleFunc("/hosts", func(w http.ResponseWriter, r *http.Request) {
170
- d.storeMu.Lock()
171
- defer d.storeMu.Unlock()
172
- list := make([]HostEntry, 0, len(d.store))
173
- for _, v := range d.store {
174
- list = append(list, v)
175
- }
176
- _ = json.NewEncoder(w).Encode(list)
177
- })
178
- mux.HandleFunc("/override", func(w http.ResponseWriter, r *http.Request) {
179
- switch r.Method {
180
- case "POST":
181
- peerID := r.URL.Query().Get("peer")
182
- dur := 30 * time.Second
183
- if s := r.URL.Query().Get("ttl"); s != "" {
184
- if v, err := time.ParseDuration(s); err == nil {
185
- dur = v
186
- }
187
- }
188
- d.pick.pin(peerID, dur)
189
- w.WriteHeader(204)
190
- case "DELETE":
191
- d.pick.unpin()
192
- w.WriteHeader(204)
193
- default:
194
- w.WriteHeader(405)
195
- }
196
- })
197
- mux.HandleFunc("/health", func(w http.ResponseWriter, r *http.Request) {
198
- type info struct {
199
- Status string `json:"status"`
200
- Addrs []string `json:"multiaddrs"`
201
- }
202
- var list []string = make([]string, 0)
203
- for _, a := range d.h.Addrs() {
204
- list = append(list, fmt.Sprintf("%s/p2p/%s", a.String(), d.h.ID().String()))
205
- }
206
- resp := info{
207
- Status: "ok",
208
- Addrs: list,
209
- }
210
- w.Header().Set("Content-Type", "application/json")
211
- _ = json.NewEncoder(w).Encode(resp)
212
- })
213
- log.Info().Msgf("director HTTP API on %s", addr)
214
- return http.ListenAndServe(addr, mux)
215
-}
1
+package relaydns
2
+
3
+import (
4
+ "context"
5
+ "encoding/json"
6
+ "fmt"
7
+ "io"
8
+ "net"
9
+ "net/http"
10
+ "sync"
11
+ "time"
12
+
13
+ pubsub "github.com/libp2p/go-libp2p-pubsub"
14
+ "github.com/libp2p/go-libp2p/core/host"
15
+ "github.com/libp2p/go-libp2p/core/peer"
16
+ ma "github.com/multiformats/go-multiaddr"
17
+ "github.com/rs/zerolog/log"
18
+)
19
+
20
+type Director struct {
21
+ ctx context.Context
22
+ h host.Host
23
+ protocol string
24
+ topicName string
25
+ sub *pubsub.Subscription
26
+
27
+ storeMu sync.Mutex
28
+ store map[string]HostEntry
29
+ ttl time.Duration
30
+ pick *Picker
31
+}
32
+
33
+func NewDirector(ctx context.Context, h host.Host, protocol, topic string) (*Director, error) {
34
+ ps, err := pubsub.NewGossipSub(ctx, h)
35
+ if err != nil {
36
+ return nil, err
37
+ }
38
+ t, err := ps.Join(topic)
39
+ if err != nil {
40
+ return nil, err
41
+ }
42
+ sub, err := t.Subscribe()
43
+ if err != nil {
44
+ return nil, err
45
+ }
46
+ d := &Director{
47
+ ctx: ctx,
48
+ h: h,
49
+ protocol: protocol,
50
+ topicName: topic,
51
+ sub: sub,
52
+ store: map[string]HostEntry{},
53
+ ttl: 45 * time.Second,
54
+ pick: &Picker{},
55
+ }
56
+ go d.collect()
57
+ go d.gc()
58
+ return d, nil
59
+}
60
+
61
+func (d *Director) Close() error {
62
+ d.sub.Cancel()
63
+ return nil
64
+}
65
+
66
+func (d *Director) collect() {
67
+ for {
68
+ msg, err := d.sub.Next(d.ctx)
69
+ if err != nil {
70
+ return
71
+ }
72
+ var ad Advertise
73
+ if err := json.Unmarshal(msg.Data, &ad); err != nil {
74
+ continue
75
+ }
76
+ var ai *peer.AddrInfo
77
+ // pick any addr that includes /p2p/ (dnsaddr resolved entries will)
78
+ for _, s := range ad.Addrs {
79
+ m, err := ma.NewMultiaddr(s)
80
+ if err != nil {
81
+ continue
82
+ }
83
+ if a, err := peer.AddrInfoFromP2pAddr(m); err == nil {
84
+ ai = a
85
+ break
86
+ }
87
+ }
88
+ if ai == nil {
89
+ continue
90
+ }
91
+ d.storeMu.Lock()
92
+ d.store[ad.Peer] = HostEntry{Info: ad, AddrInfo: ai, LastSeen: time.Now()}
93
+ // refresh picker snapshot
94
+ snap := make([]HostEntry, 0, len(d.store))
95
+ for _, v := range d.store {
96
+ snap = append(snap, v)
97
+ }
98
+ d.storeMu.Unlock()
99
+ d.pick.update(snap)
100
+ }
101
+}
102
+
103
+func (d *Director) gc() {
104
+ t := time.NewTicker(5 * time.Second)
105
+ defer t.Stop()
106
+ for {
107
+ select {
108
+ case <-d.ctx.Done():
109
+ return
110
+ case <-t.C:
111
+ now := time.Now()
112
+ d.storeMu.Lock()
113
+ for k, v := range d.store {
114
+ if now.Sub(v.LastSeen) > d.ttl {
115
+ delete(d.store, k)
116
+ }
117
+ }
118
+ snap := make([]HostEntry, 0, len(d.store))
119
+ for _, v := range d.store {
120
+ snap = append(snap, v)
121
+ }
122
+ d.storeMu.Unlock()
123
+ d.pick.update(snap)
124
+ }
125
+ }
126
+}
127
+
128
+func (d *Director) ServeTCP(addr string) error {
129
+ ln, err := net.Listen("tcp", addr)
130
+ if err != nil {
131
+ return err
132
+ }
133
+
134
+ log.Info().Msgf("director TCP listening on %s (protocol %s)", addr, d.protocol)
135
+ for {
136
+ c, err := ln.Accept()
137
+ if err != nil {
138
+ continue
139
+ }
140
+ go d.handleConn(c)
141
+ }
142
+}
143
+
144
+func (d *Director) handleConn(c net.Conn) {
145
+ defer c.Close()
146
+ entry, ok := d.pick.choose()
147
+ if !ok {
148
+ log.Warn().Msg("no backend peers available")
149
+ return
150
+ }
151
+ // ensure we're connected (AddrInfo contains addrs)
152
+ if err := d.h.Connect(d.ctx, *entry.AddrInfo); err != nil {
153
+ log.Error().Err(err).Msgf("connect %s failed", entry.AddrInfo.ID)
154
+ return
155
+ }
156
+ s, err := d.h.NewStream(d.ctx, entry.AddrInfo.ID, protocolID(d.protocol))
157
+ if err != nil {
158
+ log.Error().Err(err).Msg("new stream")
159
+ return
160
+ }
161
+ defer s.Close()
162
+ // raw byte tunnel
163
+ go io.Copy(s, c)
164
+ io.Copy(c, s)
165
+}
166
+
167
+func (d *Director) ServeHTTP(addr string) error {
168
+ mux := http.NewServeMux()
169
+ mux.HandleFunc("/hosts", func(w http.ResponseWriter, r *http.Request) {
170
+ d.storeMu.Lock()
171
+ defer d.storeMu.Unlock()
172
+ list := make([]HostEntry, 0, len(d.store))
173
+ for _, v := range d.store {
174
+ list = append(list, v)
175
+ }
176
+ _ = json.NewEncoder(w).Encode(list)
177
+ })
178
+ mux.HandleFunc("/override", func(w http.ResponseWriter, r *http.Request) {
179
+ switch r.Method {
180
+ case "POST":
181
+ peerID := r.URL.Query().Get("peer")
182
+ dur := 30 * time.Second
183
+ if s := r.URL.Query().Get("ttl"); s != "" {
184
+ if v, err := time.ParseDuration(s); err == nil {
185
+ dur = v
186
+ }
187
+ }
188
+ d.pick.pin(peerID, dur)
189
+ w.WriteHeader(204)
190
+ case "DELETE":
191
+ d.pick.unpin()
192
+ w.WriteHeader(204)
193
+ default:
194
+ w.WriteHeader(405)
195
+ }
196
+ })
197
+ mux.HandleFunc("/health", func(w http.ResponseWriter, r *http.Request) {
198
+ type info struct {
199
+ Status string `json:"status"`
200
+ Addrs []string `json:"multiaddrs"`
201
+ }
202
+ var list []string = make([]string, 0)
203
+ for _, a := range d.h.Addrs() {
204
+ list = append(list, fmt.Sprintf("%s/p2p/%s", a.String(), d.h.ID().String()))
205
+ }
206
+ resp := info{
207
+ Status: "ok",
208
+ Addrs: list,
209
+ }
210
+ w.Header().Set("Content-Type", "application/json")
211
+ _ = json.NewEncoder(w).Encode(resp)
212
+ })
213
+ log.Info().Msgf("director HTTP API on %s", addr)
214
+ return http.ListenAndServe(addr, mux)
215
+}
relaydns/host.go
+55
-55
@@ -1,59 +1,59 @@
1
-package relaydns
2
-
1
+package relaydns
2
+
3
import (
4
- "context"
5
- "fmt"
4
+ "context"
5
+ "fmt"
6
7
- "github.com/libp2p/go-libp2p"
8
- "github.com/libp2p/go-libp2p/core/host"
9
- "github.com/libp2p/go-libp2p/core/peer"
10
- ma "github.com/multiformats/go-multiaddr"
11
- "github.com/rs/zerolog/log"
7
+ "github.com/libp2p/go-libp2p"
8
+ "github.com/libp2p/go-libp2p/core/host"
9
+ "github.com/libp2p/go-libp2p/core/peer"
10
+ ma "github.com/multiformats/go-multiaddr"
11
+ "github.com/rs/zerolog/log"
12
)
13
-
14
-func MakeHost(ctx context.Context, port int, enableRelay bool) (host.Host, error) {
15
- addrs := []string{
16
- fmt.Sprintf("/ip4/0.0.0.0/tcp/%d", port),
17
- fmt.Sprintf("/ip4/0.0.0.0/udp/%d/quic-v1", port),
18
- fmt.Sprintf("/ip6/::/tcp/%d", port),
19
- fmt.Sprintf("/ip6/::/udp/%d/quic-v1", port),
20
- }
21
-
22
- opts := []libp2p.Option{
23
- libp2p.ListenAddrStrings(addrs...),
24
- libp2p.DefaultTransports, // TCP+QUIC
25
- libp2p.NATPortMap(),
26
- libp2p.EnableNATService(), // AutoNAT helper
27
- libp2p.EnableHolePunching(), // DCUtR
28
- libp2p.DefaultSecurity,
29
- libp2p.DefaultMuxers,
30
- }
31
- if enableRelay {
32
- opts = append(opts, libp2p.EnableRelay()) // circuit relay (useful both as client & svc)
33
- }
34
- h, err := libp2p.New(opts...)
35
- if err != nil {
36
- return nil, err
37
- }
38
- return h, nil
39
-}
40
-
41
-func ConnectBootstraps(ctx context.Context, h host.Host, addrs []string) {
42
- for _, s := range addrs {
43
- m, err := ma.NewMultiaddr(s)
44
- if err != nil {
45
- log.Warn().Err(err).Msgf("bootstrap bad multiaddr %q", s)
46
- continue
47
- }
48
- ai, err := peer.AddrInfoFromP2pAddr(m)
49
- if err != nil {
50
- log.Warn().Err(err).Msgf("bootstrap missing /p2p/ in %q", s)
51
- continue
52
- }
53
- if err := h.Connect(ctx, *ai); err != nil {
54
- log.Warn().Err(err).Msgf("bootstrap connect %s", ai.ID)
55
- } else {
56
- log.Info().Msgf("connected bootstrap %s", ai.ID)
57
- }
58
- }
13
+
14
+func MakeHost(ctx context.Context, port int, enableRelay bool) (host.Host, error) {
15
+ addrs := []string{
16
+ fmt.Sprintf("/ip4/0.0.0.0/tcp/%d", port),
17
+ fmt.Sprintf("/ip4/0.0.0.0/udp/%d/quic-v1", port),
18
+ fmt.Sprintf("/ip6/::/tcp/%d", port),
19
+ fmt.Sprintf("/ip6/::/udp/%d/quic-v1", port),
20
+ }
21
+
22
+ opts := []libp2p.Option{
23
+ libp2p.ListenAddrStrings(addrs...),
24
+ libp2p.DefaultTransports, // TCP+QUIC
25
+ libp2p.NATPortMap(),
26
+ libp2p.EnableNATService(), // AutoNAT helper
27
+ libp2p.EnableHolePunching(), // DCUtR
28
+ libp2p.DefaultSecurity,
29
+ libp2p.DefaultMuxers,
30
+ }
31
+ if enableRelay {
32
+ opts = append(opts, libp2p.EnableRelay()) // circuit relay (useful both as client & svc)
33
+ }
34
+ h, err := libp2p.New(opts...)
35
+ if err != nil {
36
+ return nil, err
37
+ }
38
+ return h, nil
39
+}
40
+
41
+func ConnectBootstraps(ctx context.Context, h host.Host, addrs []string) {
42
+ for _, s := range addrs {
43
+ m, err := ma.NewMultiaddr(s)
44
+ if err != nil {
45
+ log.Warn().Err(err).Msgf("bootstrap bad multiaddr %q", s)
46
+ continue
47
+ }
48
+ ai, err := peer.AddrInfoFromP2pAddr(m)
49
+ if err != nil {
50
+ log.Warn().Err(err).Msgf("bootstrap missing /p2p/ in %q", s)
51
+ continue
52
+ }
53
+ if err := h.Connect(ctx, *ai); err != nil {
54
+ log.Warn().Err(err).Msgf("bootstrap connect %s", ai.ID)
55
+ } else {
56
+ log.Info().Msgf("connected bootstrap %s", ai.ID)
57
+ }
58
+ }
59
}
relaydns/types.go
+72
-72
@@ -1,72 +1,72 @@
1
-package relaydns
2
-
3
-import (
4
- "sync"
5
- "sync/atomic"
6
- "time"
7
-
8
- "github.com/libp2p/go-libp2p/core/peer"
9
- "github.com/libp2p/go-libp2p/core/protocol"
10
-)
11
-
12
-type Advertise struct {
13
- Peer string `json:"peer"`
14
- Name string `json:"name,omitempty"`
15
- DNS string `json:"dns,omitempty"`
16
- Addrs []string `json:"addrs"`
17
- Ready bool `json:"ready"`
18
- Load float64 `json:"load"`
19
- TS time.Time `json:"ts"`
20
-}
21
-
22
-type HostEntry struct {
23
- Info Advertise
24
- AddrInfo *peer.AddrInfo
25
- LastSeen time.Time
26
-}
27
-
28
-type Picker struct {
29
- mu sync.RWMutex
30
- rr uint64
31
- list []HostEntry
32
- pinTo string
33
- pinTil time.Time
34
-}
35
-
36
-func (p *Picker) update(list []HostEntry) {
37
- p.mu.Lock()
38
- p.list = list
39
- p.mu.Unlock()
40
-}
41
-func (p *Picker) choose() (HostEntry, bool) {
42
- p.mu.RLock()
43
- defer p.mu.RUnlock()
44
- if len(p.list) == 0 {
45
- return HostEntry{}, false
46
- }
47
- if p.pinTo != "" && time.Now().Before(p.pinTil) {
48
- for _, e := range p.list {
49
- if e.Info.Peer == p.pinTo {
50
- return e, true
51
- }
52
- }
53
- }
54
- i := atomic.AddUint64(&p.rr, 1)
55
- return p.list[i%uint64(len(p.list))], true
56
-}
57
-func (p *Picker) pin(peerID string, dur time.Duration) {
58
- p.mu.Lock()
59
- p.pinTo = peerID
60
- p.pinTil = time.Now().Add(dur)
61
- p.mu.Unlock()
62
-}
63
-func (p *Picker) unpin() {
64
- p.mu.Lock()
65
- p.pinTo = ""
66
- p.pinTil = time.Time{}
67
- p.mu.Unlock()
68
-}
69
-
70
-func protocolID(s string) protocol.ID {
71
- return protocol.ID(s)
72
-}
1
+package relaydns
2
+
3
+import (
4
+ "sync"
5
+ "sync/atomic"
6
+ "time"
7
+
8
+ "github.com/libp2p/go-libp2p/core/peer"
9
+ "github.com/libp2p/go-libp2p/core/protocol"
10
+)
11
+
12
+type Advertise struct {
13
+ Peer string `json:"peer"`
14
+ Name string `json:"name,omitempty"`
15
+ DNS string `json:"dns,omitempty"`
16
+ Addrs []string `json:"addrs"`
17
+ Ready bool `json:"ready"`
18
+ Load float64 `json:"load"`
19
+ TS time.Time `json:"ts"`
20
+}
21
+
22
+type HostEntry struct {
23
+ Info Advertise
24
+ AddrInfo *peer.AddrInfo
25
+ LastSeen time.Time
26
+}
27
+
28
+type Picker struct {
29
+ mu sync.RWMutex
30
+ rr uint64
31
+ list []HostEntry
32
+ pinTo string
33
+ pinTil time.Time
34
+}
35
+
36
+func (p *Picker) update(list []HostEntry) {
37
+ p.mu.Lock()
38
+ p.list = list
39
+ p.mu.Unlock()
40
+}
41
+func (p *Picker) choose() (HostEntry, bool) {
42
+ p.mu.RLock()
43
+ defer p.mu.RUnlock()
44
+ if len(p.list) == 0 {
45
+ return HostEntry{}, false
46
+ }
47
+ if p.pinTo != "" && time.Now().Before(p.pinTil) {
48
+ for _, e := range p.list {
49
+ if e.Info.Peer == p.pinTo {
50
+ return e, true
51
+ }
52
+ }
53
+ }
54
+ i := atomic.AddUint64(&p.rr, 1)
55
+ return p.list[i%uint64(len(p.list))], true
56
+}
57
+func (p *Picker) pin(peerID string, dur time.Duration) {
58
+ p.mu.Lock()
59
+ p.pinTo = peerID
60
+ p.pinTil = time.Now().Add(dur)
61
+ p.mu.Unlock()
62
+}
63
+func (p *Picker) unpin() {
64
+ p.mu.Lock()
65
+ p.pinTo = ""
66
+ p.pinTil = time.Time{}
67
+ p.mu.Unlock()
68
+}
69
+
70
+func protocolID(s string) protocol.ID {
71
+ return protocol.ID(s)
72
+}