@cryptotaxi247 / kubo / commits / 280c7e7e0

implement diagnostics service

Jeromy committed Oct 9, 2014 at 22:51 UTC 280c7e7e06fe69942d9b0b08e0bc752ce2e1715b
13 files changed +391 -26
Godeps/Godeps.json
+1 -1
@@ -1,6 +1,6 @@
1 {
2 "ImportPath": "github.com/jbenet/go-ipfs",
3 - "GoVersion": "go1.3",
3 + "GoVersion": "go1.3.3",
4 "Packages": [
5 "./..."
6 ],
cmd/ipfs/ipfs.go
+1
@@ -81,6 +81,7 @@ func ipfsCmd(c *commander.Command, args []string) error {
81 }
82
83 func main() {
84 + u.AllLoggersOn()
85 u.Debug = false
86
87 // setup logging
core/core.go
+23 -12
@@ -13,6 +13,7 @@ import (
13 bserv "github.com/jbenet/go-ipfs/blockservice"
14 config "github.com/jbenet/go-ipfs/config"
15 ci "github.com/jbenet/go-ipfs/crypto"
16 + diag "github.com/jbenet/go-ipfs/diagnostics"
17 exchange "github.com/jbenet/go-ipfs/exchange"
18 bitswap "github.com/jbenet/go-ipfs/exchange/bitswap"
19 merkledag "github.com/jbenet/go-ipfs/merkledag"
@@ -64,6 +65,9 @@ type IpfsNode struct {
65
66 // the name system, resolves paths to hashes
67 Namesys namesys.NameSystem
68 +
69 + // the diagnostics service
70 + Diagnostics *diag.Diagnostics
71 }
72
73 // NewIpfsNode constructs a new IpfsNode based on the given config.
@@ -103,12 +107,14 @@ func NewIpfsNode(cfg *config.Config, online bool) (*IpfsNode, error) {
107 // TODO: refactor so we can use IpfsRouting interface instead of being DHT-specific
108 route *dht.IpfsDHT
109 exchangeSession exchange.Interface
110 + diagnostics *diag.Diagnostics
111 )
112
113 if online {
114
115 dhtService := netservice.NewService(nil) // nil handler for now, need to patch it
116 exchangeService := netservice.NewService(nil) // nil handler for now, need to patch it
117 + diagService := netservice.NewService(nil)
118
119 if err := dhtService.Start(ctx); err != nil {
120 return nil, err
@@ -118,14 +124,18 @@ func NewIpfsNode(cfg *config.Config, online bool) (*IpfsNode, error) {
124 }
125
126 net, err = inet.NewIpfsNetwork(context.TODO(), local, peerstore, &mux.ProtocolMap{
121 - mux.ProtocolID_Routing: dhtService,
122 - mux.ProtocolID_Exchange: exchangeService,
127 + mux.ProtocolID_Routing: dhtService,
128 + mux.ProtocolID_Exchange: exchangeService,
129 + mux.ProtocolID_Diagnostic: diagService,
130 // add protocol services here.
131 })
132 if err != nil {
133 return nil, err
134 }
135
136 + diagnostics = diag.NewDiagnostics(local, net, diagService)
137 + diagService.Handler = diagnostics
138 +
139 route = dht.NewDHT(local, peerstore, net, dhtService, d)
140 // TODO(brian): perform this inside NewDHT factory method
141 dhtService.Handler = route // wire the handler to the service.
@@ -149,16 +159,17 @@ func NewIpfsNode(cfg *config.Config, online bool) (*IpfsNode, error) {
159
160 success = true
161 return &IpfsNode{
152 - Config: cfg,
153 - Peerstore: peerstore,
154 - Datastore: d,
155 - Blocks: bs,
156 - DAG: dag,
157 - Resolver: &path.Resolver{DAG: dag},
158 - Exchange: exchangeSession,
159 - Identity: local,
160 - Routing: route,
161 - Namesys: ns,
162 + Config: cfg,
163 + Peerstore: peerstore,
164 + Datastore: d,
165 + Blocks: bs,
166 + DAG: dag,
167 + Resolver: &path.Resolver{DAG: dag},
168 + Exchange: exchangeSession,
169 + Identity: local,
170 + Routing: route,
171 + Namesys: ns,
172 + Diagnostics: diagnostics,
173 }, nil
174 }
175
daemon/daemon.go
+10
@@ -8,6 +8,7 @@ import (
8 "os"
9 "path"
10 "sync"
11 + "time"
12
13 core "github.com/jbenet/go-ipfs/core"
14 "github.com/jbenet/go-ipfs/core/commands"
@@ -136,6 +137,15 @@ func (dl *DaemonListener) handleConnection(conn net.Conn) {
137 err = commands.Publish(dl.node, command.Args, command.Opts, conn)
138 case "resolve":
139 err = commands.Resolve(dl.node, command.Args, command.Opts, conn)
140 + case "diag":
141 + log.Debug("DIAGNOSTIC!")
142 + info, err := dl.node.Diagnostics.GetDiagnostic(time.Second * 20)
143 + if err != nil {
144 + fmt.Fprintln(conn, err)
145 + return
146 + }
147 + enc := json.NewEncoder(conn)
148 + err = enc.Encode(info)
149 default:
150 err = fmt.Errorf("Invalid Command: '%s'", command.Command)
151 }
diagnostics/diag.go new
+263
@@ -0,0 +1,263 @@
1 +package diagnostic
2 +
3 +import (
4 + "bytes"
5 + "encoding/json"
6 + "errors"
7 + "io"
8 + "sync"
9 + "time"
10 +
11 + "crypto/rand"
12 +
13 + "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
14 + "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/goprotobuf/proto"
15 +
16 + net "github.com/jbenet/go-ipfs/net"
17 + msg "github.com/jbenet/go-ipfs/net/message"
18 + peer "github.com/jbenet/go-ipfs/peer"
19 + util "github.com/jbenet/go-ipfs/util"
20 +)
21 +
22 +var log = util.Logger("diagnostics")
23 +
24 +type Diagnostics struct {
25 + network net.Network
26 + sender net.Sender
27 + self *peer.Peer
28 +
29 + diagLock sync.Mutex
30 + diagMap map[string]time.Time
31 + birth time.Time
32 +}
33 +
34 +func NewDiagnostics(self *peer.Peer, inet net.Network, sender net.Sender) *Diagnostics {
35 + return &Diagnostics{
36 + network: inet,
37 + sender: sender,
38 + self: self,
39 + diagMap: make(map[string]time.Time),
40 + birth: time.Now(),
41 + }
42 +}
43 +
44 +type connDiagInfo struct {
45 + Latency time.Duration
46 + ID string
47 +}
48 +
49 +type diagInfo struct {
50 + ID string
51 + Connections []connDiagInfo
52 + Keys []string
53 + LifeSpan time.Duration
54 + CodeVersion string
55 +}
56 +
57 +func (di *diagInfo) Marshal() []byte {
58 + b, err := json.Marshal(di)
59 + if err != nil {
60 + panic(err)
61 + }
62 + //TODO: also consider compressing this. There will be a lot of these
63 + return b
64 +}
65 +
66 +func (d *Diagnostics) getPeers() []*peer.Peer {
67 + // <HACKY>
68 + n, ok := d.network.(*net.IpfsNetwork)
69 + if !ok {
70 + return nil
71 + }
72 + s := n.GetSwarm()
73 + return s.GetPeerList()
74 + // </HACKY>
75 +}
76 +
77 +func (d *Diagnostics) getDiagInfo() *diagInfo {
78 + di := new(diagInfo)
79 + di.CodeVersion = "github.com/jbenet/go-ipfs"
80 + di.ID = d.self.ID.Pretty()
81 + di.LifeSpan = time.Since(d.birth)
82 + di.Keys = nil // Currently no way to query datastore
83 +
84 + for _, p := range d.getPeers() {
85 + di.Connections = append(di.Connections, connDiagInfo{p.GetLatency(), p.ID.Pretty()})
86 + }
87 + return di
88 +}
89 +
90 +func newID() string {
91 + id := make([]byte, 4)
92 + rand.Read(id)
93 + return string(id)
94 +}
95 +
96 +func (d *Diagnostics) GetDiagnostic(timeout time.Duration) ([]*diagInfo, error) {
97 + log.Debug("Getting diagnostic.")
98 + ctx, _ := context.WithTimeout(context.TODO(), timeout)
99 +
100 + diagID := newID()
101 + d.diagLock.Lock()
102 + d.diagMap[diagID] = time.Now()
103 + d.diagLock.Unlock()
104 +
105 + log.Debug("Begin Diagnostic")
106 +
107 + peers := d.getPeers()
108 + log.Debug("Sending diagnostic request to %d peers.", len(peers))
109 +
110 + var out []*diagInfo
111 + di := d.getDiagInfo()
112 + out = append(out, di)
113 +
114 + pmes := newMessage(diagID)
115 + for _, p := range peers {
116 + log.Debug("Sending getDiagnostic to: %s", p)
117 + data, err := d.getDiagnosticFromPeer(ctx, p, pmes)
118 + if err != nil {
119 + log.Error("GetDiagnostic error: %v", err)
120 + continue
121 + }
122 + buf := bytes.NewBuffer(data)
123 + dec := json.NewDecoder(buf)
124 + for {
125 + di := new(diagInfo)
126 + err := dec.Decode(di)
127 + if err != nil {
128 + if err != io.EOF {
129 + log.Error("error decoding diagInfo: %v", err)
130 + }
131 + break
132 + }
133 + out = append(out, di)
134 + }
135 + }
136 + return out, nil
137 +}
138 +
139 +// TODO: this method no longer needed.
140 +func (d *Diagnostics) getDiagnosticFromPeer(ctx context.Context, p *peer.Peer, mes *Message) ([]byte, error) {
141 + rpmes, err := d.sendRequest(ctx, p, mes)
142 + if err != nil {
143 + return nil, err
144 + }
145 + return rpmes.GetData(), nil
146 +}
147 +
148 +func newMessage(diagID string) *Message {
149 + pmes := new(Message)
150 + pmes.DiagID = proto.String(diagID)
151 + return pmes
152 +}
153 +
154 +func (d *Diagnostics) sendRequest(ctx context.Context, p *peer.Peer, pmes *Message) (*Message, error) {
155 +
156 + mes, err := msg.FromObject(p, pmes)
157 + if err != nil {
158 + return nil, err
159 + }
160 +
161 + start := time.Now()
162 +
163 + rmes, err := d.sender.SendRequest(ctx, mes)
164 + if err != nil {
165 + return nil, err
166 + }
167 + if rmes == nil {
168 + return nil, errors.New("no response to request")
169 + }
170 +
171 + rtt := time.Since(start)
172 + log.Info("diagnostic request took: %s", rtt.String())
173 +
174 + rpmes := new(Message)
175 + if err := proto.Unmarshal(rmes.Data(), rpmes); err != nil {
176 + return nil, err
177 + }
178 +
179 + return rpmes, nil
180 +}
181 +
182 +// NOTE: not yet finished, low priority
183 +func (d *Diagnostics) handleDiagnostic(p *peer.Peer, pmes *Message) (*Message, error) {
184 + resp := newMessage(pmes.GetDiagID())
185 + d.diagLock.Lock()
186 + _, found := d.diagMap[pmes.GetDiagID()]
187 + if found {
188 + d.diagLock.Unlock()
189 + return resp, nil
190 + }
191 + d.diagMap[pmes.GetDiagID()] = time.Now()
192 + d.diagLock.Unlock()
193 +
194 + buf := new(bytes.Buffer)
195 + di := d.getDiagInfo()
196 + buf.Write(di.Marshal())
197 +
198 + ctx, _ := context.WithTimeout(context.TODO(), time.Second*10)
199 +
200 + for _, p := range d.getPeers() {
201 + log.Debug("Sending diagnostic request to peer: %s", p)
202 + out, err := d.getDiagnosticFromPeer(ctx, p, pmes)
203 + if err != nil {
204 + log.Error("getDiagnostic error: %v", err)
205 + continue
206 + }
207 + _, err = buf.Write(out)
208 + if err != nil {
209 + log.Error("getDiagnostic write output error: %v", err)
210 + continue
211 + }
212 + }
213 +
214 + resp.Data = buf.Bytes()
215 + return resp, nil
216 +}
217 +
218 +func (d *Diagnostics) HandleMessage(ctx context.Context, mes msg.NetMessage) msg.NetMessage {
219 + mData := mes.Data()
220 + if mData == nil {
221 + log.Error("message did not include Data")
222 + return nil
223 + }
224 +
225 + mPeer := mes.Peer()
226 + if mPeer == nil {
227 + log.Error("message did not include a Peer")
228 + return nil
229 + }
230 +
231 + // deserialize msg
232 + pmes := new(Message)
233 + err := proto.Unmarshal(mData, pmes)
234 + if err != nil {
235 + log.Error("Failed to decode protobuf message: %v", err)
236 + return nil
237 + }
238 +
239 + // Print out diagnostic
240 + log.Info("[peer: %s] Got message from [%s]\n",
241 + d.self.ID.Pretty(), mPeer.ID.Pretty())
242 +
243 + // dispatch handler.
244 + rpmes, err := d.handleDiagnostic(mPeer, pmes)
245 + if err != nil {
246 + log.Error("handleDiagnostic error: %s", err)
247 + return nil
248 + }
249 +
250 + // if nil response, return it before serializing
251 + if rpmes == nil {
252 + return nil
253 + }
254 +
255 + // serialize response msg
256 + rmes, err := msg.FromObject(mPeer, rpmes)
257 + if err != nil {
258 + log.Error("Failed to encode protobuf message: %v", err)
259 + return nil
260 + }
261 +
262 + return rmes
263 +}
diagnostics/message.pb.go new
+48
@@ -0,0 +1,48 @@
1 +// Code generated by protoc-gen-go.
2 +// source: message.proto
3 +// DO NOT EDIT!
4 +
5 +/*
6 +Package diagnostic is a generated protocol buffer package.
7 +
8 +It is generated from these files:
9 + message.proto
10 +
11 +It has these top-level messages:
12 + Message
13 +*/
14 +package diagnostic
15 +
16 +import proto "code.google.com/p/goprotobuf/proto"
17 +import math "math"
18 +
19 +// Reference imports to suppress errors if they are not otherwise used.
20 +var _ = proto.Marshal
21 +var _ = math.Inf
22 +
23 +type Message struct {
24 + DiagID *string `protobuf:"bytes,1,req" json:"DiagID,omitempty"`
25 + Data []byte `protobuf:"bytes,2,opt" json:"Data,omitempty"`
26 + XXX_unrecognized []byte `json:"-"`
27 +}
28 +
29 +func (m *Message) Reset() { *m = Message{} }
30 +func (m *Message) String() string { return proto.CompactTextString(m) }
31 +func (*Message) ProtoMessage() {}
32 +
33 +func (m *Message) GetDiagID() string {
34 + if m != nil && m.DiagID != nil {
35 + return *m.DiagID
36 + }
37 + return ""
38 +}
39 +
40 +func (m *Message) GetData() []byte {
41 + if m != nil {
42 + return m.Data
43 + }
44 + return nil
45 +}
46 +
47 +func init() {
48 +}
diagnostics/message.proto new
+6
@@ -0,0 +1,6 @@
1 +package diagnostic;
2 +
3 +message Message {
4 + required string DiagID = 1;
5 + optional bytes Data = 2;
6 +}
net/mux/mux.pb.go
+14 -13
@@ -1,4 +1,4 @@
1 -// Code generated by protoc-gen-gogo.
1 +// Code generated by protoc-gen-go.
2 // source: mux.proto
3 // DO NOT EDIT!
4
@@ -13,22 +13,21 @@ It has these top-level messages:
13 */
14 package mux
15
16 -import proto "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/gogoprotobuf/proto"
17 -import json "encoding/json"
16 +import proto "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/goprotobuf/proto"
17 import math "math"
18
20 -// Reference proto, json, and math imports to suppress error if they are not otherwise used.
19 +// Reference imports to suppress errors if they are not otherwise used.
20 var _ = proto.Marshal
22 -var _ = &json.SyntaxError{}
21 var _ = math.Inf
22
23 type ProtocolID int32
24
25 const (
28 - ProtocolID_Test ProtocolID = 0
29 - ProtocolID_Identify ProtocolID = 1
30 - ProtocolID_Routing ProtocolID = 2
31 - ProtocolID_Exchange ProtocolID = 3
26 + ProtocolID_Test ProtocolID = 0
27 + ProtocolID_Identify ProtocolID = 1
28 + ProtocolID_Routing ProtocolID = 2
29 + ProtocolID_Exchange ProtocolID = 3
30 + ProtocolID_Diagnostic ProtocolID = 4
31 )
32
33 var ProtocolID_name = map[int32]string{
@@ -36,12 +35,14 @@ var ProtocolID_name = map[int32]string{
35 1: "Identify",
36 2: "Routing",
37 3: "Exchange",
38 + 4: "Diagnostic",
39 }
40 var ProtocolID_value = map[string]int32{
41 - "Test": 0,
42 - "Identify": 1,
43 - "Routing": 2,
44 - "Exchange": 3,
41 + "Test": 0,
42 + "Identify": 1,
43 + "Routing": 2,
44 + "Exchange": 3,
45 + "Diagnostic": 4,
46 }
47
48 func (x ProtocolID) Enum() *ProtocolID {
net/mux/mux.proto
+1
@@ -5,6 +5,7 @@ enum ProtocolID {
5 Identify = 1; // setup
6 Routing = 2; // dht
7 Exchange = 3; // bitswap
8 + Diagnostic = 4;
9 }
10
11 message PBProtocolMessage {
net/net.go
+5
@@ -107,3 +107,8 @@ func (n *IpfsNetwork) Close() error {
107 n.cancel = nil
108 return nil
109 }
110 +
111 +// XXX
112 +func (n *IpfsNetwork) GetSwarm() *swarm.Swarm {
113 + return n.swarm
114 +}
net/swarm/conn.go
+2
@@ -122,10 +122,12 @@ func (s *Swarm) connSetup(c *conn.Conn) error {
122 // add to conns
123 s.connsLock.Lock()
124 if _, ok := s.conns[c.Peer.Key()]; ok {
125 + log.Debug("Conn already open!")
126 s.connsLock.Unlock()
127 return ErrAlreadyOpen
128 }
129 s.conns[c.Peer.Key()] = c
130 + log.Debug("Added conn to map!")
131 s.connsLock.Unlock()
132
133 // kick off reader goroutine
net/swarm/swarm.go
+10
@@ -201,5 +201,15 @@ func (s *Swarm) GetErrChan() chan error {
201 return s.errChan
202 }
203
204 +func (s *Swarm) GetPeerList() []*peer.Peer {
205 + var out []*peer.Peer
206 + s.connsLock.Lock()
207 + for _, p := range s.conns {
208 + out = append(out, p.Peer)
209 + }
210 + s.connsLock.Unlock()
211 + return out
212 +}
213 +
214 // Temporary to ensure that the Swarm always matches the Network interface as we are changing it
215 // var _ Network = &Swarm{}
util/util.go
+7
@@ -153,6 +153,13 @@ func SetupLogging() {
153 }
154 }
155
156 +func AllLoggersOn() {
157 + for n, log := range loggers {
158 + logging.SetLevel(logging.DEBUG, n)
159 + log.Error("setting logger: %s to %v\n", n, logging.DEBUG)
160 + }
161 +}
162 +
163 // Logger retrieves a particular logger + initializes it at a particular level
164 func Logger(name string) *logging.Logger {
165 log := logging.MustGetLogger(name)