// 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" "golang.org/x/exp/slices" ) type FS struct { mu sync.RWMutex log logr.Logger NoCache http.Handler src string dst string statCancel func() flushCancel func() preloadCancel func() mdf *os.File jenc *json.Encoder dc *DirCache sc *StatCache q *Quota md MetadataMap pm map[string]*Preload prec chan *Preload pref chan string preq int maxq int } type Preload struct { Name string cancel func() Status int Running bool } type Preloads []*Preload func (Preloads) Less(i, j *Preload) bool { return i.Name < j.Name } func NewFS(max int64, maxq int, 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()) preloadCtx, preloadCancel := context.WithCancel(context.Background()) if maxq < -1 { maxq = -1 } fs = &FS{ log: log, NoCache: http.FileServer(http.Dir(src)), src: src, dst: dst, mdf: mdf, statCancel: statCancel, flushCancel: flushCancel, preloadCancel: preloadCancel, dc: NewDirCache(), sc: NewStatCache(statCtx), md: make(MetadataMap), jenc: json.NewEncoder(mdf), pm: make(map[string]*Preload), prec: make(chan *Preload, 1), pref: make(chan string, 1), maxq: maxq, } 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) go fs.preload(preloadCtx) go fs.preloadFinish(preloadCtx) 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) { fs.mu.Lock() defer fs.mu.Unlock() if p, ok := fs.pm[name]; ok && p.Running { p.cancel() } } func (fs *FS) Preloads() (ps Preloads) { fs.mu.RLock() defer fs.mu.RUnlock() for _, p := range fs.pm { ps = append(ps, &Preload{ Name: p.Name, Running: p.Running, Status: fs.CacheStatus(p.Name), }) } slices.SortFunc(ps, ps.Less) return } func (fs *FS) Preload(name string) { fs.mu.Lock() defer fs.mu.Unlock() p, ok := fs.pm[name] if !ok { p = &Preload{Name: name} fs.pm[name] = p if fs.maxq == -1 || fs.preq < fs.maxq { fs.prec <- p } return } if p.Running { fs.log.WithValues("file", name).Error(errors.New("preload is running"), "error starting preload") } return } func (fs *FS) nextPreload() { if fs.preq == fs.maxq { return } for _, p := range fs.pm { if p.Running { continue } fs.prec <- p return } } func (fs *FS) preload(ctx context.Context) { var p *Preload for { select { case <-ctx.Done(): return case p = <-fs.prec: } name := p.Name file, err := fs.Open(name) if err != nil { fs.log.Error(err, "error staring next preload", "file", name) return } f, ok := file.(*File) if !ok { continue } ctx, cancel := context.WithCancel(context.Background()) fs.mu.Lock() fs.preq++ p.Running = true p.cancel = cancel fs.mu.Unlock() go f.Preload(ctx, func() { fs.pref <- name }) } } func (fs *FS) preloadFinish(ctx context.Context) { for { select { case <-ctx.Done(): return case name := <-fs.pref: fs.mu.Lock() delete(fs.pm, name) fs.preq-- fs.nextPreload() fs.mu.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.mu.Lock() fs.preloadCancel() fs.flushCancel() fs.statCancel() ctx, cancel := context.WithCancel(context.Background()) defer cancel() go func() { for { select { case <-ctx.Done(): return case <-fs.prec: case <-fs.pref: } } }() for _, p := range fs.pm { if p.Running { p.cancel() } } for _, md := range fs.md { md.f.Sync() md.f.Close() } fs.mu.Unlock() 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 }