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 88b2966a [YUNIKORN-3446] Race condition in ensureGroupTrackerForApp
creates duplicate GroupTrackers and bypasses group quota (#1154)
88b2966a is described below
commit 88b2966a81aaca3dd4f099b22e0676b991a327e0
Author: hedger9487 <[email protected]>
AuthorDate: Wed Sep 9 12:51:38 2026 +0530
[YUNIKORN-3446] Race condition in ensureGroupTrackerForApp creates
duplicate GroupTrackers and bypasses group quota (#1154)
Closes: #1154
Signed-off-by: mani <[email protected]>
---
pkg/scheduler/ugm/manager.go | 25 +++++++++------
pkg/scheduler/ugm/manager_test.go | 66 +++++++++++++++++++++++++++++++++++++++
2 files changed, 81 insertions(+), 10 deletions(-)
diff --git a/pkg/scheduler/ugm/manager.go b/pkg/scheduler/ugm/manager.go
index b7d7475b..35aca540 100644
--- a/pkg/scheduler/ugm/manager.go
+++ b/pkg/scheduler/ugm/manager.go
@@ -240,16 +240,7 @@ func (m *Manager) ensureGroupTrackerForApp(queuePath,
applicationID string, user
// something matched, get the tracker or create if it does not exist
if appGroup != common.Empty {
- groupTracker = m.GetGroupTracker(appGroup)
- if groupTracker == nil {
- log.Log(log.SchedUGM).Info("Group tracker doesn't
exists. Creating appGroup tracker",
- zap.String("queue path", queuePath),
- zap.String("group", appGroup))
- groupTracker = newGroupTracker(appGroup, m.events)
- m.Lock()
- m.groupTrackers[appGroup] = groupTracker
- m.Unlock()
- }
+ groupTracker = m.getGroupTracker(appGroup)
}
log.Log(log.SchedUGM).Info("Group tracker set for user application",
zap.String("group", appGroup),
@@ -638,6 +629,20 @@ func (m *Manager) getUserTracker(user string) *UserTracker
{
return userTracker
}
+// getGroupTracker returns the requested group tracker and creates one if it
does not exist.
+func (m *Manager) getGroupTracker(group string) *GroupTracker {
+ m.Lock()
+ defer m.Unlock()
+ if gt, ok := m.groupTrackers[group]; ok {
+ return gt
+ }
+ log.Log(log.SchedUGM).Info("Group tracker doesn't exists. Creating
group tracker.",
+ zap.String("group", group))
+ groupTracker := newGroupTracker(group, m.events)
+ m.groupTrackers[group] = groupTracker
+ return groupTracker
+}
+
func (m *Manager) getUserWildCardLimitsConfig(queuePath string) *LimitConfig {
if config, ok := m.userWildCardLimitsConfig[queuePath]; ok {
return config
diff --git a/pkg/scheduler/ugm/manager_test.go
b/pkg/scheduler/ugm/manager_test.go
index 0c3e06fc..077469d9 100644
--- a/pkg/scheduler/ugm/manager_test.go
+++ b/pkg/scheduler/ugm/manager_test.go
@@ -22,6 +22,7 @@ import (
"fmt"
"strconv"
"strings"
+ "sync"
"testing"
"gotest.tools/v3/assert"
@@ -2027,3 +2028,68 @@ func assertWildCardLimits(t *testing.T, limitsConfig
map[string]*LimitConfig, ex
}
assert.Equal(t, resources.Equals(expResource, configuredResource), true)
}
+
+func TestEnsureGroupTrackerForAppConcurrentRace(t *testing.T) {
+ groupName := "devs"
+ queuePath := queuePathLeaf
+ user := security.UserGroup{
+ User: "alice",
+ Groups: []string{groupName},
+ }
+ usage :=
resources.NewResourceFromMap(map[string]resources.Quantity{"vcores": 100})
+
+ // Run multiple rounds to consistently trigger the check-then-act race
window
+ for round := 1; round <= 10; round++ {
+ setupUGM()
+ manager := GetUserManager()
+
+ manager.Lock()
+ manager.configuredGroups = map[string][]string{
+ queuePath: {groupName},
+ }
+ manager.Unlock()
+
+ numApps := 50
+ var wg sync.WaitGroup
+ start := make(chan struct{})
+
+ for i := 0; i < numApps; i++ {
+ wg.Add(1)
+ go func(appIndex int) {
+ defer wg.Done()
+ <-start
+ appID := fmt.Sprintf("app-r%d-%d", round,
appIndex)
+ manager.IncreaseTrackedResource(queuePath,
appID, usage, user)
+ }(i)
+ }
+
+ close(start)
+ wg.Wait()
+
+ canonicalGroupTracker := manager.GetGroupTracker(groupName)
+ assert.Assert(t, canonicalGroupTracker != nil, "canonical group
tracker should exist in manager")
+
+ userTracker := manager.GetUserTracker("alice")
+ assert.Assert(t, userTracker != nil, "user tracker should
exist")
+
+ userTracker.RLock()
+ uniqueTrackers := make(map[*GroupTracker]int)
+ mismatchedApps := make([]string, 0)
+ for appID, gt := range userTracker.appGroupTrackers {
+ uniqueTrackers[gt]++
+ if gt != canonicalGroupTracker {
+ mismatchedApps = append(mismatchedApps, appID)
+ }
+ }
+ userTracker.RUnlock()
+
+ if len(uniqueTrackers) > 1 || len(mismatchedApps) > 0 {
+ t.Logf("=== CHAOS HARNESS TRIGGERED IN ROUND %d ===",
round)
+ t.Logf("Total applications launched: %d", numApps)
+ t.Logf("Distinct GroupTracker instances created in
memory: %d", len(uniqueTrackers))
+ t.Logf("Applications linked to non-canonical orphaned
GroupTrackers: %d", len(mismatchedApps))
+ t.Fatalf("CRITICAL BUG DETECTED: Multiple GroupTracker
instances created for group '%s'! Expected 1, found %d instances (Round %d)",
+ groupName, len(uniqueTrackers), round)
+ }
+ }
+}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]