implement dagmodifier and tests.
Jeromy committed
Oct 6, 2014 at 23:49 UTC
3591e10b2e4a9f40ea8c33984933d27bd0354388
9 files changed
+489
-37
.gitignore
+1
@@ -3,3 +3,4 @@
3
*.out
4
*.test
5
*.orig
6
+*~
blockservice/blockservice.go
+2
-2
@@ -36,7 +36,7 @@ func NewBlockService(d ds.Datastore, rem exchange.Interface) (*BlockService, err
36
// AddBlock adds a particular block to the service, Putting it into the datastore.
37
func (s *BlockService) AddBlock(b *blocks.Block) (u.Key, error) {
38
k := b.Key()
39
- log.Debug("storing [%s] in datastore", k)
39
+ log.Debug("blockservice: storing [%s] in datastore", k.Pretty())
40
// TODO(brian): define a block datastore with a Put method which accepts a
41
// block parameter
42
err := s.Datastore.Put(k.DsKey(), b.Data)
@@ -53,7 +53,7 @@ func (s *BlockService) AddBlock(b *blocks.Block) (u.Key, error) {
53
// GetBlock retrieves a particular block from the service,
54
// Getting it from the datastore using the key (hash).
55
func (s *BlockService) GetBlock(k u.Key) (*blocks.Block, error) {
56
- log.Debug("BlockService GetBlock: '%s'", k)
56
+ log.Debug("BlockService GetBlock: '%s'", k.Pretty())
57
datai, err := s.Datastore.Get(k.DsKey())
58
if err == nil {
59
log.Debug("Blockservice: Got data in datastore.")
importer/dagwriter/dagmodifier.go
new
+190
@@ -0,0 +1,190 @@
1
+package dagwriter
2
+
3
+import (
4
+ "errors"
5
+
6
+ "code.google.com/p/goprotobuf/proto"
7
+
8
+ imp "github.com/jbenet/go-ipfs/importer"
9
+ ft "github.com/jbenet/go-ipfs/importer/format"
10
+ mdag "github.com/jbenet/go-ipfs/merkledag"
11
+ u "github.com/jbenet/go-ipfs/util"
12
+)
13
+
14
+// DagModifier is the only struct licensed and able to correctly
15
+// perform surgery on a DAG 'file'
16
+// Dear god, please rename this to something more pleasant
17
+type DagModifier struct {
18
+ dagserv *mdag.DAGService
19
+ curNode *mdag.Node
20
+
21
+ pbdata *ft.PBData
22
+}
23
+
24
+func NewDagModifier(from *mdag.Node, serv *mdag.DAGService) (*DagModifier, error) {
25
+ pbd, err := ft.FromBytes(from.Data)
26
+ if err != nil {
27
+ return nil, err
28
+ }
29
+
30
+ return &DagModifier{
31
+ curNode: from.Copy(),
32
+ dagserv: serv,
33
+ pbdata: pbd,
34
+ }, nil
35
+}
36
+
37
+// WriteAt will modify a dag file in place
38
+// NOTE: it currently assumes only a single level of indirection
39
+func (dm *DagModifier) WriteAt(b []byte, offset uint64) (int, error) {
40
+
41
+ // Check bounds
42
+ if dm.pbdata.GetFilesize() < offset {
43
+ return 0, errors.New("Attempted to perform write starting past end of file")
44
+ }
45
+
46
+ // This shouldnt be necessary if we do subblocks sizes properly
47
+ newsize := dm.pbdata.GetFilesize()
48
+ if uint64(len(b))+offset > dm.pbdata.GetFilesize() {
49
+ newsize = uint64(len(b)) + offset
50
+ }
51
+
52
+ // First need to find where we are writing at
53
+ end := uint64(len(b)) + offset
54
+ zeroblocklen := uint64(len(dm.pbdata.Data))
55
+ origlen := len(b)
56
+
57
+ if end <= zeroblocklen {
58
+ log.Debug("Writing into zero block.")
59
+ // Replacing zeroeth data block (embedded in the root node)
60
+ //TODO: check chunking here
61
+ copy(dm.pbdata.Data[offset:], b)
62
+ return len(b), nil
63
+ }
64
+
65
+ // Find where write should start
66
+ var traversed uint64
67
+ startsubblk := len(dm.pbdata.Blocksizes)
68
+ if offset < zeroblocklen {
69
+ dm.pbdata.Data = dm.pbdata.Data[:offset]
70
+ startsubblk = 0
71
+ } else {
72
+ traversed = uint64(zeroblocklen)
73
+ for i, size := range dm.pbdata.Blocksizes {
74
+ if uint64(offset) < traversed+size {
75
+ log.Debug("Starting mod at block %d. [%d < %d + %d]", i, offset, traversed, size)
76
+ // Here is where we start
77
+ startsubblk = i
78
+ lnk := dm.curNode.Links[i]
79
+ node, err := dm.dagserv.Get(u.Key(lnk.Hash))
80
+ if err != nil {
81
+ return 0, err
82
+ }
83
+ data, err := ft.UnwrapData(node.Data)
84
+ if err != nil {
85
+ return 0, err
86
+ }
87
+ b = append(data[:offset-traversed], b...)
88
+ break
89
+ }
90
+ traversed += size
91
+ }
92
+ if startsubblk == len(dm.pbdata.Blocksizes) {
93
+ // TODO: something?
94
+ /*
95
+ if traversed < offset {
96
+ return 0, errors.New("Tried to start write outside bounds of file.")
97
+ }
98
+ */
99
+ }
100
+ }
101
+
102
+ // Find blocks that need to be overwritten
103
+ var changed []int
104
+ mid := -1
105
+ var midoff uint64
106
+ for i, size := range dm.pbdata.Blocksizes[startsubblk:] {
107
+ if end > traversed {
108
+ changed = append(changed, i+startsubblk)
109
+ } else if end == traversed {
110
+ break
111
+ } else {
112
+ break
113
+ }
114
+ traversed += size
115
+ if end < traversed {
116
+ mid = i + startsubblk
117
+ midoff = end - (traversed - size)
118
+ break
119
+ }
120
+ }
121
+
122
+ var midlnk *mdag.Link
123
+ if mid >= 0 {
124
+ midlnk = dm.curNode.Links[mid]
125
+ midnode, err := dm.dagserv.Get(u.Key(midlnk.Hash))
126
+ if err != nil {
127
+ return 0, err
128
+ }
129
+
130
+ // NOTE: this may have to be changed later when we have multiple
131
+ // layers of indirection
132
+ data, err := ft.UnwrapData(midnode.Data)
133
+ if err != nil {
134
+ return 0, err
135
+ }
136
+ b = append(b, data[midoff:]...)
137
+ }
138
+
139
+ // TODO: dont assume a splitting func here
140
+ subblocks := splitBytes(b, &imp.SizeSplitter2{512})
141
+ var links []*mdag.Link
142
+ var sizes []uint64
143
+ for _, sb := range subblocks {
144
+ n := &mdag.Node{Data: ft.WrapData(sb)}
145
+ _, err := dm.dagserv.Add(n)
146
+ if err != nil {
147
+ log.Error("Failed adding node to DAG service: %s", err)
148
+ return 0, err
149
+ }
150
+ lnk, err := mdag.MakeLink(n)
151
+ if err != nil {
152
+ return 0, err
153
+ }
154
+ links = append(links, lnk)
155
+ sizes = append(sizes, uint64(len(sb)))
156
+ }
157
+
158
+ // This is disgusting
159
+ if len(changed) > 0 {
160
+ dm.curNode.Links = append(dm.curNode.Links[:changed[0]], append(links, dm.curNode.Links[changed[len(changed)-1]+1:]...)...)
161
+ dm.pbdata.Blocksizes = append(dm.pbdata.Blocksizes[:changed[0]], append(sizes, dm.pbdata.Blocksizes[changed[len(changed)-1]+1:]...)...)
162
+ } else {
163
+ dm.curNode.Links = append(dm.curNode.Links, links...)
164
+ dm.pbdata.Blocksizes = append(dm.pbdata.Blocksizes, sizes...)
165
+ }
166
+ dm.pbdata.Filesize = proto.Uint64(newsize)
167
+
168
+ return origlen, nil
169
+}
170
+
171
+func splitBytes(b []byte, spl imp.StreamSplitter) [][]byte {
172
+ ch := make(chan []byte)
173
+ out := spl.Split(ch)
174
+ ch <- b
175
+ close(ch)
176
+ var arr [][]byte
177
+ for blk := range out {
178
+ arr = append(arr, blk)
179
+ }
180
+ return arr
181
+}
182
+
183
+func (dm *DagModifier) GetNode() (*mdag.Node, error) {
184
+ b, err := proto.Marshal(dm.pbdata)
185
+ if err != nil {
186
+ return nil, err
187
+ }
188
+ dm.curNode.Data = b
189
+ return dm.curNode.Copy(), nil
190
+}
importer/dagwriter/dagmodifier_test.go
new
+187
@@ -0,0 +1,187 @@
1
+package dagwriter
2
+
3
+import (
4
+ "fmt"
5
+ "io"
6
+ "io/ioutil"
7
+ "math/rand"
8
+ "testing"
9
+ "time"
10
+
11
+ "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/op/go-logging"
12
+ bs "github.com/jbenet/go-ipfs/blockservice"
13
+ imp "github.com/jbenet/go-ipfs/importer"
14
+ ft "github.com/jbenet/go-ipfs/importer/format"
15
+ mdag "github.com/jbenet/go-ipfs/merkledag"
16
+
17
+ ds "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/datastore.go"
18
+)
19
+
20
+type randGen struct {
21
+ src rand.Source
22
+}
23
+
24
+func newRand() *randGen {
25
+ return &randGen{rand.NewSource(time.Now().UnixNano())}
26
+}
27
+
28
+func (r *randGen) Read(p []byte) (n int, err error) {
29
+ todo := len(p)
30
+ offset := 0
31
+ for {
32
+ val := int64(r.src.Int63())
33
+ for i := 0; i < 8; i++ {
34
+ p[offset] = byte(val & 0xff)
35
+ todo--
36
+ if todo == 0 {
37
+ return len(p), nil
38
+ }
39
+ offset++
40
+ val >>= 8
41
+ }
42
+ }
43
+
44
+ panic("unreachable")
45
+}
46
+
47
+func getMockDagServ(t *testing.T) *mdag.DAGService {
48
+ dstore := ds.NewMapDatastore()
49
+ bserv, err := bs.NewBlockService(dstore, nil)
50
+ if err != nil {
51
+ t.Fatal(err)
52
+ }
53
+ return &mdag.DAGService{bserv}
54
+}
55
+
56
+func getNode(t *testing.T, dserv *mdag.DAGService, size int64) ([]byte, *mdag.Node) {
57
+ dw := NewDagWriter(dserv, &imp.SizeSplitter2{500})
58
+
59
+ n, err := io.CopyN(dw, newRand(), size)
60
+ if err != nil {
61
+ t.Fatal(err)
62
+ }
63
+ if n != size {
64
+ t.Fatal("Incorrect copy amount!")
65
+ }
66
+
67
+ dw.Close()
68
+ node := dw.GetNode()
69
+
70
+ dr, err := mdag.NewDagReader(node, dserv)
71
+ if err != nil {
72
+ t.Fatal(err)
73
+ }
74
+
75
+ b, err := ioutil.ReadAll(dr)
76
+ if err != nil {
77
+ t.Fatal(err)
78
+ }
79
+
80
+ return b, node
81
+}
82
+
83
+func testModWrite(t *testing.T, beg, size uint64, orig []byte, dm *DagModifier) []byte {
84
+ newdata := make([]byte, size)
85
+ r := newRand()
86
+ r.Read(newdata)
87
+
88
+ if size+beg > uint64(len(orig)) {
89
+ orig = append(orig, make([]byte, (size+beg)-uint64(len(orig)))...)
90
+ }
91
+ copy(orig[beg:], newdata)
92
+
93
+ nmod, err := dm.WriteAt(newdata, uint64(beg))
94
+ if err != nil {
95
+ t.Fatal(err)
96
+ }
97
+
98
+ if nmod != int(size) {
99
+ t.Fatalf("Mod length not correct! %d != %d", nmod, size)
100
+ }
101
+
102
+ nd, err := dm.GetNode()
103
+ if err != nil {
104
+ t.Fatal(err)
105
+ }
106
+
107
+ rd, err := mdag.NewDagReader(nd, dm.dagserv)
108
+ if err != nil {
109
+ t.Fatal(err)
110
+ }
111
+
112
+ after, err := ioutil.ReadAll(rd)
113
+ if err != nil {
114
+ t.Fatal(err)
115
+ }
116
+
117
+ err = arrComp(after, orig)
118
+ if err != nil {
119
+ t.Fatal(err)
120
+ }
121
+ return orig
122
+}
123
+
124
+func TestDagModifierBasic(t *testing.T) {
125
+ logging.SetLevel(logging.CRITICAL, "blockservice")
126
+ logging.SetLevel(logging.CRITICAL, "merkledag")
127
+ dserv := getMockDagServ(t)
128
+ b, n := getNode(t, dserv, 50000)
129
+
130
+ dagmod, err := NewDagModifier(n, dserv)
131
+ if err != nil {
132
+ t.Fatal(err)
133
+ }
134
+
135
+ // Within zero block
136
+ beg := uint64(15)
137
+ length := uint64(60)
138
+
139
+ t.Log("Testing mod within zero block")
140
+ b = testModWrite(t, beg, length, b, dagmod)
141
+
142
+ // Within bounds of existing file
143
+ beg = 1000
144
+ length = 4000
145
+ t.Log("Testing mod within bounds of existing file.")
146
+ b = testModWrite(t, beg, length, b, dagmod)
147
+
148
+ // Extend bounds
149
+ beg = 49500
150
+ length = 4000
151
+
152
+ t.Log("Testing mod that extends file.")
153
+ b = testModWrite(t, beg, length, b, dagmod)
154
+
155
+ // "Append"
156
+ beg = uint64(len(b))
157
+ length = 3000
158
+ b = testModWrite(t, beg, length, b, dagmod)
159
+
160
+ // Verify reported length
161
+ node, err := dagmod.GetNode()
162
+ if err != nil {
163
+ t.Fatal(err)
164
+ }
165
+
166
+ size, err := ft.DataSize(node.Data)
167
+ if err != nil {
168
+ t.Fatal(err)
169
+ }
170
+
171
+ expected := uint64(50000 + 3500 + 3000)
172
+ if size != expected {
173
+ t.Fatal("Final reported size is incorrect [%d != %d]", size, expected)
174
+ }
175
+}
176
+
177
+func arrComp(a, b []byte) error {
178
+ if len(a) != len(b) {
179
+ return fmt.Errorf("Arrays differ in length. %d != %d", len(a), len(b))
180
+ }
181
+ for i, v := range a {
182
+ if v != b[i] {
183
+ return fmt.Errorf("Arrays differ at index: %d", i)
184
+ }
185
+ }
186
+ return nil
187
+}
importer/dagwriter/dagwriter.go
+13
-4
@@ -32,10 +32,11 @@ func NewDagWriter(ds *dag.DAGService, splitter imp.StreamSplitter) *DagWriter {
32
func (dw *DagWriter) startSplitter() {
33
blkchan := dw.splitter.Split(dw.splChan)
34
first := <-blkchan
35
+ mbf := new(ft.MultiBlock)
36
root := new(dag.Node)
36
- fileSize := uint64(0)
37
+
38
for blkData := range blkchan {
38
- fileSize += uint64(len(blkData))
39
+ mbf.AddBlockSize(uint64(len(blkData)))
40
node := &dag.Node{Data: ft.WrapData(blkData)}
41
_, err := dw.dagserv.Add(node)
42
if err != nil {
@@ -50,8 +51,16 @@ func (dw *DagWriter) startSplitter() {
51
return
52
}
53
}
53
- root.Data = ft.FilePBData(first, fileSize)
54
- _, err := dw.dagserv.Add(root)
54
+ mbf.Data = first
55
+ data, err := mbf.GetBytes()
56
+ if err != nil {
57
+ dw.seterr = err
58
+ log.Critical("Failed generating bytes for multiblock file: %s", err)
59
+ return
60
+ }
61
+ root.Data = data
62
+
63
+ _, err = dw.dagserv.Add(root)
64
if err != nil {
65
dw.seterr = err
66
log.Critical("Got error adding created node to dagservice: %s", err)
importer/dagwriter/dagwriter_test.go
+25
@@ -100,3 +100,28 @@ func TestMassiveWrite(t *testing.T) {
100
}
101
dw.Close()
102
}
103
+
104
+func BenchmarkDagWriter(b *testing.B) {
105
+ dstore := ds.NewNullDatastore()
106
+ bserv, err := bs.NewBlockService(dstore, nil)
107
+ if err != nil {
108
+ b.Fatal(err)
109
+ }
110
+ dag := &mdag.DAGService{bserv}
111
+
112
+ b.ResetTimer()
113
+ nbytes := int64(b.N)
114
+ for i := 0; i < b.N; i++ {
115
+ b.SetBytes(nbytes)
116
+ dw := NewDagWriter(dag, &imp.SizeSplitter2{4096})
117
+ n, err := io.CopyN(dw, &datasource{}, nbytes)
118
+ if err != nil {
119
+ b.Fatal(err)
120
+ }
121
+ if n != nbytes {
122
+ b.Fatal("Incorrect copy size.")
123
+ }
124
+ dw.Close()
125
+ }
126
+
127
+}
importer/format/format.go
+39
@@ -8,6 +8,15 @@ import (
8
"code.google.com/p/goprotobuf/proto"
9
)
10
11
+func FromBytes(data []byte) (*PBData, error) {
12
+ pbdata := new(PBData)
13
+ err := proto.Unmarshal(data, pbdata)
14
+ if err != nil {
15
+ return nil, err
16
+ }
17
+ return pbdata, nil
18
+}
19
+
20
func FilePBData(data []byte, totalsize uint64) []byte {
21
pbfile := new(PBData)
22
typ := PBData_File
@@ -51,6 +60,15 @@ func WrapData(b []byte) []byte {
60
return out
61
}
62
63
+func UnwrapData(data []byte) ([]byte, error) {
64
+ pbdata := new(PBData)
65
+ err := proto.Unmarshal(data, pbdata)
66
+ if err != nil {
67
+ return nil, err
68
+ }
69
+ return pbdata.GetData(), nil
70
+}
71
+
72
func DataSize(data []byte) (uint64, error) {
73
pbdata := new(PBData)
74
err := proto.Unmarshal(data, pbdata)
@@ -69,3 +87,24 @@ func DataSize(data []byte) (uint64, error) {
87
return 0, errors.New("Unrecognized node data type!")
88
}
89
}
90
+
91
+type MultiBlock struct {
92
+ Data []byte
93
+ blocksizes []uint64
94
+ subtotal uint64
95
+}
96
+
97
+func (mb *MultiBlock) AddBlockSize(s uint64) {
98
+ mb.subtotal += s
99
+ mb.blocksizes = append(mb.blocksizes, s)
100
+}
101
+
102
+func (mb *MultiBlock) GetBytes() ([]byte, error) {
103
+ pbn := new(PBData)
104
+ t := PBData_File
105
+ pbn.Type = &t
106
+ pbn.Filesize = proto.Uint64(uint64(len(mb.Data)) + mb.subtotal)
107
+ pbn.Blocksizes = mb.blocksizes
108
+ pbn.Data = mb.Data
109
+ return proto.Marshal(pbn)
110
+}
importer/importer.go
+10
-6
@@ -31,19 +31,23 @@ func NewDagFromReaderWithSplitter(r io.Reader, spl BlockSplitter) (*dag.Node, er
31
first := <-blkChan
32
root := &dag.Node{}
33
34
- i := 0
35
- totalsize := uint64(len(first))
34
+ mbf := new(ft.MultiBlock)
35
for blk := range blkChan {
37
- totalsize += uint64(len(blk))
36
+ mbf.AddBlockSize(uint64(len(blk)))
37
child := &dag.Node{Data: ft.WrapData(blk)}
39
- err := root.AddNodeLink(fmt.Sprintf("%d", i), child)
38
+ err := root.AddNodeLink("", child)
39
if err != nil {
40
return nil, err
41
}
43
- i++
42
}
43
46
- root.Data = ft.FilePBData(first, totalsize)
44
+ mbf.Data = first
45
+ data, err := mbf.GetBytes()
46
+ if err != nil {
47
+ return nil, err
48
+ }
49
+
50
+ root.Data = data
51
return root, nil
52
}
53
merkledag/merkledag.go
+22
-25
@@ -34,9 +34,6 @@ type Link struct {
34
// cumulative size of target object
35
Size uint64
36
37
- // cumulative size of data stored in object
38
- DataSize uint64
39
-
37
// multihash of the target object
38
Hash mh.Multihash
39
@@ -44,45 +41,45 @@ type Link struct {
41
Node *Node
42
}
43
47
-// AddNodeLink adds a link to another node.
48
-func (n *Node) AddNodeLink(name string, that *Node) error {
49
- s, err := that.Size()
44
+func MakeLink(n *Node) (*Link, error) {
45
+ s, err := n.Size()
46
if err != nil {
51
- return err
47
+ return nil, err
48
}
49
54
- h, err := that.Multihash()
50
+ h, err := n.Multihash()
51
if err != nil {
56
- return err
52
+ return nil, err
53
}
58
-
59
- n.Links = append(n.Links, &Link{
60
- Name: name,
54
+ return &Link{
55
Size: s,
56
Hash: h,
63
- Node: that,
64
- })
65
- return nil
57
+ }, nil
58
}
59
68
-// AddNodeLink adds a link to another node. without keeping a reference to
69
-// the child node
70
-func (n *Node) AddNodeLinkClean(name string, that *Node) error {
71
- s, err := that.Size()
60
+// AddNodeLink adds a link to another node.
61
+func (n *Node) AddNodeLink(name string, that *Node) error {
62
+ lnk, err := MakeLink(that)
63
if err != nil {
64
return err
65
}
66
+ lnk.Name = name
67
+ lnk.Node = that
68
+
69
+ n.Links = append(n.Links, lnk)
70
+ return nil
71
+}
72
76
- h, err := that.Multihash()
73
+// AddNodeLink adds a link to another node. without keeping a reference to
74
+// the child node
75
+func (n *Node) AddNodeLinkClean(name string, that *Node) error {
76
+ lnk, err := MakeLink(that)
77
if err != nil {
78
return err
79
}
80
+ lnk.Name = name
81
81
- n.Links = append(n.Links, &Link{
82
- Name: name,
83
- Size: s,
84
- Hash: h,
85
- })
82
+ n.Links = append(n.Links, lnk)
83
return nil
84
}
85