feat: add a stats dht command
Currently, it just prints out the routing tables.
Steven Allen committed
Apr 26, 2020 at 19:30 UTC
5751917f4979baa3faab193f3579d4f1e7383b6f
3 files changed
+196
core/commands/commands_test.go
+1
@@ -195,6 +195,7 @@ func TestCommands(t *testing.T) {
195
"/stats",
196
"/stats/bitswap",
197
"/stats/bw",
198
+ "/stats/dht",
199
"/stats/repo",
200
"/swarm",
201
"/swarm/addrs",
core/commands/stat.go
+1
@@ -29,6 +29,7 @@ for your IPFS node.`,
29
"bw": statBwCmd,
30
"repo": repoStatCmd,
31
"bitswap": bitswapStatCmd,
32
+ "dht": statDhtCmd,
33
},
34
}
35
core/commands/stat_dht.go
new
+194
@@ -0,0 +1,194 @@
1
+package commands
2
+
3
+import (
4
+ "fmt"
5
+ "io"
6
+ "text/tabwriter"
7
+ "time"
8
+
9
+ cmdenv "github.com/ipfs/go-ipfs/core/commands/cmdenv"
10
+
11
+ cmds "github.com/ipfs/go-ipfs-cmds"
12
+ "github.com/libp2p/go-libp2p-core/network"
13
+ pstore "github.com/libp2p/go-libp2p-core/peerstore"
14
+ dht "github.com/libp2p/go-libp2p-kad-dht"
15
+ kbucket "github.com/libp2p/go-libp2p-kbucket"
16
+)
17
+
18
+type dhtPeerInfo struct {
19
+ ID string
20
+ Connected bool
21
+ AgentVersion string
22
+ LastUsefulAt string
23
+ LastQueriedAt string
24
+}
25
+
26
+type dhtStat struct {
27
+ Name string
28
+ Buckets []dhtBucket
29
+}
30
+
31
+type dhtBucket struct {
32
+ LastRefresh string
33
+ Peers []dhtPeerInfo
34
+}
35
+
36
+var statDhtCmd = &cmds.Command{
37
+ Helptext: cmds.HelpText{
38
+ Tagline: "Returns statistics about the node's DHT(s)",
39
+ ShortDescription: `
40
+Returns statistics about the DHT(s) the node is participating in.
41
+
42
+This interface is not stable and may change from release to release.
43
+`,
44
+ },
45
+ Arguments: []cmds.Argument{
46
+ cmds.StringArg("dht", false, true, "The DHT whose table should be listed (wan or lan). Defaults to both."),
47
+ },
48
+ Options: []cmds.Option{},
49
+ Run: func(req *cmds.Request, res cmds.ResponseEmitter, env cmds.Environment) error {
50
+ nd, err := cmdenv.GetNode(env)
51
+ if err != nil {
52
+ return err
53
+ }
54
+
55
+ if !nd.IsOnline {
56
+ return ErrNotOnline
57
+ }
58
+
59
+ if nd.DHT == nil {
60
+ return ErrNotDHT
61
+ }
62
+
63
+ id := kbucket.ConvertPeerID(nd.Identity)
64
+
65
+ dhts := req.Arguments
66
+ if len(dhts) == 0 {
67
+ dhts = []string{"wan", "lan"}
68
+ }
69
+
70
+ for _, name := range dhts {
71
+ var dht *dht.IpfsDHT
72
+ switch name {
73
+ case "wan":
74
+ dht = nd.DHT.WAN
75
+ case "lan":
76
+ dht = nd.DHT.LAN
77
+ default:
78
+ return cmds.Errorf(cmds.ErrClient, "unknown dht type: %s", name)
79
+ }
80
+
81
+ rt := dht.RoutingTable()
82
+ lastRefresh := rt.GetTrackedCplsForRefresh()
83
+ infos := rt.GetPeerInfos()
84
+ buckets := make([]dhtBucket, 0, len(lastRefresh))
85
+ for _, pi := range infos {
86
+ cpl := kbucket.CommonPrefixLen(id, kbucket.ConvertPeerID(pi.Id))
87
+ if len(buckets) <= cpl {
88
+ buckets = append(buckets, make([]dhtBucket, 1+cpl-len(buckets))...)
89
+ }
90
+
91
+ info := dhtPeerInfo{ID: pi.Id.String()}
92
+
93
+ if ver, err := nd.Peerstore.Get(pi.Id, "AgentVersion"); err == nil {
94
+ info.AgentVersion, _ = ver.(string)
95
+ } else if err == pstore.ErrNotFound {
96
+ // ignore
97
+ } else {
98
+ // this is a bug, usually.
99
+ log.Errorw(
100
+ "failed to get agent version from peerstore",
101
+ "error", err,
102
+ )
103
+ }
104
+ if !pi.LastUsefulAt.IsZero() {
105
+ info.LastUsefulAt = pi.LastUsefulAt.Format(time.RFC3339)
106
+ }
107
+
108
+ if !pi.LastSuccessfulOutboundQueryAt.IsZero() {
109
+ info.LastQueriedAt = pi.LastSuccessfulOutboundQueryAt.Format(time.RFC3339)
110
+ }
111
+
112
+ info.Connected = nd.PeerHost.Network().Connectedness(pi.Id) == network.Connected
113
+
114
+ buckets[cpl].Peers = append(buckets[cpl].Peers, info)
115
+ }
116
+ for i := 0; i < len(buckets) && i < len(lastRefresh); i++ {
117
+ refreshTime := lastRefresh[i]
118
+ if !refreshTime.IsZero() {
119
+ buckets[i].LastRefresh = refreshTime.Format(time.RFC3339)
120
+ }
121
+ }
122
+ if err := res.Emit(dhtStat{
123
+ Name: name,
124
+ Buckets: buckets,
125
+ }); err != nil {
126
+ return err
127
+ }
128
+ }
129
+
130
+ return nil
131
+ },
132
+ Encoders: cmds.EncoderMap{
133
+ cmds.Text: cmds.MakeTypedEncoder(func(req *cmds.Request, w io.Writer, out dhtStat) error {
134
+ tw := tabwriter.NewWriter(w, 4, 4, 2, ' ', 0)
135
+ defer tw.Flush()
136
+
137
+ // Formats a time into XX ago and remove any decimal
138
+ // parts. That is, change "2m3.00010101s" to "2m3s ago".
139
+ now := time.Now()
140
+ since := func(t time.Time) string {
141
+ return now.Sub(t).Round(time.Second).String() + " ago"
142
+ }
143
+
144
+ count := 0
145
+ for _, bucket := range out.Buckets {
146
+ count += len(bucket.Peers)
147
+ }
148
+
149
+ fmt.Fprintf(tw, "DHT %s (%d peers):\t\t\t\n", out.Name, count)
150
+
151
+ for i, bucket := range out.Buckets {
152
+ lastRefresh := "never"
153
+ if bucket.LastRefresh != "" {
154
+ t, err := time.Parse(time.RFC3339, bucket.LastRefresh)
155
+ if err != nil {
156
+ return err
157
+ }
158
+ lastRefresh = since(t)
159
+ }
160
+ fmt.Fprintf(tw, " Bucket %2d (%d peers) - refreshed %s:\t\t\t\n", i, len(bucket.Peers), lastRefresh)
161
+ fmt.Fprintln(tw, " Peer\tlast useful\tlast queried\tAgent Version")
162
+
163
+ for _, p := range bucket.Peers {
164
+ lastUseful := "never"
165
+ if p.LastUsefulAt != "" {
166
+ t, err := time.Parse(time.RFC3339, p.LastUsefulAt)
167
+ if err != nil {
168
+ return err
169
+ }
170
+ lastUseful = since(t)
171
+ }
172
+
173
+ lastQueried := "never"
174
+ if p.LastUsefulAt != "" {
175
+ t, err := time.Parse(time.RFC3339, p.LastQueriedAt)
176
+ if err != nil {
177
+ return err
178
+ }
179
+ lastQueried = since(t)
180
+ }
181
+
182
+ state := " "
183
+ if p.Connected {
184
+ state = "@"
185
+ }
186
+ fmt.Fprintf(tw, " %s %s\t%s\t%s\t%s\n", state, p.ID, lastUseful, lastQueried, p.AgentVersion)
187
+ }
188
+ fmt.Fprintln(tw, "\t\t\t")
189
+ }
190
+ return nil
191
+ }),
192
+ },
193
+ Type: dhtStat{},
194
+}