diff --git a/cmd/cachefs/main.go b/cmd/cachefs/main.go index 5ae43ea..efccf35 100644 --- a/cmd/cachefs/main.go +++ b/cmd/cachefs/main.go @@ -46,7 +46,7 @@ func main() { flag.StringVar(&listenHttp, "listen", "127.0.0.1:8080", "listen addr:port") flag.StringVar(&listenWebdav, "webdav", "", "listen addr:port for webdav") flag.StringVar(&listenCache, "cache", "", "listen addr:port for cache only") - flag.IntVar(&max, "max", -1, "max parallel preloads") + flag.IntVar(&max, "max", 1, "max parallel preloads") flag.IntVar(&bs, "bs", -1, "tune preload buffer size in bytes (default: 8192)") flag.Int64Var("a, "quota", 1, "max disk usage quota for the dst cache in GiB") flag.BoolVar(&block, "block", true, "block until sftp is connected") diff --git a/pkg/fs/file.go b/pkg/fs/file.go index ed36431..88a9979 100644 --- a/pkg/fs/file.go +++ b/pkg/fs/file.go @@ -50,13 +50,14 @@ func (p *preload) Read(data []byte) (n int, err error) { func (f *File) Preload(ctx context.Context, unlock func()) { log := f.log defer f.Close() - defer unlock() if f.offline { log.V(2).Error(errors.New("no preload in offline mode"), "error preloading file") + unlock() return } if f.md.FullyCached() { log.V(2).Info("skipped preload for fully cached file") + unlock() return } @@ -65,6 +66,7 @@ func (f *File) Preload(ctx context.Context, unlock func()) { _, err := io.Copy(Discard, p) if err == context.Canceled { log.V(2).Info("preload canceled", "skipped", p.skipped, "written", p.written) + // do not call unlock() to keep preload in list return } if err != nil && err != io.EOF { @@ -75,6 +77,7 @@ func (f *File) Preload(ctx context.Context, unlock func()) { log.Error(err, "error closing cache file") } log.V(2).Info("preload finished", "skipped", p.skipped, "written", p.written) + unlock() } func (f *File) Read(p []byte) (n int, err error) { diff --git a/pkg/fs/fs.go b/pkg/fs/fs.go index 8a8c76e..3b0892e 100644 --- a/pkg/fs/fs.go +++ b/pkg/fs/fs.go @@ -145,16 +145,16 @@ func (fs *FS) QuotaUsage() (cur, max int64) { return fs.q.Usage() } -func (fs *FS) CancelPreload(name string) { - fs.ph.CancelPreload(name) +func (fs *FS) RemovePreload(name string) { + fs.ph.RemovePreload(name) } func (fs *FS) Preloads() Preloads { return fs.ph.Preloads() } -func (fs *FS) Preload(name string) { - fs.ph.Preload(name) +func (fs *FS) Preload(name string, prio int) { + fs.ph.Preload(name, prio) } func skipLog(name string) bool { diff --git a/pkg/fs/preload.go b/pkg/fs/preload.go index 1f02dd9..0f5ed58 100644 --- a/pkg/fs/preload.go +++ b/pkg/fs/preload.go @@ -4,138 +4,160 @@ package fs import ( "context" - "errors" "sync" "github.com/go-logr/logr" "golang.org/x/exp/slices" ) +type queue []*Preload + +func (queue) LessReverse(i, j *Preload) bool { return i.Prio > j.Prio } + +func (q *queue) Head() (p *Preload) { + l := len(*q) + if l > 0 { + p = (*q)[l-1] + } + return +} + +func (q *queue) Pop() (p *Preload) { + l := len(*q) + if l > 0 { + p = q.Head() + *q = (*q)[:l-1] + } + return +} + +func (q *queue) Sort() { + slices.SortFunc(*q, q.LessReverse) +} + +func (q *queue) Add(ph *PreloadHandler, name string, prio int) { + for _, p := range *q { + if p.Name == name { + p.Prio += prio + q.Sort() + return + } + } + q.Append(&Preload{Name: name, ph: ph}) +} + +func (q *queue) Append(p *Preload) { + *q = append(*q, p) + q.Sort() +} + +func (q *queue) Remove(name string) { + for i, p := range *q { + if p.Name == name { + p.stop() + *q = slices.Delete(*q, i, i+1) + q.Sort() + return + } + } +} + type Preload struct { + ph *PreloadHandler Name string cancel func() + Prio int Status int Running bool } -type Preloads []*Preload - -func (Preloads) Less(i, j *Preload) bool { return i.Name < j.Name } +type Preloads []Preload type PreloadHandler struct { - mu sync.RWMutex - log logr.Logger - fs *FS - pm map[string]*Preload - sched chan *Preload - fin chan string - pre int - max int + mu sync.RWMutex + log logr.Logger + fs *FS + q queue + fin chan string + max int } func NewPreloadHandler(ctx context.Context, fs *FS, max int, log logr.Logger) (ph *PreloadHandler) { - if max < -1 { - max = -1 + if max <= 0 { + max = 1 } ph = &PreloadHandler{ - log: log, - fs: fs, - pm: make(map[string]*Preload), - sched: make(chan *Preload, 1), - fin: make(chan string, 1), - max: max, + log: log, + fs: fs, + q: make(queue, 0), + fin: make(chan string, 1), + max: max, } - go ph.preload(ctx) go ph.preloadFinish(ctx) return } -func (ph *PreloadHandler) CancelPreload(name string) { - ph.mu.Lock() - defer ph.mu.Unlock() - p, ok := ph.pm[name] - if ok { - if p.Running { - p.cancel() - } - delete(ph.pm, name) - } +func (ph *PreloadHandler) RemovePreload(name string) { + ph.mu.RLock() + defer ph.mu.RUnlock() + ph.q.Remove(name) + ph.schedule() } func (ph *PreloadHandler) Preloads() (ps Preloads) { ph.mu.RLock() defer ph.mu.RUnlock() - for _, p := range ph.pm { - ps = append(ps, &Preload{ - Name: p.Name, - Running: p.Running, - Status: ph.fs.CacheStatus(p.Name), - }) + ps = make(Preloads, len(ph.q)) + for i, v := range ph.q { + ps[i] = *v + ps[i].Status = ph.fs.CacheStatus(v.Name) } - slices.SortFunc(ps, ps.Less) return } -func (ph *PreloadHandler) Preload(name string) { +func (ph *PreloadHandler) Preload(name string, prio int) { ph.mu.Lock() defer ph.mu.Unlock() + ph.q.Add(ph, name, prio) + ph.schedule() +} - p, ok := ph.pm[name] - if !ok { - p = &Preload{Name: name} - ph.pm[name] = p - if ph.max == -1 || ph.pre < ph.max { - ph.sched <- p +func (ph *PreloadHandler) schedule() { + for i := 0; i < len(ph.q); i++ { + if i < ph.max { + ph.q[i].start() + } else { + ph.q[i].stop() } - return } +} +func (p *Preload) start() { if p.Running { - ph.log.WithValues("file", name).Error(errors.New("preload is running"), "error starting preload") - } - return -} - -func (ph *PreloadHandler) nextPreload() { - if ph.pre == ph.max { return } - for _, p := range ph.pm { - if p.Running { - continue - } - ph.sched <- p + name := p.Name + file, err := p.ph.fs.Open(name) + if err != nil { + p.ph.log.Error(err, "error staring next preload", "file", name) return } + f, ok := file.(*File) + if !ok { + return + } + ctx, cancel := context.WithCancel(context.Background()) + p.Running = true + p.cancel = cancel + go f.Preload(ctx, func() { + p.ph.fin <- name + }) } -func (ph *PreloadHandler) preload(ctx context.Context) { - var p *Preload - for { - select { - case <-ctx.Done(): - return - case p = <-ph.sched: - } - name := p.Name - file, err := ph.fs.Open(name) - if err != nil { - ph.log.Error(err, "error staring next preload", "file", name) - return - } - f, ok := file.(*File) - if !ok { - continue - } - ctx, cancel := context.WithCancel(context.Background()) - ph.mu.Lock() - ph.pre++ - p.Running = true - p.cancel = cancel - ph.mu.Unlock() - go f.Preload(ctx, func() { - ph.fin <- name - }) +func (p *Preload) stop() { + if p.Running { + p.cancel() + p.Running = false } } @@ -147,23 +169,19 @@ func (ph *PreloadHandler) preloadFinish(ctx context.Context) { return case name := <-ph.fin: ph.mu.Lock() - delete(ph.pm, name) - ph.pre-- - ph.nextPreload() + ph.q.Remove(name) + ph.schedule() ph.mu.Unlock() } } } func (ph *PreloadHandler) close() { - for _, p := range ph.pm { - if p.Running { - p.cancel() - } + for _, p := range ph.q { + p.stop() } for { select { - case <-ph.sched: case <-ph.fin: default: return diff --git a/pkg/srv/interceptor.go b/pkg/srv/interceptor.go index bf058a0..4300de1 100644 --- a/pkg/srv/interceptor.go +++ b/pkg/srv/interceptor.go @@ -60,6 +60,7 @@ type file struct { URI template.HTML Anchor string Status int + Prio int Running bool } @@ -136,12 +137,10 @@ func getPreloads(path string, fs *fs.FS) (dc dirContents) { Name: template.HTML(p.Name), URI: template.HTML(p.Name), Status: p.Status, + Prio: p.Prio, Running: p.Running, }) } - slices.SortFunc(dc.Files, func(i, _ file) bool { - return i.Running - }) return } diff --git a/pkg/srv/srv.go b/pkg/srv/srv.go index 6362327..f1808a5 100644 --- a/pkg/srv/srv.go +++ b/pkg/srv/srv.go @@ -85,11 +85,16 @@ func (fs *FileServer) ServeHTTP(w http.ResponseWriter, r *http.Request) { } if !d.IsDir() { - if option == "p" || option == "s" { - if option == "p" { - fs.fs.Preload(p) - } else if option == "s" { - fs.fs.CancelPreload(p) + if option == "p" || option == "i" || option == "d" || option == "s" { + switch option { + case "p": + fs.fs.Preload(p, 0) + case "i": + fs.fs.Preload(p, 1) + case "d": + fs.fs.Preload(p, -1) + case "s": + fs.fs.RemovePreload(p) } rdir := "/" if r.FormValue("r") == "preloads" { diff --git a/pkg/srv/templates/preloads.html b/pkg/srv/templates/preloads.html index cc7f072..63ab6c1 100644 --- a/pkg/srv/templates/preloads.html +++ b/pkg/srv/templates/preloads.html @@ -1,5 +1,7 @@ {{define "body" -}}
[v]: show video
+[+]: increase priority
+[-]: decrease priority
[s]: stop preloading
Quota: {{.QuotaCur}} / {{.QuotaMax}} GiB
@@ -12,7 +14,7 @@ Quota: {{.QuotaCur}} / {{.QuotaMax}} GiB