Files
egress-proxies-operator/docs/architecture.md
Jan Novak aeb4115c72 Add tracing manifests and docs; clean up branch lint findings
Manager env block (downward-API resource attrs, commented OTLP
examples), architecture §10 + Decisions entries, README section.
Lint: goconst constants, gofmt, logcheck (Setup now takes its logger
from ctx via logf.FromContext).

Co-Authored-By: Claude <noreply@anthropic.com>
2026-08-24 11:30:27 +02:00

28 KiB

Architecture

Components

Component Package Runs as Leader-elected Role
Proxy CRD + helpers api/v1alpha1 types Proxy spec/status, CEL validation, defaulting, pure helpers
Reconciler internal/controller controller yes (with the manager) the state machine: provision, replace, delete, represent health
Provider contract internal/provider library Provider interface, error taxonomy, deterministic naming, config, metrics decorator
kubernetes provider internal/provider/kubernetes library real Squid pods in this cluster (local dev/CI)
gcp provider internal/provider/gcp library Compute Engine VMs, four API calls, fire-and-forget ops
Health engine internal/health Runnable yes through-the-proxy probes, thresholds, transition events
Lease store internal/lease Runnable (expiry sweep) no in-memory leases + cooldowns, single mutex
Discovery API internal/discovery Runnable no HTTP list/lease/release/report on :8090
Orphan GC internal/gc Runnable yes deletes tagged instances whose CR is gone
Metrics internal/metrics library explicit registration, scrape-time collectors
Composition root cmd/main.go binary flags, provider registry, wires everything onto one manager

Event flow: cluster events → reconciler functions

Which functions run in response to which Kubernetes cluster events, as wired in internal/controller/proxy_controller.go.

1. How cluster events reach the reconciler

KUBERNETES CLUSTER EVENTS                                (wiring: SetupWithManager,
─────────────────────────                                 proxy_controller.go)

 Proxy CR created / spec edited /     Secret created /           health transition
 status patched / delete requested    updated / deleted          (engine, see §6)
              │                              │                          │
   watch: For(&crawlv1alpha1.Proxy{})   watch: Watches(          WatchesRawSource(
              │                         &corev1.Secret{}, ...)   source.Channel(
              │                              │                    r.HealthEvents, ...))
              │                r.proxiesForSecret(ctx, secret)          │
              │                  │ r.List(Proxies in namespace)         │
              │                  │ keeps those whose                    │
              │                  │ spec.cloudInit.secretRef matches     │
              │                  ▼                                      │
              │                [reconcile.Request per matching Proxy]   │
              ▼                              │                          │
        ┌────────────────────────────────────┴──────────────────────────┴──┐
        │                controller-runtime workqueue                      │◄── RequeueAfter
        │  (dedup by namespace/name, rate-limited,                         │    timers
        │   MaxConcurrentReconciles: 3)                                    │◄── error backoff
        └──────────────────────────┬────────────────────────────────────────┘
                                   ▼
                    ProxyReconciler.Reconcile(ctx, req)

The Secret watch makes rotating a cloud-init Secret a first-class event: it re-enqueues every Proxy referencing that Secret, which is how secret rotation triggers VM replacement even though the Proxy spec is untouched.

2. Inside Reconcile — dispatch and the single status write

Reconcile(ctx, req)
  │ r.Get(ctx, req.NamespacedName, &p)          ── fetch the Proxy (NotFound → done)
  │ base := p.DeepCopy()                        ── snapshot for the diff
  │ defer patchStatusIfChanged(ctx, base, &p) ──────────────────────────┐
  │                                                                     │
  ├─ p.DeletionTimestamp set ──► reconcileDelete(ctx, &p)               │
  ├─ p.Spec.Mode == External ──► reconcileExternal(ctx, &p)             │
  └─ otherwise (Managed) ──────► reconcileManaged(ctx, &p)              │
                                                                        ▼
                                     patchStatusIfChanged  (status.go)
                                       │ p.Status.ObservedGeneration = p.Generation
                                       │ p.Status.Phase = computePhase(&p)
                                       │ equality.Semantic.DeepEqual(base, p)?
                                       └─ changed → r.Status().Patch(...)   ◄── the ONLY
                                          unchanged → no API call               status write

