wilfred-s commented on code in PR #1043:
URL: https://github.com/apache/yunikorn-k8shim/pull/1043#discussion_r3881269168


##########
pkg/plugin/predicates/predicate_manager.go:
##########
@@ -36,22 +37,24 @@ import (
        apiConfig "k8s.io/kubernetes/pkg/scheduler/apis/config"
        "k8s.io/kubernetes/pkg/scheduler/apis/config/scheme"
        "k8s.io/kubernetes/pkg/scheduler/framework"
-       "k8s.io/kubernetes/pkg/scheduler/framework/plugins"
        "k8s.io/kubernetes/pkg/scheduler/framework/plugins/names"
        fwruntime "k8s.io/kubernetes/pkg/scheduler/framework/runtime"
        "k8s.io/kubernetes/pkg/scheduler/metrics"
 
        "github.com/apache/yunikorn-k8shim/pkg/log"
+

Review Comment:
   NIT remove empty line



##########
pkg/cache/context.go:
##########
@@ -692,6 +699,38 @@ func (ctx *Context) EventsToRegister(queueingHintFn 
fwk.QueueingHintFn) []fwk.Cl
        return ctx.predManager.EventsToRegister(queueingHintFn)
 }
 
+// PreFilter evaluates given prefilter based predicates based on current 
context
+func (ctx *Context) PreFilter(name string, allocate bool) 
*si.PreFilterPredicatesResponse {
+       ctx.lock.RLock()
+       defer ctx.lock.RUnlock()
+       pod := ctx.schedulerCache.GetPod(name)
+       if pod == nil {
+               log.Log(log.ShimContext).Error("failed running PreFilter 
plugin",
+                       zap.String("pod", name),
+                       zap.Error(ErrorPodNotFound))
+               return &si.PreFilterPredicatesResponse{
+                       FeasibleNodes: make(map[string]*si.Empty),
+                       Success:       false,
+               }
+       }
+       // if pod exists in cache, try to run predicates
+       // need to lock cache here as predicates need a stable view into the 
cache
+       ctx.schedulerCache.LockForWrites()
+       defer ctx.schedulerCache.UnlockForWrites()
+       feasibleNodes, cycleState, err := ctx.predManager.PreFilter(pod, 
allocate)

Review Comment:
   Can we read lock here and oinly pick the write lock for storing the cycle 
state?
   We can never run the pre-filter on the same pod in multiple threads, there 
is no chance of a race between unlock cache and store of the cycle state



##########
pkg/plugin/predicates/predicate_manager.go:
##########
@@ -191,80 +174,90 @@ func (p *predicateManagerImpl) 
removePodFromNodeNoFail(node fwk.NodeInfo, pod *v
        }
 }
 
-func (p *predicateManagerImpl) predicatesReserve(pod *v1.Pod, node 
*framework.NodeInfo) (string, error) {
+func (p *predicateManagerImpl) PreFilter(pod *v1.Pod, allocate bool) 
(map[string]*si.Empty, *framework.CycleState, error) {
        ctx := context.Background()
-       state := framework.NewCycleState()
-       return p.podFitsNode(ctx, state, *p.reservationPreFilters, 
*p.reservationFilters, pod, node)
-}
+       cycleState := framework.NewCycleState()
 
-func (p *predicateManagerImpl) predicatesAllocate(pod *v1.Pod, node 
*framework.NodeInfo) (string, error) {
-       ctx := context.Background()
-       state := framework.NewCycleState()
-       return p.podFitsNode(ctx, state, *p.allocationPreFilters, 
*p.allocationFilters, pod, node)
-}
-
-func (p *predicateManagerImpl) podFitsNode(ctx context.Context, state 
*framework.CycleState, preFilters []fwk.PreFilterPlugin, filters 
[]fwk.FilterPlugin, pod *v1.Pod, node *framework.NodeInfo) (string, error) {
-       // Run "prefilter" plugins.
-       status, plugin, skip := p.runPreFilterPlugins(ctx, state, preFilters, 
pod, node)
-       if !status.IsSuccess() && !status.IsSkip() {
-               return plugin, errors.New(status.Message())
+       var status *fwk.Status
+       var feasibleNodes map[string]*si.Empty
+       if allocate {
+               status, feasibleNodes = p.runPreFilterPlugins(ctx, cycleState, 
*p.allocationPreFilters, pod)
+       } else {
+               status, feasibleNodes = p.runPreFilterPlugins(ctx, cycleState, 
*p.reservationPreFilters, pod)
        }
-
-       // Run "filter" plugins on node
-       status, plugin = p.runFilterPlugins(ctx, filters, state, pod, node, 
skip)
-       if !status.IsSuccess() {
-               return plugin, errors.New(status.Message())
+       if !status.IsSuccess() && !status.IsSkip() {
+               return map[string]*si.Empty{}, cycleState, 
errors.New(status.Message())

Review Comment:
   if pre-filter fails we should not return a cycle state as we should never go 
to filter and we skip storing it



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to