// 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" ) type FS struct { mu sync.RWMutex log logr.Logger NoCache http.Handler src string dst string mdf *os.File jenc *json.Encoder mm map[string]*Metadata prem map[string]func() cancel func() } func NewFS(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 } fs = &FS{ log: log, NoCache: http.FileServer(http.Dir(src)), src: src, dst: dst, mdf: mdf, mm: make(map[string]*Metadata), prem: make(map[string]func()), jenc: json.NewEncoder(mdf), } err = json.NewDecoder(mdf).Decode(&fs.mm) if err == io.EOF { err = nil } fs.initMetadata() ctx, cancel := context.WithCancel(context.Background()) fs.cancel = cancel go fs.flusher(ctx) return fs, err } func (fs *FS) initMetadata() { log := fs.log for k, v := range fs.mm { _, err := fs.statDst(k) if errors.Is(err, os.ErrNotExist) { log.V(2).Info("removing not existing file from metadata", "file", k) delete(fs.mm, k) continue } 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) continue } delete(fs.mm, k) continue } v.fs = fs v.name = k } } func (fs *FS) Stat(name string) (fi stdfs.FileInfo, err error) { fi, _, err = fs.StatWithOffline(name) return } func (fs *FS) StatWithOffline(name string) (fi stdfs.FileInfo, offline bool, err error) { sp, dp := fs.paths(name) fi, err = os.Stat(sp) if errors.Is(err, os.ErrNotExist) { fi, err = os.Stat(dp) if err != nil { return } offline = true } 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) { _, dp := fs.paths(name) return os.Stat(dp) } func (fs *FS) CancelPreload(name string) { if cancel, ok := fs.prem[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.prem[name] = cancel unlock := func() { delete(fs.prem, 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 string) (f http.File, err error) { log := fs.log.WithValues("file", name) sp, dp := fs.paths(name) offline := false sf, err := os.Open(sp) if err != nil { if errors.Is(err, os.ErrNotExist) { sf, err = os.Open(dp) if err != nil { if !skipLog(name) { log.Error(err, "error opening source file") } return } offline = true log.V(2).Info("file offline mode", "path", dp) } else { log.Error(err, "error opening source file") return } } sfi, err := sf.Stat() if err != nil { sf.Close() log.Error(err, "error stat source file") return } if sfi.IsDir() { if empty, e := dirEmpty(sf.Name()); empty && e == nil { sf.Close() sf, err = os.Open(dp) offline = true log.V(2).Info("dir offline mode", "path", dp) } return &Dir{f: sf}, err } md := fs.metadata(name, sfi.Size()) return &File{ log: log, f: sf, md: md, offline: offline, }, nil } 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.mm[name]; ok { return md } md = &Metadata{Size: size, fs: fs, name: name} fs.mm[name] = md return } func (fs *FS) CacheStatus(name string) int { fs.mu.RLock() defer fs.mu.RUnlock() if md, ok := fs.mm[name]; ok { return int(math.Round(float64(md.Chunks.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.flushMetadata() } } func (fs *FS) flushMetadata() { log := fs.log fs.mu.Lock() defer fs.mu.Unlock() mm := make(map[string]*Metadata) for k, v := range fs.mm { if v.f != nil || len(v.Chunks) > 0 { mm[k] = v } } _, 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(mm) if err != nil { log.Error(err, "failure flushing metadata file") } } func (fs *FS) Close() { fs.cancel() for _, md := range fs.mm { 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 } defer f.Close() _, err = f.Readdir(1) if err == io.EOF { return true, nil } return false, err }