From f09c1fbcca0bc9587f570d06c80018c11927063f Mon Sep 17 00:00:00 2001 From: Claude Date: Thu, 1 Oct 2026 09:43:09 +0200 Subject: [PATCH] fix(proxy): gate /readyz on a decodable Fetch round-trip The proxy is a full Fetch codec: it decodes, merges and re-encodes broker Fetch responses. A proxy-to-broker Fetch serialization or connection-state mismatch therefore breaks consume while Metadata and ListOffsets keep answering normally. Readiness based on "a backend exists" reports ready in that state, Kubernetes keeps the pod in the Service, and every consumer silently receives zero records. /readyz now answers two questions: is a backend known, and does that backend actually serve Fetch. The deep gate sends a Fetch for one partition of a deliberately unknown topic ID on a fresh connection, so it needs no real topic, reads no data, and also detects a stale or half-open pooled connection. A correct broker answers UNKNOWN_TOPIC_ID. KAFSCALE_PROXY_READYZ_FETCH_PROBE=false falls back to the shallow check. It is an escape hatch for a false negative in the field, not a steady state. Framing is the part worth guarding. forwardToBackend prepends the frame length itself, and encodeFetchRequest already strips the size prefix that kmsg's AppendRequest emits, so the probe payload must be the bare request header plus body. Stripping a second time removes the API key and version: the request is then not a malformed Fetch but a different API, and a broker that cannot parse a request header closes the connection. The probe reports EOF and the failure reads like a broken broker rather than a broken probe. buildFetchProbePayload is split out so that encoding is unit-testable without a backend. Two tests state the invariant from both sides: the payload parses back as Fetch v13 with the probe's topic ID, and a payload that has lost its four header bytes must not parse as Fetch, so the first test cannot silently stop distinguishing the bug. Startup ordering is a consequence worth documenting: with the probe enabled the proxy stays NotReady until a broker serves Fetch, so an installer that waits for proxy readiness before creating the resource that produces the broker will deadlock. - cmd/proxy/main.go: split checkReady into haveBackend plus the deep gate, add fetchProbe, buildFetchProbePayload and envBoolDefault - cmd/proxy/main_test.go: TestBuildFetchProbePayloadIsSingleFramed, TestBuildFetchProbePayloadRejectsDoubleStrip, TestEnvBoolDefault - docs/operations.md: "Proxy readiness" section and the env var index entry --- cmd/proxy/main.go | 128 ++++++++++++++++++++++++++++++++++++++++- cmd/proxy/main_test.go | 100 ++++++++++++++++++++++++++++++++ docs/operations.md | 35 +++++++++++ 3 files changed, 262 insertions(+), 1 deletion(-) diff --git a/cmd/proxy/main.go b/cmd/proxy/main.go index 865743a6..d3fefef2 100644 --- a/cmd/proxy/main.go +++ b/cmd/proxy/main.go @@ -65,7 +65,10 @@ type proxy struct { metaFlight singleflight.Group backendRetries int backendBackoff time.Duration - lfs *lfsModule // nil when LFS disabled + // readyzFetchProbe gates /readyz on an actual Fetch round-trip to a + // backend, not merely on "a backend exists". See fetchProbe. + readyzFetchProbe bool + lfs *lfsModule // nil when LFS disabled } func main() { @@ -121,6 +124,10 @@ func main() { topicNames: make(map[[16]byte]string), backendRetries: backendRetries, backendBackoff: backendBackoff, + // Default on: readiness requires a decodable Fetch round-trip. Set + // KAFSCALE_PROXY_READYZ_FETCH_PROBE=false to fall back to the shallow + // "a backend exists" check. + readyzFetchProbe: envBoolDefault("KAFSCALE_PROXY_READYZ_FETCH_PROBE", true), } if etcdStore, ok := store.(*metadata.EtcdStore); ok { @@ -221,6 +228,23 @@ func envInt(key string, fallback int) int { return parsed } +// envBoolDefault parses a boolean environment variable, falling back when the +// variable is unset or unparseable. +func envBoolDefault(key string, fallback bool) bool { + val := strings.TrimSpace(os.Getenv(key)) + if val == "" { + return fallback + } + switch strings.ToLower(val) { + case "1", "true", "yes", "y", "on": + return true + case "0", "false", "no", "n", "off": + return false + default: + return fallback + } +} + func portFromAddr(addr string, fallback int) int { _, portStr, err := net.SplitHostPort(addr) if err != nil { @@ -349,6 +373,30 @@ func (p *proxy) cacheFresh() bool { // fetch only when the cache TTL has expired (e.g. no traffic for >60s). // The fallback uses a short timeout to prevent health probes from blocking. func (p *proxy) checkReady(ctx context.Context) bool { + if !p.haveBackend(ctx) { + return false + } + // A known backend is necessary but not sufficient. The proxy is a full + // Fetch codec: it decodes, merges and re-encodes broker Fetch responses. + // A proxy<->broker Fetch (de)serialization or connection-state mismatch + // therefore breaks consume while Metadata and ListOffsets still answer + // normally, so the shallow check above passes and the proxy would serve + // traffic that silently returns zero records. Require a Fetch round-trip + // that actually decodes before reporting Ready. + if p.readyzFetchProbe { + if err := p.fetchProbe(ctx); err != nil { + p.logger.Warn("readyz fetch probe failed; staying NotReady", "error", err) + return false + } + } + return true +} + +// haveBackend reports whether the proxy knows of at least one broker backend. +// This is the shallow readiness check: cached state when fresh, falling back to +// a live metadata fetch only when the cache TTL has expired (e.g. no traffic +// for >60s). The fallback uses a short timeout so health probes do not block. +func (p *proxy) haveBackend(ctx context.Context) bool { if len(p.backends) > 0 { return true } @@ -364,6 +412,84 @@ func (p *proxy) checkReady(ctx context.Context) bool { return err == nil && len(backends) > 0 } +// fetchProbeTopicID is a non-zero, almost-certainly-unknown topic ID. Fetch +// v13 keys topics by ID rather than by name, so the broker resolves this ID, +// finds no match and returns a decodable UNKNOWN_TOPIC_ID response. That +// exercises the full v13 Fetch codec without depending on any real topic. +var fetchProbeTopicID = [16]byte{ + 0xab, 0xcd, 0xef, 0x01, 0x23, 0x45, 0x67, 0x89, + 0xab, 0xcd, 0xef, 0x01, 0x23, 0x45, 0x67, 0x89, +} + +// fetchProbeVersion is the highest Fetch version the proxy speaks. Probing the +// highest version is deliberate: it is the one real consumers negotiate and the +// one whose codec is most likely to diverge between proxy and broker. +const fetchProbeVersion int16 = 13 + +// buildFetchProbePayload encodes the readiness probe's Fetch request as a bare +// wire payload: request header plus body, with NO length prefix. forwardToBackend +// adds the frame length itself. +// +// Split out from fetchProbe so the encoding can be unit-tested without a +// backend. It is worth testing: a payload that is off by even the four header +// bytes is not a malformed Fetch, it is a different API key, and the broker +// answers that by closing the connection. The probe then reports EOF, which +// reads like a broken broker rather than a broken probe. +func buildFetchProbePayload() []byte { + clientID := "kafscale-proxy-readyz" + header := &protocol.RequestHeader{ + APIKey: protocol.APIKeyFetch, + APIVersion: fetchProbeVersion, + CorrelationID: 0, + ClientID: &clientID, + } + req := kmsg.NewPtrFetchRequest() + req.Version = fetchProbeVersion + req.ReplicaID = -1 // ordinary consumer + req.MaxWaitMillis = 0 + req.MinBytes = 0 + probePart := kmsg.NewFetchRequestTopicPartition() + probePart.Partition = 0 + probePart.FetchOffset = 0 + probePart.PartitionMaxBytes = 1 + probeTopic := kmsg.NewFetchRequestTopic() + probeTopic.TopicID = fetchProbeTopicID + probeTopic.Partitions = []kmsg.FetchRequestTopicPartition{probePart} + req.Topics = []kmsg.FetchRequestTopic{probeTopic} + // encodeFetchRequest already strips the size prefix that kmsg's + // AppendRequest emits, so its result is exactly the bare payload + // forwardToBackend expects. Do not strip again. + return encodeFetchRequest(header, req) +} + +// fetchProbe sends a minimal Fetch to a backend on a fresh connection and +// confirms the proxy can forward it and decode the response. The fresh +// connection matters: a stale or half-open pooled connection is detected too. +// +// Readiness therefore reflects "the broker actually serves Fetch" rather than +// "a backend exists", which also enforces coordinated startup ordering. +func (p *proxy) fetchProbe(ctx context.Context) error { + probeCtx, cancel := context.WithTimeout(ctx, 3*time.Second) + defer cancel() + + conn, addr, err := p.connectBackendExcluding(probeCtx, nil) + if err != nil { + return fmt.Errorf("fetch probe: no backend: %w", err) + } + defer conn.Close() + + respBytes, err := p.forwardToBackend(probeCtx, conn, addr, buildFetchProbePayload()) + if err != nil { + return fmt.Errorf("fetch probe forward to %s: %w", addr, err) + } + // The same decode path the proxy uses for real fetches: strip the response + // header, then kmsg ReadFrom. + if _, err := parseFetchResponse(respBytes, fetchProbeVersion); err != nil { + return fmt.Errorf("fetch probe from %s: %w", addr, err) + } + return nil +} + func (p *proxy) initMetadataCache(ctx context.Context) { if p.store == nil { return diff --git a/cmd/proxy/main_test.go b/cmd/proxy/main_test.go index 071eafdf..14b95a6a 100644 --- a/cmd/proxy/main_test.go +++ b/cmd/proxy/main_test.go @@ -23,6 +23,7 @@ import ( "io" "log/slog" "net" + "os" "sort" "sync" "testing" @@ -1593,3 +1594,102 @@ func splitFetchByOwnerForTest(req *kmsg.FetchRequest, owner func(topic string, p } return groups } + +// TestBuildFetchProbePayloadIsSingleFramed is the regression guard for the +// readiness probe's wire encoding. +// +// forwardToBackend prepends the frame length itself, and encodeFetchRequest +// already strips the size prefix that kmsg's AppendRequest emits. The probe +// payload must therefore be the bare request header plus body. Stripping a +// second time removes the API key and version (four bytes), and the broker then +// reads the correlation ID as the API key: the request is not a malformed +// Fetch, it is a different API entirely. A broker that cannot parse a request +// header closes the connection, so the probe reports EOF and the failure reads +// like a broken broker rather than a broken probe. +// +// Parsing the payload back with the same decoder the broker uses catches that +// without needing a backend. +func TestBuildFetchProbePayloadIsSingleFramed(t *testing.T) { + payload := buildFetchProbePayload() + + header, req, err := protocol.ParseRequest(payload) + if err != nil { + t.Fatalf("probe payload is not parseable as a request: %v", err) + } + if header.APIKey != protocol.APIKeyFetch { + t.Fatalf("api key = %d, want %d (Fetch); the payload is framed wrongly", + header.APIKey, protocol.APIKeyFetch) + } + if header.APIVersion != fetchProbeVersion { + t.Fatalf("api version = %d, want %d", header.APIVersion, fetchProbeVersion) + } + fetchReq, ok := req.(*kmsg.FetchRequest) + if !ok { + t.Fatalf("decoded request is %T, want *kmsg.FetchRequest", req) + } + if len(fetchReq.Topics) != 1 { + t.Fatalf("topics = %d, want 1", len(fetchReq.Topics)) + } + if fetchReq.Topics[0].TopicID != fetchProbeTopicID { + t.Fatalf("topic id = %x, want %x", fetchReq.Topics[0].TopicID, fetchProbeTopicID) + } + if len(fetchReq.Topics[0].Partitions) != 1 { + t.Fatalf("partitions = %d, want 1", len(fetchReq.Topics[0].Partitions)) + } + if fetchReq.ReplicaID != -1 { + t.Fatalf("replica id = %d, want -1 (ordinary consumer)", fetchReq.ReplicaID) + } +} + +// TestBuildFetchProbePayloadRejectsDoubleStrip states the same invariant from +// the other side: the four leading bytes carry the API key and version, so a +// payload that has lost them must NOT parse as a Fetch request. Without this +// the test above could pass against an encoder that happens to be tolerant. +func TestBuildFetchProbePayloadRejectsDoubleStrip(t *testing.T) { + payload := buildFetchProbePayload() + if len(payload) < 4 { + t.Fatalf("probe payload too short: %d bytes", len(payload)) + } + doubleStripped := payload[4:] + + header, _, err := protocol.ParseRequest(doubleStripped) + if err == nil && header.APIKey == protocol.APIKeyFetch { + t.Fatal("a payload stripped of its header still parsed as Fetch; " + + "this test can no longer distinguish the framing bug") + } +} + +func TestEnvBoolDefault(t *testing.T) { + const key = "KAFSCALE_TEST_ENV_BOOL" + cases := []struct { + val string + set bool + fallback bool + want bool + }{ + {set: false, fallback: true, want: true}, + {set: false, fallback: false, want: false}, + {val: "", set: true, fallback: true, want: true}, + {val: "true", set: true, fallback: false, want: true}, + {val: "TRUE", set: true, fallback: false, want: true}, + {val: " on ", set: true, fallback: false, want: true}, + {val: "1", set: true, fallback: false, want: true}, + {val: "false", set: true, fallback: true, want: false}, + {val: "off", set: true, fallback: true, want: false}, + {val: "0", set: true, fallback: true, want: false}, + // Unparseable must fall back rather than silently flip a gate. + {val: "maybe", set: true, fallback: true, want: true}, + {val: "maybe", set: true, fallback: false, want: false}, + } + for _, c := range cases { + if c.set { + t.Setenv(key, c.val) + } else { + os.Unsetenv(key) + } + if got := envBoolDefault(key, c.fallback); got != c.want { + t.Errorf("envBoolDefault(%q set=%v, fallback=%v) = %v, want %v", + c.val, c.set, c.fallback, got, c.want) + } + } +} diff --git a/docs/operations.md b/docs/operations.md index 5413aa35..c32249f5 100644 --- a/docs/operations.md +++ b/docs/operations.md @@ -283,6 +283,40 @@ Recommended operator alerting (when using Prometheus Operator): - `KafscaleSnapshotStale` – last successful snapshot older than the staleness threshold. - `KafscaleSnapshotNeverSucceeded` – no successful snapshots recorded. +## Proxy readiness + +The proxy exposes `/readyz` and `/livez` on `KAFSCALE_PROXY_HEALTH_ADDR` +(default `:9094`). + +`/readyz` answers two questions, not one: + +1. **Is a backend known?** Cached backend state when fresh, otherwise a live + metadata fetch from etcd. +2. **Does that backend actually serve Fetch?** A Fetch round-trip on a fresh + connection that must decode cleanly. + +The second question exists because the proxy is a full Fetch codec: it decodes, +merges and re-encodes broker Fetch responses. A proxy-to-broker Fetch +serialization or connection-state mismatch therefore breaks consume while +Metadata and ListOffsets keep answering normally. Readiness based on the first +question alone reports `ready` while every consumer silently receives zero +records, and Kubernetes keeps the pod in the Service. + +The probe asks for one partition of a deliberately unknown topic ID, so it needs +no real topic and reads no data; a correct broker answers `UNKNOWN_TOPIC_ID`. +The fresh connection is deliberate too: it detects a stale or half-open pooled +connection, which a cached-state check cannot. + +Set `KAFSCALE_PROXY_READYZ_FETCH_PROBE=false` to fall back to the shallow +"a backend exists" check. That is an escape hatch for a probe that reports a +false negative in the field, not a recommended steady state: it removes the only +readiness signal that distinguishes "serving" from "running". + +**Startup ordering.** With the probe enabled the proxy stays NotReady until a +broker serves Fetch. An installer that waits for proxy readiness before creating +the cluster resource that produces the broker will therefore deadlock. Create +the broker first, or do not gate the install on proxy readiness. + ## Environment Variable Index ### Operator @@ -387,6 +421,7 @@ Cost-optimized (accepts a larger loss window after crash): - `KAFSCALE_PROXY_BACKEND_CACHE_TTL_SEC` – Seconds to cache backend metadata before refreshing from etcd (default `60`). - `KAFSCALE_PROXY_BACKEND_BACKOFF_MS` – Backoff between backend connection retries (default `500`). - `KAFSCALE_PROXY_BACKEND_RETRIES` – Backend connection retry count (default `6`). +- `KAFSCALE_PROXY_READYZ_FETCH_PROBE` – Gate `/readyz` on a real Fetch round-trip (default `true`). See [Proxy readiness](#proxy-readiness). - `KAFSCALE_PROXY_LFS_ENABLED` – Enable the LFS HTTP API on the unified proxy (default `false`). ### Console