3. reconcileManaged — the state machine

reconcileManaged(ctx, p)
  │
  ├─ controllerutil.AddFinalizer? ──► r.Update ──► return {}        (watch event re-triggers)
  ├─ permanent-failure latch (FindStatusCondition == PermanentError
  │    at this generation) ──► return {}                            (silent until spec edit)
  ├─ r.Providers[p.Spec.Provider] missing ──► setProvisioned(PermanentError) → Failed
  │
  ├─ resolveCloudInit(ctx, p) ──► r.Get(Secret) if secretRef        (error → CloudInitError + backoff)
  ├─ hash := specHash(p, cloudInit)                                  (spechash.go)
  │
  ├─ status.providerID == "" ─────────► prov.Create(CreateRequest{Name: NameFromUID(p.UID), ...})
  │                                       │ setSpecHash → r.Update (annotation)
  │                                       └ stage providerID + Provisioned=False/Provisioning
  │                                         ──► RequeueAfter: ProvisioningPoll
  │
  ├─ annotation != hash, annotation == "" ──► adopt: setSpecHash → r.Update
  │                                           ──► RequeueAfter: RequeueNow
  ├─ annotation != hash, annotation != "" ──► replaceInstance:
  │                                             prov.Get ─ NotFound → setSpecHash, clear ID/IP
  │                                             │            ──► RequeueNow (next pass creates)
  │                                             └ exists  → prov.Delete, Provisioned=False/Replacing
  │                                                          ──► RequeueAfter: DeletionPoll
  │
  └─ annotation == hash ──► prov.Get(providerID)
        ├─ ErrNotFound      ──► clear ID/IP            ──► RequeueNow (next pass creates)
        ├─ Provisioning     ──► Provisioned=False      ──► ProvisioningPoll
        ├─ Running          ──► status.ip = inst.IP,
        │                       Provisioned=True/Created,
        │                       applyHealth (see §6)   ──► DriftPoll
        └─ Stopped/Termin.  ──► prov.Delete (cattle)   ──► DeletionPoll

  any provider error ──► providerFailure(p, err) ── provider.Class(err):
        ├─ ErrQuotaExceeded ──► condition QuotaExceeded ──► RequeueAfter: QuotaRetry (nil error)
        ├─ ErrPermanent     ──► condition PermanentError ──► phase Failed, no retry
        └─ ErrTransient     ──► return err ──► workqueue exponential backoff

4. reconcileDelete and reconcileExternal

reconcileDelete(ctx, p)                          reconcileExternal(ctx, p)
  ├─ no finalizer ──► return {}                    │ status.ip = spec.endpoint.host
  ├─ providerID == "" ──► RemoveFinalizer          │ setProvisioned(True/ExternalEndpoint)
  │     → r.Update → object actually deleted       │ applyHealth (see §6)
  ├─ prov.Get → ErrNotFound ──► RemoveFinalizer    └─ return {}   (no finalizer,
  │     → r.Update → object actually deleted                       no provider calls ever)
  └─ exists ──► prov.Delete
        → Provisioned=False/Deleting
        ──► RequeueAfter: DeletionPoll (poll until gone)

5. What provider calls do in the outside world

kubernetes pod provider (internal/provider/kubernetes/)
prov.Create ──► buildPod (pure) ──► client.Create(corev1.Pod)   ─┐ these cause Pod events,
prov.Get    ──► client.Get(Pod) → phase/IP → InstanceState       │ but the operator does NOT
prov.Delete ──► client.Delete(Pod, tolerate NotFound)            │ watch Pods — it observes
prov.ListByTag ─► client.List(Pods by labels, all namespaces)   ─┘ them by polling prov.Get
                                                                    on each RequeueAfter tick

gcp provider (internal/provider/gcp/) — instances.{Insert,Get,Delete,AggregatedList}, nothing else
prov.Create ──► buildInsertRequest (pure) ──► instances.Insert   ─┐ fire-and-forget:
                  409 alreadyExists = success (idempotent retry)  │ Operation.Wait is never
