@cryptotaxi247 / kubo / commits / 881447e68

refac(bitswap) use blockstore

Brian Tiger Chow committed Sep 16, 2014 at 04:50 UTC 881447e68eb622bdfe97a75a789ecc905d9dd1d2
1 file changed +19 -26
bitswap/bitswap.go
+19 -26
@@ -11,6 +11,7 @@ import (
11 bsnet "github.com/jbenet/go-ipfs/bitswap/network"
12 notifications "github.com/jbenet/go-ipfs/bitswap/notifications"
13 blocks "github.com/jbenet/go-ipfs/blocks"
14 + blockstore "github.com/jbenet/go-ipfs/blockstore"
15 peer "github.com/jbenet/go-ipfs/peer"
16 routing "github.com/jbenet/go-ipfs/routing"
17 u "github.com/jbenet/go-ipfs/util"
@@ -35,8 +36,8 @@ type BitSwap struct {
36 // sender delivers messages on behalf of the session
37 sender bsnet.NetworkAdapter
38
38 - // datastore is the local database // Ledgers of known
39 - datastore ds.Datastore
39 + // blockstore is the local database
40 + blockstore blockstore.Blockstore
41
42 // routing interface for communication
43 routing routing.IpfsRouting
@@ -66,7 +67,7 @@ func NewSession(parent context.Context, s bsnet.NetworkService, p *peer.Peer, d
67 receiver := bsnet.Forwarder{}
68 bs := &BitSwap{
69 peer: p,
69 - datastore: d,
70 + blockstore: blockstore.NewBlockstore(d),
71 partners: LedgerMap{},
72 wantList: KeySet{},
73 routing: r,
@@ -151,7 +152,9 @@ func (bs *BitSwap) HasBlock(blk blocks.Block) error {
152 return bs.routing.Provide(blk.Key())
153 }
154
155 +// TODO(brian): get a return value
156 func (bs *BitSwap) SendBlock(p *peer.Peer, b blocks.Block) {
157 + u.DOut("Sending block to peer.\n")
158 message := bsmsg.New()
159 // TODO(brian): change interface to accept value instead of pointer
160 message.AppendBlock(b)
@@ -162,40 +165,26 @@ func (bs *BitSwap) SendBlock(p *peer.Peer, b blocks.Block) {
165 // and then if we do, check the ledger for whether or not we should send it.
166 func (bs *BitSwap) peerWantsBlock(p *peer.Peer, wanted u.Key) {
167 u.DOut("peer [%s] wants block [%s]\n", p.ID.Pretty(), wanted.Pretty())
168 +
169 ledger := bs.getLedger(p)
170
167 - blk_i, err := bs.datastore.Get(wanted.DatastoreKey())
168 - if err != nil {
169 - if err == ds.ErrNotFound {
170 - ledger.Wants(wanted)
171 - }
172 - u.PErr("datastore get error: %v\n", err)
171 + if !ledger.ShouldSend() {
172 return
173 }
174
176 - blk, ok := blk_i.([]byte)
177 - if !ok {
178 - u.PErr("data conversion error.\n")
175 + block, err := bs.blockstore.Get(wanted)
176 + if err != nil { // TODO(brian): log/return the error
177 + ledger.Wants(wanted)
178 return
179 }
181 -
182 - if ledger.ShouldSend() {
183 - u.DOut("Sending block to peer.\n")
184 - bblk, err := blocks.NewBlock(blk)
185 - if err != nil {
186 - u.PErr("newBlock error: %v\n", err)
187 - return
188 - }
189 - bs.SendBlock(p, *bblk)
190 - ledger.SentBytes(len(blk))
191 - } else {
192 - u.DOut("Decided not to send block.")
193 - }
180 + bs.SendBlock(p, *block)
181 + ledger.SentBytes(numBytes(*block))
182 }
183
184 +// TODO(brian): return error
185 func (bs *BitSwap) blockReceive(p *peer.Peer, blk blocks.Block) {
186 u.DOut("blockReceive: %s\n", blk.Key().Pretty())
198 - err := bs.datastore.Put(ds.NewKey(string(blk.Key())), blk.Data)
187 + err := bs.blockstore.Put(blk)
188 if err != nil {
189 u.PErr("blockReceive error: %v\n", err)
190 return
@@ -255,3 +244,7 @@ func (bs *BitSwap) ReceiveMessage(
244 }
245 return nil, nil, errors.New("TODO implement")
246 }
247 +
248 +func numBytes(b blocks.Block) int {
249 + return len(b.Data)
250 +}