cachefs/pkg/fs/metadata.go
2022-04-06 01:08:19 +02:00

419 lines
7.9 KiB
Go

// Copyright (C) 2022 Marius Schellenberger
package fs
import (
"cachefs/pkg/chunk"
"cachefs/pkg/provider"
"context"
"encoding/json"
"errors"
"io"
"io/fs"
stdfs "io/fs"
"math"
"os"
"path/filepath"
"strings"
"sync"
"sync/atomic"
"time"
"github.com/go-logr/logr"
"golang.org/x/exp/maps"
)
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 string, dst provider.FS) error {
if !filepath.IsAbs(file) {
return errors.New("metadata path is not absolute")
}
_, err := dst.Fstat(0)
if err == provider.ErrFstat {
return err
}
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)
fs.WalkDir(provider.WalkFS(dst), ".", func(path string, d stdfs.DirEntry, err error) error {
if err != nil {
return nil
}
if d.IsDir() {
return nil
}
f, err := dst.Open(path)
if err != nil {
return nil
}
defer f.Close()
fi, err := f.Stat()
if err != nil {
return nil
}
size := fi.Size()
blocks, err := dst.Fstat(f.Fd())
if err != nil {
return nil
}
if size <= blocks*512 {
md[strings.TrimPrefix(path, dst.Root())] = &Metadata{
Size: size,
Chunks: chunk.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.dst.Stat(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() {
fs.WalkDir(provider.WalkFS(mh.fs.dst), ".", func(path string, d stdfs.DirEntry, err error) error {
if err == nil && d.IsDir() {
if path == mh.dst {
return nil
}
empty, err := mh.dirEmpty(path)
if err != nil || !empty {
return nil
}
log := mh.log.WithValues("dir", path)
log.V(2).Info("removing empty dir")
err = mh.fs.dst.Remove(path)
if err != nil {
log.Error(err, "failed to remove empty dir")
}
}
return nil
})
}
func (mh *MetadataHandler) dirEmpty(name string) (bool, error) {
f, err := mh.fs.dst.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
}
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 {
if md.f != nil {
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 provider.File `json:"-"`
name string `json:"-"`
Size int64 `json:"s"`
Atime int64 `json:"a"`
Chunks chunk.Chunks `json:"c"`
}
func (md *Metadata) Close() error {
md.mu.Lock()
defer md.mu.Unlock()
err := md.f.Close()
md.f = nil
return err
}
func (md *Metadata) Delete() error {
md.mu.Lock()
defer md.mu.Unlock()
md.Chunks = chunk.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()
}