@cryptotaxi247 / kubo / commits / 27dc9594b

refactor(bitswap) bitswap.Network now abstracts ipfs.Network + ipfs.Routing

@jbenet @whyrusleeping the next commit will change bitswap.Network.FindProviders to only deal with IDs

Brian Tiger Chow committed Dec 23, 2014 at 08:16 UTC 27dc9594ba431b9bf3675b7652f0497ff371414a
10 files changed +71 -55
blockservice/mock.go
+2 -3
@@ -11,9 +11,8 @@ import (
11
12 // Mocks returns |n| connected mock Blockservices
13 func Mocks(t *testing.T, n int) []*BlockService {
14 - net := tn.VirtualNetwork(delay.Fixed(0))
15 - rs := mockrouting.NewServer()
16 - sg := bitswap.NewSessionGenerator(net, rs)
14 + net := tn.VirtualNetwork(mockrouting.NewServer(), delay.Fixed(0))
15 + sg := bitswap.NewSessionGenerator(net)
16
17 instances := sg.Instances(n)
18
core/core.go
+2 -2
@@ -150,9 +150,9 @@ func NewIpfsNode(ctx context.Context, cfg *config.Config, online bool) (n *IpfsN
150
151 // setup exchange service
152 const alwaysSendToPeer = true // use YesManStrategy
153 - bitswapNetwork := bsnet.NewFromIpfsNetwork(n.Network)
153 + bitswapNetwork := bsnet.NewFromIpfsNetwork(n.Network, n.Routing)
154
155 - n.Exchange = bitswap.New(ctx, n.Identity, bitswapNetwork, n.Routing, blockstore, alwaysSendToPeer)
155 + n.Exchange = bitswap.New(ctx, n.Identity, bitswapNetwork, blockstore, alwaysSendToPeer)
156
157 // TODO consider moving connection supervision into the Network. We've
158 // discussed improvements to this Node constructor. One improvement
epictest/addcat_test.go
+6 -5
@@ -87,11 +87,12 @@ func RandomBytes(n int64) []byte {
87 func AddCatBytes(data []byte, conf Config) error {
88
89 sessionGenerator := bitswap.NewSessionGenerator(
90 - tn.VirtualNetwork(delay.Fixed(conf.NetworkLatency)), // TODO rename VirtualNetwork
91 - mockrouting.NewServerWithDelay(mockrouting.DelayConfig{
92 - Query: delay.Fixed(conf.RoutingLatency),
93 - ValueVisibility: delay.Fixed(conf.RoutingLatency),
94 - }),
90 + tn.VirtualNetwork(
91 + mockrouting.NewServerWithDelay(mockrouting.DelayConfig{
92 + Query: delay.Fixed(conf.RoutingLatency),
93 + ValueVisibility: delay.Fixed(conf.RoutingLatency),
94 + }),
95 + delay.Fixed(conf.NetworkLatency)), // TODO rename VirtualNetwork
96 )
97 defer sessionGenerator.Close()
98
exchange/bitswap/bitswap.go
+4 -8
@@ -46,7 +46,7 @@ var (
46 // BitSwapNetwork. This function registers the returned instance as the network
47 // delegate.
48 // Runs until context is cancelled.
49 -func New(parent context.Context, p peer.ID, network bsnet.BitSwapNetwork, routing bsnet.Routing,
49 +func New(parent context.Context, p peer.ID, network bsnet.BitSwapNetwork,
50 bstore blockstore.Blockstore, nice bool) exchange.Interface {
51
52 ctx, cancelFunc := context.WithCancel(parent)
@@ -63,7 +63,6 @@ func New(parent context.Context, p peer.ID, network bsnet.BitSwapNetwork, routin
63 cancelFunc: cancelFunc,
64 notifications: notif,
65 engine: decision.NewEngine(ctx, bstore),
66 - routing: routing,
66 network: network,
67 wantlist: wantlist.NewThreadSafe(),
68 batchRequests: make(chan []u.Key, sizeBatchRequestChan),
@@ -85,9 +84,6 @@ type bitswap struct {
84 // NB: ensure threadsafety
85 blockstore blockstore.Blockstore
86
88 - // routing interface for communication
89 - routing bsnet.Routing
90 -
87 notifications notifications.PubSub
88
89 // Requests for a set of related blocks
@@ -165,7 +161,7 @@ func (bs *bitswap) HasBlock(ctx context.Context, blk *blocks.Block) error {
161 }
162 bs.wantlist.Remove(blk.Key())
163 bs.notifications.Publish(blk)
168 - return bs.routing.Provide(ctx, blk.Key())
164 + return bs.network.Provide(ctx, blk.Key())
165 }
166
167 func (bs *bitswap) sendWantListTo(ctx context.Context, peers <-chan peer.PeerInfo) error {
@@ -212,7 +208,7 @@ func (bs *bitswap) sendWantlistToProviders(ctx context.Context, wantlist *wantli
208 go func(k u.Key) {
209 defer wg.Done()
210 child, _ := context.WithTimeout(ctx, providerRequestTimeout)
215 - providers := bs.routing.FindProvidersAsync(child, k, maxProvidersPerRequest)
211 + providers := bs.network.FindProvidersAsync(child, k, maxProvidersPerRequest)
212 for prov := range providers {
213 bs.network.Peerstore().AddAddresses(prov.ID, prov.Addrs)
214 if set.TryAdd(prov.ID) { //Do once per peer
@@ -265,7 +261,7 @@ func (bs *bitswap) clientWorker(parent context.Context) {
261 // it. Later, this assumption may not hold as true if we implement
262 // newer bitswap strategies.
263 child, _ := context.WithTimeout(ctx, providerRequestTimeout)
268 - providers := bs.routing.FindProvidersAsync(child, ks[0], maxProvidersPerRequest)
264 + providers := bs.network.FindProvidersAsync(child, ks[0], maxProvidersPerRequest)
265 err := bs.sendWantListTo(ctx, providers)
266 if err != nil {
267 log.Errorf("error sending wantlist: %s", err)
exchange/bitswap/bitswap_test.go
+16 -23
@@ -24,9 +24,8 @@ const kNetworkDelay = 0 * time.Millisecond
24 func TestClose(t *testing.T) {
25 // TODO
26 t.Skip("TODO Bitswap's Close implementation is a WIP")
27 - vnet := tn.VirtualNetwork(delay.Fixed(kNetworkDelay))
28 - rout := mockrouting.NewServer()
29 - sesgen := NewSessionGenerator(vnet, rout)
27 + vnet := tn.VirtualNetwork(mockrouting.NewServer(), delay.Fixed(kNetworkDelay))
28 + sesgen := NewSessionGenerator(vnet)
29 defer sesgen.Close()
30 bgen := blocksutil.NewBlockGenerator()
31
@@ -39,9 +38,8 @@ func TestClose(t *testing.T) {
38
39 func TestGetBlockTimeout(t *testing.T) {
40
42 - net := tn.VirtualNetwork(delay.Fixed(kNetworkDelay))
43 - rs := mockrouting.NewServer()
44 - g := NewSessionGenerator(net, rs)
41 + net := tn.VirtualNetwork(mockrouting.NewServer(), delay.Fixed(kNetworkDelay))
42 + g := NewSessionGenerator(net)
43 defer g.Close()
44
45 self := g.Next()
@@ -55,11 +53,11 @@ func TestGetBlockTimeout(t *testing.T) {
53 }
54 }
55
58 -func TestProviderForKeyButNetworkCannotFind(t *testing.T) {
56 +func TestProviderForKeyButNetworkCannotFind(t *testing.T) { // TODO revisit this
57
60 - net := tn.VirtualNetwork(delay.Fixed(kNetworkDelay))
58 rs := mockrouting.NewServer()
62 - g := NewSessionGenerator(net, rs)
59 + net := tn.VirtualNetwork(rs, delay.Fixed(kNetworkDelay))
60 + g := NewSessionGenerator(net)
61 defer g.Close()
62
63 block := blocks.NewBlock([]byte("block"))
@@ -81,10 +79,9 @@ func TestProviderForKeyButNetworkCannotFind(t *testing.T) {
79
80 func TestGetBlockFromPeerAfterPeerAnnounces(t *testing.T) {
81
84 - net := tn.VirtualNetwork(delay.Fixed(kNetworkDelay))
85 - rs := mockrouting.NewServer()
82 + net := tn.VirtualNetwork(mockrouting.NewServer(), delay.Fixed(kNetworkDelay))
83 block := blocks.NewBlock([]byte("block"))
87 - g := NewSessionGenerator(net, rs)
84 + g := NewSessionGenerator(net)
85 defer g.Close()
86
87 hasBlock := g.Next()
@@ -136,9 +133,8 @@ func PerformDistributionTest(t *testing.T, numInstances, numBlocks int) {
133 if testing.Short() {
134 t.SkipNow()
135 }
139 - net := tn.VirtualNetwork(delay.Fixed(kNetworkDelay))
140 - rs := mockrouting.NewServer()
141 - sg := NewSessionGenerator(net, rs)
136 + net := tn.VirtualNetwork(mockrouting.NewServer(), delay.Fixed(kNetworkDelay))
137 + sg := NewSessionGenerator(net)
138 defer sg.Close()
139 bg := blocksutil.NewBlockGenerator()
140
@@ -152,10 +148,9 @@ func PerformDistributionTest(t *testing.T, numInstances, numBlocks int) {
148 var blkeys []u.Key
149 first := instances[0]
150 for _, b := range blocks {
155 - first.Blockstore().Put(b)
151 + first.Blockstore().Put(b) // TODO remove. don't need to do this. bitswap owns block
152 blkeys = append(blkeys, b.Key())
153 first.Exchange.HasBlock(context.Background(), b)
158 - rs.Client(peer.PeerInfo{ID: first.Peer}).Provide(context.Background(), b.Key())
154 }
155
156 t.Log("Distribute!")
@@ -202,9 +197,8 @@ func TestSendToWantingPeer(t *testing.T) {
197 t.SkipNow()
198 }
199
205 - net := tn.VirtualNetwork(delay.Fixed(kNetworkDelay))
206 - rs := mockrouting.NewServer()
207 - sg := NewSessionGenerator(net, rs)
200 + net := tn.VirtualNetwork(mockrouting.NewServer(), delay.Fixed(kNetworkDelay))
201 + sg := NewSessionGenerator(net)
202 defer sg.Close()
203 bg := blocksutil.NewBlockGenerator()
204
@@ -248,9 +242,8 @@ func TestSendToWantingPeer(t *testing.T) {
242 }
243
244 func TestBasicBitswap(t *testing.T) {
251 - net := tn.VirtualNetwork(delay.Fixed(kNetworkDelay))
252 - rs := mockrouting.NewServer()
253 - sg := NewSessionGenerator(net, rs)
245 + net := tn.VirtualNetwork(mockrouting.NewServer(), delay.Fixed(kNetworkDelay))
246 + sg := NewSessionGenerator(net)
247 bg := blocksutil.NewBlockGenerator()
248
249 t.Log("Test a few nodes trying to get one file with a lot of blocks")
exchange/bitswap/network/interface.go
+2
@@ -31,6 +31,8 @@ type BitSwapNetwork interface {
31 // SetDelegate registers the Reciver to handle messages received from the
32 // network.
33 SetDelegate(Receiver)
34 +
35 + Routing
36 }
37
38 // Implement Receiver to receive messages from the BitSwapNetwork
exchange/bitswap/network/ipfs_impl.go
+13 -1
@@ -13,9 +13,10 @@ var log = util.Logger("bitswap_network")
13
14 // NewFromIpfsNetwork returns a BitSwapNetwork supported by underlying IPFS
15 // Dialer & Service
16 -func NewFromIpfsNetwork(n inet.Network) BitSwapNetwork {
16 +func NewFromIpfsNetwork(n inet.Network, r Routing) BitSwapNetwork {
17 bitswapNetwork := impl{
18 network: n,
19 + routing: r,
20 }
21 n.SetHandler(inet.ProtocolBitswap, bitswapNetwork.handleNewStream)
22 return &bitswapNetwork
@@ -25,6 +26,7 @@ func NewFromIpfsNetwork(n inet.Network) BitSwapNetwork {
26 // NetMessage objects, into the bitswap network interface.
27 type impl struct {
28 network inet.Network
29 + routing Routing
30
31 // inbound messages from the network are forwarded to the receiver
32 receiver Receiver
@@ -74,6 +76,16 @@ func (bsnet *impl) Peerstore() peer.Peerstore {
76 return bsnet.Peerstore()
77 }
78
79 +// FindProvidersAsync returns a channel of providers for the given key
80 +func (bsnet *impl) FindProvidersAsync(ctx context.Context, k util.Key, max int) <-chan peer.PeerInfo { // TODO change to return ID
81 + return bsnet.routing.FindProvidersAsync(ctx, k, max)
82 +}
83 +
84 +// Provide provides the key to the network
85 +func (bsnet *impl) Provide(ctx context.Context, k util.Key) error {
86 + return bsnet.routing.Provide(ctx, k)
87 +}
88 +
89 // handleNewStream receives a new stream from the network.
90 func (bsnet *impl) handleNewStream(s inet.Stream) {
91
exchange/bitswap/testnet/network.go
+19 -3
@@ -5,10 +5,12 @@ import (
5 "fmt"
6
7 context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
8 + "github.com/jbenet/go-ipfs/routing/mock"
9
10 bsmsg "github.com/jbenet/go-ipfs/exchange/bitswap/message"
11 bsnet "github.com/jbenet/go-ipfs/exchange/bitswap/network"
12 peer "github.com/jbenet/go-ipfs/peer"
13 + "github.com/jbenet/go-ipfs/util"
14 delay "github.com/jbenet/go-ipfs/util/delay"
15 )
16
@@ -33,16 +35,18 @@ type Network interface {
35
36 // network impl
37
36 -func VirtualNetwork(d delay.D) Network {
38 +func VirtualNetwork(rs mockrouting.Server, d delay.D) Network {
39 return &network{
40 clients: make(map[peer.ID]bsnet.Receiver),
41 delay: d,
42 + routingserver: rs,
43 }
44 }
45
46 type network struct {
44 - clients map[peer.ID]bsnet.Receiver
45 - delay delay.D
47 + clients map[peer.ID]bsnet.Receiver
48 + routingserver mockrouting.Server
49 + delay delay.D
50 }
51
52 func (n *network) Adapter(p peer.ID) bsnet.BitSwapNetwork {
@@ -50,6 +54,7 @@ func (n *network) Adapter(p peer.ID) bsnet.BitSwapNetwork {
54 local: p,
55 network: n,
56 peerstore: peer.NewPeerstore(),
57 + routing: n.routingserver.Client(peer.PeerInfo{ID: p}),
58 }
59 n.clients[p] = client
60 return client
@@ -151,6 +156,7 @@ type networkClient struct {
156 bsnet.Receiver
157 network Network
158 peerstore peer.Peerstore
159 + routing bsnet.Routing
160 }
161
162 func (nc *networkClient) SendMessage(
@@ -167,6 +173,16 @@ func (nc *networkClient) SendRequest(
173 return nc.network.SendRequest(ctx, nc.local, to, message)
174 }
175
176 +// FindProvidersAsync returns a channel of providers for the given key
177 +func (nc *networkClient) FindProvidersAsync(ctx context.Context, k util.Key, max int) <-chan peer.PeerInfo { // TODO change to return ID
178 + return nc.routing.FindProvidersAsync(ctx, k, max)
179 +}
180 +
181 +// Provide provides the key to the network
182 +func (nc *networkClient) Provide(ctx context.Context, k util.Key) error {
183 + return nc.routing.Provide(ctx, k)
184 +}
185 +
186 func (nc *networkClient) DialPeer(ctx context.Context, p peer.ID) error {
187 // no need to do anything because dialing isn't a thing in this test net.
188 if !nc.network.HasPeer(p) {
exchange/bitswap/testnet/network_test.go
+3 -2
@@ -11,10 +11,11 @@ import (
11 bsnet "github.com/jbenet/go-ipfs/exchange/bitswap/network"
12 peer "github.com/jbenet/go-ipfs/peer"
13 delay "github.com/jbenet/go-ipfs/util/delay"
14 + mockrouting "github.com/jbenet/go-ipfs/routing/mock"
15 )
16
17 func TestSendRequestToCooperativePeer(t *testing.T) {
17 - net := VirtualNetwork(delay.Fixed(0))
18 + net := VirtualNetwork(mockrouting.NewServer(),delay.Fixed(0))
19
20 idOfRecipient := peer.ID("recipient")
21
@@ -65,7 +66,7 @@ func TestSendRequestToCooperativePeer(t *testing.T) {
66 }
67
68 func TestSendMessageAsyncButWaitForResponse(t *testing.T) {
68 - net := VirtualNetwork(delay.Fixed(0))
69 + net := VirtualNetwork(mockrouting.NewServer(), delay.Fixed(0))
70 idOfResponder := peer.ID("responder")
71 waiter := net.Adapter(peer.ID("waiter"))
72 responder := net.Adapter(idOfResponder)
exchange/bitswap/testutils.go
+4 -8
@@ -10,18 +10,16 @@ import (
10 exchange "github.com/jbenet/go-ipfs/exchange"
11 tn "github.com/jbenet/go-ipfs/exchange/bitswap/testnet"
12 peer "github.com/jbenet/go-ipfs/peer"
13 - mockrouting "github.com/jbenet/go-ipfs/routing/mock"
13 datastore2 "github.com/jbenet/go-ipfs/util/datastore2"
14 delay "github.com/jbenet/go-ipfs/util/delay"
15 )
16
17 func NewSessionGenerator(
19 - net tn.Network, rs mockrouting.Server) SessionGenerator {
18 + net tn.Network) SessionGenerator {
19 ctx, cancel := context.WithCancel(context.TODO())
20 return SessionGenerator{
21 ps: peer.NewPeerstore(),
22 net: net,
24 - rs: rs,
23 seq: 0,
24 ctx: ctx, // TODO take ctx as param to Next, Instances
25 cancel: cancel,
@@ -31,7 +29,6 @@ func NewSessionGenerator(
29 type SessionGenerator struct {
30 seq int
31 net tn.Network
34 - rs mockrouting.Server
32 ps peer.Peerstore
33 ctx context.Context
34 cancel context.CancelFunc
@@ -44,7 +41,7 @@ func (g *SessionGenerator) Close() error {
41
42 func (g *SessionGenerator) Next() Instance {
43 g.seq++
47 - return session(g.ctx, g.net, g.rs, g.ps, peer.ID(g.seq))
44 + return session(g.ctx, g.net, g.ps, peer.ID(g.seq))
45 }
46
47 func (g *SessionGenerator) Instances(n int) []Instance {
@@ -77,10 +74,9 @@ func (i *Instance) SetBlockstoreLatency(t time.Duration) time.Duration {
74 // NB: It's easy make mistakes by providing the same peer ID to two different
75 // sessions. To safeguard, use the SessionGenerator to generate sessions. It's
76 // just a much better idea.
80 -func session(ctx context.Context, net tn.Network, rs mockrouting.Server, ps peer.Peerstore, p peer.ID) Instance {
77 +func session(ctx context.Context, net tn.Network, ps peer.Peerstore, p peer.ID) Instance {
78
79 adapter := net.Adapter(p)
83 - htc := rs.Client(peer.PeerInfo{ID: p})
80
81 bsdelay := delay.Fixed(0)
82 const kWriteCacheElems = 100
@@ -92,7 +88,7 @@ func session(ctx context.Context, net tn.Network, rs mockrouting.Server, ps peer
88
89 const alwaysSendToPeer = true
90
95 - bs := New(ctx, p, adapter, htc, bstore, alwaysSendToPeer)
91 + bs := New(ctx, p, adapter, bstore, alwaysSendToPeer)
92
93 return Instance{
94 Peer: p,