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