prov.Get    ──► instances.Get → status/NatIP → InstanceState      │ called; readiness is
                  RUNNING without NatIP = still Provisioning      │ discovered by Get polls,
prov.Delete ──► instances.Delete (404 = success)                  │ exactly like the pod
prov.ListByTag ─► AggregatedList(label filter,                   ─┘ provider
                    ReturnPartialSuccess: true)
providerID = zones/<zone>/instances/<name> — zone-qualified, so Get/Delete
stay correct even mid-replacement after a zone edit

The reconciler never watches provider-side resources (Pods or GCP VMs). All instance-state observation is poll-based through the Provider interface, so the same flow works identically for a cloud API that has no watch mechanism at all.

6. Health engine (internal/health/) — probes and transitions

The engine is a leader-elected manager Runnable with its own goroutines, independent of the workqueue. It owns health state; the reconciler owns its representation in status — that split keeps exactly one writer of .status and makes write-only-on-transition fall out for free.

Engine.Start(ctx)                                (engine.go)
  ├─ spawns Workers (8) probe goroutines ◄─┐
  └─ ticker loop (Tick = 1s):              │ jobs channel (non-blocking send;
       tick(ctx, now, jobs)                │ saturated pool → retry next tick)
         │ Reader.List(Proxies)  ── from the manager cache
         │ per proxy: skip if no IP/host or deleting (state pruned →
         │   a replaced instance starts with fresh counters)
         │ newState: seed verdict from an existing Healthy condition
         │   (leader handover), jitter first probe across the interval
         │ due && !inFlight ──► jobs ◄── probe worker picks up
         └ prune states for proxies gone from the cache
                                           │
              probe(ctx, proxyURL, hc, tls)   (probe.go)
                │ fresh transport per probe, DisableKeepAlives=true
                │   (load-bearing: keep-alives would cache the CONNECT
                │    tunnel and later probes would never re-exercise it)
                │ https probe URL ⇒ CONNECT through the proxy + TLS inside
                └ success = err == nil AND expected status code
                                           │
              record(job, result, now)     ── under one mutex
                │ counters: consecOK/consecFail; verdict flips only at
                │   successThreshold / failureThreshold
                │ emit ONLY on: first-ever verdict │ threshold flip │
                │   latency Δ > max(20ms, 50% of reported) rate-limited
                │   to one report per MinReportInterval (60s)
                ▼
        Events chan (buffered 64, non-blocking send;
        on drop the reported markers do NOT advance → next probe retries)
                │
                ▼
        source.Channel → workqueue → Reconcile        (see §1)
                │
                ▼
        r.applyHealth(p)  ── reads Engine.Snapshot(key)   (status.go)
              stages the Healthy condition + latencyMillis +
              lastHealthCheckTime; computePhase turns Provisioned=True
              + Healthy=True/False into phase Ready / Unhealthy

Consequence worth knowing: status.lastHealthCheckTime is the time of the last status-affecting probe, not the most recent probe — suppressed probes deliberately never write status. True probe recency will live in metrics (Step 9).

7. Discovery + lease API (internal/discovery/, internal/lease/)

HTTP-driven, not cluster-event-driven: crawler clients call in; the only Kubernetes interaction is reading Proxies from the manager's cache. The server is a non-leader-elected Runnable (all replicas would serve, but the deployment ships replicas: 1 because lease state is per-process — an operator restart drops all leases and cooldowns, a documented caveat). Client-facing reference with request/response schemas and curl examples: api.md.

crawler client
  │  Authorization: Bearer $DISCOVERY_TOKEN   (empty token = auth disabled, loud startup warning)
  ▼
Server.handler()      middleware, outermost first        (server.go)
  recover → request-log → MaxBytesReader(64KiB) → bearer auth (constant-time; /healthz exempt)
  │
  ├─ GET  /healthz            ──► 200 ok (unauthenticated)
  │
  ├─ GET  /v1/proxies?attr.k=v&healthy=true            (handlers.go)
  │     Reader.List(Proxies) ── manager cache
  │     filter: attributes equality + Healthy condition
  │     + Store.Counts() for activeLeases
  │     ──► 200 {"proxies":[...], "count":N}   (empty list is 200, not 404)
  │
  ├─ POST /v1/leases  {"selector":{...},"ttlSeconds":300,"target":"..."}
  │     Reader.List → filter selector; unhealthy matches counted, not offered
  │     Store.Acquire(healthy candidates, target, ttl)   ── one lock: select+insert
  │     │    selection: fewest active leases, then latency, then name
  │     ├─ granted ──► 201 {leaseID, proxy:{...}, expiresAt, ttlSeconds}
  │     └─ ErrNoMatch ──► 409 {"error":"no_match", considered, atCapacity,
  │                            inCooldown, unhealthy}
  │
  ├─ DELETE /v1/leases/{id} ──► Store.Release ──► always 204 (idempotent)
  │
  └─ POST /v1/leases/{id}/report {"result":"ok|rate_limited|banned","target":"..."}
        Store.Report ── rate_limited/banned ⇒ cooldown[{proxy,target}] for
        │               CooldownWindow (target falls back: report → lease → global)
        ├─ 204 │ 400 invalid_result │ 404 unknown_lease
        └─ an expired lease still resolves for CooldownWindow past its TTL —
           a late report lands exactly when the proxy is being rate-limited

Store.Start(ctx)  ── manager Runnable, NOT leader-elected: sweeps expired
                     leases + cooldowns; correctness never depends on the
                     sweep (every read checks ExpiresAt against the clock)

8. Orphan GC (internal/gc/) — the crash-safety net

Timer-driven, leader-elected (destructive ⇒ single writer). Exists for the one gap the reconciler cannot close alone: a crash after a provider Create but before the status write that records the instance.

Sweeper.Start(ctx)   ── refuses to run when the cache is namespace-
  │                     restricted unless --gc-allow-namespaced is explicit
  │                     (an incomplete live set would "orphan" live VMs)
  └─ every Interval (10m; first sweep a full interval after start):
       sweep(ctx)
         │ Reader.List(Proxies) → live UID set
         │   List fails → skip the whole sweep (never guess)
         │   a CR with deletionTimestamp still counts as LIVE — its
         │   finalizer owns that deletion; GC racing it double-deletes
         └ per provider: ListByTag
             │ error → log, continue with the next provider
             └ delete only when ALL hold:
                 has the proxy-operator-uid label   (ownership proof)
                 older than MinAge (10m)            (not mid-create)
                 UID matches no existing CR         (truly orphaned)
               each kill logged loudly with provider, providerID, UID

9. Metrics (internal/metrics/)

Registered explicitly from cmd/main.go (no init(); tests use fresh registries). Two kinds:

  • Scrape-time collectorsproxy_operator_proxies{phase} and proxy_operator_leases_active read the cache / lease store at every scrape; reconcile-incremented gauges inevitably drift and leak series.
  • Fed vectorshealthcheck_duration_seconds{proxy} and healthcheck_failures_total{proxy} observe EVERY probe (status writes are transition-only; metrics carry the high-frequency signal), and the health engine deletes a proxy's series when it prunes its state; lease_requests_total{outcome} from the discovery handlers; provider_requests_total{provider,op,result} from the provider.WithMetrics decorator — the one place Class() is called purely for observability.

Each consuming package defines its own small recorder interface (health.ProbeMetrics, discovery.LeaseMetrics, provider.RequestRecorder); metrics.Metrics satisfies all of them structurally, so no package other than cmd/main.go imports the metrics package.

10. Tracing (internal/tracing/)

OpenTelemetry tracing, integrated with — not replacing — the logr/zap logging. Everything hangs off standard OTEL_* env vars: with no OTEL_EXPORTER_OTLP_ENDPOINT / OTEL_TRACES_EXPORTER set, no SDK is installed, spans are no-ops, and logs are byte-for-byte what they were. OTEL_SDK_DISABLED=true and OTEL_TRACES_EXPORTER=none force it off; console prints spans to stdout for local dev; sampling follows OTEL_TRACES_SAMPLER(_ARG).

