diff --git a/cmd/cachefs/main.go b/cmd/cachefs/main.go index a32ff9d..aad2c37 100644 --- a/cmd/cachefs/main.go +++ b/cmd/cachefs/main.go @@ -30,6 +30,7 @@ var ( listenCache string quota int64 max int + bs int log logr.Logger ) @@ -45,12 +46,17 @@ func main() { flag.StringVar(&listenWebdav, "webdav", "", "listen addr:port for webdav") flag.StringVar(&listenCache, "cache", "", "listen addr:port for cache only") 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.Parse() log = klogr.New().WithName("main") 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")) if err != nil { klog.Fatalf("init failed: %s", err) diff --git a/pkg/fs/discard.go b/pkg/fs/discard.go new file mode 100644 index 0000000..fdf049d --- /dev/null +++ b/pkg/fs/discard.go @@ -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 + } + } +} diff --git a/pkg/fs/file.go b/pkg/fs/file.go index 17356b4..33400d9 100644 --- a/pkg/fs/file.go +++ b/pkg/fs/file.go @@ -41,31 +41,33 @@ func (p *preload) Read(data []byte) (n int, err error) { _, err = p.f.Seek(int64(n), io.SeekCurrent) return } - n, err = p.f.readWithoutCache(data) + n, err = p.f.readToCache(data) p.written++ return } func (f *File) Preload(ctx context.Context, unlock func()) { log := f.log + defer unlock() if f.offline { 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 } log.V(2).Info("preload started") p := &preload{f: f, ctx: ctx} - _, err := io.Copy(io.Discard, p) + _, err := io.Copy(Discard, p) if err == context.Canceled { log.V(2).Info("preload canceled", "skipped", p.skipped, "written", p.written) - unlock() return } if err != nil && err != io.EOF { log.Error(err, "error preloading file") } log.V(2).Info("preload finished", "skipped", p.skipped, "written", p.written) - unlock() } func (f *File) Read(p []byte) (n int, err error) { @@ -73,8 +75,15 @@ func (f *File) Read(p []byte) (n int, err error) { if f.hasChunk(len(p)) { n, err = f.md.ReadAt(p, f.offset) if err != nil { - log.Error(err, "error reading cache file") - return + if !IsIOErr(err) { + log.Error(err, "error reading cache file") + return + } + n, err = f.readSource(p) + if err != nil { + log.Error(err, "error reading source file") + return + } } _, err = f.Seek(int64(n), io.SeekCurrent) if err != nil { @@ -82,7 +91,7 @@ func (f *File) Read(p []byte) (n int, err error) { } return } - return f.readWithoutCache(p) + return f.readToCache(p) } func (f *File) size() int64 { @@ -93,12 +102,16 @@ func (f *File) hasChunk(n int) bool { 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 { return 0, io.EOF } + return f.f.Read(p) +} + +func (f *File) readToCache(p []byte) (n int, err error) { log := f.log - n, err = f.f.Read(p) + n, err = f.readSource(p) if n > 0 { if n, err := f.md.WriteAt(p[:n], f.offset); err != nil { log.Error(err, "error writing cache file") diff --git a/pkg/fs/metadata.go b/pkg/fs/metadata.go index a7ee420..8e4b57a 100644 --- a/pkg/fs/metadata.go +++ b/pkg/fs/metadata.go @@ -273,6 +273,31 @@ func (mh *MetadataHandler) close() { 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 { mu sync.RWMutex `json:"-"` fs *FS `json:"-"` @@ -293,6 +318,15 @@ func (md *Metadata) Delete() error { 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 { md.mu.RLock() defer md.mu.RUnlock()