added preload and streaming buffer pool
This commit is contained in:
parent
cd0b0254c7
commit
b010b3a93e
5 changed files with 70 additions and 20 deletions
11
TODO.txt
11
TODO.txt
|
|
@ -2,11 +2,18 @@
|
||||||
* printf '\x6c' | dd seek=20 bs=1 count=1 conv=notrunc of=test.dump
|
* printf '\x6c' | dd seek=20 bs=1 count=1 conv=notrunc of=test.dump
|
||||||
|
|
||||||
* crypto
|
* 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
|
* sftp
|
||||||
* implement read ahead buffer?
|
* implement read ahead buffer?
|
||||||
|
|
||||||
* add breadcrumb
|
|
||||||
* config file?
|
* config file?
|
||||||
* whitelist allowed src paths
|
* whitelist allowed src paths
|
||||||
|
|
|
||||||
|
|
@ -4,7 +4,6 @@ package fs
|
||||||
|
|
||||||
import (
|
import (
|
||||||
"io"
|
"io"
|
||||||
"sync"
|
|
||||||
)
|
)
|
||||||
|
|
||||||
var Discard io.Writer = discard{}
|
var Discard io.Writer = discard{}
|
||||||
|
|
@ -21,30 +20,33 @@ func (discard) WriteString(s string) (int, error) {
|
||||||
return len(s), nil
|
return len(s), nil
|
||||||
}
|
}
|
||||||
|
|
||||||
const DefaultBufferSize = 8192
|
const (
|
||||||
|
StreamingBufferSize = 32768
|
||||||
|
DefaultBufferSize = 8192
|
||||||
|
)
|
||||||
|
|
||||||
var blackHolePool = sync.Pool{
|
func poolNewFunc(size int) func() *[]byte {
|
||||||
New: func() any {
|
return func() *[]byte {
|
||||||
b := make([]byte, DefaultBufferSize)
|
b := make([]byte, size)
|
||||||
return &b
|
return &b
|
||||||
},
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
var (
|
||||||
|
streamingPool = newPool(poolNewFunc(StreamingBufferSize))
|
||||||
|
blackHolePool = newPool(poolNewFunc(DefaultBufferSize))
|
||||||
|
)
|
||||||
|
|
||||||
func SetBufferSize(size int) bool {
|
func SetBufferSize(size int) bool {
|
||||||
if size < 0 || size == DefaultBufferSize {
|
if size < 0 || size == DefaultBufferSize {
|
||||||
return false
|
return false
|
||||||
}
|
}
|
||||||
blackHolePool = sync.Pool{
|
blackHolePool = newPool(poolNewFunc(size))
|
||||||
New: func() any {
|
|
||||||
b := make([]byte, size)
|
|
||||||
return &b
|
|
||||||
},
|
|
||||||
}
|
|
||||||
return true
|
return true
|
||||||
}
|
}
|
||||||
|
|
||||||
func (discard) ReadFrom(r io.Reader) (n int64, err error) {
|
func (discard) ReadFrom(r io.Reader) (n int64, err error) {
|
||||||
bufp := blackHolePool.Get().(*[]byte)
|
bufp := blackHolePool.Get()
|
||||||
readSize := 0
|
readSize := 0
|
||||||
for {
|
for {
|
||||||
readSize, err = r.Read(*bufp)
|
readSize, err = r.Read(*bufp)
|
||||||
|
|
|
||||||
|
|
@ -134,9 +134,7 @@ func (f *File) readToCache(p []byte) (n int, err error) {
|
||||||
n, err = f.readSource(p)
|
n, err = f.readSource(p)
|
||||||
if n > 0 {
|
if n > 0 {
|
||||||
// async cache write
|
// async cache write
|
||||||
buf := make([]byte, n)
|
if _, err = f.md.WriteAt(p[:n], f.offset); err != nil {
|
||||||
copy(buf, p[:n])
|
|
||||||
if _, err = f.md.WriteAt(buf, f.offset); err != nil {
|
|
||||||
log.Error(err, "error writing cache file")
|
log.Error(err, "error writing cache file")
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -414,6 +414,7 @@ func (md *Metadata) addChunk(off int64, n int) {
|
||||||
type writeAt struct {
|
type writeAt struct {
|
||||||
pos int64
|
pos int64
|
||||||
data []byte
|
data []byte
|
||||||
|
ret *[]byte
|
||||||
}
|
}
|
||||||
|
|
||||||
func (md *Metadata) writer() {
|
func (md *Metadata) writer() {
|
||||||
|
|
@ -428,6 +429,9 @@ func (md *Metadata) writer() {
|
||||||
} else {
|
} else {
|
||||||
md.err = err
|
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
|
return 0, err
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
md.wc <- writeAt{pos, data}
|
wa := writeAt{pos: pos}
|
||||||
return len(data), nil
|
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 {
|
func (md *Metadata) openCacheFile() error {
|
||||||
|
|
|
||||||
27
pkg/fs/pool.go
Normal file
27
pkg/fs/pool.go
Normal file
|
|
@ -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)
|
||||||
|
}
|
||||||
Loading…
Add table
Add a link
Reference in a new issue