// Copyright (C) 2022 Marius Schellenberger package fs import ( "context" "encoding/json" "errors" "io" stdfs "io/fs" "math" "os" "path/filepath" "strings" "sync" "sync/atomic" "time" "github.com/go-logr/logr" "golang.org/x/exp/maps" "golang.org/x/sys/unix" ) func now() int64 { return time.Now().UTC().Unix() } type MetadataHandler struct { mu sync.RWMutex log logr.Logger md map[string]*Metadata fs *FS dst string f *os.File enc *json.Encoder } func MetadataGenerator(file, dst string) error { if !filepath.IsAbs(file) { return errors.New("metadata path is not absolute") } if !filepath.IsAbs(dst) { return errors.New("dst path is not absolute") } f, err := os.OpenFile(file, os.O_RDWR|os.O_CREATE, 0o640) if err != nil { return err } defer f.Close() err = f.Truncate(0) if err != nil { return err } enc := json.NewEncoder(f) md := make(map[string]*Metadata) filepath.WalkDir(dst, func(path string, d stdfs.DirEntry, err error) error { if err != nil { return nil } if d.IsDir() { return nil } var stat unix.Stat_t f, err := os.Open(path) if err != nil { return nil } defer f.Close() fi, err := f.Stat() if err != nil { return nil } size := fi.Size() err = unix.Fstat(int(f.Fd()), &stat) if err != nil { return nil } if size <= stat.Blocks*512 { md[strings.TrimPrefix(path, dst)] = &Metadata{ Size: size, Chunks: Chunks{{0, size}}, } } return nil }) return enc.Encode(md) } func NewMetadataHandler(ctx context.Context, fs *FS, file, dst string, log logr.Logger) (mh *MetadataHandler, err error) { if !filepath.IsAbs(file) { return nil, errors.New("metadata path is not absolute") } mh = &MetadataHandler{ log: log, md: make(map[string]*Metadata), fs: fs, dst: dst, } mh.f, err = os.OpenFile(file, os.O_RDWR|os.O_CREATE, 0o640) if err != nil { return } err = json.NewDecoder(mh.f).Decode(&mh.md) if err != nil && err == io.EOF { return } mh.enc = json.NewEncoder(mh.f) err = nil mh.init() go mh.flusher(ctx) return } func (mh *MetadataHandler) DeleteOldest() (s int64) { mh.mu.Lock() defer mh.mu.Unlock() atime := int64(math.MaxInt64) var m *Metadata for _, v := range mh.md { if v.Atime != 0 && v.f != nil && v.Atime < atime { atime = v.Atime m = v } } if m != nil { log := mh.log.WithValues("file", m.name) log.Info("deleting oldest file") cs := m.Size if !m.FullyCached() { cs = m.ChunkSize() } err := m.Delete() if err != nil { log.Error(err, "error deleting oldest file") return } s = cs } return } func (mh *MetadataHandler) Size() (s int64) { mh.mu.RLock() defer mh.mu.RUnlock() for _, md := range mh.md { s += md.ChunkSize() } return } func (mh *MetadataHandler) CacheStatus(name string) int { mh.mu.RLock() md, ok := mh.md[name] mh.mu.RUnlock() if !ok { return -1 } if md.FullyCached() { return 100 } size := md.ChunkSize() if size == 0 { return -1 } return int(math.Round(float64(size) / float64(md.Size) * 100)) } func (mh *MetadataHandler) Metadata(name string, size int64) (md *Metadata) { mh.mu.Lock() defer mh.mu.Unlock() if md, ok := mh.md[name]; ok { if size > 0 && md.Size != size { log := mh.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: mh.fs, name: name} mh.md[name] = md return } func (mh *MetadataHandler) init() { log := mh.log maps.DeleteFunc(mh.md, func(k string, v *Metadata) bool { _, err := mh.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 := mh.fs.RemoveDst(k) if err != nil { log.Error(err, "error removing empty file", "file", k) return false } return true } v.fs = mh.fs v.name = k return false }) } func (mh *MetadataHandler) flusher(ctx context.Context) { for { select { case <-ctx.Done(): mh.close() return case <-time.After(time.Minute * 5): } mh.cleanupEmptyDirs() mh.flush() } } func (mh *MetadataHandler) cleanupEmptyDirs() { filepath.WalkDir(mh.dst, func(path string, d stdfs.DirEntry, err error) error { if err == nil && d.IsDir() { if path == mh.dst { return nil } empty, err := dirEmpty(path) if err != nil || !empty { return nil } log := mh.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 (mh *MetadataHandler) flush() { log := mh.log mh.mu.Lock() defer mh.mu.Unlock() log.V(2).Info("flushing metadata to disk") md := maps.Clone(mh.md) maps.DeleteFunc(mh.md, func(_ string, v *Metadata) bool { return len(v.Chunks) == 0 }) _, err := mh.f.Seek(0, io.SeekStart) if err != nil { log.Error(err, "failure seeking metadata file") return } err = mh.f.Truncate(0) if err != nil { log.Error(err, "failure truncating metadata file") return } err = mh.enc.Encode(md) if err != nil { log.Error(err, "failure flushing metadata file") } return } func (mh *MetadataHandler) close() { mh.mu.Lock() for _, md := range mh.md { md.f.Sync() md.f.Close() } mh.mu.Unlock() mh.flush() mh.f.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:"-"` f *os.File `json:"-"` name string `json:"-"` Size int64 `json:"s"` Atime int64 `json:"a"` Chunks Chunks `json:"c"` } func (md *Metadata) Delete() error { md.mu.Lock() defer md.mu.Unlock() md.Chunks = Chunks{} md.f.Close() md.f = nil atomic.StoreInt64(&md.Atime, now()) return md.fs.RemoveDst(md.name) } func (md *Metadata) FullyCached() bool { md.mu.RLock() defer md.mu.RUnlock() return len(md.Chunks) == 1 && md.Size == md.Chunks[0][1] } func (md *Metadata) HasChunk(off int64, n int) bool { md.mu.RLock() defer md.mu.RUnlock() return md.Chunks.Exists(off, n, md.Size) } func (md *Metadata) ChunkSize() int64 { md.mu.RLock() defer md.mu.RUnlock() return md.Chunks.Size() } func (md *Metadata) AddChunk(off int64, n int) { md.mu.Lock() defer md.mu.Unlock() md.Chunks.Add(off, n) } func (md *Metadata) ReadAt(p []byte, pos int64) (int, error) { if md.f == nil { err := md.openCacheFile() if err != nil { return 0, err } } atomic.StoreInt64(&md.Atime, now()) return md.f.ReadAt(p, pos) } func (md *Metadata) WriteAt(data []byte, pos int64) (int, error) { if md.f == nil { err := md.openCacheFile() if err != nil { return 0, err } } atomic.StoreInt64(&md.Atime, now()) md.fs.q.Add(len(data)) return md.f.WriteAt(data, pos) } func (md *Metadata) openCacheFile() error { md.mu.Lock() defer md.mu.Unlock() if md.f != nil { return nil } f, err := md.fs.openCacheFile(md.name, md.Size) md.f = f return err } func (md *Metadata) Stat() (os.FileInfo, error) { if md.f == nil { err := md.openCacheFile() if err != nil { return nil, err } } return md.f.Stat() }