implement a basic data format for data inside dag nodes
Jeromy committed
Sep 5, 2014 at 20:47 UTC
dad9751754e4bd08743cdb560794f00b00caa5f0
13 files changed
+333
-19
.gitignore
+1
@@ -2,3 +2,4 @@
2
.ipfsconfig
3
*.out
4
*.test
5
+*.orig
cmd/ipfs/add.go
+1
-1
@@ -85,7 +85,7 @@ func addPath(n *core.IpfsNode, fpath string, depth int) (*dag.Node, error) {
85
}
86
87
func addDir(n *core.IpfsNode, fpath string, depth int) (*dag.Node, error) {
88
- tree := &dag.Node{}
88
+ tree := &dag.Node{Data: dag.FolderPBData()}
89
90
files, err := ioutil.ReadDir(fpath)
91
if err != nil {
cmd/ipfs/cat.go
+28
-8
@@ -2,9 +2,13 @@ package main
2
3
import (
4
"fmt"
5
+ "io"
6
+ "os"
7
8
"github.com/gonuts/flag"
9
"github.com/jbenet/commander"
10
+ bserv "github.com/jbenet/go-ipfs/blockservice"
11
+ dag "github.com/jbenet/go-ipfs/merkledag"
12
u "github.com/jbenet/go-ipfs/util"
13
)
14
@@ -37,21 +41,37 @@ func catCmd(c *commander.Command, inp []string) error {
41
return err
42
}
43
40
- fmt.Println("Printing Data!")
41
- _, err = fmt.Printf("%s", nd.Data)
44
+ err = ExpandDag(nd, n.Blocks)
45
if err != nil {
46
return err
47
}
48
46
- fmt.Println("Printing child nodes:")
47
- for _, subn := range nd.Links {
48
- k := u.Key(subn.Hash)
49
- blk, err := n.Blocks.GetBlock(k)
50
- fmt.Printf("Getting link: %s\n", k.Pretty())
49
+ read, err := dag.NewDagReader(nd)
50
+ if err != nil {
51
+ fmt.Println(err)
52
+ continue
53
+ }
54
+
55
+ _, err = io.Copy(os.Stdout, read)
56
+ if err != nil {
57
+ fmt.Println(err)
58
+ continue
59
+ }
60
+ }
61
+ return nil
62
+}
63
+
64
+// Expand all subnodes in this dag so printing can occur without error
65
+//TODO: this needs to be done MUCH better in a somewhat asynchronous way.
66
+//also should be moved elsewhere.
67
+func ExpandDag(nd *dag.Node, bs *bserv.BlockService) error {
68
+ for _, lnk := range nd.Links {
69
+ if lnk.Node == nil {
70
+ blk, err := bs.GetBlock(u.Key(lnk.Hash))
71
if err != nil {
72
return err
73
}
54
- fmt.Println(string(blk.Data))
74
+ lnk.Node = &dag.Node{Data: dag.WrapData(blk.Data)}
75
}
76
}
77
return nil
dev.md
+3
-1
@@ -31,10 +31,12 @@ There are multiple subpackages:
31
- multiplexing connections (tcp atm)
32
- peer addressing
33
- dht - impl basic kademlia routing
34
+- bitswap - impl basic block exchange functionality
35
36
### What's in progress:
37
37
-- bitswap - impl basic block exchange functionality
38
+- crypto - building trust between peers in the network
39
+
40
41
### What's next:
42
fuse/readonly/readonly_unix.go
+15
-2
@@ -6,6 +6,7 @@ package readonly
6
7
import (
8
"fmt"
9
+ "io/ioutil"
10
"os"
11
"os/exec"
12
"os/signal"
@@ -119,7 +120,13 @@ func (s *Node) ReadDir(intr fs.Intr) ([]fuse.Dirent, fuse.Error) {
120
// ReadAll reads the object data as file data
121
func (s *Node) ReadAll(intr fs.Intr) ([]byte, fuse.Error) {
122
u.DOut("Read node.\n")
122
- return []byte(s.Nd.Data), nil
123
+ r, err := mdag.NewDagReader(s.Nd)
124
+ if err != nil {
125
+ return nil, err
126
+ }
127
+ // this is a terrible function... 'ReadAll'?
128
+ // what if i have a 6TB file? GG RAM.
129
+ return ioutil.ReadAll(r)
130
}
131
132
// Mount mounts an IpfsNode instance at a particular path. It
@@ -132,7 +139,13 @@ func Mount(ipfs *core.IpfsNode, fpath string) error {
139
140
go func() {
141
<-sigc
135
- Unmount(fpath)
142
+ for {
143
+ err := Unmount(fpath)
144
+ if err == nil {
145
+ return
146
+ }
147
+ time.Sleep(time.Millisecond * 10)
148
+ }
149
}()
150
151
c, err := fuse.Mount(fpath)
importer/importer.go
+7
-5
@@ -20,19 +20,21 @@ 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
- blkChan := SplitterBySize(1024 * 512)(r)
24
- root := &dag.Node{}
23
+ return NewDagFromReaderWithSplitter(r, SplitterBySize(1024*512))
24
+}
25
+
26
+func NewDagFromReaderWithSplitter(r io.Reader, spl BlockSplitter) (*dag.Node, error) {
27
+ blkChan := spl(r)
28
+ root := &dag.Node{Data: dag.FilePBData()}
29
30
for blk := range blkChan {
27
- child := &dag.Node{Data: blk}
31
+ child := &dag.Node{Data: dag.WrapData(blk)}
32
err := root.AddNodeLink("", child)
33
if err != nil {
34
return nil, err
35
}
36
}
37
34
- fmt.Println(root.Links)
35
-
38
return root, nil
39
}
40
importer/importer_test.go
new
+34
@@ -0,0 +1,34 @@
1
+package importer
2
+
3
+import (
4
+ "bytes"
5
+ "crypto/rand"
6
+ "io"
7
+ "io/ioutil"
8
+ "testing"
9
+
10
+ dag "github.com/jbenet/go-ipfs/merkledag"
11
+)
12
+
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))
18
+ if err != nil {
19
+ t.Fatal(err)
20
+ }
21
+ r, err := dag.NewDagReader(nd)
22
+ if err != nil {
23
+ t.Fatal(err)
24
+ }
25
+
26
+ out, err := ioutil.ReadAll(r)
27
+ if err != nil {
28
+ t.Fatal(err)
29
+ }
30
+
31
+ if !bytes.Equal(out, should) {
32
+ t.Fatal("Output not the same as input.")
33
+ }
34
+}
merkledag/Makefile
+4
-1
@@ -1,8 +1,11 @@
1
2
-all: node.pb.go
2
+all: node.pb.go data.pb.go
3
4
node.pb.go: node.proto
5
protoc --gogo_out=. --proto_path=../../../../:/usr/local/opt/protobuf/include:. $<
6
7
+data.pb.go: data.proto
8
+ protoc --go_out=. data.proto
9
+
10
clean:
11
rm node.pb.go
merkledag/dagreader.go
new
+98
@@ -0,0 +1,98 @@
1
+package merkledag
2
+
3
+import (
4
+ "bytes"
5
+ "errors"
6
+ "io"
7
+
8
+ "code.google.com/p/goprotobuf/proto"
9
+)
10
+
11
+var ErrIsDir = errors.New("this dag node is a directory.")
12
+
13
+// DagReader provides a way to easily read the data contained in a dag.
14
+type DagReader struct {
15
+ node *Node
16
+ position int
17
+ buf *bytes.Buffer
18
+ thisData []byte
19
+}
20
+
21
+func NewDagReader(n *Node) (io.Reader, error) {
22
+ pb := new(PBData)
23
+ err := proto.Unmarshal(n.Data, pb)
24
+ if err != nil {
25
+ return nil, err
26
+ }
27
+ switch pb.GetType() {
28
+ case PBData_Directory:
29
+ return nil, ErrIsDir
30
+ case PBData_File:
31
+ return &DagReader{
32
+ node: n,
33
+ thisData: pb.GetData(),
34
+ }, nil
35
+ case PBData_Raw:
36
+ return bytes.NewBuffer(pb.GetData()), nil
37
+ default:
38
+ panic("Unrecognized node type!")
39
+ }
40
+}
41
+
42
+func (dr *DagReader) precalcNextBuf() error {
43
+ if dr.position >= len(dr.node.Links) {
44
+ return io.EOF
45
+ }
46
+ nxtLink := dr.node.Links[dr.position]
47
+ nxt := nxtLink.Node
48
+ if nxt == nil {
49
+ //TODO: should use dagservice or something to get needed block
50
+ return errors.New("Link to nil node! Tree not fully expanded!")
51
+ }
52
+ pb := new(PBData)
53
+ err := proto.Unmarshal(nxt.Data, pb)
54
+ if err != nil {
55
+ return err
56
+ }
57
+ dr.position++
58
+
59
+ // TODO: dont assume a single layer of indirection
60
+ switch pb.GetType() {
61
+ case PBData_Directory:
62
+ panic("Why is there a directory under a file?")
63
+ case PBData_File:
64
+ //TODO: maybe have a PBData_Block type for indirect blocks?
65
+ panic("Not yet handling different layers of indirection!")
66
+ case PBData_Raw:
67
+ dr.buf = bytes.NewBuffer(pb.GetData())
68
+ return nil
69
+ default:
70
+ panic("Unrecognized node type!")
71
+ }
72
+}
73
+
74
+func (dr *DagReader) Read(b []byte) (int, error) {
75
+ if dr.buf == nil {
76
+ err := dr.precalcNextBuf()
77
+ if err != nil {
78
+ return 0, err
79
+ }
80
+ }
81
+ total := 0
82
+ for {
83
+ n, err := dr.buf.Read(b[total:])
84
+ total += n
85
+ if err != nil {
86
+ if err != io.EOF {
87
+ return total, err
88
+ }
89
+ }
90
+ if total == len(b) {
91
+ return total, nil
92
+ }
93
+ err = dr.precalcNextBuf()
94
+ if err != nil {
95
+ return total, err
96
+ }
97
+ }
98
+}
merkledag/data.pb.go
new
+85
@@ -0,0 +1,85 @@
1
+// Code generated by protoc-gen-go.
2
+// source: data.proto
3
+// DO NOT EDIT!
4
+
5
+/*
6
+Package merkledag is a generated protocol buffer package.
7
+
8
+It is generated from these files:
9
+ data.proto
10
+
11
+It has these top-level messages:
12
+ PBData
13
+*/
14
+package merkledag
15
+
16
+import proto "code.google.com/p/goprotobuf/proto"
17
+import math "math"
18
+
19
+// Reference imports to suppress errors if they are not otherwise used.
20
+var _ = proto.Marshal
21
+var _ = math.Inf
22
+
23
+type PBData_DataType int32
24
+
25
+const (
26
+ PBData_Raw PBData_DataType = 0
27
+ PBData_Directory PBData_DataType = 1
28
+ PBData_File PBData_DataType = 2
29
+)
30
+
31
+var PBData_DataType_name = map[int32]string{
32
+ 0: "Raw",
33
+ 1: "Directory",
34
+ 2: "File",
35
+}
36
+var PBData_DataType_value = map[string]int32{
37
+ "Raw": 0,
38
+ "Directory": 1,
39
+ "File": 2,
40
+}
41
+
42
+func (x PBData_DataType) Enum() *PBData_DataType {
43
+ p := new(PBData_DataType)
44
+ *p = x
45
+ return p
46
+}
47
+func (x PBData_DataType) String() string {
48
+ return proto.EnumName(PBData_DataType_name, int32(x))
49
+}
50
+func (x *PBData_DataType) UnmarshalJSON(data []byte) error {
51
+ value, err := proto.UnmarshalJSONEnum(PBData_DataType_value, data, "PBData_DataType")
52
+ if err != nil {
53
+ return err
54
+ }
55
+ *x = PBData_DataType(value)
56
+ return nil
57
+}
58
+
59
+type PBData struct {
60
+ Type *PBData_DataType `protobuf:"varint,1,req,enum=merkledag.PBData_DataType" json:"Type,omitempty"`
61
+ Data []byte `protobuf:"bytes,2,opt" json:"Data,omitempty"`
62
+ XXX_unrecognized []byte `json:"-"`
63
+}
64
+
65
+func (m *PBData) Reset() { *m = PBData{} }
66
+func (m *PBData) String() string { return proto.CompactTextString(m) }
67
+func (*PBData) ProtoMessage() {}
68
+
69
+func (m *PBData) GetType() PBData_DataType {
70
+ if m != nil && m.Type != nil {
71
+ return *m.Type
72
+ }
73
+ return PBData_Raw
74
+}
75
+
76
+func (m *PBData) GetData() []byte {
77
+ if m != nil {
78
+ return m.Data
79
+ }
80
+ return nil
81
+}
82
+
83
+func init() {
84
+ proto.RegisterEnum("merkledag.PBData_DataType", PBData_DataType_name, PBData_DataType_value)
85
+}
merkledag/data.proto
new
+12
@@ -0,0 +1,12 @@
1
+package merkledag;
2
+
3
+message PBData {
4
+ enum DataType {
5
+ Raw = 0;
6
+ Directory = 1;
7
+ File = 2;
8
+ }
9
+
10
+ required DataType Type = 1;
11
+ optional bytes Data = 2;
12
+}
merkledag/merkledag.go
+43
@@ -3,6 +3,8 @@ package merkledag
3
import (
4
"fmt"
5
6
+ "code.google.com/p/goprotobuf/proto"
7
+
8
blocks "github.com/jbenet/go-ipfs/blocks"
9
bserv "github.com/jbenet/go-ipfs/blockservice"
10
u "github.com/jbenet/go-ipfs/util"
@@ -152,3 +154,44 @@ func (n *DAGService) Get(k u.Key) (*Node, error) {
154
155
return Decoded(b.Data)
156
}
157
+
158
+func FilePBData() []byte {
159
+ pbfile := new(PBData)
160
+ typ := PBData_File
161
+ pbfile.Type = &typ
162
+
163
+ data, err := proto.Marshal(pbfile)
164
+ if err != nil {
165
+ //this really shouldnt happen, i promise
166
+ panic(err)
167
+ }
168
+ return data
169
+}
170
+
171
+func FolderPBData() []byte {
172
+ pbfile := new(PBData)
173
+ typ := PBData_Directory
174
+ pbfile.Type = &typ
175
+
176
+ data, err := proto.Marshal(pbfile)
177
+ if err != nil {
178
+ //this really shouldnt happen, i promise
179
+ panic(err)
180
+ }
181
+ return data
182
+}
183
+
184
+func WrapData(b []byte) []byte {
185
+ pbdata := new(PBData)
186
+ typ := PBData_Raw
187
+ pbdata.Data = b
188
+ pbdata.Type = &typ
189
+
190
+ out, err := proto.Marshal(pbdata)
191
+ if err != nil {
192
+ // This shouldnt happen. seriously.
193
+ panic(err)
194
+ }
195
+
196
+ return out
197
+}
swarm/conn.go
+2
-1
@@ -13,7 +13,8 @@ import (
13
// ChanBuffer is the size of the buffer in the Conn Chan
14
const ChanBuffer = 10
15
16
-const MaxMessageSize = 1 << 19
16
+// 1 MB
17
+const MaxMessageSize = 1 << 20
18
19
// Conn represents a connection to another Peer (IPFS Node).
20
type Conn struct {