Skip to content

Commit 16ac563

Browse files
committed
feat(nvca): support per-function BYOO collector resource override
Adds a byooResources field to the nvcf-workload-config WorkloadConfig so a Helm function can size its BYOO OTel collector sidecar above the cluster default. The override is validated at decode time (positive quantities, >=1Gi memory floor) and dropped with a warning if invalid so translate falls back to the safe cluster default. At translate time it is overlaid per resource key onto the cluster default, and any limit below its request is raised to keep the pod spec valid. Signed-off-by: shobham <shobham@nvidia.com>
1 parent bd89567 commit 16ac563

6 files changed

Lines changed: 220 additions & 1 deletion

File tree

src/compute-plane-services/nvca/internal/miniservice/translate_workload.go

Lines changed: 32 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -51,7 +51,7 @@ func (r *Reconciler) translateWorkload(
5151
InstanceTypeLabelSelectorKey: nodefeatures.UniformInstanceTypeLabelKey,
5252
WorkloadResources: corev1.ResourceRequirements{},
5353
Tolerations: append([]corev1.Toleration(nil), r.cfg.Workload.Tolerations...),
54-
OTelResources: k8sutil.GetContainerResourcesBYOO(r.cfg),
54+
OTelResources: mergeBYOOResources(k8sutil.GetContainerResourcesBYOO(r.cfg), ms.Spec.WorkloadConfig.GetBYOOResources()),
5555
FluentbitResources: k8sutil.GetContainerResourcesFluentBit(r.cfg),
5656
FluentbitEnabled: r.FeatureFlagFetcher.IsFeatureFlagEnabled(featureflag.BYOOFluentBit),
5757
ClusterRegion: r.ClusterRegion,
@@ -127,6 +127,37 @@ func (r *Reconciler) translateTaskWorkload(
127127
return metaToClientObjs(objs), nil
128128
}
129129

130+
// mergeBYOOResources overlays a per-workload BYOO collector resource override on top of
131+
// the cluster-level default. Only the resource keys present in the override are replaced,
132+
// so a function can bump e.g. memory while keeping the default CPU. A nil override returns
133+
// the base unchanged. As a safety net, any limit that ends up below its request is raised
134+
// to the request so the resulting container spec stays valid.
135+
func mergeBYOOResources(base corev1.ResourceRequirements, override *corev1.ResourceRequirements) corev1.ResourceRequirements {
136+
if override == nil {
137+
return base
138+
}
139+
out := *base.DeepCopy()
140+
apply := func(dst *corev1.ResourceList, src corev1.ResourceList) {
141+
if len(src) == 0 {
142+
return
143+
}
144+
if *dst == nil {
145+
*dst = corev1.ResourceList{}
146+
}
147+
for name, q := range src {
148+
(*dst)[name] = q.DeepCopy()
149+
}
150+
}
151+
apply(&out.Requests, override.Requests)
152+
apply(&out.Limits, override.Limits)
153+
for name, req := range out.Requests {
154+
if lim, ok := out.Limits[name]; ok && req.Cmp(lim) > 0 {
155+
out.Limits[name] = req.DeepCopy()
156+
}
157+
}
158+
return out
159+
}
160+
130161
func metaToClientObjs(mobjs []metav1.Object) (cobjs []client.Object) {
131162
cobjs = make([]client.Object, len(mobjs))
132163
for i, obj := range mobjs {

src/compute-plane-services/nvca/internal/miniservice/translate_workload_test.go

Lines changed: 42 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -27,6 +27,7 @@ import (
2727
"github.com/stretchr/testify/assert"
2828
"github.com/stretchr/testify/require"
2929
corev1 "k8s.io/api/core/v1"
30+
"k8s.io/apimachinery/pkg/api/resource"
3031
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
3132

3233
nvcav1alpha1 "github.com/NVIDIA/nvcf/src/compute-plane-services/nvca/pkg/apis/nvca/v1alpha1"
@@ -115,3 +116,44 @@ func TestTranslateWorkloadLLMUsesCanonicalPylonArgs(t *testing.T) {
115116
assert.Equal(t, []string{"--backend-connectivity=reverse"}, backendConnectivityArgs)
116117
assert.Equal(t, []string{"--initial-input-tps=100"}, initialInputTPSArgs)
117118
}
119+
120+
func TestMergeBYOOResources(t *testing.T) {
121+
base := corev1.ResourceRequirements{
122+
Requests: corev1.ResourceList{
123+
corev1.ResourceCPU: resource.MustParse("1"),
124+
corev1.ResourceMemory: resource.MustParse("2Gi"),
125+
},
126+
Limits: corev1.ResourceList{
127+
corev1.ResourceCPU: resource.MustParse("1"),
128+
corev1.ResourceMemory: resource.MustParse("2Gi"),
129+
},
130+
}
131+
132+
t.Run("nil override returns base unchanged", func(t *testing.T) {
133+
got := mergeBYOOResources(base, nil)
134+
assert.Equal(t, base, got)
135+
})
136+
137+
t.Run("memory-only override keeps default cpu", func(t *testing.T) {
138+
override := &corev1.ResourceRequirements{
139+
Requests: corev1.ResourceList{corev1.ResourceMemory: resource.MustParse("4Gi")},
140+
Limits: corev1.ResourceList{corev1.ResourceMemory: resource.MustParse("4Gi")},
141+
}
142+
got := mergeBYOOResources(base, override)
143+
assert.Equal(t, "1", got.Requests.Cpu().String())
144+
assert.Equal(t, "4Gi", got.Requests.Memory().String())
145+
assert.Equal(t, "1", got.Limits.Cpu().String())
146+
assert.Equal(t, "4Gi", got.Limits.Memory().String())
147+
// base must not be mutated
148+
assert.Equal(t, "2Gi", base.Limits.Memory().String())
149+
})
150+
151+
t.Run("request above limit clamps limit up", func(t *testing.T) {
152+
override := &corev1.ResourceRequirements{
153+
Requests: corev1.ResourceList{corev1.ResourceMemory: resource.MustParse("4Gi")},
154+
}
155+
got := mergeBYOOResources(base, override)
156+
assert.Equal(t, "4Gi", got.Requests.Memory().String())
157+
assert.Equal(t, "4Gi", got.Limits.Memory().String())
158+
})
159+
}

src/compute-plane-services/nvca/pkg/apis/nvca/v1alpha1/miniservice_types.go

Lines changed: 18 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,7 @@ package v1alpha1
1919

2020
import (
2121
"github.com/NVIDIA/nvcf/src/libraries/go/lib/pkg/icms-translate/translate/common"
22+
corev1 "k8s.io/api/core/v1"
2223
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
2324
)
2425

@@ -58,6 +59,14 @@ type WorkloadConfig struct {
5859
// FeatureFlags toggles per-workload behaviors, keyed by feature flag name.
5960
// +optional
6061
FeatureFlags map[string]bool `json:"featureFlags,omitempty"`
62+
63+
// BYOOResources overrides the BYOO OTel collector sidecar's CPU/memory for this
64+
// workload. When set, it takes precedence over the cluster-level default at
65+
// translate time (only the resource keys present are overridden). Intended for
66+
// large-log / high-throughput functions whose collector needs more headroom than
67+
// the cluster default.
68+
// +optional
69+
BYOOResources *corev1.ResourceRequirements `json:"byooResources,omitempty"`
6170
}
6271

6372
// IsFeatureFlagEnabled reports whether the named feature flag is enabled.
@@ -68,6 +77,15 @@ func (c *WorkloadConfig) IsFeatureFlagEnabled(key string) bool {
6877
return c.FeatureFlags[key]
6978
}
7079

80+
// GetBYOOResources returns the per-workload BYOO collector resource override, or nil
81+
// if none is configured. It is nil-safe so callers can use it on a nil WorkloadConfig.
82+
func (c *WorkloadConfig) GetBYOOResources() *corev1.ResourceRequirements {
83+
if c == nil {
84+
return nil
85+
}
86+
return c.BYOOResources
87+
}
88+
7189
type LocalObjectReference struct {
7290
Name string `json:"name"`
7391
}

src/compute-plane-services/nvca/pkg/apis/nvca/v1alpha1/zz_generated.deepcopy.go

Lines changed: 6 additions & 0 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

src/compute-plane-services/nvca/pkg/featureflag/featureflag_workload.go

Lines changed: 32 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -19,13 +19,19 @@ package featureflag
1919

2020
import (
2121
"context"
22+
"fmt"
2223

2324
"github.com/NVIDIA/nvcf/src/compute-plane-services/nvca/pkg/apis/nvca/v1alpha1"
2425
"github.com/go-logr/logr"
2526
corev1 "k8s.io/api/core/v1"
2627
"sigs.k8s.io/yaml"
2728
)
2829

30+
// minBYOOMemoryBytes is the floor for a per-workload BYOO collector memory override.
31+
// Perf testing showed the collector OOMs between 512Mi and 1Gi under burst, so an
32+
// override below 1Gi is rejected (the safe cluster default is used instead).
33+
const minBYOOMemoryBytes = 1 << 30 // 1Gi
34+
2935
const (
3036
// WorkloadConfigConfigMapName is the fixed name of the ConfigMap a chart author may include
3137
// to supply workload-specific configuration. The ConfigMap is never created on-cluster; it is
@@ -75,5 +81,31 @@ func DecodeWorkloadConfig(ctx context.Context, log logr.Logger, cm *corev1.Confi
7581
delete(cfg.FeatureFlags, key)
7682
}
7783
}
84+
85+
// Drop an invalid BYOO resource override rather than failing the deploy: the
86+
// translate path then falls back to the safe cluster-level default.
87+
if cfg.BYOOResources != nil {
88+
if err := validateBYOOResources(cfg.BYOOResources); err != nil {
89+
log.Info("Ignoring invalid byooResources in ConfigMap %q: %v", WorkloadConfigConfigMapName, err)
90+
cfg.BYOOResources = nil
91+
}
92+
}
7893
return &cfg, nil
7994
}
95+
96+
// validateBYOOResources rejects a per-workload BYOO collector resource override that
97+
// would be unsafe: every specified quantity must be positive, and any memory quantity
98+
// must be at least 1Gi (below that the collector OOMs under burst).
99+
func validateBYOOResources(rr *corev1.ResourceRequirements) error {
100+
for kind, list := range map[string]corev1.ResourceList{"requests": rr.Requests, "limits": rr.Limits} {
101+
for name, q := range list {
102+
if q.Sign() <= 0 {
103+
return fmt.Errorf("%s.%s must be positive (got %q)", kind, name, q.String())
104+
}
105+
if name == corev1.ResourceMemory && q.CmpInt64(minBYOOMemoryBytes) < 0 {
106+
return fmt.Errorf("%s.memory must be at least 1Gi (got %q)", kind, q.String())
107+
}
108+
}
109+
}
110+
return nil
111+
}

src/compute-plane-services/nvca/pkg/featureflag/featureflag_workload_test.go

Lines changed: 90 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -25,6 +25,7 @@ import (
2525
"github.com/stretchr/testify/assert"
2626
"github.com/stretchr/testify/require"
2727
corev1 "k8s.io/api/core/v1"
28+
"k8s.io/apimachinery/pkg/api/resource"
2829

2930
"github.com/NVIDIA/nvcf/src/compute-plane-services/nvca/pkg/apis/nvca/v1alpha1"
3031
)
@@ -83,6 +84,36 @@ func TestDecodeWorkloadConfig(t *testing.T) {
8384
FeatureFlags: map[string]bool{StatusByWorkerReadiness: true},
8485
},
8586
},
87+
{
88+
name: "valid byooResources override is kept",
89+
cm: workloadConfigCM("byooResources:\n" +
90+
" requests:\n cpu: \"1\"\n memory: 4Gi\n" +
91+
" limits:\n cpu: \"1\"\n memory: 4Gi\n"),
92+
want: &v1alpha1.WorkloadConfig{
93+
BYOOResources: &corev1.ResourceRequirements{
94+
Requests: corev1.ResourceList{
95+
corev1.ResourceCPU: resource.MustParse("1"),
96+
corev1.ResourceMemory: resource.MustParse("4Gi"),
97+
},
98+
Limits: corev1.ResourceList{
99+
corev1.ResourceCPU: resource.MustParse("1"),
100+
corev1.ResourceMemory: resource.MustParse("4Gi"),
101+
},
102+
},
103+
},
104+
},
105+
{
106+
name: "byooResources below 1Gi memory floor is dropped",
107+
cm: workloadConfigCM("byooResources:\n" +
108+
" limits:\n memory: 512Mi\n"),
109+
want: &v1alpha1.WorkloadConfig{},
110+
},
111+
{
112+
name: "byooResources with non-positive cpu is dropped",
113+
cm: workloadConfigCM("byooResources:\n" +
114+
" requests:\n cpu: \"0\"\n memory: 2Gi\n"),
115+
want: &v1alpha1.WorkloadConfig{},
116+
},
86117
{
87118
name: "invalid yaml returns error",
88119
cm: workloadConfigCM("featureFlags: [not-a-map"),
@@ -126,3 +157,62 @@ func TestWorkloadConfigIsFeatureFlagEnabled(t *testing.T) {
126157
assert.True(t, cfg.IsFeatureFlagEnabled(StatusByWorkerReadiness))
127158
assert.False(t, cfg.IsFeatureFlagEnabled("SomeOtherFlag"))
128159
}
160+
161+
func TestWorkloadConfigGetBYOOResources(t *testing.T) {
162+
var nilCfg *v1alpha1.WorkloadConfig
163+
assert.Nil(t, nilCfg.GetBYOOResources())
164+
165+
empty := &v1alpha1.WorkloadConfig{}
166+
assert.Nil(t, empty.GetBYOOResources())
167+
168+
rr := &corev1.ResourceRequirements{
169+
Limits: corev1.ResourceList{corev1.ResourceMemory: resource.MustParse("4Gi")},
170+
}
171+
cfg := &v1alpha1.WorkloadConfig{BYOOResources: rr}
172+
assert.Equal(t, rr, cfg.GetBYOOResources())
173+
}
174+
175+
func TestValidateBYOOResources(t *testing.T) {
176+
tests := []struct {
177+
name string
178+
rr *corev1.ResourceRequirements
179+
wantErr bool
180+
}{
181+
{
182+
name: "valid requests and limits",
183+
rr: &corev1.ResourceRequirements{
184+
Requests: corev1.ResourceList{corev1.ResourceCPU: resource.MustParse("1"), corev1.ResourceMemory: resource.MustParse("4Gi")},
185+
Limits: corev1.ResourceList{corev1.ResourceCPU: resource.MustParse("1"), corev1.ResourceMemory: resource.MustParse("4Gi")},
186+
},
187+
},
188+
{
189+
name: "memory exactly at 1Gi floor is valid",
190+
rr: &corev1.ResourceRequirements{Limits: corev1.ResourceList{corev1.ResourceMemory: resource.MustParse("1Gi")}},
191+
},
192+
{
193+
name: "memory below 1Gi floor is rejected",
194+
rr: &corev1.ResourceRequirements{Limits: corev1.ResourceList{corev1.ResourceMemory: resource.MustParse("512Mi")}},
195+
wantErr: true,
196+
},
197+
{
198+
name: "zero cpu is rejected",
199+
rr: &corev1.ResourceRequirements{Requests: corev1.ResourceList{corev1.ResourceCPU: resource.MustParse("0")}},
200+
wantErr: true,
201+
},
202+
{
203+
name: "negative memory is rejected",
204+
rr: &corev1.ResourceRequirements{Limits: corev1.ResourceList{corev1.ResourceMemory: resource.MustParse("-1Gi")}},
205+
wantErr: true,
206+
},
207+
}
208+
for _, tt := range tests {
209+
t.Run(tt.name, func(t *testing.T) {
210+
err := validateBYOOResources(tt.rr)
211+
if tt.wantErr {
212+
require.Error(t, err)
213+
return
214+
}
215+
require.NoError(t, err)
216+
})
217+
}
218+
}

0 commit comments

Comments
 (0)