@cryptotaxi247 / kubo / commits / 9c6228d18

bitswap and dht: lots of debugging logs

Juan Batiz-Benet committed Jan 3, 2015 at 08:54 UTC 9c6228d18f4c0a0a34872097206f80f5753c34a3
5 files changed +52 -7
exchange/bitswap/bitswap.go
+7
@@ -108,6 +108,7 @@ type bitswap struct {
108 // GetBlock attempts to retrieve a particular block from peers within the
109 // deadline enforced by the context.
110 func (bs *bitswap) GetBlock(parent context.Context, k u.Key) (*blocks.Block, error) {
111 + log := log.Prefix("bitswap(%s).GetBlock(%s)", bs.self, k)
112
113 // Any async work initiated by this function must end when this function
114 // returns. To ensure this, derive a new context. Note that it is okay to
@@ -120,10 +121,12 @@ func (bs *bitswap) GetBlock(parent context.Context, k u.Key) (*blocks.Block, err
121
122 ctx = eventlog.ContextWithLoggable(ctx, eventlog.Uuid("GetBlockRequest"))
123 log.Event(ctx, "GetBlockRequestBegin", &k)
124 + log.Debugf("GetBlockRequestBegin")
125
126 defer func() {
127 cancelFunc()
128 log.Event(ctx, "GetBlockRequestEnd", &k)
129 + log.Debugf("GetBlockRequestEnd")
130 }()
131
132 promise, err := bs.GetBlocks(ctx, []u.Key{k})
@@ -263,12 +266,16 @@ func (bs *bitswap) sendWantlistToProviders(ctx context.Context) {
266 }
267
268 func (bs *bitswap) taskWorker(ctx context.Context) {
269 + log := log.Prefix("bitswap(%s).taskWorker", bs.self)
270 for {
271 select {
272 case <-ctx.Done():
273 + log.Debugf("exiting")
274 return
275 case envelope := <-bs.engine.Outbox():
276 + log.Debugf("message to %s sending...", envelope.Peer)
277 bs.send(ctx, envelope.Peer, envelope.Message)
278 + log.Debugf("message to %s sent", envelope.Peer)
279 }
280 }
281 }
exchange/bitswap/decision/engine.go
+9 -1
@@ -91,6 +91,7 @@ func NewEngine(ctx context.Context, bs bstore.Blockstore) *Engine {
91 }
92
93 func (e *Engine) taskWorker(ctx context.Context) {
94 + log := log.Prefix("bitswap.Engine.taskWorker")
95 for {
96 nextTask := e.peerRequestQueue.Pop()
97 if nextTask == nil {
@@ -98,11 +99,16 @@ func (e *Engine) taskWorker(ctx context.Context) {
99 // Wait until there are!
100 select {
101 case <-ctx.Done():
102 + log.Debugf("exiting: %s", ctx.Err())
103 return
104 case <-e.workSignal:
105 + log.Debugf("woken up")
106 }
107 continue
108 }
109 + log := log.Prefix("%s", nextTask)
110 + log.Debugf("processing")
111 +
112 block, err := e.bs.Get(nextTask.Entry.Key)
113 if err != nil {
114 log.Warning("engine: task exists to send block, but block is not in blockstore")
@@ -113,10 +119,12 @@ func (e *Engine) taskWorker(ctx context.Context) {
119 m := bsmsg.New()
120 m.AddBlock(block)
121 // TODO: maybe add keys from our wantlist?
122 + log.Debugf("sending...")
123 select {
124 case <-ctx.Done():
125 return
126 case e.outbox <- Envelope{Peer: nextTask.Target, Message: m}:
127 + log.Debugf("sent")
128 }
129 }
130 }
@@ -140,7 +148,7 @@ func (e *Engine) Peers() []peer.ID {
148 // MessageReceived performs book-keeping. Returns error if passed invalid
149 // arguments.
150 func (e *Engine) MessageReceived(p peer.ID, m bsmsg.BitSwapMessage) error {
143 - log := log.Prefix("Engine.MessageReceived(%s)", p)
151 + log := log.Prefix("bitswap.Engine.MessageReceived(%s)", p)
152 log.Debugf("enter")
153 defer log.Debugf("exit")
154
exchange/bitswap/decision/taskqueue.go
+5
@@ -1,6 +1,7 @@
1 package decision
2
3 import (
4 + "fmt"
5 "sync"
6
7 wantlist "github.com/jbenet/go-ipfs/exchange/bitswap/wantlist"
@@ -30,6 +31,10 @@ type task struct {
31 Trash bool
32 }
33
34 +func (t *task) String() string {
35 + return fmt.Sprintf("<Task %s, %s, %v>", t.Target, t.Entry.Key, t.Trash)
36 +}
37 +
38 // Push currently adds a new task to the end of the list
39 func (tl *taskQueue) Push(entry wantlist.Entry, to peer.ID) {
40 tl.lock.Lock()
exchange/bitswap/network/ipfs_impl.go
+30 -5
@@ -2,15 +2,17 @@ package network
2
3 import (
4 context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
5 +
6 bsmsg "github.com/jbenet/go-ipfs/exchange/bitswap/message"
7 host "github.com/jbenet/go-ipfs/p2p/host"
8 inet "github.com/jbenet/go-ipfs/p2p/net"
9 peer "github.com/jbenet/go-ipfs/p2p/peer"
10 routing "github.com/jbenet/go-ipfs/routing"
11 util "github.com/jbenet/go-ipfs/util"
12 + eventlog "github.com/jbenet/go-ipfs/util/eventlog"
13 )
14
13 -var log = util.Logger("bitswap_network")
15 +var log = eventlog.Logger("bitswap_network")
16
17 // NewFromIpfsHost returns a BitSwapNetwork supported by underlying IPFS host
18 func NewFromIpfsHost(host host.Host, r routing.IpfsRouting) BitSwapNetwork {
@@ -41,13 +43,23 @@ func (bsnet *impl) SendMessage(
43 p peer.ID,
44 outgoing bsmsg.BitSwapMessage) error {
45
46 + log := log.Prefix("bitswap net SendMessage to %s", p)
47 +
48 + log.Debug("opening stream")
49 s, err := bsnet.host.NewStream(ProtocolBitswap, p)
50 if err != nil {
51 return err
52 }
53 defer s.Close()
54
50 - return outgoing.ToNet(s)
55 + log.Debug("sending")
56 + if err := outgoing.ToNet(s); err != nil {
57 + log.Errorf("error: %s", err)
58 + return err
59 + }
60 +
61 + log.Debug("sent")
62 + return err
63 }
64
65 func (bsnet *impl) SendRequest(
@@ -55,18 +67,30 @@ func (bsnet *impl) SendRequest(
67 p peer.ID,
68 outgoing bsmsg.BitSwapMessage) (bsmsg.BitSwapMessage, error) {
69
58 - log.Debugf("bsnet SendRequest to %s", p)
70 + log := log.Prefix("bitswap net SendRequest to %s", p)
71 +
72 + log.Debug("opening stream")
73 s, err := bsnet.host.NewStream(ProtocolBitswap, p)
74 if err != nil {
75 return nil, err
76 }
77 defer s.Close()
78
79 + log.Debug("sending")
80 if err := outgoing.ToNet(s); err != nil {
81 + log.Errorf("error: %s", err)
82 return nil, err
83 }
84
69 - return bsmsg.FromNet(s)
85 + log.Debug("sent, now receiveing")
86 + incoming, err := bsmsg.FromNet(s)
87 + if err != nil {
88 + log.Errorf("error: %s", err)
89 + return incoming, err
90 + }
91 +
92 + log.Debug("received")
93 + return incoming, nil
94 }
95
96 func (bsnet *impl) SetDelegate(r Receiver) {
@@ -106,11 +130,12 @@ func (bsnet *impl) handleNewStream(s inet.Stream) {
130 received, err := bsmsg.FromNet(s)
131 if err != nil {
132 go bsnet.receiver.ReceiveError(err)
133 + log.Errorf("bitswap net handleNewStream from %s error: %s", s.Conn().RemotePeer(), err)
134 return
135 }
136
137 p := s.Conn().RemotePeer()
138 ctx := context.Background()
114 - log.Debugf("bsnet handleNewStream from %s", s.Conn().RemotePeer())
139 + log.Debugf("bitswap net handleNewStream from %s", s.Conn().RemotePeer())
140 bsnet.receiver.ReceiveMessage(ctx, p, received)
141 }
routing/dht/handlers.go
+1 -1
@@ -148,7 +148,7 @@ func (dht *IpfsDHT) handleFindPeer(ctx context.Context, p peer.ID, pmes *pb.Mess
148 }
149
150 if closest == nil {
151 - log.Warningf("handleFindPeer: could not find anything.")
151 + log.Warningf("%s handleFindPeer %s: could not find anything.", dht.self, p)
152 return resp, nil
153 }
154