diff --git a/cmd/cachefs/main.go b/cmd/cachefs/main.go index 1f28261..5ae43ea 100644 --- a/cmd/cachefs/main.go +++ b/cmd/cachefs/main.go @@ -49,7 +49,7 @@ func main() { flag.IntVar(&max, "max", -1, "max parallel preloads") flag.IntVar(&bs, "bs", -1, "tune preload buffer size in bytes (default: 8192)") flag.Int64Var("a, "quota", 1, "max disk usage quota for the dst cache in GiB") - flag.BoolVar(&block, "block", false, "block until sftp is connected") + flag.BoolVar(&block, "block", true, "block until sftp is connected") flag.Parse() log = klogr.New().WithName("main") @@ -82,6 +82,7 @@ func main() { ), } go func() { + log.Info("starting main server", "addr", listenHttp) err := s.ListenAndServe() if err != nil && err != http.ErrServerClosed { klog.Fatalf("init failed: %s", err) @@ -92,6 +93,7 @@ func main() { cache *http.Server ) if listenWebdav != "" { + log.Info("starting webdav server", "addr", listenWebdav) dav = &http.Server{ Addr: listenWebdav, Handler: srv.NewWebDavServer( @@ -107,6 +109,7 @@ func main() { }() } if listenCache != "" { + log.Info("starting cache server", "addr", listenCache) cache = &http.Server{ Addr: listenCache, Handler: srv.NewCacheServer( diff --git a/pkg/fs/metadata.go b/pkg/fs/metadata.go index c621942..a4478dd 100644 --- a/pkg/fs/metadata.go +++ b/pkg/fs/metadata.go @@ -5,6 +5,7 @@ package fs import ( "cachefs/pkg/chunk" "cachefs/pkg/provider" + "cachefs/pkg/provider/sftp" "context" "encoding/json" "errors" @@ -80,19 +81,20 @@ func MetadataGenerator(file string, dst provider.FS, fstat, encrypted bool) erro return nil } size := fi.Size() + path = "/" + strings.TrimPrefix(path, dst.Root()) if fstat { blocks, err := dst.Fstat(f.Fd()) if err != nil { return nil } if size <= blocks*512 { - md[strings.TrimPrefix(path, dst.Root())] = &Metadata{ + md[path] = &Metadata{ Size: size, Chunks: chunk.Chunks{{0, size}}, } } } else { - md[strings.TrimPrefix(path, dst.Root())] = &Metadata{ + md[path] = &Metadata{ Size: size, Chunks: chunk.Chunks{{0, size}}, } @@ -332,7 +334,7 @@ func MapError(err error) error { } func IsIOErr(err error) bool { - return MapError(err) == IOErr + return MapError(err) == IOErr || errors.Is(err, sftp.ErrNotConnected) } type Metadata struct { diff --git a/pkg/fs/preload.go b/pkg/fs/preload.go index 664f854..1f02dd9 100644 --- a/pkg/fs/preload.go +++ b/pkg/fs/preload.go @@ -53,8 +53,12 @@ func NewPreloadHandler(ctx context.Context, fs *FS, max int, log logr.Logger) (p func (ph *PreloadHandler) CancelPreload(name string) { ph.mu.Lock() defer ph.mu.Unlock() - if p, ok := ph.pm[name]; ok && p.Running { - p.cancel() + p, ok := ph.pm[name] + if ok { + if p.Running { + p.cancel() + } + delete(ph.pm, name) } } diff --git a/pkg/provider/crypto/file.go b/pkg/provider/crypto/file.go index 027132f..1927465 100644 --- a/pkg/provider/crypto/file.go +++ b/pkg/provider/crypto/file.go @@ -5,12 +5,15 @@ import ( "io" "io/fs" "os" + "sync" ) var _ provider.File = (*file)(nil) type file struct { provider.File + rmu sync.Mutex + wmu sync.Mutex fs *FS r *reader w *writer @@ -140,10 +143,12 @@ func (f *file) Seek(offset int64, whence int) (n int64, err error) { f.r.nonce = f.nonce f.w.nonce = f.nonce f.r.off = roff + f.r.cn = cn return } func (f *file) Read(p []byte) (n int, err error) { + f.r.cn = -1 n, err = f.r.Read(p) return } @@ -154,6 +159,8 @@ func (f *file) Write(p []byte) (n int, err error) { } func (f *file) ReadAt(p []byte, pos int64) (n int, err error) { + f.rmu.Lock() + defer f.rmu.Unlock() cn, _, roff := align(pos) var last bool if f.r.cn != cn { @@ -165,7 +172,6 @@ func (f *file) ReadAt(p []byte, pos int64) (n int, err error) { if err != nil { return } - f.r.cn = cn } n = copy(p, f.r.unread[roff:]) if last && len(f.r.unread) == n+int(roff) { @@ -175,6 +181,8 @@ func (f *file) ReadAt(p []byte, pos int64) (n int, err error) { } func (f *file) WriteAt(data []byte, pos int64) (n int, err error) { + f.wmu.Lock() + defer f.wmu.Unlock() cn, off, woff := align(pos) c, ok := f.w.cm[cn] if !ok {