Skip to content

Commit 9474fa8

Browse files
committed
fix(k8s): harden client init, attach write, and quota parsing
Address three robustness issues in the Kubernetes backend: - api.go: Client() stored the init error in a function-local variable, so after a failed first call the sync.Once never re-ran and subsequent calls returned (nil, nil), masking the failure. Persist the error in a package variable and return it on every call. - container.go: SendCommand() checked IsAttached() before taking the lock that guards e.stream, so the stream could be cleared in between and cause a nil-pointer write. Nil-check e.stream under the read lock instead. - quota.go: buildResourceQuota/buildLimitRange used resource.MustParse on user-configurable quantity strings, panicking the process on invalid values. Parse with resource.ParseQuantity and propagate errors through EnsureResourceQuota/EnsureLimitRange.
1 parent 63bc480 commit 9474fa8

5 files changed

Lines changed: 95 additions & 53 deletions

File tree

‎environment/kubernetes/api.go‎

Lines changed: 9 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -14,29 +14,29 @@ import (
1414
var (
1515
_konce sync.Once
1616
_client kubernetes.Interface
17+
_kerr error
1718
)
1819

1920
// Client returns a shared Kubernetes clientset. The client is created once
2021
// and reused for all subsequent calls. It uses in-cluster config when no
2122
// kubeconfig path is specified, falling back to the provided kubeconfig file.
2223
func Client() (kubernetes.Interface, error) {
23-
var err error
2424
_konce.Do(func() {
2525
var cfg *rest.Config
2626
kubeconfig := config.Get().Kubernetes.Kubeconfig
2727
if kubeconfig != "" {
28-
cfg, err = clientcmd.BuildConfigFromFlags("", kubeconfig)
28+
cfg, _kerr = clientcmd.BuildConfigFromFlags("", kubeconfig)
2929
} else {
30-
cfg, err = rest.InClusterConfig()
30+
cfg, _kerr = rest.InClusterConfig()
3131
}
32-
if err != nil {
33-
err = errors.Wrap(err, "environment/kubernetes: failed to build config")
32+
if _kerr != nil {
33+
_kerr = errors.Wrap(_kerr, "environment/kubernetes: failed to build config")
3434
return
3535
}
36-
_client, err = kubernetes.NewForConfig(cfg)
37-
if err != nil {
38-
err = errors.Wrap(err, "environment/kubernetes: failed to create clientset")
36+
_client, _kerr = kubernetes.NewForConfig(cfg)
37+
if _kerr != nil {
38+
_kerr = errors.Wrap(_kerr, "environment/kubernetes: failed to create clientset")
3939
}
4040
})
41-
return _client, err
41+
return _client, _kerr
4242
}

‎environment/kubernetes/container.go‎

Lines changed: 6 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -348,13 +348,15 @@ func (e *Environment) Attach(ctx context.Context) error {
348348

349349
// SendCommand writes a command string to the attached Pod's stdin.
350350
func (e *Environment) SendCommand(c string) error {
351-
if !e.IsAttached() {
352-
return errors.New("environment/kubernetes: not attached to pod")
353-
}
354-
355351
e.mu.RLock()
356352
defer e.mu.RUnlock()
357353

354+
// Check the stream under the lock so it cannot be cleared between an
355+
// attachment check and the write below.
356+
if e.stream == nil {
357+
return errors.New("environment/kubernetes: not attached to pod")
358+
}
359+
358360
// If this is the stop command, mark the server as stopping.
359361
if e.meta.Stop.Type == "command" && c == e.meta.Stop.Value {
360362
e.SetState(environment.ProcessStoppingState)

‎environment/kubernetes/container_test.go‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -9,6 +9,8 @@ import (
99
"github.com/exonical/wings/config"
1010
)
1111

12+
// TestResolveImagePullPolicy verifies pull-policy resolution for remote
13+
// images, ~-prefixed local images, and configured overrides.
1214
func TestResolveImagePullPolicy(t *testing.T) {
1315
g := Goblin(t)
1416

‎environment/kubernetes/quota.go‎

Lines changed: 58 additions & 38 deletions
Original file line numberDiff line numberDiff line change
@@ -12,7 +12,7 @@ import (
1212
)
1313

1414
const (
15-
quotaName = "pelican-wings"
15+
quotaName = "pelican-wings"
1616
limitRangeName = "pelican-wings"
1717
)
1818

@@ -25,7 +25,10 @@ func (e *Environment) EnsureResourceQuota(ctx context.Context) error {
2525
}
2626

2727
ns := e.namespace()
28-
quota := buildResourceQuota(ns, &cfg.Kubernetes.ResourceQuota)
28+
quota, err := buildResourceQuota(ns, &cfg.Kubernetes.ResourceQuota)
29+
if err != nil {
30+
return err
31+
}
2932

3033
existing, err := e.client.CoreV1().ResourceQuotas(ns).Get(ctx, quotaName, metav1.GetOptions{})
3134
if err == nil {
@@ -56,7 +59,10 @@ func (e *Environment) EnsureLimitRange(ctx context.Context) error {
5659
}
5760

5861
ns := e.namespace()
59-
lr := buildLimitRange(ns, &cfg.Kubernetes.LimitRange)
62+
lr, err := buildLimitRange(ns, &cfg.Kubernetes.LimitRange)
63+
if err != nil {
64+
return err
65+
}
6066

6167
existing, err := e.client.CoreV1().LimitRanges(ns).Get(ctx, limitRangeName, metav1.GetOptions{})
6268
if err == nil {
@@ -78,31 +84,46 @@ func (e *Environment) EnsureLimitRange(ctx context.Context) error {
7884
return nil
7985
}
8086

87+
// setQuantity parses a quantity string into the resource list under key when
88+
// the value is non-empty. Unlike resource.MustParse it returns an error for
89+
// invalid (user-configurable) values instead of panicking the process.
90+
func setQuantity(list corev1.ResourceList, key corev1.ResourceName, field, value string) error {
91+
if value == "" {
92+
return nil
93+
}
94+
q, err := resource.ParseQuantity(value)
95+
if err != nil {
96+
return errors.Wrapf(err, "environment/kubernetes: invalid quantity %q for %s", value, field)
97+
}
98+
list[key] = q
99+
return nil
100+
}
101+
81102
// buildResourceQuota constructs the ResourceQuota spec from config values.
82-
func buildResourceQuota(namespace string, cfg *config.KubeResourceQuota) *corev1.ResourceQuota {
103+
func buildResourceQuota(namespace string, cfg *config.KubeResourceQuota) (*corev1.ResourceQuota, error) {
83104
hard := corev1.ResourceList{}
84105

85-
if cfg.CPULimit != "" {
86-
hard[corev1.ResourceLimitsCPU] = resource.MustParse(cfg.CPULimit)
87-
}
88-
if cfg.MemoryLimit != "" {
89-
hard[corev1.ResourceLimitsMemory] = resource.MustParse(cfg.MemoryLimit)
90-
}
91-
if cfg.CPURequest != "" {
92-
hard[corev1.ResourceRequestsCPU] = resource.MustParse(cfg.CPURequest)
93-
}
94-
if cfg.MemoryRequest != "" {
95-
hard[corev1.ResourceRequestsMemory] = resource.MustParse(cfg.MemoryRequest)
106+
for _, q := range []struct {
107+
key corev1.ResourceName
108+
field string
109+
value string
110+
}{
111+
{corev1.ResourceLimitsCPU, "cpu_limit", cfg.CPULimit},
112+
{corev1.ResourceLimitsMemory, "memory_limit", cfg.MemoryLimit},
113+
{corev1.ResourceRequestsCPU, "cpu_request", cfg.CPURequest},
114+
{corev1.ResourceRequestsMemory, "memory_request", cfg.MemoryRequest},
115+
{corev1.ResourceRequestsStorage, "max_storage", cfg.MaxStorage},
116+
} {
117+
if err := setQuantity(hard, q.key, q.field, q.value); err != nil {
118+
return nil, err
119+
}
96120
}
97121
if cfg.MaxPods > 0 {
98122
hard[corev1.ResourcePods] = *resource.NewQuantity(cfg.MaxPods, resource.DecimalSI)
99123
}
100124
if cfg.MaxPVCs > 0 {
101125
hard[corev1.ResourcePersistentVolumeClaims] = *resource.NewQuantity(cfg.MaxPVCs, resource.DecimalSI)
102126
}
103-
if cfg.MaxStorage != "" {
104-
hard[corev1.ResourceRequestsStorage] = resource.MustParse(cfg.MaxStorage)
105-
}
106127

107128
return &corev1.ResourceQuota{
108129
ObjectMeta: metav1.ObjectMeta{
@@ -115,35 +136,34 @@ func buildResourceQuota(namespace string, cfg *config.KubeResourceQuota) *corev1
115136
Spec: corev1.ResourceQuotaSpec{
116137
Hard: hard,
117138
},
118-
}
139+
}, nil
119140
}
120141

121142
// buildLimitRange constructs the LimitRange spec from config values.
122-
func buildLimitRange(namespace string, cfg *config.KubeLimitRange) *corev1.LimitRange {
143+
func buildLimitRange(namespace string, cfg *config.KubeLimitRange) (*corev1.LimitRange, error) {
123144
containerLimit := corev1.LimitRangeItem{
124145
Type: corev1.LimitTypeContainer,
125146
Default: corev1.ResourceList{},
126147
DefaultRequest: corev1.ResourceList{},
127148
Max: corev1.ResourceList{},
128149
}
129150

130-
if cfg.DefaultCPULimit != "" {
131-
containerLimit.Default[corev1.ResourceCPU] = resource.MustParse(cfg.DefaultCPULimit)
132-
}
133-
if cfg.DefaultMemoryLimit != "" {
134-
containerLimit.Default[corev1.ResourceMemory] = resource.MustParse(cfg.DefaultMemoryLimit)
135-
}
136-
if cfg.DefaultCPURequest != "" {
137-
containerLimit.DefaultRequest[corev1.ResourceCPU] = resource.MustParse(cfg.DefaultCPURequest)
138-
}
139-
if cfg.DefaultMemoryRequest != "" {
140-
containerLimit.DefaultRequest[corev1.ResourceMemory] = resource.MustParse(cfg.DefaultMemoryRequest)
141-
}
142-
if cfg.MaxCPU != "" {
143-
containerLimit.Max[corev1.ResourceCPU] = resource.MustParse(cfg.MaxCPU)
144-
}
145-
if cfg.MaxMemory != "" {
146-
containerLimit.Max[corev1.ResourceMemory] = resource.MustParse(cfg.MaxMemory)
151+
for _, q := range []struct {
152+
list corev1.ResourceList
153+
key corev1.ResourceName
154+
field string
155+
value string
156+
}{
157+
{containerLimit.Default, corev1.ResourceCPU, "default_cpu_limit", cfg.DefaultCPULimit},
158+
{containerLimit.Default, corev1.ResourceMemory, "default_memory_limit", cfg.DefaultMemoryLimit},
159+
{containerLimit.DefaultRequest, corev1.ResourceCPU, "default_cpu_request", cfg.DefaultCPURequest},
160+
{containerLimit.DefaultRequest, corev1.ResourceMemory, "default_memory_request", cfg.DefaultMemoryRequest},
161+
{containerLimit.Max, corev1.ResourceCPU, "max_cpu", cfg.MaxCPU},
162+
{containerLimit.Max, corev1.ResourceMemory, "max_memory", cfg.MaxMemory},
163+
} {
164+
if err := setQuantity(q.list, q.key, q.field, q.value); err != nil {
165+
return nil, err
166+
}
147167
}
148168

149169
return &corev1.LimitRange{
@@ -157,5 +177,5 @@ func buildLimitRange(namespace string, cfg *config.KubeLimitRange) *corev1.Limit
157177
Spec: corev1.LimitRangeSpec{
158178
Limits: []corev1.LimitRangeItem{containerLimit},
159179
},
160-
}
180+
}, nil
161181
}

‎environment/kubernetes/quota_test.go‎

Lines changed: 20 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,8 @@ import (
1515
"github.com/exonical/wings/system"
1616
)
1717

18+
// TestQuota covers building and reconciling ResourceQuota and LimitRange
19+
// objects from configuration, including invalid-quantity error handling.
1820
func TestQuota(t *testing.T) {
1921
g := Goblin(t)
2022

@@ -264,7 +266,8 @@ func TestQuota(t *testing.T) {
264266
MaxPods: 10,
265267
MaxPVCs: 0,
266268
}
267-
rq := buildResourceQuota("test-ns", cfg)
269+
rq, err := buildResourceQuota("test-ns", cfg)
270+
g.Assert(err).IsNil()
268271
g.Assert(rq.Namespace).Equal("test-ns")
269272

270273
_, hasCPU := rq.Spec.Hard[corev1.ResourceLimitsCPU]
@@ -289,7 +292,8 @@ func TestQuota(t *testing.T) {
289292
MaxCPU: "8",
290293
MaxMemory: "",
291294
}
292-
lr := buildLimitRange("test-ns", cfg)
295+
lr, err := buildLimitRange("test-ns", cfg)
296+
g.Assert(err).IsNil()
293297
g.Assert(lr.Namespace).Equal("test-ns")
294298
g.Assert(len(lr.Spec.Limits)).Equal(1)
295299

@@ -306,6 +310,20 @@ func TestQuota(t *testing.T) {
306310
_, hasMaxMem := item.Max[corev1.ResourceMemory]
307311
g.Assert(hasMaxMem).IsFalse()
308312
})
313+
314+
g.It("should return an error for an invalid quantity instead of panicking", func() {
315+
lr, err := buildLimitRange("test-ns", &config.KubeLimitRange{DefaultCPULimit: "not-a-quantity"})
316+
g.Assert(err != nil).IsTrue()
317+
g.Assert(lr == nil).IsTrue()
318+
})
319+
})
320+
321+
g.Describe("buildResourceQuota invalid input", func() {
322+
g.It("should return an error for an invalid quantity instead of panicking", func() {
323+
rq, err := buildResourceQuota("test-ns", &config.KubeResourceQuota{CPULimit: "not-a-quantity"})
324+
g.Assert(err != nil).IsTrue()
325+
g.Assert(rq == nil).IsTrue()
326+
})
309327
})
310328
})
311329
}

0 commit comments

Comments
 (0)