@cryptotaxi247 / kubo / commits / e7aa1166b

add writerAt for fuse writes

Jeromy committed Oct 3, 2014 at 23:04 UTC e7aa1166bcbaa57e20512467cb6661bca2da4aeb
2 files changed +51 -18
fuse/ipns/ipns_unix.go
+22 -18
@@ -1,14 +1,13 @@
1 package ipns
2
3 import (
4 + "bytes"
5 "fmt"
6 "io/ioutil"
7 "os"
8 "path/filepath"
9 "time"
10
10 - "bytes"
11 -
11 "github.com/jbenet/go-ipfs/Godeps/_workspace/src/bazil.org/fuse"
12 "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/goprotobuf/proto"
@@ -205,7 +204,7 @@ type Node struct {
204 cached *mdag.PBData
205
206 // For writing
208 - dataBuf *bytes.Buffer
207 + writerBuf WriteAtBuf
208 }
209
210 func (s *Node) loadData() error {
@@ -290,21 +289,25 @@ func (s *Node) ReadAll(intr fs.Intr) ([]byte, fuse.Error) {
289 }
290 // this is a terrible function... 'ReadAll'?
291 // what if i have a 6TB file? GG RAM.
293 - return ioutil.ReadAll(r)
292 + b, err := ioutil.ReadAll(r)
293 + if err != nil {
294 + log.Error("[%s] Readall error: %s", s.name, err)
295 + return nil, err
296 + }
297 + if len(b) > 4 {
298 + log.Debug("ReadAll trailing bytes: %v", b[len(b)-4:])
299 + }
300 + return b, nil
301 }
302
303 func (n *Node) Write(req *fuse.WriteRequest, resp *fuse.WriteResponse, intr fs.Intr) fuse.Error {
304 log.Debug("ipns: Node Write [%s]: flags = %s, offset = %d, size = %d", n.name, req.Flags.String(), req.Offset, len(req.Data))
298 - if n.dataBuf == nil {
299 - n.dataBuf = new(bytes.Buffer)
305 + if n.writerBuf == nil {
306 + n.writerBuf = NewWriterAtFromBytes(nil)
307 }
301 - if req.Offset == 0 {
302 - n.dataBuf.Reset()
303 - n.dataBuf.Write(req.Data)
304 - resp.Size = len(req.Data)
305 - } else {
306 - log.Error("Unhandled write to offset!")
307 - n.dataBuf = nil
308 + _, err := n.writerBuf.WriteAt(req.Data, req.Offset)
309 + if err != nil {
310 + return err
311 }
312 return nil
313 }
@@ -312,7 +315,7 @@ func (n *Node) Write(req *fuse.WriteRequest, resp *fuse.WriteResponse, intr fs.I
315 func (n *Node) Flush(req *fuse.FlushRequest, intr fs.Intr) fuse.Error {
316 log.Debug("Got flush request [%s]!", n.name)
317
315 - if n.dataBuf != nil {
318 + if n.writerBuf != nil {
319 //TODO:
320 // This operation holds everything in memory,
321 // should be changed to stream the block creation/storage
@@ -321,9 +324,10 @@ func (n *Node) Flush(req *fuse.FlushRequest, intr fs.Intr) fuse.Error {
324 //NOTE:
325 // This should only occur on a file object, if this were to be a
326 // folder, bad things would happen.
324 - newNode, err := imp.NewDagFromReader(n.dataBuf)
327 + buf := bytes.NewReader(n.writerBuf.Bytes())
328 + newNode, err := imp.NewDagFromReader(buf)
329 if err != nil {
326 - log.Critical("error creating dag from dataBuf: %s", err)
330 + log.Critical("error creating dag from writerBuf: %s", err)
331 return fuse.ENODATA
332 }
333 if n.parent != nil {
@@ -350,7 +354,7 @@ func (n *Node) Flush(req *fuse.FlushRequest, intr fs.Intr) fuse.Error {
354 fmt.Println(string(b))
355 //
356
353 - n.dataBuf = nil
357 + n.writerBuf = nil
358
359 n.wasChanged()
360 }
@@ -382,7 +386,7 @@ func (n *Node) republishRoot() error {
386 return err
387 }
388
385 - n.dataBuf = nil
389 + n.writerBuf = nil
390
391 ndkey, err := root.Nd.Key()
392 if err != nil {
fuse/ipns/writerat.go new
+29
@@ -0,0 +1,29 @@
1 +package ipns
2 +
3 +import "io"
4 +
5 +type WriteAtBuf interface {
6 + io.WriterAt
7 + Bytes() []byte
8 +}
9 +
10 +type writerAt struct {
11 + buf []byte
12 +}
13 +
14 +func NewWriterAtFromBytes(b []byte) WriteAtBuf {
15 + return &writerAt{b}
16 +}
17 +
18 +// TODO: make this better in the future, this is just a quick hack for now
19 +func (wa *writerAt) WriteAt(p []byte, off int64) (int, error) {
20 + if off+int64(len(p)) > int64(len(wa.buf)) {
21 + wa.buf = append(wa.buf, make([]byte, (int(off)+len(p))-len(wa.buf))...)
22 + }
23 + copy(wa.buf[off:], p)
24 + return len(p), nil
25 +}
26 +
27 +func (wa *writerAt) Bytes() []byte {
28 + return wa.buf
29 +}