Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
128 changes: 127 additions & 1 deletion cmd/proxy/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -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() {
Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -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
}
Expand All @@ -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
Expand Down
100 changes: 100 additions & 0 deletions cmd/proxy/main_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@ import (
"io"
"log/slog"
"net"
"os"
"sort"
"sync"
"testing"
Expand Down Expand Up @@ -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)
}
}
}
35 changes: 35 additions & 0 deletions docs/operations.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down
Loading