@cryptotaxi247 / kubo / commits / ba4ea139b

address concerns from PR

Jeromy committed Jan 27, 2015 at 01:28 UTC ba4ea139b9a6b3d10bba7ffa8b33b65ea5366a0a
4 files changed +87 -12
importer/importer_test.go
+46
@@ -6,6 +6,7 @@ import (
6 "fmt"
7 "io"
8 "io/ioutil"
9 + mrand "math/rand"
10 "os"
11 "testing"
12
@@ -307,6 +308,51 @@ func TestSeekToAlmostBegin(t *testing.T) {
308 }
309 }
310
311 +func TestSeekingStress(t *testing.T) {
312 + nbytes := int64(1024 * 1024)
313 + should := make([]byte, nbytes)
314 + u.NewTimeSeededRand().Read(should)
315 +
316 + read := bytes.NewReader(should)
317 + dnp := getDagservAndPinner(t)
318 + nd, err := BuildDagFromReader(read, dnp.ds, dnp.mp, &chunk.SizeSplitter{1000})
319 + if err != nil {
320 + t.Fatal(err)
321 + }
322 +
323 + rs, err := uio.NewDagReader(context.TODO(), nd, dnp.ds)
324 + if err != nil {
325 + t.Fatal(err)
326 + }
327 +
328 + testbuf := make([]byte, nbytes)
329 + for i := 0; i < 50; i++ {
330 + offset := mrand.Intn(int(nbytes))
331 + l := int(nbytes) - offset
332 + n, err := rs.Seek(int64(offset), os.SEEK_SET)
333 + if err != nil {
334 + t.Fatal(err)
335 + }
336 + if n != int64(offset) {
337 + t.Fatal("Seek failed to move to correct position")
338 + }
339 +
340 + nread, err := rs.Read(testbuf[:l])
341 + if err != nil {
342 + t.Fatal(err)
343 + }
344 + if nread != l {
345 + t.Fatal("Failed to read enough bytes")
346 + }
347 +
348 + err = arrComp(testbuf[:l], should[offset:offset+l])
349 + if err != nil {
350 + t.Fatal(err)
351 + }
352 + }
353 +
354 +}
355 +
356 func TestSeekingConsistency(t *testing.T) {
357 nbytes := int64(128 * 1024)
358 should := make([]byte, nbytes)
merkledag/merkledag.go
+2
@@ -177,6 +177,8 @@ func (ds *dagService) GetDAG(ctx context.Context, root *Node) []NodeGetter {
177 return ds.GetNodes(ctx, keys)
178 }
179
180 +// GetNodes returns an array of 'NodeGetter' promises, with each corresponding
181 +// to the key with the same index as the passed in keys
182 func (ds *dagService) GetNodes(ctx context.Context, keys []u.Key) []NodeGetter {
183 promises := make([]NodeGetter, len(keys))
184 sendChans := make([]chan<- *Node, len(keys))
unixfs/io/dagmodifier_test.go
+4 -4
@@ -39,7 +39,7 @@ func getNode(t *testing.T, dserv mdag.DAGService, size int64) ([]byte, *mdag.Nod
39 t.Fatal(err)
40 }
41
42 - dr, err := NewDagReader(context.TODO(), node, dserv)
42 + dr, err := NewDagReader(context.Background(), node, dserv)
43 if err != nil {
44 t.Fatal(err)
45 }
@@ -76,7 +76,7 @@ func testModWrite(t *testing.T, beg, size uint64, orig []byte, dm *DagModifier)
76 t.Fatal(err)
77 }
78
79 - rd, err := NewDagReader(context.TODO(), nd, dm.dagserv)
79 + rd, err := NewDagReader(context.Background(), nd, dm.dagserv)
80 if err != nil {
81 t.Fatal(err)
82 }
@@ -174,7 +174,7 @@ func TestMultiWrite(t *testing.T) {
174 t.Fatal(err)
175 }
176
177 - read, err := NewDagReader(context.TODO(), nd, dserv)
177 + read, err := NewDagReader(context.Background(), nd, dserv)
178 if err != nil {
179 t.Fatal(err)
180 }
@@ -216,7 +216,7 @@ func TestMultiWriteCoal(t *testing.T) {
216 t.Fatal(err)
217 }
218
219 - read, err := NewDagReader(context.TODO(), nd, dserv)
219 + read, err := NewDagReader(context.Background(), nd, dserv)
220 if err != nil {
221 t.Fatal(err)
222 }
unixfs/io/dagreader.go
+35 -8
@@ -18,19 +18,31 @@ var ErrIsDir = errors.New("this dag node is a directory")
18
19 // DagReader provides a way to easily read the data contained in a dag.
20 type DagReader struct {
21 - serv mdag.DAGService
22 - node *mdag.Node
23 - pbdata *ftpb.Data
24 - buf ReadSeekCloser
25 - promises []mdag.NodeGetter
21 + serv mdag.DAGService
22 +
23 + // the node being read
24 + node *mdag.Node
25 +
26 + // cached protobuf structure from node.Data
27 + pbdata *ftpb.Data
28 +
29 + // the current data buffer to be read from
30 + // will either be a bytes.Reader or a child DagReader
31 + buf ReadSeekCloser
32 +
33 + // NodeGetters for each of 'nodes' child links
34 + promises []mdag.NodeGetter
35 +
36 + // the index of the child link currently being read from
37 linkPosition int
27 - offset int64
38 +
39 + // current offset for the read head within the 'file'
40 + offset int64
41
42 // Our context
43 ctx context.Context
44
32 - // Context for children
33 - fctx context.Context
45 + // context cancel for children
46 cancel func()
47 }
48
@@ -145,6 +157,8 @@ func (dr *DagReader) Close() error {
157 return nil
158 }
159
160 +// Seek implements io.Seeker, and will seek to a given offset in the file
161 +// interface matches standard unix seek
162 func (dr *DagReader) Seek(offset int64, whence int) (int64, error) {
163 switch whence {
164 case os.SEEK_SET:
@@ -152,18 +166,26 @@ func (dr *DagReader) Seek(offset int64, whence int) (int64, error) {
166 return -1, errors.New("Invalid offset")
167 }
168
169 + // Grab cached protobuf object (solely to make code look cleaner)
170 pb := dr.pbdata
171 +
172 + // left represents the number of bytes remaining to seek to (from beginning)
173 left := offset
174 if int64(len(pb.Data)) > offset {
175 + // Close current buf to close potential child dagreader
176 dr.buf.Close()
177 dr.buf = NewRSNCFromBytes(pb.GetData()[offset:])
178 +
179 + // start reading links from the beginning
180 dr.linkPosition = 0
181 dr.offset = offset
182 return offset, nil
183 } else {
184 + // skip past root block data
185 left -= int64(len(pb.Data))
186 }
187
188 + // iterate through links and find where we need to be
189 for i := 0; i < len(pb.Blocksizes); i++ {
190 if pb.Blocksizes[i] > uint64(left) {
191 dr.linkPosition = i
@@ -173,15 +195,19 @@ func (dr *DagReader) Seek(offset int64, whence int) (int64, error) {
195 }
196 }
197
198 + // start sub-block request
199 err := dr.precalcNextBuf()
200 if err != nil {
201 return 0, err
202 }
203
204 + // set proper offset within child readseeker
205 n, err := dr.buf.Seek(left, os.SEEK_SET)
206 if err != nil {
207 return -1, err
208 }
209 +
210 + // sanity
211 left -= n
212 if left != 0 {
213 return -1, errors.New("failed to seek properly")
@@ -201,6 +227,7 @@ func (dr *DagReader) Seek(offset int64, whence int) (int64, error) {
227 return 0, nil
228 }
229
230 +// readSeekNopCloser wraps a bytes.Reader to implement ReadSeekCloser
231 type readSeekNopCloser struct {
232 *bytes.Reader
233 }