170 lines
2.8 KiB
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
|
|
}
|
|
}
|
|
}
|