diff --git a/TODO.txt b/TODO.txt index d7ba8ad..fadc44f 100644 --- a/TODO.txt +++ b/TODO.txt @@ -2,11 +2,18 @@ * printf '\x6c' | dd seek=20 bs=1 count=1 conv=notrunc of=test.dump * crypto - * flush last chunk in streaming mode + * flush last chunk in streaming mode - check if implemented + +* web + * listing of all files in cache + * add breadcrumb + * display size of files + +* tuning + * debug logger for writer queue size * sftp * implement read ahead buffer? -* add breadcrumb * config file? * whitelist allowed src paths diff --git a/pkg/fs/discard.go b/pkg/fs/discard.go index fdf049d..ecf2769 100644 --- a/pkg/fs/discard.go +++ b/pkg/fs/discard.go @@ -4,7 +4,6 @@ package fs import ( "io" - "sync" ) var Discard io.Writer = discard{} @@ -21,30 +20,33 @@ func (discard) WriteString(s string) (int, error) { return len(s), nil } -const DefaultBufferSize = 8192 +const ( + StreamingBufferSize = 32768 + DefaultBufferSize = 8192 +) -var blackHolePool = sync.Pool{ - New: func() any { - b := make([]byte, DefaultBufferSize) +func poolNewFunc(size int) func() *[]byte { + return func() *[]byte { + b := make([]byte, size) return &b - }, + } } +var ( + streamingPool = newPool(poolNewFunc(StreamingBufferSize)) + blackHolePool = newPool(poolNewFunc(DefaultBufferSize)) +) + func SetBufferSize(size int) bool { if size < 0 || size == DefaultBufferSize { return false } - blackHolePool = sync.Pool{ - New: func() any { - b := make([]byte, size) - return &b - }, - } + blackHolePool = newPool(poolNewFunc(size)) return true } func (discard) ReadFrom(r io.Reader) (n int64, err error) { - bufp := blackHolePool.Get().(*[]byte) + bufp := blackHolePool.Get() readSize := 0 for { readSize, err = r.Read(*bufp) diff --git a/pkg/fs/file.go b/pkg/fs/file.go index e6db254..bf7909a 100644 --- a/pkg/fs/file.go +++ b/pkg/fs/file.go @@ -134,9 +134,7 @@ func (f *File) readToCache(p []byte) (n int, err error) { n, err = f.readSource(p) if n > 0 { // async cache write - buf := make([]byte, n) - copy(buf, p[:n]) - if _, err = f.md.WriteAt(buf, f.offset); err != nil { + if _, err = f.md.WriteAt(p[:n], f.offset); err != nil { log.Error(err, "error writing cache file") return } diff --git a/pkg/fs/metadata.go b/pkg/fs/metadata.go index 4390be0..57c37fc 100644 --- a/pkg/fs/metadata.go +++ b/pkg/fs/metadata.go @@ -414,6 +414,7 @@ func (md *Metadata) addChunk(off int64, n int) { type writeAt struct { pos int64 data []byte + ret *[]byte } func (md *Metadata) writer() { @@ -428,6 +429,9 @@ func (md *Metadata) writer() { } else { md.err = err } + if wa.ret != nil { + streamingPool.Put(wa.ret) + } } }() } @@ -442,8 +446,20 @@ func (md *Metadata) WriteAt(data []byte, pos int64) (int, error) { return 0, err } } - md.wc <- writeAt{pos, data} - return len(data), nil + wa := writeAt{pos: pos} + ret := streamingPool.Get() + n := len(data) + buf := *ret + if n > len(buf) { + streamingPool.Put(ret) + buf = make([]byte, n) + } else { + wa.ret = ret + } + copy(buf, data) + wa.data = buf[:n] + md.wc <- wa + return n, nil } func (md *Metadata) openCacheFile() error { diff --git a/pkg/fs/pool.go b/pkg/fs/pool.go new file mode 100644 index 0000000..a6ee514 --- /dev/null +++ b/pkg/fs/pool.go @@ -0,0 +1,27 @@ +// Copyright (C) 2022 Marius Schellenberger + +package fs + +import ( + "sync" +) + +type pool[T any] struct { + p sync.Pool +} + +func newPool[T any](newf func() T) *pool[T] { + return &pool[T]{ + p: sync.Pool{ + New: func() any { return newf() }, + }, + } +} + +func (p *pool[T]) Get() T { + return p.p.Get().(T) +} + +func (p *pool[T]) Put(v T) { + p.p.Put(v) +}