@cryptotaxi247 / kubo / commits / 9338caa9d

working on making importer not break on large files

Jeromy committed Aug 31, 2014 at 10:10 UTC 9338caa9d8df155f2584b15ea9b1895ceb3c94a1
3 files changed +46 -33
importer/importer.go
+11 -27
@@ -3,7 +3,6 @@ package importer
3 import (
4 "fmt"
5 "io"
6 - "io/ioutil"
6 "os"
7
8 dag "github.com/jbenet/go-ipfs/merkledag"
@@ -20,32 +19,17 @@ var ErrSizeLimitExceeded = fmt.Errorf("object size limit exceeded")
19
20 // NewDagFromReader constructs a Merkle DAG from the given io.Reader.
21 // size required for block construction.
23 -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
22 +func NewDagFromReader(r io.Reader) (*dag.Node, error) {
23 + blkChan := SplitterBySize(1024 * 512)(r)
24 + root := &dag.Node{}
25 +
26 + for blk := range blkChan {
27 + child := &dag.Node{Data: blk}
28 + err := root.AddNodeLink("", child)
29 + if err != nil {
30 + return nil, err
31 + }
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 {
40 - return nil, err
41 - }
42 -
43 - if int64(len(buf)) > BlockSizeLimit {
44 - return nil, ErrSizeLimitExceeded // lying punk.
45 - }
46 -
47 - root := &dag.Node{Data: buf}
48 - // no children for now because not block splitting yet
33 return root, nil
34 }
35
@@ -66,5 +50,5 @@ func NewDagFromFile(fpath string) (*dag.Node, error) {
50 }
51 defer f.Close()
52
69 - return NewDagFromReader(f, stat.Size())
53 + return NewDagFromReader(f)
54 }
importer/split_test.go
+2 -4
@@ -1,15 +1,14 @@
1 package importer
2
3 import (
4 - "testing"
5 - "crypto/rand"
4 "bytes"
5 + "crypto/rand"
6 + "testing"
7 )
8
9 func TestDataSplitting(t *testing.T) {
10 buf := make([]byte, 16*1024*1024)
11 rand.Read(buf)
12 -
12 split := Rabin(buf)
13
14 if len(split) == 1 {
@@ -47,4 +46,3 @@ func TestDataSplitting(t *testing.T) {
46 t.Log(len(split))
47 t.Log(min, max, mxcount)
48 }
50 -
importer/splitting.go
+33 -2
@@ -1,6 +1,37 @@
1 package importer
2
3 -type BlockSplitter func([]byte) [][]byte
3 +import (
4 + "io"
5 +
6 + u "github.com/jbenet/go-ipfs/util"
7 +)
8 +
9 +type BlockSplitter func(io.Reader) chan []byte
10 +
11 +func SplitterBySize(n int) BlockSplitter {
12 + return func(r io.Reader) chan []byte {
13 + out := make(chan []byte)
14 + go func(n int) {
15 + defer close(out)
16 + for {
17 + chunk := make([]byte, n)
18 + nread, err := r.Read(chunk)
19 + if err != nil {
20 + if err == io.EOF {
21 + return
22 + }
23 + u.PErr("block split error: %v\n", err)
24 + return
25 + }
26 + if nread < n {
27 + chunk = chunk[:n]
28 + }
29 + out <- chunk
30 + }
31 + }(n)
32 + return out
33 + }
34 +}
35
36 // TODO: this should take a reader, not a byte array. what if we're splitting a 3TB file?
37 func Rabin(b []byte) [][]byte {
@@ -39,7 +70,7 @@ func Rabin(b []byte) [][]byte {
70 }
71
72 // first 13 bits of polynomial are 0
42 - if poly % 8192 == 0 && i-blk_beg_i >= min_blk_size {
73 + if poly%8192 == 0 && i-blk_beg_i >= min_blk_size {
74 // push block
75 out = append(out, b[blk_beg_i:i])
76 blk_beg_i = i