Skip to content

Commit 1b4e2bc

Browse files
authored
Merge pull request #424 from Ygnas/fix/controller-robustness-improvements
Fix: harden controllers against stale objects, orphaned resources, an…
2 parents f68d4b2 + bcde375 commit 1b4e2bc

9 files changed

Lines changed: 122 additions & 28 deletions

File tree

kagenti-operator/cmd/main.go

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -587,6 +587,7 @@ func main() {
587587
if controller.CertManagerCRDExists(mgr.GetConfig()) {
588588
if err = (&controller.SharedTrustReconciler{
589589
Client: mgr.GetClient(),
590+
Scheme: mgr.GetScheme(),
590591
Recorder: mgr.GetEventRecorderFor("shared-trust-controller"), //nolint:staticcheck
591592
}).SetupWithManager(mgr); err != nil {
592593
setupLog.Error(err, "unable to create controller", "controller", "SharedTrust")
@@ -633,6 +634,7 @@ func main() {
633634
Client: mgr.GetClient(),
634635
APIReader: mgr.GetAPIReader(),
635636
Config: mgr.GetConfig(),
637+
Scheme: mgr.GetScheme(),
636638
Namespace: getOperatorNamespace(),
637639
Log: ctrl.Log.WithName("bootstrap"),
638640
MLflowWorkspace: mlflowWorkspace,

kagenti-operator/internal/bootstrap/otel.go

Lines changed: 43 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -34,11 +34,13 @@ import (
3434
"k8s.io/apimachinery/pkg/api/errors"
3535
"k8s.io/apimachinery/pkg/api/meta"
3636
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
37+
"k8s.io/apimachinery/pkg/runtime"
3738
"k8s.io/apimachinery/pkg/types"
3839
"k8s.io/apimachinery/pkg/util/wait"
3940
"k8s.io/client-go/discovery"
4041
"k8s.io/client-go/rest"
4142
"sigs.k8s.io/controller-runtime/pkg/client"
43+
"sigs.k8s.io/controller-runtime/pkg/controller/controllerutil"
4244
"sigs.k8s.io/yaml"
4345

4446
"github.com/kagenti/operator/internal/mlflow"
@@ -76,6 +78,7 @@ type OtelBootstrapRunnable struct {
7678
Client client.Client
7779
APIReader client.Reader
7880
Config *rest.Config
81+
Scheme *runtime.Scheme
7982
Namespace string
8083
Log logr.Logger
8184

@@ -93,27 +96,61 @@ type OtelBootstrapRunnable struct {
9396
EnsureExperiment func(ctx context.Context, baseURL, workspace string) (string, error)
9497
}
9598

99+
// operatorDeploymentNames lists possible Deployment names for the operator itself,
100+
// used to set OwnerReferences on bootstrap-created resources.
101+
var operatorDeploymentNames = []string{
102+
"kagenti-controller-manager",
103+
"controller-manager",
104+
}
105+
106+
// getOperatorOwner looks up the operator's own Deployment to use as an OwnerReference.
107+
// Returns nil if the deployment cannot be found (best-effort).
108+
func (r *OtelBootstrapRunnable) getOperatorOwner(ctx context.Context, log logr.Logger) *appsv1.Deployment {
109+
for _, name := range operatorDeploymentNames {
110+
deploy := &appsv1.Deployment{}
111+
key := types.NamespacedName{Name: name, Namespace: r.Namespace}
112+
if err := r.Client.Get(ctx, key, deploy); err == nil {
113+
return deploy
114+
}
115+
}
116+
log.Info("Could not find operator Deployment for OwnerReference, ConfigMaps will be unowned")
117+
return nil
118+
}
119+
120+
// setOwnerIfAvailable sets an OwnerReference on the given object if an owner is available.
121+
func (r *OtelBootstrapRunnable) setOwnerIfAvailable(owner *appsv1.Deployment, obj client.Object, log logr.Logger) {
122+
if owner == nil || r.Scheme == nil {
123+
return
124+
}
125+
if err := controllerutil.SetOwnerReference(owner, obj, r.Scheme); err != nil {
126+
log.Error(err, "Failed to set OwnerReference on resource", "name", obj.GetName())
127+
}
128+
}
129+
96130
// Start runs the bootstrap sequence. Called by the manager after leader election
97131
// and cache sync, before controllers start processing events.
98132
func (r *OtelBootstrapRunnable) Start(ctx context.Context) error {
99133
log := r.Log.WithName("otel-bootstrap")
100134
log.Info("Starting OTel collector bootstrap")
101135

136+
// Look up operator Deployment once for OwnerReference on created resources.
137+
owner := r.getOperatorOwner(ctx, log)
138+
102139
isOCP, err := r.detectOpenShift(ctx)
103140
if err != nil {
104141
return fmt.Errorf("detecting OpenShift: %w", err)
105142
}
106143

107144
if isOCP {
108145
log.Info("OpenShift detected, reconciling ingress CA trust")
109-
if err := r.reconcileIngressCA(ctx, log); err != nil {
146+
if err := r.reconcileIngressCA(ctx, log, owner); err != nil {
110147
return fmt.Errorf("ingress CA bootstrap: %w", err)
111148
}
112149
} else {
113150
log.Info("Not running on OpenShift, skipping ingress CA trust")
114151
}
115152

116-
if err := r.reconcileCollectorConfig(ctx, log, isOCP); err != nil {
153+
if err := r.reconcileCollectorConfig(ctx, log, isOCP, owner); err != nil {
117154
return fmt.Errorf("collector config bootstrap: %w", err)
118155
}
119156

@@ -156,7 +193,7 @@ func (r *OtelBootstrapRunnable) detectOpenShift(ctx context.Context) (bool, erro
156193

157194
// reconcileIngressCA reads the OpenShift ingress CA and root CA, then creates
158195
// or updates the otel-ingress-ca ConfigMap in the operator namespace.
159-
func (r *OtelBootstrapRunnable) reconcileIngressCA(ctx context.Context, log logr.Logger) error {
196+
func (r *OtelBootstrapRunnable) reconcileIngressCA(ctx context.Context, log logr.Logger, owner *appsv1.Deployment) error {
160197
ingressCert := &corev1.ConfigMap{}
161198
key := types.NamespacedName{Name: ingressCertConfigMap, Namespace: ingressCertNamespace}
162199
if err := r.APIReader.Get(ctx, key, ingressCert); err != nil {
@@ -197,6 +234,7 @@ func (r *OtelBootstrapRunnable) reconcileIngressCA(ctx context.Context, log logr
197234
},
198235
Data: map[string]string{caBundleKey: caBundle},
199236
}
237+
r.setOwnerIfAvailable(owner, cm, log)
200238
if err := r.Client.Create(ctx, cm); err != nil {
201239
if !errors.IsAlreadyExists(err) {
202240
return fmt.Errorf("creating %s ConfigMap: %w", ingressCAConfigMap, err)
@@ -231,7 +269,7 @@ func (r *OtelBootstrapRunnable) reconcileIngressCA(ctx context.Context, log logr
231269

232270
// reconcileCollectorConfig discovers available components and assembles the
233271
// OTel collector ConfigMap from preset configurations.
234-
func (r *OtelBootstrapRunnable) reconcileCollectorConfig(ctx context.Context, log logr.Logger, isOCP bool) error {
272+
func (r *OtelBootstrapRunnable) reconcileCollectorConfig(ctx context.Context, log logr.Logger, isOCP bool, owner *appsv1.Deployment) error {
235273
mf, err := r.discoverMLflow(ctx, log)
236274
if err != nil {
237275
return err
@@ -291,6 +329,7 @@ func (r *OtelBootstrapRunnable) reconcileCollectorConfig(ctx context.Context, lo
291329
},
292330
Data: map[string]string{configMapDataKey: configStr},
293331
}
332+
r.setOwnerIfAvailable(owner, cm, log)
294333
if err := r.Client.Create(ctx, cm); err != nil {
295334
if !errors.IsAlreadyExists(err) {
296335
return fmt.Errorf("creating collector ConfigMap: %w", err)

kagenti-operator/internal/bootstrap/otel_test.go

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -52,6 +52,7 @@ func newRunnable(cl client.Client, isOCP func(context.Context) (bool, error), ml
5252
return &OtelBootstrapRunnable{
5353
Client: cl,
5454
APIReader: cl,
55+
Scheme: testScheme(),
5556
Namespace: testNamespace,
5657
Log: testLogger(),
5758
IsOpenShift: isOCP,

kagenti-operator/internal/controller/agentruntime_controller.go

Lines changed: 20 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -154,11 +154,7 @@ func (r *AgentRuntimeReconciler) Reconcile(ctx context.Context, req ctrl.Request
154154
// 4. Resolve targetRef (existence check)
155155
if err := r.resolveTargetRef(ctx, rt); err != nil {
156156
logger.Error(err, "Failed to resolve targetRef")
157-
r.setPhase(rt, agentv1alpha1.RuntimePhaseError)
158-
r.setCondition(rt, ConditionTypeTargetResolved, metav1.ConditionFalse, "TargetNotFound", err.Error())
159-
if statusErr := r.Status().Update(ctx, rt); statusErr != nil {
160-
logger.Error(statusErr, "Failed to update status")
161-
}
157+
r.updateErrorStatus(ctx, req.NamespacedName, ConditionTypeTargetResolved, "TargetNotFound", err.Error())
162158
if r.Recorder != nil {
163159
r.Recorder.Event(rt, corev1.EventTypeWarning, "TargetNotFound", err.Error())
164160
}
@@ -210,11 +206,7 @@ func (r *AgentRuntimeReconciler) Reconcile(ctx context.Context, req ctrl.Request
210206
configResult, err := ComputeConfigHash(ctx, r.Client, rt.Namespace)
211207
if err != nil {
212208
logger.Error(err, "Failed to compute config hash")
213-
r.setPhase(rt, agentv1alpha1.RuntimePhaseError)
214-
r.setCondition(rt, ConditionTypeReady, metav1.ConditionFalse, "ConfigHashError", err.Error())
215-
if statusErr := r.Status().Update(ctx, rt); statusErr != nil {
216-
logger.Error(statusErr, "Failed to update status")
217-
}
209+
r.updateErrorStatus(ctx, req.NamespacedName, ConditionTypeReady, "ConfigHashError", err.Error())
218210
return ctrl.Result{RequeueAfter: 30 * time.Second}, nil
219211
}
220212

@@ -238,11 +230,7 @@ func (r *AgentRuntimeReconciler) Reconcile(ctx context.Context, req ctrl.Request
238230
// 6. Apply labels and annotations to the target workload
239231
if err := r.applyWorkloadConfig(ctx, rt, configResult.Hash); err != nil {
240232
logger.Error(err, "Failed to apply workload config")
241-
r.setPhase(rt, agentv1alpha1.RuntimePhaseError)
242-
r.setCondition(rt, ConditionTypeReady, metav1.ConditionFalse, "ConfigApplyError", err.Error())
243-
if statusErr := r.Status().Update(ctx, rt); statusErr != nil {
244-
logger.Error(statusErr, "Failed to update status")
245-
}
233+
r.updateErrorStatus(ctx, req.NamespacedName, ConditionTypeReady, "ConfigApplyError", err.Error())
246234
return ctrl.Result{RequeueAfter: 30 * time.Second}, nil
247235
}
248236

@@ -777,6 +765,23 @@ func (r *AgentRuntimeReconciler) setCondition(rt *agentv1alpha1.AgentRuntime, co
777765
})
778766
}
779767

768+
// updateErrorStatus sets the AgentRuntime phase to Error and updates a condition
769+
// with retry-on-conflict semantics, re-fetching the object on each attempt.
770+
func (r *AgentRuntimeReconciler) updateErrorStatus(ctx context.Context, key types.NamespacedName, condType, reason, message string) {
771+
logger := log.FromContext(ctx)
772+
if statusErr := retry.RetryOnConflict(retry.DefaultRetry, func() error {
773+
latest := &agentv1alpha1.AgentRuntime{}
774+
if err := r.Get(ctx, key, latest); err != nil {
775+
return err
776+
}
777+
r.setPhase(latest, agentv1alpha1.RuntimePhaseError)
778+
r.setCondition(latest, condType, metav1.ConditionFalse, reason, message)
779+
return r.Status().Update(ctx, latest)
780+
}); statusErr != nil {
781+
logger.Error(statusErr, "Failed to update error status", "condition", condType, "reason", reason)
782+
}
783+
}
784+
780785
// fetchAndUpdateCard discovers the agent card from the workload's Service endpoint
781786
// and populates status.card. Skips fetch when the feature flag is disabled or
782787
// when the workload's change-detection key has not changed.

kagenti-operator/internal/controller/mlflow_controller.go

Lines changed: 20 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -223,22 +223,35 @@ func (r *MLflowReconciler) configureDeployment(ctx context.Context, dep *appsv1.
223223
if annotations == nil {
224224
annotations = make(map[string]string)
225225
}
226-
annotations[AnnotationMLflowExperimentID] = experimentID
227-
annotations[AnnotationMLflowExperimentName] = experimentName
228-
annotations[AnnotationMLflowTrackingURI] = trackingURI
229-
annotations[AnnotationMLflowTrackingAuth] = "kubernetes-namespaced"
226+
227+
annotationsChanged := false
228+
for k, v := range map[string]string{
229+
AnnotationMLflowExperimentID: experimentID,
230+
AnnotationMLflowExperimentName: experimentName,
231+
AnnotationMLflowTrackingURI: trackingURI,
232+
AnnotationMLflowTrackingAuth: "kubernetes-namespaced",
233+
} {
234+
if annotations[k] != v {
235+
annotations[k] = v
236+
annotationsChanged = true
237+
}
238+
}
230239
latest.Spec.Template.Annotations = annotations
231240

232-
changed := false
241+
envChanged := false
233242
for i := range latest.Spec.Template.Spec.Containers {
234243
for name, value := range desired {
235244
if setEnvVar(&latest.Spec.Template.Spec.Containers[i], name, value) {
236-
changed = true
245+
envChanged = true
237246
}
238247
}
239248
}
240249

241-
if changed {
250+
if !envChanged && !annotationsChanged {
251+
return nil
252+
}
253+
254+
if envChanged {
242255
logger.Info("Injected MLflow env vars into Deployment containers", "deployment", dep.Name)
243256
}
244257

kagenti-operator/internal/controller/sharedtrust_controller.go

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -26,6 +26,7 @@ import (
2626
appsv1 "k8s.io/api/apps/v1"
2727
corev1 "k8s.io/api/core/v1"
2828
apierrors "k8s.io/apimachinery/pkg/api/errors"
29+
"k8s.io/apimachinery/pkg/runtime"
2930
"k8s.io/apimachinery/pkg/types"
3031
"k8s.io/apimachinery/pkg/util/wait"
3132
"k8s.io/client-go/discovery"
@@ -34,6 +35,7 @@ import (
3435
"k8s.io/client-go/util/retry"
3536
ctrl "sigs.k8s.io/controller-runtime"
3637
"sigs.k8s.io/controller-runtime/pkg/client"
38+
"sigs.k8s.io/controller-runtime/pkg/controller/controllerutil"
3739
"sigs.k8s.io/controller-runtime/pkg/handler"
3840
"sigs.k8s.io/controller-runtime/pkg/log"
3941
"sigs.k8s.io/controller-runtime/pkg/reconcile"
@@ -100,6 +102,7 @@ var (
100102

101103
type SharedTrustReconciler struct {
102104
client.Client
105+
Scheme *runtime.Scheme
103106
Recorder record.EventRecorder
104107
}
105108

@@ -224,6 +227,9 @@ func (r *SharedTrustReconciler) reconcileCacertsSecrets(ctx context.Context) (bo
224227
secret.Namespace = ic.Namespace
225228
secret.Type = corev1.SecretTypeOpaque
226229
secret.Data = desired
230+
if err := controllerutil.SetOwnerReference(intSecret, secret, r.Scheme); err != nil {
231+
return false, fmt.Errorf("setting owner reference for cacerts secret in %s: %w", ic.Namespace, err)
232+
}
227233
if err := r.Create(ctx, secret); err != nil {
228234
return false, fmt.Errorf("creating cacerts secret in %s: %w", ic.Namespace, err)
229235
}

kagenti-operator/internal/controller/sharedtrust_controller_test.go

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -187,6 +187,7 @@ func newReconciler(t *testing.T, objs ...runtime.Object) *SharedTrustReconciler
187187
cb := fake.NewClientBuilder().WithScheme(scheme).WithRuntimeObjects(clientObjs...)
188188
return &SharedTrustReconciler{
189189
Client: cb.Build(),
190+
Scheme: scheme,
190191
Recorder: record.NewFakeRecorder(10),
191192
}
192193
}

kagenti-operator/internal/signature/verifier.go

Lines changed: 20 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -353,8 +353,26 @@ func removeEmptyFields(m map[string]interface{}) map[string]interface{} {
353353
result[k] = cleaned
354354
}
355355
case []interface{}:
356-
if len(val) > 0 {
357-
result[k] = val
356+
var cleaned []interface{}
357+
for _, elem := range val {
358+
switch e := elem.(type) {
359+
case map[string]interface{}:
360+
c := removeEmptyFields(e)
361+
if len(c) > 0 {
362+
cleaned = append(cleaned, c)
363+
}
364+
case string:
365+
if e != "" {
366+
cleaned = append(cleaned, e)
367+
}
368+
case nil:
369+
// skip nil elements
370+
default:
371+
cleaned = append(cleaned, e)
372+
}
373+
}
374+
if len(cleaned) > 0 {
375+
result[k] = cleaned
358376
}
359377
case string:
360378
if val != "" {

scripts/kind-with-registry.yaml

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -19,9 +19,18 @@ spec:
1919
labels:
2020
app: registry
2121
spec:
22+
securityContext:
23+
runAsNonRoot: true
24+
seccompProfile:
25+
type: RuntimeDefault
2226
containers:
2327
- name: registry
2428
image: public.ecr.aws/docker/library/registry:3.0.0-rc.4
29+
securityContext:
30+
allowPrivilegeEscalation: false
31+
capabilities:
32+
drop:
33+
- ALL
2534
ports:
2635
- containerPort: 5000
2736
name: registry

0 commit comments

Comments
 (0)