Skip to content

Commit

Permalink
Merge pull request #29 from valentin-carl/26-rework-deadline-scheduler
Browse files Browse the repository at this point in the history
refactor deadline scheduler
  • Loading branch information
KonsumGandalf authored Jan 4, 2024
2 parents ccfc1a3 + 01c072a commit 56ba63c
Show file tree
Hide file tree
Showing 2 changed files with 19 additions and 14 deletions.
21 changes: 10 additions & 11 deletions pkg/nexus/deadline/scheduler/deadline-scheduler.go
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
package deadline

import (
"log"
"fmt"
"time"

"github.com/nuclio/nuclio/pkg/nexus/common/models/interfaces"
Expand All @@ -24,6 +24,7 @@ func NewScheduler(baseNexusScheduler *common.BaseNexusScheduler, deadlineConfig
}

func NewDefaultScheduler(baseNexusScheduler *common.BaseNexusScheduler) *DeadlineScheduler {

return NewScheduler(baseNexusScheduler, *models.NewDefaultDeadlineSchedulerConfig())
}

Expand All @@ -48,19 +49,17 @@ func (ds *DeadlineScheduler) GetStatus() interfaces.SchedulerStatus {
// TODO: fix this please sleep -> something todo until next awakening (do it) -> sleep
func (ds *DeadlineScheduler) executeSchedule() {
for ds.RunFlag {
if ds.Queue.Len() == 0 {
println("Sleeping for ", ds.SleepDuration.Milliseconds(), " milliseconds")
time.Sleep(ds.SleepDuration)
continue
nextWakeUpTime := time.Now().Add(ds.SleepDuration)

}
removeUntil := time.Now().Add(ds.DeadlineRemovalThreshold)

for ds.Queue.Len() > 0 &&
ds.Queue.Peek().Deadline.Before(removeUntil) {

log.Println("Checking for expired deadlines...")
timeUntilDeadline := ds.Queue.Peek().Deadline.Sub(time.Now())
log.Println(timeUntilDeadline)
if timeUntilDeadline < ds.DeadlineRemovalThreshold {
println("Removing item from queue")
ds.Pop()
}

fmt.Println("Sleeping:", time.Until(nextWakeUpTime).Seconds(), "seconds")
time.Sleep(time.Until(nextWakeUpTime))
}
}
12 changes: 9 additions & 3 deletions pkg/nexus/nexus/nexus.go
Original file line number Diff line number Diff line change
Expand Up @@ -3,8 +3,10 @@ package nexus
import (
"log"
"sync"
"time"

bulk "github.com/nuclio/nuclio/pkg/nexus/bulk/scheduler"
"github.com/nuclio/nuclio/pkg/nexus/common/models/configs"
"github.com/nuclio/nuclio/pkg/nexus/common/models/interfaces"
common "github.com/nuclio/nuclio/pkg/nexus/common/models/structs"
queue "github.com/nuclio/nuclio/pkg/nexus/common/queue"
Expand All @@ -25,10 +27,14 @@ func Initialize() (nexus *Nexus) {
Queue: &nexusQueue,
}

baseScheduler := scheduler.NewDefaultBaseNexusScheduler(&nexusQueue)
defaultBaseScheduler := scheduler.NewDefaultBaseNexusScheduler(&nexusQueue)

deadlineScheduler := deadline.NewDefaultScheduler(baseScheduler)
bulkScheduler := bulk.NewDefaultScheduler(baseScheduler)
// deadline scheduler config
deadlineBaseSchedulerConfig := configs.NewBaseNexusSchedulerConfig(false, 10*time.Second)
deadlineBaseScheduler := scheduler.NewBaseNexusScheduler(&nexusQueue, deadlineBaseSchedulerConfig)
deadlineScheduler := deadline.NewDefaultScheduler(deadlineBaseScheduler)

bulkScheduler := bulk.NewDefaultScheduler(defaultBaseScheduler)

nexus.schedulers = map[string]interfaces.INexusScheduler{
"deadline": deadlineScheduler,
Expand Down

0 comments on commit 56ba63c

Please sign in to comment.