added disk quota
This commit is contained in:
parent
4499e88245
commit
9cc2b04342
26 changed files with 1376 additions and 61 deletions
|
|
@ -2,15 +2,15 @@
|
|||
|
||||
package fs
|
||||
|
||||
import "sort"
|
||||
import (
|
||||
"golang.org/x/exp/slices"
|
||||
)
|
||||
|
||||
type Chunk [2]int64
|
||||
|
||||
type Chunks []Chunk
|
||||
|
||||
func (cs Chunks) Len() int { return len(cs) }
|
||||
func (cs Chunks) Less(i, j int) bool { return cs[i][0] < cs[j][0] }
|
||||
func (cs Chunks) Swap(i, j int) { cs[i], cs[j] = cs[j], cs[i] }
|
||||
func (Chunks) Less(i, j Chunk) bool { return i[0] < j[0] }
|
||||
|
||||
func (cs Chunks) Exists(off int64, n int, size int64) bool {
|
||||
end := off + int64(n)
|
||||
|
|
@ -42,7 +42,7 @@ func (cs *Chunks) merge() {
|
|||
if len(c) < 2 {
|
||||
return
|
||||
}
|
||||
sort.Sort(*cs)
|
||||
slices.SortFunc(c, c.Less)
|
||||
for i := 0; i < len(c); i++ {
|
||||
if i+1 == len(c) {
|
||||
break
|
||||
|
|
@ -51,7 +51,7 @@ func (cs *Chunks) merge() {
|
|||
if c[i+1][1] > c[i][1] {
|
||||
c[i][1] = c[i+1][1]
|
||||
}
|
||||
c = append(c[:i+1], c[i+2:]...)
|
||||
c = slices.Delete(c, i+1, i+2)
|
||||
i--
|
||||
}
|
||||
}
|
||||
|
|
|
|||
80
pkg/fs/fs.go
80
pkg/fs/fs.go
|
|
@ -18,6 +18,7 @@ import (
|
|||
"time"
|
||||
|
||||
"github.com/go-logr/logr"
|
||||
"golang.org/x/exp/maps"
|
||||
)
|
||||
|
||||
type FS struct {
|
||||
|
|
@ -29,12 +30,13 @@ type FS struct {
|
|||
mdf *os.File
|
||||
jenc *json.Encoder
|
||||
dc *DirCache
|
||||
mm map[string]*Metadata
|
||||
prem map[string]func()
|
||||
q *Quota
|
||||
md MetadataMap
|
||||
pm map[string]func()
|
||||
cancel func()
|
||||
}
|
||||
|
||||
func NewFS(src, dst, metadata string, log logr.Logger) (fs *FS, err error) {
|
||||
func NewFS(max int64, src, dst, metadata string, log logr.Logger) (fs *FS, err error) {
|
||||
if !filepath.IsAbs(src) {
|
||||
return nil, errors.New("src path is not absolute")
|
||||
}
|
||||
|
|
@ -58,15 +60,20 @@ func NewFS(src, dst, metadata string, log logr.Logger) (fs *FS, err error) {
|
|||
dst: dst,
|
||||
mdf: mdf,
|
||||
dc: NewDirCache(),
|
||||
mm: make(map[string]*Metadata),
|
||||
prem: make(map[string]func()),
|
||||
md: make(MetadataMap),
|
||||
pm: make(map[string]func()),
|
||||
jenc: json.NewEncoder(mdf),
|
||||
}
|
||||
err = json.NewDecoder(mdf).Decode(&fs.mm)
|
||||
fs.q, err = NewQuota(max, fs, fs.log.WithName("quota"))
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
err = json.NewDecoder(mdf).Decode(&fs.md)
|
||||
if err == io.EOF {
|
||||
err = nil
|
||||
}
|
||||
fs.initMetadata()
|
||||
fs.q.Init()
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
fs.cancel = cancel
|
||||
go fs.flusher(ctx)
|
||||
|
|
@ -75,26 +82,25 @@ func NewFS(src, dst, metadata string, log logr.Logger) (fs *FS, err error) {
|
|||
|
||||
func (fs *FS) initMetadata() {
|
||||
log := fs.log
|
||||
for k, v := range fs.mm {
|
||||
maps.DeleteFunc(fs.md, func(k string, v *Metadata) bool {
|
||||
_, 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
|
||||
return true
|
||||
}
|
||||
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
|
||||
return false
|
||||
}
|
||||
delete(fs.mm, k)
|
||||
continue
|
||||
return true
|
||||
}
|
||||
v.fs = fs
|
||||
v.name = k
|
||||
}
|
||||
return false
|
||||
})
|
||||
}
|
||||
|
||||
func (fs *FS) Stat(name string) (fi stdfs.FileInfo, err error) {
|
||||
|
|
@ -126,7 +132,7 @@ func (fs *FS) statDst(name string) (fi stdfs.FileInfo, err error) {
|
|||
}
|
||||
|
||||
func (fs *FS) CancelPreload(name string) {
|
||||
if cancel, ok := fs.prem[name]; ok {
|
||||
if cancel, ok := fs.pm[name]; ok {
|
||||
cancel()
|
||||
}
|
||||
}
|
||||
|
|
@ -146,9 +152,9 @@ func (fs *FS) Preload(name string) {
|
|||
return
|
||||
}
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
fs.prem[name] = cancel
|
||||
fs.pm[name] = cancel
|
||||
unlock := func() {
|
||||
delete(fs.prem, name)
|
||||
delete(fs.pm, name)
|
||||
f.md.UnlockPreload()
|
||||
}
|
||||
go f.Preload(ctx, unlock)
|
||||
|
|
@ -199,15 +205,17 @@ func (fs *FS) Open(name string) (f http.File, err error) {
|
|||
offline = true
|
||||
log.V(2).Info("dir offline mode", "path", dp)
|
||||
}
|
||||
return &Dir{f: sf, dc: fs.dc}, err
|
||||
f = &Dir{f: sf, dc: fs.dc}
|
||||
return
|
||||
}
|
||||
md := fs.metadata(name, sfi.Size())
|
||||
return &File{
|
||||
f = &File{
|
||||
log: log,
|
||||
f: sf,
|
||||
md: md,
|
||||
offline: offline,
|
||||
}, nil
|
||||
}
|
||||
return
|
||||
}
|
||||
|
||||
func (fs *FS) openCacheFile(name string, size int64) (df *os.File, err error) {
|
||||
|
|
@ -248,19 +256,29 @@ func (fs *FS) openCacheFile(name string, size int64) (df *os.File, err error) {
|
|||
|
||||
func (fs *FS) metadata(name string, size int64) (md *Metadata) {
|
||||
fs.mu.Lock()
|
||||
fs.log.Info("fs.mu.Lock()")
|
||||
defer fs.mu.Unlock()
|
||||
if md, ok := fs.mm[name]; ok {
|
||||
defer fs.log.Info("fs.mu.Unlock()")
|
||||
if md, ok := fs.md[name]; ok {
|
||||
if size > 0 && md.Size != size {
|
||||
log := fs.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: fs, name: name}
|
||||
fs.mm[name] = md
|
||||
fs.md[name] = md
|
||||
return
|
||||
}
|
||||
|
||||
func (fs *FS) CacheStatus(name string) int {
|
||||
fs.mu.RLock()
|
||||
defer fs.mu.RUnlock()
|
||||
if md, ok := fs.mm[name]; ok {
|
||||
if md, ok := fs.md[name]; ok {
|
||||
return int(math.Round(float64(md.Chunks.Size()) / float64(md.Size) * 100))
|
||||
}
|
||||
return -1
|
||||
|
|
@ -281,12 +299,11 @@ 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
|
||||
}
|
||||
}
|
||||
log.V(2).Info("flushing metadata to disk")
|
||||
md := maps.Clone(fs.md)
|
||||
maps.DeleteFunc(md, func(_ string, v *Metadata) bool {
|
||||
return len(v.Chunks) == 0
|
||||
})
|
||||
_, err := fs.mdf.Seek(0, io.SeekStart)
|
||||
if err != nil {
|
||||
log.Error(err, "failure seeking metadata file")
|
||||
|
|
@ -297,7 +314,7 @@ func (fs *FS) flushMetadata() {
|
|||
log.Error(err, "failure truncating metadata file")
|
||||
return
|
||||
}
|
||||
err = fs.jenc.Encode(mm)
|
||||
err = fs.jenc.Encode(md)
|
||||
if err != nil {
|
||||
log.Error(err, "failure flushing metadata file")
|
||||
}
|
||||
|
|
@ -305,7 +322,10 @@ func (fs *FS) flushMetadata() {
|
|||
|
||||
func (fs *FS) Close() {
|
||||
fs.cancel()
|
||||
for _, md := range fs.mm {
|
||||
for _, cancel := range fs.pm {
|
||||
cancel()
|
||||
}
|
||||
for _, md := range fs.md {
|
||||
md.f.Sync()
|
||||
md.f.Close()
|
||||
}
|
||||
|
|
|
|||
|
|
@ -5,14 +5,23 @@ package fs
|
|||
import (
|
||||
"os"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"time"
|
||||
)
|
||||
|
||||
func now() int64 {
|
||||
return time.Now().UTC().Unix()
|
||||
}
|
||||
|
||||
type MetadataMap map[string]*Metadata
|
||||
|
||||
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"`
|
||||
preload bool `json:"-"`
|
||||
}
|
||||
|
|
@ -23,6 +32,7 @@ func (md *Metadata) Delete() error {
|
|||
md.Chunks = Chunks{}
|
||||
md.f.Close()
|
||||
md.f = nil
|
||||
atomic.StoreInt64(&md.Atime, now())
|
||||
return md.fs.RemoveDst(md.name)
|
||||
}
|
||||
|
||||
|
|
@ -45,6 +55,7 @@ func (md *Metadata) ReadAt(p []byte, pos int64) (int, error) {
|
|||
return 0, err
|
||||
}
|
||||
}
|
||||
atomic.StoreInt64(&md.Atime, now())
|
||||
return md.f.ReadAt(p, pos)
|
||||
}
|
||||
|
||||
|
|
@ -55,6 +66,8 @@ func (md *Metadata) WriteAt(data []byte, pos int64) (int, error) {
|
|||
return 0, err
|
||||
}
|
||||
}
|
||||
atomic.StoreInt64(&md.Atime, now())
|
||||
md.fs.q.Add(len(data))
|
||||
return md.f.WriteAt(data, pos)
|
||||
}
|
||||
|
||||
|
|
|
|||
70
pkg/fs/quota.go
Normal file
70
pkg/fs/quota.go
Normal file
|
|
@ -0,0 +1,70 @@
|
|||
// Copyright (C) 2022 Marius Schellenberger
|
||||
|
||||
package fs
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"math"
|
||||
"sync/atomic"
|
||||
|
||||
"github.com/go-logr/logr"
|
||||
)
|
||||
|
||||
var InvalidMaxQuota = errors.New("invalid max quota")
|
||||
|
||||
type Quota struct {
|
||||
log logr.Logger
|
||||
fs *FS
|
||||
max int64
|
||||
cur int64
|
||||
}
|
||||
|
||||
func NewQuota(max int64, fs *FS, log logr.Logger) (q *Quota, err error) {
|
||||
if max <= 0 {
|
||||
err = InvalidMaxQuota
|
||||
return
|
||||
}
|
||||
q = &Quota{log: log, fs: fs, max: max}
|
||||
return
|
||||
}
|
||||
|
||||
func (q *Quota) Init() {
|
||||
for _, v := range q.fs.md {
|
||||
q.cur += v.Chunks.Size()
|
||||
}
|
||||
q.log.Info("quota usage", "current", q.cur, "max", q.max)
|
||||
}
|
||||
|
||||
func (q *Quota) cleanup() {
|
||||
if !q.fs.mu.TryLock() {
|
||||
return
|
||||
}
|
||||
q.log.Info("!q.fs.mu.TryLock()")
|
||||
q.log.Info("quota usage", "current", q.cur, "max", q.max)
|
||||
defer q.fs.mu.Unlock()
|
||||
defer q.log.Info("q.fs.mu.Unlock()")
|
||||
atime := int64(math.MaxInt64)
|
||||
var m *Metadata
|
||||
for _, v := range q.fs.md {
|
||||
if v.Atime != 0 && v.f != nil && v.Atime < atime {
|
||||
atime = v.Atime
|
||||
m = v
|
||||
}
|
||||
}
|
||||
if m != nil {
|
||||
log := q.log.WithValues("file", m.name)
|
||||
log.Info("deleting oldest file")
|
||||
err := m.Delete()
|
||||
if err != nil {
|
||||
q.log.Error(err, "error deleting oldest file")
|
||||
return
|
||||
}
|
||||
atomic.AddInt64(&q.cur, -m.Size)
|
||||
}
|
||||
}
|
||||
|
||||
func (q *Quota) Add(n int) {
|
||||
if q.max < atomic.AddInt64(&q.cur, int64(n)) {
|
||||
q.cleanup()
|
||||
}
|
||||
}
|
||||
|
|
@ -10,6 +10,8 @@ import (
|
|||
"strings"
|
||||
|
||||
"cachefs/pkg/fs"
|
||||
|
||||
"golang.org/x/exp/slices"
|
||||
)
|
||||
|
||||
type statusInterceptor struct {
|
||||
|
|
@ -51,9 +53,7 @@ type file struct {
|
|||
|
||||
type files []file
|
||||
|
||||
func (f files) Len() int { return len(f) }
|
||||
func (f files) Less(i, j int) bool { return f[i].Name < f[j].Name }
|
||||
func (f files) Swap(i, j int) { f[i], f[j] = f[j], f[i] }
|
||||
func (files) Less(i, j file) bool { return i.Name < j.Name }
|
||||
|
||||
type responseInterceptor struct {
|
||||
buf bytes.Buffer
|
||||
|
|
@ -98,7 +98,7 @@ func (r *responseInterceptor) GetPaths(path string, fs *fs.FS) (dir dirContents,
|
|||
}
|
||||
}
|
||||
sort.Strings(dir.Dirs)
|
||||
sort.Sort(dir.Files)
|
||||
slices.SortFunc(dir.Files, dir.Files.Less)
|
||||
return
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -7,7 +7,9 @@ import (
|
|||
"io"
|
||||
stdfs "io/fs"
|
||||
"net/http"
|
||||
"os"
|
||||
"path"
|
||||
"runtime"
|
||||
"strings"
|
||||
|
||||
"cachefs/pkg/fs"
|
||||
|
|
@ -15,6 +17,21 @@ import (
|
|||
"github.com/go-logr/logr"
|
||||
)
|
||||
|
||||
func PrintStack() {
|
||||
os.Stderr.Write(Stack())
|
||||
}
|
||||
|
||||
func Stack() []byte {
|
||||
buf := make([]byte, 1024)
|
||||
for {
|
||||
n := runtime.Stack(buf, true)
|
||||
if n < len(buf) {
|
||||
return buf[:n]
|
||||
}
|
||||
buf = make([]byte, 2*len(buf))
|
||||
}
|
||||
}
|
||||
|
||||
type FileServer struct {
|
||||
log logr.Logger
|
||||
fs *fs.FS
|
||||
|
|
@ -51,13 +68,17 @@ func (fs *FileServer) ServeHTTP(w http.ResponseWriter, r *http.Request) {
|
|||
r.URL.Path = upath
|
||||
}
|
||||
p := path.Clean(upath)
|
||||
option := r.FormValue("o")
|
||||
if option == "debug" {
|
||||
PrintStack()
|
||||
}
|
||||
d, _, err := fs.fs.StatWithOffline(p)
|
||||
if err != nil {
|
||||
msg, code := toHTTPError(err)
|
||||
http.Error(w, msg, code)
|
||||
return
|
||||
}
|
||||
option := r.FormValue("o")
|
||||
//option := r.FormValue("o")
|
||||
if option == "v" {
|
||||
err = video.Execute(w, r.URL.Path)
|
||||
if err != nil {
|
||||
|
|
|
|||
|
|
@ -30,7 +30,7 @@ var (
|
|||
{{end -}}
|
||||
{{range $s := .Files -}}
|
||||
<tr><td><a id="{{$s.Anchor}}" href="{{$s.Name}}">{{$s.Name}}</a></td>
|
||||
<td><a href="{{$s.Name}}?o=v">[v]</a> <a href="{{$s.Name}}?o=n">[n]</a> <a href="{{$s.Name}}?o=p">[p]</a> <a href="{{$s.Name}}?o=s">[s]</a> {{if ge $s.Status 0}}{{$s.Status}}%{{end}}</td></tr>
|
||||
<td><a href="{{$s.Name}}?o=v">[v]</a> <a href="{{$s.Name}}?o=n">[n]</a> <a href="{{$s.Name}}?o=p">[p]</a> <a href="{{$s.Name}}?o=s">[s]</a> {{if ne $s.Status -1}}{{$s.Status}}%{{end}}</td></tr>
|
||||
{{end -}}
|
||||
</table>
|
||||
</article>
|
||||
|
|
|
|||
|
|
@ -9,6 +9,7 @@ import (
|
|||
"cachefs/pkg/fs"
|
||||
|
||||
"github.com/go-logr/logr"
|
||||
"golang.org/x/exp/slices"
|
||||
"golang.org/x/net/webdav"
|
||||
)
|
||||
|
||||
|
|
@ -65,10 +66,5 @@ func skipDavLog(path string) bool {
|
|||
if i > 0 {
|
||||
path = path[i+1:]
|
||||
}
|
||||
for _, v := range skipDavFiles {
|
||||
if path == v {
|
||||
return true
|
||||
}
|
||||
}
|
||||
return false
|
||||
return slices.Contains(skipDavFiles, path)
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue