@cryptotaxi247 / kubo / commits / 092df586d

clear out memory after reads from the dagreader

License: MIT Signed-off-by: Jeromy <jeromyj@gmail.com>

Jeromy committed Dec 26, 2017 at 14:30 UTC 092df586d13a131edebe0e822f083c2f99601f0a
1 file changed +56 -4
unixfs/io/pbdagreader.go
+56 -4
@@ -10,7 +10,9 @@ import (
10 ft "github.com/ipfs/go-ipfs/unixfs"
11 ftpb "github.com/ipfs/go-ipfs/unixfs/pb"
12
13 + node "gx/ipfs/QmNwUEK7QbwSqyKBu3mMtToo8SUc6wQJ7gdZq4gGGJqfnf/go-ipld-format"
14 proto "gx/ipfs/QmZ4Qi3GaRbjcx28Sme5eMH7RQjGkt8wHxt2a65oLaeFEV/gogo-protobuf/proto"
15 + cid "gx/ipfs/QmeSrf6pzut73u6zLQkRFQ3ygt3k6XFT2kjdYP8Tnkwwyg/go-cid"
16 )
17
18 // DagReader provides a way to easily read the data contained in a dag.
@@ -30,6 +32,9 @@ type pbDagReader struct {
32 // NodeGetters for each of 'nodes' child links
33 promises []mdag.NodeGetter
34
35 + // the cid of each child of the current node
36 + links []*cid.Cid
37 +
38 // the index of the child link currently being read from
39 linkPosition int
40
@@ -47,30 +52,54 @@ var _ DagReader = (*pbDagReader)(nil)
52
53 func NewPBFileReader(ctx context.Context, n *mdag.ProtoNode, pb *ftpb.Data, serv mdag.DAGService) *pbDagReader {
54 fctx, cancel := context.WithCancel(ctx)
50 - promises := mdag.GetDAG(fctx, serv, n)
55 + curLinks := getLinkCids(n)
56 return &pbDagReader{
57 node: n,
58 serv: serv,
59 buf: NewBufDagReader(pb.GetData()),
55 - promises: promises,
60 + promises: make([]mdag.NodeGetter, len(curLinks)),
61 + links: curLinks,
62 ctx: fctx,
63 cancel: cancel,
64 pbdata: pb,
65 }
66 }
67
68 +const preloadSize = 10
69 +
70 +func (dr *pbDagReader) preloadNextNodes(ctx context.Context) {
71 + beg := dr.linkPosition
72 + end := beg + preloadSize
73 + if end >= len(dr.links) {
74 + end = len(dr.links)
75 + }
76 +
77 + for i, p := range mdag.GetNodes(ctx, dr.serv, dr.links[beg:end]) {
78 + dr.promises[beg+i] = p
79 + }
80 +}
81 +
82 // precalcNextBuf follows the next link in line and loads it from the
83 // DAGService, setting the next buffer to read from
84 func (dr *pbDagReader) precalcNextBuf(ctx context.Context) error {
65 - dr.buf.Close() // Just to make sure
85 + if dr.buf != nil {
86 + dr.buf.Close() // Just to make sure
87 + dr.buf = nil
88 + }
89 +
90 if dr.linkPosition >= len(dr.promises) {
91 return io.EOF
92 }
93
94 + if dr.promises[dr.linkPosition] == nil {
95 + dr.preloadNextNodes(ctx)
96 + }
97 +
98 nxt, err := dr.promises[dr.linkPosition].Get(ctx)
99 if err != nil {
100 return err
101 }
102 + dr.promises[dr.linkPosition] = nil
103 dr.linkPosition++
104
105 switch nxt := nxt.(type) {
@@ -105,6 +134,15 @@ func (dr *pbDagReader) precalcNextBuf(ctx context.Context) error {
134 }
135 }
136
137 +func getLinkCids(n node.Node) []*cid.Cid {
138 + links := n.Links()
139 + out := make([]*cid.Cid, 0, len(links))
140 + for _, l := range links {
141 + out = append(out, l.Cid)
142 + }
143 + return out
144 +}
145 +
146 // Size return the total length of the data from the DAG structured file.
147 func (dr *pbDagReader) Size() uint64 {
148 return dr.pbdata.GetFilesize()
@@ -117,6 +155,12 @@ func (dr *pbDagReader) Read(b []byte) (int, error) {
155
156 // CtxReadFull reads data from the DAG structured file
157 func (dr *pbDagReader) CtxReadFull(ctx context.Context, b []byte) (int, error) {
158 + if dr.buf == nil {
159 + if err := dr.precalcNextBuf(ctx); err != nil {
160 + return 0, err
161 + }
162 + }
163 +
164 // If no cached buffer, load one
165 total := 0
166 for {
@@ -145,6 +189,12 @@ func (dr *pbDagReader) CtxReadFull(ctx context.Context, b []byte) (int, error) {
189 }
190
191 func (dr *pbDagReader) WriteTo(w io.Writer) (int64, error) {
192 + if dr.buf == nil {
193 + if err := dr.precalcNextBuf(dr.ctx); err != nil {
194 + return 0, err
195 + }
196 + }
197 +
198 // If no cached buffer, load one
199 total := int64(0)
200 for {
@@ -199,7 +249,9 @@ func (dr *pbDagReader) Seek(offset int64, whence int) (int64, error) {
249 left := offset
250 if int64(len(pb.Data)) >= offset {
251 // Close current buf to close potential child dagreader
202 - dr.buf.Close()
252 + if dr.buf != nil {
253 + dr.buf.Close()
254 + }
255 dr.buf = NewBufDagReader(pb.GetData()[offset:])
256
257 // start reading links from the beginning