@cryptotaxi247 / kubo / commits / 5be35a83f

beginnings of a bitswap refactor

Jeromy committed Nov 18, 2014 at 21:31 UTC 5be35a83f18d06e31723646aa9cbf0f85f7813c2
6 files changed +117 -52
blockservice/blockservice.go
+1 -1
@@ -85,7 +85,7 @@ func (s *BlockService) GetBlock(ctx context.Context, k u.Key) (*blocks.Block, er
85 }, nil
86 } else if err == ds.ErrNotFound && s.Remote != nil {
87 log.Debug("Blockservice: Searching bitswap.")
88 - blk, err := s.Remote.Block(ctx, k)
88 + blk, err := s.Remote.GetBlock(ctx, k)
89 if err != nil {
90 return nil, err
91 }
exchange/bitswap/bitswap.go
+101 -38
@@ -42,8 +42,10 @@ func New(ctx context.Context, p peer.Peer,
42 routing: routing,
43 sender: network,
44 wantlist: u.NewKeySet(),
45 + blockReq: make(chan u.Key, 32),
46 }
47 network.SetDelegate(bs)
48 + go bs.run(ctx)
49
50 return bs
51 }
@@ -63,6 +65,8 @@ type bitswap struct {
65
66 notifications notifications.PubSub
67
68 + blockReq chan u.Key
69 +
70 // strategy listens to network traffic and makes decisions about how to
71 // interact with partners.
72 // TODO(brian): save the strategy's state to the datastore
@@ -75,7 +79,7 @@ type bitswap struct {
79 // deadline enforced by the context
80 //
81 // TODO ensure only one active request per key
78 -func (bs *bitswap) Block(parent context.Context, k u.Key) (*blocks.Block, error) {
82 +func (bs *bitswap) GetBlock(parent context.Context, k u.Key) (*blocks.Block, error) {
83 log.Debugf("Get Block %v", k)
84 now := time.Now()
85 defer func() {
@@ -88,42 +92,11 @@ func (bs *bitswap) Block(parent context.Context, k u.Key) (*blocks.Block, error)
92 bs.wantlist.Add(k)
93 promise := bs.notifications.Subscribe(ctx, k)
94
91 - const maxProviders = 20
92 - peersToQuery := bs.routing.FindProvidersAsync(ctx, k, maxProviders)
93 -
94 - go func() {
95 - message := bsmsg.New()
96 - for _, wanted := range bs.wantlist.Keys() {
97 - message.AddWanted(wanted)
98 - }
99 - for peerToQuery := range peersToQuery {
100 - log.Debugf("bitswap got peersToQuery: %s", peerToQuery)
101 - go func(p peer.Peer) {
102 -
103 - log.Debugf("bitswap dialing peer: %s", p)
104 - err := bs.sender.DialPeer(ctx, p)
105 - if err != nil {
106 - log.Errorf("Error sender.DialPeer(%s)", p)
107 - return
108 - }
109 -
110 - response, err := bs.sender.SendRequest(ctx, p, message)
111 - if err != nil {
112 - log.Errorf("Error sender.SendRequest(%s) = %s", p, err)
113 - return
114 - }
115 - // FIXME ensure accounting is handled correctly when
116 - // communication fails. May require slightly different API to
117 - // get better guarantees. May need shared sequence numbers.
118 - bs.strategy.MessageSent(p, message)
119 -
120 - if response == nil {
121 - return
122 - }
123 - bs.ReceiveMessage(ctx, p, response)
124 - }(peerToQuery)
125 - }
126 - }()
95 + select {
96 + case bs.blockReq <- k:
97 + case <-parent.Done():
98 + return nil, parent.Err()
99 + }
100
101 select {
102 case block := <-promise:
@@ -134,6 +107,96 @@ func (bs *bitswap) Block(parent context.Context, k u.Key) (*blocks.Block, error)
107 }
108 }
109
110 +func (bs *bitswap) GetBlocks(parent context.Context, ks []u.Key) (*blocks.Block, error) {
111 + // TODO: something smart
112 + return nil, nil
113 +}
114 +
115 +func (bs *bitswap) sendWantListTo(ctx context.Context, peers <-chan peer.Peer) error {
116 + message := bsmsg.New()
117 + for _, wanted := range bs.wantlist.Keys() {
118 + message.AddWanted(wanted)
119 + }
120 + for peerToQuery := range peers {
121 + log.Debugf("bitswap got peersToQuery: %s", peerToQuery)
122 + go func(p peer.Peer) {
123 +
124 + log.Debugf("bitswap dialing peer: %s", p)
125 + err := bs.sender.DialPeer(ctx, p)
126 + if err != nil {
127 + log.Errorf("Error sender.DialPeer(%s)", p)
128 + return
129 + }
130 +
131 + response, err := bs.sender.SendRequest(ctx, p, message)
132 + if err != nil {
133 + log.Errorf("Error sender.SendRequest(%s) = %s", p, err)
134 + return
135 + }
136 + // FIXME ensure accounting is handled correctly when
137 + // communication fails. May require slightly different API to
138 + // get better guarantees. May need shared sequence numbers.
139 + bs.strategy.MessageSent(p, message)
140 +
141 + if response == nil {
142 + return
143 + }
144 + bs.ReceiveMessage(ctx, p, response)
145 + }(peerToQuery)
146 + }
147 + return nil
148 +}
149 +
150 +func (bs *bitswap) run(ctx context.Context) {
151 + var sendlist <-chan peer.Peer
152 +
153 + // Every so often, we should resend out our current want list
154 + rebroadcastTime := time.Second * 5
155 +
156 + // Time to wait before sending out wantlists to better batch up requests
157 + bufferTime := time.Millisecond * 3
158 + peersPerSend := 6
159 +
160 + timeout := time.After(rebroadcastTime)
161 + threshold := 10
162 + unsent := 0
163 + for {
164 + select {
165 + case <-timeout:
166 + if sendlist == nil {
167 + // rely on semi randomness of maps
168 + firstKey := bs.wantlist.Keys()[0]
169 + sendlist = bs.routing.FindProvidersAsync(ctx, firstKey, 6)
170 + }
171 + err := bs.sendWantListTo(ctx, sendlist)
172 + if err != nil {
173 + log.Error("error sending wantlist: %s", err)
174 + }
175 + sendlist = nil
176 + timeout = time.After(rebroadcastTime)
177 + case k := <-bs.blockReq:
178 + if unsent == 0 {
179 + sendlist = bs.routing.FindProvidersAsync(ctx, k, peersPerSend)
180 + }
181 + unsent++
182 +
183 + if unsent >= threshold {
184 + // send wantlist to sendlist
185 + bs.sendWantListTo(ctx, sendlist)
186 + unsent = 0
187 + timeout = time.After(rebroadcastTime)
188 + sendlist = nil
189 + } else {
190 + // set a timeout to wait for more blocks or send current wantlist
191 +
192 + timeout = time.After(bufferTime)
193 + }
194 + case <-ctx.Done():
195 + return
196 + }
197 + }
198 +}
199 +
200 // HasBlock announces the existance of a block to this bitswap service. The
201 // service will potentially notify its peers.
202 func (bs *bitswap) HasBlock(ctx context.Context, blk blocks.Block) error {
@@ -192,8 +255,8 @@ func (bs *bitswap) ReceiveMessage(ctx context.Context, p peer.Peer, incoming bsm
255 }
256 }
257 }
195 - defer bs.strategy.MessageSent(p, message)
258
259 + bs.strategy.MessageSent(p, message)
260 log.Debug("Returning message.")
261 return p, message
262 }
exchange/bitswap/bitswap_test.go
+10 -8
@@ -31,7 +31,7 @@ func TestGetBlockTimeout(t *testing.T) {
31
32 ctx, _ := context.WithTimeout(context.Background(), time.Nanosecond)
33 block := blocks.NewBlock([]byte("block"))
34 - _, err := self.exchange.Block(ctx, block.Key())
34 + _, err := self.exchange.GetBlock(ctx, block.Key())
35
36 if err != context.DeadlineExceeded {
37 t.Fatal("Expected DeadlineExceeded error")
@@ -50,7 +50,7 @@ func TestProviderForKeyButNetworkCannotFind(t *testing.T) {
50 solo := g.Next()
51
52 ctx, _ := context.WithTimeout(context.Background(), time.Nanosecond)
53 - _, err := solo.exchange.Block(ctx, block.Key())
53 + _, err := solo.exchange.GetBlock(ctx, block.Key())
54
55 if err != context.DeadlineExceeded {
56 t.Fatal("Expected DeadlineExceeded error")
@@ -78,7 +78,7 @@ func TestGetBlockFromPeerAfterPeerAnnounces(t *testing.T) {
78 wantsBlock := g.Next()
79
80 ctx, _ := context.WithTimeout(context.Background(), time.Second)
81 - received, err := wantsBlock.exchange.Block(ctx, block.Key())
81 + received, err := wantsBlock.exchange.GetBlock(ctx, block.Key())
82 if err != nil {
83 t.Log(err)
84 t.Fatal("Expected to succeed")
@@ -100,7 +100,7 @@ func TestSwarm(t *testing.T) {
100
101 t.Log("Create a ton of instances, and just a few blocks")
102
103 - numInstances := 500
103 + numInstances := 5
104 numBlocks := 2
105
106 instances := sg.Instances(numInstances)
@@ -142,7 +142,7 @@ func TestSwarm(t *testing.T) {
142
143 func getOrFail(bitswap instance, b *blocks.Block, t *testing.T, wg *sync.WaitGroup) {
144 if _, err := bitswap.blockstore.Get(b.Key()); err != nil {
145 - _, err := bitswap.exchange.Block(context.Background(), b.Key())
145 + _, err := bitswap.exchange.GetBlock(context.Background(), b.Key())
146 if err != nil {
147 t.Fatal(err)
148 }
@@ -171,7 +171,7 @@ func TestSendToWantingPeer(t *testing.T) {
171
172 t.Logf("Peer %v attempts to get %v. NB: not available\n", w.peer, alpha.Key())
173 ctx, _ := context.WithTimeout(context.Background(), timeout)
174 - _, err := w.exchange.Block(ctx, alpha.Key())
174 + _, err := w.exchange.GetBlock(ctx, alpha.Key())
175 if err == nil {
176 t.Fatalf("Expected %v to NOT be available", alpha.Key())
177 }
@@ -186,7 +186,7 @@ func TestSendToWantingPeer(t *testing.T) {
186
187 t.Logf("%v gets %v from %v and discovers it wants %v\n", me.peer, beta.Key(), w.peer, alpha.Key())
188 ctx, _ = context.WithTimeout(context.Background(), timeout)
189 - if _, err := me.exchange.Block(ctx, beta.Key()); err != nil {
189 + if _, err := me.exchange.GetBlock(ctx, beta.Key()); err != nil {
190 t.Fatal(err)
191 }
192
@@ -199,7 +199,7 @@ func TestSendToWantingPeer(t *testing.T) {
199
200 t.Logf("%v requests %v\n", me.peer, alpha.Key())
201 ctx, _ = context.WithTimeout(context.Background(), timeout)
202 - if _, err := me.exchange.Block(ctx, alpha.Key()); err != nil {
202 + if _, err := me.exchange.GetBlock(ctx, alpha.Key()); err != nil {
203 t.Fatal(err)
204 }
205
@@ -290,8 +290,10 @@ func session(net tn.Network, rs mock.RoutingServer, id peer.ID) instance {
290 routing: htc,
291 sender: adapter,
292 wantlist: util.NewKeySet(),
293 + blockReq: make(chan util.Key, 32),
294 }
295 adapter.SetDelegate(bs)
296 + go bs.run(context.TODO())
297 return instance{
298 peer: p,
299 exchange: bs,
exchange/interface.go
+2 -2
@@ -12,8 +12,8 @@ import (
12 // exchange protocol.
13 type Interface interface {
14
15 - // Block returns the block associated with a given key.
16 - Block(context.Context, u.Key) (*blocks.Block, error)
15 + // GetBlock returns the block associated with a given key.
16 + GetBlock(context.Context, u.Key) (*blocks.Block, error)
17
18 // TODO Should callers be concerned with whether the block was made
19 // available on the network?
exchange/offline/offline.go
+2 -2
@@ -23,10 +23,10 @@ func NewOfflineExchange() exchange.Interface {
23 type offlineExchange struct {
24 }
25
26 -// Block returns nil to signal that a block could not be retrieved for the
26 +// GetBlock returns nil to signal that a block could not be retrieved for the
27 // given key.
28 // NB: This function may return before the timeout expires.
29 -func (_ *offlineExchange) Block(context.Context, u.Key) (*blocks.Block, error) {
29 +func (_ *offlineExchange) GetBlock(context.Context, u.Key) (*blocks.Block, error) {
30 return nil, OfflineMode
31 }
32
exchange/offline/offline_test.go
+1 -1
@@ -11,7 +11,7 @@ import (
11
12 func TestBlockReturnsErr(t *testing.T) {
13 off := NewOfflineExchange()
14 - _, err := off.Block(context.Background(), u.Key("foo"))
14 + _, err := off.GetBlock(context.Background(), u.Key("foo"))
15 if err != nil {
16 return // as desired
17 }