Span topology — each unit of work is a new trace (watch events carry no incoming trace context):

  • Reconcile Proxy (root, via a reconcile.Reconciler decorator) → reconcile.managed / reconcile.replaceInstance / reconcile.delete, plus status.patch from the deferred flush. reconcileExternal does no I/O and is unspanned.
  • provider.create|get|delete|list (client spans, provider.WithTracing decorator wired outermost around WithMetrics); the GCP SDK's own otelhttp transport contributes HTTP child spans automatically.
  • Kubernetes API calls become client child spans via a wrapped rest.Config transport — gated on an existing parent span, so informer list/watch long-polls and leader-election renewals never create root spans. The kube-apiserver ignores incoming traceparent by design; these spans are leaves.
  • Discovery API: one server span per request (named from the route pattern), incoming W3C traceparent honored so clients' traces continue into the operator; /healthz excluded.
  • gc.sweep per GC pass; health.probe per probe only with --trace-health-probes (default off — probes run ~1/s per proxy).

Log correlation: logr sinks never see a context, so trace IDs ride on the logger — tracing.Start re-derives the ctx logger from a captured base with traceID/spanID values (lowerCamel, matching reconcileID; Grafana/Loki derived-field regexes must match traceID, not trace_id). Re-deriving from the base rather than layering keeps zap from emitting duplicate keys on nested spans. GCP V(5) wire logs get the same keys from the slog handler's ctx.

Decisions

Judgment calls the spec left open, and deliberate deviations — recorded so they read as choices, not accidents. Chronological by build step.

  • Registry takes its constructor map as a parameter instead of holding a package-level map: avoids the provider⇄registry import cycle and puts the wiring at the composition root, where it is visible.
  • Mock provider replaced by the kubernetes-pod provider (user decision, mid-build): a simulated in-memory provider was too far from the real system to build confidence in. Local dev/CI now runs real ubuntu/squid pods (Canonical's actively maintained image, verified 50M+ pulls, pinned tag) in the operator's own cluster. Trade-off accepted: envtest has no kubelet, so end-to-end proof lives in the kind quickstart, and cluster pods share one egress IP — distinct egress paths remain the GCP provider's job.
  • RequeueAfter: RequeueNow instead of the plan's Requeue: true: ctrl.Result.Requeue is deprecated in controller-runtime v0.24; a fifth configurable interval (default 1s) keeps identical semantics and stays shrinkable in tests.
  • Quota exhaustion is a wait, not a failure: ErrQuotaExceeded sets a condition and requeues slowly (5m) with a nil error — off the backoff curve, out of the error log, and never phase: Failed. Only ErrPermanent latches Failed, keyed to the generation so a spec edit auto-recovers.
  • The finalizer path never latches permanent failures: a permanent error during deletion keeps retrying visibly instead — latching there would wedge the object forever with no path out but manual finalizer surgery.
  • Health transitions travel reconciler-ward over a channel (source.Channel), not direct status patches: phase derives from both provisioning and health, so two status writers would race and flap. One writer of status; the engine owns health state, the reconciler its representation; write-only-on-transition falls out for free.
  • Health state seeds from the existing Healthy condition on leader handover (verdict kept, counters zeroed, first probe jittered), so a healthy fleet doesn't flap to Unknown on restart — but a real transition still needs a full threshold run. A never-probed proxy skips the jitter and probes on the next tick: startup spread matters for restarts, not for a single new proxy.
  • Latency suppression is max(20ms, 50%) + a 60s rate limit, and only while the verdict is healthy. The spec's bare ">50% change" is undefined at 0 and lets a proxy jittering 40↔61ms write status forever; the healthy-only guard (found by test) stops a below-threshold success streak from emitting latency updates for a proxy still reported unhealthy. Consequence: status.lastHealthCheckTime means "last status-affecting probe" — true probe recency is in the metrics.
  • Deterministic instance names are proxy- + 16 chars of base32(SHA-256(CR UID)): legal for both GCP ([a-z2-7][-a-z0-9], 22 ≤ 63 chars) and Pod names, 80 bits against birthday collisions at a fleet of tens. The replacement VM therefore has the same name as the one being deleted — which is why replacement polls to NotFound before recreating instead of racing a 409.
  • banned and rate_limited share one cooldown window: a second duration knob the spec doesn't ask for; the report's semantic difference is preserved in the API but not the store.
  • Report targets fall back report → lease → global, so a client that leased with a target can't accidentally poison the proxy's global pool by omitting the target in its report.
  • The 409 body's considered counts unhealthy matches too (the store only ever sees healthy candidates): considered = atCapacity + inCooldown + unhealthy + eligible-but-outranked, keeping the numbers additive for a human debugging "why no proxy?".
  • TTLs above --max-lease-ttl are a 400, not a silent clamp — a client asking for a week should find out.
  • Discovery is not leader-elected and ships replicas: 1: caches start before non-leader-election runnables (verified in controller-runtime's ordering), and a leader-elected server would leave non-leader replicas as broken Service endpoints. One replica because lease state is per-process.
  • GCP Create requires zone, machineType, and image and fails ErrPermanent naming the missing field — inventing machine-type defaults would silently create billable VMs of arbitrary shape.
  • Unknown GCP instance statuses map to Stopped: the reconciler's answer to Stopped is delete-and-recreate, the always-safe move for cattle when the API grows a new state.
  • Kubernetes 403s classify as ErrPermanent even though quota exhaustion also surfaces as 403 (indistinguishable from RBAC denial in apierrors): not hammering an API server that may never allow the request is the safer default; a real ResourceQuota 403 forgoes the gentler quota backoff. Documented at the classification site.
  • GC kills log at Info with a WARNING: prefix — logr has no Warn level; the plan's "log at Warn" is met in spirit with provider, providerID, and UID always attached. Same convention as the discovery server's empty-token warning.
  • GC trusts only provable orphans: instances without the UID label are never deleted, a CR with a deletionTimestamp still counts as live (its finalizer owns that deletion), and an unreadable Proxy list skips the whole sweep. The namespace guard refuses to sweep a namespace-restricted cache without --gc-allow-namespaced.
  • Cloud-init Secrets must carry crawl.example.com/cloud-init: "true": the manager caches only labelled Secrets (the operator holds cluster-wide Secret read RBAC — an unrestricted cache would hold every Secret in scope). Unlabelled referenced Secrets are invisible by construction, surfacing as CloudInitError.
  • Events RBAC from the plan is omitted: nothing wires an EventRecorder in the prototype, and granting verbs nothing uses would be RBAC lint noise. Add the marker together with the recorder if events land later.
  • logr, not slog, inside controller paths: the repo convention says slog, but log.FromContext(ctx) hands controller-runtime's logr logger to everything running under the manager — fighting that would mean two logging systems in one process. Noted as a deviation rather than silently ignored.
  • Tracing is env-gated, not flag-gated: it activates only when OTEL_EXPORTER_OTLP_ENDPOINT / OTEL_TRACES_EXPORTER is set (user decision). No collector configured means no SDK installed, no-op spans, unchanged logs — the safe default for every existing deployment. The one flag is --trace-health-probes, off by default, because probes at ~1/s per proxy would dominate trace volume.
  • Trace keys are traceID/spanID, lowerCamel like reconcileID and providerID, deliberately not the trace_id many Grafana derived-field examples assume — configure the derived-field regex accordingly. logr sinks can't read ctx, so the IDs ride on the ctx logger, re-derived from a captured base per span so zap never emits duplicate keys.
  • Probe→reconcile trace links are not attempted: the health engine's GenericEvent carries only namespace/name (no ctx), and the workqueue coalesces events, so any link would be a guess. A health-triggered reconcile starts a fresh trace; the probe that caused it is findable via its own (opt-in) span and shared proxy attributes.
  • No traceparent toward probe targets: probe transports stay uninstrumented so trace headers can never leak through a proxy to external sites. The kube-apiserver ignores incoming traceparent by design (public endpoint), so k8s client spans are leaves — in-process traces, not cross-process ones.