From 5b43a527e7aa3e6cf48dc949ab321518b6e77bbe Mon Sep 17 00:00:00 2001 From: Ed Bartosh Date: Thu, 1 Oct 2026 17:05:48 +0300 Subject: [PATCH 1/2] topology-aware: pin containers to their DRA claims. Signed-off-by: Ed Bartosh --- cmd/plugins/topology-aware/policy/dra.go | 31 ++ cmd/plugins/topology-aware/policy/dra_test.go | 375 +++++++++++++++++- cmd/plugins/topology-aware/policy/pools.go | 68 +++- .../topology-aware/policy/resources.go | 41 ++ 4 files changed, 501 insertions(+), 14 deletions(-) diff --git a/cmd/plugins/topology-aware/policy/dra.go b/cmd/plugins/topology-aware/policy/dra.go index 1b9934f72..52b13821e 100644 --- a/cmd/plugins/topology-aware/policy/dra.go +++ b/cmd/plugins/topology-aware/policy/dra.go @@ -23,6 +23,8 @@ package topologyaware // memory. Like other grants, it is accounted, kept across a reconfiguration // and saved to the cache. The claimed CPUs leave the free supply, so the // policy does not hand them out again. +// A container holding claims runs on their CPUs, and also on its own CPUs +// if it has a CPU request. import ( "fmt" @@ -35,6 +37,7 @@ import ( "k8s.io/utils/ptr" specs "tags.cncf.io/container-device-interface/specs-go" + "github.com/containers/nri-plugins/pkg/resmgr/cache" system "github.com/containers/nri-plugins/pkg/sysfs" "github.com/containers/nri-plugins/pkg/utils/cpuset" idset "github.com/intel/goresctrl/pkg/utils" @@ -174,6 +177,9 @@ func (p *policy) ReleaseClaim(uid types.UID) error { return nil } + // The kubelet unprepares a claim only after all of its containers have + // stopped, and their grants went with them, so no grant holds the + // claim's CPUs by now. grant.Release() delete(p.allocations.claims, string(uid)) p.saveAllocations() @@ -182,6 +188,31 @@ func (p *policy) ReleaseClaim(uid types.UID) error { return nil } +// containerClaims returns the CPUs of the claims a container holds. It holds +// a claim if its environment has the variable the claim's edits set, naming +// the claim's CPUs. A variable naming a claim we do not hold, or other CPUs +// than the claim's, gets the container nothing. The value is compared as a +// cpuset, so it need not be in the form claimEdits wrote it in. +// +// The environment is what every runtime shows us today. It is set by the pod +// spec too, so a container knowing a claim's UID and CPUs can run on them +// without holding it. The CDI devices of a container name its claims and come +// from the runtime alone, but not all maintained containerd and CRI-O releases +// report them to NRI yet; once they do, they should replace the environment. +func (p *policy) containerClaims(c cache.Container) cpuset.CPUSet { + cpus := cpuset.New() + for uid, claim := range p.allocations.claims { + v, ok := c.GetEnv(draEnvPrefix + uid) + if !ok { + continue + } + if held, err := cpuset.Parse(v); err == nil && held.Equals(claim.ExclusiveCPUs()) { + cpus = cpus.Union(held) + } + } + return cpus +} + // draRequest is what one allocation result requests: CPUs of a NUMA node. type draRequest struct { node system.Node diff --git a/cmd/plugins/topology-aware/policy/dra_test.go b/cmd/plugins/topology-aware/policy/dra_test.go index 647610f33..4f7fd43d5 100644 --- a/cmd/plugins/topology-aware/policy/dra_test.go +++ b/cmd/plugins/topology-aware/policy/dra_test.go @@ -17,13 +17,17 @@ package topologyaware import ( "fmt" "reflect" + "slices" "strings" "sync" "testing" cfgapi "github.com/containers/nri-plugins/pkg/apis/config/v1alpha1/resmgr/policy/topologyaware" + "github.com/containers/nri-plugins/pkg/resmgr/cache" + libmem "github.com/containers/nri-plugins/pkg/resmgr/lib/memory" policyapi "github.com/containers/nri-plugins/pkg/resmgr/policy" system "github.com/containers/nri-plugins/pkg/sysfs" + "github.com/containers/nri-plugins/pkg/topology" "github.com/containers/nri-plugins/pkg/utils/cpuset" v1 "k8s.io/api/core/v1" resourceapi "k8s.io/api/resource/v1" @@ -264,16 +268,37 @@ func freeCPUs(p *policy) cpuset.CPUSet { return grantableCPUs(p.root.FreeSupply()) } -// pinnedContainer keeps the cpuset the policy pinned it to last. +// pinnedContainer keeps the cpuset and CPU shares the policy set last. type pinnedContainer struct { mockContainer - cpus string + env map[string]string + hints topology.Hints + strict bool + cpus string + shares int64 +} + +func (c *pinnedContainer) GetTopologyHints() topology.Hints { + return c.hints +} + +func (c *pinnedContainer) StrictTopologyHints() bool { + return c.strict } func (c *pinnedContainer) SetCpusetCpus(cpus string) { c.cpus = cpus } +func (c *pinnedContainer) SetCPUShares(shares int64) { + c.shares = shares +} + +func (c *pinnedContainer) GetEnv(key string) (string, bool) { + v, ok := c.env[key] + return v, ok +} + // allocateShared gives a container a shared grant of cpu, in pool if one is // given. func allocateShared(t *testing.T, p *policy, id, cpu, pool string) *pinnedContainer { @@ -597,6 +622,352 @@ func TestReleaseClaim(t *testing.T) { } } +// A container holding claims runs on their CPUs, and also on its own CPUs if +// it has a CPU request: exclusive, shared or reserved ones. It keeps its +// claimed CPUs as other claims come and go. A variable naming other CPUs or a +// claim we do not hold is ignored. +func TestPinClaimedCPUs(t *testing.T) { + sys := testServerSystem(t) + + // Every case gets a policy of its own holding these claims, a on node 0 + // and b on node 1. + setup := func(t *testing.T) (p *policy, a, b cpuset.CPUSet) { + cfg := draTestConfig("", draTestReserved(sys)) + cfg.PinCPU = true + p, _ = draTestPolicy(t, sys, cfg) + a = allocateClaim(t, p, "a", draTestResult(0, "2")) + b = allocateClaim(t, p, "b", draTestResult(1, "2")) + return p, a, b + } + + tcs := []struct { + name string + namespace string + qos v1.PodQOSClass + cpu string + env func(a, b cpuset.CPUSet) map[string]string + claims []string + exclusive int + portion int // shared or reserved milli-CPU the container gets + reserved bool // the container gets reserved CPUs + }{ + { + name: "besteffort", + qos: v1.PodQOSBestEffort, + env: func(a, _ cpuset.CPUSet) map[string]string { + return map[string]string{draEnvPrefix + "a": a.String()} + }, + claims: []string{"a"}, + }, + { + name: "shared", + qos: v1.PodQOSBurstable, + cpu: "100m", + env: func(a, _ cpuset.CPUSet) map[string]string { + return map[string]string{draEnvPrefix + "a": a.String()} + }, + claims: []string{"a"}, + portion: 100, + }, + { + name: "reserved", + namespace: "kube-system", + qos: v1.PodQOSBurstable, + cpu: "100m", + env: func(a, _ cpuset.CPUSet) map[string]string { + return map[string]string{draEnvPrefix + "a": a.String()} + }, + claims: []string{"a"}, + portion: 100, + reserved: true, + }, + { + // A reserved container with no CPU request runs on its claims + // only, as a besteffort one does. + name: "reserved, no request", + namespace: "kube-system", + qos: v1.PodQOSBestEffort, + env: func(a, _ cpuset.CPUSet) map[string]string { + return map[string]string{draEnvPrefix + "a": a.String()} + }, + claims: []string{"a"}, + reserved: true, + }, + { + // 1500m gives one exclusive CPU and a 500m shared portion. + name: "exclusive and shared", + qos: v1.PodQOSGuaranteed, + cpu: "1500m", + env: func(a, _ cpuset.CPUSet) map[string]string { + return map[string]string{draEnvPrefix + "a": a.String()} + }, + claims: []string{"a"}, + exclusive: 1, + portion: 500, + }, + { + name: "exclusive", + qos: v1.PodQOSGuaranteed, + cpu: "2", + env: func(a, _ cpuset.CPUSet) map[string]string { + return map[string]string{draEnvPrefix + "a": a.String()} + }, + claims: []string{"a"}, + exclusive: 2, + }, + { + name: "two claims", + qos: v1.PodQOSBurstable, + cpu: "100m", + env: func(a, b cpuset.CPUSet) map[string]string { + return map[string]string{draEnvPrefix + "a": a.String(), draEnvPrefix + "b": b.String()} + }, + claims: []string{"a", "b"}, + portion: 100, + }, + { + // The same CPUs in another form, listed one by one in reverse. + name: "same CPUs, other form", + qos: v1.PodQOSBurstable, + cpu: "100m", + env: func(a, _ cpuset.CPUSet) map[string]string { + cpus := []string{} + for _, cpu := range slices.Backward(a.List()) { + cpus = append(cpus, fmt.Sprint(cpu)) + } + return map[string]string{draEnvPrefix + "a": strings.Join(cpus, ",")} + }, + claims: []string{"a"}, + portion: 100, + }, + { + name: "other CPUs", + qos: v1.PodQOSBurstable, + cpu: "100m", + env: func(_, b cpuset.CPUSet) map[string]string { + return map[string]string{draEnvPrefix + "a": b.String()} + }, + }, + { + name: "unknown claim", + qos: v1.PodQOSBurstable, + cpu: "100m", + env: func(a, _ cpuset.CPUSet) map[string]string { + return map[string]string{draEnvPrefix + "x": a.String()} + }, + }, + } + + for _, tc := range tcs { + t.Run(tc.name, func(t *testing.T) { + p, a, b := setup(t) + claimed := cpuset.New() + for _, uid := range tc.claims { + claimed = claimed.Union(map[string]cpuset.CPUSet{"a": a, "b": b}[uid]) + } + + c := &pinnedContainer{ + mockContainer: mockContainer{ + name: tc.name, + namespace: tc.namespace, + returnValueForGetID: tc.name, + returnValueForQOSClass: tc.qos, + }, + env: tc.env(a, b), + } + if tc.cpu != "" { + c.returnValueForGetResourceRequirements = v1.ResourceRequirements{ + Requests: v1.ResourceList{v1.ResourceCPU: resource.MustParse(tc.cpu)}, + Limits: v1.ResourceList{v1.ResourceCPU: resource.MustParse(tc.cpu)}, + } + } + grant, err := p.allocatePool(c, "") + if err != nil { + t.Fatalf("failed to allocate: %v", err) + } + p.applyGrant(grant) + + exclusive := grant.ExclusiveCPUs() + if exclusive.Size() != tc.exclusive { + t.Fatalf("got exclusive CPUs %s, expected %d", exclusive, tc.exclusive) + } + pinned := cpuset.MustParse(c.cpus) + if claimed.IsEmpty() { + if pinned.IsEmpty() || !pinned.Intersection(a.Union(b)).IsEmpty() { + t.Errorf("pinned to %s, expected shared CPUs, none of %s", pinned, a.Union(b)) + } + return + } + own := exclusive + switch { + case tc.portion > 0 && tc.reserved: + own = grant.ReservedCPUs() + case tc.portion > 0: + own = own.Union(grant.SharedCPUs()) + } + if want := own.Union(claimed); !pinned.Equals(want) { + t.Errorf("pinned to %s, expected %s", pinned, want) + } + + // A container with normal CPUs goes to the pool of its claims, + // so its memory comes from there. A container with reserved CPUs + // stays in the root pool, where the reserved CPUs are. + pool := p.poolForCPUs(claimed) + if tc.reserved { + pool = p.root + } + if got := grant.GetCPUNode(); !got.IsSameNode(pool) { + t.Errorf("granted from %s, expected %s", got.Name(), pool.Name()) + } + mems := libmem.NewNodeMask(pool.GetMemset(memoryAll).Members()...) + if zone := grant.GetMemoryZone(); zone == 0 || zone.And(mems) != zone { + t.Errorf("got memory zone %s, expected one within %s", zone, mems) + } + if portion := grant.CPUPortion(); portion != tc.portion { + t.Errorf("granted %dm shared or reserved CPU, expected %dm", portion, tc.portion) + } + // As for mixed exclusive and shared CPUs, the CPU shares come from + // the shared or reserved portion if there is one, otherwise from + // all CPUs the container runs on. + milliCPU := tc.portion + if milliCPU == 0 { + milliCPU = 1000 * pinned.Size() + } + if want := cache.MilliCPUToShares(int64(milliCPU)); c.shares != int64(want) { + t.Errorf("got CPU shares %d, expected %d", c.shares, want) + } + + // While another claim is held, the container keeps its claimed + // CPUs and does not run on the other claim's CPUs. Once that + // claim is released, the container is back where it started. + // Node 0 has two isolated CPUs, so a claim for four takes some + // shared CPUs too. + cpus := allocateClaim(t, p, "c", draTestResult(0, "4")) + if got := cpuset.MustParse(c.cpus); !claimed.IsSubsetOf(got) || !got.Intersection(cpus).IsEmpty() { + t.Errorf("pinned to %s while claim c holds %s, expected all of %s and none of %s", + got, cpus, claimed, cpus) + } + if err := p.ReleaseClaim("c"); err != nil { + t.Fatalf("failed to release claim c: %v", err) + } + if got := cpuset.MustParse(c.cpus); !got.Equals(pinned) { + t.Errorf("pinned to %s after claim c came and went, expected %s", got, pinned) + } + }) + } +} + +// A besteffort container holding a claim does not run on shared CPUs, so it +// holds back no shared capacity of its pool: the rest of its node can still +// be claimed. +func TestClaimingContainerSparesSharedCPUs(t *testing.T) { + sys := testServerSystem(t) + cfg := draTestConfig("", draTestReserved(sys)) + cfg.PinCPU = true + p, _ := draTestPolicy(t, sys, cfg) + + a := allocateClaim(t, p, "a", draTestResult(0, "2")) + c := &pinnedContainer{ + mockContainer: mockContainer{ + name: "claimer", + returnValueForGetID: "claimer", + returnValueForQOSClass: v1.PodQOSBestEffort, + }, + env: map[string]string{draEnvPrefix + "a": a.String()}, + } + grant, err := p.allocatePool(c, "") + if err != nil { + t.Fatalf("failed to allocate: %v", err) + } + p.applyGrant(grant) + + left := freeCPUs(p).Intersection(sys.Node(0).CPUSet()).Size() + allocateClaim(t, p, "b", draTestResult(0, fmt.Sprint(left))) +} + +// The claimed CPUs of a container must follow its strict topology hints. The +// hint names CPUs of node 0. The claim is on node 0 too, so memory passes the +// hint in both cases. Only the claimed CPUs decide the result. +func TestClaimStrictTopologyHint(t *testing.T) { + sys := testServerSystem(t) + + for _, tc := range []struct { + name string + aligned bool + }{ + {name: "aligned claim", aligned: true}, + {name: "misaligned claim"}, + } { + t.Run(tc.name, func(t *testing.T) { + cfg := draTestConfig("", draTestReserved(sys)) + cfg.PinCPU = true + p, _ := draTestPolicy(t, sys, cfg) + + a := allocateClaim(t, p, "a", draTestResult(0, "2")) + hinted := sys.Node(0).CPUSet() + if !tc.aligned { + hinted = hinted.Difference(a) + } + c := &pinnedContainer{ + mockContainer: mockContainer{ + name: "claimer", + returnValueForGetID: "claimer", + returnValueForQOSClass: v1.PodQOSBestEffort, + }, + env: map[string]string{draEnvPrefix + "a": a.String()}, + hints: topology.Hints{"test": {CPUs: hinted.String()}}, + strict: true, + } + + _, err := p.allocatePool(c, "") + switch { + case tc.aligned && err != nil: + t.Errorf("failed to allocate a container whose claim follows its hint: %v", err) + case !tc.aligned && (err == nil || !strings.Contains(err.Error(), "claimed CPUs")): + t.Errorf("got error %v, expected claimed CPUs %s to fail the strict hint", err, a) + } + }) + } +} + +// A container holding a claim stays on it across a reconfiguration. +func TestReconfigureKeepsClaimer(t *testing.T) { + sys := testServerSystem(t) + cfg := draTestConfig("", draTestReserved(sys)) + cfg.PinCPU = true + p, _ := draTestPolicy(t, sys, cfg) + + a := allocateClaim(t, p, "a", draTestResult(0, "2")) + c := &pinnedContainer{ + mockContainer: mockContainer{ + name: "claimer", + returnValueForGetID: "claimer", + returnValueForQOSClass: v1.PodQOSBestEffort, + }, + env: map[string]string{draEnvPrefix + "a": a.String()}, + } + grant, err := p.allocatePool(c, "") + if err != nil { + t.Fatalf("failed to allocate: %v", err) + } + p.applyGrant(grant) + if pinned := cpuset.MustParse(c.cpus); !pinned.Equals(a) { + t.Fatalf("pinned to %s, expected the claim's %s", pinned, a) + } + + cfg = draTestConfig("", draTestReserved(sys)) + cfg.PinCPU = true + cfg.ColocatePods = true + if err := p.Reconfigure(cfg); err != nil { + t.Fatalf("failed to reconfigure policy: %v", err) + } + grant, _ = p.allocations.getGrant("claimer") + if pinned := cpuset.MustParse(c.cpus); !pinned.Equals(a) || !grant.ClaimedCPUs().Equals(a) { + t.Errorf("pinned to %s, holding %s after reconfiguring, expected the claim's %s", + pinned, grant.ClaimedCPUs(), a) + } +} + // TestClaimAndContainerWithOneID verifies that a claim and a container // with the same ID both survive a reconfiguration and a round trip through the cache. func TestClaimAndContainerWithOneID(t *testing.T) { diff --git a/cmd/plugins/topology-aware/policy/pools.go b/cmd/plugins/topology-aware/policy/pools.go index 2eb1a7c4c..2aa9fd994 100644 --- a/cmd/plugins/topology-aware/policy/pools.go +++ b/cmd/plugins/topology-aware/policy/pools.go @@ -429,7 +429,19 @@ func (p *policy) allocatePool(container cache.Container, poolHint string) (Grant request.SetCPUType(cpuNormal) } - if request.CPUType() == cpuReserved || request.CPUType() == cpuPreserve { + if claimed := request.ClaimedCPUs(); !claimed.IsEmpty() && request.CPUType() == cpuNormal { + // A container holding DRA claims goes to the pool of their CPUs, + // so that its memory and its own CPUs come from there. A container + // with reserved CPUs stays in the root pool below: the claim's pool + // usually has no reserved CPUs, and the container would get normal + // CPUs instead. + pool = p.poolForCPUs(claimed) + o, err := p.getMemOffer(pool, request) + if err != nil { + return nil, policyError("failed to get offer for request %s: %v", request, err) + } + offer = o + } else if request.CPUType() == cpuReserved || request.CPUType() == cpuPreserve { pool = p.root o, err := p.getMemOffer(pool, request) if err != nil { @@ -513,8 +525,9 @@ func (p *policy) allocatePool(container cache.Container, poolHint string) (Grant // setPreferredCpusetCpus pins container's CPUs according to what has been // allocated for it, taking into account if the container should run -// with hyperthreads hidden. -func (p *policy) setPreferredCpusetCpus(container cache.Container, allocated cpuset.CPUSet, info string) { +// with hyperthreads hidden. The CPUs of its DRA claims are added as they +// are, hyperthreads included: they are what the claims asked for. +func (p *policy) setPreferredCpusetCpus(container cache.Container, allocated, claimed cpuset.CPUSet, info string) { allow := allocated hidingInfo := "" pod, ok := container.GetPod() @@ -527,7 +540,7 @@ func (p *policy) setPreferredCpusetCpus(container cache.Container, allocated cpu } } log.Infof("%s%s", info, hidingInfo) - container.SetCpusetCpus(allow.String()) + container.SetCpusetCpus(allow.Union(claimed).String()) } // Apply the result of allocation to the requesting container. @@ -573,12 +586,26 @@ func (p *policy) applyGrant(grant Grant) { mems = grant.GetMemoryZone() } + claimed := grant.ClaimedCPUs() + if opt.PinCPU { if cpuType == cpuPreserve { log.Infof(" => preserving %s cpuset %s", container.PrettyName(), container.GetCpusetCpus()) + } else if !claimed.IsEmpty() { + // The claimed CPUs are added to the container's own CPUs. A + // container with no CPU request of its own runs on its claims only. + own, label := exclusive, "claimed" + if cpuPortion > 0 { + own, label = cpus, kind+"+claimed" + } else if !exclusive.IsEmpty() { + label = "exclusive+claimed" + } + p.setPreferredCpusetCpus(container, own, claimed, + fmt.Sprintf(" => pinning %s to (%s) cpuset %s", + container.PrettyName(), label, own.Union(claimed))) } else { if cpus.Size() > 0 { - p.setPreferredCpusetCpus(container, cpus, + p.setPreferredCpusetCpus(container, cpus, cpuset.New(), fmt.Sprintf(" => pinning %s to (%s) cpuset %s", container.PrettyName(), kind, cpus)) } else { @@ -609,7 +636,7 @@ func (p *policy) applyGrant(grant Grant) { // as long as that allocation is genuinely system-wide exclusive. milliCPU := cpuPortion if milliCPU == 0 { - milliCPU = 1000 * grant.ExclusiveCPUs().Size() + milliCPU = 1000 * (grant.ExclusiveCPUs().Size() + claimed.Size()) } container.SetCPUShares(int64(cache.MilliCPUToShares(int64(milliCPU)))) @@ -698,17 +725,27 @@ func (p *policy) updateSharedAllocations(grant *Grant) { continue } + if other.SharedPortion() == 0 && !other.ClaimedCPUs().IsEmpty() { + log.Infof(" => %s not affected (only claimed CPUs)...", other) + continue + } + if opt.PinCPU { shared := other.GetCPUNode().FreeSupply().SharableCPUs() exclusive := other.ExclusiveCPUs() + claimed := other.ClaimedCPUs() + withClaimed := "" + if !claimed.IsEmpty() { + withClaimed = ", claimed CPUs " + claimed.String() + } if exclusive.IsEmpty() { - p.setPreferredCpusetCpus(other.GetContainer(), shared, - fmt.Sprintf(" => updating %s with shared CPUs of %s: %s...", - other, other.GetCPUNode().Name(), shared.String())) + p.setPreferredCpusetCpus(other.GetContainer(), shared, claimed, + fmt.Sprintf(" => updating %s with shared CPUs of %s: %s%s...", + other, other.GetCPUNode().Name(), shared.String(), withClaimed)) } else { - p.setPreferredCpusetCpus(other.GetContainer(), exclusive.Union(shared), - fmt.Sprintf(" => updating %s with exclusive+shared CPUs of %s: %s+%s...", - other, other.GetCPUNode().Name(), exclusive.String(), shared.String())) + p.setPreferredCpusetCpus(other.GetContainer(), exclusive.Union(shared), claimed, + fmt.Sprintf(" => updating %s with exclusive+shared CPUs of %s: %s+%s%s...", + other, other.GetCPUNode().Name(), exclusive.String(), shared.String(), withClaimed)) } } } @@ -786,7 +823,14 @@ func (p *policy) hasZeroCpuReqContainer(pool Node) bool { return false } + // A container holding DRA claims and no CPU request runs on its claims + // only, so it holds back no shared CPU. + if !g.ClaimedCPUs().IsEmpty() { + return false + } + ctr := g.GetContainer() + switch ctr.GetQOSClass() { case corev1.PodQOSBestEffort: found = true diff --git a/cmd/plugins/topology-aware/policy/resources.go b/cmd/plugins/topology-aware/policy/resources.go index b400efe2d..5d9d98722 100644 --- a/cmd/plugins/topology-aware/policy/resources.go +++ b/cmd/plugins/topology-aware/policy/resources.go @@ -103,6 +103,8 @@ type Supply interface { type Request interface { // GetContainer returns the container requesting CPU capacity. GetContainer() cache.Container + // ClaimedCPUs returns the CPUs of the DRA claims the container holds. + ClaimedCPUs() cpuset.CPUSet // String returns a printable representation of this request. String() string // CPUType returns the type of requested CPU. @@ -178,6 +180,10 @@ type Grant interface { SetMemoryZone(libmem.NodeMask) // SetMemorySize sets the amount of memory to allocate. SetMemorySize(int64) + // ClaimedCPUs returns the CPUs of the DRA claims the container holds. + ClaimedCPUs() cpuset.CPUSet + // SetClaimedCPUs sets the CPUs of the DRA claims the container holds. + SetClaimedCPUs(cpuset.CPUSet) // IrqAffinity returns the IRQ affinity for this grant. IrqAffinity() *IrqAffinity @@ -255,6 +261,7 @@ type request struct { memType memoryType // requested types of memory pickByHints bool // preference to pick resources by hints irqs *IrqAffinity // IRQ affinity for this request + claimed cpuset.CPUSet // CPUs of the DRA claims the container holds // coldStart tells the timeout (in milliseconds) how long to wait until // a DRAM memory controller should be added to a container asking for a @@ -280,6 +287,7 @@ type grant struct { memZone libmem.NodeMask // allocated memory zone cpuClass string // CPU class to apply to exclusive CPUs irqs *IrqAffinity // IRQ affinity for this request + claimed cpuset.CPUSet // CPUs of the DRA claims the container holds } var _ Grant = &grant{} @@ -455,6 +463,7 @@ func (cs *supply) AllocateCPU(r Request) (Grant, error) { } grant := newGrant(cs.node, cr.GetContainer(), cpuType, cpuClass, exclusive, 0, 0, irqs, 0) + grant.SetClaimedCPUs(cr.claimed) grant.AccountAllocateCPU() // allocate the shared fraction of CPUs @@ -821,6 +830,14 @@ func (p *policy) newRequest(container cache.Container, types libmem.TypeMask) (R pod, _ := container.GetPod() full, fraction, cpuLimit, isolate, cpuType, prio := cpuAllocationPreferences(pod, container) req, lim, mtype := memoryAllocationPreference(pod, container) + + // The DRA claims of a container are fixed for its lifetime, so they are + // looked up once, here. + claimed := cpuset.New() + if cpuType != cpuPreserve { + claimed = p.containerClaims(container) + } + coldStart := time.Duration(0) cpuClass, isCtrScoped, err := p.resolveCpuClass(container) @@ -884,6 +901,7 @@ func (p *policy) newRequest(container cache.Container, types libmem.TypeMask) (R prio: prio, pickByHints: pickByHintsPreference(pod, container), irqs: irqs, + claimed: claimed, }, nil } @@ -892,6 +910,11 @@ func (cr *request) GetContainer() cache.Container { return cr.container } +// ClaimedCPUs returns the CPUs of the DRA claims the container holds. +func (cr *request) ClaimedCPUs() cpuset.CPUSet { + return cr.claimed +} + // String returns aprintable representation of the CPU request. func (cr *request) String() string { mem := fmt.Sprintf("", @@ -1023,6 +1046,13 @@ func (cr *request) verifyStrictTopologyHints(g Grant) error { } } + if claimed := g.ClaimedCPUs(); claimed.Size() > 0 { + if cpus := hint.MisalignedCPUSet(claimed); cpus.Size() > 0 { + return policyError("claimed CPUs %q fail strict hint %v", + cpus.String(), h) + } + } + if mems := hint.MisalignedMems(g.GetMemoryZone()); mems.Size() > 0 { return policyError("granted memory zones %s fail strict hint %v", mems.String(), h) @@ -1407,6 +1437,7 @@ func (cg *grant) Clone() Grant { memZone: cg.GetMemoryZone(), memSize: cg.GetMemorySize(), coldStart: cg.ColdStart(), + claimed: cg.ClaimedCPUs(), } } @@ -1425,6 +1456,16 @@ func (cg *grant) GetContainer() cache.Container { return cg.container } +// ClaimedCPUs returns the CPUs of the DRA claims the container holds. +func (cg *grant) ClaimedCPUs() cpuset.CPUSet { + return cg.claimed +} + +// SetClaimedCPUs sets the CPUs of the DRA claims the container holds. +func (cg *grant) SetClaimedCPUs(cpus cpuset.CPUSet) { + cg.claimed = cpus +} + // GetNode returns the Node this grant gets its CPU allocation from. func (cg *grant) GetCPUNode() Node { return cg.node From fb05d6ad8cf185cb59d30ea1e15a4631080c119c Mon Sep 17 00:00:00 2001 From: Ed Bartosh Date: Thu, 1 Oct 2026 17:39:25 +0300 Subject: [PATCH 2/2] e2e: test pinning to topology-aware DRA claims. Signed-off-by: Ed Bartosh --- .../n4c16/test42-dra-devices/code.var.sh | 42 ++++++++++++++++--- .../n4c16/test42-dra-devices/dra-pod.yaml.in | 8 +++- 2 files changed, 42 insertions(+), 8 deletions(-) diff --git a/test/e2e/policies.test-suite/topology-aware/n4c16/test42-dra-devices/code.var.sh b/test/e2e/policies.test-suite/topology-aware/n4c16/test42-dra-devices/code.var.sh index 26d93dee6..4a0175219 100644 --- a/test/e2e/policies.test-suite/topology-aware/n4c16/test42-dra-devices/code.var.sh +++ b/test/e2e/policies.test-suite/topology-aware/n4c16/test42-dra-devices/code.var.sh @@ -1,5 +1,6 @@ # This test verifies that the topology-aware policy publishes one DRA -# device per NUMA node and allocates CPUs to claims. +# device per NUMA node, allocates CPUs to claims and pins containers to +# the CPUs of their claim. # cleanup() { @@ -124,7 +125,9 @@ wait="" create deviceclass # CPU 15. claim pod0 3 3 cpus0=$(claimed-cpus pod0) -verify "cpuset('$cpus0') == cpuset('$(node-cpus 3)') - cpuset('15')" +verify "cpuset('$cpus0') == cpuset('$(node-cpus 3)') - cpuset('15')" \ + "cpus['pod0c0'] == cpuset('$cpus0')" \ + "node_ids(mems['pod0c0']) == {3}" # Two CPUs are the two threads of one core, and a second claim on the same # node gets the other core. @@ -132,10 +135,14 @@ claim pod1 1 2 cpus1=$(claimed-cpus pod1) siblings=$(vm-command-q "cat /sys/devices/system/cpu/cpu${cpus1%%[-,]*}/topology/thread_siblings_list") verify "cpuset('$cpus1') == cpuset('$siblings')" \ - "cpuset('$cpus1') <= cpuset('$(node-cpus 1)')" + "cpuset('$cpus1') <= cpuset('$(node-cpus 1)')" \ + "cpus['pod1c0'] == cpuset('$cpus1')" \ + "node_ids(mems['pod1c0']) == {1}" claim pod2 1 2 cpus2=$(claimed-cpus pod2) -verify "cpuset('$cpus1') | cpuset('$cpus2') == cpuset('$(node-cpus 1)')" +verify "cpuset('$cpus1') | cpuset('$cpus2') == cpuset('$(node-cpus 1)')" \ + "cpus['pod2c0'] == cpuset('$cpus2')" \ + "node_ids(mems['pod2c0']) == {1}" # No CPUs are left on node 1, so the scheduler must not place a third claim # there, until deleting pod1 releases its claim, whose CPUs the third gets. @@ -147,7 +154,9 @@ vm-command "kubectl delete pod pod1 --now --wait" vm-command "kubectl wait --for=condition=Ready pod/pod3 --timeout=120s" || error "pod3 did not start after pod1 released its claim" cpus3=$(claimed-cpus pod3) -verify "cpuset('$cpus3') == cpuset('$cpus1')" +verify "cpuset('$cpus3') == cpuset('$cpus1')" \ + "cpus['pod3c0'] == cpuset('$cpus3')" \ + "node_ids(mems['pod3c0']) == {1}" # A shared container does not run on claimed CPUs. It moves off the CPUs a # claim takes of its node, and back once the claim is released. @@ -159,6 +168,8 @@ shared=$(pyexec 'print(sorted(cpus["pod4c0"]))') claim pod5 "$node" 2 cpus5=$(claimed-cpus pod5) verify "cpuset('$cpus5') <= set($shared)" \ + "cpus['pod5c0'] == cpuset('$cpus5')" \ + "node_ids(mems['pod5c0']) == {$node}" \ "cpus['pod4c0'] == set($shared) - cpuset('$cpus5')" vm-command "kubectl delete pod pod5 --now --wait" # Deleting the pod does not wait for the kubelet to unprepare its claim. @@ -166,8 +177,27 @@ retry-until --timeout 30 --message "pod4 to get back the CPUs of pod5's claim" \ '[ "$(report allowed >/dev/null; pyexec "print(cpus[\"pod4c0\"] == set($shared))")" == True ]' verify "cpus['pod4c0'] == set($shared)" +# Two containers of one pod share its claim. Both run on the claim's CPUs. +CONTCOUNT=2 claim pod6 0 2 +cpus6=$(claimed-cpus pod6) +verify "cpus['pod6c0'] == cpuset('$cpus6')" \ + "cpus['pod6c1'] == cpuset('$cpus6')" \ + "node_ids(mems['pod6c0']) == {0}" \ + "node_ids(mems['pod6c1']) == {0}" \ + "cpus['pod4c0'].isdisjoint(cpuset('$cpus6'))" + +# A container with a CPU request of its own also gets shared CPUs. Node 2 has +# no other claims, so pod7 runs on all its CPUs: the 2 claimed ones and the 2 +# shared ones. +OWNCPU=100m claim pod7 2 2 +cpus7=$(claimed-cpus pod7) +verify "cpuset('$cpus7') <= cpuset('$(node-cpus 2)')" \ + "cpus['pod7c0'] == cpuset('$(node-cpus 2)')" \ + "node_ids(mems['pod7c0']) == {2}" \ + "cpus['pod4c0'].isdisjoint(cpuset('$cpus7'))" + vm-command "kubectl get deviceclasses,resourceslices,resourceclaims -o yaml" || : cleanup -echo "OK: the topology-aware policy published one DRA device per NUMA node, and gave each claim CPUs of its own" +echo "OK: the topology-aware policy published one DRA device per NUMA node, gave each claim CPUs of its own, and pinned containers to their claim's CPUs and to the memory of its node" diff --git a/test/e2e/policies.test-suite/topology-aware/n4c16/test42-dra-devices/dra-pod.yaml.in b/test/e2e/policies.test-suite/topology-aware/n4c16/test42-dra-devices/dra-pod.yaml.in index 971c19491..f736d47a2 100644 --- a/test/e2e/policies.test-suite/topology-aware/n4c16/test42-dra-devices/dra-pod.yaml.in +++ b/test/e2e/policies.test-suite/topology-aware/n4c16/test42-dra-devices/dra-pod.yaml.in @@ -7,14 +7,18 @@ spec: - name: cpus resourceClaimTemplateName: ${NAME} containers: - - name: ${NAME}c0 + $(for contnum in $(seq 1 ${CONTCOUNT}); do echo " + - name: ${NAME}c$(( contnum - 1 )) image: quay.io/prometheus/busybox imagePullPolicy: IfNotPresent command: - sh - -c - - echo ${NAME}c0 \$(sleep inf) + - echo ${NAME}c$(( contnum - 1 )) \$(sleep inf) resources: claims: - name: cpus + $( [ -n "${OWNCPU}" ] && echo "requests: + cpu: ${OWNCPU}" ) + "; done ) terminationGracePeriodSeconds: 1