@cryptotaxi247 / kubo / commits / 09854c064

Move proto-dag reader to separate file and change name

License: MIT Signed-off-by: Jakub Sztandera <kubuxu@protonmail.ch>

Jakub Sztandera committed Nov 21, 2016 at 17:18 UTC 09854c0645d80aa23efbaa8cf5fc8fafeaff2aad
3 files changed +264 -250
unixfs/archive/tar/writer.go
+1 -1
@@ -64,7 +64,7 @@ func (w *Writer) writeFile(nd *mdag.ProtoNode, pb *upb.Data, fpath string) error
64 return err
65 }
66
67 - dagr := uio.NewDataFileReader(w.ctx, nd, pb, w.Dag)
67 + dagr := uio.NewPBFileReader(w.ctx, nd, pb, w.Dag)
68 if _, err := dagr.WriteTo(w.TarW); err != nil {
69 return err
70 }
unixfs/io/dagreader.go
+1 -249
@@ -1,12 +1,10 @@
1 package io
2
3 import (
4 - "bytes"
4 "context"
5 "errors"
6 "fmt"
7 "io"
9 - "os"
8
9 mdag "github.com/ipfs/go-ipfs/merkledag"
10 ft "github.com/ipfs/go-ipfs/unixfs"
@@ -27,36 +25,6 @@ type DagReader interface {
25 Offset() int64
26 }
27
30 -// DagReader provides a way to easily read the data contained in a dag.
31 -type pbDagReader struct {
32 - serv mdag.DAGService
33 -
34 - // the node being read
35 - node *mdag.ProtoNode
36 -
37 - // cached protobuf structure from node.Data
38 - pbdata *ftpb.Data
39 -
40 - // the current data buffer to be read from
41 - // will either be a bytes.Reader or a child DagReader
42 - buf ReadSeekCloser
43 -
44 - // NodeGetters for each of 'nodes' child links
45 - promises []mdag.NodeGetter
46 -
47 - // the index of the child link currently being read from
48 - linkPosition int
49 -
50 - // current offset for the read head within the 'file'
51 - offset int64
52 -
53 - // Our context
54 - ctx context.Context
55 -
56 - // context cancel for children
57 - cancel func()
58 -}
59 -
28 type ReadSeekCloser interface {
29 io.Reader
30 io.Seeker
@@ -83,7 +51,7 @@ func NewDagReader(ctx context.Context, n node.Node, serv mdag.DAGService) (DagRe
51 // Dont allow reading directories
52 return nil, ErrIsDir
53 case ftpb.Data_File, ftpb.Data_Raw:
86 - return NewDataFileReader(ctx, n, pb, serv), nil
54 + return NewPBFileReader(ctx, n, pb, serv), nil
55 case ftpb.Data_Metadata:
56 if len(n.Links()) == 0 {
57 return nil, errors.New("incorrectly formatted metadata object")
@@ -107,219 +75,3 @@ func NewDagReader(ctx context.Context, n node.Node, serv mdag.DAGService) (DagRe
75 return nil, fmt.Errorf("unrecognized node type")
76 }
77 }
110 -
111 -func NewDataFileReader(ctx context.Context, n *mdag.ProtoNode, pb *ftpb.Data, serv mdag.DAGService) *pbDagReader {
112 - fctx, cancel := context.WithCancel(ctx)
113 - promises := mdag.GetDAG(fctx, serv, n)
114 - return &pbDagReader{
115 - node: n,
116 - serv: serv,
117 - buf: NewRSNCFromBytes(pb.GetData()),
118 - promises: promises,
119 - ctx: fctx,
120 - cancel: cancel,
121 - pbdata: pb,
122 - }
123 -}
124 -
125 -// precalcNextBuf follows the next link in line and loads it from the
126 -// DAGService, setting the next buffer to read from
127 -func (dr *pbDagReader) precalcNextBuf(ctx context.Context) error {
128 - dr.buf.Close() // Just to make sure
129 - if dr.linkPosition >= len(dr.promises) {
130 - return io.EOF
131 - }
132 -
133 - nxt, err := dr.promises[dr.linkPosition].Get(ctx)
134 - if err != nil {
135 - return err
136 - }
137 - dr.linkPosition++
138 -
139 - switch nxt := nxt.(type) {
140 - case *mdag.ProtoNode:
141 - pb := new(ftpb.Data)
142 - err = proto.Unmarshal(nxt.Data(), pb)
143 - if err != nil {
144 - return fmt.Errorf("incorrectly formatted protobuf: %s", err)
145 - }
146 -
147 - switch pb.GetType() {
148 - case ftpb.Data_Directory:
149 - // A directory should not exist within a file
150 - return ft.ErrInvalidDirLocation
151 - case ftpb.Data_File:
152 - dr.buf = NewDataFileReader(dr.ctx, nxt, pb, dr.serv)
153 - return nil
154 - case ftpb.Data_Raw:
155 - dr.buf = NewRSNCFromBytes(pb.GetData())
156 - return nil
157 - case ftpb.Data_Metadata:
158 - return errors.New("shouldnt have had metadata object inside file")
159 - case ftpb.Data_Symlink:
160 - return errors.New("shouldnt have had symlink inside file")
161 - default:
162 - return ft.ErrUnrecognizedType
163 - }
164 - case *mdag.RawNode:
165 - dr.buf = NewRSNCFromBytes(nxt.RawData())
166 - return nil
167 - default:
168 - return errors.New("unrecognized node type in pbDagReader")
169 - }
170 -}
171 -
172 -// Size return the total length of the data from the DAG structured file.
173 -func (dr *pbDagReader) Size() uint64 {
174 - return dr.pbdata.GetFilesize()
175 -}
176 -
177 -// Read reads data from the DAG structured file
178 -func (dr *pbDagReader) Read(b []byte) (int, error) {
179 - return dr.CtxReadFull(dr.ctx, b)
180 -}
181 -
182 -// CtxReadFull reads data from the DAG structured file
183 -func (dr *pbDagReader) CtxReadFull(ctx context.Context, b []byte) (int, error) {
184 - // If no cached buffer, load one
185 - total := 0
186 - for {
187 - // Attempt to fill bytes from cached buffer
188 - n, err := dr.buf.Read(b[total:])
189 - total += n
190 - dr.offset += int64(n)
191 - if err != nil {
192 - // EOF is expected
193 - if err != io.EOF {
194 - return total, err
195 - }
196 - }
197 -
198 - // If weve read enough bytes, return
199 - if total == len(b) {
200 - return total, nil
201 - }
202 -
203 - // Otherwise, load up the next block
204 - err = dr.precalcNextBuf(ctx)
205 - if err != nil {
206 - return total, err
207 - }
208 - }
209 -}
210 -
211 -func (dr *pbDagReader) WriteTo(w io.Writer) (int64, error) {
212 - // If no cached buffer, load one
213 - total := int64(0)
214 - for {
215 - // Attempt to write bytes from cached buffer
216 - n, err := dr.buf.WriteTo(w)
217 - total += n
218 - dr.offset += n
219 - if err != nil {
220 - if err != io.EOF {
221 - return total, err
222 - }
223 - }
224 -
225 - // Otherwise, load up the next block
226 - err = dr.precalcNextBuf(dr.ctx)
227 - if err != nil {
228 - if err == io.EOF {
229 - return total, nil
230 - }
231 - return total, err
232 - }
233 - }
234 -}
235 -
236 -func (dr *pbDagReader) Close() error {
237 - dr.cancel()
238 - return nil
239 -}
240 -
241 -func (dr *pbDagReader) Offset() int64 {
242 - return dr.offset
243 -}
244 -
245 -// Seek implements io.Seeker, and will seek to a given offset in the file
246 -// interface matches standard unix seek
247 -// TODO: check if we can do relative seeks, to reduce the amount of dagreader
248 -// recreations that need to happen.
249 -func (dr *pbDagReader) Seek(offset int64, whence int) (int64, error) {
250 - switch whence {
251 - case os.SEEK_SET:
252 - if offset < 0 {
253 - return -1, errors.New("Invalid offset")
254 - }
255 -
256 - // Grab cached protobuf object (solely to make code look cleaner)
257 - pb := dr.pbdata
258 -
259 - // left represents the number of bytes remaining to seek to (from beginning)
260 - left := offset
261 - if int64(len(pb.Data)) >= offset {
262 - // Close current buf to close potential child dagreader
263 - dr.buf.Close()
264 - dr.buf = NewRSNCFromBytes(pb.GetData()[offset:])
265 -
266 - // start reading links from the beginning
267 - dr.linkPosition = 0
268 - dr.offset = offset
269 - return offset, nil
270 - } else {
271 - // skip past root block data
272 - left -= int64(len(pb.Data))
273 - }
274 -
275 - // iterate through links and find where we need to be
276 - for i := 0; i < len(pb.Blocksizes); i++ {
277 - if pb.Blocksizes[i] > uint64(left) {
278 - dr.linkPosition = i
279 - break
280 - } else {
281 - left -= int64(pb.Blocksizes[i])
282 - }
283 - }
284 -
285 - // start sub-block request
286 - err := dr.precalcNextBuf(dr.ctx)
287 - if err != nil {
288 - return 0, err
289 - }
290 -
291 - // set proper offset within child readseeker
292 - n, err := dr.buf.Seek(left, os.SEEK_SET)
293 - if err != nil {
294 - return -1, err
295 - }
296 -
297 - // sanity
298 - left -= n
299 - if left != 0 {
300 - return -1, errors.New("failed to seek properly")
301 - }
302 - dr.offset = offset
303 - return offset, nil
304 - case os.SEEK_CUR:
305 - // TODO: be smarter here
306 - noffset := dr.offset + offset
307 - return dr.Seek(noffset, os.SEEK_SET)
308 - case os.SEEK_END:
309 - noffset := int64(dr.pbdata.GetFilesize()) - offset
310 - return dr.Seek(noffset, os.SEEK_SET)
311 - default:
312 - return 0, errors.New("invalid whence")
313 - }
314 -}
315 -
316 -// readSeekNopCloser wraps a bytes.Reader to implement ReadSeekCloser
317 -type readSeekNopCloser struct {
318 - *bytes.Reader
319 -}
320 -
321 -func NewRSNCFromBytes(b []byte) ReadSeekCloser {
322 - return &readSeekNopCloser{bytes.NewReader(b)}
323 -}
324 -
325 -func (r *readSeekNopCloser) Close() error { return nil }
unixfs/io/pbdagreader.go new
+262
@@ -0,0 +1,262 @@
1 +package io
2 +
3 +import (
4 + "bytes"
5 + "context"
6 + "errors"
7 + "fmt"
8 + "io"
9 + "os"
10 +
11 + mdag "github.com/ipfs/go-ipfs/merkledag"
12 + ft "github.com/ipfs/go-ipfs/unixfs"
13 + ftpb "github.com/ipfs/go-ipfs/unixfs/pb"
14 +
15 + proto "gx/ipfs/QmZ4Qi3GaRbjcx28Sme5eMH7RQjGkt8wHxt2a65oLaeFEV/gogo-protobuf/proto"
16 +)
17 +
18 +// DagReader provides a way to easily read the data contained in a dag.
19 +type pbDagReader struct {
20 + serv mdag.DAGService
21 +
22 + // the node being read
23 + node *mdag.ProtoNode
24 +
25 + // cached protobuf structure from node.Data
26 + pbdata *ftpb.Data
27 +
28 + // the current data buffer to be read from
29 + // will either be a bytes.Reader or a child DagReader
30 + buf ReadSeekCloser
31 +
32 + // NodeGetters for each of 'nodes' child links
33 + promises []mdag.NodeGetter
34 +
35 + // the index of the child link currently being read from
36 + linkPosition int
37 +
38 + // current offset for the read head within the 'file'
39 + offset int64
40 +
41 + // Our context
42 + ctx context.Context
43 +
44 + // context cancel for children
45 + cancel func()
46 +}
47 +
48 +func NewPBFileReader(ctx context.Context, n *mdag.ProtoNode, pb *ftpb.Data, serv mdag.DAGService) *pbDagReader {
49 + fctx, cancel := context.WithCancel(ctx)
50 + promises := mdag.GetDAG(fctx, serv, n)
51 + return &pbDagReader{
52 + node: n,
53 + serv: serv,
54 + buf: NewRSNCFromBytes(pb.GetData()),
55 + promises: promises,
56 + ctx: fctx,
57 + cancel: cancel,
58 + pbdata: pb,
59 + }
60 +}
61 +
62 +// precalcNextBuf follows the next link in line and loads it from the
63 +// DAGService, setting the next buffer to read from
64 +func (dr *pbDagReader) precalcNextBuf(ctx context.Context) error {
65 + dr.buf.Close() // Just to make sure
66 + if dr.linkPosition >= len(dr.promises) {
67 + return io.EOF
68 + }
69 +
70 + nxt, err := dr.promises[dr.linkPosition].Get(ctx)
71 + if err != nil {
72 + return err
73 + }
74 + dr.linkPosition++
75 +
76 + switch nxt := nxt.(type) {
77 + case *mdag.ProtoNode:
78 + pb := new(ftpb.Data)
79 + err = proto.Unmarshal(nxt.Data(), pb)
80 + if err != nil {
81 + return fmt.Errorf("incorrectly formatted protobuf: %s", err)
82 + }
83 +
84 + switch pb.GetType() {
85 + case ftpb.Data_Directory:
86 + // A directory should not exist within a file
87 + return ft.ErrInvalidDirLocation
88 + case ftpb.Data_File:
89 + dr.buf = NewPBFileReader(dr.ctx, nxt, pb, dr.serv)
90 + return nil
91 + case ftpb.Data_Raw:
92 + dr.buf = NewRSNCFromBytes(pb.GetData())
93 + return nil
94 + case ftpb.Data_Metadata:
95 + return errors.New("shouldnt have had metadata object inside file")
96 + case ftpb.Data_Symlink:
97 + return errors.New("shouldnt have had symlink inside file")
98 + default:
99 + return ft.ErrUnrecognizedType
100 + }
101 + case *mdag.RawNode:
102 + dr.buf = NewRSNCFromBytes(nxt.RawData())
103 + return nil
104 + default:
105 + return errors.New("unrecognized node type in pbDagReader")
106 + }
107 +}
108 +
109 +// Size return the total length of the data from the DAG structured file.
110 +func (dr *pbDagReader) Size() uint64 {
111 + return dr.pbdata.GetFilesize()
112 +}
113 +
114 +// Read reads data from the DAG structured file
115 +func (dr *pbDagReader) Read(b []byte) (int, error) {
116 + return dr.CtxReadFull(dr.ctx, b)
117 +}
118 +
119 +// CtxReadFull reads data from the DAG structured file
120 +func (dr *pbDagReader) CtxReadFull(ctx context.Context, b []byte) (int, error) {
121 + // If no cached buffer, load one
122 + total := 0
123 + for {
124 + // Attempt to fill bytes from cached buffer
125 + n, err := dr.buf.Read(b[total:])
126 + total += n
127 + dr.offset += int64(n)
128 + if err != nil {
129 + // EOF is expected
130 + if err != io.EOF {
131 + return total, err
132 + }
133 + }
134 +
135 + // If weve read enough bytes, return
136 + if total == len(b) {
137 + return total, nil
138 + }
139 +
140 + // Otherwise, load up the next block
141 + err = dr.precalcNextBuf(ctx)
142 + if err != nil {
143 + return total, err
144 + }
145 + }
146 +}
147 +
148 +func (dr *pbDagReader) WriteTo(w io.Writer) (int64, error) {
149 + // If no cached buffer, load one
150 + total := int64(0)
151 + for {
152 + // Attempt to write bytes from cached buffer
153 + n, err := dr.buf.WriteTo(w)
154 + total += n
155 + dr.offset += n
156 + if err != nil {
157 + if err != io.EOF {
158 + return total, err
159 + }
160 + }
161 +
162 + // Otherwise, load up the next block
163 + err = dr.precalcNextBuf(dr.ctx)
164 + if err != nil {
165 + if err == io.EOF {
166 + return total, nil
167 + }
168 + return total, err
169 + }
170 + }
171 +}
172 +
173 +func (dr *pbDagReader) Close() error {
174 + dr.cancel()
175 + return nil
176 +}
177 +
178 +func (dr *pbDagReader) Offset() int64 {
179 + return dr.offset
180 +}
181 +
182 +// Seek implements io.Seeker, and will seek to a given offset in the file
183 +// interface matches standard unix seek
184 +// TODO: check if we can do relative seeks, to reduce the amount of dagreader
185 +// recreations that need to happen.
186 +func (dr *pbDagReader) Seek(offset int64, whence int) (int64, error) {
187 + switch whence {
188 + case os.SEEK_SET:
189 + if offset < 0 {
190 + return -1, errors.New("Invalid offset")
191 + }
192 +
193 + // Grab cached protobuf object (solely to make code look cleaner)
194 + pb := dr.pbdata
195 +
196 + // left represents the number of bytes remaining to seek to (from beginning)
197 + left := offset
198 + if int64(len(pb.Data)) >= offset {
199 + // Close current buf to close potential child dagreader
200 + dr.buf.Close()
201 + dr.buf = NewRSNCFromBytes(pb.GetData()[offset:])
202 +
203 + // start reading links from the beginning
204 + dr.linkPosition = 0
205 + dr.offset = offset
206 + return offset, nil
207 + } else {
208 + // skip past root block data
209 + left -= int64(len(pb.Data))
210 + }
211 +
212 + // iterate through links and find where we need to be
213 + for i := 0; i < len(pb.Blocksizes); i++ {
214 + if pb.Blocksizes[i] > uint64(left) {
215 + dr.linkPosition = i
216 + break
217 + } else {
218 + left -= int64(pb.Blocksizes[i])
219 + }
220 + }
221 +
222 + // start sub-block request
223 + err := dr.precalcNextBuf(dr.ctx)
224 + if err != nil {
225 + return 0, err
226 + }
227 +
228 + // set proper offset within child readseeker
229 + n, err := dr.buf.Seek(left, os.SEEK_SET)
230 + if err != nil {
231 + return -1, err
232 + }
233 +
234 + // sanity
235 + left -= n
236 + if left != 0 {
237 + return -1, errors.New("failed to seek properly")
238 + }
239 + dr.offset = offset
240 + return offset, nil
241 + case os.SEEK_CUR:
242 + // TODO: be smarter here
243 + noffset := dr.offset + offset
244 + return dr.Seek(noffset, os.SEEK_SET)
245 + case os.SEEK_END:
246 + noffset := int64(dr.pbdata.GetFilesize()) - offset
247 + return dr.Seek(noffset, os.SEEK_SET)
248 + default:
249 + return 0, errors.New("invalid whence")
250 + }
251 +}
252 +
253 +// readSeekNopCloser wraps a bytes.Reader to implement ReadSeekCloser
254 +type readSeekNopCloser struct {
255 + *bytes.Reader
256 +}
257 +
258 +func NewRSNCFromBytes(b []byte) ReadSeekCloser {
259 + return &readSeekNopCloser{bytes.NewReader(b)}
260 +}
261 +
262 +func (r *readSeekNopCloser) Close() error { return nil }