cachefs/pkg/fs/preload.go

170 lines
2.8 KiB
Go

// Copyright (C) 2022 Marius Schellenberger
package fs
import (
"context"
"errors"
"sync"
"github.com/go-logr/logr"
"golang.org/x/exp/slices"
)
type Preload struct {
Name string
cancel func()
Status int
Running bool
}
type Preloads []*Preload
func (Preloads) Less(i, j *Preload) bool { return i.Name < j.Name }
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
}
func NewPreloadHandler(ctx context.Context, fs *FS, max int, log logr.Logger) (ph *PreloadHandler) {
if max < -1 {
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,
}
go ph.preload(ctx)
go ph.preloadFinish(ctx)
return
}
func (ph *PreloadHandler) CancelPreload(name string) {
ph.mu.Lock()
defer ph.mu.Unlock()
if p, ok := ph.pm[name]; ok && p.Running {
p.cancel()
}
}
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),
})
}
slices.SortFunc(ps, ps.Less)
return
}
func (ph *PreloadHandler) Preload(name string) {
ph.mu.Lock()
defer ph.mu.Unlock()
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
}
return
}
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
return
}
}
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() {
file := f
file.Close()
ph.fin <- name
})
}
}
func (ph *PreloadHandler) preloadFinish(ctx context.Context) {
for {
select {
case <-ctx.Done():
ph.close()
return
case name := <-ph.fin:
ph.mu.Lock()
delete(ph.pm, name)
ph.pre--
ph.nextPreload()
ph.mu.Unlock()
}
}
}
func (ph *PreloadHandler) close() {
for _, p := range ph.pm {
if p.Running {
p.cancel()
}
}
for {
select {
case <-ph.sched:
case <-ph.fin:
default:
return
}
}
}