implement something like rabin fingerprinting
Jeromy committed
Sep 15, 2014 at 02:04 UTC
1a7c083850d2d34031304cce7286e3bd4eca5ef4
4 files changed
+221
-113
importer/importer_test.go
+78
-15
@@ -3,42 +3,66 @@ package importer
3
import (
4
"bytes"
5
"crypto/rand"
6
+ "fmt"
7
"io"
8
"io/ioutil"
9
+ "os"
10
"testing"
11
12
dag "github.com/jbenet/go-ipfs/merkledag"
13
)
14
13
-func TestFileConsistency(t *testing.T) {
14
- buf := new(bytes.Buffer)
15
- io.CopyN(buf, rand.Reader, 512*32)
16
- should := buf.Bytes()
17
- nd, err := NewDagFromReaderWithSplitter(buf, SplitterBySize(512))
15
+func TestBuildDag(t *testing.T) {
16
+ td := os.TempDir()
17
+ fi, err := os.Create(td + "/tmpfi")
18
if err != nil {
19
t.Fatal(err)
20
}
21
- r, err := dag.NewDagReader(nd, nil)
21
+
22
+ _, err = io.CopyN(fi, rand.Reader, 1024*1024)
23
if err != nil {
24
t.Fatal(err)
25
}
26
26
- out, err := ioutil.ReadAll(r)
27
+ fi.Close()
28
+
29
+ _, err = NewDagFromFile(td + "/tmpfi")
30
if err != nil {
31
t.Fatal(err)
32
}
33
+}
34
31
- if !bytes.Equal(out, should) {
32
- t.Fatal("Output not the same as input.")
35
+//Test where calls to read are smaller than the chunk size
36
+func TestSizeBasedSplit(t *testing.T) {
37
+ bs := SplitterBySize(512)
38
+ testFileConsistency(t, bs, 32*512)
39
+ bs = SplitterBySize(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
36
-//Test where calls to read are smaller than the chunk size
37
-func TestFileConsistencyLargeBlocks(t *testing.T) {
61
+func testFileConsistency(t *testing.T, bs BlockSplitter, nbytes int) {
62
buf := new(bytes.Buffer)
39
- io.CopyN(buf, rand.Reader, 4096*32)
63
+ io.CopyN(buf, rand.Reader, int64(nbytes))
64
should := buf.Bytes()
41
- nd, err := NewDagFromReaderWithSplitter(buf, SplitterBySize(4096))
65
+ nd, err := NewDagFromReaderWithSplitter(buf, bs)
66
if err != nil {
67
t.Fatal(err)
68
}
@@ -52,7 +76,46 @@ func TestFileConsistencyLargeBlocks(t *testing.T) {
76
t.Fatal(err)
77
}
78
55
- if !bytes.Equal(out, should) {
56
- t.Fatal("Output not the same as input.")
79
+ err = arrComp(out, should)
80
+ if err != nil {
81
+ t.Fatal(err)
82
+ }
83
+}
84
+
85
+func arrComp(a, b []byte) error {
86
+ if len(a) != len(b) {
87
+ return fmt.Errorf("Arrays differ in length. %d != %d", len(a), len(b))
88
+ }
89
+ for i, v := range a {
90
+ if v != b[i] {
91
+ return fmt.Errorf("Arrays differ at index: %d", i)
92
+ }
93
+ }
94
+ return nil
95
+}
96
+
97
+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
}
121
}
importer/rabin.go
new
+143
@@ -0,0 +1,143 @@
1
+package importer
2
+
3
+import (
4
+ "bufio"
5
+ "bytes"
6
+ "fmt"
7
+ "io"
8
+)
9
+
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
+ }
33
+}
34
+
35
+func ThisMightBeRabin(r io.Reader) chan []byte {
36
+ out := make(chan []byte)
37
+ go func() {
38
+ inbuf := bufio.NewReader(r)
39
+ blkbuf := new(bytes.Buffer)
40
+
41
+ // some bullshit numbers
42
+ a := 10
43
+ mask := 0xfff //make this smaller for smaller blocks
44
+ MOD := 33554383 //randomly chosen
45
+ windowSize := 16
46
+ an := 1
47
+ rollingHash := 0
48
+
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 }
52
+ dup := func(b []byte) []byte {
53
+ d := make([]byte, len(b))
54
+ copy(d, b)
55
+ return d
56
+ }
57
+
58
+ i := 0
59
+ for ; i < windowSize; i++ {
60
+ b, err := inbuf.ReadByte()
61
+ if err != nil {
62
+ fmt.Println(err)
63
+ return
64
+ }
65
+ blkbuf.WriteByte(b)
66
+ window[i] = b
67
+ rollingHash = (rollingHash*a + int(b)) % MOD
68
+ an = (an * a) % MOD
69
+ }
70
+ /* This is too short for a block
71
+ if rollingHash&mask == mask {
72
+ // "match"
73
+ fmt.Println("match")
74
+ }
75
+ */
76
+ for ; true; i++ {
77
+ b, err := inbuf.ReadByte()
78
+ if err != nil {
79
+ break
80
+ }
81
+ outval := get(i)
82
+ set(i, b)
83
+ blkbuf.WriteByte(b)
84
+ rollingHash = (rollingHash*a + get(i) - an*outval) % MOD
85
+ if rollingHash&mask == mask {
86
+ //print "match"
87
+ out <- dup(blkbuf.Bytes())
88
+ blkbuf.Reset()
89
+ }
90
+ peek, err := inbuf.Peek(windowSize)
91
+ if err != nil {
92
+ break
93
+ }
94
+ if len(peek) != windowSize {
95
+ break
96
+ }
97
+ }
98
+ io.Copy(blkbuf, inbuf)
99
+ out <- blkbuf.Bytes()
100
+ close(out)
101
+ }()
102
+ return out
103
+}
104
+
105
+/*
106
+func WhyrusleepingCantImplementRabin(r io.Reader) chan []byte {
107
+ out := make(chan []byte, 4)
108
+ go func() {
109
+ buf := bufio.NewReader(r)
110
+ blkbuf := new(bytes.Buffer)
111
+ window := make([]byte, 16)
112
+ var val uint64
113
+ prime := uint64(61)
114
+
115
+ get := func(i int) uint64 {
116
+ return uint64(window[i%len(window)])
117
+ }
118
+
119
+ set := func(i int, val byte) {
120
+ window[i%len(window)] = val
121
+ }
122
+
123
+ for i := 0; ; i++ {
124
+ curb, err := buf.ReadByte()
125
+ if err != nil {
126
+ break
127
+ }
128
+ set(i, curb)
129
+ blkbuf.WriteByte(curb)
130
+
131
+ hash := md5.Sum(window)
132
+ if hash[0] == 0 && hash[1] == 0 {
133
+ out <- blkbuf.Bytes()
134
+ blkbuf.Reset()
135
+ }
136
+ }
137
+ out <- blkbuf.Bytes()
138
+ close(out)
139
+ }()
140
+
141
+ return out
142
+}
143
+*/
importer/split_test.go
deleted
-48
@@ -1,48 +0,0 @@
1
-package importer
2
-
3
-import (
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
- split := Rabin(buf)
13
-
14
- if len(split) == 1 {
15
- t.Fatal("No split occurred!")
16
- }
17
-
18
- min := 2 << 15
19
- max := 0
20
-
21
- mxcount := 0
22
-
23
- n := 0
24
- for _, b := range split {
25
- if !bytes.Equal(b, buf[n:n+len(b)]) {
26
- t.Fatal("Split lost data!")
27
- }
28
- n += len(b)
29
-
30
- if len(b) < min {
31
- min = len(b)
32
- }
33
-
34
- if len(b) > max {
35
- max = len(b)
36
- }
37
-
38
- if len(b) == 16384 {
39
- mxcount++
40
- }
41
- }
42
-
43
- if n != len(buf) {
44
- t.Fatal("missing some bytes!")
45
- }
46
- t.Log(len(split))
47
- t.Log(min, max, mxcount)
48
-}
importer/splitting.go
-50
@@ -32,53 +32,3 @@ func SplitterBySize(n int) BlockSplitter {
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
-//Rabin Fingerprinting for file chunking
38
-func Rabin(b []byte) [][]byte {
39
- var out [][]byte
40
- windowsize := uint64(48)
41
- chunkMax := 1024 * 16
42
- minBlkSize := 2048
43
- blkBegI := 0
44
- prime := uint64(61)
45
-
46
- var poly uint64
47
- var curchecksum uint64
48
-
49
- // Smaller than a window? Get outa here!
50
- if len(b) <= int(windowsize) {
51
- return [][]byte{b}
52
- }
53
-
54
- i := 0
55
- for n := i; i < n+int(windowsize); i++ {
56
- cur := uint64(b[i])
57
- curchecksum = (curchecksum * prime) + cur
58
- poly = (poly * prime) + cur
59
- }
60
-
61
- for ; i < len(b); i++ {
62
- cur := uint64(b[i])
63
- curchecksum = (curchecksum * prime) + cur
64
- poly = (poly * prime) + cur
65
- curchecksum -= (uint64(b[i-1]) * prime)
66
-
67
- if i-blkBegI >= chunkMax {
68
- // push block
69
- out = append(out, b[blkBegI:i])
70
- blkBegI = i
71
- }
72
-
73
- // first 13 bits of polynomial are 0
74
- if poly%8192 == 0 && i-blkBegI >= minBlkSize {
75
- // push block
76
- out = append(out, b[blkBegI:i])
77
- blkBegI = i
78
- }
79
- }
80
- if i > blkBegI {
81
- out = append(out, b[blkBegI:])
82
- }
83
- return out
84
-}