@cryptotaxi247 / kubo / commits / 0075a47df

fix(bs) remove concrete refs to swarm and dht

Brian Tiger Chow committed Sep 14, 2014 at 14:30 UTC 0075a47df02d6c14a121dd434c7b672e34349d19
3 files changed +16 -49
bitswap/bitswap.go
+14 -48
@@ -11,13 +11,13 @@ import (
11 notifications "github.com/jbenet/go-ipfs/bitswap/notifications"
12 tx "github.com/jbenet/go-ipfs/bitswap/transmission"
13 blocks "github.com/jbenet/go-ipfs/blocks"
14 - swarm "github.com/jbenet/go-ipfs/net/swarm"
14 peer "github.com/jbenet/go-ipfs/peer"
15 routing "github.com/jbenet/go-ipfs/routing"
17 - dht "github.com/jbenet/go-ipfs/routing/dht"
16 u "github.com/jbenet/go-ipfs/util"
17 )
18
19 +// TODO(brian): ensure messages are being received
20 +
21 // PartnerWantListMax is the bound for the number of keys we'll store per
22 // partner. These are usually taken from the top of the Partner's WantList
23 // advertisements. WantLists are sorted in terms of priority.
@@ -32,16 +32,14 @@ type BitSwap struct {
32 // peer is the identity of this (local) node.
33 peer *peer.Peer
34
35 - // net holds the connections to all peers.
36 - sender tx.Sender
37 - net swarm.Network
38 - meschan *swarm.Chan
35 + // sender delivers messages on behalf of the session
36 + sender tx.Sender
37
38 // datastore is the local database // Ledgers of known
39 datastore ds.Datastore
40
41 // routing interface for communication
44 - routing *dht.IpfsDHT
42 + routing routing.IpfsRouting
43
44 notifications notifications.PubSub
45
@@ -63,25 +61,21 @@ type BitSwap struct {
61 }
62
63 // NewBitSwap creates a new BitSwap instance. It does not check its parameters.
66 -func NewBitSwap(p *peer.Peer, net swarm.Network, d ds.Datastore, r routing.IpfsRouting) *BitSwap {
64 +func NewBitSwap(p *peer.Peer, d ds.Datastore, r routing.IpfsRouting) *BitSwap {
65 receiver := tx.Forwarder{}
66 sender := tx.NewBSNetService(context.Background(), &receiver)
67 bs := &BitSwap{
70 - peer: p,
71 - net: net,
72 - datastore: d,
73 - partners: LedgerMap{},
74 - wantList: KeySet{},
75 - routing: r.(*dht.IpfsDHT),
76 - // TODO(brian): replace |meschan| with |sender| in BitSwap impl
77 - meschan: net.GetChannel(swarm.PBWrapper_BITSWAP),
68 + peer: p,
69 + datastore: d,
70 + partners: LedgerMap{},
71 + wantList: KeySet{},
72 + routing: r,
73 sender: sender,
74 haltChan: make(chan struct{}),
75 notifications: notifications.New(),
76 }
77 receiver.Delegate(bs)
78
84 - go bs.handleMessages()
79 return bs
80 }
81
@@ -130,7 +124,7 @@ func (bs *BitSwap) getBlock(k u.Key, p *peer.Peer, timeout time.Duration) (*bloc
124
125 message := bsmsg.New()
126 message.AppendWanted(k)
133 - bs.meschan.Outgoing <- message.ToSwarm(p)
127 + bs.sender.SendMessage(ctx, p, message)
128
129 block, ok := <-blockChannel
130 if !ok {
@@ -159,35 +153,7 @@ func (bs *BitSwap) HaveBlock(blk *blocks.Block) error {
153 func (bs *BitSwap) SendBlock(p *peer.Peer, b *blocks.Block) {
154 message := bsmsg.New()
155 message.AppendBlock(b)
162 - bs.meschan.Outgoing <- message.ToSwarm(p)
163 -}
164 -
165 -func (bs *BitSwap) handleMessages() {
166 - for {
167 - select {
168 - case mes := <-bs.meschan.Incoming:
169 - bsmsg, err := bsmsg.FromSwarm(*mes)
170 - if err != nil {
171 - u.PErr("%v\n", err)
172 - continue
173 - }
174 -
175 - if bsmsg.Blocks() != nil {
176 - for _, blk := range bsmsg.Blocks() {
177 - go bs.blockReceive(mes.Peer, blk)
178 - }
179 - }
180 -
181 - if bsmsg.Wantlist() != nil {
182 - for _, want := range bsmsg.Wantlist() {
183 - go bs.peerWantsBlock(mes.Peer, want)
184 - }
185 - }
186 - case <-bs.haltChan:
187 - bs.notifications.Shutdown()
188 - return
189 - }
190 - }
156 + bs.sender.SendMessage(context.Background(), p, message)
157 }
158
159 // peerWantsBlock will check if we have the block in question,
@@ -260,7 +226,7 @@ func (bs *BitSwap) SendWantList(wl KeySet) error {
226
227 // Lets just ping everybody all at once
228 for _, ledger := range bs.partners {
263 - bs.meschan.Outgoing <- message.ToSwarm(ledger.Partner)
229 + bs.sender.SendMessage(context.TODO(), ledger.Partner, message)
230 }
231
232 return nil
core/core.go
+1 -1
@@ -99,7 +99,7 @@ func NewIpfsNode(cfg *config.Config, online bool) (*IpfsNode, error) {
99 route.Start()
100
101 // TODO(brian): pass a context to bs for its async operations
102 - swap = bitswap.NewBitSwap(local, net, d, route)
102 + swap = bitswap.NewBitSwap(local, d, route)
103 swap.SetStrategy(bitswap.YesManStrategy)
104
105 // TODO(brian): pass a context to initConnections
routing/routing.go
+1
@@ -10,6 +10,7 @@ import (
10 // IpfsRouting is the routing module interface
11 // It is implemented by things like DHTs, etc.
12 type IpfsRouting interface {
13 + FindProvidersAsync(u.Key, int, time.Duration) <-chan *peer.Peer
14
15 // Basic Put/Get
16