@cryptotaxi247 / kubo / commits / bae8b9f4c

starting to move important events over to EventBegin/Done

Jeromy committed Jan 15, 2015 at 04:17 UTC bae8b9f4c052aa8c556e89c1b476f299e07b1091
4 files changed +23 -8
exchange/bitswap/bitswap.go
+2 -2
@@ -120,12 +120,12 @@ func (bs *bitswap) GetBlock(parent context.Context, k u.Key) (*blocks.Block, err
120 ctx, cancelFunc := context.WithCancel(parent)
121
122 ctx = eventlog.ContextWithLoggable(ctx, eventlog.Uuid("GetBlockRequest"))
123 - log.Event(ctx, "GetBlockRequestBegin", &k)
123 + e := log.EventBegin(ctx, "GetBlockRequest", &k)
124 log.Debugf("GetBlockRequestBegin")
125
126 defer func() {
127 cancelFunc()
128 - log.Event(ctx, "GetBlockRequestEnd", &k)
128 + e.Done()
129 log.Debugf("GetBlockRequestEnd")
130 }()
131
p2p/crypto/secio/protocol.go
+3 -1
@@ -81,7 +81,9 @@ func (s *secureSession) handshake(ctx context.Context, insecure io.ReadWriter) e
81 }
82
83 log.Debugf("handshake: %s <--start--> %s", s.localPeer, s.remotePeer)
84 - log.Event(ctx, "secureHandshakeStart", s.localPeer)
84 + e := log.EventBegin(ctx, "secureHandshake", s.localPeer)
85 + defer e.Done()
86 +
87 s.local.permanentPubKey = s.localKey.GetPublic()
88 myPubKeyBytes, err := s.local.permanentPubKey.Bytes()
89 if err != nil {
routing/dht/dht.go
+8
@@ -196,6 +196,8 @@ func (dht *IpfsDHT) getValueOrPeers(ctx context.Context, p peer.ID,
196 // getValueSingle simply performs the get value RPC with the given parameters
197 func (dht *IpfsDHT) getValueSingle(ctx context.Context, p peer.ID,
198 key u.Key) (*pb.Message, error) {
199 + e := log.EventBegin(ctx, "getValueSingle", p, &key)
200 + defer e.Done()
201
202 pmes := pb.NewMessage(pb.Message_GET_VALUE, string(key), 0)
203 return dht.sendRequest(ctx, p, pmes)
@@ -265,11 +267,17 @@ func (dht *IpfsDHT) FindLocal(id peer.ID) peer.PeerInfo {
267
268 // findPeerSingle asks peer 'p' if they know where the peer with id 'id' is
269 func (dht *IpfsDHT) findPeerSingle(ctx context.Context, p peer.ID, id peer.ID) (*pb.Message, error) {
270 + e := log.EventBegin(ctx, "findPeerSingle", p, id)
271 + defer e.Done()
272 +
273 pmes := pb.NewMessage(pb.Message_FIND_NODE, string(id), 0)
274 return dht.sendRequest(ctx, p, pmes)
275 }
276
277 func (dht *IpfsDHT) findProvidersSingle(ctx context.Context, p peer.ID, key u.Key) (*pb.Message, error) {
278 + e := log.EventBegin(ctx, "findProvidersSingle", p, &key)
279 + defer e.Done()
280 +
281 pmes := pb.NewMessage(pb.Message_GET_PROVIDERS, string(key), 0)
282 return dht.sendRequest(ctx, p, pmes)
283 }
routing/dht/routing.go
+10 -5
@@ -122,10 +122,12 @@ func (dht *IpfsDHT) GetValue(ctx context.Context, key u.Key) ([]byte, error) {
122 // Provide makes this node announce that it can provide a value for the given key
123 func (dht *IpfsDHT) Provide(ctx context.Context, key u.Key) error {
124 log := dht.log().Prefix("Provide(%s)", key)
125 +
126 log.Debugf("start", key)
126 - log.Event(ctx, "provideBegin", &key)
127 defer log.Debugf("end", key)
128 - defer log.Event(ctx, "provideEnd", &key)
128 +
129 + e := log.EventBegin(ctx, "provide", &key)
130 + defer e.Done()
131
132 // add self locally
133 dht.providers.AddProvider(key, dht.self)
@@ -163,6 +165,7 @@ func (dht *IpfsDHT) FindProviders(ctx context.Context, key u.Key) ([]peer.PeerIn
165 // Kademlia 'node lookup' operation. Returns a channel of the K closest peers
166 // to the given key
167 func (dht *IpfsDHT) getClosestPeers(ctx context.Context, key u.Key) (<-chan peer.ID, error) {
168 + e := log.EventBegin(ctx, "getClosestPeers", &key)
169 tablepeers := dht.routingTable.NearestPeers(kb.ConvertKey(key), AlphaValue)
170 if len(tablepeers) == 0 {
171 return nil, errors.Wrap(kb.ErrLookupFailure)
@@ -204,6 +207,7 @@ func (dht *IpfsDHT) getClosestPeers(ctx context.Context, key u.Key) (<-chan peer
207
208 go func() {
209 defer close(out)
210 + defer e.Done()
211 // run it!
212 _, err := query.Run(ctx, tablepeers)
213 if err != nil {
@@ -242,10 +246,9 @@ func (dht *IpfsDHT) FindProvidersAsync(ctx context.Context, key u.Key, count int
246 func (dht *IpfsDHT) findProvidersAsyncRoutine(ctx context.Context, key u.Key, count int, peerOut chan peer.PeerInfo) {
247 log := dht.log().Prefix("FindProviders(%s)", key)
248
249 + e := log.EventBegin(ctx, "findProvidersAsync", &key)
250 + defer e.Done()
251 defer close(peerOut)
246 - defer log.Event(ctx, "findProviders end", &key)
247 - log.Debug("begin")
248 - defer log.Debug("begin")
252
253 ps := pset.NewLimited(count)
254 provs := dht.providers.GetProviders(ctx, key)
@@ -314,6 +317,8 @@ func (dht *IpfsDHT) findProvidersAsyncRoutine(ctx context.Context, key u.Key, co
317
318 // FindPeer searches for a peer with given ID.
319 func (dht *IpfsDHT) FindPeer(ctx context.Context, id peer.ID) (peer.PeerInfo, error) {
320 + e := log.EventBegin(ctx, "FindPeer", id)
321 + defer e.Done()
322
323 // Check if were already connected to them
324 if pi := dht.FindLocal(id); pi.ID != "" {