Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 7 additions & 0 deletions internal/common/resource/resource.go
Original file line number Diff line number Diff line change
Expand Up @@ -264,6 +264,13 @@ func TotalPodResourceRequest(podSpec *v1.PodSpec) ComputeResources {
totalResources.Max(containerResource)
}
}

// Pod-level resources (KEP-2837): the effective request is
// max(sum of container requests, pod-level request) per resource. Inert when
// podSpec.Resources is unset.
if podSpec.Resources != nil {
totalResources.Max(FromResourceList(podSpec.Resources.Requests))
}
return totalResources
}

Expand Down
47 changes: 47 additions & 0 deletions internal/common/resource/resource_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@ import (
"github.com/stretchr/testify/require"
v1 "k8s.io/api/core/v1"
"k8s.io/apimachinery/pkg/api/resource"
"k8s.io/utils/ptr"
)

func TestComputeResources_String(t *testing.T) {
Expand Down Expand Up @@ -148,6 +149,52 @@ func TestTotalResourceRequest_ShouldCombineMaxInitContainerResourcesWithSummedCo
assert.Equal(t, result, FromResourceList(expectedResult))
}

// Pod-level resources (KEP-2837): TotalPodResourceRequest uses the effective
// request max(sum(containers), pod-level).
func TestTotalResourceRequest_PodLevelResources(t *testing.T) {
Comment thread
washcycle marked this conversation as resolved.
tests := map[string]struct {
containerResources []*v1.ResourceList
podLevelResources *v1.ResourceList
expected v1.ResourceList
}{
"pod-level only (empty containers) is accounted at the pod-level value": {
containerResources: []*v1.ResourceList{},
podLevelResources: ptr.To(makeContainerResource(4, 16)),
expected: makeContainerResource(4, 16),
},
// cpu: max(2, 4) = 4 ; memory: max(4, 2) = 4
"effective request is max of container-sum and pod-level": {
containerResources: []*v1.ResourceList{ptr.To(makeContainerResource(2, 4))},
podLevelResources: ptr.To(makeContainerResource(4, 2)),
expected: makeContainerResource(4, 4),
},
"container-sum wins when it exceeds pod-level": {
containerResources: []*v1.ResourceList{ptr.To(makeContainerResource(8, 8))},
podLevelResources: ptr.To(makeContainerResource(4, 4)),
expected: makeContainerResource(8, 8),
},
"nil pod-level leaves upstream behaviour unchanged": {
containerResources: []*v1.ResourceList{ptr.To(makeContainerResource(2, 4))},
podLevelResources: nil,
expected: makeContainerResource(2, 4),
},
}
for name, tc := range tests {
t.Run(name, func(t *testing.T) {
pod := makePodWithResource(tc.containerResources, []*v1.ResourceList{})
if tc.podLevelResources != nil {
pod.Spec.Resources = &v1.ResourceRequirements{
Requests: *tc.podLevelResources,
Limits: *tc.podLevelResources,
}
}

result := TotalPodResourceRequest(&pod.Spec)
assert.Equal(t, FromResourceList(tc.expected), result)
})
}
}

