added custom buffer size error handling
This commit is contained in:
parent
28232d976b
commit
711f020c18
4 changed files with 123 additions and 10 deletions
|
|
@ -30,6 +30,7 @@ var (
|
||||||
listenCache string
|
listenCache string
|
||||||
quota int64
|
quota int64
|
||||||
max int
|
max int
|
||||||
|
bs int
|
||||||
|
|
||||||
log logr.Logger
|
log logr.Logger
|
||||||
)
|
)
|
||||||
|
|
@ -45,12 +46,17 @@ func main() {
|
||||||
flag.StringVar(&listenWebdav, "webdav", "", "listen addr:port for webdav")
|
flag.StringVar(&listenWebdav, "webdav", "", "listen addr:port for webdav")
|
||||||
flag.StringVar(&listenCache, "cache", "", "listen addr:port for cache only")
|
flag.StringVar(&listenCache, "cache", "", "listen addr:port for cache only")
|
||||||
flag.IntVar(&max, "max", -1, "max parallel preloads")
|
flag.IntVar(&max, "max", -1, "max parallel preloads")
|
||||||
|
flag.IntVar(&bs, "bs", -1, "tune reload buffer size in bytes (default: 8192)")
|
||||||
flag.Int64Var("a, "quota", 1, "max disk usage quota in GiB")
|
flag.Int64Var("a, "quota", 1, "max disk usage quota in GiB")
|
||||||
flag.Parse()
|
flag.Parse()
|
||||||
|
|
||||||
log = klogr.New().WithName("main")
|
log = klogr.New().WithName("main")
|
||||||
log.Info("starting cachefs", "version", version)
|
log.Info("starting cachefs", "version", version)
|
||||||
|
|
||||||
|
if fs.SetBufferSize(bs) {
|
||||||
|
log.V(2).Info("changed preload buffer size", "size", bs)
|
||||||
|
}
|
||||||
|
|
||||||
filesystem, err := fs.NewFS(quota*gib, max, src, dst, metadata, klogr.New().WithName("fs"))
|
filesystem, err := fs.NewFS(quota*gib, max, src, dst, metadata, klogr.New().WithName("fs"))
|
||||||
if err != nil {
|
if err != nil {
|
||||||
klog.Fatalf("init failed: %s", err)
|
klog.Fatalf("init failed: %s", err)
|
||||||
|
|
|
||||||
60
pkg/fs/discard.go
Normal file
60
pkg/fs/discard.go
Normal file
|
|
@ -0,0 +1,60 @@
|
||||||
|
// Copyright (C) 2022 Marius Schellenberger
|
||||||
|
|
||||||
|
package fs
|
||||||
|
|
||||||
|
import (
|
||||||
|
"io"
|
||||||
|
"sync"
|
||||||
|
)
|
||||||
|
|
||||||
|
var Discard io.Writer = discard{}
|
||||||
|
|
||||||
|
type discard struct{}
|
||||||
|
|
||||||
|
var _ io.ReaderFrom = discard{}
|
||||||
|
|
||||||
|
func (discard) Write(p []byte) (int, error) {
|
||||||
|
return len(p), nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func (discard) WriteString(s string) (int, error) {
|
||||||
|
return len(s), nil
|
||||||
|
}
|
||||||
|
|
||||||
|
const DefaultBufferSize = 8192
|
||||||
|
|
||||||
|
var blackHolePool = sync.Pool{
|
||||||
|
New: func() any {
|
||||||
|
b := make([]byte, DefaultBufferSize)
|
||||||
|
return &b
|
||||||
|
},
|
||||||
|
}
|
||||||
|
|
||||||
|
func SetBufferSize(size int) bool {
|
||||||
|
if size < 0 || size == DefaultBufferSize {
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
blackHolePool = sync.Pool{
|
||||||
|
New: func() any {
|
||||||
|
b := make([]byte, size)
|
||||||
|
return &b
|
||||||
|
},
|
||||||
|
}
|
||||||
|
return true
|
||||||
|
}
|
||||||
|
|
||||||
|
func (discard) ReadFrom(r io.Reader) (n int64, err error) {
|
||||||
|
bufp := blackHolePool.Get().(*[]byte)
|
||||||
|
readSize := 0
|
||||||
|
for {
|
||||||
|
readSize, err = r.Read(*bufp)
|
||||||
|
n += int64(readSize)
|
||||||
|
if err != nil {
|
||||||
|
blackHolePool.Put(bufp)
|
||||||
|
if err == io.EOF {
|
||||||
|
return n, nil
|
||||||
|
}
|
||||||
|
return
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
@ -41,31 +41,33 @@ func (p *preload) Read(data []byte) (n int, err error) {
|
||||||
_, err = p.f.Seek(int64(n), io.SeekCurrent)
|
_, err = p.f.Seek(int64(n), io.SeekCurrent)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
n, err = p.f.readWithoutCache(data)
|
n, err = p.f.readToCache(data)
|
||||||
p.written++
|
p.written++
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
func (f *File) Preload(ctx context.Context, unlock func()) {
|
func (f *File) Preload(ctx context.Context, unlock func()) {
|
||||||
log := f.log
|
log := f.log
|
||||||
|
defer unlock()
|
||||||
if f.offline {
|
if f.offline {
|
||||||
log.V(2).Error(errors.New("no preload in offline mode"), "error preloading file")
|
log.V(2).Error(errors.New("no preload in offline mode"), "error preloading file")
|
||||||
unlock()
|
return
|
||||||
|
}
|
||||||
|
if f.md.FullyCached() {
|
||||||
|
log.V(2).Info("skipped preload for fully cached file")
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
log.V(2).Info("preload started")
|
log.V(2).Info("preload started")
|
||||||
p := &preload{f: f, ctx: ctx}
|
p := &preload{f: f, ctx: ctx}
|
||||||
_, err := io.Copy(io.Discard, p)
|
_, err := io.Copy(Discard, p)
|
||||||
if err == context.Canceled {
|
if err == context.Canceled {
|
||||||
log.V(2).Info("preload canceled", "skipped", p.skipped, "written", p.written)
|
log.V(2).Info("preload canceled", "skipped", p.skipped, "written", p.written)
|
||||||
unlock()
|
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
if err != nil && err != io.EOF {
|
if err != nil && err != io.EOF {
|
||||||
log.Error(err, "error preloading file")
|
log.Error(err, "error preloading file")
|
||||||
}
|
}
|
||||||
log.V(2).Info("preload finished", "skipped", p.skipped, "written", p.written)
|
log.V(2).Info("preload finished", "skipped", p.skipped, "written", p.written)
|
||||||
unlock()
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func (f *File) Read(p []byte) (n int, err error) {
|
func (f *File) Read(p []byte) (n int, err error) {
|
||||||
|
|
@ -73,16 +75,23 @@ func (f *File) Read(p []byte) (n int, err error) {
|
||||||
if f.hasChunk(len(p)) {
|
if f.hasChunk(len(p)) {
|
||||||
n, err = f.md.ReadAt(p, f.offset)
|
n, err = f.md.ReadAt(p, f.offset)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
if !IsIOErr(err) {
|
||||||
log.Error(err, "error reading cache file")
|
log.Error(err, "error reading cache file")
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
n, err = f.readSource(p)
|
||||||
|
if err != nil {
|
||||||
|
log.Error(err, "error reading source file")
|
||||||
|
return
|
||||||
|
}
|
||||||
|
}
|
||||||
_, err = f.Seek(int64(n), io.SeekCurrent)
|
_, err = f.Seek(int64(n), io.SeekCurrent)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
log.Error(err, "error seeking source file")
|
log.Error(err, "error seeking source file")
|
||||||
}
|
}
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
return f.readWithoutCache(p)
|
return f.readToCache(p)
|
||||||
}
|
}
|
||||||
|
|
||||||
func (f *File) size() int64 {
|
func (f *File) size() int64 {
|
||||||
|
|
@ -93,12 +102,16 @@ func (f *File) hasChunk(n int) bool {
|
||||||
return f.md.HasChunk(f.offset, n)
|
return f.md.HasChunk(f.offset, n)
|
||||||
}
|
}
|
||||||
|
|
||||||
func (f *File) readWithoutCache(p []byte) (n int, err error) {
|
func (f *File) readSource(p []byte) (n int, err error) {
|
||||||
if f.offline {
|
if f.offline {
|
||||||
return 0, io.EOF
|
return 0, io.EOF
|
||||||
}
|
}
|
||||||
|
return f.f.Read(p)
|
||||||
|
}
|
||||||
|
|
||||||
|
func (f *File) readToCache(p []byte) (n int, err error) {
|
||||||
log := f.log
|
log := f.log
|
||||||
n, err = f.f.Read(p)
|
n, err = f.readSource(p)
|
||||||
if n > 0 {
|
if n > 0 {
|
||||||
if n, err := f.md.WriteAt(p[:n], f.offset); err != nil {
|
if n, err := f.md.WriteAt(p[:n], f.offset); err != nil {
|
||||||
log.Error(err, "error writing cache file")
|
log.Error(err, "error writing cache file")
|
||||||
|
|
|
||||||
|
|
@ -273,6 +273,31 @@ func (mh *MetadataHandler) close() {
|
||||||
close(mh.fs.done)
|
close(mh.fs.done)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
type Error string
|
||||||
|
|
||||||
|
func (e Error) Error() string {
|
||||||
|
return string(e)
|
||||||
|
}
|
||||||
|
|
||||||
|
const (
|
||||||
|
IOErr Error = "input/output error"
|
||||||
|
)
|
||||||
|
|
||||||
|
func MapError(err error) error {
|
||||||
|
if err == nil {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
s := err.Error()
|
||||||
|
if strings.HasSuffix(s, IOErr.Error()) {
|
||||||
|
return IOErr
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func IsIOErr(err error) bool {
|
||||||
|
return MapError(err) == IOErr
|
||||||
|
}
|
||||||
|
|
||||||
type Metadata struct {
|
type Metadata struct {
|
||||||
mu sync.RWMutex `json:"-"`
|
mu sync.RWMutex `json:"-"`
|
||||||
fs *FS `json:"-"`
|
fs *FS `json:"-"`
|
||||||
|
|
@ -293,6 +318,15 @@ func (md *Metadata) Delete() error {
|
||||||
return md.fs.RemoveDst(md.name)
|
return md.fs.RemoveDst(md.name)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (md *Metadata) FullyCached() bool {
|
||||||
|
md.mu.RLock()
|
||||||
|
defer md.mu.RUnlock()
|
||||||
|
if len(md.Chunks) == 0 {
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
return md.Size == md.Chunks[0][1]
|
||||||
|
}
|
||||||
|
|
||||||
func (md *Metadata) HasChunk(off int64, n int) bool {
|
func (md *Metadata) HasChunk(off int64, n int) bool {
|
||||||
md.mu.RLock()
|
md.mu.RLock()
|
||||||
defer md.mu.RUnlock()
|
defer md.mu.RUnlock()
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue