539 lines
10 KiB
Go
539 lines
10 KiB
Go
// Copyright (C) 2022 Marius Schellenberger
|
|
|
|
package fs
|
|
|
|
import (
|
|
"cachefs/pkg/chunk"
|
|
"cachefs/pkg/provider"
|
|
"cachefs/pkg/provider/sftp"
|
|
"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
|
|
done chan struct{}
|
|
log logr.Logger
|
|
md map[string]*Metadata
|
|
fs *FS
|
|
dst string
|
|
f *os.File
|
|
enc *json.Encoder
|
|
}
|
|
|
|
func MetadataGenerator(file string, dst provider.FS, fstat, encrypted bool) error {
|
|
if !filepath.IsAbs(file) {
|
|
return errors.New("metadata path is not absolute")
|
|
}
|
|
if fstat {
|
|
_, 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)
|
|
base := "."
|
|
if encrypted {
|
|
base = ""
|
|
}
|
|
fs.WalkDir(provider.WalkFS(dst), base, 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()
|
|
path = "/" + strings.TrimPrefix(path, dst.Root())
|
|
if fstat {
|
|
blocks, err := dst.Fstat(f.Fd())
|
|
if err != nil {
|
|
return nil
|
|
}
|
|
if size <= blocks*512 {
|
|
md[path] = &Metadata{
|
|
Size: size,
|
|
Chunks: chunk.Chunks{{0, size}},
|
|
}
|
|
}
|
|
} else {
|
|
md[path] = &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,
|
|
done: make(chan struct{}),
|
|
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() (free 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.Atime < atime {
|
|
atime = v.Atime
|
|
m = v
|
|
}
|
|
}
|
|
if m != nil {
|
|
log := mh.log.WithValues("file", m.name)
|
|
log.Info("deleting oldest file")
|
|
s := m.Size
|
|
if !m.FullyCached() {
|
|
s = m.ChunkSize()
|
|
}
|
|
err := m.Delete()
|
|
if err != nil {
|
|
log.Error(err, "error deleting oldest file")
|
|
return
|
|
}
|
|
free = s
|
|
}
|
|
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
|
|
}
|
|
ctx, cancel := context.WithCancel(context.Background())
|
|
md = &Metadata{
|
|
Size: size,
|
|
fs: mh.fs,
|
|
name: name,
|
|
ctx: ctx,
|
|
cancel: cancel,
|
|
}
|
|
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.ctx, v.cancel = context.WithCancel(context.Background())
|
|
v.fs = mh.fs
|
|
v.name = k
|
|
return false
|
|
})
|
|
}
|
|
|
|
func (mh *MetadataHandler) flusher(ctx context.Context) {
|
|
t := time.NewTicker(time.Minute * 5)
|
|
for {
|
|
select {
|
|
case <-ctx.Done():
|
|
t.Stop()
|
|
mh.close()
|
|
close(mh.done)
|
|
return
|
|
case <-t.C:
|
|
}
|
|
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 {
|
|
md.Close()
|
|
//if md.f != nil {
|
|
// close(md.wc)
|
|
// md.f.Sync()
|
|
// md.f.Close()
|
|
//}
|
|
}
|
|
mh.mu.Unlock()
|
|
mh.flush()
|
|
mh.f.Close()
|
|
}
|
|
|
|
func (mh *MetadataHandler) Done() <-chan struct{} {
|
|
return mh.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 || errors.Is(err, sftp.ErrNotConnected)
|
|
}
|
|
|
|
type Metadata struct {
|
|
mu sync.Mutex `json:"-"`
|
|
cmu sync.RWMutex `json:"-"`
|
|
fs *FS `json:"-"`
|
|
f provider.File `json:"-"`
|
|
wc chan writeAt `json:"-"`
|
|
done chan struct{} `json:"-"`
|
|
ctx context.Context `json:"-"`
|
|
cancel func() `json:"-"`
|
|
err atomic.Pointer[mdErr] `json:"-"`
|
|
name string `json:"-"`
|
|
Size int64 `json:"s"`
|
|
Atime int64 `json:"a"`
|
|
Chunks chunk.Chunks `json:"c"`
|
|
}
|
|
|
|
func (md *Metadata) Close() (err error) {
|
|
md.mu.Lock()
|
|
defer md.mu.Unlock()
|
|
return md.close()
|
|
}
|
|
|
|
func (md *Metadata) close() (err error) {
|
|
md.cancel()
|
|
if md.wc != nil {
|
|
close(md.wc)
|
|
md.wc = nil
|
|
}
|
|
if md.f != nil {
|
|
md.f.Sync()
|
|
err = md.f.Close()
|
|
md.f = nil
|
|
}
|
|
return
|
|
}
|
|
|
|
func (md *Metadata) Delete() error {
|
|
md.mu.Lock()
|
|
defer md.mu.Unlock()
|
|
md.cmu.Lock()
|
|
md.Chunks = chunk.Chunks{}
|
|
md.cmu.Unlock()
|
|
md.close()
|
|
atomic.StoreInt64(&md.Atime, now())
|
|
return md.fs.RemoveDst(md.name)
|
|
}
|
|
|
|
func (md *Metadata) FullyCached() bool {
|
|
md.cmu.RLock()
|
|
defer md.cmu.RUnlock()
|
|
return len(md.Chunks) == 1 && md.Size == md.Chunks[0][1]
|
|
}
|
|
|
|
func (md *Metadata) HasChunk(off int64, n int) bool {
|
|
md.cmu.RLock()
|
|
defer md.cmu.RUnlock()
|
|
return md.Chunks.Exists(off, n, md.Size)
|
|
}
|
|
|
|
func (md *Metadata) ChunkSize() int64 {
|
|
md.cmu.RLock()
|
|
defer md.cmu.RUnlock()
|
|
return md.Chunks.Size()
|
|
}
|
|
|
|
func (md *Metadata) addChunk(off int64, n int) {
|
|
md.cmu.Lock()
|
|
defer md.cmu.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)
|
|
//}
|
|
|
|
type mdErr struct {
|
|
err error
|
|
}
|
|
|
|
type writeAt struct {
|
|
pos int64
|
|
data []byte
|
|
ret *[]byte
|
|
}
|
|
|
|
func (md *Metadata) resetChans() {
|
|
md.wc = make(chan writeAt, 10)
|
|
}
|
|
|
|
func (md *Metadata) writer() {
|
|
md.resetChans()
|
|
md.err.Store(nil)
|
|
go func() {
|
|
for {
|
|
select {
|
|
case <-md.ctx.Done():
|
|
return
|
|
case wa := <-md.wc:
|
|
atomic.StoreInt64(&md.Atime, now())
|
|
n, err := md.f.WriteAt(wa.data, wa.pos)
|
|
if err == nil {
|
|
md.addChunk(wa.pos, n)
|
|
md.fs.q.Add(n)
|
|
} else {
|
|
md.err.Store(&mdErr{err})
|
|
}
|
|
if wa.ret != nil {
|
|
streamingPool.Put(wa.ret)
|
|
}
|
|
}
|
|
}
|
|
}()
|
|
}
|
|
|
|
func (md *Metadata) WriteAt(data []byte, pos int64) (int, error) {
|
|
err := md.err.Load()
|
|
if err != nil && err.err != nil {
|
|
return 0, err.err
|
|
}
|
|
if md.f == nil {
|
|
err := md.openCacheFile()
|
|
if err != nil {
|
|
return 0, err
|
|
}
|
|
}
|
|
if md.ctx.Err() == context.Canceled {
|
|
return 0, context.Canceled
|
|
}
|
|
wa := writeAt{pos: pos}
|
|
ret := streamingPool.Get()
|
|
n := len(data)
|
|
buf := *ret
|
|
if n > len(buf) {
|
|
streamingPool.Put(ret)
|
|
buf = make([]byte, n)
|
|
} else {
|
|
wa.ret = ret
|
|
}
|
|
copy(buf, data)
|
|
wa.data = buf[:n]
|
|
select {
|
|
case <-md.ctx.Done():
|
|
case md.wc <- wa:
|
|
}
|
|
return n, nil
|
|
}
|
|
|
|
func (md *Metadata) Stat() (os.FileInfo, error) {
|
|
if md.f == nil {
|
|
err := md.openCacheFile()
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
}
|
|
return md.f.Stat()
|
|
}
|
|
|
|
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)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
md.f = f
|
|
if md.ctx.Err() != nil {
|
|
md.cancel()
|
|
md.ctx, md.cancel = context.WithCancel(context.Background())
|
|
}
|
|
md.writer()
|
|
return nil
|
|
}
|