func TestTotalResourceRequest_NativeSidecarsShouldBeSummed(t *testing.T) {
mainResource := makeContainerResource(2, 1)
sidecarResource := makeContainerResource(1, 1)
Expand Down
19 changes: 11 additions & 8 deletions internal/lookoutingester/instructions/instructions.go
Original file line number Diff line number Diff line change
Expand Up @@ -562,18 +562,21 @@ func getJobResources(job *api.Job) jobResources {

podSpec := job.GetMainPodSpec()

for _, container := range podSpec.Containers {
resources.Cpu += getResource(container, v1.ResourceCPU, true)
resources.Memory += getResource(container, v1.ResourceMemory, false)
resources.EphemeralStorage += getResource(container, v1.ResourceEphemeralStorage, false)
resources.Gpu += getResource(container, "nvidia.com/gpu", false)
}
// Use the canonical effective-request computation so the reported footprint
// matches what the scheduler bills: sum of main containers + native sidecars,
// max over classic init containers, and max with the pod-level block (KEP-2837).
// This also fixes the prior undercount that summed only main containers.
requests := api.SchedulingResourceRequirementsFromPodSpec(podSpec).Requests
resources.Cpu = getResourceFromList(requests, v1.ResourceCPU, true)
resources.Memory = getResourceFromList(requests, v1.ResourceMemory, false)
resources.EphemeralStorage = getResourceFromList(requests, v1.ResourceEphemeralStorage, false)
resources.Gpu = getResourceFromList(requests, "nvidia.com/gpu", false)

return resources
}

func getResource(container v1.Container, resourceName v1.ResourceName, useMillis bool) int64 {
resource, ok := container.Resources.Requests[resourceName]
func getResourceFromList(rl v1.ResourceList, resourceName v1.ResourceName, useMillis bool) int64 {
resource, ok := rl[resourceName]
if !ok {
return 0
}
Expand Down
6 changes: 6 additions & 0 deletions internal/scheduler/api.go
Original file line number Diff line number Diff line change
Expand Up @@ -238,6 +238,12 @@ func (srv *ExecutorApi) dropDisallowedResources(pod *v1.PodSpec) {
}
srv.dropDisallowedResourcesFromContainers(pod.InitContainers)
srv.dropDisallowedResourcesFromContainers(pod.Containers)
// Pod-level resources (KEP-2837) must be filtered by the same allow-list, else
// a disallowed resource could bypass it via the pod-level block.
if pod.Resources != nil {
removeDisallowedKeys(pod.Resources.Limits, srv.allowedResources)
removeDisallowedKeys(pod.Resources.Requests, srv.allowedResources)
}
}
Comment thread
greptile-apps[bot] marked this conversation as resolved.

func (srv *ExecutorApi) dropDisallowedResourcesFromContainers(containers []v1.Container) {
Expand Down
5 changes: 5 additions & 0 deletions internal/server/configuration/types.go
Original file line number Diff line number Diff line change
Expand Up @@ -99,6 +99,11 @@ type SubmissionConfig struct {
AddGangIdLabel bool
// Controls whether custom service names are allowed
AllowCustomServiceNames bool
// When true, honour Kubernetes pod-level resources (KEP-2837, podSpec.Resources):
// a container may omit its own resources if the pod-level block is set, and
// accounting uses max(sum(container requests), podLevel.Requests). Default false
// preserves container-only behaviour.
PodLevelResources bool
}

// TODO: we can probably just typedef this to map[string]string
Expand Down
9 changes: 9 additions & 0 deletions internal/server/submit/conversion/post_process.go
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,7 @@ var (
addGangIdLabel,
}
podLevelProcessors = []podProcessor{
dropPodLevelResourcesIfDisabled,
defaultActiveDeadlineSeconds,
defaultPriorityClass,
defaultResource,
Expand Down Expand Up @@ -110,6 +111,14 @@ func defaultPriorityClass(spec *v1.PodSpec, config configuration.SubmissionConfi
}
}

// Clears the pod-level resources block (KEP-2837) unless the feature is enabled,
// so all downstream accounting ignores it when the feature is off.
func dropPodLevelResourcesIfDisabled(spec *v1.PodSpec, config configuration.SubmissionConfig) {
if !config.PodLevelResources {
spec.Resources = nil
}
}

// Adds resources defined in config.DefaultJobLimits to all containers in the podspec if that container is missing
// requests/limits for that particular resource. This can be used to e.g. ensure that all jobs define at least some
// ephemeral storage.
Expand Down
33 changes: 33 additions & 0 deletions internal/server/submit/conversion/post_process_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -756,3 +756,36 @@ func submitMsgFromAnnotations(annotations map[string]string) *armadaevents.Submi
},
}
}

