diff --git a/.claude/settings.json b/.claude/settings.json index 229e0f5..1d6c96a 100644 --- a/.claude/settings.json +++ b/.claude/settings.json @@ -59,14 +59,19 @@ "Bash(git restore *)", "Bash(make manifests *)", "Bash(make test *)", - "Bash(KUBEBUILDER_ASSETS=\"/Users/jan.novak/srv/go/egress-proxies-operator/bin/k8s/1.36.2-darwin-arm64\" go test -race ./internal/controller/)" + "Bash(KUBEBUILDER_ASSETS=\"/Users/jan.novak/srv/go/egress-proxies-operator/bin/k8s/1.36.2-darwin-arm64\" go test -race ./internal/controller/)", + "Bash(grep -n 'func Channel' -A8 __CMDSUB_OUTPUT__/sigs.k8s.io/controller-runtime@v0.24.1/pkg/source/source.go)", + "Bash(grep -n 'type GenericEvent' __CMDSUB_OUTPUT__/sigs.k8s.io/controller-runtime@v0.24.1/pkg/event/event.go)", + "Bash(KUBEBUILDER_ASSETS=__TRACKED_VAR__/bin/k8s/1.36.2-darwin-arm64 go test -race ./...)", + "Bash(cat >> *)" ], "additionalDirectories": [ "/Users/jan.novak/srv/go/egress-proxies-operator/.claude", "/Users/jan.novak/srv/go/egress-proxies-operator/docs/plans", "/Users/jan.novak/srv/go/egress-proxies-operator/docs", "/Users/jan.novak/srv/go/egress-proxies-operator/docs/prompts", - "/tmp" + "/tmp", + "/Users/jan.novak/srv/go/egress-proxies-operator/docs/plans-executions" ] } } diff --git a/docs/architecture.md b/docs/architecture.md index 6145d54..6159760 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -1,11 +1,10 @@ # Architecture -> **Status:** the operator is built through Step 4 (reconciler) of +> **Status:** the operator is built through Step 5 (health engine) of > [docs/plans/2026-08-07-1747-proxy-operator.md](plans/2026-08-07-1747-proxy-operator.md). > This document currently covers the event/reconcile flow; the components > table and the Decisions section arrive with Step 10, and the diagrams -> below grow as the health engine, lease store, discovery API, and orphan -> GC land. +> below grow as the lease store, discovery API, and orphan GC land. ## Event flow: cluster events → reconciler functions @@ -18,23 +17,24 @@ wired in `internal/controller/proxy_controller.go`. KUBERNETES CLUSTER EVENTS (wiring: SetupWithManager, ───────────────────────── proxy_controller.go) - Proxy CR created / spec edited / Secret created / updated / deleted - status patched / delete requested │ - │ │ - watch: For(&crawlv1alpha1.Proxy{}) watch: Watches(&corev1.Secret{}, ...) - │ │ - │ r.proxiesForSecret(ctx, secret) - │ │ r.List(Proxies in secret's namespace) - │ │ keeps those whose - │ │ spec.cloudInit.secretRef.name matches - │ ▼ - │ [reconcile.Request per matching Proxy] - ▼ │ - ┌─────────────────────────────────────────────────┴──┐ - │ controller-runtime workqueue │◄── RequeueAfter timers - │ (dedup by namespace/name, rate-limited, │ (from prior reconciles) - │ MaxConcurrentReconciles: 3) │◄── error backoff retries - └──────────────────────────┬──────────────────────────┘ + 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) ``` @@ -93,7 +93,8 @@ reconcileManaged(ctx, p) ├─ ErrNotFound ──► clear ID/IP ──► RequeueNow (next pass creates) ├─ Provisioning ──► Provisioned=False ──► ProvisioningPoll ├─ Running ──► status.ip = inst.IP, - │ Provisioned=True/Created ──► DriftPoll + │ Provisioned=True/Created, + │ applyHealth (see §6) ──► DriftPoll └─ Stopped/Termin. ──► prov.Delete (cattle) ──► DeletionPoll any provider error ──► providerFailure(p, err) ── provider.Class(err): @@ -108,10 +109,10 @@ reconcileManaged(ctx, p) reconcileDelete(ctx, p) reconcileExternal(ctx, p) ├─ no finalizer ──► return {} │ status.ip = spec.endpoint.host ├─ providerID == "" ──► RemoveFinalizer │ setProvisioned(True/ExternalEndpoint) - │ → r.Update → object actually deleted └─ return {} (no finalizer, no - ├─ prov.Get → ErrNotFound ──► RemoveFinalizer provider calls ever; - │ → r.Update → object actually deleted the health engine — - └─ exists ──► prov.Delete Step 5 — drives the rest) + │ → 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) ``` @@ -129,4 +130,56 @@ prov.ListByTag ─► client.List(Pods by labels, all namespaces) ─┘ them The reconciler never watches provider-side resources (Pods now, GCP VMs later). 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. \ No newline at end of file +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. + +```text +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). diff --git a/docs/plans-executions/2026-08-07-1747-proxy-operator.md b/docs/plans-executions/2026-08-07-1747-proxy-operator.md index 7287b69..2938982 100644 --- a/docs/plans-executions/2026-08-07-1747-proxy-operator.md +++ b/docs/plans-executions/2026-08-07-1747-proxy-operator.md @@ -9,7 +9,7 @@ Pairs with [docs/plans/2026-08-07-1747-proxy-operator.md](../plans/2026-08-07-17 - [x] Step 2 — Provider contract (`internal/provider/`) - [x] Step 3 — Kubernetes pod provider (`internal/provider/kubernetes/`; first built as an in-memory mock, then replaced — see the two Step 3 sections below) - [x] Step 4 — Reconciler (`internal/controller/`) -- [ ] Step 5 — Health engine (`internal/health/`) +- [x] Step 5 — Health engine (`internal/health/`) - [ ] Step 6 — Lease store (`internal/lease/`) - [ ] Step 7 — Discovery API (`internal/discovery/`) - [ ] Step 8 — GCP provider (`internal/provider/gcp/`) @@ -626,3 +626,71 @@ deterministic and fast, at the cost of not exercising watch-driven requeues; Step 11's manager-driven cases cover that. The Secret watch is wired in `SetupWithManager` but the label-restricted Secret cache it assumes arrives with `cmd/main.go` in Step 10. + +## Step 5 — Health engine (`internal/health/`) + +Implemented the engine per the plan's design: `probe.go` (through-the-proxy +probe with the plan's exact transport — fresh per probe, +`DisableKeepAlives: true` so every probe re-exercises CONNECT) and +`engine.go` (leader-elected manager Runnable: 1 s scheduler tick + a pool +of 8 workers, per-proxy threshold state under one mutex, transition-only +emission over a buffered `chan event.GenericEvent`). The reconciler side +landed in the same step: a `HealthSnapshotter` interface + `applyHealth` +staging the Healthy condition/latency/lastHealthCheckTime from +`Engine.Snapshot`, and a conditional +`WatchesRawSource(source.Channel(...))` in `SetupWithManager`. Everything +tolerates nil (engine unwired) until `cmd/main.go` connects the two in +Step 10. + +Deviations and judgment calls beyond the plan text: + +- **First-probe scheduling is split by whether a verdict was seeded.** The + plan's startup jitter (`nextDue = now + rand(0, interval)`) applies only + to proxies whose state was seeded from an existing Healthy condition — + the restart case it exists for. A never-probed proxy is probed on the + next tick instead; making a brand-new proxy wait up to a full interval + for its first verdict would be pure lag with no thundering-herd benefit. +- **The latency-change emission rule only applies while the verdict is + healthy.** Caught by the first test run, not foreseen: a success streak + still below `successThreshold` (verdict unhealthy, reported unhealthy) + satisfied the plan's rule (c) — latency delta vs a stale reported value, + rate window open — and emitted a pointless latency-only update for a + proxy still reported as unhealthy. Guarded with `res.ok && *st.healthy`. +- **`ProbeTLSConfig` field added to the engine** (nil = system roots). The + probe function needs a CA override to be testable against + `httptest.NewTLSServer`, and the same knob is genuinely useful for + probing targets signed by a private CA. Not a test-only backdoor. +- **State pruning doubles as replacement hygiene:** any proxy with no + probeable host (provisioning, mid-replacement, deleting) has its state + dropped each tick, so a replacement instance always starts with fresh + counters. Complementarily, the reconciler's create branch removes the + stale Healthy condition and latency fields — a new VM shouldn't wear its + predecessor's verdict. + +Tests: a real CONNECT-capable proxy stub (hijack + bidirectional +`io.Copy`) probing a real `httptest.NewTLSServer` — CONNECT success, +refused CONNECT, unexpected status, dead proxy, plain-http forwarding; +table-driven threshold/suppression/seeding/pruning tests driving +`record`/`tick` directly; an end-to-end `Start` test (fake reader, fake +probeFn, 5 ms tick) asserting event delivery, snapshot content, and clean +shutdown on context cancel; and controller-side tests with a +`fakeSnapshotter` proving Running+healthy ⇒ `Ready`, Running+unhealthy ⇒ +`Unhealthy`, no-verdict ⇒ no condition, and stale-verdict cleanup on +replacement. + +Verification (all green): + +```bash +make test # envtest + units; health 93.4%, controller 77.4% +KUBEBUILDER_ASSETS="$PWD/bin/k8s/1.36.2-darwin-arm64" go test -race ./... +go test -race -count=2 ./internal/health/ # shook out the emission-rule bug above +``` + +Worth noting: `docs/architecture.md` (created between Steps 4 and 5 on +user request) gained a §6 for the engine and now shows the third workqueue +feed (`source.Channel`). The `Healthy` condition reasons live in the +controller package (`ReasonProbeSucceeded`/`ReasonProbeFailed`) — the +engine deliberately knows nothing about conditions except reading one at +seed time, keeping the state/representation split honest. The +`hint`-driven `wg.Go` idiom (Go 1.25+) replaced the classic +`wg.Add/defer wg.Done` in the worker pool. diff --git a/internal/controller/health_test.go b/internal/controller/health_test.go new file mode 100644 index 0000000..af84c46 --- /dev/null +++ b/internal/controller/health_test.go @@ -0,0 +1,175 @@ +package controller + +import ( + "strings" + "testing" + "time" + + apimeta "k8s.io/apimachinery/pkg/api/meta" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/types" + + 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" +) + +type fakeSnapshotter struct { + snap health.Snapshot + ok bool +} + +func (f fakeSnapshotter) Snapshot(types.NamespacedName) (health.Snapshot, bool) { + return f.snap, f.ok +} + +// TestReconcile_healthRepresentation covers the reconciler's half of the +// health split: turning the engine's Snapshot into the Healthy condition, +// the latency fields, and ultimately the Ready/Unhealthy phases. +func TestReconcile_healthRepresentation(t *testing.T) { + t.Parallel() + + probeTime := time.Now() + freshHash := specHash(managedProxy(), "") + runningStub := func() *stubProvider { + return &stubProvider{ + getInst: &provider.Instance{ID: "stub-id-1", IP: "10.1.2.3", State: provider.StateRunning}, + } + } + + tests := []struct { + name string + proxy *crawlv1alpha1.Proxy + stub *stubProvider + health HealthSnapshotter + wantPhase crawlv1alpha1.ProxyPhase + verify func(t *testing.T, r *ProxyReconciler) + }{ + { + name: "running and healthy becomes Ready", + proxy: managedProxy(withProviderID("stub-id-1"), withSpecHashAnnotation(freshHash)), + stub: runningStub(), + health: fakeSnapshotter{ok: true, snap: health.Snapshot{ + Healthy: true, Latency: 37 * time.Millisecond, LastProbe: probeTime, + }}, + wantPhase: crawlv1alpha1.PhaseReady, + verify: func(t *testing.T, r *ProxyReconciler) { + p := getProxy(t, r) + assertCondition(t, p, crawlv1alpha1.ConditionHealthy, metav1.ConditionTrue, ReasonProbeSucceeded) + if p.Status.LatencyMillis != 37 { + t.Errorf("latencyMillis = %d, want 37", p.Status.LatencyMillis) + } + if p.Status.LastHealthCheckTime == nil { + t.Error("lastHealthCheckTime not set") + } + }, + }, + { + name: "running but unhealthy becomes Unhealthy", + proxy: managedProxy(withProviderID("stub-id-1"), withSpecHashAnnotation(freshHash)), + stub: runningStub(), + health: fakeSnapshotter{ok: true, snap: health.Snapshot{ + Healthy: false, LastProbe: probeTime, + LastError: "CONNECT refused", ConsecutiveFailures: 3, + }}, + wantPhase: crawlv1alpha1.PhaseUnhealthy, + verify: func(t *testing.T, r *ProxyReconciler) { + p := getProxy(t, r) + assertCondition(t, p, crawlv1alpha1.ConditionHealthy, metav1.ConditionFalse, ReasonProbeFailed) + cond := apimeta.FindStatusCondition(p.Status.Conditions, crawlv1alpha1.ConditionHealthy) + if !strings.Contains(cond.Message, "CONNECT refused") { + t.Errorf("condition message %q does not carry the probe error", cond.Message) + } + }, + }, + { + name: "no verdict yet stays Provisioning without a Healthy condition", + proxy: managedProxy(withProviderID("stub-id-1"), withSpecHashAnnotation(freshHash)), + stub: runningStub(), + health: fakeSnapshotter{ok: false}, + wantPhase: crawlv1alpha1.PhaseProvisioning, + verify: func(t *testing.T, r *ProxyReconciler) { + p := getProxy(t, r) + if apimeta.FindStatusCondition(p.Status.Conditions, crawlv1alpha1.ConditionHealthy) != nil { + t.Error("Healthy condition present without an engine verdict") + } + }, + }, + { + name: "external proxy with a healthy verdict becomes Ready", + proxy: managedProxy(func(p *crawlv1alpha1.Proxy) { + p.Finalizers = nil + p.Spec = crawlv1alpha1.ProxySpec{ + Mode: crawlv1alpha1.ModeExternal, + Endpoint: &crawlv1alpha1.EndpointSpec{Host: "203.0.113.7"}, + } + }), + stub: &stubProvider{}, + health: fakeSnapshotter{ok: true, snap: health.Snapshot{ + Healthy: true, Latency: 5 * time.Millisecond, LastProbe: probeTime, + }}, + wantPhase: crawlv1alpha1.PhaseReady, + verify: func(t *testing.T, r *ProxyReconciler) { + assertCondition(t, getProxy(t, r), crawlv1alpha1.ConditionHealthy, metav1.ConditionTrue, ReasonProbeSucceeded) + }, + }, + { + name: "creating a replacement clears the stale Healthy verdict", + proxy: managedProxy(func(p *crawlv1alpha1.Proxy) { + p.Status.Conditions = []metav1.Condition{{ + Type: crawlv1alpha1.ConditionHealthy, Status: metav1.ConditionTrue, + Reason: ReasonProbeSucceeded, LastTransitionTime: metav1.Now(), + }} + p.Status.LatencyMillis = 42 + p.Status.LastHealthCheckTime = &metav1.Time{Time: probeTime} + }), + stub: &stubProvider{createID: "stub-id-2"}, + health: fakeSnapshotter{ok: false}, + wantPhase: crawlv1alpha1.PhaseProvisioning, + verify: func(t *testing.T, r *ProxyReconciler) { + p := getProxy(t, r) + if apimeta.FindStatusCondition(p.Status.Conditions, crawlv1alpha1.ConditionHealthy) != nil { + t.Error("stale Healthy condition survived instance creation") + } + if p.Status.LatencyMillis != 0 || p.Status.LastHealthCheckTime != nil { + t.Error("stale latency fields survived instance creation") + } + }, + }, + } + + for _, tc := range tests { + t.Run(tc.name, func(t *testing.T) { + t.Parallel() + r := newTestReconciler(t, tc.stub, tc.proxy) + r.Health = tc.health + if _, err := doReconcile(t, r); err != nil { + t.Fatalf("Reconcile: %v", err) + } + if got := getProxy(t, r).Status.Phase; got != tc.wantPhase { + t.Errorf("phase = %s, want %s", got, tc.wantPhase) + } + tc.verify(t, r) + }) + } +} + +// A nil Health snapshotter must disable representation entirely. +func TestReconcile_nilHealthSnapshotter(t *testing.T) { + t.Parallel() + freshHash := specHash(managedProxy(), "") + r := newTestReconciler(t, + &stubProvider{getInst: &provider.Instance{ID: "stub-id-1", IP: "10.1.2.3", State: provider.StateRunning}}, + managedProxy(withProviderID("stub-id-1"), withSpecHashAnnotation(freshHash))) + + if _, err := doReconcile(t, r); err != nil { + t.Fatalf("Reconcile: %v", err) + } + p := getProxy(t, r) + if apimeta.FindStatusCondition(p.Status.Conditions, crawlv1alpha1.ConditionHealthy) != nil { + t.Error("Healthy condition written with no snapshotter configured") + } + if p.Status.Phase != crawlv1alpha1.PhaseProvisioning { + t.Errorf("phase = %s, want Provisioning", p.Status.Phase) + } +} diff --git a/internal/controller/proxy_controller.go b/internal/controller/proxy_controller.go index 3cf3529..ec344d0 100644 --- a/internal/controller/proxy_controller.go +++ b/internal/controller/proxy_controller.go @@ -27,18 +27,30 @@ import ( 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" ) +// 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, @@ -50,6 +62,13 @@ type ProxyReconciler struct { // 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 @@ -142,6 +161,12 @@ func (r *ProxyReconciler) reconcileManaged(ctx context.Context, p *crawlv1alpha1 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 } @@ -176,6 +201,7 @@ func (r *ProxyReconciler) reconcileManaged(ctx context.Context, p *crawlv1alpha1 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 { @@ -262,6 +288,7 @@ func (r *ProxyReconciler) reconcileExternal(_ context.Context, p *crawlv1alpha1. } p.Status.IP = p.Spec.Endpoint.Host setProvisioned(p, metav1.ConditionTrue, ReasonExternalEndpoint, "tracking an external endpoint") + r.applyHealth(p) return ctrl.Result{}, nil } @@ -371,12 +398,15 @@ func (r *ProxyReconciler) proxiesForSecret(ctx context.Context, obj client.Objec // (cmd/main.go) restricts that cache to labelled cloud-init Secrets. func (r *ProxyReconciler) SetupWithManager(mgr ctrl.Manager) error { r.applyDefaults() - return ctrl.NewControllerManagedBy(mgr). + b := ctrl.NewControllerManagedBy(mgr). For(&crawlv1alpha1.Proxy{}). Named("proxy"). Watches(&corev1.Secret{}, handler.EnqueueRequestsFromMapFunc(r.proxiesForSecret)). - WithOptions(controller.Options{MaxConcurrentReconciles: 3}). - Complete(r) + WithOptions(controller.Options{MaxConcurrentReconciles: 3}) + if r.HealthEvents != nil { + b = b.WatchesRawSource(source.Channel(r.HealthEvents, &handler.EnqueueRequestForObject{})) + } + return b.Complete(r) } func (r *ProxyReconciler) applyDefaults() { diff --git a/internal/controller/status.go b/internal/controller/status.go index 3940c85..a2ec60f 100644 --- a/internal/controller/status.go +++ b/internal/controller/status.go @@ -2,6 +2,7 @@ package controller import ( "context" + "fmt" "k8s.io/apimachinery/pkg/api/equality" apimeta "k8s.io/apimachinery/pkg/api/meta" @@ -11,8 +12,9 @@ import ( crawlv1alpha1 "gitea.home.hrajfrisbee.cz/kacerr/egress-proxies-operator/api/v1alpha1" ) -// Reasons used on the Provisioned condition. The Healthy condition is owned -// by the health engine (internal/health) and only represented here. +// Reasons used on the Provisioned condition, plus the two the reconciler +// writes on the Healthy condition when representing the health engine's +// verdict (the engine owns the state; only the reconciler writes status). const ( ReasonProvisioning = "Provisioning" ReasonCreated = "Created" @@ -23,6 +25,9 @@ const ( ReasonCloudInitError = "CloudInitError" ReasonExternalEndpoint = "ExternalEndpoint" ReasonDeleting = "Deleting" + + ReasonProbeSucceeded = "ProbeSucceeded" + ReasonProbeFailed = "ProbeFailed" ) // setProvisioned stages the Provisioned condition on p. Nothing is written @@ -39,6 +44,39 @@ func setProvisioned(p *crawlv1alpha1.Proxy, status metav1.ConditionStatus, reaso }) } +// applyHealth stages the Healthy condition and the latency fields from the +// health engine's current snapshot. Called only from states where the proxy +// is reachable (Running, External); everywhere else the condition is either +// left as-is or removed by the create branch. +func (r *ProxyReconciler) applyHealth(p *crawlv1alpha1.Proxy) { + if r.Health == nil { + return + } + snap, ok := r.Health.Snapshot(client.ObjectKeyFromObject(p)) + if !ok { + return + } + cond := metav1.Condition{ + Type: crawlv1alpha1.ConditionHealthy, + ObservedGeneration: p.Generation, + } + if snap.Healthy { + cond.Status = metav1.ConditionTrue + cond.Reason = ReasonProbeSucceeded + cond.Message = "probe succeeded through the proxy" + } else { + cond.Status = metav1.ConditionFalse + cond.Reason = ReasonProbeFailed + cond.Message = fmt.Sprintf("%d consecutive probe failures; last: %s", + snap.ConsecutiveFailures, snap.LastError) + } + apimeta.SetStatusCondition(&p.Status.Conditions, cond) + p.Status.LatencyMillis = snap.Latency.Milliseconds() + if !snap.LastProbe.IsZero() { + p.Status.LastHealthCheckTime = &metav1.Time{Time: snap.LastProbe} + } +} + // computePhase derives status.phase from deletionTimestamp and the // Provisioned/Healthy conditions. Pure, so the truth table is unit-testable. func computePhase(p *crawlv1alpha1.Proxy) crawlv1alpha1.ProxyPhase { diff --git a/internal/health/engine.go b/internal/health/engine.go new file mode 100644 index 0000000..4377e4e --- /dev/null +++ b/internal/health/engine.go @@ -0,0 +1,340 @@ +package health + +import ( + "context" + "crypto/tls" + "math/rand/v2" + "net" + "net/url" + "strconv" + "sync" + "time" + + apimeta "k8s.io/apimachinery/pkg/api/meta" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/types" + "sigs.k8s.io/controller-runtime/pkg/client" + "sigs.k8s.io/controller-runtime/pkg/event" + logf "sigs.k8s.io/controller-runtime/pkg/log" + + crawlv1alpha1 "gitea.home.hrajfrisbee.cz/kacerr/egress-proxies-operator/api/v1alpha1" +) + +// Snapshot is the engine's current verdict for one proxy, read by the +// reconciler when it represents health in the Proxy's status. +type Snapshot struct { + Healthy bool + // Latency is the wall time of the most recent successful probe. + Latency time.Duration + // LastProbe is when the most recent probe (of either outcome) finished. + LastProbe time.Time + // LastError is the most recent probe failure; empty after a success. + LastError string + // ConsecutiveFailures is the current failure streak. + ConsecutiveFailures int32 +} + +// state is the engine's threshold bookkeeping for one proxy. The reported* +// fields track what has been delivered over Events; they only advance when a +// send succeeds, so a dropped event is retried after the next probe. +type state struct { + uid types.UID + inFlight bool + nextDue time.Time + + healthy *bool // nil until a first verdict exists + consecOK int32 + consecFail int32 + latency time.Duration + lastProbe time.Time + lastErr string + + reportedHealthy *bool + reportedLatency time.Duration + lastReport time.Time +} + +type probeJob struct { + key types.NamespacedName + uid types.UID + proxyURL *url.URL + hc crawlv1alpha1.HealthCheckSpec + interval time.Duration +} + +// Engine runs the probe scheduler and worker pool as a manager Runnable. It +// never writes Proxy status itself — keeping the reconciler the single +// status writer — and instead emits a GenericEvent per status-affecting +// transition, which the reconciler consumes via source.Channel. +type Engine struct { + // Reader lists Proxies from the manager's cache each tick. + Reader client.Reader + // Events carries one enqueue-request per status-affecting transition. + // Sends are non-blocking: a wedged reconciler must never stall probing. + Events chan event.GenericEvent + + // Workers is the probe worker pool size (default 8). + Workers int + // Tick is the scheduler interval (default 1s). At tens of proxies a + // per-second list scan is free; a timer wheel would be unjustified. + Tick time.Duration + // MinReportInterval rate-limits latency-only status reports (default 60s). + MinReportInterval time.Duration + // LatencyFloor is the absolute change below which a latency move is + // never status-affecting (default 20ms), so a proxy jittering around a + // small latency doesn't write status forever. + LatencyFloor time.Duration + // ProbeTLSConfig overrides TLS verification for https probe URLs; nil + // means system roots. Needed for private CAs (and tests). + ProbeTLSConfig *tls.Config + + probeFn func(context.Context, *url.URL, crawlv1alpha1.HealthCheckSpec, *tls.Config) probeResult + + mu sync.Mutex + states map[types.NamespacedName]*state +} + +// NewEngine returns an Engine with a buffered Events channel, ready to be +// handed to both mgr.Add and the reconciler (Health + HealthEvents fields). +func NewEngine(reader client.Reader) *Engine { + e := &Engine{Reader: reader} + e.applyDefaults() + return e +} + +func (e *Engine) applyDefaults() { + if e.Events == nil { + e.Events = make(chan event.GenericEvent, 64) + } + if e.Workers == 0 { + e.Workers = 8 + } + if e.Tick == 0 { + e.Tick = time.Second + } + if e.MinReportInterval == 0 { + e.MinReportInterval = time.Minute + } + if e.LatencyFloor == 0 { + e.LatencyFloor = 20 * time.Millisecond + } + if e.probeFn == nil { + e.probeFn = probe + } + if e.states == nil { + e.states = map[types.NamespacedName]*state{} + } +} + +// NeedLeaderElection makes the engine run only on the leader: probing from +// every replica would multiply load on the proxies, and only the leader's +// reconciler can represent the results anyway. +func (e *Engine) NeedLeaderElection() bool { return true } + +// Start runs the scheduler tick loop and the worker pool until ctx ends. +func (e *Engine) Start(ctx context.Context) error { + e.applyDefaults() + jobs := make(chan probeJob) + var wg sync.WaitGroup + for range e.Workers { + wg.Go(func() { + for { + select { + case <-ctx.Done(): + return + case job := <-jobs: + res := e.probeFn(ctx, job.proxyURL, job.hc, e.ProbeTLSConfig) + e.record(job, res, time.Now()) + } + } + }) + } + + ticker := time.NewTicker(e.Tick) + defer ticker.Stop() + for { + select { + case <-ctx.Done(): + wg.Wait() + return nil + case now := <-ticker.C: + e.tick(ctx, now, jobs) + } + } +} + +// tick lists proxies from the cache, refreshes the state map (create, seed, +// prune, UID-mismatch reset), and hands due proxies to the worker pool. +func (e *Engine) tick(ctx context.Context, now time.Time, jobs chan<- probeJob) { + var list crawlv1alpha1.ProxyList + if err := e.Reader.List(ctx, &list); err != nil { + logf.FromContext(ctx).Error(err, "health engine: listing proxies") + return + } + + e.mu.Lock() + defer e.mu.Unlock() + + probeable := make(map[types.NamespacedName]struct{}, len(list.Items)) + for i := range list.Items { + p := &list.Items[i] + host := p.EffectiveHost() + if host == "" || !p.DeletionTimestamp.IsZero() { + // Not probeable (provisioning, being replaced, or deleting). + // Its state gets pruned below, so a replacement instance starts + // with fresh counters. + continue + } + key := client.ObjectKeyFromObject(p) + probeable[key] = struct{}{} + + hc := p.HealthCheckOrDefault() + interval := time.Duration(hc.IntervalSeconds) * time.Second + + st := e.states[key] + if st == nil || st.uid != p.UID { + // New proxy, or a delete+recreate under the same name — never + // inherit the old object's counters. + st = newState(p, now, interval) + e.states[key] = st + } + if st.inFlight || now.Before(st.nextDue) { + continue + } + job := probeJob{ + key: key, + uid: p.UID, + proxyURL: &url.URL{Scheme: "http", Host: net.JoinHostPort(host, strconv.Itoa(int(p.EffectivePort())))}, + hc: hc, + interval: interval, + } + select { + case jobs <- job: + st.inFlight = true + default: + // Worker pool saturated; the proxy stays due and is retried on + // the next tick. + } + } + + for key := range e.states { + if _, ok := probeable[key]; !ok { + delete(e.states, key) + } + } +} + +// newState seeds bookkeeping for a proxy the engine hasn't tracked yet. If +// the CR already carries a Healthy verdict (leader handover, operator +// restart), the verdict is kept — so a healthy proxy doesn't flap to +// unknown — counters stay at zero so a real transition still needs a full +// threshold run, and the first probe is jittered across the interval so a +// restart doesn't fire the whole fleet's probes at once. A proxy with no +// prior verdict is probed immediately. +func newState(p *crawlv1alpha1.Proxy, now time.Time, interval time.Duration) *state { + st := &state{uid: p.UID, nextDue: now} + cond := apimeta.FindStatusCondition(p.Status.Conditions, crawlv1alpha1.ConditionHealthy) + if cond == nil || cond.Status == metav1.ConditionUnknown { + return st + } + healthy := cond.Status == metav1.ConditionTrue + reported := healthy + st.healthy = &healthy + st.reportedHealthy = &reported + st.latency = time.Duration(p.Status.LatencyMillis) * time.Millisecond + st.reportedLatency = st.latency + st.nextDue = now.Add(rand.N(interval)) + return st +} + +// record folds one probe result into the proxy's threshold state and emits +// an event when the result is status-affecting: a first-ever verdict, a +// threshold-crossing flip, or a material latency change (beyond +// max(LatencyFloor, 50% of reported) and rate-limited by MinReportInterval). +func (e *Engine) record(job probeJob, res probeResult, now time.Time) { + e.mu.Lock() + defer e.mu.Unlock() + + st := e.states[job.key] + if st == nil || st.uid != job.uid { + return // pruned or replaced while the probe was in flight + } + st.inFlight = false + st.nextDue = now.Add(job.interval) + st.lastProbe = now + + if res.ok { + st.consecOK++ + st.consecFail = 0 + st.latency = res.latency + st.lastErr = "" + } else { + st.consecFail++ + st.consecOK = 0 + st.lastErr = res.err.Error() + } + + switch { + case st.healthy == nil: + healthy := res.ok + st.healthy = &healthy + case *st.healthy && st.consecFail >= job.hc.FailureThreshold: + healthy := false + st.healthy = &healthy + case !*st.healthy && st.consecOK >= job.hc.SuccessThreshold: + healthy := true + st.healthy = &healthy + } + + var emit bool + switch { + case st.reportedHealthy == nil: + emit = true + case *st.reportedHealthy != *st.healthy: + emit = true + case res.ok && *st.healthy: + // Latency-only updates matter only for a healthy verdict; a success + // streak still below successThreshold must stay silent. + delta := st.latency - st.reportedLatency + if delta < 0 { + delta = -delta + } + emit = delta > max(e.LatencyFloor, st.reportedLatency/2) && + now.Sub(st.lastReport) > e.MinReportInterval + } + if !emit { + return + } + + evt := event.GenericEvent{Object: &crawlv1alpha1.Proxy{ + ObjectMeta: metav1.ObjectMeta{Namespace: job.key.Namespace, Name: job.key.Name}, + }} + select { + case e.Events <- evt: + reported := *st.healthy + st.reportedHealthy = &reported + st.reportedLatency = st.latency + st.lastReport = now + default: + // Channel full (reconciler wedged): drop, and deliberately do not + // advance the reported markers, so the next probe retries the emit. + } +} + +// Snapshot returns the engine's current verdict for key; ok is false while +// no verdict exists (never probed, or state was reset). +func (e *Engine) Snapshot(key types.NamespacedName) (Snapshot, bool) { + e.mu.Lock() + defer e.mu.Unlock() + st := e.states[key] + if st == nil || st.healthy == nil { + return Snapshot{}, false + } + return Snapshot{ + Healthy: *st.healthy, + Latency: st.latency, + LastProbe: st.lastProbe, + LastError: st.lastErr, + ConsecutiveFailures: st.consecFail, + }, true +} diff --git a/internal/health/engine_test.go b/internal/health/engine_test.go new file mode 100644 index 0000000..7802331 --- /dev/null +++ b/internal/health/engine_test.go @@ -0,0 +1,372 @@ +package health + +import ( + "context" + "crypto/tls" + "errors" + "net/url" + "strconv" + "strings" + "testing" + "time" + + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/runtime" + "k8s.io/apimachinery/pkg/types" + "sigs.k8s.io/controller-runtime/pkg/client/fake" + "sigs.k8s.io/controller-runtime/pkg/event" + + crawlv1alpha1 "gitea.home.hrajfrisbee.cz/kacerr/egress-proxies-operator/api/v1alpha1" +) + +var testKey = types.NamespacedName{Namespace: "default", Name: "p1"} + +func testEngine() *Engine { + e := &Engine{Events: make(chan event.GenericEvent, 8)} + e.applyDefaults() + return e +} + +func testJob(failureThreshold, successThreshold int32) probeJob { + return probeJob{ + key: testKey, + uid: "uid-1", + hc: crawlv1alpha1.HealthCheckSpec{ + FailureThreshold: failureThreshold, + SuccessThreshold: successThreshold, + }, + interval: 30 * time.Second, + } +} + +func drainOneEvent(t *testing.T, e *Engine) event.GenericEvent { + t.Helper() + select { + case evt := <-e.Events: + return evt + default: + t.Fatal("expected an event, channel is empty") + return event.GenericEvent{} + } +} + +func assertNoEvent(t *testing.T, e *Engine) { + t.Helper() + select { + case <-e.Events: + t.Fatal("unexpected event emitted") + default: + } +} + +func boolPtr(b bool) *bool { return &b } + +func TestRecord_firstResultEmits(t *testing.T) { + t.Parallel() + e := testEngine() + e.states[testKey] = &state{uid: "uid-1"} + + e.record(testJob(3, 1), probeResult{ok: true, latency: 30 * time.Millisecond}, time.Now()) + + evt := drainOneEvent(t, e) + if got := evt.Object.GetName(); got != "p1" { + t.Errorf("event object name = %q, want p1", got) + } + snap, ok := e.Snapshot(testKey) + if !ok || !snap.Healthy { + t.Errorf("Snapshot = %+v, %v; want healthy verdict", snap, ok) + } + if snap.Latency != 30*time.Millisecond { + t.Errorf("latency = %v, want 30ms", snap.Latency) + } +} + +func TestRecord_failureThresholdFlips(t *testing.T) { + t.Parallel() + e := testEngine() + e.states[testKey] = &state{uid: "uid-1", healthy: boolPtr(true), reportedHealthy: boolPtr(true)} + job := testJob(3, 1) + probeErr := probeResult{err: errors.New("connect refused")} + + e.record(job, probeErr, time.Now()) + e.record(job, probeErr, time.Now()) + assertNoEvent(t, e) + if snap, _ := e.Snapshot(testKey); !snap.Healthy { + t.Fatal("flipped unhealthy before failureThreshold was reached") + } + + e.record(job, probeErr, time.Now()) + drainOneEvent(t, e) + snap, _ := e.Snapshot(testKey) + if snap.Healthy { + t.Error("still healthy after failureThreshold consecutive failures") + } + if snap.ConsecutiveFailures != 3 { + t.Errorf("ConsecutiveFailures = %d, want 3", snap.ConsecutiveFailures) + } + if !strings.Contains(snap.LastError, "connect refused") { + t.Errorf("LastError = %q, want the probe error", snap.LastError) + } +} + +func TestRecord_successThresholdFlips(t *testing.T) { + t.Parallel() + e := testEngine() + e.states[testKey] = &state{uid: "uid-1", healthy: boolPtr(false), reportedHealthy: boolPtr(false)} + job := testJob(3, 2) + success := probeResult{ok: true, latency: 25 * time.Millisecond} + + e.record(job, success, time.Now()) + assertNoEvent(t, e) + + e.record(job, success, time.Now()) + drainOneEvent(t, e) + if snap, _ := e.Snapshot(testKey); !snap.Healthy { + t.Error("not healthy after successThreshold consecutive successes") + } +} + +func TestRecord_latencySuppression(t *testing.T) { + t.Parallel() + now := time.Now() + + tests := []struct { + name string + reportedLatency time.Duration + lastReport time.Time + newLatency time.Duration + wantEmit bool + }{ + { + name: "small change under the relative floor is suppressed", + reportedLatency: 100 * time.Millisecond, + lastReport: now.Add(-2 * time.Minute), + newLatency: 110 * time.Millisecond, + wantEmit: false, + }, + { + name: "small absolute jitter at low latency is suppressed", + reportedLatency: 5 * time.Millisecond, + lastReport: now.Add(-2 * time.Minute), + newLatency: 20 * time.Millisecond, // >50% but under the 20ms floor + wantEmit: false, + }, + { + name: "material change after the rate window emits", + reportedLatency: 100 * time.Millisecond, + lastReport: now.Add(-2 * time.Minute), + newLatency: 200 * time.Millisecond, + wantEmit: true, + }, + { + name: "material change inside the rate window is suppressed", + reportedLatency: 100 * time.Millisecond, + lastReport: now.Add(-10 * time.Second), + newLatency: 400 * time.Millisecond, + wantEmit: false, + }, + } + + for _, tc := range tests { + t.Run(tc.name, func(t *testing.T) { + t.Parallel() + e := testEngine() + e.states[testKey] = &state{ + uid: "uid-1", + healthy: boolPtr(true), + reportedHealthy: boolPtr(true), + reportedLatency: tc.reportedLatency, + lastReport: tc.lastReport, + } + e.record(testJob(3, 1), probeResult{ok: true, latency: tc.newLatency}, now) + if tc.wantEmit { + drainOneEvent(t, e) + } else { + assertNoEvent(t, e) + } + }) + } +} + +func TestRecord_droppedEventIsRetried(t *testing.T) { + t.Parallel() + e := testEngine() + e.Events = make(chan event.GenericEvent) // unbuffered, nobody reading + e.states[testKey] = &state{uid: "uid-1"} + success := probeResult{ok: true, latency: 30 * time.Millisecond} + + e.record(testJob(3, 1), success, time.Now()) + e.mu.Lock() + reported := e.states[testKey].reportedHealthy + e.mu.Unlock() + if reported != nil { + t.Fatal("reported marker advanced although the event was dropped") + } + + // Channel drains (reconciler recovers): the next probe re-emits. + e.Events = make(chan event.GenericEvent, 1) + e.record(testJob(3, 1), success, time.Now()) + drainOneEvent(t, e) +} + +func TestRecord_staleJobIsIgnored(t *testing.T) { + t.Parallel() + e := testEngine() + e.states[testKey] = &state{uid: "uid-NEW"} + + job := testJob(3, 1) + job.uid = "uid-OLD" + e.record(job, probeResult{ok: true, latency: time.Millisecond}, time.Now()) + + assertNoEvent(t, e) + if _, ok := e.Snapshot(testKey); ok { + t.Error("stale probe produced a verdict for the new object") + } +} + +func externalProxy(name, host string, mut ...func(*crawlv1alpha1.Proxy)) *crawlv1alpha1.Proxy { + p := &crawlv1alpha1.Proxy{ + ObjectMeta: metav1.ObjectMeta{ + Name: name, Namespace: "default", UID: types.UID("uid-" + name), + }, + Spec: crawlv1alpha1.ProxySpec{ + Mode: crawlv1alpha1.ModeExternal, + Endpoint: &crawlv1alpha1.EndpointSpec{Host: host, Port: 3128}, + }, + } + for _, m := range mut { + m(p) + } + return p +} + +func TestTick_schedulingAndPruning(t *testing.T) { + t.Parallel() + s := runtime.NewScheme() + if err := crawlv1alpha1.AddToScheme(s); err != nil { + t.Fatalf("scheme: %v", err) + } + + probeable := externalProxy("probeable", "10.0.0.1") + seeded := externalProxy("seeded", "10.0.0.2", func(p *crawlv1alpha1.Proxy) { + p.Status.Conditions = []metav1.Condition{{ + Type: crawlv1alpha1.ConditionHealthy, Status: metav1.ConditionTrue, + Reason: "ProbeSucceeded", LastTransitionTime: metav1.Now(), + }} + p.Status.LatencyMillis = 42 + }) + noIP := &crawlv1alpha1.Proxy{ + ObjectMeta: metav1.ObjectMeta{Name: "no-ip", Namespace: "default", UID: "uid-no-ip"}, + Spec: crawlv1alpha1.ProxySpec{Mode: crawlv1alpha1.ModeManaged, Provider: "stub"}, + } + + e := testEngine() + e.Reader = fake.NewClientBuilder().WithScheme(s). + WithObjects(probeable, seeded, noIP).Build() + // Stale entries: one for a proxy that no longer exists, one under a key + // that now belongs to a different UID (delete + recreate). + e.states[types.NamespacedName{Namespace: "default", Name: "gone"}] = &state{uid: "uid-gone"} + e.states[types.NamespacedName{Namespace: "default", Name: "probeable"}] = &state{ + uid: "uid-previous-incarnation", healthy: boolPtr(false), + } + + jobs := make(chan probeJob, 8) + e.tick(context.Background(), time.Now(), jobs) + + var dispatched []probeJob + for { + select { + case j := <-jobs: + dispatched = append(dispatched, j) + continue + default: + } + break + } + + if len(dispatched) != 1 { + t.Fatalf("dispatched %d jobs, want exactly 1 (only the fresh probeable proxy)", len(dispatched)) + } + j := dispatched[0] + if j.key.Name != "probeable" || j.uid != "uid-probeable" { + t.Errorf("dispatched job = %+v, want the recreated probeable proxy", j) + } + if want := "http://" + "10.0.0.1:" + strconv.Itoa(3128); j.proxyURL.String() != want { + t.Errorf("proxyURL = %s, want %s", j.proxyURL, want) + } + + e.mu.Lock() + defer e.mu.Unlock() + if _, ok := e.states[types.NamespacedName{Namespace: "default", Name: "gone"}]; ok { + t.Error("state for a deleted proxy was not pruned") + } + if _, ok := e.states[types.NamespacedName{Namespace: "default", Name: "no-ip"}]; ok { + t.Error("state was created for a proxy with no IP") + } + st := e.states[types.NamespacedName{Namespace: "default", Name: "probeable"}] + if st == nil || st.uid != "uid-probeable" { + t.Fatalf("state for recreated proxy = %+v, want fresh state with the new UID", st) + } + if st.healthy != nil && !*st.healthy { + t.Error("recreated proxy inherited the previous incarnation's unhealthy verdict") + } + seededSt := e.states[types.NamespacedName{Namespace: "default", Name: "seeded"}] + if seededSt == nil { + t.Fatal("no state created for the seeded proxy") + } + if seededSt.healthy == nil || !*seededSt.healthy { + t.Error("seeded proxy did not inherit its Healthy condition") + } + if seededSt.reportedHealthy == nil || !*seededSt.reportedHealthy { + t.Error("seeded verdict must count as already reported, or restart would re-emit for the whole fleet") + } + if seededSt.reportedLatency != 42*time.Millisecond { + t.Errorf("seeded reportedLatency = %v, want 42ms", seededSt.reportedLatency) + } + if seededSt.consecOK != 0 || seededSt.consecFail != 0 { + t.Error("seeded counters must start at zero") + } +} + +func TestEngine_StartEndToEnd(t *testing.T) { + t.Parallel() + s := runtime.NewScheme() + if err := crawlv1alpha1.AddToScheme(s); err != nil { + t.Fatalf("scheme: %v", err) + } + + e := testEngine() + e.Tick = 5 * time.Millisecond + e.Reader = fake.NewClientBuilder().WithScheme(s). + WithObjects(externalProxy("p1", "192.0.2.1")).Build() + e.probeFn = func(context.Context, *url.URL, crawlv1alpha1.HealthCheckSpec, *tls.Config) probeResult { + return probeResult{ok: true, latency: 12 * time.Millisecond} + } + + ctx, cancel := context.WithCancel(context.Background()) + done := make(chan error, 1) + go func() { done <- e.Start(ctx) }() + + select { + case evt := <-e.Events: + if evt.Object.GetName() != "p1" { + t.Errorf("event for %q, want p1", evt.Object.GetName()) + } + case <-time.After(5 * time.Second): + t.Fatal("no health event within 5s") + } + snap, ok := e.Snapshot(types.NamespacedName{Namespace: "default", Name: "p1"}) + if !ok || !snap.Healthy || snap.Latency != 12*time.Millisecond { + t.Errorf("Snapshot = %+v, %v; want healthy at 12ms", snap, ok) + } + + cancel() + select { + case err := <-done: + if err != nil { + t.Errorf("Start returned %v, want nil on context cancel", err) + } + case <-time.After(5 * time.Second): + t.Fatal("Start did not stop within 5s of cancel") + } +} diff --git a/internal/health/probe.go b/internal/health/probe.go new file mode 100644 index 0000000..4dbb1ea --- /dev/null +++ b/internal/health/probe.go @@ -0,0 +1,66 @@ +// Package health actively probes every proxy by fetching a URL through the +// proxy itself, keeps per-proxy threshold state, and pushes status-affecting +// transitions to the reconciler over a channel. The engine owns health +// state; the reconciler owns its representation in the Proxy's status. +package health + +import ( + "context" + "crypto/tls" + "fmt" + "net" + "net/http" + "net/url" + "slices" + "time" + + crawlv1alpha1 "gitea.home.hrajfrisbee.cz/kacerr/egress-proxies-operator/api/v1alpha1" +) + +// probeResult is the outcome of a single through-the-proxy probe. +type probeResult struct { + ok bool + latency time.Duration + err error +} + +// probe fetches hc.ProbeURL through the proxy at proxyURL. For an https +// probe URL the transport issues CONNECT to the proxy and TLS-handshakes +// through the tunnel; a proxy that accepts TCP but cannot egress answers +// CONNECT with a non-200, which client.Do surfaces as an error, not a +// response — so success requires err == nil AND an expected status code. +// tlsCfg is nil in production (system roots); tests and private-CA setups +// inject their own. +func probe(ctx context.Context, proxyURL *url.URL, hc crawlv1alpha1.HealthCheckSpec, tlsCfg *tls.Config) probeResult { + timeout := time.Duration(hc.TimeoutSeconds) * time.Second + transport := &http.Transport{ + Proxy: http.ProxyURL(proxyURL), + // Load-bearing: with keep-alives on, net/http caches the established + // CONNECT tunnel and later probes would never re-exercise CONNECT — + // exactly the failure this probe exists to catch. + DisableKeepAlives: true, + ForceAttemptHTTP2: false, + TLSHandshakeTimeout: timeout, + ResponseHeaderTimeout: timeout, + TLSClientConfig: tlsCfg, + DialContext: (&net.Dialer{Timeout: timeout}).DialContext, + } + defer transport.CloseIdleConnections() + + client := &http.Client{Transport: transport, Timeout: timeout} + req, err := http.NewRequestWithContext(ctx, http.MethodGet, hc.ProbeURL, nil) + if err != nil { + return probeResult{err: fmt.Errorf("building probe request: %w", err)} + } + start := time.Now() + resp, err := client.Do(req) + latency := time.Since(start) + if err != nil { + return probeResult{latency: latency, err: err} + } + defer func() { _ = resp.Body.Close() }() + if !slices.Contains(hc.ExpectedStatusCodes, int32(resp.StatusCode)) { + return probeResult{latency: latency, err: fmt.Errorf("unexpected status %d", resp.StatusCode)} + } + return probeResult{ok: true, latency: latency} +} diff --git a/internal/health/probe_test.go b/internal/health/probe_test.go new file mode 100644 index 0000000..26fc2cf --- /dev/null +++ b/internal/health/probe_test.go @@ -0,0 +1,169 @@ +package health + +import ( + "context" + "crypto/tls" + "crypto/x509" + "io" + "net" + "net/http" + "net/http/httptest" + "net/url" + "testing" + "time" + + crawlv1alpha1 "gitea.home.hrajfrisbee.cz/kacerr/egress-proxies-operator/api/v1alpha1" +) + +// startConnectProxy runs a minimal but real HTTP proxy: CONNECT tunneling +// for https targets, absolute-URI forwarding for plain http ones. With +// refuseConnect it answers CONNECT with 502 — the "accepts TCP but cannot +// egress" failure mode the probe must classify as unhealthy. +func startConnectProxy(t *testing.T, refuseConnect bool) *httptest.Server { + t.Helper() + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.Method == http.MethodConnect { + if refuseConnect { + http.Error(w, "no egress", http.StatusBadGateway) + return + } + dst, err := net.DialTimeout("tcp", r.Host, time.Second) + if err != nil { + http.Error(w, err.Error(), http.StatusBadGateway) + return + } + conn, bufrw, err := http.NewResponseController(w).Hijack() + if err != nil { + _ = dst.Close() + t.Errorf("hijack: %v", err) + return + } + _, _ = bufrw.WriteString("HTTP/1.1 200 Connection established\r\n\r\n") + _ = bufrw.Flush() + done := make(chan struct{}, 2) + go func() { _, _ = io.Copy(dst, bufrw); done <- struct{}{} }() + go func() { _, _ = io.Copy(conn, dst); done <- struct{}{} }() + <-done + _ = conn.Close() + _ = dst.Close() + return + } + out := r.Clone(r.Context()) + out.RequestURI = "" + resp, err := http.DefaultTransport.RoundTrip(out) + if err != nil { + http.Error(w, err.Error(), http.StatusBadGateway) + return + } + defer func() { _ = resp.Body.Close() }() + for k, vv := range resp.Header { + for _, v := range vv { + w.Header().Add(k, v) + } + } + w.WriteHeader(resp.StatusCode) + _, _ = io.Copy(w, resp.Body) + })) + t.Cleanup(srv.Close) + return srv +} + +func startTLSTarget(t *testing.T, status int) (*httptest.Server, *tls.Config) { + t.Helper() + target := httptest.NewTLSServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + w.WriteHeader(status) + })) + t.Cleanup(target.Close) + pool := x509.NewCertPool() + pool.AddCert(target.Certificate()) + return target, &tls.Config{RootCAs: pool} +} + +func proxyURL(t *testing.T, srv *httptest.Server) *url.URL { + t.Helper() + u, err := url.Parse(srv.URL) + if err != nil { + t.Fatalf("parsing proxy URL: %v", err) + } + return u +} + +func testHC(probeTarget string) crawlv1alpha1.HealthCheckSpec { + return crawlv1alpha1.HealthCheckSpec{ + ProbeURL: probeTarget, + IntervalSeconds: 30, + TimeoutSeconds: 5, + FailureThreshold: 3, + SuccessThreshold: 1, + ExpectedStatusCodes: []int32{200, 204}, + } +} + +func TestProbe_connectTunnelSucceeds(t *testing.T) { + t.Parallel() + target, tlsCfg := startTLSTarget(t, http.StatusNoContent) + proxy := startConnectProxy(t, false) + + res := probe(context.Background(), proxyURL(t, proxy), testHC(target.URL), tlsCfg) + if !res.ok { + t.Fatalf("probe failed through working proxy: %v", res.err) + } + if res.latency <= 0 { + t.Errorf("latency = %v, want > 0", res.latency) + } +} + +func TestProbe_refusedConnectFails(t *testing.T) { + t.Parallel() + target, tlsCfg := startTLSTarget(t, http.StatusNoContent) + proxy := startConnectProxy(t, true) + + res := probe(context.Background(), proxyURL(t, proxy), testHC(target.URL), tlsCfg) + if res.ok { + t.Fatal("probe succeeded through a proxy that refuses CONNECT") + } + if res.err == nil { + t.Error("expected an error from the refused CONNECT") + } +} + +func TestProbe_unexpectedStatusFails(t *testing.T) { + t.Parallel() + target, tlsCfg := startTLSTarget(t, http.StatusInternalServerError) + proxy := startConnectProxy(t, false) + + res := probe(context.Background(), proxyURL(t, proxy), testHC(target.URL), tlsCfg) + if res.ok { + t.Fatal("probe succeeded on a 500 response") + } +} + +func TestProbe_unreachableProxyFails(t *testing.T) { + t.Parallel() + // A listener that is immediately closed: guaranteed-refused port. + l, err := net.Listen("tcp", "127.0.0.1:0") + if err != nil { + t.Fatalf("reserving port: %v", err) + } + dead := &url.URL{Scheme: "http", Host: l.Addr().String()} + _ = l.Close() + + res := probe(context.Background(), dead, testHC("https://example.invalid/"), nil) + if res.ok { + t.Fatal("probe succeeded against a dead proxy") + } +} + +func TestProbe_plainHTTPForwardSucceeds(t *testing.T) { + t.Parallel() + target := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + w.WriteHeader(http.StatusNoContent) + })) + t.Cleanup(target.Close) + proxy := startConnectProxy(t, false) + + res := probe(context.Background(), proxyURL(t, proxy), testHC(target.URL), nil) + if !res.ok { + t.Fatalf("plain-http probe failed: %v", res.err) + } +}