add more tests and rework a lot of utility structures
Jeromy committed
Oct 7, 2014 at 05:55 UTC
5c802ae8524004bf37db96e98ae518bf1cf90252
10 files changed
+175
-128
importer/dagwriter/dagmodifier.go
+11
-12
@@ -1,6 +1,7 @@
1
package dagwriter
2
3
import (
4
+ "bytes"
5
"errors"
6
7
"code.google.com/p/goprotobuf/proto"
@@ -18,19 +19,21 @@ type DagModifier struct {
19
dagserv *mdag.DAGService
20
curNode *mdag.Node
21
21
- pbdata *ft.PBData
22
+ pbdata *ft.PBData
23
+ splitter imp.BlockSplitter
24
}
25
24
-func NewDagModifier(from *mdag.Node, serv *mdag.DAGService) (*DagModifier, error) {
26
+func NewDagModifier(from *mdag.Node, serv *mdag.DAGService, spl imp.BlockSplitter) (*DagModifier, error) {
27
pbd, err := ft.FromBytes(from.Data)
28
if err != nil {
29
return nil, err
30
}
31
32
return &DagModifier{
31
- curNode: from.Copy(),
32
- dagserv: serv,
33
- pbdata: pbd,
33
+ curNode: from.Copy(),
34
+ dagserv: serv,
35
+ pbdata: pbd,
36
+ splitter: spl,
37
}, nil
38
}
39
@@ -136,8 +139,7 @@ func (dm *DagModifier) WriteAt(b []byte, offset uint64) (int, error) {
139
b = append(b, data[midoff:]...)
140
}
141
139
- // TODO: dont assume a splitting func here
140
- subblocks := splitBytes(b, &imp.SizeSplitter2{512})
142
+ subblocks := splitBytes(b, dm.splitter)
143
var links []*mdag.Link
144
var sizes []uint64
145
for _, sb := range subblocks {
@@ -168,11 +170,8 @@ func (dm *DagModifier) WriteAt(b []byte, offset uint64) (int, error) {
170
return origlen, nil
171
}
172
171
-func splitBytes(b []byte, spl imp.StreamSplitter) [][]byte {
172
- ch := make(chan []byte)
173
- out := spl.Split(ch)
174
- ch <- b
175
- close(ch)
173
+func splitBytes(b []byte, spl imp.BlockSplitter) [][]byte {
174
+ out := spl.Split(bytes.NewReader(b))
175
var arr [][]byte
176
for blk := range out {
177
arr = append(arr, blk)
importer/dagwriter/dagmodifier_test.go
+5
-33
@@ -4,46 +4,18 @@ import (
4
"fmt"
5
"io"
6
"io/ioutil"
7
- "math/rand"
7
"testing"
9
- "time"
8
9
"github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/op/go-logging"
10
bs "github.com/jbenet/go-ipfs/blockservice"
11
imp "github.com/jbenet/go-ipfs/importer"
12
ft "github.com/jbenet/go-ipfs/importer/format"
13
mdag "github.com/jbenet/go-ipfs/merkledag"
14
+ u "github.com/jbenet/go-ipfs/util"
15
16
ds "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/datastore.go"
17
)
18
20
-type randGen struct {
21
- src rand.Source
22
-}
23
-
24
-func newRand() *randGen {
25
- return &randGen{rand.NewSource(time.Now().UnixNano())}
26
-}
27
-
28
-func (r *randGen) Read(p []byte) (n int, err error) {
29
- todo := len(p)
30
- offset := 0
31
- for {
32
- val := int64(r.src.Int63())
33
- for i := 0; i < 8; i++ {
34
- p[offset] = byte(val & 0xff)
35
- todo--
36
- if todo == 0 {
37
- return len(p), nil
38
- }
39
- offset++
40
- val >>= 8
41
- }
42
- }
43
-
44
- panic("unreachable")
45
-}
46
-
19
func getMockDagServ(t *testing.T) *mdag.DAGService {
20
dstore := ds.NewMapDatastore()
21
bserv, err := bs.NewBlockService(dstore, nil)
@@ -54,9 +26,9 @@ func getMockDagServ(t *testing.T) *mdag.DAGService {
26
}
27
28
func getNode(t *testing.T, dserv *mdag.DAGService, size int64) ([]byte, *mdag.Node) {
57
- dw := NewDagWriter(dserv, &imp.SizeSplitter2{500})
29
+ dw := NewDagWriter(dserv, &imp.SizeSplitter{500})
30
59
- n, err := io.CopyN(dw, newRand(), size)
31
+ n, err := io.CopyN(dw, u.NewFastRand(), size)
32
if err != nil {
33
t.Fatal(err)
34
}
@@ -82,7 +54,7 @@ func getNode(t *testing.T, dserv *mdag.DAGService, size int64) ([]byte, *mdag.No
54
55
func testModWrite(t *testing.T, beg, size uint64, orig []byte, dm *DagModifier) []byte {
56
newdata := make([]byte, size)
85
- r := newRand()
57
+ r := u.NewFastRand()
58
r.Read(newdata)
59
60
if size+beg > uint64(len(orig)) {
@@ -127,7 +99,7 @@ func TestDagModifierBasic(t *testing.T) {
99
dserv := getMockDagServ(t)
100
b, n := getNode(t, dserv, 50000)
101
130
- dagmod, err := NewDagModifier(n, dserv)
102
+ dagmod, err := NewDagModifier(n, dserv, &imp.SizeSplitter{512})
103
if err != nil {
104
t.Fatal(err)
105
}
importer/dagwriter/dagwriter.go
+4
-3
@@ -15,11 +15,11 @@ type DagWriter struct {
15
totalSize int64
16
splChan chan []byte
17
done chan struct{}
18
- splitter imp.StreamSplitter
18
+ splitter imp.BlockSplitter
19
seterr error
20
}
21
22
-func NewDagWriter(ds *dag.DAGService, splitter imp.StreamSplitter) *DagWriter {
22
+func NewDagWriter(ds *dag.DAGService, splitter imp.BlockSplitter) *DagWriter {
23
dw := new(DagWriter)
24
dw.dagserv = ds
25
dw.splChan = make(chan []byte, 8)
@@ -30,7 +30,8 @@ func NewDagWriter(ds *dag.DAGService, splitter imp.StreamSplitter) *DagWriter {
30
}
31
32
func (dw *DagWriter) startSplitter() {
33
- blkchan := dw.splitter.Split(dw.splChan)
33
+ r := util.NewByteChanReader(dw.splChan)
34
+ blkchan := dw.splitter.Split(r)
35
first := <-blkchan
36
mbf := new(ft.MultiBlock)
37
root := new(dag.Node)
importer/dagwriter/dagwriter_test.go
+3
-3
@@ -54,7 +54,7 @@ func TestDagWriter(t *testing.T) {
54
t.Fatal(err)
55
}
56
dag := &mdag.DAGService{bserv}
57
- dw := NewDagWriter(dag, &imp.SizeSplitter2{4096})
57
+ dw := NewDagWriter(dag, &imp.SizeSplitter{4096})
58
59
nbytes := int64(1024 * 1024 * 2)
60
n, err := io.CopyN(dw, &datasource{}, nbytes)
@@ -88,7 +88,7 @@ func TestMassiveWrite(t *testing.T) {
88
t.Fatal(err)
89
}
90
dag := &mdag.DAGService{bserv}
91
- dw := NewDagWriter(dag, &imp.SizeSplitter2{4096})
91
+ dw := NewDagWriter(dag, &imp.SizeSplitter{4096})
92
93
nbytes := int64(1024 * 1024 * 1024 * 16)
94
n, err := io.CopyN(dw, &datasource{}, nbytes)
@@ -113,7 +113,7 @@ func BenchmarkDagWriter(b *testing.B) {
113
nbytes := int64(b.N)
114
for i := 0; i < b.N; i++ {
115
b.SetBytes(nbytes)
116
- dw := NewDagWriter(dag, &imp.SizeSplitter2{4096})
116
+ dw := NewDagWriter(dag, &imp.SizeSplitter{4096})
117
n, err := io.CopyN(dw, &datasource{}, nbytes)
118
if err != nil {
119
b.Fatal(err)
importer/format/format_test.go
new
+36
@@ -0,0 +1,36 @@
1
+package format
2
+
3
+import (
4
+ "testing"
5
+
6
+ "code.google.com/p/goprotobuf/proto"
7
+)
8
+
9
+func TestMultiBlock(t *testing.T) {
10
+ mbf := new(MultiBlock)
11
+ for i := 0; i < 15; i++ {
12
+ mbf.AddBlockSize(100)
13
+ }
14
+
15
+ mbf.Data = make([]byte, 128)
16
+
17
+ b, err := mbf.GetBytes()
18
+ if err != nil {
19
+ t.Fatal(err)
20
+ }
21
+
22
+ pbn := new(PBData)
23
+ err = proto.Unmarshal(b, pbn)
24
+ if err != nil {
25
+ t.Fatal(err)
26
+ }
27
+
28
+ ds, err := DataSize(b)
29
+ if err != nil {
30
+ t.Fatal(err)
31
+ }
32
+
33
+ if ds != (100*15)+128 {
34
+ t.Fatal("Datasize calculations incorrect!")
35
+ }
36
+}
importer/importer.go
+3
@@ -7,8 +7,11 @@ import (
7
8
ft "github.com/jbenet/go-ipfs/importer/format"
9
dag "github.com/jbenet/go-ipfs/merkledag"
10
+ "github.com/jbenet/go-ipfs/util"
11
)
12
13
+var log = util.Logger("importer")
14
+
15
// BlockSizeLimit specifies the maximum size an imported block can have.
16
var BlockSizeLimit = int64(1048576) // 1 MB
17
importer/rabin.go
-40
@@ -92,43 +92,3 @@ func (mr *MaybeRabin) Split(r io.Reader) chan []byte {
92
}()
93
return out
94
}
95
-
96
-/*
97
-func WhyrusleepingCantImplementRabin(r io.Reader) chan []byte {
98
- out := make(chan []byte, 4)
99
- go func() {
100
- buf := bufio.NewReader(r)
101
- blkbuf := new(bytes.Buffer)
102
- window := make([]byte, 16)
103
- var val uint64
104
- prime := uint64(61)
105
-
106
- get := func(i int) uint64 {
107
- return uint64(window[i%len(window)])
108
- }
109
-
110
- set := func(i int, val byte) {
111
- window[i%len(window)] = val
112
- }
113
-
114
- for i := 0; ; i++ {
115
- curb, err := buf.ReadByte()
116
- if err != nil {
117
- break
118
- }
119
- set(i, curb)
120
- blkbuf.WriteByte(curb)
121
-
122
- hash := md5.Sum(window)
123
- if hash[0] == 0 && hash[1] == 0 {
124
- out <- blkbuf.Bytes()
125
- blkbuf.Reset()
126
- }
127
- }
128
- out <- blkbuf.Bytes()
129
- close(out)
130
- }()
131
-
132
- return out
133
-}
134
-*/
importer/splitting.go
+3
-36
@@ -1,19 +1,9 @@
1
package importer
2
3
-import (
4
- "io"
3
+import "io"
4
6
- u "github.com/jbenet/go-ipfs/util"
7
-)
8
-
9
-// OLD
5
type BlockSplitter interface {
11
- Split(io.Reader) chan []byte
12
-}
13
-
14
-// NEW
15
-type StreamSplitter interface {
16
- Split(chan []byte) chan []byte
6
+ Split(r io.Reader) chan []byte
7
}
8
9
type SizeSplitter struct {
@@ -34,7 +24,7 @@ func (ss *SizeSplitter) Split(r io.Reader) chan []byte {
24
}
25
return
26
}
37
- u.PErr("block split error: %v\n", err)
27
+ log.Error("Block split error: %s", err)
28
return
29
}
30
if nread < ss.Size {
@@ -45,26 +35,3 @@ func (ss *SizeSplitter) Split(r io.Reader) chan []byte {
35
}()
36
return out
37
}
48
-
49
-type SizeSplitter2 struct {
50
- Size int
51
-}
52
-
53
-func (ss *SizeSplitter2) Split(in chan []byte) chan []byte {
54
- out := make(chan []byte)
55
- go func() {
56
- defer close(out)
57
- var buf []byte
58
- for b := range in {
59
- buf = append(buf, b...)
60
- for len(buf) > ss.Size {
61
- out <- buf[:ss.Size]
62
- buf = buf[ss.Size:]
63
- }
64
- }
65
- if len(buf) > 0 {
66
- out <- buf
67
- }
68
- }()
69
- return out
70
-}
util/util.go
+76
@@ -3,10 +3,13 @@ package util
3
import (
4
"errors"
5
"fmt"
6
+ "io"
7
+ "math/rand"
8
"os"
9
"os/user"
10
"path/filepath"
11
"strings"
12
+ "time"
13
14
ds "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/datastore.go"
15
b58 "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-base58"
@@ -154,3 +157,76 @@ func ExpandPathnames(paths []string) ([]string, error) {
157
}
158
return out, nil
159
}
160
+
161
+// byteChanReader wraps a byte chan in a reader
162
+type byteChanReader struct {
163
+ in chan []byte
164
+ buf []byte
165
+}
166
+
167
+func NewByteChanReader(in chan []byte) io.Reader {
168
+ return &byteChanReader{in: in}
169
+}
170
+
171
+func (bcr *byteChanReader) Read(b []byte) (int, error) {
172
+ if len(bcr.buf) == 0 {
173
+ data, ok := <-bcr.in
174
+ if !ok {
175
+ return 0, io.EOF
176
+ }
177
+ bcr.buf = data
178
+ }
179
+
180
+ if len(bcr.buf) >= len(b) {
181
+ copy(b, bcr.buf)
182
+ bcr.buf = bcr.buf[len(b):]
183
+ return len(b), nil
184
+ }
185
+
186
+ copy(b, bcr.buf)
187
+ b = b[len(bcr.buf):]
188
+ totread := len(bcr.buf)
189
+
190
+ for data := range bcr.in {
191
+ if len(data) > len(b) {
192
+ totread += len(b)
193
+ copy(b, data[:len(b)])
194
+ bcr.buf = data[len(b):]
195
+ return totread, nil
196
+ }
197
+ copy(b, data)
198
+ totread += len(data)
199
+ b = b[len(data):]
200
+ if len(b) == 0 {
201
+ return totread, nil
202
+ }
203
+ }
204
+ return totread, io.EOF
205
+}
206
+
207
+type randGen struct {
208
+ src rand.Source
209
+}
210
+
211
+func NewFastRand() io.Reader {
212
+ return &randGen{rand.NewSource(time.Now().UnixNano())}
213
+}
214
+
215
+func (r *randGen) Read(p []byte) (n int, err error) {
216
+ todo := len(p)
217
+ offset := 0
218
+ for {
219
+ val := int64(r.src.Int63())
220
+ for i := 0; i < 8; i++ {
221
+ p[offset] = byte(val & 0xff)
222
+ todo--
223
+ if todo == 0 {
224
+ return len(p), nil
225
+ }
226
+ offset++
227
+ val >>= 8
228
+ }
229
+ }
230
+
231
+ panic("unreachable")
232
+}
util/util_test.go
+34
-1
@@ -2,8 +2,11 @@ package util
2
3
import (
4
"bytes"
5
- mh "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multihash"
5
+ "io/ioutil"
6
+ "math/rand"
7
"testing"
8
+
9
+ mh "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multihash"
10
)
11
12
func TestKey(t *testing.T) {
@@ -25,3 +28,33 @@ func TestKey(t *testing.T) {
28
t.Error("Keys not equal.")
29
}
30
}
31
+
32
+func TestByteChanReader(t *testing.T) {
33
+ data := make([]byte, 1024*1024)
34
+ r := NewFastRand()
35
+ r.Read(data)
36
+ dch := make(chan []byte, 8)
37
+
38
+ go func() {
39
+ beg := 0
40
+ for i := 0; i < len(data); {
41
+ i += rand.Intn(100) + 1
42
+ if i > len(data) {
43
+ i = len(data)
44
+ }
45
+ dch <- data[beg:i]
46
+ beg = i
47
+ }
48
+ close(dch)
49
+ }()
50
+
51
+ read := NewByteChanReader(dch)
52
+ out, err := ioutil.ReadAll(read)
53
+ if err != nil {
54
+ t.Fatal(err)
55
+ }
56
+
57
+ if !bytes.Equal(out, data) {
58
+ t.Fatal("Reader failed to stream correct bytes")
59
+ }
60
+}