implement seeking in the dagreader
Jeromy committed
Jan 25, 2015 at 08:54 UTC
26826bd55eb2c12001e65d458d6ed1d89b01f176
12 files changed
+312
-63
core/commands/cat.go
+1
-1
@@ -73,7 +73,7 @@ func cat(node *core.IpfsNode, paths []string) ([]io.Reader, uint64, error) {
73
}
74
length += nodeLength
75
76
- read, err := uio.NewDagReader(dagnode, node.DAG)
76
+ read, err := uio.NewDagReader(node.Context(), dagnode, node.DAG)
77
if err != nil {
78
return nil, 0, err
79
}
core/corehttp/gateway_handler.go
+1
-1
@@ -77,7 +77,7 @@ func (i *gatewayHandler) AddNodeToDAG(nd *dag.Node) (u.Key, error) {
77
}
78
79
func (i *gatewayHandler) NewDagReader(nd *dag.Node) (io.Reader, error) {
80
- return uio.NewDagReader(nd, i.node.DAG)
80
+ return uio.NewDagReader(i.node.Context(), nd, i.node.DAG)
81
}
82
83
func (i *gatewayHandler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
core/coreunix/cat.go
+1
-1
@@ -12,5 +12,5 @@ func Cat(n *core.IpfsNode, path string) (io.Reader, error) {
12
if err != nil {
13
return nil, err
14
}
15
- return uio.NewDagReader(dagNode, n.DAG)
15
+ return uio.NewDagReader(n.ContextGroup.Context(), dagNode, n.DAG)
16
}
fuse/ipns/ipns_unix.go
+2
-1
@@ -11,6 +11,7 @@ import (
11
12
fuse "github.com/jbenet/go-ipfs/Godeps/_workspace/src/bazil.org/fuse"
13
fs "github.com/jbenet/go-ipfs/Godeps/_workspace/src/bazil.org/fuse/fs"
14
+ "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
15
proto "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/goprotobuf/proto"
16
17
core "github.com/jbenet/go-ipfs/core"
@@ -337,7 +338,7 @@ func (s *Node) ReadDir(intr fs.Intr) ([]fuse.Dirent, fuse.Error) {
338
// ReadAll reads the object data as file data
339
func (s *Node) ReadAll(intr fs.Intr) ([]byte, fuse.Error) {
340
log.Debugf("ipns: ReadAll [%s]", s.name)
340
- r, err := uio.NewDagReader(s.Nd, s.Ipfs.DAG)
341
+ r, err := uio.NewDagReader(context.TODO(), s.Nd, s.Ipfs.DAG)
342
if err != nil {
343
return nil, err
344
}
fuse/readonly/readonly_unix.go
+2
-1
@@ -10,6 +10,7 @@ import (
10
11
fuse "github.com/jbenet/go-ipfs/Godeps/_workspace/src/bazil.org/fuse"
12
fs "github.com/jbenet/go-ipfs/Godeps/_workspace/src/bazil.org/fuse/fs"
13
+ "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
14
proto "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/goprotobuf/proto"
15
16
core "github.com/jbenet/go-ipfs/core"
@@ -145,7 +146,7 @@ func (s *Node) ReadDir(intr fs.Intr) ([]fuse.Dirent, fuse.Error) {
146
// ReadAll reads the object data as file data
147
func (s *Node) ReadAll(intr fs.Intr) ([]byte, fuse.Error) {
148
log.Debug("Read node.")
148
- r, err := uio.NewDagReader(s.Nd, s.Ipfs.DAG)
149
+ r, err := uio.NewDagReader(context.TODO(), s.Nd, s.Ipfs.DAG)
150
if err != nil {
151
return nil, err
152
}
importer/importer_test.go
+129
-3
@@ -6,8 +6,10 @@ import (
6
"fmt"
7
"io"
8
"io/ioutil"
9
+ "os"
10
"testing"
11
12
+ "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
13
ds "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-datastore"
14
dssync "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-datastore/sync"
15
bstore "github.com/jbenet/go-ipfs/blocks/blockstore"
@@ -51,7 +53,7 @@ func testFileConsistency(t *testing.T, bs chunk.BlockSplitter, nbytes int) {
53
t.Fatal(err)
54
}
55
54
- r, err := uio.NewDagReader(nd, dnp.ds)
56
+ r, err := uio.NewDagReader(context.TODO(), nd, dnp.ds)
57
if err != nil {
58
t.Fatal(err)
59
}
@@ -77,7 +79,7 @@ func TestBuilderConsistency(t *testing.T) {
79
if err != nil {
80
t.Fatal(err)
81
}
80
- r, err := uio.NewDagReader(nd, dagserv)
82
+ r, err := uio.NewDagReader(context.TODO(), nd, dagserv)
83
if err != nil {
84
t.Fatal(err)
85
}
@@ -165,7 +167,7 @@ func TestIndirectBlocks(t *testing.T) {
167
t.Fatal(err)
168
}
169
168
- reader, err := uio.NewDagReader(dag, dnp.ds)
170
+ reader, err := uio.NewDagReader(context.TODO(), dag, dnp.ds)
171
if err != nil {
172
t.Fatal(err)
173
}
@@ -179,3 +181,127 @@ func TestIndirectBlocks(t *testing.T) {
181
t.Fatal("Not equal!")
182
}
183
}
184
+
185
+func TestSeekingBasic(t *testing.T) {
186
+ nbytes := int64(10 * 1024)
187
+ should := make([]byte, nbytes)
188
+ u.NewTimeSeededRand().Read(should)
189
+
190
+ read := bytes.NewReader(should)
191
+ dnp := getDagservAndPinner(t)
192
+ nd, err := BuildDagFromReader(read, dnp.ds, dnp.mp, &chunk.SizeSplitter{500})
193
+ if err != nil {
194
+ t.Fatal(err)
195
+ }
196
+
197
+ rs, err := uio.NewDagReader(context.TODO(), nd, dnp.ds)
198
+ if err != nil {
199
+ t.Fatal(err)
200
+ }
201
+
202
+ start := int64(4000)
203
+ n, err := rs.Seek(start, os.SEEK_SET)
204
+ if err != nil {
205
+ t.Fatal(err)
206
+ }
207
+ if n != start {
208
+ t.Fatal("Failed to seek to correct offset")
209
+ }
210
+
211
+ out, err := ioutil.ReadAll(rs)
212
+ if err != nil {
213
+ t.Fatal(err)
214
+ }
215
+
216
+ err = arrComp(out, should[start:])
217
+ if err != nil {
218
+ t.Fatal(err)
219
+ }
220
+}
221
+
222
+func TestSeekToBegin(t *testing.T) {
223
+ nbytes := int64(10 * 1024)
224
+ should := make([]byte, nbytes)
225
+ u.NewTimeSeededRand().Read(should)
226
+
227
+ read := bytes.NewReader(should)
228
+ dnp := getDagservAndPinner(t)
229
+ nd, err := BuildDagFromReader(read, dnp.ds, dnp.mp, &chunk.SizeSplitter{500})
230
+ if err != nil {
231
+ t.Fatal(err)
232
+ }
233
+
234
+ rs, err := uio.NewDagReader(context.TODO(), nd, dnp.ds)
235
+ if err != nil {
236
+ t.Fatal(err)
237
+ }
238
+
239
+ n, err := io.CopyN(ioutil.Discard, rs, 1024*4)
240
+ if err != nil {
241
+ t.Fatal(err)
242
+ }
243
+ if n != 4096 {
244
+ t.Fatal("Copy didnt copy enough bytes")
245
+ }
246
+
247
+ seeked, err := rs.Seek(0, os.SEEK_SET)
248
+ if err != nil {
249
+ t.Fatal(err)
250
+ }
251
+ if seeked != 0 {
252
+ t.Fatal("Failed to seek to beginning")
253
+ }
254
+
255
+ out, err := ioutil.ReadAll(rs)
256
+ if err != nil {
257
+ t.Fatal(err)
258
+ }
259
+
260
+ err = arrComp(out, should)
261
+ if err != nil {
262
+ t.Fatal(err)
263
+ }
264
+}
265
+
266
+func TestSeekingConsistency(t *testing.T) {
267
+ nbytes := int64(128 * 1024)
268
+ should := make([]byte, nbytes)
269
+ u.NewTimeSeededRand().Read(should)
270
+
271
+ read := bytes.NewReader(should)
272
+ dnp := getDagservAndPinner(t)
273
+ nd, err := BuildDagFromReader(read, dnp.ds, dnp.mp, &chunk.SizeSplitter{500})
274
+ if err != nil {
275
+ t.Fatal(err)
276
+ }
277
+
278
+ rs, err := uio.NewDagReader(context.TODO(), nd, dnp.ds)
279
+ if err != nil {
280
+ t.Fatal(err)
281
+ }
282
+
283
+ out := make([]byte, nbytes)
284
+
285
+ for coff := nbytes - 4096; coff >= 0; coff -= 4096 {
286
+ t.Log(coff)
287
+ n, err := rs.Seek(coff, os.SEEK_SET)
288
+ if err != nil {
289
+ t.Fatal(err)
290
+ }
291
+ if n != coff {
292
+ t.Fatal("wasnt able to seek to the right position")
293
+ }
294
+ nread, err := rs.Read(out[coff : coff+4096])
295
+ if err != nil {
296
+ t.Fatal(err)
297
+ }
298
+ if nread != 4096 {
299
+ t.Fatal("didnt read the correct number of bytes")
300
+ }
301
+ }
302
+
303
+ err = arrComp(out, should)
304
+ if err != nil {
305
+ t.Fatal(err)
306
+ }
307
+}
merkledag/merkledag.go
+44
-31
@@ -2,7 +2,6 @@
2
package merkledag
3
4
import (
5
- "bytes"
5
"fmt"
6
"sync"
7
"time"
@@ -27,6 +26,7 @@ type DAGService interface {
26
// GetDAG returns, in order, all the single leve child
27
// nodes of the passed in node.
28
GetDAG(context.Context, *Node) <-chan *Node
29
+ GetNodes(context.Context, []u.Key) <-chan *Node
30
}
31
32
func NewDAGService(bs *bserv.BlockService) DAGService {
@@ -155,11 +155,10 @@ func FetchGraph(ctx context.Context, root *Node, serv DAGService) chan struct{}
155
156
// FindLinks searches this nodes links for the given key,
157
// returns the indexes of any links pointing to it
158
-func FindLinks(n *Node, k u.Key, start int) []int {
158
+func FindLinks(links []u.Key, k u.Key, start int) []int {
159
var out []int
160
- keybytes := []byte(k)
161
- for i, lnk := range n.Links[start:] {
162
- if bytes.Equal([]byte(lnk.Hash), keybytes) {
160
+ for i, lnk_k := range links[start:] {
161
+ if k == lnk_k {
162
out = append(out, i+start)
163
}
164
}
@@ -170,40 +169,54 @@ func FindLinks(n *Node, k u.Key, start int) []int {
169
// It returns a channel of nodes, which the caller can receive
170
// all the child nodes of 'root' on, in proper order.
171
func (ds *dagService) GetDAG(ctx context.Context, root *Node) <-chan *Node {
172
+ var keys []u.Key
173
+ for _, lnk := range root.Links {
174
+ keys = append(keys, u.Key(lnk.Hash))
175
+ }
176
+
177
+ return ds.GetNodes(ctx, keys)
178
+}
179
+
180
+func (ds *dagService) GetNodes(ctx context.Context, keys []u.Key) <-chan *Node {
181
sig := make(chan *Node)
182
go func() {
183
defer close(sig)
176
-
177
- var keys []u.Key
178
- for _, lnk := range root.Links {
179
- keys = append(keys, u.Key(lnk.Hash))
180
- }
184
blkchan := ds.Blocks.GetBlocks(ctx, keys)
185
183
- nodes := make([]*Node, len(root.Links))
186
+ nodes := make([]*Node, len(keys))
187
next := 0
185
- for blk := range blkchan {
186
- nd, err := Decoded(blk.Data)
187
- if err != nil {
188
- // NB: can occur in normal situations, with improperly formatted
189
- // input data
190
- log.Error("Got back bad block!")
191
- break
192
- }
193
- is := FindLinks(root, blk.Key(), next)
194
- for _, i := range is {
195
- nodes[i] = nd
196
- }
197
-
198
- for ; next < len(nodes) && nodes[next] != nil; next++ {
199
- sig <- nodes[next]
188
+ for {
189
+ select {
190
+ case blk, ok := <-blkchan:
191
+ if !ok {
192
+ if next < len(nodes) {
193
+ log.Errorf("Did not receive correct number of nodes!")
194
+ }
195
+ return
196
+ }
197
+ nd, err := Decoded(blk.Data)
198
+ if err != nil {
199
+ // NB: can occur in normal situations, with improperly formatted
200
+ // input data
201
+ log.Error("Got back bad block!")
202
+ break
203
+ }
204
+ is := FindLinks(keys, blk.Key(), next)
205
+ for _, i := range is {
206
+ nodes[i] = nd
207
+ }
208
+
209
+ for ; next < len(nodes) && nodes[next] != nil; next++ {
210
+ select {
211
+ case sig <- nodes[next]:
212
+ case <-ctx.Done():
213
+ return
214
+ }
215
+ }
216
+ case <-ctx.Done():
217
+ return
218
}
219
}
202
- if next < len(nodes) {
203
- // TODO: bubble errors back up.
204
- log.Errorf("Did not receive correct number of nodes!")
205
- }
220
}()
207
-
221
return sig
222
}
merkledag/merkledag_test.go
+3
-2
@@ -8,6 +8,7 @@ import (
8
"sync"
9
"testing"
10
11
+ "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
12
ds "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-datastore"
13
dssync "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-datastore/sync"
14
bstore "github.com/jbenet/go-ipfs/blocks/blockstore"
@@ -162,7 +163,7 @@ func runBatchFetchTest(t *testing.T, read io.Reader) {
163
164
t.Log("finished setup.")
165
165
- dagr, err := uio.NewDagReader(root, dagservs[0])
166
+ dagr, err := uio.NewDagReader(context.TODO(), root, dagservs[0])
167
if err != nil {
168
t.Fatal(err)
169
}
@@ -195,7 +196,7 @@ func runBatchFetchTest(t *testing.T, read io.Reader) {
196
}
197
fmt.Println("Got first node back.")
198
198
- read, err := uio.NewDagReader(first, dagservs[i])
199
+ read, err := uio.NewDagReader(context.TODO(), first, dagservs[i])
200
if err != nil {
201
t.Fatal(err)
202
}
server/http/ipfs.go
+2
-1
@@ -3,6 +3,7 @@ package http
3
import (
4
"io"
5
6
+ "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
7
core "github.com/jbenet/go-ipfs/core"
8
"github.com/jbenet/go-ipfs/importer"
9
chunk "github.com/jbenet/go-ipfs/importer/chunk"
@@ -36,5 +37,5 @@ func (i *ipfsHandler) AddNodeToDAG(nd *dag.Node) (u.Key, error) {
37
}
38
39
func (i *ipfsHandler) NewDagReader(nd *dag.Node) (io.Reader, error) {
39
- return uio.NewDagReader(nd, i.node.DAG)
40
+ return uio.NewDagReader(context.TODO(), nd, i.node.DAG)
41
}
unixfs/io/dagmodifier_test.go
+5
-4
@@ -6,6 +6,7 @@ import (
6
"io/ioutil"
7
"testing"
8
9
+ "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
10
"github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-datastore/sync"
11
"github.com/jbenet/go-ipfs/blocks/blockstore"
12
bs "github.com/jbenet/go-ipfs/blockservice"
@@ -38,7 +39,7 @@ func getNode(t *testing.T, dserv mdag.DAGService, size int64) ([]byte, *mdag.Nod
39
t.Fatal(err)
40
}
41
41
- dr, err := NewDagReader(node, dserv)
42
+ dr, err := NewDagReader(context.TODO(), node, dserv)
43
if err != nil {
44
t.Fatal(err)
45
}
@@ -75,7 +76,7 @@ func testModWrite(t *testing.T, beg, size uint64, orig []byte, dm *DagModifier)
76
t.Fatal(err)
77
}
78
78
- rd, err := NewDagReader(nd, dm.dagserv)
79
+ rd, err := NewDagReader(context.TODO(), nd, dm.dagserv)
80
if err != nil {
81
t.Fatal(err)
82
}
@@ -173,7 +174,7 @@ func TestMultiWrite(t *testing.T) {
174
t.Fatal(err)
175
}
176
176
- read, err := NewDagReader(nd, dserv)
177
+ read, err := NewDagReader(context.TODO(), nd, dserv)
178
if err != nil {
179
t.Fatal(err)
180
}
@@ -215,7 +216,7 @@ func TestMultiWriteCoal(t *testing.T) {
216
t.Fatal(err)
217
}
218
218
- read, err := NewDagReader(nd, dserv)
219
+ read, err := NewDagReader(context.TODO(), nd, dserv)
220
if err != nil {
221
t.Fatal(err)
222
}
unixfs/io/dagreader.go
+120
-16
@@ -3,7 +3,10 @@ package io
3
import (
4
"bytes"
5
"errors"
6
+ "fmt"
7
"io"
8
+ "io/ioutil"
9
+ "os"
10
11
"github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
12
@@ -11,6 +14,7 @@ import (
14
mdag "github.com/jbenet/go-ipfs/merkledag"
15
ft "github.com/jbenet/go-ipfs/unixfs"
16
ftpb "github.com/jbenet/go-ipfs/unixfs/pb"
17
+ u "github.com/jbenet/go-ipfs/util"
18
)
19
20
var ErrIsDir = errors.New("this dag node is a directory")
@@ -19,14 +23,28 @@ var ErrIsDir = errors.New("this dag node is a directory")
23
type DagReader struct {
24
serv mdag.DAGService
25
node *mdag.Node
22
- buf io.Reader
26
+ buf ReadSeekCloser
27
fetchChan <-chan *mdag.Node
28
linkPosition int
29
+ offset int64
30
+
31
+ // Our context
32
+ ctx context.Context
33
+
34
+ // Context for children
35
+ fctx context.Context
36
+ cancel func()
37
+}
38
+
39
+type ReadSeekCloser interface {
40
+ io.Reader
41
+ io.Seeker
42
+ io.Closer
43
}
44
45
// NewDagReader creates a new reader object that reads the data represented by the given
46
// node, using the passed in DAGService for data retreival
29
-func NewDagReader(n *mdag.Node, serv mdag.DAGService) (io.Reader, error) {
47
+func NewDagReader(ctx context.Context, n *mdag.Node, serv mdag.DAGService) (ReadSeekCloser, error) {
48
pb := new(ftpb.Data)
49
err := proto.Unmarshal(n.Data, pb)
50
if err != nil {
@@ -38,16 +56,20 @@ func NewDagReader(n *mdag.Node, serv mdag.DAGService) (io.Reader, error) {
56
// Dont allow reading directories
57
return nil, ErrIsDir
58
case ftpb.Data_File:
41
- fetchChan := serv.GetDAG(context.TODO(), n)
59
+ fctx, cancel := context.WithCancel(ctx)
60
+ fetchChan := serv.GetDAG(fctx, n)
61
return &DagReader{
62
node: n,
63
serv: serv,
45
- buf: bytes.NewBuffer(pb.GetData()),
64
+ buf: NewRSNCFromBytes(pb.GetData()),
65
fetchChan: fetchChan,
66
+ ctx: ctx,
67
+ fctx: fctx,
68
+ cancel: cancel,
69
}, nil
70
case ftpb.Data_Raw:
71
// Raw block will just be a single level, return a byte buffer
50
- return bytes.NewBuffer(pb.GetData()), nil
72
+ return NewRSNCFromBytes(pb.GetData()), nil
73
default:
74
return nil, ft.ErrUnrecognizedType
75
}
@@ -70,6 +92,8 @@ func (dr *DagReader) precalcNextBuf() error {
92
if !ok {
93
return io.EOF
94
}
95
+ case <-dr.ctx.Done():
96
+ return dr.ctx.Err()
97
}
98
99
pb := new(ftpb.Data)
@@ -85,20 +109,37 @@ func (dr *DagReader) precalcNextBuf() error {
109
case ftpb.Data_File:
110
//TODO: this *should* work, needs testing first
111
log.Warning("Running untested code for multilayered indirect FS reads.")
88
- subr, err := NewDagReader(nxt, dr.serv)
112
+ subr, err := NewDagReader(dr.fctx, nxt, dr.serv)
113
if err != nil {
114
return err
115
}
116
dr.buf = subr
117
return nil
118
case ftpb.Data_Raw:
95
- dr.buf = bytes.NewBuffer(pb.GetData())
119
+ dr.buf = NewRSNCFromBytes(pb.GetData())
120
return nil
121
default:
122
return ft.ErrUnrecognizedType
123
}
124
}
125
126
+func (dr *DagReader) resetBlockFetch(nlinkpos int) {
127
+ dr.cancel()
128
+ dr.fetchChan = nil
129
+ dr.linkPosition = nlinkpos
130
+
131
+ var keys []u.Key
132
+ for _, lnk := range dr.node.Links[dr.linkPosition:] {
133
+ keys = append(keys, u.Key(lnk.Hash))
134
+ }
135
+
136
+ fctx, cancel := context.WithCancel(dr.ctx)
137
+ dr.cancel = cancel
138
+ dr.fctx = fctx
139
+ fch := dr.serv.GetNodes(fctx, keys)
140
+ dr.fetchChan = fch
141
+}
142
+
143
// Read reads data from the DAG structured file
144
func (dr *DagReader) Read(b []byte) (int, error) {
145
// If no cached buffer, load one
@@ -113,6 +154,7 @@ func (dr *DagReader) Read(b []byte) (int, error) {
154
// Attempt to fill bytes from cached buffer
155
n, err := dr.buf.Read(b[total:])
156
total += n
157
+ dr.offset += int64(n)
158
if err != nil {
159
// EOF is expected
160
if err != io.EOF {
@@ -133,28 +175,90 @@ func (dr *DagReader) Read(b []byte) (int, error) {
175
}
176
}
177
136
-/*
178
+func (dr *DagReader) Close() error {
179
+ if dr.fctx != nil {
180
+ dr.cancel()
181
+ }
182
+ return nil
183
+}
184
+
185
func (dr *DagReader) Seek(offset int64, whence int) (int64, error) {
186
switch whence {
187
case os.SEEK_SET:
140
- for i := 0; i < len(dr.node.Links); i++ {
141
- nsize := dr.node.Links[i].Size - 8
142
- if offset > nsize {
143
- offset -= nsize
144
- } else {
188
+ if offset < 0 {
189
+ return -1, errors.New("Invalid offset")
190
+ }
191
+ //TODO: this pb should be cached
192
+ pb := new(ftpb.Data)
193
+ err := proto.Unmarshal(dr.node.Data, pb)
194
+ if err != nil {
195
+ return -1, err
196
+ }
197
+
198
+ if offset == 0 {
199
+ dr.resetBlockFetch(0)
200
+ dr.buf = NewRSNCFromBytes(pb.GetData())
201
+ return 0, nil
202
+ }
203
+
204
+ left := offset
205
+ if int64(len(pb.Data)) > offset {
206
+ dr.buf = NewRSNCFromBytes(pb.GetData()[offset:])
207
+ dr.linkPosition = 0
208
+ dr.offset = offset
209
+ return offset, nil
210
+ } else {
211
+ left -= int64(len(pb.Data))
212
+ }
213
+
214
+ i := 0
215
+ for ; i < len(pb.Blocksizes); i++ {
216
+ if pb.Blocksizes[i] > uint64(left) {
217
break
218
+ } else {
219
+ left -= int64(pb.Blocksizes[i])
220
}
221
}
148
- dr.position = i
149
- err := dr.precalcNextBuf()
222
+ dr.resetBlockFetch(i)
223
+ err = dr.precalcNextBuf()
224
if err != nil {
225
return 0, err
226
}
227
+
228
+ n, err := io.CopyN(ioutil.Discard, dr.buf, left)
229
+ if err != nil {
230
+ fmt.Printf("the copy failed: %s - [%d]\n", err, n)
231
+ return -1, err
232
+ }
233
+ left -= n
234
+ if left != 0 {
235
+ return -1, errors.New("failed to seek properly")
236
+ }
237
+ dr.offset = offset
238
+ return offset, nil
239
case os.SEEK_CUR:
240
+ noffset := dr.offset + offset
241
+ return dr.Seek(noffset, os.SEEK_SET)
242
case os.SEEK_END:
243
+ pb := new(ftpb.Data)
244
+ err := proto.Unmarshal(dr.node.Data, pb)
245
+ if err != nil {
246
+ return -1, err
247
+ }
248
+ noffset := int64(pb.GetFilesize()) - offset
249
+ return dr.Seek(noffset, os.SEEK_SET)
250
default:
251
return 0, errors.New("invalid whence")
252
}
253
return 0, nil
254
}
160
-*/
255
+
256
+type readSeekNopCloser struct {
257
+ *bytes.Reader
258
+}
259
+
260
+func NewRSNCFromBytes(b []byte) ReadSeekCloser {
261
+ return &readSeekNopCloser{bytes.NewReader(b)}
262
+}
263
+
264
+func (r *readSeekNopCloser) Close() error { return nil }
unixfs/tar/reader.go
+2
-1
@@ -4,6 +4,7 @@ import (
4
"archive/tar"
5
"bytes"
6
"compress/gzip"
7
+ "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
8
"io"
9
gopath "path"
10
"strings"
@@ -114,7 +115,7 @@ func (i *Reader) writeToBuf(dagnode *mdag.Node, path string, depth int) {
115
}
116
i.flush()
117
117
- reader, err := uio.NewDagReader(dagnode, i.dag)
118
+ reader, err := uio.NewDagReader(context.TODO(), dagnode, i.dag)
119
if err != nil {
120
i.emitError(err)
121
return