@cryptotaxi247 / kubo / commits / c2b497e31

switch over to using sendMessage vs sendRequest

Jeromy committed Dec 1, 2014 at 21:38 UTC c2b497e3157222187a4f65b4fec03ce6f26bf2c4
3 files changed +6 -22
exchange/bitswap/bitswap.go
+3 -7
@@ -151,6 +151,7 @@ func (bs *bitswap) sendWantListTo(ctx context.Context, peers <-chan peer.Peer) e
151 message.AddWanted(wanted)
152 }
153 for peerToQuery := range peers {
154 + log.Debug("sending query to: %s", peerToQuery)
155 log.Event(ctx, "PeerToQuery", peerToQuery)
156 go func(p peer.Peer) {
157
@@ -161,20 +162,15 @@ func (bs *bitswap) sendWantListTo(ctx context.Context, peers <-chan peer.Peer) e
162 return
163 }
164
164 - response, err := bs.sender.SendRequest(ctx, p, message)
165 + err = bs.sender.SendMessage(ctx, p, message)
166 if err != nil {
166 - log.Errorf("Error sender.SendRequest(%s) = %s", p, err)
167 + log.Errorf("Error sender.SendMessage(%s) = %s", p, err)
168 return
169 }
170 // FIXME ensure accounting is handled correctly when
171 // communication fails. May require slightly different API to
172 // get better guarantees. May need shared sequence numbers.
173 bs.strategy.MessageSent(p, message)
173 -
174 - if response == nil {
175 - return
176 - }
177 - bs.ReceiveMessage(ctx, p, response)
174 }(peerToQuery)
175 }
176 return nil
exchange/bitswap/network/ipfs_impl.go
+2 -15
@@ -48,21 +48,8 @@ func (bsnet *impl) HandleMessage(
48 return nil
49 }
50
51 - p, bsmsg := bsnet.receiver.ReceiveMessage(ctx, incoming.Peer(), received)
52 -
53 - // TODO(brian): put this in a helper function
54 - if bsmsg == nil || p == nil {
55 - return nil
56 - }
57 -
58 - outgoing, err := bsmsg.ToNet(p)
59 - if err != nil {
60 - go bsnet.receiver.ReceiveError(err)
61 - return nil
62 - }
63 -
64 - log.Debugf("Message size: %d", len(outgoing.Data()))
65 - return outgoing
51 + bsnet.receiver.ReceiveMessage(ctx, incoming.Peer(), received)
52 + return nil
53 }
54
55 func (bsnet *impl) DialPeer(ctx context.Context, p peer.Peer) error {
routing/dht/routing.go
+1
@@ -126,6 +126,7 @@ func (dht *IpfsDHT) Provide(ctx context.Context, key u.Key) error {
126 }
127
128 func (dht *IpfsDHT) FindProvidersAsync(ctx context.Context, key u.Key, count int) <-chan peer.Peer {
129 + log.Debug("Find Providers: %s", key)
130 peerOut := make(chan peer.Peer, count)
131 go func() {
132 ps := newPeerSet()