func TestDropPodLevelResourcesIfDisabled(t *testing.T) {
podLevel := &v1.ResourceRequirements{
Requests: v1.ResourceList{"cpu": resource.MustParse("2")},
Limits: v1.ResourceList{"cpu": resource.MustParse("2")},
}

tests := map[string]struct {
initialResources *v1.ResourceRequirements
enabled bool
expectedResources *v1.ResourceRequirements
}{
"disabled clears the pod-level block": {
initialResources: podLevel,
enabled: false,
},
"enabled preserves the pod-level block": {
initialResources: podLevel,
enabled: true,
expectedResources: podLevel,
},
"unset pod-level block is left unset": {
enabled: true,
},
}
for name, tc := range tests {
t.Run(name, func(t *testing.T) {
spec := &v1.PodSpec{Resources: tc.initialResources.DeepCopy()}
dropPodLevelResourcesIfDisabled(spec, configuration.SubmissionConfig{PodLevelResources: tc.enabled})
assert.Equal(t, tc.expectedResources, spec.Resources)
})
}
}
70 changes: 70 additions & 0 deletions internal/server/submit/validation/submit_request.go
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@ import (
"strings"

"github.com/pkg/errors"
v1 "k8s.io/api/core/v1"
"k8s.io/component-helpers/scheduling/corev1/nodeaffinity"

"github.com/armadaproject/armada/internal/common/constants"
Expand Down Expand Up @@ -252,8 +253,21 @@ func validateResources(j *api.JobSubmitRequestItem, config configuration.Submiss
if maxOversubscriptionByResource == nil {
maxOversubscriptionByResource = map[string]float64{}
}
// Pod-level resources (KEP-2837): when enabled and a pod-level block is set, a
// container may omit its own resources. Containers that do set resources are
// still validated below, as is the pod-level block itself.
podLevelResourcesEnabled := config.PodLevelResources && spec.Resources != nil
if podLevelResourcesEnabled {
if err := validatePodLevelResources(spec, maxOversubscriptionByResource, config); err != nil {
return err
}
}
for _, container := range armadaslices.Concatenate(spec.Containers, spec.InitContainers) {
if len(container.Resources.Requests) == 0 && len(container.Resources.Limits) == 0 {
if podLevelResourcesEnabled {
// Budget is supplied at the pod level; nothing to validate for this container.
continue
}
return fmt.Errorf("container %v has no resources specified", container.Name)
}

Expand Down Expand Up @@ -308,6 +322,62 @@ func validateResources(j *api.JobSubmitRequestItem, config configuration.Submiss
return nil
}

// validatePodLevelResources validates a pod-level resources block (KEP-2837),
// mirroring the per-container checks: requests and limits must be non-negative,
// cover the same resource set, satisfy limit >= request within the
// max-oversubscription ratio, and meet MinJobResources.
func validatePodLevelResources(
spec *v1.PodSpec,
maxOversubscriptionByResource map[string]float64,
config configuration.SubmissionConfig,
) error {
resources := spec.Resources
if len(resources.Requests) == 0 && len(resources.Limits) == 0 {
return fmt.Errorf("pod-level resources block is empty")
}
if len(resources.Requests) != len(resources.Limits) {
return fmt.Errorf("pod-level resources define different resources for requests and limits")
}
for resourceName, request := range resources.Requests {
if request.Sign() < 0 {
return fmt.Errorf("pod-level resources define negative request (%s) for resource %s", request.String(), resourceName)
}
}
for resourceName, limit := range resources.Limits {
if limit.Sign() < 0 {
return fmt.Errorf("pod-level resources define negative limit (%s) for resource %s", limit.String(), resourceName)
}
}
for resourceName, request := range resources.Requests {
limit, ok := resources.Limits[resourceName]
if !ok {
return fmt.Errorf("pod-level resources define %s for requests but not limits", resourceName)
}
if limit.MilliValue() < request.MilliValue() {
return fmt.Errorf("pod-level resources define %s with limits smaller than requests", resourceName)
}
maxOversubscription, ok := maxOversubscriptionByResource[resourceName.String()]
if !ok {
maxOversubscription = 1.0
}
if float64(limit.MilliValue()) > maxOversubscription*float64(request.MilliValue()) {
return fmt.Errorf("pod-level resources define %s with limits greater than %.2f*requests", resourceName, maxOversubscription)
}
}
// MinJobResources is checked against the effective request
// (max of the pod-level request and the summed container requests), since that
// is what the scheduler reserves; checking the raw pod-level value alone would
// wrongly reject a pod whose container total already meets the minimum.
effective := api.SchedulingResourceRequirementsFromPodSpec(spec).Requests
for rc, serverRsc := range config.MinJobResources {
eff := effective[rc]
if eff.Value() < serverRsc.Value() {
return fmt.Errorf("effective %s requests (%s) below server minimum (%s)", rc, &eff, &serverRsc)
}
}
return nil
}

// jobAdapter turns JobSubmitRequestItem into a MinimalJob
// This is needed for gang information to be extracted
type jobAdapter struct {
Expand Down
Loading
Loading