@cryptotaxi247 / kubo / commits / 1e93ee00c

clean up benchmarks, implement WriterTo on DAGReader, and optimize DagReader

Jeromy committed Feb 4, 2015 at 19:43 UTC 1e93ee00c0751bbd42188ea91fd4f5b5723e49a8
2 files changed +60 -66
importer/importer_test.go
+16 -50
@@ -65,21 +65,9 @@ func BenchmarkBalancedReadSmallBlock(b *testing.B) {
65 nbytes := int64(10000000)
66 nd, ds := getBalancedDag(b, nbytes, 4096)
67
68 - b.StartTimer()
69 - for i := 0; i < b.N; i++ {
70 - read, err := uio.NewDagReader(context.TODO(), nd, ds)
71 - if err != nil {
72 - b.Fatal(err)
73 - }
74 - n, err := io.Copy(ioutil.Discard, read)
75 - if err != nil {
76 - b.Fatal(err)
77 - }
78 - if n != nbytes {
79 - b.Fatal("Failed to read correct amount")
80 - }
81 - }
68 b.SetBytes(nbytes)
69 + b.StartTimer()
70 + runReadBench(b, nd, ds)
71 }
72
73 func BenchmarkTrickleReadSmallBlock(b *testing.B) {
@@ -87,22 +75,9 @@ func BenchmarkTrickleReadSmallBlock(b *testing.B) {
75 nbytes := int64(10000000)
76 nd, ds := getTrickleDag(b, nbytes, 4096)
77
90 - b.StartTimer()
91 - for i := 0; i < b.N; i++ {
92 - read, err := uio.NewDagReader(context.TODO(), nd, ds)
93 - if err != nil {
94 - b.Fatal(err)
95 - }
96 -
97 - n, err := io.Copy(ioutil.Discard, read)
98 - if err != nil {
99 - b.Fatal(err)
100 - }
101 - if n != nbytes {
102 - b.Fatal("Failed to read correct amount")
103 - }
104 - }
78 b.SetBytes(nbytes)
79 + b.StartTimer()
80 + runReadBench(b, nd, ds)
81 }
82
83 func BenchmarkBalancedReadFull(b *testing.B) {
@@ -110,21 +85,9 @@ func BenchmarkBalancedReadFull(b *testing.B) {
85 nbytes := int64(10000000)
86 nd, ds := getBalancedDag(b, nbytes, chunk.DefaultBlockSize)
87
113 - b.StartTimer()
114 - for i := 0; i < b.N; i++ {
115 - read, err := uio.NewDagReader(context.TODO(), nd, ds)
116 - if err != nil {
117 - b.Fatal(err)
118 - }
119 - n, err := io.Copy(ioutil.Discard, read)
120 - if err != nil {
121 - b.Fatal(err)
122 - }
123 - if n != nbytes {
124 - b.Fatal("Failed to read correct amount")
125 - }
126 - }
88 b.SetBytes(nbytes)
89 + b.StartTimer()
90 + runReadBench(b, nd, ds)
91 }
92
93 func BenchmarkTrickleReadFull(b *testing.B) {
@@ -132,20 +95,23 @@ func BenchmarkTrickleReadFull(b *testing.B) {
95 nbytes := int64(10000000)
96 nd, ds := getTrickleDag(b, nbytes, chunk.DefaultBlockSize)
97
98 + b.SetBytes(nbytes)
99 b.StartTimer()
100 + runReadBench(b, nd, ds)
101 +}
102 +
103 +func runReadBench(b *testing.B, nd *dag.Node, ds dag.DAGService) {
104 for i := 0; i < b.N; i++ {
137 - read, err := uio.NewDagReader(context.TODO(), nd, ds)
105 + ctx, cancel := context.WithCancel(context.TODO())
106 + read, err := uio.NewDagReader(ctx, nd, ds)
107 if err != nil {
108 b.Fatal(err)
109 }
110
142 - n, err := io.Copy(ioutil.Discard, read)
143 - if err != nil {
111 + _, err = read.WriteTo(ioutil.Discard)
112 + if err != nil && err != io.EOF {
113 b.Fatal(err)
114 }
146 - if n != nbytes {
147 - b.Fatal("Failed to read correct amount")
148 - }
115 + cancel()
116 }
150 - b.SetBytes(nbytes)
117 }
unixfs/io/dagreader.go
+44 -16
@@ -50,6 +50,7 @@ type ReadSeekCloser interface {
50 io.Reader
51 io.Seeker
52 io.Closer
53 + io.WriterTo
54 }
55
56 // NewDagReader creates a new reader object that reads the data represented by the given
@@ -68,22 +69,26 @@ func NewDagReader(ctx context.Context, n *mdag.Node, serv mdag.DAGService) (*Dag
69 case ftpb.Data_Raw:
70 fallthrough
71 case ftpb.Data_File:
71 - fctx, cancel := context.WithCancel(ctx)
72 - promises := serv.GetDAG(fctx, n)
73 - return &DagReader{
74 - node: n,
75 - serv: serv,
76 - buf: NewRSNCFromBytes(pb.GetData()),
77 - promises: promises,
78 - ctx: fctx,
79 - cancel: cancel,
80 - pbdata: pb,
81 - }, nil
72 + return newDataFileReader(ctx, n, pb, serv), nil
73 default:
74 return nil, ft.ErrUnrecognizedType
75 }
76 }
77
78 +func newDataFileReader(ctx context.Context, n *mdag.Node, pb *ftpb.Data, serv mdag.DAGService) *DagReader {
79 + fctx, cancel := context.WithCancel(ctx)
80 + promises := serv.GetDAG(fctx, n)
81 + return &DagReader{
82 + node: n,
83 + serv: serv,
84 + buf: NewRSNCFromBytes(pb.GetData()),
85 + promises: promises,
86 + ctx: fctx,
87 + cancel: cancel,
88 + pbdata: pb,
89 + }
90 +}
91 +
92 // precalcNextBuf follows the next link in line and loads it from the DAGService,
93 // setting the next buffer to read from
94 func (dr *DagReader) precalcNextBuf() error {
@@ -108,11 +113,7 @@ func (dr *DagReader) precalcNextBuf() error {
113 // A directory should not exist within a file
114 return ft.ErrInvalidDirLocation
115 case ftpb.Data_File:
111 - subr, err := NewDagReader(dr.ctx, nxt, dr.serv)
112 - if err != nil {
113 - return err
114 - }
115 - dr.buf = subr
116 + dr.buf = newDataFileReader(dr.ctx, nxt, pb, dr.serv)
117 return nil
118 case ftpb.Data_Raw:
119 dr.buf = NewRSNCFromBytes(pb.GetData())
@@ -156,6 +157,31 @@ func (dr *DagReader) Read(b []byte) (int, error) {
157 }
158 }
159
160 +func (dr *DagReader) WriteTo(w io.Writer) (int64, error) {
161 + // If no cached buffer, load one
162 + total := int64(0)
163 + for {
164 + // Attempt to write bytes from cached buffer
165 + n, err := dr.buf.WriteTo(w)
166 + total += n
167 + dr.offset += n
168 + if err != nil {
169 + if err != io.EOF {
170 + return total, err
171 + }
172 + }
173 +
174 + // Otherwise, load up the next block
175 + err = dr.precalcNextBuf()
176 + if err != nil {
177 + if err == io.EOF {
178 + return total, nil
179 + }
180 + return total, err
181 + }
182 + }
183 +}
184 +
185 func (dr *DagReader) Close() error {
186 dr.cancel()
187 return nil
@@ -163,6 +189,8 @@ func (dr *DagReader) Close() error {
189
190 // Seek implements io.Seeker, and will seek to a given offset in the file
191 // interface matches standard unix seek
192 +// TODO: check if we can do relative seeks, to reduce the amount of dagreader
193 +// recreations that need to happen.
194 func (dr *DagReader) Seek(offset int64, whence int) (int64, error) {
195 switch whence {
196 case os.SEEK_SET: