cachefs/pkg/fs/fs.go
2022-03-14 22:58:36 +01:00

304 lines
6.1 KiB
Go

// Copyright (C) 2022 Marius Schellenberger
package fs
import (
"context"
"encoding/json"
"errors"
"io"
stdfs "io/fs"
"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
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),
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) Open(name string) (http.File, error) {
return fs.open(name)
}
func (fs *FS) OpenFile(name string) (*File, error) {
return fs.open(name)
}
func (fs *FS) Preload(name string) {
f, err := fs.open(name)
if err != nil {
return
}
if f.md.Preload() {
fs.log.WithValues("file", name).Error(errors.New("preload is running"), "error starting preload")
return
}
unlock := func() { f.md.UnlockPreload() }
go f.Preload(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 *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 &File{f: sf, offline: offline}, 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) 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
}