@cryptotaxi247 / kubo / commits / 9ae420abb

SizeSplitter fix: keep-reading until chunk full

if the underlying reader is buffered with a smaller buffer it would force the chunk sizes to come out smaller than intended. cc @whyrusleeping @mappum

Juan Batiz-Benet committed Nov 18, 2014 at 05:48 UTC 9ae420abb80538d09e9777827ee7bcf75d810ec3
2 files changed +70 -10
importer/chunk/splitting.go
+21 -10
@@ -23,23 +23,34 @@ func (ss *SizeSplitter) Split(r io.Reader) chan []byte {
23 out := make(chan []byte)
24 go func() {
25 defer close(out)
26 +
27 + // all-chunks loop (keep creating chunks)
28 for {
29 + // log.Infof("making chunk with size: %d", ss.Size)
30 chunk := make([]byte, ss.Size)
28 - nread, err := r.Read(chunk)
29 - if err != nil {
31 + sofar := 0
32 +
33 + // this-chunk loop (keep reading until this chunk full)
34 + for {
35 + nread, err := r.Read(chunk[sofar:])
36 + sofar += nread
37 if err == io.EOF {
31 - if nread > 0 {
32 - out <- chunk[:nread]
38 + if sofar > 0 {
39 + // log.Infof("sending out chunk with size: %d", sofar)
40 + out <- chunk[:sofar]
41 }
42 return
43 }
36 - log.Errorf("Block split error: %s", err)
37 - return
38 - }
39 - if nread < ss.Size {
40 - chunk = chunk[:nread]
44 + if err != nil {
45 + log.Errorf("Block split error: %s", err)
46 + return
47 + }
48 + if sofar == ss.Size {
49 + // log.Infof("sending out chunk with size: %d", sofar)
50 + out <- chunk[:sofar]
51 + break // break out of this-chunk loop
52 + }
53 }
42 - out <- chunk
54 }
55 }()
56 return out
importer/chunk/splitting_test.go
+49
@@ -3,6 +3,7 @@ package chunk
3 import (
4 "bytes"
5 "crypto/rand"
6 + "io"
7 "testing"
8 )
9
@@ -54,3 +55,51 @@ func TestSizeSplitterIsDeterministic(t *testing.T) {
55 test()
56 }
57 }
58 +
59 +func TestSizeSplitterFillsChunks(t *testing.T) {
60 + if testing.Short() {
61 + t.SkipNow()
62 + }
63 +
64 + max := 10000000
65 + b := randBuf(t, max)
66 + r := &clipReader{r: bytes.NewReader(b), size: 4000}
67 + s := SizeSplitter{Size: 1024 * 256}
68 + c := s.Split(r)
69 +
70 + sofar := 0
71 + whole := make([]byte, max)
72 + for chunk := range c {
73 +
74 + bc := b[sofar : sofar+len(chunk)]
75 + if !bytes.Equal(bc, chunk) {
76 + t.Fatalf("chunk not correct: (sofar: %d) %d != %d, %v != %v", sofar, len(bc), len(chunk), bc[:100], chunk[:100])
77 + }
78 +
79 + copy(whole[sofar:], chunk)
80 +
81 + sofar += len(chunk)
82 + if sofar != max && len(chunk) < s.Size {
83 + t.Fatal("sizesplitter split at a smaller size")
84 + }
85 + }
86 +
87 + if !bytes.Equal(b, whole) {
88 + t.Fatal("splitter did not split right")
89 + }
90 +}
91 +
92 +type clipReader struct {
93 + size int
94 + r io.Reader
95 +}
96 +
97 +func (s *clipReader) Read(buf []byte) (int, error) {
98 +
99 + // clip the incoming buffer to produce smaller chunks
100 + if len(buf) > s.size {
101 + buf = buf[:s.size]
102 + }
103 +
104 + return s.r.Read(buf)
105 +}