@cryptotaxi247 / kubo / commits / 78454884d

clean up code, make it have a nicer interface

Jeromy committed Sep 15, 2014 at 04:17 UTC 78454884db316d58c37e0a6addd3a7ed155a4c69
4 files changed +67 -114
importer/importer.go
+2 -2
@@ -20,11 +20,11 @@ var ErrSizeLimitExceeded = fmt.Errorf("object size limit exceeded")
20 // NewDagFromReader constructs a Merkle DAG from the given io.Reader.
21 // size required for block construction.
22 func NewDagFromReader(r io.Reader) (*dag.Node, error) {
23 - return NewDagFromReaderWithSplitter(r, SplitterBySize(1024*512))
23 + return NewDagFromReaderWithSplitter(r, &SizeSplitter{1024 * 512})
24 }
25
26 func NewDagFromReaderWithSplitter(r io.Reader, spl BlockSplitter) (*dag.Node, error) {
27 - blkChan := spl(r)
27 + blkChan := spl.Split(r)
28 root := &dag.Node{Data: dag.FilePBData()}
29
30 for blk := range blkChan {
importer/importer_test.go
+3 -40
@@ -34,30 +34,15 @@ func TestBuildDag(t *testing.T) {
34
35 //Test where calls to read are smaller than the chunk size
36 func TestSizeBasedSplit(t *testing.T) {
37 - bs := SplitterBySize(512)
37 + bs := &SizeSplitter{512}
38 testFileConsistency(t, bs, 32*512)
39 - bs = SplitterBySize(4096)
39 + bs = &SizeSplitter{4096}
40 testFileConsistency(t, bs, 32*4096)
41
42 // Uneven offset
43 testFileConsistency(t, bs, 31*4095)
44 }
45
46 -func TestOtherSplit(t *testing.T) {
47 - //split := WhyrusleepingCantImplementRabin
48 - //testFileConsistency(t, split, 4096*64)
49 -}
50 -
51 -type testData struct{ n uint64 }
52 -
53 -func (t *testData) Read(b []byte) (int, error) {
54 - for i, _ := range b {
55 - b[i] = byte(t.n % 256)
56 - t.n++
57 - }
58 - return len(b), nil
59 -}
60 -
46 func testFileConsistency(t *testing.T, bs BlockSplitter, nbytes int) {
47 buf := new(bytes.Buffer)
48 io.CopyN(buf, rand.Reader, int64(nbytes))
@@ -95,27 +80,5 @@ func arrComp(a, b []byte) error {
80 }
81
82 func TestMaybeRabinConsistency(t *testing.T) {
98 - testFileConsistency(t, ThisMightBeRabin, 256*4096)
99 -}
100 -
101 -func TestRabinSplit(t *testing.T) {
102 -
103 - //Generate some random data
104 - nbytes := 256 * 4096
105 - buf := new(bytes.Buffer)
106 - io.CopyN(buf, rand.Reader, int64(nbytes))
107 - good := buf.Bytes()
108 -
109 - // Get block generator for random data
110 - ch := ThisMightBeRabin(buf)
111 -
112 - i := 0
113 - var blocks [][]byte
114 - for blk := range ch {
115 - if !bytes.Equal(blk, good[i:len(blk)+i]) {
116 - t.Fatalf("bad block! %v", blk[:32])
117 - }
118 - i += len(blk)
119 - blocks = append(blocks, blk)
120 - }
83 + testFileConsistency(t, NewMaybeRabin(4096), 256*4096)
84 }
importer/rabin.go
+37 -51
@@ -5,93 +5,79 @@ import (
5 "bytes"
6 "fmt"
7 "io"
8 + "math"
9 )
10
10 -//pseudocode stolen from the internet
11 -func rollhash(S []byte) {
12 - a := 10
13 - mask := 0xfff
14 - MOD := 33554383 //randomly chosen
15 - windowSize := 16
16 - an := 1
17 - rollingHash := 0
18 - for i := 0; i < windowSize; i++ {
19 - rollingHash = (rollingHash*a + int(S[i])) % MOD
20 - an = (an * a) % MOD
21 - }
22 - if rollingHash&mask == mask {
23 - // "match"
24 - fmt.Println("match")
25 - }
26 - for i := 1; i < len(S)-windowSize; i++ {
27 - rollingHash = (rollingHash*a + int(S[i+windowSize-1]) - an*int(S[i-1])) % MOD
28 - if rollingHash&mask == mask {
29 - //print "match"
30 - fmt.Println("match")
31 - }
32 - }
11 +type MaybeRabin struct {
12 + mask int
13 + windowSize int
14 +}
15 +
16 +func NewMaybeRabin(avgBlkSize int) *MaybeRabin {
17 + blkbits := uint(math.Log2(float64(avgBlkSize)))
18 + rb := new(MaybeRabin)
19 + rb.mask = (1 << blkbits) - 1
20 + rb.windowSize = 16 // probably a good number...
21 + return rb
22 }
23
35 -func ThisMightBeRabin(r io.Reader) chan []byte {
36 - out := make(chan []byte)
24 +func (mr *MaybeRabin) Split(r io.Reader) chan []byte {
25 + out := make(chan []byte, 16)
26 go func() {
27 inbuf := bufio.NewReader(r)
28 blkbuf := new(bytes.Buffer)
29
41 - // some bullshit numbers
42 - a := 10
43 - mask := 0xfff //make this smaller for smaller blocks
44 - MOD := 33554383 //randomly chosen
45 - windowSize := 16
30 + // some bullshit numbers i made up
31 + a := 10 // honestly, no idea what this is
32 + MOD := 33554383 // randomly chosen (seriously)
33 an := 1
34 rollingHash := 0
35
49 - window := make([]byte, windowSize)
50 - get := func(i int) int { return int(window[i%len(window)]) }
51 - set := func(i int, val byte) { window[i%len(window)] = val }
36 + // Window is a circular buffer
37 + window := make([]byte, mr.windowSize)
38 + push := func(i int, val byte) (outval int) {
39 + outval = int(window[i%len(window)])
40 + window[i%len(window)] = val
41 + return
42 + }
43 +
44 + // Duplicate byte slice
45 dup := func(b []byte) []byte {
46 d := make([]byte, len(b))
47 copy(d, b)
48 return d
49 }
50
51 + // Fill up the window
52 i := 0
59 - for ; i < windowSize; i++ {
53 + for ; i < mr.windowSize; i++ {
54 b, err := inbuf.ReadByte()
55 if err != nil {
56 fmt.Println(err)
57 return
58 }
59 blkbuf.WriteByte(b)
66 - window[i] = b
60 + push(i, b)
61 rollingHash = (rollingHash*a + int(b)) % MOD
62 an = (an * a) % MOD
63 }
70 - /* This is too short for a block
71 - if rollingHash&mask == mask {
72 - // "match"
73 - fmt.Println("match")
74 - }
75 - */
64 +
65 for ; true; i++ {
66 b, err := inbuf.ReadByte()
67 if err != nil {
68 break
69 }
81 - outval := get(i)
82 - set(i, b)
70 + outval := push(i, b)
71 blkbuf.WriteByte(b)
84 - rollingHash = (rollingHash*a + get(i) - an*outval) % MOD
85 - if rollingHash&mask == mask {
86 - //print "match"
72 + rollingHash = (rollingHash*a + int(b) - an*outval) % MOD
73 + if rollingHash&mr.mask == mr.mask {
74 out <- dup(blkbuf.Bytes())
75 blkbuf.Reset()
76 }
90 - peek, err := inbuf.Peek(windowSize)
91 - if err != nil {
92 - break
93 - }
94 - if len(peek) != windowSize {
77 +
78 + // Check if there are enough remaining
79 + peek, err := inbuf.Peek(mr.windowSize)
80 + if err != nil || len(peek) != mr.windowSize {
81 break
82 }
83 }
importer/splitting.go
+25 -21
@@ -6,29 +6,33 @@ import (
6 u "github.com/jbenet/go-ipfs/util"
7 )
8
9 -type BlockSplitter func(io.Reader) chan []byte
9 +type BlockSplitter interface {
10 + Split(io.Reader) chan []byte
11 +}
12 +
13 +type SizeSplitter struct {
14 + Size int
15 +}
16
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)
17 +func (ss *SizeSplitter) Split(r io.Reader) chan []byte {
18 + out := make(chan []byte)
19 + go func() {
20 + defer close(out)
21 + for {
22 + chunk := make([]byte, ss.Size)
23 + nread, err := r.Read(chunk)
24 + if err != nil {
25 + if err == io.EOF {
26 return
27 }
26 - if nread < n {
27 - chunk = chunk[:nread]
28 - }
29 - out <- chunk
28 + u.PErr("block split error: %v\n", err)
29 + return
30 + }
31 + if nread < ss.Size {
32 + chunk = chunk[:nread]
33 }
31 - }(n)
32 - return out
33 - }
34 + out <- chunk
35 + }
36 + }()
37 + return out
38 }