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]

Reply via email to