// Copyright (C) 2022 Marius Schellenberger package fs import ( "context" "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) 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 if p.Prio < 0 { p.Prio = 0 } q.Sort() return } } *q = append(*q, &Preload{Name: name, ph: ph}) 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 type PreloadHandler struct { 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 <= 0 { max = 1 } ph = &PreloadHandler{ log: log, fs: fs, q: make(queue, 0), fin: make(chan string, 1), max: max, } go ph.preloadFinish(ctx) return } 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() ps = make(Preloads, len(ph.q)) for i, v := range ph.q { ps[i] = *v ps[i].Status = ph.fs.CacheStatus(v.Name) } return } func (ph *PreloadHandler) Preload(name string, prio int) { ph.mu.Lock() defer ph.mu.Unlock() ph.q.Add(ph, name, prio) ph.schedule() } 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() } } } func (p *Preload) start() { if p.Running { return } 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 (p *Preload) stop() { if p.Running { p.cancel() p.Running = false } } func (ph *PreloadHandler) preloadFinish(ctx context.Context) { for { select { case <-ctx.Done(): ph.close() return case name := <-ph.fin: ph.mu.Lock() ph.q.Remove(name) ph.schedule() ph.mu.Unlock() } } } func (ph *PreloadHandler) close() { for _, p := range ph.q { p.stop() } for { select { case <-ph.fin: default: return } } }