pbacsko commented on code in PR #1043:
URL: https://github.com/apache/yunikorn-k8shim/pull/1043#discussion_r3657302580
##########
pkg/plugin/predicates/predicate_manager.go:
##########
@@ -187,80 +171,92 @@ 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())
}
- return "", nil
+ return 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, map[string]*si.Empty) {
+ skipPlugins := sets.New[string]()
+ feasibleNodes := make(map[string]*si.Empty)
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
}
+ var mergedPreFilterResults *fwk.PreFilterResult
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]*si.Empty{}
}
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)), map[string]*si.Empty{}
+ } else {
+ if mergedPreFilterResults == nil {
+ mergedPreFilterResults = &fwk.PreFilterResult{}
+ }
+ mergedPreFilterResults.Merge(nodes)
Review Comment:
BUG: `Merge()` has a return value which is discarded. It is a pure function
which does not mutate anything.
Simply:
```
mergedPreFilterResults = mergedPreFilterResults.Merge(nodes)
```
##########
pkg/cache/context.go:
##########
@@ -706,9 +739,14 @@ func (ctx *Context) IsPodFitNode(name, node string,
allocate bool) error {
return ErrorNodeNotFound
}
// need to lock cache here as predicates need a stable view into the
cache
- ctx.schedulerCache.LockForReads()
- defer ctx.schedulerCache.UnlockForReads()
- plugin, err := ctx.predManager.Predicates(pod, targetNode, allocate)
+ ctx.schedulerCache.LockForWrites()
+ defer ctx.schedulerCache.UnlockForWrites()
+ cycleState := ctx.schedulerCache.GetCycleState(pod)
+ if cycleState == nil {
+ return ErrorCycleStateNotFound
+ }
+ plugin, err := ctx.predManager.Filter(pod, targetNode, cycleState,
allocate)
+ ctx.schedulerCache.DeleteCycleState(pod)
Review Comment:
BUG: the real flow is one `PreFilter()`, then `Filter()` once per candidate
node (core's `tryNodes()` calls `preAllocateConditions()` once, then
`tryNode()` -> `Predicates()` per node). After the first node's `Filter()` the
state is gone, so every subsequent node returns `ErrorCycleStateNotFound`,
which core treats as a predicate failure.
##########
pkg/cache/context.go:
##########
@@ -706,9 +739,14 @@ func (ctx *Context) IsPodFitNode(name, node string,
allocate bool) error {
return ErrorNodeNotFound
}
// need to lock cache here as predicates need a stable view into the
cache
- ctx.schedulerCache.LockForReads()
- defer ctx.schedulerCache.UnlockForReads()
- plugin, err := ctx.predManager.Predicates(pod, targetNode, allocate)
+ ctx.schedulerCache.LockForWrites()
+ defer ctx.schedulerCache.UnlockForWrites()
Review Comment:
I think removing `DeleteCycleState()` is sufficient. So these can go back to
read locks.
--
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]