master
go 252 lines 6.4 KB
Raw
1 package commands
2
3 import (
4 "fmt"
5 "io"
6 "text/tabwriter"
7 "time"
8
9 cmdenv "github.com/ipfs/kubo/core/commands/cmdenv"
10 "github.com/ipfs/kubo/core/commands/cmdutils"
11
12 cmds "github.com/ipfs/go-ipfs-cmds"
13 dht "github.com/libp2p/go-libp2p-kad-dht"
14 "github.com/libp2p/go-libp2p-kad-dht/fullrt"
15 kbucket "github.com/libp2p/go-libp2p-kbucket"
16 "github.com/libp2p/go-libp2p/core/network"
17 pstore "github.com/libp2p/go-libp2p/core/peerstore"
18 )
19
20 type dhtPeerInfo struct {
21 ID string
22 Connected bool
23 AgentVersion string
24 LastUsefulAt string
25 LastQueriedAt string
26 }
27
28 type dhtStat struct {
29 Name string
30 Buckets []dhtBucket
31 }
32
33 type dhtBucket struct {
34 LastRefresh string
35 Peers []dhtPeerInfo
36 }
37
38 var statDhtCmd = &cmds.Command{
39 Helptext: cmds.HelpText{
40 Tagline: "Returns statistics about the node's DHT(s).",
41 ShortDescription: `
42 Returns statistics about the DHT(s) the node is participating in.
43
44 This interface is not stable and may change from release to release.
45 `,
46 },
47 Arguments: []cmds.Argument{
48 cmds.StringArg("dht", false, true, "The DHT whose table should be listed (wanserver, lanserver, wan, lan). "+
49 "wan and lan refer to client routing tables. When using the experimental DHT client only WAN is supported. Defaults to wan and lan."),
50 },
51 Options: []cmds.Option{},
52 Run: func(req *cmds.Request, res cmds.ResponseEmitter, env cmds.Environment) error {
53 nd, err := cmdenv.GetNode(env)
54 if err != nil {
55 return err
56 }
57
58 if !nd.IsOnline {
59 return ErrNotOnline
60 }
61
62 if nd.DHT == nil {
63 return ErrNotDHT
64 }
65
66 id := kbucket.ConvertPeerID(nd.Identity)
67
68 dhts := req.Arguments
69 if len(dhts) == 0 {
70 dhts = []string{"wan", "lan"}
71 }
72
73 dhttypeloop:
74 for _, name := range dhts {
75 var dht *dht.IpfsDHT
76
77 var separateClient bool
78 // Check if using separate DHT client (e.g., accelerated DHT)
79 if nd.HasActiveDHTClient() && nd.DHTClient != nd.DHT {
80 separateClient = true
81 }
82
83 switch name {
84 case "wan":
85 if separateClient {
86 client, ok := nd.DHTClient.(*fullrt.FullRT)
87 if !ok {
88 return cmds.Errorf(cmds.ErrClient, "could not generate stats for the WAN DHT client type")
89 }
90 peerMap := client.Stat()
91 buckets := make([]dhtBucket, 1)
92 b := &dhtBucket{}
93 for _, p := range peerMap {
94 info := dhtPeerInfo{ID: p.String()}
95
96 if ver, err := nd.Peerstore.Get(p, "AgentVersion"); err == nil {
97 if vs, ok := ver.(string); ok {
98 info.AgentVersion = cmdutils.CleanAndTrim(vs)
99 }
100 } else if err == pstore.ErrNotFound {
101 // ignore
102 } else {
103 // this is a bug, usually.
104 log.Errorw(
105 "failed to get agent version from peerstore",
106 "error", err,
107 )
108 }
109
110 info.Connected = nd.PeerHost.Network().Connectedness(p) == network.Connected
111 b.Peers = append(b.Peers, info)
112 }
113 buckets[0] = *b
114
115 if err := res.Emit(dhtStat{
116 Name: name,
117 Buckets: buckets,
118 }); err != nil {
119 return err
120 }
121 continue dhttypeloop
122 }
123 fallthrough
124 case "wanserver":
125 dht = nd.DHT.WAN
126 case "lan":
127 if separateClient {
128 return cmds.Errorf(cmds.ErrClient, "no LAN client found")
129 }
130 fallthrough
131 case "lanserver":
132 dht = nd.DHT.LAN
133 default:
134 return cmds.Errorf(cmds.ErrClient, "unknown dht type: %s", name)
135 }
136
137 rt := dht.RoutingTable()
138 lastRefresh := rt.GetTrackedCplsForRefresh()
139 infos := rt.GetPeerInfos()
140 buckets := make([]dhtBucket, 0, len(lastRefresh))
141 for _, pi := range infos {
142 cpl := kbucket.CommonPrefixLen(id, kbucket.ConvertPeerID(pi.Id))
143 if len(buckets) <= cpl {
144 buckets = append(buckets, make([]dhtBucket, 1+cpl-len(buckets))...)
145 }
146
147 info := dhtPeerInfo{ID: pi.Id.String()}
148
149 if ver, err := nd.Peerstore.Get(pi.Id, "AgentVersion"); err == nil {
150 if vs, ok := ver.(string); ok {
151 info.AgentVersion = cmdutils.CleanAndTrim(vs)
152 }
153 } else if err == pstore.ErrNotFound {
154 // ignore
155 } else {
156 // this is a bug, usually.
157 log.Errorw(
158 "failed to get agent version from peerstore",
159 "error", err,
160 )
161 }
162 if !pi.LastUsefulAt.IsZero() {
163 info.LastUsefulAt = pi.LastUsefulAt.Format(time.RFC3339)
164 }
165
166 if !pi.LastSuccessfulOutboundQueryAt.IsZero() {
167 info.LastQueriedAt = pi.LastSuccessfulOutboundQueryAt.Format(time.RFC3339)
168 }
169
170 info.Connected = nd.PeerHost.Network().Connectedness(pi.Id) == network.Connected
171
172 buckets[cpl].Peers = append(buckets[cpl].Peers, info)
173 }
174 for i := 0; i < len(buckets) && i < len(lastRefresh); i++ {
175 refreshTime := lastRefresh[i]
176 if !refreshTime.IsZero() {
177 buckets[i].LastRefresh = refreshTime.Format(time.RFC3339)
178 }
179 }
180 if err := res.Emit(dhtStat{
181 Name: name,
182 Buckets: buckets,
183 }); err != nil {
184 return err
185 }
186 }
187
188 return nil
189 },
190 Encoders: cmds.EncoderMap{
191 cmds.Text: cmds.MakeTypedEncoder(func(req *cmds.Request, w io.Writer, out dhtStat) error {
192 tw := tabwriter.NewWriter(w, 4, 4, 2, ' ', 0)
193 defer tw.Flush()
194
195 // Formats a time into XX ago and remove any decimal
196 // parts. That is, change "2m3.00010101s" to "2m3s ago".
197 now := time.Now()
198 since := func(t time.Time) string {
199 return now.Sub(t).Round(time.Second).String() + " ago"
200 }
201
202 count := 0
203 for _, bucket := range out.Buckets {
204 count += len(bucket.Peers)
205 }
206
207 fmt.Fprintf(tw, "DHT %s (%d peers):\t\t\t\n", out.Name, count)
208
209 for i, bucket := range out.Buckets {
210 lastRefresh := "never"
211 if bucket.LastRefresh != "" {
212 t, err := time.Parse(time.RFC3339, bucket.LastRefresh)
213 if err != nil {
214 return err
215 }
216 lastRefresh = since(t)
217 }
218 fmt.Fprintf(tw, " Bucket %2d (%d peers) - refreshed %s:\t\t\t\n", i, len(bucket.Peers), lastRefresh)
219 fmt.Fprintln(tw, " Peer\tlast useful\tlast queried\tAgent Version")
220
221 for _, p := range bucket.Peers {
222 lastUseful := "never"
223 if p.LastUsefulAt != "" {
224 t, err := time.Parse(time.RFC3339, p.LastUsefulAt)
225 if err != nil {
226 return err
227 }
228 lastUseful = since(t)
229 }
230
231 lastQueried := "never"
232 if p.LastUsefulAt != "" {
233 t, err := time.Parse(time.RFC3339, p.LastQueriedAt)
234 if err != nil {
235 return err
236 }
237 lastQueried = since(t)
238 }
239
240 state := " "
241 if p.Connected {
242 state = "@"
243 }
244 fmt.Fprintf(tw, " %s %s\t%s\t%s\t%s\n", state, p.ID, lastUseful, lastQueried, p.AgentVersion)
245 }
246 fmt.Fprintln(tw, "\t\t\t")
247 }
248 return nil
249 }),
250 },
251 Type: dhtStat{},
252 }