This is an automated email from the ASF dual-hosted git repository.
wilfreds 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 5716f462 [YUNIKORN-2519] Remove bypass ACL check from placement rules
(#829)
5716f462 is described below
commit 5716f4627dc948182a5268c7c12974421cefa571
Author: Wilfred Spiegelenburg <[email protected]>
AuthorDate: Wed Apr 3 22:41:59 2024 +1100
[YUNIKORN-2519] Remove bypass ACL check from placement rules (#829)
Instead of returning a flag to not bypass the ACL check by all rules
except for the recovery rule special case the recovery rule to bypass
checks.
Change logger handler to SchedApplication for placements
Closes: #829
Signed-off-by: Wilfred Spiegelenburg <[email protected]>
---
pkg/scheduler/partition.go | 14 ++-
pkg/scheduler/placement/fixed_rule.go | 37 ++++----
pkg/scheduler/placement/fixed_rule_test.go | 29 ++----
pkg/scheduler/placement/placement.go | 121 ++++++++++++++------------
pkg/scheduler/placement/placement_test.go | 87 ++++++++++++++++++
pkg/scheduler/placement/provided_rule.go | 26 +++---
pkg/scheduler/placement/provided_rule_test.go | 32 +++----
pkg/scheduler/placement/recovery_rule.go | 16 ++--
pkg/scheduler/placement/recovery_rule_test.go | 7 +-
pkg/scheduler/placement/rule.go | 3 +-
pkg/scheduler/placement/rule_test.go | 13 +--
pkg/scheduler/placement/tag_rule.go | 26 +++---
pkg/scheduler/placement/tag_rule_test.go | 35 +++-----
pkg/scheduler/placement/testrule.go | 8 +-
pkg/scheduler/placement/user_rule.go | 24 ++---
pkg/scheduler/placement/user_rule_test.go | 29 ++----
16 files changed, 274 insertions(+), 233 deletions(-)
diff --git a/pkg/scheduler/partition.go b/pkg/scheduler/partition.go
index 26ae861c..207a0ad5 100644
--- a/pkg/scheduler/partition.go
+++ b/pkg/scheduler/partition.go
@@ -289,7 +289,9 @@ func (pc *PartitionContext) getPlacementManager()
*placement.AppPlacementManager
return pc.placementManager
}
-// Add a new application to the partition.
+// AddApplication adds a new application to the partition.
+// Runs the placement rules for the queue resolution. Creates a new dynamic
queue if the queue does not yet
+// exists.
// NOTE: this is a lock free call. It must NOT be called holding the
PartitionContext lock.
func (pc *PartitionContext) AddApplication(app *objects.Application) error {
if pc.isDraining() || pc.isStopped() {
@@ -302,16 +304,13 @@ func (pc *PartitionContext) AddApplication(app
*objects.Application) error {
return fmt.Errorf("adding application %s to partition %s, but
application already existed", appID, pc.Name)
}
- // Put app under the queue
- pm := pc.getPlacementManager()
- err := pm.PlaceApplication(app)
+ // Resolve the queue for this app using the placement rules
+ // We either have an error or a queue name is set on the application.
+ err := pc.getPlacementManager().PlaceApplication(app)
if err != nil {
return fmt.Errorf("failed to place application %s: %v", appID,
err)
}
queueName := app.GetQueuePath()
- if queueName == "" {
- return fmt.Errorf("application rejected by placement rules:
%s", appID)
- }
// lock the partition and make the last change: we need to do this
before creating the queues.
// queue cleanup might otherwise remove the queue again before we can
add the application
@@ -322,7 +321,6 @@ func (pc *PartitionContext) AddApplication(app
*objects.Application) error {
// create the queue if necessary
if queue == nil {
- var err error
if common.IsRecoveryQueue(queueName) {
queue, err = pc.createRecoveryQueue()
if err != nil {
diff --git a/pkg/scheduler/placement/fixed_rule.go
b/pkg/scheduler/placement/fixed_rule.go
index 0dce3847..b62abd51 100644
--- a/pkg/scheduler/placement/fixed_rule.go
+++ b/pkg/scheduler/placement/fixed_rule.go
@@ -30,16 +30,16 @@ import (
"github.com/apache/yunikorn-core/pkg/scheduler/placement/types"
)
+// A rule to place an application based on the queue in the configuration.
+// If the queue provided is fully qualified, starts with "root.", the parent
rule is skipped and the queue is created as
+// configured. If the queue is not qualified all "." characters will be
replaced and the parent rule run before making
+// the queue name fully qualified.
type fixedRule struct {
basicRule
queue string
qualified bool
}
-// A rule to place an application based on the queue in the configuration.
-// If the queue provided is fully qualified, starts with "root.", the parent
rule is skipped and the queue is created as
-// configured. If the queue is not qualified all "." characters will be
replaced and the parent rule run before making
-// the queue name fully qualified.
func (fr *fixedRule) getName() string {
return types.Fixed
}
@@ -63,31 +63,30 @@ func (fr *fixedRule) initialise(conf configs.PlacementRule)
error {
return err
}
-func (fr *fixedRule) placeApplication(app *objects.Application, queueFn
func(string) *objects.Queue) (string, bool, error) {
+func (fr *fixedRule) placeApplication(app *objects.Application, queueFn
func(string) *objects.Queue) (string, error) {
// before anything run the filter
if !fr.filter.allowUser(app.GetUser()) {
- log.Log(log.Config).Debug("Fixed rule filtered",
+ log.Log(log.SchedApplication).Debug("Fixed rule filtered",
zap.String("application", app.ApplicationID),
zap.Any("user", app.GetUser()),
zap.String("queueName", fr.queue))
- return "", true, nil
+ return "", nil
}
- var parentName string
- var aclCheck = true
- var err error
queueName := fr.queue
// if the fixed queue is already fully qualified skip the parent check
if !fr.qualified {
+ var parentName string
+ var err error
// run the parent rule if set
if fr.parent != nil {
- parentName, aclCheck, err =
fr.parent.placeApplication(app, queueFn)
+ parentName, err = fr.parent.placeApplication(app,
queueFn)
// failed parent rule, fail this rule
if err != nil {
- return "", aclCheck, err
+ return "", err
}
// rule did not return a parent: this could be filter
or create flag related
if parentName == "" {
- return "", aclCheck, nil
+ return "", nil
}
// check if this is a parent queue and qualify it
if !strings.HasPrefix(parentName,
configs.RootQueue+configs.DOT) {
@@ -96,7 +95,7 @@ func (fr *fixedRule) placeApplication(app
*objects.Application, queueFn func(str
// if the parent queue exists it cannot be a leaf
parentQueue := queueFn(parentName)
if parentQueue != nil && parentQueue.IsLeafQueue() {
- return "", aclCheck, fmt.Errorf("parent rule
returned a leaf queue: %s", parentName)
+ return "", fmt.Errorf("parent rule returned a
leaf queue: %s", parentName)
}
}
// the parent is set from the rule otherwise set it to the root
@@ -105,18 +104,18 @@ func (fr *fixedRule) placeApplication(app
*objects.Application, queueFn func(str
}
queueName = parentName + configs.DOT + fr.queue
}
- // Log the result before we really create
- log.Log(log.Config).Debug("Fixed rule intermediate result",
+ // Log the result before we check the create flag
+ log.Log(log.SchedApplication).Debug("Fixed rule intermediate result",
zap.String("application", app.ApplicationID),
zap.String("queue", queueName))
// get the queue object
queue := queueFn(queueName)
// if we cannot create the queue must exist
if !fr.create && queue == nil {
- return "", aclCheck, nil
+ return "", nil
}
- log.Log(log.Config).Info("Fixed rule application placed",
+ log.Log(log.SchedApplication).Info("Fixed rule application placed",
zap.String("application", app.ApplicationID),
zap.String("queue", queueName))
- return queueName, aclCheck, nil
+ return queueName, nil
}
diff --git a/pkg/scheduler/placement/fixed_rule_test.go
b/pkg/scheduler/placement/fixed_rule_test.go
index c3e4e951..5067f93e 100644
--- a/pkg/scheduler/placement/fixed_rule_test.go
+++ b/pkg/scheduler/placement/fixed_rule_test.go
@@ -90,12 +90,10 @@ partitions:
t.Errorf("fixed rule create failed with queue name, err %v",
err)
}
var queue string
- var aclCheck bool
- queue, aclCheck, err = fr.placeApplication(app, queueFunc)
+ queue, err = fr.placeApplication(app, queueFunc)
if queue != "root.testqueue" || err != nil {
t.Errorf("fixed rule failed to place queue in correct queue
'%s', err %v", queue, err)
}
- assert.Check(t, aclCheck, "acls should be checked")
// fixed queue that exists directly in hierarchy
conf = configs.PlacementRule{
@@ -106,11 +104,10 @@ partitions:
if err != nil || fr == nil {
t.Errorf("fixed rule create failed with queue name, err %v",
err)
}
- queue, aclCheck, err = fr.placeApplication(app, queueFunc)
+ queue, err = fr.placeApplication(app, queueFunc)
if queue != "root.testparent.testchild" || err != nil {
t.Errorf("fixed rule failed to place queue in correct queue
'%s', err %v", queue, err)
}
- assert.Check(t, aclCheck, "acls should be checked")
// fixed queue that does not exists
conf = configs.PlacementRule{
@@ -122,11 +119,10 @@ partitions:
if err != nil || fr == nil {
t.Errorf("fixed rule create failed with queue name, err %v",
err)
}
- queue, aclCheck, err = fr.placeApplication(app, queueFunc)
+ queue, err = fr.placeApplication(app, queueFunc)
if queue != "root.newqueue" || err != nil {
t.Errorf("fixed rule failed to place queue in to be created
queue '%s', err %v", queue, err)
}
- assert.Check(t, aclCheck, "acls should be checked")
// trying to place in a parent queue should not fail: failure happens
on create in this case
conf = configs.PlacementRule{
@@ -137,11 +133,10 @@ partitions:
if err != nil || fr == nil {
t.Errorf("fixed rule create failed with queue name, err %v",
err)
}
- queue, aclCheck, err = fr.placeApplication(app, queueFunc)
+ queue, err = fr.placeApplication(app, queueFunc)
if queue != "root.testparent" || err != nil {
t.Errorf("fixed rule did fail with parent queue '%s', error
%v", queue, err)
}
- assert.Check(t, aclCheck, "acls should be checked")
// trying to place in a child using a parent
conf = configs.PlacementRule{
@@ -156,11 +151,10 @@ partitions:
if err != nil || fr == nil {
t.Errorf("fixed rule create failed with queue name, err %v",
err)
}
- queue, aclCheck, err = fr.placeApplication(app, queueFunc)
+ queue, err = fr.placeApplication(app, queueFunc)
if queue != "root.testparent.testchild" || err != nil {
t.Errorf("fixed rule with parent queue should not have failed
'%s', error %v", queue, err)
}
- assert.Check(t, aclCheck, "acls should be checked")
}
func TestFixedRuleParent(t *testing.T) {
@@ -190,12 +184,10 @@ func TestFixedRuleParent(t *testing.T) {
t.Errorf("fixed rule create failed with queue name, err %v",
err)
}
var queue string
- var aclCheck bool
- queue, aclCheck, err = fr.placeApplication(app, queueFunc)
+ queue, err = fr.placeApplication(app, queueFunc)
if queue != "" || err != nil {
t.Errorf("fixed rule with create false for child should have
failed and gave '%s', error %v", queue, err)
}
- assert.Check(t, aclCheck, "acls should be checked")
// trying to place in a child using a non creatable parent
conf = configs.PlacementRule{
@@ -212,11 +204,10 @@ func TestFixedRuleParent(t *testing.T) {
if err != nil || fr == nil {
t.Errorf("fixed rule create failed with queue name, err %v",
err)
}
- queue, aclCheck, err = fr.placeApplication(app, queueFunc)
+ queue, err = fr.placeApplication(app, queueFunc)
if queue != "" || err != nil {
t.Errorf("fixed rule with non existing parent queue should have
failed '%s', error %v", queue, err)
}
- assert.Check(t, aclCheck, "acls should be checked")
// trying to place in a child using a creatable parent
conf = configs.PlacementRule{
@@ -233,11 +224,10 @@ func TestFixedRuleParent(t *testing.T) {
if err != nil || fr == nil {
t.Errorf("fixed rule create failed with queue name, err %v",
err)
}
- queue, aclCheck, err = fr.placeApplication(app, queueFunc)
+ queue, err = fr.placeApplication(app, queueFunc)
if queue != nameParentChild || err != nil {
t.Errorf("fixed rule with non existing parent queue should
created '%s', error %v", queue, err)
}
- assert.Check(t, aclCheck, "acls should be checked")
// trying to place in a child using a parent which is defined as a leaf
conf = configs.PlacementRule{
@@ -253,9 +243,8 @@ func TestFixedRuleParent(t *testing.T) {
if err != nil || fr == nil {
t.Errorf("fixed rule create failed with queue name, err %v",
err)
}
- queue, aclCheck, err = fr.placeApplication(app, queueFunc)
+ queue, err = fr.placeApplication(app, queueFunc)
if queue != "" || err == nil {
t.Errorf("fixed rule with parent declared as leaf should have
failed '%s', error %v", queue, err)
}
- assert.Check(t, aclCheck, "acls should be checked")
}
diff --git a/pkg/scheduler/placement/placement.go
b/pkg/scheduler/placement/placement.go
index b74a4865..ef942464 100644
--- a/pkg/scheduler/placement/placement.go
+++ b/pkg/scheduler/placement/placement.go
@@ -19,18 +19,22 @@
package placement
import (
- "fmt"
+ "errors"
"strings"
"sync"
"go.uber.org/zap"
+ "github.com/apache/yunikorn-core/pkg/common"
"github.com/apache/yunikorn-core/pkg/common/configs"
"github.com/apache/yunikorn-core/pkg/log"
"github.com/apache/yunikorn-core/pkg/scheduler/objects"
"github.com/apache/yunikorn-core/pkg/scheduler/placement/types"
)
+// RejectedError is the standard error returned if placement has failed
+var RejectedError = errors.New("application rejected: no placement rule
matched")
+
type AppPlacementManager struct {
rules []rule
queueFn func(string) *objects.Queue
@@ -113,75 +117,84 @@ func (m *AppPlacementManager) PlaceApplication(app
*objects.Application) error {
defer m.RUnlock()
var queueName string
- var aclCheck bool
var err error
for _, checkRule := range m.rules {
- log.Log(log.Config).Debug("Executing rule for placing
application",
+ log.Log(log.SchedApplication).Debug("Executing rule for placing
application",
zap.String("ruleName", checkRule.getName()),
zap.String("application", app.ApplicationID))
- queueName, aclCheck, err = checkRule.placeApplication(app,
m.queueFn)
+ queueName, err = checkRule.placeApplication(app, m.queueFn)
if err != nil {
- log.Log(log.Config).Error("rule execution failed",
+ log.Log(log.SchedApplication).Error("rule execution
failed",
zap.String("ruleName", checkRule.getName()),
zap.Error(err))
app.SetQueuePath("")
return err
}
- // queueName returned make sure ACL allows access and create
the queueName if not exist
- if queueName != "" {
- // get the queue object
- queue := m.queueFn(queueName)
- // walk up the tree if the queue does not exist
- if queue == nil {
- current := queueName
- for queue == nil {
- current =
current[0:strings.LastIndex(current, configs.DOT)]
- // check if the queue exist
- queue = m.queueFn(current)
- }
- // Check if the user is allowed to submit to
this queueName, if not next rule
- if aclCheck &&
!queue.CheckSubmitAccess(app.GetUser()) {
- log.Log(log.Config).Debug("Submit
access denied on queue",
- zap.String("queueName",
queue.GetQueuePath()),
- zap.String("ruleName",
checkRule.getName()),
- zap.String("application",
app.ApplicationID))
- // reset the queue name for the last
rule in the chain
- queueName = ""
- continue
- }
- } else {
- // Check if this final queue is a leaf queue,
if not next rule
- if !queue.IsLeafQueue() {
- log.Log(log.Config).Debug("Rule
returned parent queue",
- zap.String("queueName",
queueName),
- zap.String("ruleName",
checkRule.getName()),
- zap.String("application",
app.ApplicationID))
- // reset the queue name for the last
rule in the chain
- queueName = ""
- continue
- }
- // Check if the user is allowed to submit to
this queueName, if not next rule
- if aclCheck &&
!queue.CheckSubmitAccess(app.GetUser()) {
- log.Log(log.Config).Debug("Submit
access denied on queue",
- zap.String("queueName",
queueName),
- zap.String("ruleName",
checkRule.getName()),
- zap.String("application",
app.ApplicationID))
- // reset the queue name for the last
rule in the chain
- queueName = ""
- continue
- }
- }
- // we have a queue that allows submitting and can be
created: app placed
+ // no queue name next rule
+ if queueName == "" {
+ continue
+ }
+ // We have the recovery queue bail out: only if we are doing
forced placement
+ // Recovery rule is last in the list. Recovery queue cannot be
returned by other rules.
+ // We do not want to trigger any checks for this queue.
+ if queueName == common.RecoveryQueueFull &&
app.IsCreateForced() {
+ log.Log(log.SchedApplication).Info("Placing application
in recovery queue",
+ zap.String("application", app.ApplicationID))
break
}
+ // queueName returned make sure ACL allows access and set the
queueName in the app
+ queue := m.queueFn(queueName)
+ // walk up the tree if the queue does not exist
+ if queue == nil {
+ current := queueName
+ for queue == nil {
+ current = current[0:strings.LastIndex(current,
configs.DOT)]
+ // check if the queue exist
+ queue = m.queueFn(current)
+ }
+ // Check if the user is allowed to submit to this
queueName, if not next rule
+ if !queue.CheckSubmitAccess(app.GetUser()) {
+ log.Log(log.SchedApplication).Debug("Submit
access denied on queue",
+ zap.String("queueName",
queue.GetQueuePath()),
+ zap.String("ruleName",
checkRule.getName()),
+ zap.String("application",
app.ApplicationID))
+ // reset the queue name for the last rule in
the chain
+ queueName = ""
+ continue
+ }
+ } else {
+ // Check if this final queue is a leaf queue, if not
next rule
+ if !queue.IsLeafQueue() {
+ log.Log(log.SchedApplication).Debug("Rule
returned parent queue",
+ zap.String("queueName", queueName),
+ zap.String("ruleName",
checkRule.getName()),
+ zap.String("application",
app.ApplicationID))
+ // reset the queue name for the last rule in
the chain
+ queueName = ""
+ continue
+ }
+ // Check if the user is allowed to submit to this
queueName, if not next rule
+ if !queue.CheckSubmitAccess(app.GetUser()) {
+ log.Log(log.SchedApplication).Debug("Submit
access denied on queue",
+ zap.String("queueName", queueName),
+ zap.String("ruleName",
checkRule.getName()),
+ zap.String("application",
app.ApplicationID))
+ // reset the queue name for the last rule in
the chain
+ queueName = ""
+ continue
+ }
+ }
+ // we have a queue that allows submitting and can be created:
app placed
+ log.Log(log.SchedApplication).Info("Rule result for placing
application",
+ zap.String("application", app.ApplicationID),
+ zap.String("ruleName", checkRule.getName()),
+ zap.String("queueName", queueName))
+ break
}
- log.Log(log.Config).Debug("Rule result for placing application",
- zap.String("application", app.ApplicationID),
- zap.String("queueName", queueName))
// no more rules to check no queueName found reject placement
if queueName == "" {
app.SetQueuePath("")
- return fmt.Errorf("application rejected: no placement rule
matched")
+ return RejectedError
}
// Add the queue into the application, overriding what was submitted
app.SetQueuePath(queueName)
diff --git a/pkg/scheduler/placement/placement_test.go
b/pkg/scheduler/placement/placement_test.go
index 947fb94f..ac297734 100644
--- a/pkg/scheduler/placement/placement_test.go
+++ b/pkg/scheduler/placement/placement_test.go
@@ -19,13 +19,16 @@
package placement
import (
+ "errors"
"testing"
"gotest.tools/v3/assert"
+ "github.com/apache/yunikorn-core/pkg/common"
"github.com/apache/yunikorn-core/pkg/common/configs"
"github.com/apache/yunikorn-core/pkg/common/security"
"github.com/apache/yunikorn-core/pkg/scheduler/placement/types"
+ siCommon "github.com/apache/yunikorn-scheduler-interface/lib/go/common"
)
// basic test to check if no rules leave the manager unusable
@@ -256,3 +259,87 @@ partitions:
t.Errorf("parent queue: app should not have been placed, queue:
'%s', error: %v", queueName, err)
}
}
+
+func TestForcePlaceApp(t *testing.T) {
+ const (
+ def = "default"
+ defQ = "root.default"
+ )
+
+ // Create the structure for the test
+ // specifically no acl to allow on root
+ data := `
+partitions:
+ - name: default
+ queues:
+ - name: root
+ submitacl: "any-user"
+ queues:
+ - name: default
+ submitacl: "*"
+ - name: acldeny
+ submitacl: " "
+ - name: parent
+ parent: true
+ submitacl: "*"
+`
+ err := initQueueStructure([]byte(data))
+ assert.NilError(t, err, "setting up the queue config failed")
+ // update the manager
+ rules := []configs.PlacementRule{
+ {Name: "provided",
+ Create: false},
+ {Name: "tag",
+ Value: "namespace",
+ Create: true},
+ }
+ man := NewPlacementManager(rules, queueFunc)
+ if man == nil {
+ t.Fatal("placement manager create failed")
+ }
+
+ tags := make(map[string]string)
+ user := security.UserGroup{
+ User: "any-user",
+ Groups: []string{},
+ }
+ deny := security.UserGroup{
+ User: "deny-user",
+ Groups: []string{},
+ }
+ var tests = []struct {
+ name string
+ queue string
+ placed string
+ tags map[string]string
+ user security.UserGroup
+ }{
+ {"empty", "", "", tags, user},
+ {"provided unqualified", def, defQ, tags, user},
+ {"provided qualified", defQ, defQ, tags, user},
+ {"provided not exist", "unknown", "", tags, user},
+ {"provided parent", "root.parent", "", tags, user},
+ {"acl deny", "root.acldeny", "", tags, deny},
+ {"create", "unknown", "root.namespace",
map[string]string{"namespace": "namespace"}, user},
+ {"deny create", "unknown", "", map[string]string{"namespace":
"namespace"}, deny},
+ {"forced exist", defQ, defQ,
map[string]string{siCommon.AppTagCreateForce: "true"}, user},
+ {"forced and create", "unknown", "root.namespace",
map[string]string{siCommon.AppTagCreateForce: "true", "namespace":
"namespace"}, user},
+ {"forced and deny create", "unknown", common.RecoveryQueueFull,
map[string]string{siCommon.AppTagCreateForce: "true", "namespace":
"namespace"}, deny},
+ {"forced parent", "root.parent", common.RecoveryQueueFull,
map[string]string{siCommon.AppTagCreateForce: "true"}, user},
+ {"forced acl deny", "root.acldeny", common.RecoveryQueueFull,
map[string]string{siCommon.AppTagCreateForce: "true"}, deny},
+ {"forced not exist", "unknown", common.RecoveryQueueFull,
map[string]string{siCommon.AppTagCreateForce: "true"}, user},
+ {"forced not exist acl deny", "unknown",
common.RecoveryQueueFull, map[string]string{siCommon.AppTagCreateForce:
"true"}, deny},
+ }
+ for _, tt := range tests {
+ t.Run(tt.name, func(t *testing.T) {
+ app := newApplication("app1", "default", tt.queue,
tt.user, tt.tags, nil, "")
+ err = man.PlaceApplication(app)
+ if tt.placed == "" {
+ assert.Assert(t, errors.Is(err, RejectedError),
"unexpected error or no error returned")
+ } else {
+ assert.NilError(t, err, "unexpected placement
failure")
+ assert.Equal(t, tt.placed, app.GetQueuePath(),
"incorrect queue set")
+ }
+ })
+ }
+}
diff --git a/pkg/scheduler/placement/provided_rule.go
b/pkg/scheduler/placement/provided_rule.go
index e7fa15f5..27974a78 100644
--- a/pkg/scheduler/placement/provided_rule.go
+++ b/pkg/scheduler/placement/provided_rule.go
@@ -52,35 +52,34 @@ func (pr *providedRule) initialise(conf
configs.PlacementRule) error {
return err
}
-func (pr *providedRule) placeApplication(app *objects.Application, queueFn
func(string) *objects.Queue) (string, bool, error) {
+func (pr *providedRule) placeApplication(app *objects.Application, queueFn
func(string) *objects.Queue) (string, error) {
// since this is the provided rule we must have a queue in the info
already
queueName := app.GetQueuePath()
if queueName == "" {
- return "", true, nil
+ return "", nil
}
// before anything run the filter
if !pr.filter.allowUser(app.GetUser()) {
- log.Log(log.Config).Debug("Provided rule filtered",
+ log.Log(log.SchedApplication).Debug("Provided rule filtered",
zap.String("application", app.ApplicationID),
zap.Any("user", app.GetUser()))
- return "", true, nil
+ return "", nil
}
var parentName string
- var aclCheck = true
var err error
// if we have a fully qualified queue passed in do not run the parent
rule
if !strings.HasPrefix(queueName, configs.RootQueue+configs.DOT) {
// run the parent rule if set
if pr.parent != nil {
- parentName, aclCheck, err =
pr.parent.placeApplication(app, queueFn)
+ parentName, err = pr.parent.placeApplication(app,
queueFn)
// failed parent rule, fail this rule
if err != nil {
- return "", aclCheck, err
+ return "", err
}
// rule did not return a parent: this could be filter
or create flag related
if parentName == "" {
- return "", aclCheck, nil
+ return "", nil
}
// check if this is a parent queue and qualify it
if !strings.HasPrefix(parentName,
configs.RootQueue+configs.DOT) {
@@ -89,7 +88,7 @@ func (pr *providedRule) placeApplication(app
*objects.Application, queueFn func(
// if the parent queue exists it cannot be a leaf
parentQueue := queueFn(parentName)
if parentQueue != nil && parentQueue.IsLeafQueue() {
- return "", aclCheck, fmt.Errorf("parent rule
returned a leaf queue: %s", parentName)
+ return "", fmt.Errorf("parent rule returned a
leaf queue: %s", parentName)
}
}
// the parent is set from the rule otherwise set it to the root
@@ -99,17 +98,18 @@ func (pr *providedRule) placeApplication(app
*objects.Application, queueFn func(
// Make it a fully qualified queue
queueName = parentName + configs.DOT + replaceDot(queueName)
}
- log.Log(log.Config).Debug("Provided rule intermediate result",
+ // Log the result before we check the create flag
+ log.Log(log.SchedApplication).Debug("Provided rule intermediate result",
zap.String("application", app.ApplicationID),
zap.String("queue", queueName))
// get the queue object
queue := queueFn(queueName)
// if we cannot create the queue must exist
if !pr.create && queue == nil {
- return "", aclCheck, nil
+ return "", nil
}
- log.Log(log.Config).Info("Provided rule application placed",
+ log.Log(log.SchedApplication).Info("Provided rule application placed",
zap.String("application", app.ApplicationID),
zap.String("queue", queueName))
- return queueName, aclCheck, nil
+ return queueName, nil
}
diff --git a/pkg/scheduler/placement/provided_rule_test.go
b/pkg/scheduler/placement/provided_rule_test.go
index e400c922..5077c834 100644
--- a/pkg/scheduler/placement/provided_rule_test.go
+++ b/pkg/scheduler/placement/provided_rule_test.go
@@ -57,26 +57,22 @@ partitions:
// queue that does not exists directly under the root
appInfo := newApplication("app1", "default", "unknown", user, tags,
nil, "")
var queue string
- var aclCheck bool
- queue, aclCheck, err = pr.placeApplication(appInfo, queueFunc)
+ queue, err = pr.placeApplication(appInfo, queueFunc)
if queue != "" || err != nil {
t.Errorf("provided rule placed app in incorrect queue '%s', err
%v", queue, err)
}
- assert.Check(t, aclCheck, "acls should be checked")
// trying to place when no queue provided in the app
appInfo = newApplication("app1", "default", "", user, tags, nil, "")
- queue, aclCheck, err = pr.placeApplication(appInfo, queueFunc)
+ queue, err = pr.placeApplication(appInfo, queueFunc)
if queue != "" || err != nil {
t.Errorf("provided rule placed app in incorrect queue '%s',
error %v", queue, err)
}
- assert.Check(t, aclCheck, "acls should be checked")
// trying to place in a qualified queue that does not exist
appInfo = newApplication("app1", "default", "root.unknown", user, tags,
nil, "")
- queue, aclCheck, err = pr.placeApplication(appInfo, queueFunc)
+ queue, err = pr.placeApplication(appInfo, queueFunc)
if queue != "" || err != nil {
t.Errorf("provided rule placed app in incorrect queue '%s',
error %v", queue, err)
}
- assert.Check(t, aclCheck, "acls should be checked")
// same queue now with create flag
conf = configs.PlacementRule{
Name: "provided",
@@ -86,11 +82,10 @@ partitions:
if err != nil || pr == nil {
t.Errorf("provided rule create failed, err %v", err)
}
- queue, aclCheck, err = pr.placeApplication(appInfo, queueFunc)
+ queue, err = pr.placeApplication(appInfo, queueFunc)
if queue != "root.unknown" || err != nil {
t.Errorf("provided rule placed app in incorrect queue '%s',
error %v", queue, err)
}
- assert.Check(t, aclCheck, "acls should be checked")
conf = configs.PlacementRule{
Name: "provided",
@@ -106,20 +101,18 @@ partitions:
// unqualified queue with parent rule that exists directly in hierarchy
appInfo = newApplication("app1", "default", "testchild", user, tags,
nil, "")
- queue, aclCheck, err = pr.placeApplication(appInfo, queueFunc)
+ queue, err = pr.placeApplication(appInfo, queueFunc)
if queue != "root.testparent.testchild" || err != nil {
t.Errorf("provided rule failed to place queue in correct queue
'%s', err %v", queue, err)
}
- assert.Check(t, aclCheck, "acls should be checked")
// qualified queue with parent rule (parent rule ignored)
appInfo = newApplication("app1", "default", "root.testparent", user,
tags, nil, "")
- queue, aclCheck, err = pr.placeApplication(appInfo, queueFunc)
+ queue, err = pr.placeApplication(appInfo, queueFunc)
if queue != "root.testparent" || err != nil {
t.Errorf("provided rule placed in to be created queue with
create false '%s', err %v", queue, err)
}
- assert.Check(t, aclCheck, "acls should be checked")
}
func TestProvidedRuleParent(t *testing.T) {
@@ -149,12 +142,10 @@ func TestProvidedRuleParent(t *testing.T) {
appInfo := newApplication("app1", "default", "unknown", user, tags,
nil, "")
var queue string
- var aclCheck bool
- queue, aclCheck, err = pr.placeApplication(appInfo, queueFunc)
+ queue, err = pr.placeApplication(appInfo, queueFunc)
if queue != "" || err != nil {
t.Errorf("provided rule placed app in incorrect queue '%s', err
%v", queue, err)
}
- assert.Check(t, aclCheck, "acls should be checked")
// trying to place in a child using a non creatable parent
conf = configs.PlacementRule{
@@ -172,11 +163,10 @@ func TestProvidedRuleParent(t *testing.T) {
}
appInfo = newApplication("app1", "default", "testchild", user, tags,
nil, "")
- queue, aclCheck, err = pr.placeApplication(appInfo, queueFunc)
+ queue, err = pr.placeApplication(appInfo, queueFunc)
if queue != "" || err != nil {
t.Errorf("provided rule placed app in incorrect queue '%s', err
%v", queue, err)
}
- assert.Check(t, aclCheck, "acls should be checked")
// trying to place in a child using a creatable parent
conf = configs.PlacementRule{
@@ -192,11 +182,10 @@ func TestProvidedRuleParent(t *testing.T) {
if err != nil || pr == nil {
t.Errorf("provided rule create failed, err %v", err)
}
- queue, aclCheck, err = pr.placeApplication(appInfo, queueFunc)
+ queue, err = pr.placeApplication(appInfo, queueFunc)
if queue != nameParentChild || err != nil {
t.Errorf("provided rule with non existing parent queue should
create '%s', error %v", queue, err)
}
- assert.Check(t, aclCheck, "acls should be checked")
// trying to place in a child using a parent which is defined as a leaf
conf = configs.PlacementRule{
@@ -213,9 +202,8 @@ func TestProvidedRuleParent(t *testing.T) {
}
appInfo = newApplication("app1", "default", "unknown", user, tags, nil,
"")
- queue, aclCheck, err = pr.placeApplication(appInfo, queueFunc)
+ queue, err = pr.placeApplication(appInfo, queueFunc)
if queue != "" || err == nil {
t.Errorf("provided rule placed app in incorrect queue '%s', err
%v", queue, err)
}
- assert.Check(t, aclCheck, "acls should be checked")
}
diff --git a/pkg/scheduler/placement/recovery_rule.go
b/pkg/scheduler/placement/recovery_rule.go
index 71f7dc23..4e2101b8 100644
--- a/pkg/scheduler/placement/recovery_rule.go
+++ b/pkg/scheduler/placement/recovery_rule.go
@@ -28,31 +28,31 @@ import (
"github.com/apache/yunikorn-core/pkg/scheduler/placement/types"
)
+// A rule to place an application into the recovery queue if no other rules
matched and application submission is forced.
+// This rule will be run implicitly after all other placement rules are
evaluated to ensure that an application
+// corresponding to an already-executing workload can be accepted successfully.
type recoveryRule struct {
basicRule
}
-// A rule to place an application into the recovery queue if no other rules
matched and application submission is forced.
-// This rule will be run implicitly after all other placement rules are
evaluated to ensure that an application
-// corresponding to an already-executing workload can be accepted successfully.
func (rr *recoveryRule) getName() string {
return types.Recovery
}
-func (rr *recoveryRule) initialise(conf configs.PlacementRule) error {
+func (rr *recoveryRule) initialise(_ configs.PlacementRule) error {
// no configuration needed for the recovery rule
return nil
}
-func (rr *recoveryRule) placeApplication(app *objects.Application, _
func(string) *objects.Queue) (string, bool, error) {
+func (rr *recoveryRule) placeApplication(app *objects.Application, _
func(string) *objects.Queue) (string, error) {
// only forced applications should resolve to the recovery queue
if !app.IsCreateForced() {
- return "", false, nil
+ return "", nil
}
queueName := common.RecoveryQueueFull
- log.Log(log.Config).Info("Recovery rule application placed",
+ log.Log(log.SchedApplication).Info("Recovery rule application placed",
zap.String("application", app.ApplicationID),
zap.String("queue", queueName))
- return queueName, false, nil
+ return queueName, nil
}
diff --git a/pkg/scheduler/placement/recovery_rule_test.go
b/pkg/scheduler/placement/recovery_rule_test.go
index 7185396d..75f005da 100644
--- a/pkg/scheduler/placement/recovery_rule_test.go
+++ b/pkg/scheduler/placement/recovery_rule_test.go
@@ -63,18 +63,15 @@ partitions:
app := newApplication("app1", "default", "ignored", user, tags, nil, "")
var queue string
- var aclCheck bool
- queue, aclCheck, err = rr.placeApplication(app, queueFunc)
+ queue, err = rr.placeApplication(app, queueFunc)
if queue != "" || err != nil {
t.Errorf("recovery rule did not bypass non-forced application,
resolved queue '%s', err %v ", queue, err)
}
- assert.Check(t, !aclCheck, "acl check should not be set for recovery
rule")
tags[siCommon.AppTagCreateForce] = "true"
app = newApplication("app1", "default", "ignored", user, tags, nil, "")
- queue, aclCheck, err = rr.placeApplication(app, queueFunc)
+ queue, err = rr.placeApplication(app, queueFunc)
if queue != common.RecoveryQueueFull || err != nil {
t.Errorf("recovery rule did not place forced application into
recovery queue, resolved queue '%s', err %v ", queue, err)
}
- assert.Check(t, !aclCheck, "acl check should not be set for recovery
rule")
}
diff --git a/pkg/scheduler/placement/rule.go b/pkg/scheduler/placement/rule.go
index 84e1d31a..1816a52d 100644
--- a/pkg/scheduler/placement/rule.go
+++ b/pkg/scheduler/placement/rule.go
@@ -38,9 +38,8 @@ type rule interface {
// Execute the rule and return the queue getName the application is
placed in.
// Returns the fully qualified queue getName if the rule finds a queue
or an empty string if the rule did not match.
- // Additionally, returns true if ACLs should be checked or false
otherwise.
// The error must only be set if there is a failure while executing the
rule not if the rule did not match.
- placeApplication(app *objects.Application, queueFn func(string)
*objects.Queue) (string, bool, error)
+ placeApplication(app *objects.Application, queueFn func(string)
*objects.Queue) (string, error)
// Return the getName of the rule which is defined in the rule.
// The basicRule provides a "unnamed rule" implementation.
diff --git a/pkg/scheduler/placement/rule_test.go
b/pkg/scheduler/placement/rule_test.go
index c6d4e49f..f4e741c8 100644
--- a/pkg/scheduler/placement/rule_test.go
+++ b/pkg/scheduler/placement/rule_test.go
@@ -61,36 +61,31 @@ func TestPlaceApp(t *testing.T) {
}
nr, err := newRule(conf)
assert.NilError(t, err, "unexpected rule initialisation error")
- var aclCheck bool
// place application that should fail
- _, aclCheck, err = nr.placeApplication(nil, nil)
+ _, err = nr.placeApplication(nil, nil)
if err == nil {
t.Error("test rule place application did not fail as expected")
}
- assert.Check(t, aclCheck, "acls should be checked")
var queue string
// place application that should not fail and return "test"
- queue, aclCheck, err = nr.placeApplication(&objects.Application{}, nil)
+ queue, err = nr.placeApplication(&objects.Application{}, nil)
if err != nil || queue != "test" {
t.Errorf("test rule place application did not fail, err: %v, ",
err)
}
- assert.Check(t, aclCheck, "acls should be checked")
// place application that should not fail and return the queue in the
object
app := &objects.Application{}
app.SetQueuePath("passedin")
- queue, aclCheck, err = nr.placeApplication(app, nil)
+ queue, err = nr.placeApplication(app, nil)
if err != nil || queue != "passedin" {
t.Errorf("test rule place application did not fail, err: %v, ",
err)
}
- assert.Check(t, aclCheck, "acls should be checked")
// place application that should not fail and return the queue in the
object
app = &objects.Application{}
app.SetQueuePath("user.name")
- queue, aclCheck, err = nr.placeApplication(app, nil)
+ queue, err = nr.placeApplication(app, nil)
if err != nil || queue != "user_dot_name" {
t.Errorf("test rule place application did not fail, err: %v, ",
err)
}
- assert.Check(t, aclCheck, "acls should be checked")
}
func TestReplaceDot(t *testing.T) {
diff --git a/pkg/scheduler/placement/tag_rule.go
b/pkg/scheduler/placement/tag_rule.go
index 1282a86a..95f38ea3 100644
--- a/pkg/scheduler/placement/tag_rule.go
+++ b/pkg/scheduler/placement/tag_rule.go
@@ -57,36 +57,35 @@ func (tr *tagRule) initialise(conf configs.PlacementRule)
error {
return err
}
-func (tr *tagRule) placeApplication(app *objects.Application, queueFn
func(string) *objects.Queue) (string, bool, error) {
+func (tr *tagRule) placeApplication(app *objects.Application, queueFn
func(string) *objects.Queue) (string, error) {
// if the tag is not present we can skipp all other processing
tagVal := app.GetTag(tr.tagName)
if tagVal == "" {
- return "", true, nil
+ return "", nil
}
// before anything run the filter
if !tr.filter.allowUser(app.GetUser()) {
- log.Log(log.Config).Debug("Tag rule filtered",
+ log.Log(log.SchedApplication).Debug("Tag rule filtered",
zap.String("application", app.ApplicationID),
zap.Any("user", app.GetUser()),
zap.String("tagName", tr.tagName))
- return "", true, nil
+ return "", nil
}
var parentName string
- var aclCheck = true
var err error
queueName := tagVal
// if we have a fully qualified queue in the value do not run the
parent rule
if !strings.HasPrefix(queueName, configs.RootQueue+configs.DOT) {
// run the parent rule if set
if tr.parent != nil {
- parentName, aclCheck, err =
tr.parent.placeApplication(app, queueFn)
+ parentName, err = tr.parent.placeApplication(app,
queueFn)
// failed parent rule, fail this rule
if err != nil {
- return "", aclCheck, err
+ return "", err
}
// rule did not match: this could be filter or create
flag related
if parentName == "" {
- return "", aclCheck, nil
+ return "", nil
}
// check if this is a parent queue and qualify it
if !strings.HasPrefix(parentName,
configs.RootQueue+configs.DOT) {
@@ -95,7 +94,7 @@ func (tr *tagRule) placeApplication(app *objects.Application,
queueFn func(strin
// if the parent queue exists it cannot be a leaf
parentQueue := queueFn(parentName)
if parentQueue != nil && parentQueue.IsLeafQueue() {
- return "", aclCheck, fmt.Errorf("parent rule
returned a leaf queue: %s", parentName)
+ return "", fmt.Errorf("parent rule returned a
leaf queue: %s", parentName)
}
}
// the parent is set from the rule otherwise set it to the root
@@ -104,17 +103,18 @@ func (tr *tagRule) placeApplication(app
*objects.Application, queueFn func(strin
}
queueName = parentName + configs.DOT + replaceDot(tagVal)
}
- log.Log(log.Config).Debug("Tag rule intermediate result",
+ // Log the result before we check the create flag
+ log.Log(log.SchedApplication).Debug("Tag rule intermediate result",
zap.String("application", app.ApplicationID),
zap.String("queue", queueName))
// get the queue object
queue := queueFn(queueName)
// if we cannot create the queue it must exist, rule does not match
otherwise
if !tr.create && queue == nil {
- return "", aclCheck, nil
+ return "", nil
}
- log.Log(log.Config).Info("Tag rule application placed",
+ log.Log(log.SchedApplication).Info("Tag rule application placed",
zap.String("application", app.ApplicationID),
zap.String("queue", queueName))
- return queueName, aclCheck, nil
+ return queueName, nil
}
diff --git a/pkg/scheduler/placement/tag_rule_test.go
b/pkg/scheduler/placement/tag_rule_test.go
index 251cf81a..6eb8e79e 100644
--- a/pkg/scheduler/placement/tag_rule_test.go
+++ b/pkg/scheduler/placement/tag_rule_test.go
@@ -90,48 +90,42 @@ partitions:
tags := make(map[string]string)
appInfo := newApplication("app1", "default", "ignored", user, tags,
nil, "")
var queue string
- var aclCheck bool
- queue, aclCheck, err = tr.placeApplication(appInfo, queueFunc)
+ queue, err = tr.placeApplication(appInfo, queueFunc)
if queue != "" || err != nil {
t.Errorf("tag rule failed with no tag value '%s', err %v",
queue, err)
}
- assert.Check(t, aclCheck, "acls should be checked")
// tag queue that exists directly in hierarchy
tags = map[string]string{"label1": "testqueue"}
appInfo = newApplication("app1", "default", "ignored", user, tags, nil,
"")
- queue, aclCheck, err = tr.placeApplication(appInfo, queueFunc)
+ queue, err = tr.placeApplication(appInfo, queueFunc)
if queue != "root.testqueue" || err != nil {
t.Errorf("tag rule failed to place queue in correct queue '%s',
err %v", queue, err)
}
- assert.Check(t, aclCheck, "acls should be checked")
// tag queue that does not exists
tags = map[string]string{"label1": "unknown"}
appInfo = newApplication("app1", "default", "ignored", user, tags, nil,
"")
- queue, aclCheck, err = tr.placeApplication(appInfo, queueFunc)
+ queue, err = tr.placeApplication(appInfo, queueFunc)
if queue != "" || err != nil {
t.Errorf("tag rule placed in queue that does not exists '%s',
err %v", queue, err)
}
- assert.Check(t, aclCheck, "acls should be checked")
// tag queue fully qualified
tags = map[string]string{"label1": "root.testparent.testchild"}
appInfo = newApplication("app1", "default", "ignored", user, tags, nil,
"")
- queue, aclCheck, err = tr.placeApplication(appInfo, queueFunc)
+ queue, err = tr.placeApplication(appInfo, queueFunc)
if queue != "root.testparent.testchild" || err != nil {
t.Errorf("tag rule did fail with qualified queue '%s', error
%v", queue, err)
}
- assert.Check(t, aclCheck, "acls should be checked")
// tag queue references recovery
tags = map[string]string{"label1": common.RecoveryQueueFull}
appInfo = newApplication("app1", "default", "ignored", user, tags, nil,
"")
- queue, aclCheck, err = tr.placeApplication(appInfo, queueFunc)
+ queue, err = tr.placeApplication(appInfo, queueFunc)
if queue != "" || err != nil {
t.Errorf("tag rule failed with explicit recovery queue: queue
'%s', error %v", queue, err)
}
- assert.Check(t, aclCheck, "acls should be checked")
// trying to place in a child using a parent
conf = configs.PlacementRule{
@@ -148,18 +142,16 @@ partitions:
}
tags = map[string]string{"label1": "testchild"}
appInfo = newApplication("app1", "default", "ignored", user, tags, nil,
"")
- queue, aclCheck, err = tr.placeApplication(appInfo, queueFunc)
+ queue, err = tr.placeApplication(appInfo, queueFunc)
if queue != "" || err != nil {
t.Errorf("tag rule with parent queue should have failed value
not set '%s', error %v", queue, err)
}
- assert.Check(t, aclCheck, "acls should be checked")
tags = map[string]string{"label1": "testchild", "label2": "testparent"}
appInfo = newApplication("app1", "default", "ignored", user, tags, nil,
"")
- queue, aclCheck, err = tr.placeApplication(appInfo, queueFunc)
+ queue, err = tr.placeApplication(appInfo, queueFunc)
if queue != "root.testparent.testchild" || err != nil {
t.Errorf("tag rule with parent queue incorrect queue '%s',
error %v", queue, err)
}
- assert.Check(t, aclCheck, "acls should be checked")
}
func TestTagRuleParent(t *testing.T) {
@@ -190,12 +182,10 @@ func TestTagRuleParent(t *testing.T) {
tags := map[string]string{"label1": "testchild", "label2": "testparent"}
appInfo := newApplication("app1", "default", "unknown", user, tags,
nil, "")
var queue string
- var aclCheck bool
- queue, aclCheck, err = ur.placeApplication(appInfo, queueFunc)
+ queue, err = ur.placeApplication(appInfo, queueFunc)
if queue != "" || err != nil {
t.Errorf("tag rule placed app in incorrect queue '%s', err %v",
queue, err)
}
- assert.Check(t, aclCheck, "acls should be checked")
// trying to place in a child using a non creatable parent
conf = configs.PlacementRule{
@@ -215,11 +205,10 @@ func TestTagRuleParent(t *testing.T) {
tags = map[string]string{"label1": "testchild", "label2":
"testparentnew"}
appInfo = newApplication("app1", "default", "unknown", user, tags, nil,
"")
- queue, aclCheck, err = ur.placeApplication(appInfo, queueFunc)
+ queue, err = ur.placeApplication(appInfo, queueFunc)
if queue != "" || err != nil {
t.Errorf("tag rule placed app in incorrect queue '%s', err %v",
queue, err)
}
- assert.Check(t, aclCheck, "acls should be checked")
// trying to place in a child using a creatable parent
conf = configs.PlacementRule{
@@ -236,11 +225,10 @@ func TestTagRuleParent(t *testing.T) {
if err != nil || ur == nil {
t.Errorf("tag rule create failed with queue name, err %v", err)
}
- queue, aclCheck, err = ur.placeApplication(appInfo, queueFunc)
+ queue, err = ur.placeApplication(appInfo, queueFunc)
if queue != nameParentChild || err != nil {
t.Errorf("user rule with non existing parent queue should
create '%s', error %v", queue, err)
}
- assert.Check(t, aclCheck, "acls should be checked")
// trying to place in a child using a parent which is defined as a leaf
conf = configs.PlacementRule{
@@ -258,9 +246,8 @@ func TestTagRuleParent(t *testing.T) {
}
appInfo = newApplication("app1", "default", "unknown", user, tags, nil,
"")
- queue, aclCheck, err = ur.placeApplication(appInfo, queueFunc)
+ queue, err = ur.placeApplication(appInfo, queueFunc)
if queue != "" || err == nil {
t.Errorf("tag rule placed app in incorrect queue '%s', err %v",
queue, err)
}
- assert.Check(t, aclCheck, "acls should be checked")
}
diff --git a/pkg/scheduler/placement/testrule.go
b/pkg/scheduler/placement/testrule.go
index 7554d42d..fce27c68 100644
--- a/pkg/scheduler/placement/testrule.go
+++ b/pkg/scheduler/placement/testrule.go
@@ -48,12 +48,12 @@ func (tr *testRule) initialise(conf configs.PlacementRule)
error {
}
// Simple test rule that just checks the app passed in and returns fixed queue
names.
-func (tr *testRule) placeApplication(app *objects.Application, queueFn
func(string) *objects.Queue) (string, bool, error) {
+func (tr *testRule) placeApplication(app *objects.Application, queueFn
func(string) *objects.Queue) (string, error) {
if app == nil {
- return "", true, fmt.Errorf("nil app passed in")
+ return "", fmt.Errorf("nil app passed in")
}
if queuePath := app.GetQueuePath(); queuePath != "" {
- return replaceDot(queuePath), true, nil
+ return replaceDot(queuePath), nil
}
- return types.Test, true, nil
+ return types.Test, nil
}
diff --git a/pkg/scheduler/placement/user_rule.go
b/pkg/scheduler/placement/user_rule.go
index f7113cf2..2c195e4e 100644
--- a/pkg/scheduler/placement/user_rule.go
+++ b/pkg/scheduler/placement/user_rule.go
@@ -49,28 +49,27 @@ func (ur *userRule) initialise(conf configs.PlacementRule)
error {
return err
}
-func (ur *userRule) placeApplication(app *objects.Application, queueFn
func(string) *objects.Queue) (string, bool, error) {
+func (ur *userRule) placeApplication(app *objects.Application, queueFn
func(string) *objects.Queue) (string, error) {
// before anything run the filter
userName := app.GetUser().User
if !ur.filter.allowUser(app.GetUser()) {
- log.Log(log.Config).Debug("User rule filtered",
+ log.Log(log.SchedApplication).Debug("User rule filtered",
zap.String("application", app.ApplicationID),
zap.Any("user", app.GetUser()))
- return "", true, nil
+ return "", nil
}
var parentName string
- var aclCheck = true
var err error
// run the parent rule if set
if ur.parent != nil {
- parentName, aclCheck, err = ur.parent.placeApplication(app,
queueFn)
+ parentName, err = ur.parent.placeApplication(app, queueFn)
// failed parent rule, fail this rule
if err != nil {
- return "", aclCheck, err
+ return "", err
}
// rule did not match: this could be filter or create flag
related
if parentName == "" {
- return "", aclCheck, nil
+ return "", nil
}
// check if this is a parent queue and qualify it
if !strings.HasPrefix(parentName,
configs.RootQueue+configs.DOT) {
@@ -79,7 +78,7 @@ func (ur *userRule) placeApplication(app
*objects.Application, queueFn func(stri
// if the parent queue exists it cannot be a leaf
parentQueue := queueFn(parentName)
if parentQueue != nil && parentQueue.IsLeafQueue() {
- return "", aclCheck, fmt.Errorf("parent rule returned a
leaf queue: %s", parentName)
+ return "", fmt.Errorf("parent rule returned a leaf
queue: %s", parentName)
}
}
// the parent is set from the rule otherwise set it to the root
@@ -87,17 +86,18 @@ func (ur *userRule) placeApplication(app
*objects.Application, queueFn func(stri
parentName = configs.RootQueue
}
queueName := parentName + configs.DOT + replaceDot(userName)
- log.Log(log.Config).Debug("User rule intermediate result",
+ // Log the result before we check the create flag
+ log.Log(log.SchedApplication).Debug("User rule intermediate result",
zap.String("application", app.ApplicationID),
zap.String("queue", queueName))
// get the queue object
queue := queueFn(queueName)
// if we cannot create the queue it must exist, rule does not match
otherwise
if !ur.create && queue == nil {
- return "", aclCheck, nil
+ return "", nil
}
- log.Log(log.Config).Info("User rule application placed",
+ log.Log(log.SchedApplication).Info("User rule application placed",
zap.String("application", app.ApplicationID),
zap.String("queue", queueName))
- return queueName, aclCheck, nil
+ return queueName, nil
}
diff --git a/pkg/scheduler/placement/user_rule_test.go
b/pkg/scheduler/placement/user_rule_test.go
index c38a89ff..9e7126e3 100644
--- a/pkg/scheduler/placement/user_rule_test.go
+++ b/pkg/scheduler/placement/user_rule_test.go
@@ -59,34 +59,30 @@ partitions:
t.Errorf("user rule create failed, err %v", err)
}
var queue string
- var aclCheck bool
- queue, aclCheck, err = ur.placeApplication(appInfo, queueFunc)
+ queue, err = ur.placeApplication(appInfo, queueFunc)
if queue != "root.testchild" || err != nil {
t.Errorf("user rule failed to place queue in correct queue
'%s', err %v", queue, err)
}
- assert.Check(t, aclCheck, "acls should be checked")
// trying to place in a parent queue should fail on queue create not in
the rule
user = security.UserGroup{
User: "testparent",
Groups: []string{},
}
appInfo = newApplication("app1", "default", "ignored", user, tags, nil,
"")
- queue, aclCheck, err = ur.placeApplication(appInfo, queueFunc)
+ queue, err = ur.placeApplication(appInfo, queueFunc)
if queue != "root.testparent" || err != nil {
t.Errorf("user rule failed with parent queue '%s', error %v",
queue, err)
}
- assert.Check(t, aclCheck, "acls should be checked")
user = security.UserGroup{
User: "test.user",
Groups: []string{},
}
appInfo = newApplication("app1", "default", "ignored", user, tags, nil,
"")
- queue, aclCheck, err = ur.placeApplication(appInfo, queueFunc)
+ queue, err = ur.placeApplication(appInfo, queueFunc)
if queue == "" || err != nil {
t.Errorf("user rule with dotted user should not have failed
'%s', error %v", queue, err)
}
- assert.Check(t, aclCheck, "acls should be checked")
// user queue that exists directly in hierarchy
conf = configs.PlacementRule{
@@ -105,11 +101,10 @@ partitions:
if err != nil || ur == nil {
t.Errorf("user rule create failed with queue name, err %v", err)
}
- queue, aclCheck, err = ur.placeApplication(appInfo, queueFunc)
+ queue, err = ur.placeApplication(appInfo, queueFunc)
if queue != "root.testparent.testchild" || err != nil {
t.Errorf("user rule failed to place queue in correct queue
'%s', err %v", queue, err)
}
- assert.Check(t, aclCheck, "acls should be checked")
// user queue that does not exists
user = security.UserGroup{
@@ -126,11 +121,10 @@ partitions:
if err != nil || ur == nil {
t.Errorf("user rule create failed with queue name, err %v", err)
}
- queue, aclCheck, err = ur.placeApplication(appInfo, queueFunc)
+ queue, err = ur.placeApplication(appInfo, queueFunc)
if queue != "root.unknown" || err != nil {
t.Errorf("user rule placed in to be created queue with create
false '%s', err %v", queue, err)
}
- assert.Check(t, aclCheck, "acls should be checked")
}
func TestUserRuleParent(t *testing.T) {
@@ -160,12 +154,10 @@ func TestUserRuleParent(t *testing.T) {
appInfo := newApplication("app1", "default", "unknown", user, tags,
nil, "")
var queue string
- var aclCheck bool
- queue, aclCheck, err = ur.placeApplication(appInfo, queueFunc)
+ queue, err = ur.placeApplication(appInfo, queueFunc)
if queue != "" || err != nil {
t.Errorf("user rule placed app in incorrect queue '%s', err
%v", queue, err)
}
- assert.Check(t, aclCheck, "acls should be checked")
// trying to place in a child using a non creatable parent
conf = configs.PlacementRule{
@@ -183,11 +175,10 @@ func TestUserRuleParent(t *testing.T) {
}
appInfo = newApplication("app1", "default", "unknown", user, tags, nil,
"")
- queue, aclCheck, err = ur.placeApplication(appInfo, queueFunc)
+ queue, err = ur.placeApplication(appInfo, queueFunc)
if queue != "" || err != nil {
t.Errorf("user rule placed app in incorrect queue '%s', err
%v", queue, err)
}
- assert.Check(t, aclCheck, "acls should be checked")
// trying to place in a child using a creatable parent
conf = configs.PlacementRule{
@@ -203,11 +194,10 @@ func TestUserRuleParent(t *testing.T) {
if err != nil || ur == nil {
t.Errorf("user rule create failed with queue name, err %v", err)
}
- queue, aclCheck, err = ur.placeApplication(appInfo, queueFunc)
+ queue, err = ur.placeApplication(appInfo, queueFunc)
if queue != nameParentChild || err != nil {
t.Errorf("user rule with non existing parent queue should
create '%s', error %v", queue, err)
}
- assert.Check(t, aclCheck, "acls should be checked")
// trying to place in a child using a parent which is defined as a leaf
conf = configs.PlacementRule{
@@ -224,9 +214,8 @@ func TestUserRuleParent(t *testing.T) {
}
appInfo = newApplication("app1", "default", "unknown", user, tags, nil,
"")
- queue, aclCheck, err = ur.placeApplication(appInfo, queueFunc)
+ queue, err = ur.placeApplication(appInfo, queueFunc)
if queue != "" || err == nil {
t.Errorf("user rule placed app in incorrect queue '%s', err
%v", queue, err)
}
- assert.Check(t, aclCheck, "acls should be checked")
}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]