Skip to content
Merged
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
2 changes: 2 additions & 0 deletions kagenti-operator/cmd/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -587,6 +587,7 @@ func main() {
if controller.CertManagerCRDExists(mgr.GetConfig()) {
if err = (&controller.SharedTrustReconciler{
Client: mgr.GetClient(),
Scheme: mgr.GetScheme(),
Recorder: mgr.GetEventRecorderFor("shared-trust-controller"), //nolint:staticcheck
}).SetupWithManager(mgr); err != nil {
setupLog.Error(err, "unable to create controller", "controller", "SharedTrust")
Expand Down Expand Up @@ -633,6 +634,7 @@ func main() {
Client: mgr.GetClient(),
APIReader: mgr.GetAPIReader(),
Config: mgr.GetConfig(),
Scheme: mgr.GetScheme(),
Namespace: getOperatorNamespace(),
Log: ctrl.Log.WithName("bootstrap"),
MLflowWorkspace: mlflowWorkspace,
Expand Down
47 changes: 43 additions & 4 deletions kagenti-operator/internal/bootstrap/otel.go
Original file line number Diff line number Diff line change
Expand Up @@ -34,11 +34,13 @@ import (
"k8s.io/apimachinery/pkg/api/errors"
"k8s.io/apimachinery/pkg/api/meta"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/runtime"
"k8s.io/apimachinery/pkg/types"
"k8s.io/apimachinery/pkg/util/wait"
"k8s.io/client-go/discovery"
"k8s.io/client-go/rest"
"sigs.k8s.io/controller-runtime/pkg/client"
"sigs.k8s.io/controller-runtime/pkg/controller/controllerutil"
"sigs.k8s.io/yaml"

"github.com/kagenti/operator/internal/mlflow"
Expand Down Expand Up @@ -76,6 +78,7 @@ type OtelBootstrapRunnable struct {
Client client.Client
APIReader client.Reader
Config *rest.Config
Scheme *runtime.Scheme
Namespace string
Log logr.Logger

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

// operatorDeploymentNames lists possible Deployment names for the operator itself,
// used to set OwnerReferences on bootstrap-created resources.
var operatorDeploymentNames = []string{
"kagenti-controller-manager",
"controller-manager",
}

// getOperatorOwner looks up the operator's own Deployment to use as an OwnerReference.
// Returns nil if the deployment cannot be found (best-effort).
func (r *OtelBootstrapRunnable) getOperatorOwner(ctx context.Context, log logr.Logger) *appsv1.Deployment {
for _, name := range operatorDeploymentNames {
deploy := &appsv1.Deployment{}
key := types.NamespacedName{Name: name, Namespace: r.Namespace}
if err := r.Client.Get(ctx, key, deploy); err == nil {
return deploy
}
}
log.Info("Could not find operator Deployment for OwnerReference, ConfigMaps will be unowned")
return nil
}

// setOwnerIfAvailable sets an OwnerReference on the given object if an owner is available.
func (r *OtelBootstrapRunnable) setOwnerIfAvailable(owner *appsv1.Deployment, obj client.Object, log logr.Logger) {
if owner == nil || r.Scheme == nil {
return
}
if err := controllerutil.SetOwnerReference(owner, obj, r.Scheme); err != nil {
log.Error(err, "Failed to set OwnerReference on resource", "name", obj.GetName())
}
}

// Start runs the bootstrap sequence. Called by the manager after leader election
// and cache sync, before controllers start processing events.
func (r *OtelBootstrapRunnable) Start(ctx context.Context) error {
log := r.Log.WithName("otel-bootstrap")
log.Info("Starting OTel collector bootstrap")

// Look up operator Deployment once for OwnerReference on created resources.
owner := r.getOperatorOwner(ctx, log)

isOCP, err := r.detectOpenShift(ctx)
if err != nil {
return fmt.Errorf("detecting OpenShift: %w", err)
}

if isOCP {
log.Info("OpenShift detected, reconciling ingress CA trust")
if err := r.reconcileIngressCA(ctx, log); err != nil {
if err := r.reconcileIngressCA(ctx, log, owner); err != nil {
return fmt.Errorf("ingress CA bootstrap: %w", err)
}
} else {
log.Info("Not running on OpenShift, skipping ingress CA trust")
}

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

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

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

// reconcileCollectorConfig discovers available components and assembles the
// OTel collector ConfigMap from preset configurations.
func (r *OtelBootstrapRunnable) reconcileCollectorConfig(ctx context.Context, log logr.Logger, isOCP bool) error {
func (r *OtelBootstrapRunnable) reconcileCollectorConfig(ctx context.Context, log logr.Logger, isOCP bool, owner *appsv1.Deployment) error {
mf, err := r.discoverMLflow(ctx, log)
if err != nil {
return err
Expand Down Expand Up @@ -291,6 +329,7 @@ func (r *OtelBootstrapRunnable) reconcileCollectorConfig(ctx context.Context, lo
},
Data: map[string]string{configMapDataKey: configStr},
}
r.setOwnerIfAvailable(owner, cm, log)
if err := r.Client.Create(ctx, cm); err != nil {
if !errors.IsAlreadyExists(err) {
return fmt.Errorf("creating collector ConfigMap: %w", err)
Expand Down
1 change: 1 addition & 0 deletions kagenti-operator/internal/bootstrap/otel_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -52,6 +52,7 @@ func newRunnable(cl client.Client, isOCP func(context.Context) (bool, error), ml
return &OtelBootstrapRunnable{
Client: cl,
APIReader: cl,
Scheme: testScheme(),
Namespace: testNamespace,
Log: testLogger(),
IsOpenShift: isOCP,
Expand Down
35 changes: 20 additions & 15 deletions kagenti-operator/internal/controller/agentruntime_controller.go
Original file line number Diff line number Diff line change
Expand Up @@ -154,11 +154,7 @@ func (r *AgentRuntimeReconciler) Reconcile(ctx context.Context, req ctrl.Request
// 4. Resolve targetRef (existence check)
if err := r.resolveTargetRef(ctx, rt); err != nil {
logger.Error(err, "Failed to resolve targetRef")
r.setPhase(rt, agentv1alpha1.RuntimePhaseError)
r.setCondition(rt, ConditionTypeTargetResolved, metav1.ConditionFalse, "TargetNotFound", err.Error())
if statusErr := r.Status().Update(ctx, rt); statusErr != nil {
logger.Error(statusErr, "Failed to update status")
}
r.updateErrorStatus(ctx, req.NamespacedName, ConditionTypeTargetResolved, "TargetNotFound", err.Error())
if r.Recorder != nil {
r.Recorder.Event(rt, corev1.EventTypeWarning, "TargetNotFound", err.Error())
}
Expand Down Expand Up @@ -210,11 +206,7 @@ func (r *AgentRuntimeReconciler) Reconcile(ctx context.Context, req ctrl.Request
configResult, err := ComputeConfigHash(ctx, r.Client, rt.Namespace)
if err != nil {
logger.Error(err, "Failed to compute config hash")
r.setPhase(rt, agentv1alpha1.RuntimePhaseError)
r.setCondition(rt, ConditionTypeReady, metav1.ConditionFalse, "ConfigHashError", err.Error())
if statusErr := r.Status().Update(ctx, rt); statusErr != nil {
logger.Error(statusErr, "Failed to update status")
}
r.updateErrorStatus(ctx, req.NamespacedName, ConditionTypeReady, "ConfigHashError", err.Error())
return ctrl.Result{RequeueAfter: 30 * time.Second}, nil
}

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

Expand Down Expand Up @@ -777,6 +765,23 @@ func (r *AgentRuntimeReconciler) setCondition(rt *agentv1alpha1.AgentRuntime, co
})
}

// updateErrorStatus sets the AgentRuntime phase to Error and updates a condition
// with retry-on-conflict semantics, re-fetching the object on each attempt.
func (r *AgentRuntimeReconciler) updateErrorStatus(ctx context.Context, key types.NamespacedName, condType, reason, message string) {
logger := log.FromContext(ctx)
if statusErr := retry.RetryOnConflict(retry.DefaultRetry, func() error {
latest := &agentv1alpha1.AgentRuntime{}
if err := r.Get(ctx, key, latest); err != nil {
return err
}
r.setPhase(latest, agentv1alpha1.RuntimePhaseError)
r.setCondition(latest, condType, metav1.ConditionFalse, reason, message)
return r.Status().Update(ctx, latest)
}); statusErr != nil {
logger.Error(statusErr, "Failed to update error status", "condition", condType, "reason", reason)
}
}

// fetchAndUpdateCard discovers the agent card from the workload's Service endpoint
// and populates status.card. Skips fetch when the feature flag is disabled or
// when the workload's change-detection key has not changed.
Expand Down
27 changes: 20 additions & 7 deletions kagenti-operator/internal/controller/mlflow_controller.go
Original file line number Diff line number Diff line change
Expand Up @@ -223,22 +223,35 @@ func (r *MLflowReconciler) configureDeployment(ctx context.Context, dep *appsv1.
if annotations == nil {
annotations = make(map[string]string)
}
annotations[AnnotationMLflowExperimentID] = experimentID
annotations[AnnotationMLflowExperimentName] = experimentName
annotations[AnnotationMLflowTrackingURI] = trackingURI
annotations[AnnotationMLflowTrackingAuth] = "kubernetes-namespaced"

annotationsChanged := false
for k, v := range map[string]string{
AnnotationMLflowExperimentID: experimentID,
AnnotationMLflowExperimentName: experimentName,
AnnotationMLflowTrackingURI: trackingURI,
AnnotationMLflowTrackingAuth: "kubernetes-namespaced",
} {
if annotations[k] != v {
annotations[k] = v
annotationsChanged = true
}
}
latest.Spec.Template.Annotations = annotations

changed := false
envChanged := false
for i := range latest.Spec.Template.Spec.Containers {
for name, value := range desired {
if setEnvVar(&latest.Spec.Template.Spec.Containers[i], name, value) {
changed = true
envChanged = true
}
}
}

if changed {
if !envChanged && !annotationsChanged {
return nil
}

if envChanged {
logger.Info("Injected MLflow env vars into Deployment containers", "deployment", dep.Name)
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@ import (
appsv1 "k8s.io/api/apps/v1"
corev1 "k8s.io/api/core/v1"
apierrors "k8s.io/apimachinery/pkg/api/errors"
"k8s.io/apimachinery/pkg/runtime"
"k8s.io/apimachinery/pkg/types"
"k8s.io/apimachinery/pkg/util/wait"
"k8s.io/client-go/discovery"
Expand All @@ -34,6 +35,7 @@ import (
"k8s.io/client-go/util/retry"
ctrl "sigs.k8s.io/controller-runtime"
"sigs.k8s.io/controller-runtime/pkg/client"
"sigs.k8s.io/controller-runtime/pkg/controller/controllerutil"
"sigs.k8s.io/controller-runtime/pkg/handler"
"sigs.k8s.io/controller-runtime/pkg/log"
"sigs.k8s.io/controller-runtime/pkg/reconcile"
Expand Down Expand Up @@ -100,6 +102,7 @@ var (

type SharedTrustReconciler struct {
client.Client
Scheme *runtime.Scheme
Recorder record.EventRecorder
}

Expand Down Expand Up @@ -224,6 +227,9 @@ func (r *SharedTrustReconciler) reconcileCacertsSecrets(ctx context.Context) (bo
secret.Namespace = ic.Namespace
secret.Type = corev1.SecretTypeOpaque
secret.Data = desired
if err := controllerutil.SetOwnerReference(intSecret, secret, r.Scheme); err != nil {
return false, fmt.Errorf("setting owner reference for cacerts secret in %s: %w", ic.Namespace, err)
}
if err := r.Create(ctx, secret); err != nil {
return false, fmt.Errorf("creating cacerts secret in %s: %w", ic.Namespace, err)
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -187,6 +187,7 @@ func newReconciler(t *testing.T, objs ...runtime.Object) *SharedTrustReconciler
cb := fake.NewClientBuilder().WithScheme(scheme).WithRuntimeObjects(clientObjs...)
return &SharedTrustReconciler{
Client: cb.Build(),
Scheme: scheme,
Recorder: record.NewFakeRecorder(10),
}
}
Expand Down
22 changes: 20 additions & 2 deletions kagenti-operator/internal/signature/verifier.go
Original file line number Diff line number Diff line change
Expand Up @@ -353,8 +353,26 @@ func removeEmptyFields(m map[string]interface{}) map[string]interface{} {
result[k] = cleaned
}
case []interface{}:
if len(val) > 0 {
result[k] = val
var cleaned []interface{}
for _, elem := range val {
switch e := elem.(type) {
case map[string]interface{}:
c := removeEmptyFields(e)
if len(c) > 0 {
cleaned = append(cleaned, c)
}
case string:
if e != "" {
cleaned = append(cleaned, e)
}
case nil:
// skip nil elements
default:
cleaned = append(cleaned, e)
}
}
if len(cleaned) > 0 {
result[k] = cleaned
}
case string:
if val != "" {
Expand Down
9 changes: 9 additions & 0 deletions scripts/kind-with-registry.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -19,9 +19,18 @@ spec:
labels:
app: registry
spec:
securityContext:
runAsNonRoot: true
seccompProfile:
type: RuntimeDefault
containers:
- name: registry
image: public.ecr.aws/docker/library/registry:3.0.0-rc.4
securityContext:
allowPrivilegeEscalation: false
capabilities:
drop:
- ALL
ports:
- containerPort: 5000
name: registry
Expand Down
Loading