@cryptotaxi247 / kubo / commits / 7cbfbbc0a

add blockset and bloomfilter and beginnings of pinning service

Jeromy committed Oct 16, 2014 at 00:20 UTC 7cbfbbc0ad7e47fe4ed077cbaef4c817887f81a9
7 files changed +288
blocks/bloom/filter.go new
+70
@@ -0,0 +1,70 @@
1 +package bloom
2 +
3 +import (
4 + "fmt"
5 + "hash"
6 + "hash/adler32"
7 + "hash/crc32"
8 + "hash/fnv"
9 + "math/big"
10 +)
11 +
12 +type Filter interface {
13 + Add([]byte)
14 + Find([]byte) bool
15 +}
16 +
17 +func BasicFilter() Filter {
18 + // Non crypto hashes, because speed
19 + return NewFilter(2048, adler32.New(), fnv.New32(), crc32.NewIEEE())
20 +}
21 +
22 +func NewFilter(size int, hashes ...hash.Hash) Filter {
23 + return &filter{
24 + filter: make([]byte, size),
25 + hashes: hashes,
26 + }
27 +}
28 +
29 +type filter struct {
30 + filter []byte
31 + hashes []hash.Hash
32 +}
33 +
34 +func (f *filter) Add(k []byte) {
35 + for _, h := range f.hashes {
36 + i := bytesMod(h.Sum(k), int64(len(f.filter)*8))
37 + f.setBit(i)
38 + }
39 +}
40 +
41 +func (f *filter) Find(k []byte) bool {
42 + for _, h := range f.hashes {
43 + i := bytesMod(h.Sum(k), int64(len(f.filter)*8))
44 + if !f.getBit(i) {
45 + return false
46 + }
47 + }
48 + return true
49 +}
50 +
51 +func (f *filter) setBit(i int64) {
52 + fmt.Printf("setting bit %d\n", i)
53 + f.filter[i/8] |= (1 << byte(i%8))
54 +}
55 +
56 +func (f *filter) getBit(i int64) bool {
57 + fmt.Printf("getting bit %d\n", i)
58 + return f.filter[i/8]&(1<<byte(i%8)) != 0
59 +}
60 +
61 +func bytesMod(b []byte, modulo int64) int64 {
62 + i := big.NewInt(0)
63 + i = i.SetBytes(b)
64 +
65 + bigmod := big.NewInt(int64(modulo))
66 + result := big.NewInt(0)
67 + result.Mod(i, bigmod)
68 +
69 + return result.Int64()
70 +}
blocks/bloom/filter.proto new
+10
@@ -0,0 +1,10 @@
1 +package bloom;
2 +
3 +message PackedFilter {
4 + enum HashType {
5 +
6 + }
7 + optional bool compressed;
8 + optional bytes data;
9 + repeated HashType hashes;
10 +}
blocks/bloom/filter_test.go new
+30
@@ -0,0 +1,30 @@
1 +package bloom
2 +
3 +import "testing"
4 +
5 +func TestFilter(t *testing.T) {
6 + f := BasicFilter()
7 + keys := [][]byte{
8 + []byte("hello"),
9 + []byte("fish"),
10 + []byte("ipfsrocks"),
11 + }
12 +
13 + f.Add(keys[0])
14 + if !f.Find(keys[0]) {
15 + t.Fatal("Failed to find single inserted key!")
16 + }
17 +
18 + f.Add(keys[1])
19 + if !f.Find(keys[1]) {
20 + t.Fatal("Failed to find key!")
21 + }
22 +
23 + f.Add(keys[2])
24 +
25 + for _, k := range keys {
26 + if !f.Find(k) {
27 + t.Fatal("Couldnt find one of three keys")
28 + }
29 + }
30 +}
blocks/set/set.go new
+42
@@ -0,0 +1,42 @@
1 +package set
2 +
3 +import (
4 + "github.com/jbenet/go-ipfs/blocks/bloom"
5 + "github.com/jbenet/go-ipfs/util"
6 +)
7 +
8 +type BlockSet interface {
9 + AddBlock(util.Key)
10 + RemoveBlock(util.Key)
11 + HasKey(util.Key) bool
12 + GetBloomFilter() bloom.Filter
13 +}
14 +
15 +func NewSimpleBlockSet() BlockSet {
16 + return &simpleBlockSet{blocks: make(map[util.Key]struct{})}
17 +}
18 +
19 +type simpleBlockSet struct {
20 + blocks map[util.Key]struct{}
21 +}
22 +
23 +func (b *simpleBlockSet) AddBlock(k util.Key) {
24 + b.blocks[k] = struct{}{}
25 +}
26 +
27 +func (b *simpleBlockSet) RemoveBlock(k util.Key) {
28 + delete(b.blocks, k)
29 +}
30 +
31 +func (b *simpleBlockSet) HasKey(k util.Key) bool {
32 + _, has := b.blocks[k]
33 + return has
34 +}
35 +
36 +func (b *simpleBlockSet) GetBloomFilter() bloom.Filter {
37 + f := bloom.BasicFilter()
38 + for k, _ := range b.blocks {
39 + f.Add([]byte(k))
40 + }
41 + return f
42 +}
merkledag/merkledag.go
+8
@@ -59,6 +59,14 @@ func MakeLink(n *Node) (*Link, error) {
59 }, nil
60 }
61
62 +func (l *Link) GetNode(serv *DAGService) (*Node, error) {
63 + if l.Node != nil {
64 + return l.Node, nil
65 + }
66 +
67 + return serv.Get(u.Key(l.Hash))
68 +}
69 +
70 // AddNodeLink adds a link to another node.
71 func (n *Node) AddNodeLink(name string, that *Node) error {
72 lnk, err := MakeLink(that)
pin/pin.go new
+88
@@ -0,0 +1,88 @@
1 +package pin
2 +
3 +import (
4 + ds "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/datastore.go"
5 + "github.com/jbenet/go-ipfs/blocks/set"
6 + mdag "github.com/jbenet/go-ipfs/merkledag"
7 + "github.com/jbenet/go-ipfs/util"
8 +)
9 +
10 +type Pinner interface {
11 + Pin(*mdag.Node, bool) error
12 + Unpin(util.Key, bool) error
13 +}
14 +
15 +type pinner struct {
16 + recursePin set.BlockSet
17 + directPin set.BlockSet
18 + indirPin set.BlockSet
19 + dserv *mdag.DAGService
20 +}
21 +
22 +func NewPinner(dstore ds.Datastore, serv *mdag.DAGService) Pinner {
23 + return &pinner{
24 + recursePin: set.NewSimpleBlockSet(),
25 + directPin: set.NewSimpleBlockSet(),
26 + indirPin: NewRefCountBlockSet(),
27 + dserv: serv,
28 + }
29 +}
30 +
31 +func (p *pinner) Pin(node *mdag.Node, recurse bool) error {
32 + k, err := node.Key()
33 + if err != nil {
34 + return err
35 + }
36 +
37 + if recurse {
38 + if p.recursePin.HasKey(k) {
39 + return nil
40 + }
41 +
42 + p.recursePin.AddBlock(k)
43 +
44 + err := p.pinLinks(node)
45 + if err != nil {
46 + return err
47 + }
48 + } else {
49 + p.directPin.AddBlock(k)
50 + }
51 + return nil
52 +}
53 +
54 +func (p *pinner) Unpin(k util.Key, recurse bool) error {
55 + panic("not yet implemented!")
56 + return nil
57 +}
58 +
59 +func (p *pinner) pinIndirectRecurse(node *mdag.Node) error {
60 + k, err := node.Key()
61 + if err != nil {
62 + return err
63 + }
64 +
65 + p.indirPin.AddBlock(k)
66 + return p.pinLinks(node)
67 +}
68 +
69 +func (p *pinner) pinLinks(node *mdag.Node) error {
70 + for _, l := range node.Links {
71 + subnode, err := l.GetNode(p.dserv)
72 + if err != nil {
73 + // TODO: Maybe just log and continue?
74 + return err
75 + }
76 + err = p.pinIndirectRecurse(subnode)
77 + if err != nil {
78 + return err
79 + }
80 + }
81 + return nil
82 +}
83 +
84 +func (p *pinner) IsPinned(key util.Key) bool {
85 + return p.recursePin.HasKey(key) ||
86 + p.directPin.HasKey(key) ||
87 + p.indirPin.HasKey(key)
88 +}
pin/refset.go new
+40
@@ -0,0 +1,40 @@
1 +package pin
2 +
3 +import (
4 + "github.com/jbenet/go-ipfs/blocks/bloom"
5 + "github.com/jbenet/go-ipfs/blocks/set"
6 + "github.com/jbenet/go-ipfs/util"
7 +)
8 +
9 +type refCntBlockSet struct {
10 + blocks map[util.Key]int
11 +}
12 +
13 +func NewRefCountBlockSet() set.BlockSet {
14 + return &refCntBlockSet{blocks: make(map[util.Key]int)}
15 +}
16 +
17 +func (r *refCntBlockSet) AddBlock(k util.Key) {
18 + r.blocks[k]++
19 +}
20 +
21 +func (r *refCntBlockSet) RemoveBlock(k util.Key) {
22 + v, ok := r.blocks[k]
23 + if !ok {
24 + return
25 + }
26 + if v <= 1 {
27 + delete(r.blocks, k)
28 + } else {
29 + r.blocks[k] = v - 1
30 + }
31 +}
32 +
33 +func (r *refCntBlockSet) HasKey(k util.Key) bool {
34 + _, ok := r.blocks[k]
35 + return ok
36 +}
37 +
38 +func (r *refCntBlockSet) GetBloomFilter() bloom.Filter {
39 + return nil
40 +}