fixed preload resumption
This commit is contained in:
parent
304135be2e
commit
caf12d9ca5
5 changed files with 49 additions and 5 deletions
|
|
@ -5,6 +5,7 @@ package fs
|
|||
import (
|
||||
"context"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"github.com/go-logr/logr"
|
||||
"golang.org/x/exp/slices"
|
||||
|
|
@ -44,12 +45,23 @@ func (q *queue) Remove(name string) {
|
|||
}
|
||||
}
|
||||
|
||||
func (q *queue) GetPreload(name string) *Preload {
|
||||
for _, p := range *q {
|
||||
if p.Name == name {
|
||||
return p
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
type Preload struct {
|
||||
ph *PreloadHandler
|
||||
delay time.Duration
|
||||
Name string
|
||||
cancel func()
|
||||
Prio int
|
||||
Status int
|
||||
Errc int
|
||||
Running bool
|
||||
}
|
||||
|
||||
|
|
@ -113,7 +125,7 @@ func (ph *PreloadHandler) Preload(name string, prio int) {
|
|||
func (ph *PreloadHandler) schedule() {
|
||||
for i := 0; i < len(ph.q); i++ {
|
||||
if i < ph.max {
|
||||
ph.q[i].start()
|
||||
go ph.q[i].start()
|
||||
} else {
|
||||
ph.q[i].stop(true)
|
||||
}
|
||||
|
|
@ -124,10 +136,14 @@ func (p *Preload) start() {
|
|||
if p.Running {
|
||||
return
|
||||
}
|
||||
if p.delay != 0 {
|
||||
time.Sleep(p.delay)
|
||||
}
|
||||
name := p.Name
|
||||
file, err := p.ph.fs.Open(name)
|
||||
if err != nil {
|
||||
p.ph.log.Error(err, "error staring next preload", "file", name)
|
||||
p.ph.err <- name
|
||||
return
|
||||
}
|
||||
f, ok := file.(*File)
|
||||
|
|
@ -173,8 +189,21 @@ func (ph *PreloadHandler) preloadStatus(ctx context.Context) {
|
|||
ph.mu.Unlock()
|
||||
case name := <-ph.err:
|
||||
ph.mu.Lock()
|
||||
ph.q.Add(ph, name, -1)
|
||||
ph.schedule()
|
||||
p := ph.q.GetPreload(name)
|
||||
if p != nil {
|
||||
p.stop(false)
|
||||
p.delay = time.Second * 2
|
||||
if p.Errc >= 10 {
|
||||
if p.Prio >= 0 {
|
||||
p.Prio = 0
|
||||
}
|
||||
p.Prio -= 1
|
||||
ph.schedule()
|
||||
} else {
|
||||
p.Errc += 1
|
||||
go p.start()
|
||||
}
|
||||
}
|
||||
ph.mu.Unlock()
|
||||
}
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue