implemented priority preloading

This commit is contained in:
ston1th 2022-05-03 22:14:47 +02:00
commit eb6272e411
7 changed files with 134 additions and 107 deletions

View file

@ -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