/* Copyright 2026. Licensed under the Apache License, Version 2.0 (the "License"); you may not use this file except in compliance with the License. You may obtain a copy of the License at http://www.apache.org/licenses/LICENSE-2.0 Unless required by applicable law or agreed to in writing, software distributed under the License is distributed on an "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the License for the specific language governing permissions and limitations under the License. */ package controller import ( "context" "errors" "fmt" "time" corev1 "k8s.io/api/core/v1" apierrors "k8s.io/apimachinery/pkg/api/errors" apimeta "k8s.io/apimachinery/pkg/api/meta" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/runtime" "k8s.io/apimachinery/pkg/types" ctrl "sigs.k8s.io/controller-runtime" "sigs.k8s.io/controller-runtime/pkg/client" "sigs.k8s.io/controller-runtime/pkg/controller" "sigs.k8s.io/controller-runtime/pkg/controller/controllerutil" "sigs.k8s.io/controller-runtime/pkg/event" "sigs.k8s.io/controller-runtime/pkg/handler" logf "sigs.k8s.io/controller-runtime/pkg/log" "sigs.k8s.io/controller-runtime/pkg/reconcile" "sigs.k8s.io/controller-runtime/pkg/source" crawlv1alpha1 "gitea.home.hrajfrisbee.cz/kacerr/egress-proxies-operator/api/v1alpha1" "gitea.home.hrajfrisbee.cz/kacerr/egress-proxies-operator/internal/health" "gitea.home.hrajfrisbee.cz/kacerr/egress-proxies-operator/internal/provider" "gitea.home.hrajfrisbee.cz/kacerr/egress-proxies-operator/internal/tracing" ) // HealthSnapshotter provides the current probe verdict for a proxy. The // health engine implements it; the reconciler is its only consumer, turning // snapshots into the Healthy condition — the engine owns health state, the // reconciler owns its representation. type HealthSnapshotter interface { Snapshot(key types.NamespacedName) (health.Snapshot, bool) } // ProxyReconciler reconciles Proxy objects as a state machine: every // reconcile derives exactly one action from (spec, status, provider Get), // performs it, and requeues. Status is written at most once per reconcile, // by the deferred patch in Reconcile. type ProxyReconciler struct { client.Client Scheme *runtime.Scheme // Providers maps spec.provider values to configured backends. Providers map[string]provider.Provider // Health supplies probe verdicts; nil disables health representation // (the Healthy condition simply never appears). Health HealthSnapshotter // HealthEvents, when non-nil, is watched as a raw source so the health // engine can enqueue proxies on status-affecting transitions. HealthEvents <-chan event.GenericEvent // Poll intervals are struct fields, never consts, so tests can shrink // them to milliseconds. ProvisioningPoll time.Duration // while waiting for an instance to reach Running DriftPoll time.Duration // between re-checks of a Running instance DeletionPoll time.Duration // while waiting for an instance to disappear QuotaRetry time.Duration // after ErrQuotaExceeded; slow, off the backoff curve RequeueNow time.Duration // "process the next state promptly" (Result.Requeue is deprecated) } // +kubebuilder:rbac:groups=crawl.example.com,resources=proxies,verbs=get;list;watch;create;update;patch;delete // +kubebuilder:rbac:groups=crawl.example.com,resources=proxies/status,verbs=get;update;patch // +kubebuilder:rbac:groups=crawl.example.com,resources=proxies/finalizers,verbs=update // +kubebuilder:rbac:groups="",resources=secrets,verbs=get;list;watch // The pod verbs are for the kubernetes-pod provider; cluster-scoped, since // its ListByTag enumerates the operator's Pods across all namespaces. // +kubebuilder:rbac:groups="",resources=pods,verbs=get;list;watch;create;delete // Reconcile fetches the Proxy named by req into p (r.Get fills the struct // through the pointer), dispatches to the delete/external/managed state // machines, and flushes any status change exactly once on the way out. func (r *ProxyReconciler) Reconcile(ctx context.Context, req ctrl.Request) (res ctrl.Result, err error) { var p crawlv1alpha1.Proxy if err := r.Get(ctx, req.NamespacedName, &p); err != nil { return ctrl.Result{}, client.IgnoreNotFound(err) } base := p.DeepCopy() defer func() { // Runs inside the root reconcile span (this defer fires before the // tracing.NewReconciler wrapper sees the return), and its error is // folded into err, which that wrapper records. pctx, span := tracing.Start(ctx, "status.patch") defer span.End() // NotFound is expected when this reconcile just removed the last // finalizer and the object is already gone. if perr := r.patchStatusIfChanged(pctx, base, &p); perr != nil && !apierrors.IsNotFound(perr) { err = errors.Join(err, perr) } }() switch { case !p.DeletionTimestamp.IsZero(): return r.reconcileDelete(ctx, &p) case p.Spec.Mode == crawlv1alpha1.ModeExternal: return r.reconcileExternal(ctx, &p) default: return r.reconcileManaged(ctx, &p) } } func (r *ProxyReconciler) reconcileManaged(ctx context.Context, p *crawlv1alpha1.Proxy) (ctrl.Result, error) { ctx, span := tracing.Start(ctx, "reconcile.managed") defer span.End() log := logf.FromContext(ctx) if controllerutil.AddFinalizer(p, crawlv1alpha1.FinalizerName) { // The Update event re-triggers reconciliation; provisioning starts // on the next pass, with the finalizer safely persisted first. return ctrl.Result{}, r.Update(ctx, p) } // Permanent-failure latch: once this generation has failed permanently, // stop calling the provider until the spec changes. if cond := apimeta.FindStatusCondition(p.Status.Conditions, crawlv1alpha1.ConditionProvisioned); cond != nil && cond.Status == metav1.ConditionFalse && cond.Reason == ReasonPermanentError && cond.ObservedGeneration == p.Generation { return ctrl.Result{}, nil } prov, ok := r.Providers[p.Spec.Provider] if !ok { setProvisioned(p, metav1.ConditionFalse, ReasonPermanentError, fmt.Sprintf("provider %q is not configured", p.Spec.Provider)) return ctrl.Result{}, nil } cloudInit, err := r.resolveCloudInit(ctx, p) if err != nil { setProvisioned(p, metav1.ConditionFalse, ReasonCloudInitError, err.Error()) return ctrl.Result{}, err } hash := specHash(p, cloudInit) if p.Status.ProviderID == "" { id, err := prov.Create(ctx, provider.CreateRequest{ Name: provider.NameFromUID(p.UID), UID: string(p.UID), Namespace: p.Namespace, ProxyName: p.Name, Placement: placementFrom(p.Spec.Placement), CloudInit: cloudInit, Port: p.EffectivePort(), }) if err != nil { return r.providerFailure(p, err) } log.Info("created instance", "provider", p.Spec.Provider, "providerID", id) if err := r.setSpecHash(ctx, p, hash); err != nil { return ctrl.Result{}, err } p.Status.ProviderID = id p.Status.IP = "" setProvisioned(p, metav1.ConditionFalse, ReasonProvisioning, "instance created; waiting for it to run") // Any Healthy verdict belonged to the previous instance; the health // engine starts fresh for the new one (its state was pruned while // the proxy had no IP), and so must the status. apimeta.RemoveStatusCondition(&p.Status.Conditions, crawlv1alpha1.ConditionHealthy) p.Status.LatencyMillis = 0 p.Status.LastHealthCheckTime = nil return ctrl.Result{RequeueAfter: r.ProvisioningPoll}, nil } if ann := p.Annotations[crawlv1alpha1.AnnotationSpecHash]; ann != hash { if ann == "" { // Adopt: an instance provisioned before the hash-input struct // gained a field (or by an older operator version) keeps its // instance; replacing the whole fleet on upgrade would be wrong. if err := r.setSpecHash(ctx, p, hash); err != nil { return ctrl.Result{}, err } return ctrl.Result{RequeueAfter: r.RequeueNow}, nil } return r.replaceInstance(ctx, p, prov, hash) } inst, err := prov.Get(ctx, p.Status.ProviderID) if provider.Class(err) == provider.ErrNotFound { p.Status.ProviderID = "" p.Status.IP = "" return ctrl.Result{RequeueAfter: r.RequeueNow}, nil } if err != nil { return r.providerFailure(p, err) } switch inst.State { case provider.StateProvisioning: p.Status.IP = "" setProvisioned(p, metav1.ConditionFalse, ReasonProvisioning, "waiting for the instance to run") return ctrl.Result{RequeueAfter: r.ProvisioningPoll}, nil case provider.StateRunning: p.Status.IP = inst.IP setProvisioned(p, metav1.ConditionTrue, ReasonCreated, "instance is running") r.applyHealth(p) return ctrl.Result{RequeueAfter: r.DriftPoll}, nil default: // Stopped, Terminated: cattle, not pets — delete and recreate. if err := prov.Delete(ctx, p.Status.ProviderID); err != nil { return r.providerFailure(p, err) } log.Info("deleting instance for recreation", "providerID", p.Status.ProviderID, "state", inst.State) p.Status.IP = "" setProvisioned(p, metav1.ConditionFalse, ReasonRecreating, fmt.Sprintf("instance is %s; deleting it for recreation", inst.State)) return ctrl.Result{RequeueAfter: r.DeletionPoll}, nil } } // replaceInstance handles a spec-hash mismatch. The replacement instance has // the same deterministic name as the old one (both derive from the CR UID), // so recreating before the old instance is fully gone would hit "already // exists" — hence: delete, poll to NotFound, only then advance the hash and // let the create branch run. func (r *ProxyReconciler) replaceInstance(ctx context.Context, p *crawlv1alpha1.Proxy, prov provider.Provider, hash string) (ctrl.Result, error) { ctx, span := tracing.Start(ctx, "reconcile.replaceInstance") defer span.End() _, err := prov.Get(ctx, p.Status.ProviderID) if provider.Class(err) == provider.ErrNotFound { // Old instance is gone. The Update inside setSpecHash refreshes p // from the server — including status — so the status clear must be // staged after it, or it would be silently overwritten. A crash // between the two writes recovers either way: the create branch's // Create is idempotent by name, and a stale ID resolves to NotFound // again. if err := r.setSpecHash(ctx, p, hash); err != nil { return ctrl.Result{}, err } p.Status.ProviderID = "" p.Status.IP = "" return ctrl.Result{RequeueAfter: r.RequeueNow}, nil } if err != nil { return r.providerFailure(p, err) } if err := prov.Delete(ctx, p.Status.ProviderID); err != nil { return r.providerFailure(p, err) } logf.FromContext(ctx).Info("replacing instance after spec change", "providerID", p.Status.ProviderID) p.Status.IP = "" setProvisioned(p, metav1.ConditionFalse, ReasonReplacing, "spec changed; deleting the old instance before recreating") return ctrl.Result{RequeueAfter: r.DeletionPoll}, nil } func (r *ProxyReconciler) reconcileDelete(ctx context.Context, p *crawlv1alpha1.Proxy) (ctrl.Result, error) { ctx, span := tracing.Start(ctx, "reconcile.delete") defer span.End() if !controllerutil.ContainsFinalizer(p, crawlv1alpha1.FinalizerName) { return ctrl.Result{}, nil } if p.Status.ProviderID == "" { // Nothing was ever recorded as created; orphan GC reaps any stray // instance a crashed create might have left behind. controllerutil.RemoveFinalizer(p, crawlv1alpha1.FinalizerName) return ctrl.Result{}, r.Update(ctx, p) } prov, ok := r.Providers[p.Spec.Provider] if !ok { return ctrl.Result{}, fmt.Errorf( "provider %q is not configured; cannot clean up instance %s", p.Spec.Provider, p.Status.ProviderID) } _, err := prov.Get(ctx, p.Status.ProviderID) if provider.Class(err) == provider.ErrNotFound { controllerutil.RemoveFinalizer(p, crawlv1alpha1.FinalizerName) return ctrl.Result{}, r.Update(ctx, p) } if err != nil { return r.deletionFailure(err) } if err := prov.Delete(ctx, p.Status.ProviderID); err != nil { return r.deletionFailure(err) } logf.FromContext(ctx).Info("deleting instance", "providerID", p.Status.ProviderID) setProvisioned(p, metav1.ConditionFalse, ReasonDeleting, "deleting the instance before removing the finalizer") return ctrl.Result{RequeueAfter: r.DeletionPoll}, nil } func (r *ProxyReconciler) reconcileExternal(_ context.Context, p *crawlv1alpha1.Proxy) (ctrl.Result, error) { if p.Spec.Endpoint == nil { // CEL guarantees an endpoint on any object that went through the API // server; tolerate its absence instead of panicking. setProvisioned(p, metav1.ConditionFalse, ReasonPermanentError, "external proxy has no endpoint") return ctrl.Result{}, nil } p.Status.IP = p.Spec.Endpoint.Host setProvisioned(p, metav1.ConditionTrue, ReasonExternalEndpoint, "tracking an external endpoint") r.applyHealth(p) return ctrl.Result{}, nil } // providerFailure translates a classified provider error into the // state-machine's reaction: transient errors ride the workqueue's // exponential backoff, quota errors back off slowly without counting as // errors, and permanent errors latch Failed and stop retrying. func (r *ProxyReconciler) providerFailure(p *crawlv1alpha1.Proxy, err error) (ctrl.Result, error) { switch provider.Class(err) { case provider.ErrQuotaExceeded: setProvisioned(p, metav1.ConditionFalse, ReasonQuotaExceeded, err.Error()) return ctrl.Result{RequeueAfter: r.QuotaRetry}, nil case provider.ErrPermanent: setProvisioned(p, metav1.ConditionFalse, ReasonPermanentError, err.Error()) return ctrl.Result{}, nil default: return ctrl.Result{}, err } } // deletionFailure is providerFailure for the finalizer path, where latching // a permanent failure would wedge the object forever with no retry — keep // retrying instead, visibly, until cleanup succeeds or an operator // intervenes. func (r *ProxyReconciler) deletionFailure(err error) (ctrl.Result, error) { if provider.Class(err) == provider.ErrQuotaExceeded { return ctrl.Result{RequeueAfter: r.QuotaRetry}, nil } return ctrl.Result{}, err } // resolveCloudInit returns the effective cloud-init user-data, reading the // referenced Secret if one is used. Both the spec hash and CreateRequest see // only resolved content, so rotating a Secret triggers replacement. func (r *ProxyReconciler) resolveCloudInit(ctx context.Context, p *crawlv1alpha1.Proxy) (string, error) { ci := p.Spec.CloudInit if ci == nil { return "", nil } if ci.Inline != "" { return ci.Inline, nil } if ci.SecretRef == nil { return "", nil } key := ci.SecretRef.Key if key == "" { key = crawlv1alpha1.DefaultCloudInitSecretKey } var sec corev1.Secret if err := r.Get(ctx, client.ObjectKey{Namespace: p.Namespace, Name: ci.SecretRef.Name}, &sec); err != nil { return "", fmt.Errorf("resolving cloudInit secret %q: %w", ci.SecretRef.Name, err) } data, ok := sec.Data[key] if !ok { return "", fmt.Errorf("cloudInit secret %q has no key %q", ci.SecretRef.Name, key) } return string(data), nil } // setSpecHash persists the spec-hash annotation. Status changes staged on p // are untouched by the Update (they live on the status subresource) and are // flushed by the deferred patch in Reconcile. func (r *ProxyReconciler) setSpecHash(ctx context.Context, p *crawlv1alpha1.Proxy, hash string) error { if p.Annotations[crawlv1alpha1.AnnotationSpecHash] == hash { return nil } if p.Annotations == nil { p.Annotations = map[string]string{} } p.Annotations[crawlv1alpha1.AnnotationSpecHash] = hash return r.Update(ctx, p) } func placementFrom(ps *crawlv1alpha1.PlacementSpec) provider.Placement { if ps == nil { return provider.Placement{} } return provider.Placement{ Region: ps.Region, Zone: ps.Zone, MachineType: ps.MachineType, Image: ps.Image, } } // proxiesForSecret maps a Secret event to the Proxies whose cloudInit // references it, so rotating a Secret re-triggers the replacement check. func (r *ProxyReconciler) proxiesForSecret(ctx context.Context, obj client.Object) []reconcile.Request { var list crawlv1alpha1.ProxyList if err := r.List(ctx, &list, client.InNamespace(obj.GetNamespace())); err != nil { logf.FromContext(ctx).Error(err, "listing proxies for secret event", "secret", obj.GetName()) return nil } var reqs []reconcile.Request for i := range list.Items { p := &list.Items[i] if ci := p.Spec.CloudInit; ci != nil && ci.SecretRef != nil && ci.SecretRef.Name == obj.GetName() { reqs = append(reqs, reconcile.Request{NamespacedName: client.ObjectKeyFromObject(p)}) } } return reqs } // SetupWithManager sets up the controller with the Manager. The Secret watch // only fires for Secrets the manager's cache holds; the composition root // (cmd/main.go) restricts that cache to labelled cloud-init Secrets. func (r *ProxyReconciler) SetupWithManager(mgr ctrl.Manager) error { r.applyDefaults() b := ctrl.NewControllerManagedBy(mgr). For(&crawlv1alpha1.Proxy{}). Named("proxy"). Watches(&corev1.Secret{}, handler.EnqueueRequestsFromMapFunc(r.proxiesForSecret)). WithOptions(controller.Options{MaxConcurrentReconciles: 3}) if r.HealthEvents != nil { b = b.WatchesRawSource(source.Channel(r.HealthEvents, &handler.EnqueueRequestForObject{})) } // Root span per reconcile; sub-reconcilers and the status flush hang // their spans off it. Tests calling r.Reconcile directly bypass the // wrapper and see no-op spans — the global tracer is never set there. return b.Complete(tracing.NewReconciler("Proxy", r)) } func (r *ProxyReconciler) applyDefaults() { if r.ProvisioningPoll == 0 { r.ProvisioningPoll = 10 * time.Second } if r.DriftPoll == 0 { r.DriftPoll = 2 * time.Minute } if r.DeletionPoll == 0 { r.DeletionPoll = 10 * time.Second } if r.QuotaRetry == 0 { r.QuotaRetry = 5 * time.Minute } if r.RequeueNow == 0 { r.RequeueNow = time.Second } }