@cryptotaxi247 / kubo / commits / 3eafbea26

tag peers associated with a bitswap session

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

Jeromy committed Oct 14, 2017 at 08:33 UTC 3eafbea262589c9cb9c592143bd6a46859182104
4 files changed +28 -2
exchange/bitswap/network/interface.go
+4
@@ -4,8 +4,10 @@ import (
4 "context"
5
6 bsmsg "github.com/ipfs/go-ipfs/exchange/bitswap/message"
7 +
8 cid "gx/ipfs/QmNp85zy9RLrQ5oQD4hPyS39ezrrXpcaa7R4Y9kxdWQLLQ/go-cid"
9 peer "gx/ipfs/QmXYjuNuxVzXKJCfWasQk1RqkhVLDM9jtUKhqc2WPQmFSB/go-libp2p-peer"
10 + ifconnmgr "gx/ipfs/QmYkCrTwivapqdB3JbwvwvxymseahVkcm46ThRMAA24zCr/go-libp2p-interface-connmgr"
11 protocol "gx/ipfs/QmZNkThpqfVXs9GNbexPrfBbXSLNYeKrE7jwFM2oqHbyqN/go-libp2p-protocol"
12 )
13
@@ -34,6 +36,8 @@ type BitSwapNetwork interface {
36
37 NewMessageSender(context.Context, peer.ID) (MessageSender, error)
38
39 + ConnectionManager() ifconnmgr.ConnManager
40 +
41 Routing
42 }
43
exchange/bitswap/network/ipfs_impl.go
+5
@@ -15,6 +15,7 @@ import (
15 logging "gx/ipfs/QmSpJByNKFX1sCsHBEp3R73FL4NF6FnQTEGyNAXHm2GS52/go-log"
16 ma "gx/ipfs/QmXY77cVe7rVRQXZZQRioukUM7aRW3BTcAgJe12MCtb3Ji/go-multiaddr"
17 peer "gx/ipfs/QmXYjuNuxVzXKJCfWasQk1RqkhVLDM9jtUKhqc2WPQmFSB/go-libp2p-peer"
18 + ifconnmgr "gx/ipfs/QmYkCrTwivapqdB3JbwvwvxymseahVkcm46ThRMAA24zCr/go-libp2p-interface-connmgr"
19 ggio "gx/ipfs/QmZ4Qi3GaRbjcx28Sme5eMH7RQjGkt8wHxt2a65oLaeFEV/gogo-protobuf/io"
20 host "gx/ipfs/Qmc1XhrFEiSeBNn3mpfg6gEuYCt5im2gYmNVmncsvmpeAk/go-libp2p-host"
21 )
@@ -212,6 +213,10 @@ func (bsnet *impl) handleNewStream(s inet.Stream) {
213 }
214 }
215
216 +func (bsnet *impl) ConnectionManager() ifconnmgr.ConnManager {
217 + return bsnet.host.ConnManager()
218 +}
219 +
220 type netNotifiee impl
221
222 func (nn *netNotifiee) impl() *impl {
exchange/bitswap/session.go
+13 -1
@@ -2,6 +2,7 @@ package bitswap
2
3 import (
4 "context"
5 + "fmt"
6 "time"
7
8 notifications "github.com/ipfs/go-ipfs/exchange/bitswap/notifications"
@@ -44,7 +45,8 @@ type Session struct {
45
46 uuid logging.Loggable
47
47 - id uint64
48 + id uint64
49 + tag string
50 }
51
52 // NewSession creates a new bitswap session whose lifetime is bounded by the
@@ -66,6 +68,8 @@ func (bs *Bitswap) NewSession(ctx context.Context) *Session {
68 id: bs.getNextSessionID(),
69 }
70
71 + s.tag = fmt.Sprint("bs-ses-", s.id)
72 +
73 cache, _ := lru.New(2048)
74 s.interest = cache
75
@@ -139,6 +143,9 @@ func (s *Session) addActivePeer(p peer.ID) {
143 if _, ok := s.activePeers[p]; !ok {
144 s.activePeers[p] = struct{}{}
145 s.activePeersArr = append(s.activePeersArr, p)
146 +
147 + cmgr := s.bs.network.ConnectionManager()
148 + cmgr.TagPeer(p, s.tag, 10)
149 }
150 }
151
@@ -216,6 +223,11 @@ func (s *Session) run(ctx context.Context) {
223 case <-ctx.Done():
224 s.tick.Stop()
225 s.bs.removeSession(s)
226 +
227 + cmgr := s.bs.network.ConnectionManager()
228 + for _, p := range s.activePeersArr {
229 + cmgr.UntagPeer(p, s.tag)
230 + }
231 return
232 }
233 }
exchange/bitswap/testnet/virtual.go
+6 -1
@@ -8,12 +8,13 @@ import (
8 bsnet "github.com/ipfs/go-ipfs/exchange/bitswap/network"
9 mockrouting "github.com/ipfs/go-ipfs/routing/mock"
10 delay "github.com/ipfs/go-ipfs/thirdparty/delay"
11 - testutil "gx/ipfs/QmWRCn8vruNAzHx8i6SAXinuheRitKEGu8c7m26stKvsYx/go-testutil"
11
12 cid "gx/ipfs/QmNp85zy9RLrQ5oQD4hPyS39ezrrXpcaa7R4Y9kxdWQLLQ/go-cid"
13 routing "gx/ipfs/QmPR2JzfKd9poHx9XBhzoFeBBC31ZM3W5iUPKJZWyaoZZm/go-libp2p-routing"
14 logging "gx/ipfs/QmSpJByNKFX1sCsHBEp3R73FL4NF6FnQTEGyNAXHm2GS52/go-log"
15 + testutil "gx/ipfs/QmWRCn8vruNAzHx8i6SAXinuheRitKEGu8c7m26stKvsYx/go-testutil"
16 peer "gx/ipfs/QmXYjuNuxVzXKJCfWasQk1RqkhVLDM9jtUKhqc2WPQmFSB/go-libp2p-peer"
17 + ifconnmgr "gx/ipfs/QmYkCrTwivapqdB3JbwvwvxymseahVkcm46ThRMAA24zCr/go-libp2p-interface-connmgr"
18 )
19
20 var log = logging.Logger("bstestnet")
@@ -118,6 +119,10 @@ func (nc *networkClient) FindProvidersAsync(ctx context.Context, k *cid.Cid, max
119 return out
120 }
121
122 +func (nc *networkClient) ConnectionManager() ifconnmgr.ConnManager {
123 + return &ifconnmgr.NullConnMgr{}
124 +}
125 +
126 type messagePasser struct {
127 net *network
128 target peer.ID