@cryptotaxi247 / kubo / commits / 678db4fa4

more work on bitswap and other code cleanup

Jeromy committed Aug 25, 2014 at 09:44 UTC 678db4fa40483915977603536ecaf0123a0315a8
7 files changed +121 -30
.gitignore new
+4
@@ -0,0 +1,4 @@
1 +*.swp
2 +.ipfsconfig
3 +*.out
4 +*.test
README.md
+9
@@ -62,6 +62,15 @@ Guidelines:
62 - if you'd like to work on ipfs part-time (20+ hrs/wk) or full-time (40+ hrs/wk), contact [@jbenet](https://github.com/jbenet)
63 - have fun!
64
65 +## Todo
66 +
67 +Ipfs is still under heavy development, there is a lot to be done!
68 +
69 +- [ ] Finish Bitswap
70 +- [ ] Connect fuse interface to Blockservice
71 +- [ ] Write tests for bitswap
72 +- [ ] Come up with more TODO items
73 +
74 ## Development Dependencies
75
76 If you make changes to the protocol buffers, you will need to install the [protoc compiler](https://code.google.com/p/protobuf/downloads/list).
bitswap/bitswap.go
+29 -24
@@ -5,6 +5,7 @@ import (
5 blocks "github.com/jbenet/go-ipfs/blocks"
6 peer "github.com/jbenet/go-ipfs/peer"
7 routing "github.com/jbenet/go-ipfs/routing"
8 + dht "github.com/jbenet/go-ipfs/routing/dht"
9 swarm "github.com/jbenet/go-ipfs/swarm"
10 u "github.com/jbenet/go-ipfs/util"
11
@@ -36,7 +37,7 @@ type BitSwap struct {
37 datastore ds.Datastore
38
39 // routing interface for communication
39 - routing routing.IpfsRouting
40 + routing *dht.IpfsDHT
41
42 listener *swarm.MesListener
43
@@ -63,7 +64,7 @@ func NewBitSwap(p *peer.Peer, net swarm.Network, d ds.Datastore, r routing.IpfsR
64 datastore: d,
65 partners: LedgerMap{},
66 wantList: KeySet{},
66 - routing: r,
67 + routing: r.(*dht.IpfsDHT),
68 meschan: net.GetChannel(swarm.PBWrapper_BITSWAP),
69 haltChan: make(chan struct{}),
70 }
@@ -76,32 +77,32 @@ func NewBitSwap(p *peer.Peer, net swarm.Network, d ds.Datastore, r routing.IpfsR
77 func (bs *BitSwap) GetBlock(k u.Key, timeout time.Duration) (
78 *blocks.Block, error) {
79 begin := time.Now()
79 - provs, err := bs.routing.FindProviders(k, timeout)
80 - if err != nil {
81 - u.PErr("GetBlock error: %s\n", err)
82 - return nil, err
83 - }
80 tleft := timeout - time.Now().Sub(begin)
81 + provs_ch := bs.routing.FindProvidersAsync(k, 20, timeout)
82
83 valchan := make(chan []byte)
84 after := time.After(tleft)
88 - for _, p := range provs {
89 - go func(pr *peer.Peer) {
90 - ledger := bs.GetLedger(pr.Key())
91 - blk, err := bs.getBlock(k, pr, tleft)
92 - if err != nil {
93 - u.PErr("%v\n", err)
94 - return
95 - }
96 - // NOTE: this credits everyone who sends us a block,
97 - // even if we dont use it
98 - ledger.ReceivedBytes(uint64(len(blk)))
99 - select {
100 - case valchan <- blk:
101 - default:
102 - }
103 - }(p)
104 - }
85 +
86 + // TODO: when the data is received, shut down this for loop
87 + go func() {
88 + for p := range provs_ch {
89 + go func(pr *peer.Peer) {
90 + ledger := bs.GetLedger(pr.Key())
91 + blk, err := bs.getBlock(k, pr, tleft)
92 + if err != nil {
93 + u.PErr("%v\n", err)
94 + return
95 + }
96 + // NOTE: this credits everyone who sends us a block,
97 + // even if we dont use it
98 + ledger.ReceivedBytes(uint64(len(blk)))
99 + select {
100 + case valchan <- blk:
101 + default:
102 + }
103 + }(p)
104 + }
105 + }()
106
107 select {
108 case blkdata := <-valchan:
@@ -213,3 +214,7 @@ func (bs *BitSwap) GetLedger(k u.Key) *Ledger {
214 bs.partners[k] = l
215 return l
216 }
217 +
218 +func (bs *BitSwap) Halt() {
219 + bs.haltChan <- struct{}{}
220 +}
blockservice/blocks_test.go renamed
+5 -3
@@ -1,11 +1,13 @@
1 -package blocks
1 +package blockservice
2
3 import (
4 "bytes"
5 "fmt"
6 + "testing"
7 +
8 ds "github.com/jbenet/datastore.go"
9 + blocks "github.com/jbenet/go-ipfs/blocks"
10 u "github.com/jbenet/go-ipfs/util"
8 - "testing"
11 )
12
13 func TestBlocks(t *testing.T) {
@@ -17,7 +19,7 @@ func TestBlocks(t *testing.T) {
19 return
20 }
21
20 - b, err := NewBlock([]byte("beep boop"))
22 + b, err := blocks.NewBlock([]byte("beep boop"))
23 if err != nil {
24 t.Error("failed to construct block", err)
25 return
importer/importer.go
+7 -2
@@ -2,10 +2,11 @@ package importer
2
3 import (
4 "fmt"
5 - dag "github.com/jbenet/go-ipfs/merkledag"
5 "io"
6 "io/ioutil"
7 "os"
8 +
9 + dag "github.com/jbenet/go-ipfs/merkledag"
10 )
11
12 // BlockSizeLimit specifies the maximum size an imported block can have.
@@ -23,12 +24,16 @@ func NewDagFromReader(r io.Reader, size int64) (*dag.Node, error) {
24 // todo: block-splitting based on rabin fingerprinting
25 // todo: block-splitting with user-defined function
26 // todo: block-splitting at all. :P
27 + // todo: write mote todos
28
29 // totally just trusts the reported size. fix later.
30 if size > BlockSizeLimit { // 1 MB limit for now.
31 return nil, ErrSizeLimitExceeded
32 }
33
34 + // Ensure that we dont get stuck reading way too much data
35 + r = io.LimitReader(r, BlockSizeLimit)
36 +
37 // we're doing it live!
38 buf, err := ioutil.ReadAll(r)
39 if err != nil {
@@ -52,7 +57,7 @@ func NewDagFromFile(fpath string) (*dag.Node, error) {
57 }
58
59 if stat.IsDir() {
55 - return nil, fmt.Errorf("`fpath` is a directory")
60 + return nil, fmt.Errorf("`%s` is a directory", fpath)
61 }
62
63 f, err := os.Open(fpath)
routing/dht/routing.go
+55 -1
@@ -22,7 +22,7 @@ import (
22
23 // PutValue adds value corresponding to given Key.
24 // This is the top level "Store" operation of the DHT
25 -func (dht *IpfsDHT) PutValue(key u.Key, value []byte) {
25 +func (dht *IpfsDHT) PutValue(key u.Key, value []byte) error {
26 complete := make(chan struct{})
27 count := 0
28 for _, route := range dht.routingTables {
@@ -45,6 +45,7 @@ func (dht *IpfsDHT) PutValue(key u.Key, value []byte) {
45 for i := 0; i < count; i++ {
46 <-complete
47 }
48 + return nil
49 }
50
51 // GetValue searches for the value corresponding to given Key.
@@ -183,6 +184,59 @@ func (dht *IpfsDHT) Provide(key u.Key) error {
184 return nil
185 }
186
187 +func (dht *IpfsDHT) FindProvidersAsync(key u.Key, count int, timeout time.Duration) chan *peer.Peer {
188 + peerOut := make(chan *peer.Peer, count)
189 + go func() {
190 + ps := newPeerSet()
191 + provs := dht.providers.GetProviders(key)
192 + for _, p := range provs {
193 + count--
194 + // NOTE: assuming that the list of peers is unique
195 + ps.Add(p)
196 + peerOut <- p
197 + if count <= 0 {
198 + return
199 + }
200 + }
201 +
202 + peers := dht.routingTables[0].NearestPeers(kb.ConvertKey(key), AlphaValue)
203 + for _, pp := range peers {
204 + go func() {
205 + pmes, err := dht.findProvidersSingle(pp, key, 0, timeout)
206 + if err != nil {
207 + u.PErr("%v\n", err)
208 + return
209 + }
210 + dht.addPeerListAsync(key, pmes.GetPeers(), ps, count, peerOut)
211 + }()
212 + }
213 +
214 + }()
215 + return peerOut
216 +}
217 +
218 +//TODO: this function could also be done asynchronously
219 +func (dht *IpfsDHT) addPeerListAsync(k u.Key, peers []*PBDHTMessage_PBPeer, ps *peerSet, count int, out chan *peer.Peer) {
220 + for _, pbp := range peers {
221 + maddr, err := ma.NewMultiaddr(pbp.GetAddr())
222 + if err != nil {
223 + u.PErr("%v\n", err)
224 + continue
225 + }
226 + p, err := dht.network.GetConnection(peer.ID(pbp.GetId()), maddr)
227 + if err != nil {
228 + u.PErr("%v\n", err)
229 + continue
230 + }
231 + dht.providers.AddProvider(k, p)
232 + if ps.AddIfSmallerThan(p, count) {
233 + out <- p
234 + } else if ps.Size() >= count {
235 + return
236 + }
237 + }
238 +}
239 +
240 // FindProviders searches for peers who can provide the value for given key.
241 func (dht *IpfsDHT) FindProviders(key u.Key, timeout time.Duration) ([]*peer.Peer, error) {
242 ll := startNewRPC("FindProviders")
routing/dht/util.go
+12
@@ -40,6 +40,7 @@ func (c *counter) Size() (s int) {
40 return
41 }
42
43 +// peerSet is a threadsafe set of peers
44 type peerSet struct {
45 ps map[string]bool
46 lk sync.RWMutex
@@ -69,3 +70,14 @@ func (ps *peerSet) Size() int {
70 defer ps.lk.RUnlock()
71 return len(ps.ps)
72 }
73 +
74 +func (ps *peerSet) AddIfSmallerThan(p *peer.Peer, maxsize int) bool {
75 + var success bool
76 + ps.lk.Lock()
77 + if _, ok := ps.ps[string(p.ID)]; !ok && len(ps.ps) < maxsize {
78 + success = true
79 + ps.ps[string(p.ID)] = true
80 + }
81 + ps.lk.Unlock()
82 + return success
83 +}