use raw primitives
Brian Tiger Chow committed
Dec 25, 2014 at 01:18 UTC
217611c2370e5fb99793fe5575dadb37f52939f4
1 file changed
+77
-41
epictest/addcat_test.go
+77
-41
@@ -10,19 +10,27 @@ import (
10
"time"
11
12
context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
13
- "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-random"
13
+ datastore "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-datastore"
14
+ sync "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-datastore/sync"
15
+ random "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-random"
16
+ blockstore "github.com/jbenet/go-ipfs/blocks/blockstore"
17
blockservice "github.com/jbenet/go-ipfs/blockservice"
18
+ exchange "github.com/jbenet/go-ipfs/exchange"
19
bitswap "github.com/jbenet/go-ipfs/exchange/bitswap"
16
- tn "github.com/jbenet/go-ipfs/exchange/bitswap/testnet"
20
+ bsnet "github.com/jbenet/go-ipfs/exchange/bitswap/network"
21
importer "github.com/jbenet/go-ipfs/importer"
22
chunk "github.com/jbenet/go-ipfs/importer/chunk"
23
merkledag "github.com/jbenet/go-ipfs/merkledag"
24
+ net "github.com/jbenet/go-ipfs/net"
25
mocknet "github.com/jbenet/go-ipfs/net/mock"
26
path "github.com/jbenet/go-ipfs/path"
22
- mockrouting "github.com/jbenet/go-ipfs/routing/mock"
27
+ peer "github.com/jbenet/go-ipfs/peer"
28
+ dht "github.com/jbenet/go-ipfs/routing/dht"
29
uio "github.com/jbenet/go-ipfs/unixfs/io"
30
util "github.com/jbenet/go-ipfs/util"
31
+ "github.com/jbenet/go-ipfs/util/datastore2"
32
errors "github.com/jbenet/go-ipfs/util/debugerror"
33
+ delay "github.com/jbenet/go-ipfs/util/delay"
34
)
35
36
const kSeed = 1
@@ -34,7 +42,7 @@ func Test1KBInstantaneous(t *testing.T) {
42
BlockstoreLatency: 0,
43
}
44
37
- if err := AddCatBytes(RandomBytes(1*KB), conf); err != nil {
45
+ if err := AddCatBytes(RandomBytes(100*MB), conf); err != nil {
46
t.Fatal(err)
47
}
48
}
@@ -88,55 +96,83 @@ func RandomBytes(n int64) []byte {
96
return data.Bytes()
97
}
98
99
+type instance struct {
100
+ ID peer.ID
101
+ Network net.Network
102
+ Blockstore blockstore.Blockstore
103
+ Datastore datastore.ThreadSafeDatastore
104
+ DHT *dht.IpfsDHT
105
+ Exchange exchange.Interface
106
+ BitSwapNetwork bsnet.BitSwapNetwork
107
+
108
+ datastoreDelay delay.D
109
+}
110
+
111
func AddCatBytes(data []byte, conf Config) error {
112
ctx, cancel := context.WithCancel(context.Background())
113
defer cancel()
94
- mn := mocknet.New(ctx)
95
- // defer mn.Close() FIXME does mocknet require clean-up
114
+ const numPeers = 2
115
+ instances := make(map[peer.ID]*instance, numPeers)
116
+
117
+ // create network
118
+ mn, err := mocknet.FullMeshLinked(ctx, numPeers)
119
+ if err != nil {
120
+ return errors.Wrap(err)
121
+ }
122
mn.SetLinkDefaults(mocknet.LinkOptions{
97
- Latency: conf.NetworkLatency,
123
+ Latency: conf.NetworkLatency,
124
// TODO add to conf. This is tricky because we want 0 values to be functional.
125
Bandwidth: math.MaxInt32,
126
})
101
- dhtNetwork := mockrouting.NewDHTNetwork(mn)
102
- net, err := tn.StreamNet(ctx, mn, dhtNetwork)
103
- if err != nil {
104
- return errors.Wrap(err)
127
+ for _, p := range mn.Peers() {
128
+ instances[p] = &instance{
129
+ ID: p,
130
+ Network: mn.Net(p),
131
+ }
132
+ }
133
+
134
+ // create dht network
135
+ for _, p := range mn.Peers() {
136
+ dsDelay := delay.Fixed(conf.BlockstoreLatency)
137
+ instances[p].Datastore = sync.MutexWrap(datastore2.WithDelay(datastore.NewMapDatastore(), dsDelay))
138
+ instances[p].datastoreDelay = dsDelay
139
+ }
140
+ for _, p := range mn.Peers() {
141
+ instances[p].DHT = dht.NewDHT(ctx, p, instances[p].Network, instances[p].Datastore)
142
}
106
- sessionGenerator := bitswap.NewSessionGenerator(net)
107
- defer sessionGenerator.Close()
108
-
109
- adder := sessionGenerator.Next()
110
- catter := sessionGenerator.Next()
111
- // catter.Routing.Update(context.TODO(), adder.Peer)
112
-
113
- peers := mn.Peers()
114
- if len(peers) != 2 {
115
- return errors.New("peers not in network")
116
- }
117
-
118
- for _, i := range peers {
119
- for _, j := range peers {
120
- if i == j {
121
- continue
122
- }
123
- if _, err := mn.LinkPeers(i, j); err != nil {
124
- return err
125
- }
126
- if err := mn.ConnectPeers(i, j); err != nil {
127
- return err
128
- }
143
+ // create two bitswap network clients
144
+ for _, p := range mn.Peers() {
145
+ instances[p].BitSwapNetwork = bsnet.NewFromIpfsNetwork(instances[p].Network, instances[p].DHT)
146
+ }
147
+ for _, p := range mn.Peers() {
148
+ const kWriteCacheElems = 100
149
+ const alwaysSendToPeer = true
150
+ adapter := instances[p].BitSwapNetwork
151
+ dstore := instances[p].Datastore
152
+ instances[p].Blockstore, err = blockstore.WriteCached(blockstore.NewBlockstore(dstore), kWriteCacheElems)
153
+ if err != nil {
154
+ return err
155
}
156
+ instances[p].Exchange = bitswap.New(ctx, p, adapter, instances[p].Blockstore, alwaysSendToPeer)
157
}
158
+ var peers []peer.ID
159
+ for _, p := range mn.Peers() {
160
+ peers = append(peers, p)
161
+ }
162
+
163
+ adder := instances[peers[0]]
164
+ catter := instances[peers[1]]
165
132
- catter.SetBlockstoreLatency(conf.BlockstoreLatency)
166
+ // bootstrap the DHTs
167
+ adder.DHT.Connect(ctx, catter.ID)
168
+ catter.DHT.Connect(ctx, adder.ID)
169
134
- adder.SetBlockstoreLatency(0) // disable blockstore latency during add operation
170
+ adder.datastoreDelay.Set(0) // disable blockstore latency during add operation
171
keyAdded, err := add(adder, bytes.NewReader(data))
172
if err != nil {
173
return err
174
}
139
- adder.SetBlockstoreLatency(conf.BlockstoreLatency) // add some blockstore delay to make the catter wait
175
+ adder.datastoreDelay.Set(conf.BlockstoreLatency) // add some blockstore delay to make the catter wait
176
177
readerCatted, err := cat(catter, keyAdded)
178
if err != nil {
@@ -152,8 +188,8 @@ func AddCatBytes(data []byte, conf Config) error {
188
return nil
189
}
190
155
-func cat(catter bitswap.Instance, k util.Key) (io.Reader, error) {
156
- catterdag := merkledag.NewDAGService(&blockservice.BlockService{catter.Blockstore(), catter.Exchange})
191
+func cat(catter *instance, k util.Key) (io.Reader, error) {
192
+ catterdag := merkledag.NewDAGService(&blockservice.BlockService{catter.Blockstore, catter.Exchange})
193
nodeCatted, err := (&path.Resolver{catterdag}).ResolvePath(k.String())
194
if err != nil {
195
return nil, err
@@ -161,10 +197,10 @@ func cat(catter bitswap.Instance, k util.Key) (io.Reader, error) {
197
return uio.NewDagReader(nodeCatted, catterdag)
198
}
199
164
-func add(adder bitswap.Instance, r io.Reader) (util.Key, error) {
200
+func add(adder *instance, r io.Reader) (util.Key, error) {
201
nodeAdded, err := importer.BuildDagFromReader(
202
r,
167
- merkledag.NewDAGService(&blockservice.BlockService{adder.Blockstore(), adder.Exchange}),
203
+ merkledag.NewDAGService(&blockservice.BlockService{adder.Blockstore, adder.Exchange}),
204
nil,
205
chunk.DefaultSplitter,
206
)