This is an automated email from the ASF dual-hosted git repository.
manirajv06 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/yunikorn-core.git
The following commit(s) were added to refs/heads/master by this push:
new 7d26b8dc [YUNIKORN-3300] Use feasible nodes returned from K8s
PreFilter plugin in core itself (#1096)
7d26b8dc is described below
commit 7d26b8dc0aeb10f6b558d304a0b585bb47c66d14
Author: mani <[email protected]>
AuthorDate: Thu Aug 27 17:56:57 2026 +0530
[YUNIKORN-3300] Use feasible nodes returned from K8s PreFilter plugin in
core itself (#1096)
K8s PreFilter plugin return list of feasible nodes based on the pod spec.
Re-use the same in core itself for every next node being picked from node
iterator before doing the predicate checks to improve the overall efficiency.
Closes: #1096
Signed-off-by: mani <[email protected]>
---
go.mod | 35 ++--
go.sum | 88 +++++-----
pkg/common/errors.go | 3 +-
pkg/mock/predicate_plugin.go | 72 +++++---
pkg/mock/preemption_predicate_plugin.go | 77 ++++----
pkg/mock/rm_callback.go | 8 +
pkg/plugins/plugins_test.go | 4 +
pkg/scheduler/objects/allocation.go | 22 +++
pkg/scheduler/objects/application.go | 152 +++++++++++-----
pkg/scheduler/objects/application_test.go | 280 ++++++++++++++++++------------
pkg/scheduler/objects/node.go | 15 ++
pkg/scheduler/objects/node_test.go | 65 +++++--
pkg/scheduler/objects/preemption.go | 16 ++
pkg/scheduler/objects/preemption_test.go | 234 +++++++++++++++----------
pkg/scheduler/partition_test.go | 4 +-
15 files changed, 697 insertions(+), 378 deletions(-)
diff --git a/go.mod b/go.mod
index 2d3ee6f5..fbafcbb8 100644
--- a/go.mod
+++ b/go.mod
@@ -22,22 +22,22 @@ module github.com/apache/yunikorn-core
go 1.25.0
require (
- github.com/apache/yunikorn-scheduler-interface
v0.0.0-20260727092410-674338955bdf
- github.com/go-ldap/ldap/v3 v3.4.13
+ github.com/apache/yunikorn-scheduler-interface
v0.0.0-20260727104803-9a5c60e5c879
+ github.com/go-ldap/ldap/v3 v3.4.14
github.com/google/go-cmp v0.7.0
github.com/google/uuid v1.6.0
github.com/julienschmidt/httprouter v1.3.0
github.com/looplab/fsm v1.0.3
- github.com/prometheus/client_golang v1.23.2
+ github.com/prometheus/client_golang v1.24.1
github.com/prometheus/client_model v0.6.2
- github.com/prometheus/common v0.67.5
+ github.com/prometheus/common v0.70.1
github.com/sasha-s/go-deadlock v0.3.9
github.com/tidwall/btree v1.8.1
- go.uber.org/zap v1.27.1
- go.yaml.in/yaml/v3 v3.0.4
- golang.org/x/exp v0.0.0-20260312153236-7ab1446f8b90
+ go.uber.org/zap v1.28.0
+ go.yaml.in/yaml/v3 v3.0.5
+ golang.org/x/exp v0.0.0-20260727155853-b88d891fe743
golang.org/x/time v0.15.0
- google.golang.org/grpc v1.82.1
+ google.golang.org/grpc v1.83.0
gotest.tools/v3 v3.5.2
)
@@ -45,18 +45,17 @@ require (
github.com/Azure/go-ntlmssp v0.1.1 // indirect
github.com/beorn7/perks v1.0.1 // indirect
github.com/cespare/xxhash/v2 v2.3.0 // indirect
- github.com/go-asn1-ber/asn1-ber v1.5.8-0.20250403174932-29230038a667 //
indirect
+ github.com/go-asn1-ber/asn1-ber v1.5.8 // indirect
github.com/kylelemons/godebug v1.1.0 // indirect
github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 //
indirect
- github.com/petermattis/goid v0.0.0-20250813065127-a731cc31b4fe //
indirect
- github.com/prometheus/procfs v0.16.1 // indirect
- go.uber.org/multierr v1.10.0 // indirect
- go.yaml.in/yaml/v2 v2.4.3 // indirect
- golang.org/x/crypto v0.51.0 // indirect
- golang.org/x/net v0.54.0 // indirect
- golang.org/x/sys v0.45.0 // indirect
- golang.org/x/text v0.37.0 // indirect
- google.golang.org/genproto/googleapis/rpc
v0.0.0-20260414002931-afd174a4e478 // indirect
+ github.com/petermattis/goid v0.0.0-20260725062400-500c67a39b75 //
indirect
+ github.com/prometheus/procfs v0.21.1 // indirect
+ go.uber.org/multierr v1.11.0 // indirect
+ golang.org/x/crypto v0.54.0 // indirect
+ golang.org/x/net v0.57.0 // indirect
+ golang.org/x/sys v0.47.0 // indirect
+ golang.org/x/text v0.40.0 // indirect
+ google.golang.org/genproto/googleapis/rpc
v0.0.0-20260803160001-6ac0973c030d // indirect
google.golang.org/protobuf v1.36.11 // indirect
)
diff --git a/go.sum b/go.sum
index 82d63445..d68b2e40 100644
--- a/go.sum
+++ b/go.sum
@@ -2,18 +2,18 @@ github.com/Azure/go-ntlmssp v0.1.1
h1:l+FM/EEMb0U9QZE7mKNEDw5Mu3mFiaa2GKOoTSsNDP
github.com/Azure/go-ntlmssp v0.1.1/go.mod
h1:NYqdhxd/8aAct/s4qSYZEerdPuH1liG2/X9DiVTbhpk=
github.com/alexbrainman/sspi v0.0.0-20250919150558-7d374ff0d59e
h1:4dAU9FXIyQktpoUAgOJK3OTFc/xug0PCXYCqU0FgDKI=
github.com/alexbrainman/sspi v0.0.0-20250919150558-7d374ff0d59e/go.mod
h1:cEWa1LVoE5KvSD9ONXsZrj0z6KqySlCCNKHlLzbqAt4=
-github.com/apache/yunikorn-scheduler-interface
v0.0.0-20260727092410-674338955bdf
h1:IXEpAeqZgCXODJtvB6Ib5CCcpRyvF0mqcicz80rA3Cc=
-github.com/apache/yunikorn-scheduler-interface
v0.0.0-20260727092410-674338955bdf/go.mod
h1:qb739Bdm82PH7gsfEYabulGF90xKGNQ1hWmf197rDfw=
+github.com/apache/yunikorn-scheduler-interface
v0.0.0-20260727104803-9a5c60e5c879
h1:uNpBybWm0opxgMyDs+vh9JNpDSHbZ9o9aKmKQw17JVw=
+github.com/apache/yunikorn-scheduler-interface
v0.0.0-20260727104803-9a5c60e5c879/go.mod
h1:qb739Bdm82PH7gsfEYabulGF90xKGNQ1hWmf197rDfw=
github.com/beorn7/perks v1.0.1 h1:VlbKKnNfV8bJzeqoa4cOKqO6bYr3WgKZxO8Z16+hsOM=
github.com/beorn7/perks v1.0.1/go.mod
h1:G2ZrVWU2WbWT9wwq4/hrbKbnv/1ERSJQ0ibhJ6rlkpw=
github.com/cespare/xxhash/v2 v2.3.0
h1:UL815xU9SqsFlibzuggzjXhog7bL6oX9BbNZnL2UFvs=
github.com/cespare/xxhash/v2 v2.3.0/go.mod
h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs=
github.com/davecgh/go-spew v1.1.1
h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c=
github.com/davecgh/go-spew v1.1.1/go.mod
h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
-github.com/go-asn1-ber/asn1-ber v1.5.8-0.20250403174932-29230038a667
h1:BP4M0CvQ4S3TGls2FvczZtj5Re/2ZzkV9VwqPHH/3Bo=
-github.com/go-asn1-ber/asn1-ber v1.5.8-0.20250403174932-29230038a667/go.mod
h1:hEBeB/ic+5LoWskz+yKT7vGhhPYkProFKoKdwZRWMe0=
-github.com/go-ldap/ldap/v3 v3.4.13
h1:+x1nG9h+MZN7h/lUi5Q3UZ0fJ1GyDQYbPvbuH38baDQ=
-github.com/go-ldap/ldap/v3 v3.4.13/go.mod
h1:LxsGZV6vbaK0sIvYfsv47rfh4ca0JXokCoKjZxsszv0=
+github.com/go-asn1-ber/asn1-ber v1.5.8
h1:H9AZkK22UOmfX8J84ubyaZxKJZ3FMHVwn8swoMML7iQ=
+github.com/go-asn1-ber/asn1-ber v1.5.8/go.mod
h1:hEBeB/ic+5LoWskz+yKT7vGhhPYkProFKoKdwZRWMe0=
+github.com/go-ldap/ldap/v3 v3.4.14
h1:D6PYdEgsaVzsXyr6w/yDC06Ria4uUhWm+Rb+er8lfAs=
+github.com/go-ldap/ldap/v3 v3.4.14/go.mod
h1:S4eJUMUNjDkE0ZJtIZdybwyb03sGGLW6gxXT1Hs8VKA=
github.com/go-logr/logr v1.4.3 h1:CjnDlHq8ikf6E492q6eKboGOC0T8CDaOvkHCIg8idEI=
github.com/go-logr/logr v1.4.3/go.mod
h1:9T104GzyrTigFIr8wt5mBrctHMim0Nb2HLGrmQ40KvY=
github.com/go-logr/stdr v1.2.2 h1:hSWxHoqTgW2S2qGc0LTAI563KZ5YKYRhT3MFKZMbjag=
@@ -40,32 +40,27 @@ github.com/jcmturner/rpc/v2 v2.0.3
h1:7FXXj8Ti1IaVFpSAziCZWNzbNuZmnvw/i6CqLNdWfZ
github.com/jcmturner/rpc/v2 v2.0.3/go.mod
h1:VUJYCIDm3PVOEHw8sgt091/20OJjskO/YJki3ELg/Hc=
github.com/julienschmidt/httprouter v1.3.0
h1:U0609e9tgbseu3rBINet9P48AI/D3oJs4dN7jwJOQ1U=
github.com/julienschmidt/httprouter v1.3.0/go.mod
h1:JR6WtHb+2LUe8TCKY3cZOxFyyO8IZAc4RVcycCCAKdM=
-github.com/klauspost/compress v1.18.0
h1:c/Cqfb0r+Yi+JtIEq73FWXVkRonBlf0CRNYc8Zttxdo=
-github.com/klauspost/compress v1.18.0/go.mod
h1:2Pp+KzxcywXVXMr50+X0Q/Lsb43OQHYWRCY2AiWywWQ=
-github.com/kr/pretty v0.3.1 h1:flRD4NNwYAUpkphVc1HcthR4KEIFJ65n8Mw5qdRn3LE=
-github.com/kr/pretty v0.3.1/go.mod
h1:hoEshYVHaxMs3cyo3Yncou5ZscifuDolrwPKZanG3xk=
-github.com/kr/text v0.2.0 h1:5Nx0Ya0ZqY2ygV366QzturHI13Jq95ApcVaJBhpS+AY=
-github.com/kr/text v0.2.0/go.mod
h1:eLer722TekiGuMkidMxC/pM04lWEeraHUUmBw8l2grE=
+github.com/klauspost/compress v1.19.1
h1:VsB4HPswih7mmZ8WleSFQ75c/Ui1M4trX5oAsJnhSlk=
+github.com/klauspost/compress v1.19.1/go.mod
h1:cwPg85FWrGar70rWktvGQj8/hthj3wpl0PGDogxkrSQ=
github.com/kylelemons/godebug v1.1.0
h1:RPNrshWIDI6G2gRW9EHilWtl7Z6Sb1BR0xunSBf0SNc=
github.com/kylelemons/godebug v1.1.0/go.mod
h1:9/0rRGxNHcop5bhtWyNeEfOS8JIWk580+fNqagV/RAw=
github.com/looplab/fsm v1.0.3 h1:qtxBsa2onOs0qFOtkqwf5zE0uP0+Te+wlIvXctPKpcw=
github.com/looplab/fsm v1.0.3/go.mod
h1:PmD3fFvQEIsjMEfvZdrCDZ6y8VwKTwWNjlpEr6IKPO4=
github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822
h1:C3w9PqII01/Oq1c1nUAm88MOHcQC9l5mIlSMApZMrHA=
github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822/go.mod
h1:+n7T8mK8HuQTcFwEeznm/DIxMOiR9yIdICNftLE1DvQ=
-github.com/petermattis/goid v0.0.0-20250813065127-a731cc31b4fe
h1:vHpqOnPlnkba8iSxU4j/CvDSS9J4+F4473esQsYLGoE=
github.com/petermattis/goid v0.0.0-20250813065127-a731cc31b4fe/go.mod
h1:pxMtw7cyUw6B2bRH0ZBANSPg+AoSud1I1iyJHI69jH4=
+github.com/petermattis/goid v0.0.0-20260725062400-500c67a39b75
h1:VmZ6mKVkxavKEhEy4ZYyV7BwBYBFBP0TwIqmLk84fpU=
+github.com/petermattis/goid v0.0.0-20260725062400-500c67a39b75/go.mod
h1:pxMtw7cyUw6B2bRH0ZBANSPg+AoSud1I1iyJHI69jH4=
github.com/pmezard/go-difflib v1.0.0
h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM=
github.com/pmezard/go-difflib v1.0.0/go.mod
h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4=
-github.com/prometheus/client_golang v1.23.2
h1:Je96obch5RDVy3FDMndoUsjAhG5Edi49h0RJWRi/o0o=
-github.com/prometheus/client_golang v1.23.2/go.mod
h1:Tb1a6LWHB3/SPIzCoaDXI4I8UHKeFTEQ1YCr+0Gyqmg=
+github.com/prometheus/client_golang v1.24.1
h1:JnJkREXzWxUdCuPFpIWZiPispT9xVV59uiuyR2bPlnU=
+github.com/prometheus/client_golang v1.24.1/go.mod
h1:F+oSRECHg4sse5ucfYpYDeIv/hu68Zo0uoHKetWnzcE=
github.com/prometheus/client_model v0.6.2
h1:oBsgwpGs7iVziMvrGhE53c/GrLUsZdHnqNwqPLxwZyk=
github.com/prometheus/client_model v0.6.2/go.mod
h1:y3m2F6Gdpfy6Ut/GBsUqTWZqCUvMVzSfMLjcu6wAwpE=
-github.com/prometheus/common v0.67.5
h1:pIgK94WWlQt1WLwAC5j2ynLaBRDiinoAb86HZHTUGI4=
-github.com/prometheus/common v0.67.5/go.mod
h1:SjE/0MzDEEAyrdr5Gqc6G+sXI67maCxzaT3A2+HqjUw=
-github.com/prometheus/procfs v0.16.1
h1:hZ15bTNuirocR6u0JZ6BAHHmwS1p8B4P6MRqxtzMyRg=
-github.com/prometheus/procfs v0.16.1/go.mod
h1:teAbpZRB1iIAJYREa1LsoWUXykVXA1KlTmWl8x/U+Is=
-github.com/rogpeppe/go-internal v1.10.0
h1:TMyTOH3F/DB16zRVcYyreMH6GnZZrwQVAoYjRBZyWFQ=
-github.com/rogpeppe/go-internal v1.10.0/go.mod
h1:UQnix2H7Ngw/k4C5ijL5+65zddjncjaFoBhdsK/akog=
+github.com/prometheus/common v0.70.1
h1:1HvjP4D5oL3t8RsPlwxA9onvvStjtIHYE5XuuwOi/PY=
+github.com/prometheus/common v0.70.1/go.mod
h1:VdFUQDMZK3VLkurFUVhia6uys/0suUp86TJz5qbJRhc=
+github.com/prometheus/procfs v0.21.1
h1:GljZCt+zSTS+NZq88cyQ1LjZ+RCHp3uVuabBWA5+OJI=
+github.com/prometheus/procfs v0.21.1/go.mod
h1:aB55Cww9pdSJVHk0hUf0inxWyyjPogFIjmHKYgMKmtY=
github.com/sasha-s/go-deadlock v0.3.9
h1:fiaT9rB7g5sr5ddNZvlwheclN9IP86eFW9WgqlEQV+w=
github.com/sasha-s/go-deadlock v0.3.9/go.mod
h1:KuZj51ZFmx42q/mPaYbRk0P1xcwe697zsJKE03vD4/Y=
github.com/stretchr/testify v1.11.1
h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu7U=
@@ -74,30 +69,30 @@ github.com/tidwall/btree v1.8.1
h1:27ehoXvm5AG/g+1VxLS1SD3vRhp/H7LuEfwNvddEdmA=
github.com/tidwall/btree v1.8.1/go.mod
h1:jBbTdUWhSZClZWoDg54VnvV7/54modSOzDN7VXftj1A=
go.opentelemetry.io/auto/sdk v1.2.1
h1:jXsnJ4Lmnqd11kwkBV2LgLoFMZKizbCi5fNZ/ipaZ64=
go.opentelemetry.io/auto/sdk v1.2.1/go.mod
h1:KRTj+aOaElaLi+wW1kO/DZRXwkF4C5xPbEe3ZiIhN7Y=
-go.opentelemetry.io/otel v1.43.0
h1:mYIM03dnh5zfN7HautFE4ieIig9amkNANT+xcVxAj9I=
-go.opentelemetry.io/otel v1.43.0/go.mod
h1:JuG+u74mvjvcm8vj8pI5XiHy1zDeoCS2LB1spIq7Ay0=
-go.opentelemetry.io/otel/metric v1.43.0
h1:d7638QeInOnuwOONPp4JAOGfbCEpYb+K6DVWvdxGzgM=
-go.opentelemetry.io/otel/metric v1.43.0/go.mod
h1:RDnPtIxvqlgO8GRW18W6Z/4P462ldprJtfxHxyKd2PY=
-go.opentelemetry.io/otel/sdk v1.43.0
h1:pi5mE86i5rTeLXqoF/hhiBtUNcrAGHLKQdhg4h4V9Dg=
-go.opentelemetry.io/otel/sdk v1.43.0/go.mod
h1:P+IkVU3iWukmiit/Yf9AWvpyRDlUeBaRg6Y+C58QHzg=
-go.opentelemetry.io/otel/sdk/metric v1.43.0
h1:S88dyqXjJkuBNLeMcVPRFXpRw2fuwdvfCGLEo89fDkw=
-go.opentelemetry.io/otel/sdk/metric v1.43.0/go.mod
h1:C/RJtwSEJ5hzTiUz5pXF1kILHStzb9zFlIEe85bhj6A=
-go.opentelemetry.io/otel/trace v1.43.0
h1:BkNrHpup+4k4w+ZZ86CZoHHEkohws8AY+WTX09nk+3A=
-go.opentelemetry.io/otel/trace v1.43.0/go.mod
h1:/QJhyVBUUswCphDVxq+8mld+AvhXZLhe+8WVFxiFff0=
+go.opentelemetry.io/otel v1.44.0
h1:JjwHmHpA4iZ3wBxluu2fbbE7j4kqlE8jXyAyPXH7HqU=
+go.opentelemetry.io/otel v1.44.0/go.mod
h1:BMgjTHL9WPRlRjL2oZCBTL4whCGtXch2H4BhOPIAyYc=
+go.opentelemetry.io/otel/metric v1.44.0
h1:1w0gILTcHdr3YI+ixLyjemwrVnsMURbTZFrSYCdDdmc=
+go.opentelemetry.io/otel/metric v1.44.0/go.mod
h1:8O7hanEPBNgEMmybD3s2VBKcgWOCsA6tzHBPODAiquo=
+go.opentelemetry.io/otel/sdk v1.44.0
h1:nHYwb9lK+fJPU/dnT6s7W7Z8itMWyqrnVfbheVYrZ58=
+go.opentelemetry.io/otel/sdk v1.44.0/go.mod
h1:Osuydd3Se74nqjAKxid74N5eC+jfEqfTegHRnq58oK0=
+go.opentelemetry.io/otel/sdk/metric v1.44.0
h1:3LlKgI+VjbVsjNRFZJZAJ30WjXC5VkNRks6si09iEfI=
+go.opentelemetry.io/otel/sdk/metric v1.44.0/go.mod
h1:5B5pMARnXxKhltooO4xUuCBorl65a4EpnTalObqOigA=
+go.opentelemetry.io/otel/trace v1.44.0
h1:jxF5CsGYCe74MCRx2X4g7WsY/VBKRqqpNvXlX/6gtIk=
+go.opentelemetry.io/otel/trace v1.44.0/go.mod
h1:oLl1jrMQAVo6v3GAggN+1VH9VIz9iUSvW53sW1Q8PIE=
go.uber.org/goleak v1.3.0 h1:2K3zAYmnTNqV73imy9J1T3WC+gmCePx2hEGkimedGto=
go.uber.org/goleak v1.3.0/go.mod
h1:CoHD4mav9JJNrW/WLlf7HGZPjdw8EucARQHekz1X6bE=
-go.uber.org/multierr v1.10.0 h1:S0h4aNzvfcFsC3dRF1jLoaov7oRaKqRGC/pUEJ2yvPQ=
-go.uber.org/multierr v1.10.0/go.mod
h1:20+QtiLqy0Nd6FdQB9TLXag12DsQkrbs3htMFfDN80Y=
-go.uber.org/zap v1.27.1 h1:08RqriUEv8+ArZRYSTXy1LeBScaMpVSTBhCeaZYfMYc=
-go.uber.org/zap v1.27.1/go.mod h1:GB2qFLM7cTU87MWRP2mPIjqfIDnGu+VIO4V/SdhGo2E=
-go.yaml.in/yaml/v2 v2.4.3 h1:6gvOSjQoTB3vt1l+CU+tSyi/HOjfOjRLJ4YwYZGwRO0=
-go.yaml.in/yaml/v2 v2.4.3/go.mod
h1:zSxWcmIDjOzPXpjlTTbAsKokqkDNAVtZO0WOMiT90s8=
-go.yaml.in/yaml/v3 v3.0.4 h1:tfq32ie2Jv2UxXFdLJdh3jXuOzWiL1fo0bu/FbuKpbc=
-go.yaml.in/yaml/v3 v3.0.4/go.mod
h1:DhzuOOF2ATzADvBadXxruRBLzYTpT36CKvDb3+aBEFg=
+go.uber.org/multierr v1.11.0 h1:blXXJkSxSSfBVBlC76pxqeO+LN3aDfLQo+309xJstO0=
+go.uber.org/multierr v1.11.0/go.mod
h1:20+QtiLqy0Nd6FdQB9TLXag12DsQkrbs3htMFfDN80Y=
+go.uber.org/zap v1.28.0 h1:IZzaP1Fv73/T/pBMLk4VutPl36uNC+OSUh3JLG3FIjo=
+go.uber.org/zap v1.28.0/go.mod h1:rDLpOi171uODNm/mxFcuYWxDsqWSAVkFdX4XojSKg/Q=
+go.yaml.in/yaml/v2 v2.4.4 h1:tuyd0P+2Ont/d6e2rl3be67goVK4R6deVxCUX5vyPaQ=
+go.yaml.in/yaml/v2 v2.4.4/go.mod
h1:gMZqIpDtDqOfM0uNfy0SkpRhvUryYH0Z6wdMYcacYXQ=
+go.yaml.in/yaml/v3 v3.0.5 h1:N6y/pJk8buWs9NY5ERU2HSMfm+IuD/OtfdAnq6kESPw=
+go.yaml.in/yaml/v3 v3.0.5/go.mod
h1:HVTZu1O7/Vkt2N+BFy8Zza+lnLsABggaTM2ZpNIGuKg=
golang.org/x/crypto v0.52.0 h1:RMs7fP2rXdep0CftQlK8Uf+kibLm7qkCcradZWYz988=
golang.org/x/crypto v0.52.0/go.mod
h1:1QgfPxDqh0T2M/elOJtp9RvuR95kVjir0e6/BvEmGbc=
-golang.org/x/exp v0.0.0-20260312153236-7ab1446f8b90
h1:jiDhWWeC7jfWqR9c/uplMOqJ0sbNlNWv0UkzE0vX1MA=
-golang.org/x/exp v0.0.0-20260312153236-7ab1446f8b90/go.mod
h1:xE1HEv6b+1SCZ5/uscMRjUBKtIxworgEcEi+/n9NQDQ=
+golang.org/x/exp v0.0.0-20260727155853-b88d891fe743
h1:ex206bKw+v3K0dm3andkrIF+ijyQKJG1pLgwQ2PYdQM=
+golang.org/x/exp v0.0.0-20260727155853-b88d891fe743/go.mod
h1:EdfpwwqSu+0Li0mzskwHU6FWDV3t9Q+RZDo3QMUtL3Q=
golang.org/x/net v0.55.0 h1:bcvxaJn3e1U6InsFWt1JUq1aSjnRxLzT2rtD2KfkDF8=
golang.org/x/net v0.55.0/go.mod h1:L5U2KuzuOe1lY7Z+aWVIKK6qEeJXnXV9yzGA+WCHJww=
golang.org/x/sys v0.45.0 h1:dO4czNzziLiiXplLQgBCEpCvXQ3dnkn0SdaZSYdQ+FY=
@@ -108,15 +103,12 @@ golang.org/x/time v0.15.0
h1:bbrp8t3bGUeFOx08pvsMYRTCVSMk89u4tKbNOZbp88U=
golang.org/x/time v0.15.0/go.mod
h1:Y4YMaQmXwGQZoFaVFk4YpCt4FLQMYKZe9oeV/f4MSno=
gonum.org/v1/gonum v0.17.0 h1:VbpOemQlsSMrYmn7T2OUvQ4dqxQXU+ouZFQsZOx50z4=
gonum.org/v1/gonum v0.17.0/go.mod
h1:El3tOrEuMpv2UdMrbNlKEh9vd86bmQ6vqIcDwxEOc1E=
-google.golang.org/genproto/googleapis/rpc v0.0.0-20260414002931-afd174a4e478
h1:RmoJA1ujG+/lRGNfUnOMfhCy5EipVMyvUE+KNbPbTlw=
-google.golang.org/genproto/googleapis/rpc
v0.0.0-20260414002931-afd174a4e478/go.mod
h1:4Hqkh8ycfw05ld/3BWL7rJOSfebL2Q+DVDeRgYgxUU8=
-google.golang.org/grpc v1.82.1 h1:NnAxzGRA0677vCa4BUkOAnO5+FfQqVl9iUXeD0IqcGE=
-google.golang.org/grpc v1.82.1/go.mod
h1:yzTZ1TB1Z3SG+LIYaI+WiE8D5+PZ3ArnrSp8zF3+/ZA=
+google.golang.org/genproto/googleapis/rpc v0.0.0-20260803160001-6ac0973c030d
h1:IL4hdHzcUv2l/gcg98/Rj3FbtE6axwqslOW8SW0C+S0=
+google.golang.org/genproto/googleapis/rpc
v0.0.0-20260803160001-6ac0973c030d/go.mod
h1:4Hqkh8ycfw05ld/3BWL7rJOSfebL2Q+DVDeRgYgxUU8=
+google.golang.org/grpc v1.83.0 h1:JeNZEKJFbQxArAMl+hiytHauacDNqJUllNfmIMmpqnQ=
+google.golang.org/grpc v1.83.0/go.mod
h1:kDyl6SKsiHKt0uylY5gtn5cEjkrIOhQOGDgIc4JGwzQ=
google.golang.org/protobuf v1.36.11
h1:fV6ZwhNocDyBLK0dj+fg8ektcVegBBuEolpbTQyBNVE=
google.golang.org/protobuf v1.36.11/go.mod
h1:HTf+CrKn2C3g5S8VImy6tdcUvCska2kB7j23XfzDpco=
-gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod
h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0=
-gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c
h1:Hei/4ADfdWqJk1ZMxUNpqntNwaWcugrBjAiHlqqRiVk=
-gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c/go.mod
h1:JHkPIbrfpd72SG/EVd6muEfDQjcINNoR0C8j2r3qZ4Q=
gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA=
gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM=
gotest.tools/v3 v3.5.2 h1:7koQfIKdy+I8UTetycgUqXWSDwpgv193Ka+qRsmBY8Q=
diff --git a/pkg/common/errors.go b/pkg/common/errors.go
index 210c4089..fc7bd167 100644
--- a/pkg/common/errors.go
+++ b/pkg/common/errors.go
@@ -30,7 +30,8 @@ var (
// ErrorNodeAlreadyReserved returned when the node is already reserved,
failing the reservation
ErrorNodeAlreadyReserved = errors.New("node is already reserved")
// ErrorNodeNotFitReserve returned when the allocation does not fit on
an empty node, failing the reservation
- ErrorNodeNotFitReserve = errors.New("reservation does not fit on node")
+ ErrorNodeNotFitReserve = errors.New("reservation does not fit on node")
+ ErrorPreFilterPredicate = errors.New("predicate checks failed. Unable
to find the suitable node to run this pod")
)
// Constant messages for AllocationLog entries
diff --git a/pkg/mock/predicate_plugin.go b/pkg/mock/predicate_plugin.go
index 27aa3822..1f56c244 100644
--- a/pkg/mock/predicate_plugin.go
+++ b/pkg/mock/predicate_plugin.go
@@ -29,41 +29,63 @@ import (
type PredicatePlugin struct {
ResourceManagerCallback
- mustFail bool
- nodes map[string]int
+ mustPreFilterFail bool
+ mustFilterFail bool
+ nodes map[string]int
}
-func (f *PredicatePlugin) Predicates(args *si.PredicatesArgs) error {
- if f.mustFail {
- log.Log(log.Test).Info("fake predicate plugin fail: must fail
set")
- return fmt.Errorf("fake predicate plugin failed")
+func (f *PredicatePlugin) PreFilterPredicates(args
*si.PreFilterPredicatesArgs) *si.PreFilterPredicatesResponse {
+ feasibleNodes := make(map[string]*si.Empty)
+ result := &si.PreFilterPredicatesResponse{
+ Success: false,
+ FeasibleNodes: map[string]*si.Empty{},
}
- if fail, ok := f.nodes[args.NodeID]; ok {
- if args.Allocate && fail >= 0 {
- log.Log(log.Test).Info("fake predicate plugin node
allocate fail",
- zap.String("node", args.NodeID),
- zap.Int("fail mode", fail))
- return fmt.Errorf("fake predicate plugin failed")
- }
- if !args.Allocate && fail <= 0 {
- log.Log(log.Test).Info("fake predicate plugin node
reserve fail",
- zap.String("node", args.NodeID),
- zap.Int("fail mode", fail))
- return fmt.Errorf("fake predicate plugin failed")
+ if f.mustPreFilterFail {
+ log.Log(log.Test).Info("fake predicate prefilter plugin fail:
must fail set")
+ return result
+ }
+ for k, v := range f.nodes {
+ if args.Allocate {
+ if v < 0 {
+ feasibleNodes[k] = &si.Empty{}
+ }
+ } else {
+ if v > 0 {
+ feasibleNodes[k] = &si.Empty{}
+ }
}
}
- log.Log(log.Test).Info("fake predicate plugin pass",
+ log.Log(log.Test).Info("fake predicate prefilter plugin pass",
+ zap.Bool("allocate", args.Allocate),
+ zap.String("allocationKey", args.AllocationKey),
+ zap.Any("feasibleNodes", feasibleNodes),
+ zap.Int("feasibleNodes", len(feasibleNodes)))
+ result.Success = true
+ result.FeasibleNodes = feasibleNodes
+ return result
+}
+
+func (f *PredicatePlugin) Predicates(args *si.PredicatesArgs) error {
+ if f.mustFilterFail {
+ log.Log(log.Test).Info("fake predicate filter plugin fail: must
fail set")
+ return fmt.Errorf("fake predicate plugin failed")
+ }
+ log.Log(log.Test).Info("fake predicate filter plugin pass",
+ zap.Bool("allocate", args.Allocate),
+ zap.String("allocationKey", args.AllocationKey),
zap.String("node", args.NodeID))
return nil
}
// NewPredicatePlugin returns a mock that can either always fail or fail based
on the node that is checked.
-// mustFail will cause the predicate check to always fail
-// nodes allows specifying which node to fail for which check using the nodeID:
-// possible values: -1 fail reserve, 0 fail always, 1 fail alloc (defaults to
always)
-func NewPredicatePlugin(mustFail bool, nodes map[string]int) *PredicatePlugin {
+// mustPreFilterFail will cause the predicate prefilter check to fail always
+// mustFilterFail will cause the predicate filter check to fail always
+// nodes allows specifying which node to make it to feasibleNodes list based
on its own value:
+// possible values: 1 - Feasible in case of reserve, -1 Feasible in case of
allow, 0 Not feasible always
+func NewPredicatePlugin(mustPreFilterFail bool, mustFilterFail bool, nodes
map[string]int) *PredicatePlugin {
return &PredicatePlugin{
- mustFail: mustFail,
- nodes: nodes,
+ mustPreFilterFail: mustPreFilterFail,
+ mustFilterFail: mustFilterFail,
+ nodes: nodes,
}
}
diff --git a/pkg/mock/preemption_predicate_plugin.go
b/pkg/mock/preemption_predicate_plugin.go
index d41b9ac9..407a7f4e 100644
--- a/pkg/mock/preemption_predicate_plugin.go
+++ b/pkg/mock/preemption_predicate_plugin.go
@@ -19,20 +19,19 @@
package mock
import (
- "errors"
"fmt"
"github.com/apache/yunikorn-core/pkg/locking"
-
"github.com/apache/yunikorn-scheduler-interface/lib/go/si"
)
type PreemptionPredicatePlugin struct {
ResourceManagerCallback
- reservations map[string]string
- allocations map[string]string
- preemptions []Preemption
- errHolder *errHolder
+ mustPreFilterFail bool
+ mustFilterFail bool
+ preemptions []Preemption
+ errHolder *errHolder
+ nodes map[string]int
locking.RWMutex
}
@@ -50,28 +49,37 @@ type errHolder struct {
err error
}
-func (m *PreemptionPredicatePlugin) Predicates(args *si.PredicatesArgs) error {
+func (m *PreemptionPredicatePlugin) PreFilterPredicates(args
*si.PreFilterPredicatesArgs) *si.PreFilterPredicatesResponse {
m.RLock()
defer m.RUnlock()
- if args.Allocate {
- nodeID, ok := m.allocations[args.AllocationKey]
- if !ok {
- return errors.New("no allocation found")
- }
- if nodeID != args.NodeID {
- return errors.New("wrong node")
- }
- return nil
- } else {
- nodeID, ok := m.reservations[args.AllocationKey]
- if !ok {
- return errors.New("no allocation found")
- }
- if nodeID != args.NodeID {
- return errors.New("wrong node")
+ feasibleNodes := make(map[string]*si.Empty)
+ result := &si.PreFilterPredicatesResponse{
+ Success: false,
+ FeasibleNodes: map[string]*si.Empty{},
+ }
+
+ if m.mustPreFilterFail {
+ m.errHolder.err = fmt.Errorf("fake preemption predicate
prefilter plugin failed")
+ return result
+ }
+ for k, v := range m.nodes {
+ if v > 0 {
+ feasibleNodes[k] = &si.Empty{}
}
- return nil
}
+ result.Success = true
+ result.FeasibleNodes = feasibleNodes
+ return result
+}
+
+func (m *PreemptionPredicatePlugin) Predicates(args *si.PredicatesArgs) error {
+ m.RLock()
+ defer m.RUnlock()
+ if m.mustFilterFail {
+ m.errHolder.err = fmt.Errorf("fake predicate filter plugin
failed")
+ return m.errHolder.err
+ }
+ return nil
}
func (m *PreemptionPredicatePlugin) PreemptionPredicates(args
*si.PreemptionPredicatesArgs) *si.PreemptionPredicatesResponse {
@@ -81,6 +89,10 @@ func (m *PreemptionPredicatePlugin)
PreemptionPredicates(args *si.PreemptionPred
Success: false,
Index: -1,
}
+ if m.mustFilterFail {
+ m.errHolder.err = fmt.Errorf("fake preemption predicate filter
plugin failed")
+ return result
+ }
for _, preemption := range m.preemptions {
if preemption.expectedAllocationKey != args.AllocationKey {
continue
@@ -123,15 +135,18 @@ func (m *PreemptionPredicatePlugin) GetPredicateError()
error {
}
// NewPreemptionPredicatePlugin returns a mock plugin that can handle multiple
predicate scenarios.
-// reservations: provide a list of allocations and node IDs for which the
reservation predicate succeeds
-// allocs: provide a list of allocations and node IDs for which the allocation
predicate succeeds
// preempt: a slice of preemption scenarios configured for the plugin to check
-func NewPreemptionPredicatePlugin(reservations, allocs map[string]string,
preempt []Preemption) *PreemptionPredicatePlugin {
+// mustPreFilterFail will cause the predicate prefilter check to fail always
+// mustFilterFail will cause the predicate filter check to fail always
+// nodes allows specifying which node to make it to feasibleNodes list based
on its own value:
+// possible values: 1 - Feasible always, 0 or -1 - Not feasible always
+func NewPreemptionPredicatePlugin(preempt []Preemption, nodes map[string]int,
mustPreFilterFail, mustFilterFail bool) *PreemptionPredicatePlugin {
return &PreemptionPredicatePlugin{
- reservations: reservations,
- allocations: allocs,
- preemptions: preempt,
- errHolder: &errHolder{},
+ preemptions: preempt,
+ errHolder: &errHolder{},
+ nodes: nodes,
+ mustPreFilterFail: mustPreFilterFail,
+ mustFilterFail: mustFilterFail,
}
}
diff --git a/pkg/mock/rm_callback.go b/pkg/mock/rm_callback.go
index ccaab4a2..9adad7ce 100644
--- a/pkg/mock/rm_callback.go
+++ b/pkg/mock/rm_callback.go
@@ -41,6 +41,14 @@ func (f *ResourceManagerCallback) Predicates(_
*si.PredicatesArgs) error {
return nil
}
+func (f *ResourceManagerCallback) PreFilterPredicates(_
*si.PreFilterPredicatesArgs) *si.PreFilterPredicatesResponse {
+ // simulate "ideal" preemption check
+ return &si.PreFilterPredicatesResponse{
+ Success: true,
+ FeasibleNodes: map[string]*si.Empty{},
+ }
+}
+
func (f *ResourceManagerCallback) PreemptionPredicates(args
*si.PreemptionPredicatesArgs) *si.PreemptionPredicatesResponse {
// simulate "ideal" preemption check
return &si.PreemptionPredicatesResponse{
diff --git a/pkg/plugins/plugins_test.go b/pkg/plugins/plugins_test.go
index f992f8bc..4f318573 100644
--- a/pkg/plugins/plugins_test.go
+++ b/pkg/plugins/plugins_test.go
@@ -53,6 +53,10 @@ func (f *RMPluginImplemented) Predicates(_
*si.PredicatesArgs) error {
return nil
}
+func (f *RMPluginImplemented) PreFilterPredicates(_
*si.PreFilterPredicatesArgs) *si.PreFilterPredicatesResponse {
+ return nil
+}
+
func (f *RMPluginImplemented) PreemptionPredicates(_
*si.PreemptionPredicatesArgs) *si.PreemptionPredicatesResponse {
return nil
}
diff --git a/pkg/scheduler/objects/allocation.go
b/pkg/scheduler/objects/allocation.go
index 92d2b900..71111e22 100644
--- a/pkg/scheduler/objects/allocation.go
+++ b/pkg/scheduler/objects/allocation.go
@@ -31,6 +31,7 @@ import (
"github.com/apache/yunikorn-core/pkg/events"
"github.com/apache/yunikorn-core/pkg/locking"
"github.com/apache/yunikorn-core/pkg/log"
+ "github.com/apache/yunikorn-core/pkg/plugins"
schedEvt "github.com/apache/yunikorn-core/pkg/scheduler/objects/events"
siCommon "github.com/apache/yunikorn-scheduler-interface/lib/go/common"
"github.com/apache/yunikorn-scheduler-interface/lib/go/si"
@@ -620,3 +621,24 @@ func (a *Allocation) IsPreemptable() bool {
func (a *Allocation) GetAllocationName() string {
return a.tags[siCommon.DomainYuniKorn+siCommon.KeyPodName]
}
+
+func (a *Allocation) preAllocateConditions(allocate bool)
(map[string]*si.Empty, bool) {
+ var prefilterResult *si.PreFilterPredicatesResponse
+ if plugin := plugins.GetResourceManagerCallbackPlugin(); plugin != nil {
+ if prefilterResult =
plugin.PreFilterPredicates(&si.PreFilterPredicatesArgs{
+ AllocationKey: a.allocationKey,
+ Allocate: allocate,
+ }); prefilterResult != nil && !prefilterResult.Success {
+ log.Log(log.SchedNode).Debug("running prefilter
predicates failed",
+ zap.String("allocationKey", a.allocationKey),
+ zap.Bool("allocate", allocate))
+
a.LogAllocationFailure(common.ErrorPreFilterPredicate.Error(), allocate)
+ podPredicateErrors := make(map[string]int, 1)
+
podPredicateErrors[common.ErrorPreFilterPredicate.Error()]++
+ a.SendPredicatesFailedEvent(podPredicateErrors)
+ return prefilterResult.GetFeasibleNodes(), false
+ }
+ }
+ // all predicate plugins passed
+ return prefilterResult.GetFeasibleNodes(), true
+}
diff --git a/pkg/scheduler/objects/application.go
b/pkg/scheduler/objects/application.go
index 71decadc..6d7bd6fd 100644
--- a/pkg/scheduler/objects/application.go
+++ b/pkg/scheduler/objects/application.go
@@ -1209,6 +1209,7 @@ func (sa *Application) tryAllocate(headRoom
*resources.Resource, allowPreemption
request.setHeadroomCheckPassed(sa.queuePath)
requiredNode := request.GetRequiredNode()
+
// does request have any constraint to run on specific node?
if requiredNode != "" {
result := sa.tryRequiredNode(request, getNodeFn)
@@ -1268,8 +1269,7 @@ func (sa *Application) tryRequiredNode(request
*Allocation, getNodeFn func(strin
num = sa.cancelReservations(reservations)
}
_, thisReserved := sa.reservations[allocationKey]
- // now try the request, we don't care about predicate error messages
here
- result, _ := sa.tryNode(node, request) //nolint:errcheck
+ result, _ := sa.tryNode(node, request, false) //nolint:errcheck
if result != nil {
result.CancelledReservations = num
// check if the node was reserved and we allocated after a
release
@@ -1443,16 +1443,24 @@ func (sa *Application)
tryPlaceholderAllocate(nodeIterator func() NodeIterator,
}
}
}
+
// cannot allocate if the iterator is not giving us any schedulable
nodes
iterator := nodeIterator()
if iterator == nil {
return nil
}
+
// we checked all placeholders and asks nothing worked as yet
// pick the first fit and try all nodes if that fails give up
var allocResult *AllocationResult
if phFit != nil && reqFit != nil {
resKey := reqFit.GetAllocationKey()
+
+ // run predicates for this pod before in hand and fetch
feasible nodes
+ feasibleNodes, predicatesResult :=
reqFit.preAllocateConditions(true)
+ if !predicatesResult {
+ return nil
+ }
iterator.ForEachNode(func(node *Node) bool {
if !node.IsSchedulable() {
log.Log(log.SchedApplication).Debug("skipping
node for placeholder alloc as state is unschedulable",
@@ -1463,10 +1471,21 @@ func (sa *Application)
tryPlaceholderAllocate(nodeIterator func() NodeIterator,
if
!node.preAllocateCheck(reqFit.GetAllocatedResource(), resKey) {
return true
}
+ // Is this node suitable to run the pod?
+ if len(feasibleNodes) > 0 {
+ if _, ok := feasibleNodes[node.NodeID]; !ok {
+ getRateLimitedAppLog().Info("skipping
node as it is not feasible to run the pod",
+ zap.String("allocationKey",
resKey),
+ zap.String("node", node.NodeID))
+ return true
+ }
+ }
+
// skip the node if conditions can not be satisfied
if err := node.preAllocateConditions(reqFit); err !=
nil {
return true
}
+
// update just the node to make sure we keep its spot
// no queue update as we're releasing the placeholder
and are just temp over the size
if !node.TryAddAllocation(reqFit) {
@@ -1589,8 +1608,25 @@ func (sa *Application) tryReservedAllocate(headRoom
*resources.Resource, nodeIte
}
}
// check allocation possibility
- // we don't care about predicate error messages here
- result, _ := sa.tryNode(reserve.node, ask) //nolint:errcheck
+ skipReservedNode := false
+ feasibleNodes, predicatesResult :=
ask.preAllocateConditions(true)
+ if predicatesResult {
+ // Is this node suitable to run the pod?
+ if len(feasibleNodes) > 0 {
+ if _, ok := feasibleNodes[reserve.node.NodeID];
!ok {
+ skipReservedNode = true
+ }
+ }
+ } else {
+ skipReservedNode = true
+ }
+ if skipReservedNode {
+ getRateLimitedAppLog().Info("skipping reserved node as
it is not feasible to run the pod",
+ zap.String("allocationKey",
ask.GetAllocationKey()),
+ zap.String("reserved node",
reserve.node.NodeID))
+ continue
+ }
+ result, _ := sa.tryNode(reserve.node, ask, true)
//nolint:errcheck
// allocation worked fix the resultType and return
if result != nil {
@@ -1650,6 +1686,13 @@ func (sa *Application) tryPreemption(headRoom
*resources.Resource, preemptionDel
// This should never result in a reservation as the allocation is already
reserved
func (sa *Application) tryNodesNoReserve(ask *Allocation, iterator
NodeIterator, reservedNode string) *AllocationResult {
var allocResult *AllocationResult
+
+ // run predicates for this pod before in hand and fetch feasible nodes
+ feasibleNodes, predicatesResult := ask.preAllocateConditions(true)
+ if !predicatesResult {
+ return nil
+ }
+
iterator.ForEachNode(func(node *Node) bool {
if !node.IsSchedulable() {
log.Log(log.SchedApplication).Debug("skipping node for
reserved ask as state is unschedulable",
@@ -1661,8 +1704,17 @@ func (sa *Application) tryNodesNoReserve(ask
*Allocation, iterator NodeIterator,
if !node.FitInNode(ask.GetAllocatedResource()) || node.NodeID
== reservedNode {
return true
}
- // we don't care about predicate error messages here
- result, _ := sa.tryNode(node, ask) //nolint:errcheck
+
+ // Is this node suitable to run the pod?
+ if len(feasibleNodes) > 0 {
+ if _, ok := feasibleNodes[node.NodeID]; !ok {
+ getRateLimitedAppLog().Info("skipping node as
it is not feasible to run the pod",
+ zap.String("allocationKey",
ask.GetAllocationKey()),
+ zap.String("node", node.NodeID))
+ return true
+ }
+ }
+ result, _ := sa.tryNode(node, ask, true) //nolint:errcheck
// allocation worked: update resultType and return
if result != nil {
result.ResultType = AllocatedReserved
@@ -1688,6 +1740,10 @@ func (sa *Application) tryNodes(ask *Allocation,
iterator NodeIterator) *Allocat
var allocResult *AllocationResult
var predicateErrors map[string]int
tryNodeCycleStart := time.Now()
+
+ // run predicates for this pod before in hand and fetch feasible nodes
+ feasibleNodes, predicatesResult := ask.preAllocateConditions(true)
+
iterator.ForEachNode(func(node *Node) bool {
// skip the node if the node is not schedulable
if !node.IsSchedulable() {
@@ -1700,41 +1756,56 @@ func (sa *Application) tryNodes(ask *Allocation,
iterator NodeIterator) *Allocat
if !node.FitInNode(ask.GetAllocatedResource()) {
return true
}
- tryNodeStart := time.Now()
- result, err := sa.tryNode(node, ask)
- if err != nil {
- if predicateErrors == nil {
- predicateErrors = make(map[string]int)
+
+ // Is there any prefilter predicate failures? No node would be
picked up for allocation in case of any failures in predicates results
+ // and better to get into process of picking up a node for
reservation right away
+ if predicatesResult {
+ // Is this node suitable to run the pod?
+ if len(feasibleNodes) > 0 {
+ if _, ok := feasibleNodes[node.NodeID]; !ok {
+ getRateLimitedAppLog().Info("skipping
node as it is not feasible to run the pod",
+ zap.String("allocationKey",
allocKey),
+ zap.String("node", node.NodeID))
+ return true
+ }
}
- predicateErrors[err.Error()]++
- }
- // allocation worked so return
- if result != nil {
-
metrics.GetSchedulerMetrics().ObserveTryNodeLatency(tryNodeStart)
- // check if the alloc had a reservation: if it has set
the resultType and return
- if reserved != nil {
- if reserved.nodeID != node.NodeID {
- // we have a different node reserved
for this alloc
-
log.Log(log.SchedApplication).Debug("allocate picking reserved alloc during non
reserved allocate",
- zap.String("appID",
sa.ApplicationID),
- zap.String("reserved nodeID",
reserved.nodeID),
- zap.String("allocationKey",
allocKey))
- result.ReservedNodeID = reserved.nodeID
- } else {
- // NOTE: this is a safeguard as
reserved nodes should never be part of the iterator
-
log.Log(log.SchedApplication).Debug("allocate found reserved alloc during non
reserved allocate",
- zap.String("appID",
sa.ApplicationID),
- zap.String("nodeID",
node.NodeID),
- zap.String("allocationKey",
allocKey))
+ tryNodeStart := time.Now()
+ result, err := sa.tryNode(node, ask, true)
+ if err != nil {
+ if predicateErrors == nil {
+ predicateErrors = make(map[string]int)
+ }
+ predicateErrors[err.Error()]++
+ }
+ // allocation worked so return
+ if result != nil {
+
metrics.GetSchedulerMetrics().ObserveTryNodeLatency(tryNodeStart)
+ // check if the alloc had a reservation: if it
has set the resultType and return
+ if reserved != nil {
+ if reserved.nodeID != node.NodeID {
+ // we have a different node
reserved for this alloc
+
log.Log(log.SchedApplication).Debug("allocate picking reserved alloc during non
reserved allocate",
+ zap.String("appID",
sa.ApplicationID),
+ zap.String("reserved
nodeID", reserved.nodeID),
+
zap.String("allocationKey", allocKey))
+ result.ReservedNodeID =
reserved.nodeID
+ } else {
+ // NOTE: this is a safeguard as
reserved nodes should never be part of the iterator
+
log.Log(log.SchedApplication).Debug("allocate found reserved alloc during non
reserved allocate",
+ zap.String("appID",
sa.ApplicationID),
+ zap.String("nodeID",
node.NodeID),
+
zap.String("allocationKey", allocKey))
+ }
+ result.ResultType = AllocatedReserved
+ allocResult = result
+ return false
}
- result.ResultType = AllocatedReserved
+ // nothing reserved just return this as a
normal alloc
allocResult = result
return false
}
- // nothing reserved just return this as a normal alloc
- allocResult = result
- return false
}
+
// nothing allocated should we look at a reservation?
askAge := time.Since(ask.GetCreateTime())
if reserved == nil && askAge > reservationDelay {
@@ -1782,7 +1853,7 @@ func (sa *Application) tryNodes(ask *Allocation, iterator
NodeIterator) *Allocat
}
// tryNode tries allocating on one specific node
-func (sa *Application) tryNode(node *Node, ask *Allocation)
(*AllocationResult, error) {
+func (sa *Application) tryNode(node *Node, ask *Allocation, doPredicateChecks
bool) (*AllocationResult, error) {
toAllocate := ask.GetAllocatedResource()
allocationKey := ask.GetAllocationKey()
// create the key for the reservation
@@ -1790,11 +1861,12 @@ func (sa *Application) tryNode(node *Node, ask
*Allocation) (*AllocationResult,
// skip schedule onto node
return nil, nil
}
- // skip the node if conditions can not be satisfied
- if err := node.preAllocateConditions(ask); err != nil {
- return nil, err
+ if doPredicateChecks {
+ // skip the node if conditions can not be satisfied
+ if err := node.preAllocateConditions(ask); err != nil {
+ return nil, err
+ }
}
-
// everything OK really allocate
if node.TryAddAllocation(ask) {
if err :=
sa.queue.TryIncAllocatedResource(ask.GetAllocatedResource()); err != nil {
diff --git a/pkg/scheduler/objects/application_test.go
b/pkg/scheduler/objects/application_test.go
index 37e37123..6a28ddc5 100644
--- a/pkg/scheduler/objects/application_test.go
+++ b/pkg/scheduler/objects/application_test.go
@@ -3555,40 +3555,74 @@ func TestUpdateRunnableStatus(t *testing.T) {
func TestPredicateFailedEvents(t *testing.T) {
setupUGM()
+ node := newNode("node1", map[string]resources.Quantity{"first": 20})
+ node2 := newNode("node2", map[string]resources.Quantity{"first": 20})
+ nodeMap := map[string]*Node{"node1": node, "node2": node2}
+ iterator := getNodeIteratorFn(node)
+ getNode := func(nodeID string) *Node {
+ return nodeMap[nodeID]
+ }
- res, err := resources.NewResourceFromConf(map[string]string{"first":
"1"})
+ rootQ, err := createRootQueue(map[string]string{"first": "20"})
assert.NilError(t, err)
- headroom, err :=
resources.NewResourceFromConf(map[string]string{"first": "40"})
+ childQ, err := createManagedQueue(rootQ, "child", false,
map[string]string{"first": "20"})
assert.NilError(t, err)
- ask := newAllocationAsk("alloc-0", "app-1", res)
- app := newApplication(appID1, "default", "root")
- eventSystem := mock.NewEventSystem()
- ask.askEvents = schedEvt.NewAskEvents(eventSystem)
- app.disableStateChangeEvents()
- app.resetAppEvents()
- queue, err := createRootQueue(nil)
- assert.NilError(t, err, "queue create failed")
- app.queue = queue
- sr := sortedRequests{}
- sr.insert(ask)
- app.sortedRequests = sr
+
+ app := newApplication(appID1, "default", "root.child")
+ app.SetQueue(childQ)
+ childQ.applications[appID1] = app
+
attempts := 0
+ wrongNodes := make(map[string]int, 1)
+ wrongNodes["node3"] = 0
+ rightNodes := make(map[string]int, 1)
+ rightNodes["node1"] = 1
- mockPlugin := mockCommon.NewPredicatePlugin(true, nil)
- plugins.RegisterSchedulerPlugin(mockPlugin)
- defer plugins.UnregisterSchedulerPlugins()
+ tests := []struct {
+ name string
+ mockPlugin *mockCommon.PredicatePlugin
+ allocKey string
+ expectedFailedEvents int
+ }{
+ {"prefilter pass", mockCommon.NewPredicatePlugin(false, false,
nil), "alloc-1", 0},
+ {"prefilter passes but none of the node from iterator is
available in feasible nodes", mockCommon.NewPredicatePlugin(false, false,
wrongNodes), "alloc-2", 0},
+ {"prefilter fails", mockCommon.NewPredicatePlugin(true, false,
nil), "alloc-2", 1},
+ {"prefilter pass with expected feasible nodes, filter fails",
mockCommon.NewPredicatePlugin(false, true, rightNodes), "alloc-3", 1},
+ {"both prefilter and filter passes with correct feasible
nodes", mockCommon.NewPredicatePlugin(false, false, rightNodes), "alloc-3", 0},
+ }
+ for _, tt := range tests {
+ t.Run(tt.name, func(t *testing.T) {
+ ask := newAllocationAsk(tt.allocKey, appID1,
resources.NewResourceFromMap(map[string]resources.Quantity{"first": 5}))
+ err = app.AddAllocationAsk(ask)
+ assert.NilError(t, err)
+
+ eventSystem := mock.NewEventSystem()
+ ask.askEvents = schedEvt.NewAskEvents(eventSystem)
+ app.disableStateChangeEvents()
+
+ plugins.RegisterSchedulerPlugin(tt.mockPlugin)
+ app.tryAllocate(node.GetAvailableResource(), false,
time.Second, &attempts, iterator, iterator, getNode)
+ assert.Equal(t, tt.expectedFailedEvents,
len(eventSystem.Events))
+ if tt.expectedFailedEvents > 0 {
+ for _, log := range ask.GetAllocationLog() {
+ assert.Check(t,
strings.Contains(log.Message, "failed"))
+ assert.Equal(t, log.Count, int32(1))
+ }
+ assertEventsForPredicateFailures(t,
tt.allocKey, eventSystem.Events[0])
+ }
+ plugins.UnregisterSchedulerPlugins()
+ app.resetAppEvents()
+ })
+ }
+}
- app.tryAllocate(headroom, false, time.Second, &attempts, func()
NodeIterator {
- return &testIterator{}
- }, nilNodeIterator, nilGetNode)
- assert.Equal(t, 1, len(eventSystem.Events))
- event := eventSystem.Events[0]
+func assertEventsForPredicateFailures(t *testing.T, allocKey string, event
*si.EventRecord) {
assert.Equal(t, si.EventRecord_REQUEST, event.Type)
assert.Equal(t, si.EventRecord_NONE, event.EventChangeType)
assert.Equal(t, si.EventRecord_DETAILS_NONE, event.EventChangeDetail)
assert.Equal(t, "app-1", event.ReferenceID)
- assert.Equal(t, "alloc-0", event.ObjectID)
- assert.Equal(t, "Unschedulable request 'alloc-0': fake predicate plugin
failed (2x); ", event.Message)
+ assert.Equal(t, allocKey, event.ObjectID)
+ assert.Check(t, strings.Contains(event.Message, "Unschedulable request
'"+allocKey+"':"))
}
func TestRequiredNodePreemption(t *testing.T) {
@@ -3759,15 +3793,6 @@ func TestRequiredNodePreemptionFailed(t *testing.T) {
assert.Equal(t, int32(4),
ask2.allocLog[common.NoVictimForRequiredNode].Count, "incorrect number of entry
count")
}
-type testIterator struct{}
-
-func (testIterator) ForEachNode(fn func(*Node) bool) {
- node1 := newNode(nodeID1, map[string]resources.Quantity{"first": 20})
- node2 := newNode(nodeID2, map[string]resources.Quantity{"first": 20})
- fn(node1)
- fn(node2)
-}
-
func TestGetMaxResourceFromTag(t *testing.T) {
app := newApplication(appID0, "default", "root.unknown")
testGetResourceFromTag(t, siCommon.AppTagNamespaceResourceQuota,
app.tags, app.GetMaxResource)
@@ -4046,39 +4071,63 @@ func TestTryPlaceHolderAllocateDifferentNodes(t
*testing.T) {
getNode := func(nodeID string) *Node {
return nodeMap[nodeID]
}
-
- app := newApplication(appID0, "default", "root.default")
-
queue, err := createRootQueue(nil)
assert.NilError(t, err, "queue create failed")
- app.queue = queue
- res :=
resources.NewResourceFromMap(map[string]resources.Quantity{"first": 5})
- ph := newPlaceholderAlloc(appID0, nodeID1, res, tg1)
- app.AddAllocation(ph)
- app.addPlaceholderData(ph)
- assertPlaceholderData(t, app, tg1, 1, 0, 0, res)
-
- // predicate check fails on node1 and passes on node2
- mockPlugin := mockCommon.NewPredicatePlugin(false,
map[string]int{nodeID1: 0})
- plugins.RegisterSchedulerPlugin(mockPlugin)
- defer plugins.UnregisterSchedulerPlugins()
-
- // should allocate on node2
- ask := newAllocationAsk(aKey, appID0, res)
- ask.taskGroupName = tg1
- err = app.AddAllocationAsk(ask)
- assert.NilError(t, err, "ask should have been added to app")
+ wrongNodes := make(map[string]int)
+ wrongNodes[nodeID3] = 1
+ wrongNodes[nodeID4] = -1
+ rightNodes := make(map[string]int)
+ rightNodes[nodeID1] = 0
+ rightNodes[nodeID2] = -1
+ rightNodes[nodeID3] = 1
- result := app.tryPlaceholderAllocate(iterator, getNode)
- assert.Assert(t, result != nil, "result should not be nil")
- assert.Equal(t, Replaced, result.ResultType, "result type should be
Replaced")
- assert.Equal(t, nodeID2, result.NodeID, "result should be on node2")
- assert.Equal(t, ask, result.Request, "result should contain the ask")
- assert.Equal(t, ph, result.Request.GetRelease(), "real allocation
should link to placeholder")
- assert.Equal(t, result.Request, ph.GetRelease(), "placeholder should
link to real allocation")
- // placeholder data remains unchanged until RM confirms the replacement
- assertPlaceholderData(t, app, tg1, 1, 0, 0, res)
+ tests := []struct {
+ name string
+ mockPlugin *mockCommon.PredicatePlugin
+ allocResult bool
+ expectedNode string
+ }{
+ {"prefilter pass with empty feasible nodes, so original node
(node1) itself is feasible for replacement",
mockCommon.NewPredicatePlugin(false, false, nil), true, nodeID1},
+ {"prefilter passes but none of the node from iterator is
available in feasible nodes", mockCommon.NewPredicatePlugin(false, false,
wrongNodes), false, "NA"},
+ {"prefilter fails", mockCommon.NewPredicatePlugin(true, false,
nil), false, "NA"},
+ {"prefilter pass with expected feasible nodes, filter fails",
mockCommon.NewPredicatePlugin(false, true, rightNodes), false, "NA"},
+ {"both prefilter and filter passes with correct feasible nodes.
original node (node1) is not feasible, hence node2 opted for replacement",
mockCommon.NewPredicatePlugin(false, false, rightNodes), true, nodeID2},
+ }
+ for _, tt := range tests {
+ t.Run(tt.name, func(t *testing.T) {
+ app := newApplication(appID0, "default", "root.default")
+ app.queue = queue
+
+ res :=
resources.NewResourceFromMap(map[string]resources.Quantity{"first": 5})
+ ph := newPlaceholderAlloc(appID0, nodeID1, res, tg1)
+ app.AddAllocation(ph)
+ app.addPlaceholderData(ph)
+ assertPlaceholderData(t, app, tg1, 1, 0, 0, res)
+ plugins.RegisterSchedulerPlugin(tt.mockPlugin)
+
+ ask := newAllocationAsk(aKey, appID0, res)
+ ask.taskGroupName = tg1
+ err = app.AddAllocationAsk(ask)
+ assert.NilError(t, err, "ask should have been added to
app")
+
+ result := app.tryPlaceholderAllocate(iterator, getNode)
+ if tt.allocResult {
+ assert.Assert(t, result != nil, "result should
not be nil")
+ assert.Equal(t, Replaced, result.ResultType,
"result type should be Replaced")
+ assert.Equal(t, tt.expectedNode, result.NodeID,
"result should be on node2")
+ assert.Equal(t, ask, result.Request, "result
should contain the ask")
+ assert.Equal(t, ph,
result.Request.GetRelease(), "real allocation should link to placeholder")
+ assert.Equal(t, result.Request,
ph.GetRelease(), "placeholder should link to real allocation")
+ // placeholder data remains unchanged until RM
confirms the replacement
+ assertPlaceholderData(t, app, tg1, 1, 0, 0, res)
+ } else {
+ assert.Assert(t, result == nil, "result should
be nil")
+ }
+ queue.RemoveApplication(app)
+ plugins.UnregisterSchedulerPlugins()
+ })
+ }
}
// revertTriggerPredicatePlugin is a test-local predicate plugin whose
Predicates call has the side
@@ -4163,57 +4212,72 @@ func
TestTryPlaceHolderAllocateRevertsOnPreemptedPlaceholder(t *testing.T) {
}
func TestTryNodesNoReserve(t *testing.T) {
- app := newApplication(appID0, "default", "root.default")
-
- queue, err := createRootQueue(map[string]string{"first": "5"})
+ queue, err := createRootQueue(map[string]string{"first": "100"})
assert.NilError(t, err, "queue create failed")
- app.queue = queue
- res :=
resources.NewResourceFromMap(map[string]resources.Quantity{"first": 5})
- ask := newAllocationAsk(aKey, appID0, res)
- err = app.AddAllocationAsk(ask)
- assert.NilError(t, err, "ask should have been added to app")
+ wrongNodes := make(map[string]int)
+ wrongNodes[nodeID3] = 1
+ wrongNodes[nodeID4] = -1
+ rightNodes := make(map[string]int)
+ rightNodes[nodeID1] = 0
+ rightNodes[nodeID2] = 0
+ rightNodes[nodeID3] = -1
+ rightNodes[nodeID4] = 1
- // reserve the allocation on node1
- node1 := newNode(nodeID1, map[string]resources.Quantity{"first": 5})
- err = app.Reserve(node1, ask)
- assert.NilError(t, err, "reservation failed")
+ node := newNode(nodeID1, map[string]resources.Quantity{"first": 5})
- // case 1: node is the reserved node
- iterator := getNodeIteratorFn(node1)
- result := app.tryNodesNoReserve(ask, iterator(), node1.NodeID)
- assert.Assert(t, result == nil, "result should be nil since node1 is
the reserved node")
+ otherNode1 := newNode(nodeID2, map[string]resources.Quantity{"first":
5})
+ otherNode1.schedulable = false
+ otherNode2 := newNode(nodeID2, map[string]resources.Quantity{"first":
1})
+ otherNode3 := newNode(nodeID3, map[string]resources.Quantity{"first":
50})
- // case 2: node is unschedulable
- node2 := newNode(nodeID2, map[string]resources.Quantity{"first": 5})
- node2.schedulable = false
- iterator = getNodeIteratorFn(node2)
- result = app.tryNodesNoReserve(ask, iterator(), node1.NodeID)
- assert.Assert(t, result == nil, "result should be nil since node2 is
unschedulable")
-
- // case 3: node does not have enough resources
- node3 := newNode(nodeID3, map[string]resources.Quantity{"first": 1})
- iterator = getNodeIteratorFn(node3)
- result = app.tryNodesNoReserve(ask, iterator(), node1.NodeID)
- assert.Assert(t, result == nil, "result should be nil since node3 does
not have enough resources")
-
- // case 4: node fails predicate
- mockPlugin := mockCommon.NewPredicatePlugin(false,
map[string]int{nodeID4: 1})
- plugins.RegisterSchedulerPlugin(mockPlugin)
- defer plugins.UnregisterSchedulerPlugins()
- node4 := newNode(nodeID4, map[string]resources.Quantity{"first": 5})
- iterator = getNodeIteratorFn(node4)
- result = app.tryNodesNoReserve(ask, iterator(), node1.NodeID)
- assert.Assert(t, result == nil, "result should be nil since node4 fails
predicate")
-
- // case 5: success
- node5 := newNode(nodeID5, map[string]resources.Quantity{"first": 5})
- iterator = getNodeIteratorFn(node5)
- result = app.tryNodesNoReserve(ask, iterator(), node1.NodeID)
- assert.Assert(t, result != nil, "result should not be nil")
- assert.Equal(t, node5.NodeID, result.NodeID, "result should be on
node5")
- assert.Equal(t, result.ResultType, AllocatedReserved, "result type
should be AllocatedReserved")
- assert.Equal(t, result.ReservedNodeID, node1.NodeID, "reserved node
should be node1")
+ tests := []struct {
+ name string
+ mockPlugin *mockCommon.PredicatePlugin
+ otherNode *Node
+ allocResult bool
+ }{
+ {"prefilter pass, reserved node itself is being used",
mockCommon.NewPredicatePlugin(false, false, nil), node, false},
+ {"prefilter pass, reserved node is different and other node is
no schedulable", mockCommon.NewPredicatePlugin(false, false, nil), otherNode1,
false},
+ {"prefilter pass, reserved node is different and other node
doesn't have sufficient resource", mockCommon.NewPredicatePlugin(false, false,
nil), otherNode2, false},
+ {"prefilter pass with empty feasible nodes, so other node
(node3) is selected", mockCommon.NewPredicatePlugin(false, false, nil),
otherNode3, true},
+ {"prefilter passes but none of the other node from iterator is
available in feasible nodes", mockCommon.NewPredicatePlugin(false, false,
wrongNodes), otherNode3, false},
+ {"prefilter fails", mockCommon.NewPredicatePlugin(true, false,
nil), otherNode3, false},
+ {"prefilter pass with expected feasible nodes, filter fails",
mockCommon.NewPredicatePlugin(false, true, rightNodes), otherNode3, false},
+ {"both prefilter and filter passes with correct feasible nodes.
so other node (node3) is selected", mockCommon.NewPredicatePlugin(false, false,
rightNodes), otherNode3, true},
+ }
+ for _, tt := range tests {
+ t.Run(tt.name, func(t *testing.T) {
+ plugins.RegisterSchedulerPlugin(tt.mockPlugin)
+
+ app := newApplication(appID0, "default", "root.default")
+ app.queue = queue
+
+ res :=
resources.NewResourceFromMap(map[string]resources.Quantity{"first": 5})
+ ask := newAllocationAsk(aKey, appID0, res)
+ err = app.AddAllocationAsk(ask)
+ assert.NilError(t, err, "ask should have been added to
app")
+
+ // reserve the allocation on node1
+ err = app.Reserve(node, ask)
+ assert.NilError(t, err, "reservation failed")
+
+ iterator := getNodeIteratorFn(tt.otherNode)
+ result := app.tryNodesNoReserve(ask, iterator(),
node.NodeID)
+ if tt.allocResult {
+ assert.Assert(t, result != nil, "result should
not be nil")
+ assert.Equal(t, tt.otherNode.NodeID,
result.NodeID, "result should be on other node")
+ assert.Equal(t, result.ResultType,
AllocatedReserved, "result type should be AllocatedReserved")
+ assert.Equal(t, result.ReservedNodeID,
node.NodeID, "reserved node should be node1")
+ } else {
+ assert.Assert(t, result == nil, "result should
be not nil")
+ }
+ node.unReserve(ask)
+ app.UnReserve(node, ask)
+ queue.RemoveApplication(app)
+ plugins.UnregisterSchedulerPlugins()
+ })
+ }
}
func TestAppSubmissionTime(t *testing.T) {
diff --git a/pkg/scheduler/objects/node.go b/pkg/scheduler/objects/node.go
index abb5d4c4..eb987ac7 100644
--- a/pkg/scheduler/objects/node.go
+++ b/pkg/scheduler/objects/node.go
@@ -19,6 +19,7 @@
package objects
import (
+ "errors"
"fmt"
"go.uber.org/zap"
@@ -487,6 +488,20 @@ func (sn *Node) preAllocateConditions(ask *Allocation)
error {
// Checking pre-conditions in the shim for a reservation.
func (sn *Node) preReserveConditions(ask *Allocation) error {
+ // run predicates for this pod before in hand and fetch feasible nodes
+ feasibleNodes, predicatesResult := ask.preAllocateConditions(false)
+ if !predicatesResult {
+ return common.ErrorPreFilterPredicate
+ }
+ // Is this node suitable to run the pod?
+ if len(feasibleNodes) > 0 {
+ if _, ok := feasibleNodes[sn.NodeID]; !ok {
+ log.Log(log.SchedNode).Debug("skipping node as it is
not feasible to run the pod",
+ zap.String("allocationKey",
ask.GetAllocationKey()),
+ zap.String("node", sn.NodeID))
+ return errors.New("skipping node as it is not feasible
to run the pod")
+ }
+ }
return sn.preConditions(ask, false)
}
diff --git a/pkg/scheduler/objects/node_test.go
b/pkg/scheduler/objects/node_test.go
index 93fb6147..32eda82a 100644
--- a/pkg/scheduler/objects/node_test.go
+++ b/pkg/scheduler/objects/node_test.go
@@ -99,8 +99,32 @@ func TestCheckConditions(t *testing.T) {
// Check if we can allocate on scheduling node (no plugins)
res :=
resources.NewResourceFromMap(map[string]resources.Quantity{"first": 1})
ask := newAllocationAsk("test", "app001", res)
- if node.preAllocateConditions(ask) != nil {
- t.Error("node with scheduling set to true no plugins should
allow allocation")
+ tests := []struct {
+ name string
+ mockPlugin *mock.PredicatePlugin
+ predicateResult bool
+ }{
+ {"no plugins", nil, true},
+ {"filter fails", mock.NewPredicatePlugin(false, true, nil),
false},
+ {"filter passes", mock.NewPredicatePlugin(false, false, nil),
true},
+ }
+ for _, tt := range tests {
+ t.Run(tt.name, func(t *testing.T) {
+ if tt.mockPlugin != nil {
+ plugins.RegisterSchedulerPlugin(tt.mockPlugin)
+ }
+ err := node.preAllocateConditions(ask)
+ if tt.predicateResult {
+ assert.NilError(t, err, "node with scheduling
set to true no plugins should allow allocation")
+ } else {
+ assert.ErrorContains(t, err, "fake predicate
plugin failed")
+ assert.Equal(t, 1, len(ask.allocLog))
+ assert.Equal(t, "fake predicate plugin failed",
ask.allocLog["fake predicate plugin failed"].Message)
+ }
+ if tt.mockPlugin != nil {
+ plugins.UnregisterSchedulerPlugins()
+ }
+ })
}
}
@@ -920,9 +944,6 @@ func TestNode_FitInNode(t *testing.T) {
}
func TestPreconditions(t *testing.T) {
- defer plugins.UnregisterSchedulerPlugins()
-
- plugins.RegisterSchedulerPlugin(mock.NewPredicatePlugin(true,
map[string]int{}))
total :=
resources.NewResourceFromMap(map[string]resources.Quantity{"cpu": 100,
"memory": 100})
proto := newProto(testNode, total, map[string]string{
"ready": "true",
@@ -931,17 +952,29 @@ func TestPreconditions(t *testing.T) {
ask := newAllocationAsk("test", "app001", res)
node := NewNode(proto)
- // failure
- err := node.preConditions(ask, true)
- assert.ErrorContains(t, err, "fake predicate plugin failed")
- assert.Equal(t, 1, len(ask.allocLog))
- assert.Equal(t, "fake predicate plugin failed", ask.allocLog["fake
predicate plugin failed"].Message)
-
- // pass
- plugins.RegisterSchedulerPlugin(mock.NewPredicatePlugin(false,
map[string]int{}))
- err = node.preConditions(ask, true)
- assert.NilError(t, err)
- assert.Equal(t, 1, len(ask.allocLog))
+ tests := []struct {
+ name string
+ mockPlugin *mock.PredicatePlugin
+ predicateResult bool
+ }{
+ {"filter fails", mock.NewPredicatePlugin(false, true, nil),
false},
+ {"filter passes", mock.NewPredicatePlugin(false, false, nil),
true},
+ }
+ for _, tt := range tests {
+ t.Run(tt.name, func(t *testing.T) {
+ plugins.RegisterSchedulerPlugin(tt.mockPlugin)
+ err := node.preConditions(ask, true)
+ if tt.predicateResult {
+ assert.NilError(t, err)
+ assert.Equal(t, 1, len(ask.allocLog))
+ } else {
+ assert.ErrorContains(t, err, "fake predicate
plugin failed")
+ assert.Equal(t, 1, len(ask.allocLog))
+ assert.Equal(t, "fake predicate plugin failed",
ask.allocLog["fake predicate plugin failed"].Message)
+ }
+ plugins.UnregisterSchedulerPlugins()
+ })
+ }
}
func TestGetAllocations(t *testing.T) {
diff --git a/pkg/scheduler/objects/preemption.go
b/pkg/scheduler/objects/preemption.go
index 788bbbdd..0ca782c5 100644
--- a/pkg/scheduler/objects/preemption.go
+++ b/pkg/scheduler/objects/preemption.go
@@ -550,7 +550,23 @@ func (p *Preemptor) tryNodes() (string, []*Allocation,
bool) {
// calculate victim list for each node
predicateChecks := make([]*si.PreemptionPredicatesArgs, 0)
victimsByNode := make(map[string][]*Allocation)
+
+ // run predicates for this pod before in hand and fetch feasible nodes
+ feasibleNodes, predicatesResult := p.ask.preAllocateConditions(true)
+ if !predicatesResult {
+ return "", nil, false
+ }
+
for nodeID, nodeAvailable := range p.nodeAvailableMap {
+ // Is this node suitable to run the pod?
+ if len(feasibleNodes) > 0 {
+ if _, ok := feasibleNodes[nodeID]; !ok {
+ log.Log(log.SchedApplication).Debug("skipping
node as it is not feasible to run the pod",
+ zap.String("allocationKey",
p.ask.GetAllocationKey()),
+ zap.String("node", nodeID))
+ continue
+ }
+ }
allocations, ok := p.allocationsByNode[nodeID]
if !ok {
// no allocations present, but node may still be
available for scheduling
diff --git a/pkg/scheduler/objects/preemption_test.go
b/pkg/scheduler/objects/preemption_test.go
index addfb387..a6f586b8 100644
--- a/pkg/scheduler/objects/preemption_test.go
+++ b/pkg/scheduler/objects/preemption_test.go
@@ -19,6 +19,7 @@
package objects
import (
+ "errors"
"fmt"
"strconv"
"testing"
@@ -120,6 +121,24 @@ func creatApp2(
return app2, ask3, nil
}
+func resetNode(node *Node) {
+ for _, v := range node.allocations {
+ node.RemoveAllocation(v.allocationKey)
+ }
+}
+
+func resetQ(t *testing.T, queue *Queue) {
+ for _, v := range queue.applications {
+ for _, a := range v.allocations {
+ v.RemoveAllocationAsk(a.allocationKey)
+ err :=
queue.DecAllocatedResource(a.GetAllocatedResource())
+ assert.NilError(t, err)
+ }
+ queue.RemoveApplication(v)
+ }
+ queue.applications = make(map[string]*Application)
+}
+
func TestCheckPreconditions(t *testing.T) {
node := newNode("node1", map[string]resources.Quantity{"first": 5})
iterator := getNodeIteratorFn(node)
@@ -334,32 +353,58 @@ func TestTryPreemption(t *testing.T) {
childQ2, err := createManagedQueueGuaranteed(parentQ, "child2", false,
map[string]string{"first": "10"}, map[string]string{"first": "5"},
appQueueMapping)
assert.NilError(t, err)
- alloc1, alloc2, err := creatApp1(childQ1, node, nil,
map[string]resources.Quantity{"first": 5, "pods": 1}, appQueueMapping)
- assert.NilError(t, err)
-
- app2, ask3, err := creatApp2(childQ2,
map[string]resources.Quantity{"first": 5, "pods": 1}, "alloc3", appQueueMapping)
- assert.NilError(t, err)
- childQ2.incPendingResource(ask3.GetAllocatedResource())
-
- headRoom :=
resources.NewResourceFromMap(map[string]resources.Quantity{"first": 10, "pods":
3})
- preemptor := NewPreemptor(app2, headRoom, 30*time.Second, ask3,
iterator(), false)
-
// register predicate handler
preemptions := []mock.Preemption{
mock.NewPreemption(true, "alloc3", nodeID1, []string{"alloc1"},
0, 0),
}
- plugin := mock.NewPreemptionPredicatePlugin(nil, nil, preemptions)
- plugins.RegisterSchedulerPlugin(plugin)
- defer plugins.UnregisterSchedulerPlugins()
-
- result, ok := preemptor.TryPreemption()
- assert.Assert(t, result != nil, "no result")
- assert.NilError(t, plugin.GetPredicateError())
- assert.Assert(t, ok, "no victims found")
- assert.Equal(t, "alloc3", result.Request.GetAllocationKey(), "wrong
alloc")
- assert.Check(t, alloc1.IsPreempted(), "alloc1 not preempted")
- assert.Check(t, !alloc2.IsPreempted(), "alloc2 preempted")
- assert.Equal(t, len(ask3.GetAllocationLog()), 0)
+ wrongNodes := make(map[string]int, 1)
+ wrongNodes[nodeID2] = 10
+ rightNodes := make(map[string]int, 1)
+ rightNodes[nodeID1] = 10
+ tests := []struct {
+ name string
+ mockPlugin *mock.PreemptionPredicatePlugin
+ result bool
+ mockPluginError error
+ }{
+ {"prefilter fails",
mock.NewPreemptionPredicatePlugin(preemptions, nil, true, false), false,
errors.New("fail")},
+ {"prefilter passes but none of the node from iterator is
available in feasible nodes", mock.NewPreemptionPredicatePlugin(preemptions,
wrongNodes, false, false), false, nil},
+ {"prefilter pass with expected feasible nodes, filter fails",
mock.NewPreemptionPredicatePlugin(preemptions, rightNodes, false, true), false,
errors.New("fail")},
+ {"both prefilter and filter passes with correct feasible
nodes", mock.NewPreemptionPredicatePlugin(preemptions, rightNodes, false,
false), true, nil},
+ {"both prefilter and filter passes with empty feasible nodes",
mock.NewPreemptionPredicatePlugin(preemptions, nil, false, false), true, nil},
+ }
+ for _, tt := range tests {
+ t.Run(tt.name, func(t *testing.T) {
+ alloc1, alloc2, err := creatApp1(childQ1, node, nil,
map[string]resources.Quantity{"first": 5, "pods": 1}, appQueueMapping)
+ assert.NilError(t, err)
+ app2, ask3, err := creatApp2(childQ2,
map[string]resources.Quantity{"first": 5, "pods": 1}, "alloc3", appQueueMapping)
+ assert.NilError(t, err)
+ childQ2.incPendingResource(ask3.GetAllocatedResource())
+ headRoom :=
resources.NewResourceFromMap(map[string]resources.Quantity{"first": 10, "pods":
3})
+ preemptor := NewPreemptor(app2, headRoom,
30*time.Second, ask3, iterator(), false)
+ plugins.RegisterSchedulerPlugin(tt.mockPlugin)
+ result, ok := preemptor.TryPreemption()
+ if tt.result {
+ assert.Assert(t, result != nil, "no result")
+ assert.NilError(t,
tt.mockPlugin.GetPredicateError())
+ assert.Assert(t, ok, "no victims found")
+ assert.Equal(t, "alloc3",
result.Request.GetAllocationKey(), "wrong alloc")
+ assert.Check(t, alloc1.IsPreempted(), "alloc1
not preempted")
+ assert.Check(t, !alloc2.IsPreempted(), "alloc2
preempted")
+ assert.Equal(t, len(ask3.GetAllocationLog()), 0)
+ } else {
+ assert.Assert(t, result == nil, "no result")
+ }
+ if tt.mockPluginError != nil {
+ assert.ErrorContains(t,
tt.mockPlugin.GetPredicateError(), tt.mockPluginError.Error())
+ }
+ // reset
+ resetNode(node)
+ resetQ(t, childQ1)
+ resetQ(t, childQ2)
+ plugins.UnregisterSchedulerPlugins()
+ })
+ }
}
func TestTryPreemption_SendEvent(t *testing.T) {
@@ -393,7 +438,7 @@ func TestTryPreemption_SendEvent(t *testing.T) {
preemptions := []mock.Preemption{
mock.NewPreemption(true, "alloc3", nodeID1, []string{"alloc1"},
0, 0),
}
- plugin := mock.NewPreemptionPredicatePlugin(nil, nil, preemptions)
+ plugin := mock.NewPreemptionPredicatePlugin(preemptions, nil, false,
false)
plugins.RegisterSchedulerPlugin(plugin)
defer plugins.UnregisterSchedulerPlugins()
@@ -451,7 +496,7 @@ func TestTryPreemptionOnNode(t *testing.T) {
preemptions := []mock.Preemption{
mock.NewPreemption(true, "alloc3", nodeID2, []string{"alloc2"},
0, 0),
}
- plugin := mock.NewPreemptionPredicatePlugin(nil, nil, preemptions)
+ plugin := mock.NewPreemptionPredicatePlugin(preemptions, nil, false,
false)
plugins.RegisterSchedulerPlugin(plugin)
defer plugins.UnregisterSchedulerPlugins()
@@ -560,7 +605,7 @@ func TestTryPreemptionOnNodeWithOGParentAndUGPreemptor(t
*testing.T) {
preemptions := []mock.Preemption{
mock.NewPreemption(true, "alloc7", nodeID2, []string{"alloc1"},
0, 0),
}
- plugin := mock.NewPreemptionPredicatePlugin(nil, nil, preemptions)
+ plugin := mock.NewPreemptionPredicatePlugin(preemptions, nil, false,
false)
plugins.RegisterSchedulerPlugin(plugin)
defer plugins.UnregisterSchedulerPlugins()
@@ -583,43 +628,71 @@ func TestTryPreemptionOnNodeWithOGParentAndUGPreemptor(t
*testing.T) {
// root.parent.child2. Guaranteed set, first: 5. Request of first:5 is waiting
for resources.
// 1 Allocation on root.parent.child1 should be preempted to free up resources
for ask arrived in root.parent.child2.
func TestTryPreemptionOnQueue(t *testing.T) {
- appQueneMapping := NewAppQueueMapping()
+ appQueueMapping := NewAppQueueMapping()
node1 := newNode(nodeID1, map[string]resources.Quantity{"first": 10,
"pods": 2})
node2 := newNode(nodeID2, map[string]resources.Quantity{"first": 10,
"pods": 2})
iterator := getNodeIteratorFn(node1, node2)
rootQ, err := createRootQueue(map[string]string{"first": "20", "pods":
"4"})
assert.NilError(t, err)
- parentQ, err := createManagedQueueGuaranteed(rootQ, "parent", true,
map[string]string{"first": "10"}, nil, appQueneMapping)
+ parentQ, err := createManagedQueueGuaranteed(rootQ, "parent", true,
map[string]string{"first": "10"}, nil, appQueueMapping)
assert.NilError(t, err)
- childQ1, err := createManagedQueueGuaranteed(parentQ, "child1", false,
nil, map[string]string{"first": "5"}, appQueneMapping)
- assert.NilError(t, err)
- childQ2, err := createManagedQueueGuaranteed(parentQ, "child2", false,
nil, map[string]string{"first": "5"}, appQueneMapping)
- assert.NilError(t, err)
-
- alloc1, alloc2, err := creatApp1(childQ1, node1, node2,
map[string]resources.Quantity{"first": 5, "pods": 1}, appQueneMapping)
+ childQ1, err := createManagedQueueGuaranteed(parentQ, "child1", false,
nil, map[string]string{"first": "5"}, appQueueMapping)
assert.NilError(t, err)
-
- app2, ask3, err := creatApp2(childQ2,
map[string]resources.Quantity{"first": 5, "pods": 1}, "alloc3", appQueneMapping)
+ childQ2, err := createManagedQueueGuaranteed(parentQ, "child2", false,
nil, map[string]string{"first": "5"}, appQueueMapping)
assert.NilError(t, err)
- headRoom :=
resources.NewResourceFromMap(map[string]resources.Quantity{"first": 10, "pods":
3})
- preemptor := NewPreemptor(app2, headRoom, 30*time.Second, ask3,
iterator(), false)
-
- allocs := map[string]string{}
- allocs["alloc3"] = nodeID2
-
- plugin := mock.NewPreemptionPredicatePlugin(nil, allocs, nil)
- plugins.RegisterSchedulerPlugin(plugin)
- defer plugins.UnregisterSchedulerPlugins()
-
- result, ok := preemptor.TryPreemption()
- assert.Assert(t, result != nil, "no result")
- assert.Assert(t, ok, "no victims found")
- assert.Equal(t, "alloc3", result.Request.GetAllocationKey(), "wrong
alloc")
- assert.Equal(t, nodeID2, result.NodeID, "wrong node")
- assert.Check(t, !alloc1.IsPreempted(), "alloc1 preempted")
- assert.Check(t, alloc2.IsPreempted(), "alloc2 not preempted")
- assert.Equal(t, len(ask3.GetAllocationLog()), 0)
+ wrongNodes := make(map[string]int, 1)
+ wrongNodes[nodeID3] = 10
+ rightNodes := make(map[string]int, 1)
+ rightNodes[nodeID1] = 10
+ tests := []struct {
+ name string
+ mockPlugin *mock.PreemptionPredicatePlugin
+ result bool
+ mockPluginError error
+ }{
+ {"prefilter fails", mock.NewPreemptionPredicatePlugin(nil, nil,
true, false), false, errors.New("fail")},
+ {"prefilter passes but none of the node from iterator is
available in feasible nodes", mock.NewPreemptionPredicatePlugin(nil,
wrongNodes, false, false), false, nil},
+ {"prefilter pass with expected feasible nodes, filter fails",
mock.NewPreemptionPredicatePlugin(nil, rightNodes, false, true), false,
errors.New("fail")},
+ {"both prefilter and filter passes with correct feasible
nodes", mock.NewPreemptionPredicatePlugin(nil, rightNodes, false, false), true,
nil},
+ {"both prefilter and filter passes with empty feasible nodes",
mock.NewPreemptionPredicatePlugin(nil, nil, false, false), true, nil},
+ }
+ for _, tt := range tests {
+ t.Run(tt.name, func(t *testing.T) {
+ alloc1, alloc2, err := creatApp1(childQ1, node1, node2,
map[string]resources.Quantity{"first": 5, "pods": 1}, appQueueMapping)
+ assert.NilError(t, err)
+ app2, ask3, err := creatApp2(childQ2,
map[string]resources.Quantity{"first": 5, "pods": 1}, "alloc3", appQueueMapping)
+ assert.NilError(t, err)
+ headRoom :=
resources.NewResourceFromMap(map[string]resources.Quantity{"first": 10, "pods":
3})
+ preemptor := NewPreemptor(app2, headRoom,
30*time.Second, ask3, iterator(), false)
+ plugins.RegisterSchedulerPlugin(tt.mockPlugin)
+ result, ok := preemptor.TryPreemption()
+ if tt.result {
+ assert.Assert(t, result != nil, "no result")
+ assert.Assert(t, ok, "no victims found")
+ assert.Equal(t, "alloc3",
result.Request.GetAllocationKey(), "wrong alloc")
+ nodes := make(map[string]int)
+ nodes[nodeID1] = 1
+ nodes[nodeID2] = 1
+ _, ok = nodes[result.NodeID]
+ assert.Check(t, ok == true, "either node1 or
node2 chosen")
+ assert.Check(t, alloc1.IsPreempted() ||
alloc2.IsPreempted(), "either alloc1 or alloc2 preempted, but not both")
+ assert.Check(t, !alloc1.IsPreempted() ||
!alloc2.IsPreempted(), "either alloc1 or alloc2 not preempted, but not both")
+ assert.Equal(t, len(ask3.GetAllocationLog()), 0)
+ } else {
+ assert.Assert(t, result == nil, "no result")
+ }
+ if tt.mockPluginError != nil {
+ assert.ErrorContains(t,
tt.mockPlugin.GetPredicateError(), tt.mockPluginError.Error())
+ }
+ // reset
+ resetNode(node1)
+ resetNode(node2)
+ resetQ(t, childQ1)
+ resetQ(t, childQ2)
+ plugins.UnregisterSchedulerPlugins()
+ })
+ }
}
// TestTryPreemption_VictimsAvailable_InsufficientResource Test try preemption
on queue with simple queue hierarchy. Since Node has enough resources to
accomodate, preemption happens because of queue resource constraint.
@@ -699,8 +772,7 @@ func
TestTryPreemption_VictimsOnDifferentNodes_InsufficientResource(t *testing.T
mock.NewPreemption(true, "alloc3", nodeID1, []string{"alloc1"},
0, 0),
mock.NewPreemption(true, "alloc3", nodeID2, []string{"alloc2"},
0, 0),
}
-
- plugin := mock.NewPreemptionPredicatePlugin(nil, nil, preemptions)
+ plugin := mock.NewPreemptionPredicatePlugin(preemptions, nil, false,
false)
plugins.RegisterSchedulerPlugin(plugin)
defer plugins.UnregisterSchedulerPlugins()
@@ -772,10 +844,7 @@ func
TestTryPreemption_VictimReleased_InsufficientResource(t *testing.T) {
err = alloc4.SetReleased(true)
assert.NilError(t, err)
- allocs := map[string]string{}
- allocs["alloc3"] = nodeID2
-
- plugin := mock.NewPreemptionPredicatePlugin(nil, allocs, nil)
+ plugin := mock.NewPreemptionPredicatePlugin(nil, nil, false, false)
plugins.RegisterSchedulerPlugin(plugin)
defer plugins.UnregisterSchedulerPlugins()
@@ -844,7 +913,7 @@ func TestTryPreemption_VictimsAvailableOnDifferentNodes(t
*testing.T) {
mock.NewPreemption(true, "alloc3", nodeID1, []string{"alloc1"},
0, 0),
mock.NewPreemption(true, "alloc3", nodeID2, []string{"alloc2"},
0, 0),
}
- plugin := mock.NewPreemptionPredicatePlugin(nil, nil, preemptions)
+ plugin := mock.NewPreemptionPredicatePlugin(preemptions, nil, false,
false)
plugins.RegisterSchedulerPlugin(plugin)
defer plugins.UnregisterSchedulerPlugins()
@@ -908,10 +977,9 @@ func TestTryPreemption_OnQueue_VictimsOnDifferentNodes(t
*testing.T) {
headRoom :=
resources.NewResourceFromMap(map[string]resources.Quantity{"first": 10, "pods":
3})
preemptor := NewPreemptor(app2, headRoom, 30*time.Second, ask3,
iterator(), false)
- allocs := map[string]string{}
- allocs["alloc3"] = nodeID2
-
- plugin := mock.NewPreemptionPredicatePlugin(nil, allocs, nil)
+ feasibleNodes := map[string]int{}
+ feasibleNodes[nodeID2] = 1
+ plugin := mock.NewPreemptionPredicatePlugin(nil, feasibleNodes, false,
false)
plugins.RegisterSchedulerPlugin(plugin)
defer plugins.UnregisterSchedulerPlugins()
@@ -998,10 +1066,7 @@ func
TestTryPreemption_OnQueue_VictimsAvailable_LowerPriority(t *testing.T) {
headRoom :=
resources.NewResourceFromMap(map[string]resources.Quantity{"first": 10, "pods":
3})
preemptor := NewPreemptor(app2, headRoom, 30*time.Second, ask3,
iterator(), false)
- allocs := map[string]string{}
- allocs["alloc3"] = nodeID1
-
- plugin := mock.NewPreemptionPredicatePlugin(nil, allocs, nil)
+ plugin := mock.NewPreemptionPredicatePlugin(nil, nil, false, false)
plugins.RegisterSchedulerPlugin(plugin)
defer plugins.UnregisterSchedulerPlugins()
@@ -1077,10 +1142,7 @@ func
TestTryPreemption_AskResTypesDifferent_GuaranteedSetOnPreemptorSide(t *test
headRoom :=
resources.NewResourceFromMap(map[string]resources.Quantity{"vcores": 2})
preemptor := NewPreemptor(app2, headRoom, 30*time.Second, ask3,
iterator(), false)
- allocs := map[string]string{}
- allocs["alloc3"] = nodeID1
-
- plugin := mock.NewPreemptionPredicatePlugin(nil, allocs, nil)
+ plugin := mock.NewPreemptionPredicatePlugin(nil, nil, false, false)
plugins.RegisterSchedulerPlugin(plugin)
defer plugins.UnregisterSchedulerPlugins()
@@ -1155,7 +1217,7 @@ func
TestTryPreemption_OnNode_AskResTypesDifferent_GuaranteedSetOnPreemptorSide(
// register predicate handler
preemptions := []mock.Preemption{mock.NewPreemption(true, "alloc3",
nodeID1, []string{"alloc2"}, 0, 0)}
- plugin := mock.NewPreemptionPredicatePlugin(nil, nil, preemptions)
+ plugin := mock.NewPreemptionPredicatePlugin(preemptions, nil, false,
false)
plugins.RegisterSchedulerPlugin(plugin)
defer plugins.UnregisterSchedulerPlugins()
@@ -1250,9 +1312,7 @@ func
TestTryPreemption_AskResTypesDifferent_GuaranteedSetOnVictimAndPreemptorSid
preemptor := NewPreemptor(app4, headRoom, 30*time.Second, ask4,
iterator(), false)
// register predicate handler
- allocs := map[string]string{}
- allocs["alloc4"] = nodeID1
- plugin := mock.NewPreemptionPredicatePlugin(nil, allocs, nil)
+ plugin := mock.NewPreemptionPredicatePlugin(nil, nil, false, false)
plugins.RegisterSchedulerPlugin(plugin)
defer plugins.UnregisterSchedulerPlugins()
@@ -1349,7 +1409,7 @@ func
TestTryPreemption_OnNode_AskResTypesDifferent_GuaranteedSetOnVictimAndPreem
// register predicate handler
preemptions := []mock.Preemption{mock.NewPreemption(true, "alloc4",
nodeID1, []string{"alloc3", "alloc2"}, 1, 1)}
- plugin := mock.NewPreemptionPredicatePlugin(nil, nil, preemptions)
+ plugin := mock.NewPreemptionPredicatePlugin(preemptions, nil, false,
false)
plugins.RegisterSchedulerPlugin(plugin)
defer plugins.UnregisterSchedulerPlugins()
@@ -1465,9 +1525,7 @@ func
TestTryPreemption_AskResTypesSame_GuaranteedSetOnPreemptorSide(t *testing.T
preemptor := NewPreemptor(app4, headRoom, 30*time.Second, ask4,
iterator(), false)
// register predicate handler
- allocs := map[string]string{}
- allocs["alloc4"] = nodeID1
- plugin := mock.NewPreemptionPredicatePlugin(nil, allocs, nil)
+ plugin := mock.NewPreemptionPredicatePlugin(nil, nil, false, false)
plugins.RegisterSchedulerPlugin(plugin)
defer plugins.UnregisterSchedulerPlugins()
@@ -1564,7 +1622,7 @@ func
TestTryPreemption_OnNode_AskResTypesSame_GuaranteedSetOnPreemptorSide(t *te
// register predicate handler
preemptions := []mock.Preemption{mock.NewPreemption(true, "alloc4",
nodeID1, []string{"alloc3", "alloc2"}, 1, 1)}
- plugin := mock.NewPreemptionPredicatePlugin(nil, nil, preemptions)
+ plugin := mock.NewPreemptionPredicatePlugin(preemptions, nil, false,
false)
plugins.RegisterSchedulerPlugin(plugin)
defer plugins.UnregisterSchedulerPlugins()
@@ -1663,9 +1721,7 @@ func
TestTryPreemption_AskResTypesSame_GuaranteedSetOnVictimAndPreemptorSides(t
preemptor := NewPreemptor(app4, headRoom, 30*time.Second, ask4,
iterator(), false)
// register predicate handler
- allocs := map[string]string{}
- allocs["alloc4"] = nodeID1
- plugin := mock.NewPreemptionPredicatePlugin(nil, allocs, nil)
+ plugin := mock.NewPreemptionPredicatePlugin(nil, nil, false, false)
plugins.RegisterSchedulerPlugin(plugin)
defer plugins.UnregisterSchedulerPlugins()
@@ -1763,7 +1819,7 @@ func
TestTryPreemption_OnNode_AskResTypesSame_GuaranteedSetOnVictimAndPreemptorS
// register predicate handler
preemptions := []mock.Preemption{mock.NewPreemption(true, "alloc4",
nodeID1, []string{"alloc3", "alloc2"}, 1, 1)}
- plugin := mock.NewPreemptionPredicatePlugin(nil, nil, preemptions)
+ plugin := mock.NewPreemptionPredicatePlugin(preemptions, nil, false,
false)
plugins.RegisterSchedulerPlugin(plugin)
defer plugins.UnregisterSchedulerPlugins()
@@ -1848,7 +1904,7 @@ func
TestTryPreemption_OnNode_UGParent_With_UGPreemptorChild_GNotSetOnVictimChil
// register predicate handler
preemptions := []mock.Preemption{mock.NewPreemption(true, "alloc3",
nodeID1, []string{"alloc2"}, 0, 0)}
- plugin := mock.NewPreemptionPredicatePlugin(nil, nil, preemptions)
+ plugin := mock.NewPreemptionPredicatePlugin(preemptions, nil, false,
false)
plugins.RegisterSchedulerPlugin(plugin)
defer plugins.UnregisterSchedulerPlugins()
@@ -1928,7 +1984,7 @@ func
TestTryPreemption_OnNode_UGParent_With_GNotSetOnBothChilds(t *testing.T) {
preemptor := NewPreemptor(app3, headRoom, 30*time.Second, ask3,
iterator(), false)
// register predicate handler
- plugin := mock.NewPreemptionPredicatePlugin(nil, nil, nil)
+ plugin := mock.NewPreemptionPredicatePlugin(nil, nil, false, false)
plugins.RegisterSchedulerPlugin(plugin)
defer plugins.UnregisterSchedulerPlugins()
@@ -2004,7 +2060,7 @@ func
TestTryPreemption_OnNode_UGParent_With_UGPreemptorChild_OGVictimChild_As_Si
// register predicate handler
preemptions := []mock.Preemption{mock.NewPreemption(true, "alloc3",
nodeID1, []string{"alloc2"}, 0, 0)}
- plugin := mock.NewPreemptionPredicatePlugin(nil, nil, preemptions)
+ plugin := mock.NewPreemptionPredicatePlugin(preemptions, nil, false,
false)
plugins.RegisterSchedulerPlugin(plugin)
defer plugins.UnregisterSchedulerPlugins()
@@ -2366,7 +2422,7 @@ func Test_PreemptReleasesReservationsOnSuccess(t
*testing.T) {
preemptions := []mock.Preemption{
mock.NewPreemption(true, "alloc3", nodeID2, []string{"alloc2"},
0, 0),
}
- plugin := mock.NewPreemptionPredicatePlugin(nil, nil, preemptions)
+ plugin := mock.NewPreemptionPredicatePlugin(preemptions, nil, false,
false)
plugins.RegisterSchedulerPlugin(plugin)
defer plugins.UnregisterSchedulerPlugins()
diff --git a/pkg/scheduler/partition_test.go b/pkg/scheduler/partition_test.go
index e3187301..03d0495a 100644
--- a/pkg/scheduler/partition_test.go
+++ b/pkg/scheduler/partition_test.go
@@ -3649,10 +3649,10 @@ func TestFailReplacePlaceholder(t *testing.T) {
t.Fatalf("empty cluster placeholder allocate returned
allocation: %s", result)
}
// plugin to let the pre-check fail on node-1 only, means we cannot
replace the placeholder
- plugin := mock.NewPredicatePlugin(false, map[string]int{nodeID1: -1})
+ plugin := mock.NewPredicatePlugin(false, false, map[string]int{nodeID1:
0, nodeID2: 1})
plugins.RegisterSchedulerPlugin(plugin)
defer func() {
- passPlugin := mock.NewPredicatePlugin(false, nil)
+ passPlugin := mock.NewPredicatePlugin(false, false,
map[string]int{nodeID1: 0, nodeID2: 1})
plugins.RegisterSchedulerPlugin(passPlugin)
}()
var tgRes, res *resources.Resource
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]