@cryptotaxi247 / kubo / commits / ff1bf3058

add in some events to bitswap to emit worker information

License: MIT Signed-off-by: Jeromy <jeromyj@gmail.com>

Jeromy committed Jul 7, 2015 at 12:14 UTC ff1bf3058e6b6587debec973c2e2316e87bb425f
2 files changed +27 -5
exchange/bitswap/bitswap.go
+9 -1
@@ -150,7 +150,8 @@ func (bs *Bitswap) GetBlock(parent context.Context, k key.Key) (*blocks.Block, e
150 ctx, cancelFunc := context.WithCancel(parent)
151
152 ctx = eventlog.ContextWithLoggable(ctx, eventlog.Uuid("GetBlockRequest"))
153 - defer log.EventBegin(ctx, "GetBlockRequest", &k).Done()
153 + log.Event(ctx, "Bitswap.GetBlockRequest.Start", &k)
154 + defer log.Event(ctx, "Bitswap.GetBlockRequest.End", &k)
155
156 defer func() {
157 cancelFunc()
@@ -200,6 +201,10 @@ func (bs *Bitswap) GetBlocks(ctx context.Context, keys []key.Key) (<-chan *block
201 }
202 promise := bs.notifications.Subscribe(ctx, keys...)
203
204 + for _, k := range keys {
205 + log.Event(ctx, "Bitswap.GetBlockRequest.Start", &k)
206 + }
207 +
208 bs.wm.WantBlocks(keys)
209
210 req := &blockRequest{
@@ -310,6 +315,9 @@ func (bs *Bitswap) ReceiveMessage(ctx context.Context, p peer.ID, incoming bsmsg
315 return
316 }
317
318 + k := b.Key()
319 + log.Event(ctx, "Bitswap.GetBlockRequest.End", &k)
320 +
321 log.Debugf("got block %s from %s (%d,%d)", b, p, brecvd, bdup)
322 hasBlockCtx, cancel := context.WithTimeout(ctx, hasBlockTimeout)
323 if err := bs.HasBlock(hasBlockCtx, b); err != nil {
exchange/bitswap/workers.go
+18 -4
@@ -7,7 +7,9 @@ import (
7
8 process "github.com/ipfs/go-ipfs/Godeps/_workspace/src/github.com/jbenet/goprocess"
9 context "github.com/ipfs/go-ipfs/Godeps/_workspace/src/golang.org/x/net/context"
10 +
11 key "github.com/ipfs/go-ipfs/blocks/key"
12 + eventlog "github.com/ipfs/go-ipfs/thirdparty/eventlog"
13 )
14
15 var TaskWorkerCount = 8
@@ -36,8 +38,9 @@ func (bs *Bitswap) startWorkers(px process.Process, ctx context.Context) {
38
39 // Start up workers to handle requests from other nodes for the data on this node
40 for i := 0; i < TaskWorkerCount; i++ {
41 + i := i
42 px.Go(func(px process.Process) {
40 - bs.taskWorker(ctx)
43 + bs.taskWorker(ctx, i)
44 })
45 }
46
@@ -55,15 +58,18 @@ func (bs *Bitswap) startWorkers(px process.Process, ctx context.Context) {
58 // consider increasing number if providing blocks bottlenecks
59 // file transfers
60 for i := 0; i < provideWorkers; i++ {
61 + i := i
62 px.Go(func(px process.Process) {
59 - bs.provideWorker(ctx)
63 + bs.provideWorker(ctx, i)
64 })
65 }
66 }
67
64 -func (bs *Bitswap) taskWorker(ctx context.Context) {
68 +func (bs *Bitswap) taskWorker(ctx context.Context, id int) {
69 + idmap := eventlog.LoggableMap{"ID": id}
70 defer log.Info("bitswap task worker shutting down...")
71 for {
72 + log.Event(ctx, "Bitswap.TaskWorker.Loop", idmap)
73 select {
74 case nextEnvelope := <-bs.engine.Outbox():
75 select {
@@ -71,6 +77,7 @@ func (bs *Bitswap) taskWorker(ctx context.Context) {
77 if !ok {
78 continue
79 }
80 + log.Event(ctx, "Bitswap.TaskWorker.Work", eventlog.LoggableMap{"ID": id, "Target": envelope.Peer.Pretty(), "Block": envelope.Block.Multihash.B58String()})
81
82 bs.wm.SendBlock(ctx, envelope)
83 case <-ctx.Done():
@@ -82,10 +89,13 @@ func (bs *Bitswap) taskWorker(ctx context.Context) {
89 }
90 }
91
85 -func (bs *Bitswap) provideWorker(ctx context.Context) {
92 +func (bs *Bitswap) provideWorker(ctx context.Context, id int) {
93 + idmap := eventlog.LoggableMap{"ID": id}
94 for {
95 + log.Event(ctx, "Bitswap.ProvideWorker.Loop", idmap)
96 select {
97 case k, ok := <-bs.provideKeys:
98 + log.Event(ctx, "Bitswap.ProvideWorker.Work", idmap, &k)
99 if !ok {
100 log.Debug("provideKeys channel closed")
101 return
@@ -139,6 +149,7 @@ func (bs *Bitswap) providerConnector(parent context.Context) {
149 defer log.Info("bitswap client worker shutting down...")
150
151 for {
152 + log.Event(parent, "Bitswap.ProviderConnector.Loop")
153 select {
154 case req := <-bs.findKeys:
155 keys := req.keys
@@ -146,6 +157,7 @@ func (bs *Bitswap) providerConnector(parent context.Context) {
157 log.Warning("Received batch request for zero blocks")
158 continue
159 }
160 + log.Event(parent, "Bitswap.ProviderConnector.Work", eventlog.LoggableMap{"Keys": keys})
161
162 // NB: Optimization. Assumes that providers of key[0] are likely to
163 // be able to provide for all keys. This currently holds true in most
@@ -174,6 +186,7 @@ func (bs *Bitswap) rebroadcastWorker(parent context.Context) {
186 defer tick.Stop()
187
188 for {
189 + log.Event(ctx, "Bitswap.Rebroadcast.idle")
190 select {
191 case <-tick.C:
192 n := bs.wm.wl.Len()
@@ -181,6 +194,7 @@ func (bs *Bitswap) rebroadcastWorker(parent context.Context) {
194 log.Debug(n, "keys in bitswap wantlist")
195 }
196 case <-broadcastSignal.C: // resend unfulfilled wantlist keys
197 + log.Event(ctx, "Bitswap.Rebroadcast.active")
198 entries := bs.wm.wl.Entries()
199 if len(entries) > 0 {
200 bs.connectToProviders(ctx, entries)