pbacsko commented on code in PR #1043:
URL: https://github.com/apache/yunikorn-k8shim/pull/1043#discussion_r3505846274
##########
pkg/plugin/predicates/predicate_manager.go:
##########
@@ -187,80 +169,87 @@ 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) (string,
map[string]struct{}, *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 plugin string
+ var feasibleNodes map[string]struct{}
+ if allocate {
+ status, plugin, feasibleNodes = p.runPreFilterPlugins(ctx,
cycleState, *p.allocationPreFilters, pod)
+ } else {
+ status, plugin, 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 plugin, map[string]struct{}{}, cycleState,
errors.New(status.Message())
}
- return "", nil
+ return plugin, feasibleNodes, cycleState, nil
}
-func (p *predicateManagerImpl) runPreFilterPlugins(ctx context.Context, state
*framework.CycleState, plugins []fwk.PreFilterPlugin, pod *v1.Pod, node
*framework.NodeInfo) (*fwk.Status, string, map[string]bool) {
- var mergedNodes *fwk.PreFilterResult
- skip := make(map[string]bool)
+func (p *predicateManagerImpl) runPreFilterPlugins(ctx context.Context,
cycleState *framework.CycleState, plugins []fwk.PreFilterPlugin, pod *v1.Pod)
(*fwk.Status, string, map[string]struct{}) {
+ skipPlugins := sets.New[string]()
+ feasibleNodes := make(map[string]struct{})
allNodes, err := p.sharedLister.NodeInfos().List()
if err != nil {
log.Log(log.ShimPredicates).Error("failed to list nodes",
zap.Error(err))
- return fwk.AsStatus(err), "", skip
+ return fwk.AsStatus(err), "", feasibleNodes
}
for _, pl := range plugins {
plugin := pl.Name()
- nodes, status := p.runPreFilterPlugin(ctx, pl, state, pod,
allNodes)
+ nodes, status := pl.PreFilter(ctx, cycleState, pod, allNodes)
if status.IsSkip() {
- skip[plugin] = true
+ skipPlugins.Insert(plugin)
} else if !status.IsSuccess() {
if status.IsRejected() {
- return status, "", skip
+ return status, "", map[string]struct{}{}
}
err := errors.New(status.Message())
log.Log(log.ShimPredicates).Error("failed running
PreFilter plugin",
zap.String("pluginName", plugin),
zap.String("pod", fmt.Sprintf("%s/%s",
pod.Namespace, pod.Name)),
zap.Error(err))
- return fwk.AsStatus(errors.Join(fmt.Errorf("running
PreFilter plugin %q: ", plugin), err)), plugin, skip
+ return fwk.AsStatus(errors.Join(fmt.Errorf("running
PreFilter plugin %q: ", plugin), err)), plugin, map[string]struct{}{}
}
- // Merge is nil safe and returns a new PreFilterResult result
if mergedNodes was nil
- mergedNodes = mergedNodes.Merge(nodes)
- if !mergedNodes.AllNodes() &&
!mergedNodes.NodeNames.Has(node.Node().Name) {
- return fwk.NewStatus(fwk.UnschedulableAndUnresolvable,
"node not eligible"), plugin, skip
+ if nodes != nil {
+ for n := range nodes.NodeNames {
+ feasibleNodes[n] = struct{}{}
+ }
Review Comment:
Merging logic looks incorrect.
In Kubernetes, PreFilterResult.Merge() intersects node sets. With union, if
plugin A returns {n1,n2} and plugin B returns {n2,n3}, now it yields {n1,n2,n3}
instead of {n2}.
If I'm right, then make sure that this is covered properly with an unit test.
##########
pkg/cache/context.go:
##########
@@ -682,6 +683,28 @@ 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)
(map[string]struct{}, error) {
+ ctx.lock.RLock()
+ defer ctx.lock.RUnlock()
+ pod := ctx.schedulerCache.GetPod(name)
+ if pod == nil {
+ return map[string]struct{}{}, ErrorPodNotFound
+ }
+ // 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.LockForReads()
+ plugin, feasibleNodes, cycleState, err :=
ctx.predManager.PreFilter(pod, allocate)
+ ctx.schedulerCache.UnlockForReads()
+ ctx.schedulerCache.UpdateCycleState(pod, cycleState)
Review Comment:
I think this should be collapsed under single write lock - that's better
from the point of view of concurrency.
##########
pkg/cache/context.go:
##########
@@ -682,6 +683,28 @@ 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)
(map[string]struct{}, error) {
Review Comment:
nit: `Prefilter` or `PreFilter` ?
--
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]