@cryptotaxi247 / kubo / commits / a7b69500b

address concerns in PR and make log stuff more fun

Jeromy committed Oct 10, 2014 at 11:07 UTC a7b69500b14912c77a690393ca8c1f5520b6a1a1
7 files changed +78 -50
cmd/ipfs/diag.go
+3 -3
@@ -8,16 +8,16 @@ import (
8 )
9
10 var cmdIpfsDiag = &commander.Command{
11 - UsageLine: "diag",
11 + UsageLine: "net-diag",
12 Short: "Generate a diagnostics report",
13 - Long: `ipfs diag - Generate a diagnostics report.
13 + Long: `ipfs net-diag - Generate a diagnostics report.
14
15 Sends out a message to each node in the network recursively
16 requesting a listing of data about them including number of
17 connected peers and latencies between them.
18 `,
19 Run: diagCmd,
20 - Flag: *flag.NewFlagSet("ipfs-diag", flag.ExitOnError),
20 + Flag: *flag.NewFlagSet("ipfs-net-diag", flag.ExitOnError),
21 }
22
23 func init() {
core/commands/diag.go
+4 -1
@@ -18,7 +18,10 @@ func Diag(n *core.IpfsNode, args []string, opts map[string]interface{}, out io.W
18 if err != nil {
19 return err
20 }
21 - raw := opts["raw"].(bool)
21 + raw, ok := opts["raw"].(bool)
22 + if !ok {
23 + return errors.New("incorrect value to parameter 'raw'")
24 + }
25 if raw {
26 enc := json.NewEncoder(out)
27 err = enc.Encode(info)
diagnostics/diag.go
+55 -31
@@ -21,6 +21,8 @@ import (
21
22 var log = util.Logger("diagnostics")
23
24 +const ResponseTimeout = time.Second * 10
25 +
26 type Diagnostics struct {
27 network net.Network
28 sender net.Sender
@@ -64,14 +66,7 @@ func (di *diagInfo) Marshal() []byte {
66 }
67
68 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>
69 + return d.network.GetPeerList()
70 }
71
72 func (d *Diagnostics) getDiagInfo() *diagInfo {
@@ -112,30 +107,50 @@ func (d *Diagnostics) GetDiagnostic(timeout time.Duration) ([]*diagInfo, error)
107 out = append(out, di)
108
109 pmes := newMessage(diagID)
110 +
111 + respdata := make(chan []byte)
112 + sends := 0
113 for _, p := range peers {
114 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)
115 + sends++
116 + go func(p *peer.Peer) {
117 + data, err := d.getDiagnosticFromPeer(ctx, p, pmes)
118 if err != nil {
128 - if err != io.EOF {
129 - log.Error("error decoding diagInfo: %v", err)
130 - }
131 - break
119 + log.Error("GetDiagnostic error: %v", err)
120 + respdata <- nil
121 + return
122 }
133 - out = append(out, di)
123 + respdata <- data
124 + }(p)
125 + }
126 +
127 + for i := 0; i < sends; i++ {
128 + data := <-respdata
129 + if data == nil {
130 + continue
131 }
132 + out = AppendDiagnostics(data, out)
133 }
134 return out, nil
135 }
136
137 +func AppendDiagnostics(data []byte, cur []*diagInfo) []*diagInfo {
138 + buf := bytes.NewBuffer(data)
139 + dec := json.NewDecoder(buf)
140 + for {
141 + di := new(diagInfo)
142 + err := dec.Decode(di)
143 + if err != nil {
144 + if err != io.EOF {
145 + log.Error("error decoding diagInfo: %v", err)
146 + }
147 + break
148 + }
149 + cur = append(cur, di)
150 + }
151 + return cur
152 +}
153 +
154 // TODO: this method no longer needed.
155 func (d *Diagnostics) getDiagnosticFromPeer(ctx context.Context, p *peer.Peer, mes *Message) ([]byte, error) {
156 rpmes, err := d.sendRequest(ctx, p, mes)
@@ -179,7 +194,6 @@ func (d *Diagnostics) sendRequest(ctx context.Context, p *peer.Peer, pmes *Messa
194 return rpmes, nil
195 }
196
182 -// NOTE: not yet finished, low priority
197 func (d *Diagnostics) handleDiagnostic(p *peer.Peer, pmes *Message) (*Message, error) {
198 resp := newMessage(pmes.GetDiagID())
199 d.diagLock.Lock()
@@ -195,16 +209,26 @@ func (d *Diagnostics) handleDiagnostic(p *peer.Peer, pmes *Message) (*Message, e
209 di := d.getDiagInfo()
210 buf.Write(di.Marshal())
211
198 - ctx, _ := context.WithTimeout(context.TODO(), time.Second*10)
212 + ctx, _ := context.WithTimeout(context.TODO(), ResponseTimeout)
213
214 + respdata := make(chan []byte)
215 + sendcount := 0
216 for _, p := range d.getPeers() {
217 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)
218 + go func(p *peer.Peer) {
219 + out, err := d.getDiagnosticFromPeer(ctx, p, pmes)
220 + if err != nil {
221 + log.Error("getDiagnostic error: %v", err)
222 + respdata <- nil
223 + return
224 + }
225 + respdata <- out
226 + }(p)
227 + }
228 +
229 + for i := 0; i < sendcount; i++ {
230 + out := <-respdata
231 + _, err := buf.Write(out)
232 if err != nil {
233 log.Error("getDiagnostic write output error: %v", err)
234 continue
net/interface.go
+3
@@ -28,6 +28,9 @@ type Network interface {
28 // GetProtocols returns the protocols registered in the network.
29 GetProtocols() *mux.ProtocolMap
30
31 + // GetPeerList returns the list of peers currently connected in this network.
32 + GetPeerList() []*peer.Peer
33 +
34 // SendMessage sends given Message out
35 SendMessage(msg.NetMessage) error
36
net/net.go
+2 -3
@@ -108,7 +108,6 @@ func (n *IpfsNetwork) Close() error {
108 return nil
109 }
110
111 -// XXX
112 -func (n *IpfsNetwork) GetSwarm() *swarm.Swarm {
113 - return n.swarm
111 +func (n *IpfsNetwork) GetPeerList() []*peer.Peer {
112 + return n.swarm.GetPeerList()
113 }
net/swarm/swarm.go
+2 -2
@@ -203,11 +203,11 @@ func (s *Swarm) GetErrChan() chan error {
203
204 func (s *Swarm) GetPeerList() []*peer.Peer {
205 var out []*peer.Peer
206 - s.connsLock.Lock()
206 + s.connsLock.RLock()
207 for _, p := range s.conns {
208 out = append(out, p.Peer)
209 }
210 - s.connsLock.Unlock()
210 + s.connsLock.RUnlock()
211 return out
212 }
213
util/util.go
+9 -10
@@ -138,25 +138,24 @@ func SetupLogging() {
138 logging.SetBackend(backend)
139 logging.SetFormatter(logging.MustStringFormatter(LogFormat))
140
141 - // just uncomment Debug = True right here for all logging.
142 - // but please don't commit that.
143 - // Debug = True
144 - if Debug {
145 - logging.SetLevel(logging.DEBUG, "")
146 - } else {
147 - logging.SetLevel(logging.ERROR, "")
148 - }
141 + logging.SetLevel(logging.ERROR, "")
142
143 for n, log := range loggers {
144 logging.SetLevel(logging.ERROR, n)
152 - log.Error("setting logger: %s to %v\n", n, logging.ERROR)
145 + log.Error("setting logger: %s to %v", n, logging.ERROR)
146 + }
147 +
148 + logenv := os.Getenv("IPFS_LOGGING")
149 + if logenv == "all" {
150 + AllLoggersOn()
151 }
152 }
153
154 func AllLoggersOn() {
155 + logging.SetLevel(logging.DEBUG, "")
156 for n, log := range loggers {
157 logging.SetLevel(logging.DEBUG, n)
159 - log.Error("setting logger: %s to %v\n", n, logging.DEBUG)
158 + log.Error("setting logger: %s to %v", n, logging.DEBUG)
159 }
160 }
161