// Copyright (C) 2022 Marius Schellenberger package fs import ( "context" "encoding/json" "errors" "io" stdfs "io/fs" "math" "net/http" "os" "path" "path/filepath" "strings" "sync" "time" "github.com/go-logr/logr" "golang.org/x/exp/maps" ) type FS struct { mu sync.RWMutex log logr.Logger NoCache http.Handler src string dst string statCancel func() flushCancel func() mdf *os.File jenc *json.Encoder dc *DirCache sc *StatCache q *Quota md MetadataMap pm map[string]func() } func NewFS(max int64, src, dst, metadata string, log logr.Logger) (fs *FS, err error) { if !filepath.IsAbs(src) { return nil, errors.New("src path is not absolute") } if !filepath.IsAbs(dst) { return nil, errors.New("dst path is not absolute") } if !filepath.IsAbs(metadata) { return nil, errors.New("metadata path is not absolute") } if src == dst { return nil, errors.New("src and dst path can not be equal") } mdf, err := os.OpenFile(metadata, os.O_RDWR|os.O_CREATE, 0o640) if err != nil { return } statCtx, statCancel := context.WithCancel(context.Background()) flushCtx, flushCancel := context.WithCancel(context.Background()) fs = &FS{ log: log, NoCache: http.FileServer(http.Dir(src)), src: src, dst: dst, mdf: mdf, statCancel: statCancel, flushCancel: flushCancel, dc: NewDirCache(), sc: NewStatCache(statCtx), md: make(MetadataMap), pm: make(map[string]func()), jenc: json.NewEncoder(mdf), } fs.q, err = NewQuota(max, fs, fs.log.WithName("quota")) if err != nil { return } err = json.NewDecoder(mdf).Decode(&fs.md) if err == io.EOF { err = nil } fs.initMetadata() fs.q.Init() go fs.flusher(flushCtx) return fs, err } func (fs *FS) initMetadata() { log := fs.log maps.DeleteFunc(fs.md, func(k string, v *Metadata) bool { _, err := fs.statDst(k) if errors.Is(err, os.ErrNotExist) { log.V(2).Info("removing not existing file from metadata", "file", k) return true } if len(v.Chunks) == 0 { log.V(2).Info("removing empty file", "file", k) err := fs.RemoveDst(k) if err != nil { log.Error(err, "error removing empty file", "file", k) return false } return true } v.fs = fs v.name = k return false }) } func (fs *FS) Stat(name string) (fi stdfs.FileInfo, err error) { sp, dp := fs.paths(name) fi, err = fs.sc.Get(name) if err == nil { return } fi, err = os.Stat(sp) if err == nil { fs.sc.Set(name, fi) return } fi, err = os.Stat(dp) return } func (fs *FS) RemoveDst(name string) error { _, dp := fs.paths(name) return os.Remove(dp) } func (fs *FS) StatDst(name string) (fi stdfs.FileInfo, err error) { fi, err = fs.sc.Get(name) if err == nil { return } return fs.statDst(name) } func (fs *FS) statDst(name string) (stdfs.FileInfo, error) { _, dp := fs.paths(name) return os.Stat(dp) } func (fs *FS) OpenDst(name string) (f http.File, err error) { log := fs.log.WithValues("file", name) _, dp := fs.paths(name) fi, err := fs.StatDst(name) if err != nil { if !skipLog(name) { log.Error(err, "error stat cache file") } return } file, err := os.Open(dp) if err != nil { if !skipLog(name) { log.Error(err, "error opening cache file") } return } if fi.IsDir() { f = &Dir{log: log.WithName("dir"), f: file, dc: fs.dc} return } md := fs.metadata(name, fi.Size()) f = &File{ log: log, f: file, md: md, offline: true, } return } func (fs *FS) CancelPreload(name string) { if cancel, ok := fs.pm[name]; ok { cancel() } } func (fs *FS) Preload(name string) { file, err := fs.Open(name) if err != nil { return } f, ok := file.(*File) if !ok { return } if f.md.Preload() { fs.log.WithValues("file", name).Error(errors.New("preload is running"), "error starting preload") return } ctx, cancel := context.WithCancel(context.Background()) fs.pm[name] = cancel unlock := func() { delete(fs.pm, name) f.md.UnlockPreload() } go f.Preload(ctx, unlock) } func skipLog(name string) bool { return strings.HasSuffix(name, "index.html") || strings.HasSuffix(name, "favicon.ico") } func (fs *FS) paths(name string) (sp string, dp string) { p := filepath.FromSlash(path.Clean("/" + name)) sp = filepath.Join(fs.src, p) dp = filepath.Join(fs.dst, p) return } func (fs *FS) open(name, sp, dp string) (f *os.File, fi os.FileInfo, offline bool, err error) { f, err = os.Open(sp) if err == nil { fi, err = f.Stat() if err == nil { return } f.Close() } f, err = os.Open(dp) if err != nil { return } fi, err = f.Stat() if err != nil { f.Close() } offline = true return } func (fs *FS) Open(name string) (f http.File, err error) { log := fs.log.WithValues("file", name) sp, dp := fs.paths(name) file, fi, offline, err := fs.open(name, sp, dp) if err != nil { if !skipLog(name) { log.Error(err, "error opening source file") } return } if fi.IsDir() { if offline { log.V(2).Info("dir offline mode", "path", dp) } f = &Dir{log: log.WithName("dir"), f: file, dc: fs.dc} return } if offline { log.V(2).Info("file offline mode", "path", dp) } md := fs.metadata(name, fi.Size()) f = &File{ log: log, f: file, md: md, offline: offline, } return } func (fs *FS) openCacheFile(name string, size int64) (df *os.File, err error) { _, dp := fs.paths(name) log := fs.log if i := strings.LastIndex(dp, "/"); i > 0 { dir := dp[:i] _, err = os.Stat(dir) if errors.Is(err, os.ErrNotExist) { err = os.MkdirAll(dir, 0o755) if err != nil { log.Error(err, "error creating cache dir", "dir", dp) return } } } _, err = os.Stat(dp) truncate := errors.Is(err, os.ErrNotExist) df, err = os.OpenFile(dp, os.O_RDWR|os.O_CREATE, 0o644) if err != nil { log.Error(err, "error opening cache file") return } if truncate { err = df.Truncate(size) if err != nil { log.Error(err, "error truncating cache file") return } err = df.Sync() if err != nil { log.Error(err, "error syncing cache file") return } } return } func (fs *FS) metadata(name string, size int64) (md *Metadata) { fs.mu.Lock() defer fs.mu.Unlock() if md, ok := fs.md[name]; ok { if size > 0 && md.Size != size { log := fs.log.WithValues("file", name) log.V(2).Info("source file changed: deleting cache file", "src", size, "dst", md.Size) err := md.Delete() if err != nil { log.Error(err, "error deleting cache file") } } return md } md = &Metadata{Size: size, fs: fs, name: name} fs.md[name] = md return } func (fs *FS) CacheStatus(name string) int { fs.mu.RLock() defer fs.mu.RUnlock() if md, ok := fs.md[name]; ok { size := md.Chunks.Size() if size == 0 { return -1 } return int(math.Round(float64(size) / float64(md.Size) * 100)) } return -1 } func (fs *FS) flusher(ctx context.Context) { for { select { case <-ctx.Done(): return case <-time.After(time.Minute * 5): } fs.cleanupEmptyDirs() fs.flushMetadata() } } func (fs *FS) cleanupEmptyDirs() { filepath.WalkDir(fs.dst, func(path string, d stdfs.DirEntry, err error) error { if err == nil && d.IsDir() { if path == fs.dst { return nil } empty, err := dirEmpty(path) if err != nil || !empty { return nil } log := fs.log.WithValues("dir", path) log.V(2).Info("removing empty dir") err = os.Remove(path) if err != nil { log.Error(err, "failed to remove empty dir") } } return nil }) } func (fs *FS) flushMetadata() { log := fs.log fs.mu.Lock() defer fs.mu.Unlock() log.V(2).Info("flushing metadata to disk") md := maps.Clone(fs.md) maps.DeleteFunc(md, func(_ string, v *Metadata) bool { return len(v.Chunks) == 0 }) _, err := fs.mdf.Seek(0, io.SeekStart) if err != nil { log.Error(err, "failure seeking metadata file") return } err = fs.mdf.Truncate(0) if err != nil { log.Error(err, "failure truncating metadata file") return } err = fs.jenc.Encode(md) if err != nil { log.Error(err, "failure flushing metadata file") } } func (fs *FS) Close() { fs.flushCancel() fs.statCancel() for _, cancel := range fs.pm { cancel() } for _, md := range fs.md { md.f.Sync() md.f.Close() } fs.flushMetadata() } func dirEmpty(name string) (bool, error) { f, err := os.Open(name) if err != nil { return false, err } _, err = f.ReadDir(1) if err == io.EOF { f.Close() return true, nil } f.Close() return false, err }