diff --git a/.mise/tasks/michael/dev b/.mise/tasks/michael/dev index a5595121..e2aca6df 100755 --- a/.mise/tasks/michael/dev +++ b/.mise/tasks/michael/dev @@ -3,4 +3,8 @@ set -e source "$(dirname "$0")/../restic-api/env" +# No load balancer in front of a local run, so there is nothing to drain away +# from — skip the fleet's 15s SIGTERM drain and let Ctrl-C exit immediately. +export DRAIN_DELAY_MS=${DRAIN_DELAY_MS:-0} + cd packages/michael && go run . diff --git a/charts/apps/michael/values.yaml b/charts/apps/michael/values.yaml index 24638619..18ca6183 100644 --- a/charts/apps/michael/values.yaml +++ b/charts/apps/michael/values.yaml @@ -3,6 +3,27 @@ # Stateless S3 proxy, safe to run N-way. replicas: 2 +# Never dip below the declared replica count mid-rollout. The Deployment default +# (maxUnavailable 25%) would take 3 of father's 12 gateways down at once, on top +# of the 3 it surges in — a quarter of the restic data plane gone while the +# replacements are still probing their RGW backends. +strategy: + type: RollingUpdate + rollingUpdate: + maxSurge: 25% + maxUnavailable: 0 + +# Half the fleet, not the lib chart's minAvailable: 1 — a node drain may +# otherwise evict 11 of 12 gateways simultaneously. +pdb: + minAvailable: 50% + +# Must exceed DRAIN_DELAY_MS + SHUTDOWN_TIMEOUT_MS (15s + 120s in the base +# HelmRelease), or the kubelet SIGKILLs the pod while a restic blob is still on +# the wire — precisely the failure the drain exists to avoid. It is a ceiling, +# not a wait: an idle pod exits as soon as the drain delay elapses. +terminationGracePeriodSeconds: 150 + # Generous requests, deliberately huge memory limits (never CPU-limit). resources: requests: { cpu: 250m, memory: 256Mi } @@ -52,6 +73,10 @@ env: value: "3010" - name: LOG_LEVEL value: debug + # Dev only: the fleet default is 15s, which would make every `tilt down` and + # pod delete sit for a quarter minute with nothing to drain. + - name: DRAIN_DELAY_MS + value: "1000" # S3 object store = Rook-Ceph RGW. michael creates one bucket per restic # repository, so it needs a full RGW user (CephObjectStoreUser via # charts/ceph-objectuser), NOT a bucket-scoped ObjectBucketClaim. Rook writes @@ -115,6 +140,11 @@ startupProbe: tcpSocket: { port: http } periodSeconds: 2 failureThreshold: 90 +# Readiness is an HTTP endpoint, not a TCP connect, because it has to report +# one thing a bound socket cannot: that the process has begun draining. Probed +# every 2s so a SIGTERMed pod leaves the Service's endpoints well inside its +# drain delay, taking new restic requests with it. readinessProbe: - tcpSocket: { port: http } - periodSeconds: 10 + httpGet: { path: /readyz, port: http } + periodSeconds: 2 + failureThreshold: 2 diff --git a/charts/lib/yucca-common/templates/_deployment.tpl b/charts/lib/yucca-common/templates/_deployment.tpl index c0cfc598..3c28ff27 100644 --- a/charts/lib/yucca-common/templates/_deployment.tpl +++ b/charts/lib/yucca-common/templates/_deployment.tpl @@ -7,6 +7,10 @@ metadata: {{- include "yucca-common.labels" . | nindent 4 }} spec: replicas: {{ .Values.replicas | default 1 }} + {{- with .Values.strategy }} + strategy: + {{- toYaml . | nindent 4 }} + {{- end }} selector: matchLabels: {{- include "yucca-common.selectorLabels" . | nindent 6 }} @@ -19,6 +23,9 @@ spec: {{- toYaml . | nindent 8 }} {{- end }} spec: + {{- with .Values.terminationGracePeriodSeconds }} + terminationGracePeriodSeconds: {{ . }} + {{- end }} {{- with .Values.serviceAccountName }} serviceAccountName: {{ . }} {{- end }} diff --git a/kubernetes/apps/base/michael/helmrelease.yaml b/kubernetes/apps/base/michael/helmrelease.yaml index 5250c450..2c7d54e5 100644 --- a/kubernetes/apps/base/michael/helmrelease.yaml +++ b/kubernetes/apps/base/michael/helmrelease.yaml @@ -6,6 +6,10 @@ metadata: name: michael spec: interval: 1h + # Well above helm's 5m default: with maxUnavailable 0 a terminating pod holds + # its slot until its in-flight restic requests drain, so father's 12 gateways + # roll at the pace of DRAIN_DELAY_MS + SHUTDOWN_TIMEOUT_MS per wave. + timeout: 15m chart: spec: chart: charts/apps/michael @@ -42,6 +46,17 @@ spec: value: "3010" - name: LOG_LEVEL value: info + # Rollout drain. DRAIN_DELAY_MS is the window between /readyz going 503 + # and the listener closing — it has to outlast readiness detection (2s + # probe x 2 failures) plus EndpointSlice and Envoy EDS propagation, or the + # gateway is still steering fresh restic requests at a closing socket. + # SHUTDOWN_TIMEOUT_MS then caps how long an in-flight blob transfer gets + # to finish. The chart's terminationGracePeriodSeconds must stay above + # their sum. + - name: DRAIN_DELAY_MS + value: "15000" + - name: SHUTDOWN_TIMEOUT_MS + value: "120000" # Must match the cluster's STORAGE_CLUSTER_CODE (topology ConfigMap): # yucca-api mints restic tokens with a storageCluster claim, and michael # fails closed on cluster codes it does not front. diff --git a/kubernetes/components/apps/michael.yaml b/kubernetes/components/apps/michael.yaml index b164effa..89cf15cd 100644 --- a/kubernetes/components/apps/michael.yaml +++ b/kubernetes/components/apps/michael.yaml @@ -17,7 +17,8 @@ spec: namespace: yucca interval: 1h retryInterval: 2m - timeout: 10m + # Must outlast the HelmRelease's own 15m rollout budget (draining gateways). + timeout: 20m path: ./kubernetes/apps/base/michael prune: true wait: true diff --git a/packages/michael/internal/config/config.go b/packages/michael/internal/config/config.go index 8ee6baf3..460a459e 100644 --- a/packages/michael/internal/config/config.go +++ b/packages/michael/internal/config/config.go @@ -64,6 +64,18 @@ type Config struct { // address into. Only its LAST entry is trusted — see geoip.ClientAddr. ClientIPHeader string + // DrainDelay is how long michael keeps serving normally after SIGTERM while + // /readyz already answers 503. It must outlast readiness detection plus + // EndpointSlice and Envoy EDS propagation: cut short, the gateway is still + // routing fresh restic requests at a socket that has begun closing. + DrainDelay time.Duration + // ShutdownTimeout caps how long in-flight requests get to finish once the + // drain delay has elapsed. One restic blob can be tens of megabytes over a + // slow client uplink, so this is minutes — and the pod's + // terminationGracePeriodSeconds must exceed DrainDelay + ShutdownTimeout, or + // the kubelet SIGKILLs mid-upload anyway. + ShutdownTimeout time.Duration + OTLPMetricsEndpoint string OTLPMetricsURLPath string OTLPMetricsInterval time.Duration @@ -239,6 +251,9 @@ func LoadConfig() Config { asnDatabasePath := envOr("ASN_DB_PATH", "/etc/michael/asn.mmdb") clientIPHeader := envOr("CLIENT_IP_HEADER", "X-Forwarded-For") + drainDelay := envDurationMS("DRAIN_DELAY_MS", 15*time.Second) + shutdownTimeout := envDurationMS("SHUTDOWN_TIMEOUT_MS", 2*time.Minute) + otlpEndpoint := os.Getenv("OTLP_METRICS_ENDPOINT") otlpURLPath := os.Getenv("OTLP_METRICS_URL_PATH") otlpInterval := 1000 * time.Millisecond @@ -296,6 +311,8 @@ func LoadConfig() Config { S3DefaultCluster: defaultCluster, ASNDatabasePath: asnDatabasePath, ClientIPHeader: clientIPHeader, + DrainDelay: drainDelay, + ShutdownTimeout: shutdownTimeout, OTLPMetricsEndpoint: otlpEndpoint, OTLPMetricsURLPath: otlpURLPath, OTLPMetricsInterval: otlpInterval, @@ -610,6 +627,18 @@ func boolOr(v *bool, fallback bool) bool { return fallback } +func envDurationMS(key string, fallback time.Duration) time.Duration { + v := os.Getenv(key) + if v == "" { + return fallback + } + ms, err := strconv.Atoi(v) + if err != nil || ms < 0 { + log.Fatal().Msgf("%s must be a number >= 0", key) + } + return time.Duration(ms) * time.Millisecond +} + func envOr(key, fallback string) string { if v := os.Getenv(key); v != "" { return v diff --git a/packages/michael/internal/handlers/server.go b/packages/michael/internal/handlers/server.go index 794be089..7cebfbb1 100644 --- a/packages/michael/internal/handlers/server.go +++ b/packages/michael/internal/handlers/server.go @@ -5,6 +5,7 @@ import ( "crypto/ecdsa" "fmt" "net/http" + "sync/atomic" "time" "michael/internal/auth" @@ -33,8 +34,21 @@ type Server struct { // metrics and the access log. Optional: nil leaves both unattributed, which // is what the tests and any deployment without an ASN database get. ResolveClient func(*http.Request) geoip.Client + + draining atomic.Bool } +// readyPath is the kubelet's readiness probe. A repository path is always a +// UUID (auth rejects anything else), so no repository can collide with it. +const readyPath = "/readyz" + +// BeginDrain fails readiness from here on. Nothing already running is touched +// and requests that still arrive are still served in full — the point is only +// to get this replica out of the gateway's endpoint set before the listener +// closes, so a rollout moves restic to another replica instead of resetting it +// mid-backup. +func (s *Server) BeginDrain() { s.draining.Store(true) } + // NewServer builds a single-cluster server: everything is served from s under // the default cluster code. func NewServer(s storage.Storage, jwtPublicKey *ecdsa.PublicKey, m *metrics.Metrics) *Server { @@ -170,7 +184,28 @@ func (s *Server) Handler() http.Handler { }) }) - return r + return s.withReadiness(r) +} + +// withReadiness answers the readiness probe ahead of the router, keeping a +// probe every couple of seconds out of the access log and the request metrics. +// +// Backend health is deliberately not part of the answer: an RGW outage is +// shared fate across every replica, so failing readiness on it would empty the +// Service's endpoints and turn a degraded data plane into an unreachable one. +// The pool sheds with 503 + Retry-After for that case instead. +func (s *Server) withReadiness(next http.Handler) http.Handler { + return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.URL.Path != readyPath { + next.ServeHTTP(w, r) + return + } + if s.draining.Load() { + http.Error(w, "draining", http.StatusServiceUnavailable) + return + } + w.WriteHeader(http.StatusOK) + }) } // op stamps the operation name and the resolved route pattern onto the request diff --git a/packages/michael/internal/handlers/server_test.go b/packages/michael/internal/handlers/server_test.go index ca548db7..4d03fbe9 100644 --- a/packages/michael/internal/handlers/server_test.go +++ b/packages/michael/internal/handlers/server_test.go @@ -418,3 +418,46 @@ func TestClientNetworkLogOutput(t *testing.T) { } } } + +func TestReadyzReportsDrainState(t *testing.T) { + srv := newTestServer(&mockStorage{}) + handler := srv.Handler() + + rec := httptest.NewRecorder() + handler.ServeHTTP(rec, httptest.NewRequest(http.MethodGet, readyPath, nil)) + if rec.Code != http.StatusOK { + t.Fatalf("expected 200 before draining, got %d", rec.Code) + } + + srv.BeginDrain() + + rec = httptest.NewRecorder() + handler.ServeHTTP(rec, httptest.NewRequest(http.MethodGet, readyPath, nil)) + if rec.Code != http.StatusServiceUnavailable { + t.Fatalf("expected 503 while draining, got %d", rec.Code) + } +} + +// A draining instance must keep serving: readiness only steers the gateway +// away, it never rejects a request that still arrives. +func TestDrainingStillServesRequests(t *testing.T) { + store := &mockStorage{ + checkBucketFn: func(context.Context, string) (bool, error) { return false, nil }, + createBucketFn: func(context.Context, string) error { return nil }, + } + srv := newTestServer(store) + srv.BeginDrain() + + req := httptest.NewRequest(http.MethodPost, "/"+testRepository+"/?create=true", nil) + req.Header.Set("Authorization", makeBasicAuth(makeJWT(t, jwt.MapClaims{ + "user": testUser, + "repository": testRepository, + "writeOnce": false, + }))) + + rec := httptest.NewRecorder() + srv.Handler().ServeHTTP(rec, req) + if rec.Code != http.StatusOK { + t.Fatalf("expected the repository create to succeed while draining, got %d: %s", rec.Code, rec.Body) + } +} diff --git a/packages/michael/main.go b/packages/michael/main.go index 2318c2e2..d7315e52 100644 --- a/packages/michael/main.go +++ b/packages/michael/main.go @@ -119,13 +119,15 @@ func main() { Handler: srv.Handler(), } - // Graceful shutdown ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM) defer stop() - // Drive each backend pool's reconcile/probe loop until shutdown. + // The pools outlive the signal context on purpose: requests still draining + // after SIGTERM need a backend set someone is still resolving and probing. + poolCtx, stopPools := context.WithCancel(context.Background()) + defer stopPools() for _, pool := range pools { - go pool.Run(ctx) + go pool.Run(poolCtx) } go func() { @@ -136,27 +138,33 @@ func main() { }() <-ctx.Done() - log.Info().Msg("shutting down") + // Restore default signal handling: a second SIGTERM/SIGINT from an impatient + // operator now kills the process instead of being swallowed by the drain. + stop() - shutdownCtx, cancel := context.WithTimeout(context.Background(), 10*time.Second) - defer cancel() - - if err := httpSrv.Shutdown(shutdownCtx); err != nil { - log.Fatal().Err(err).Msg("shutdown error") + if err := drain(srv, httpSrv, cfg.DrainDelay, cfg.ShutdownTimeout); err != nil { + log.Error().Err(err).Msg("shutdown deadline hit with requests still in flight") } + stopPools() if err := asnDB.Close(); err != nil { log.Error().Err(err).Msg("ASN database close error") } + // Telemetry gets its own budget: a drain that used up shutdownCtx would + // otherwise hand the exporters an already-expired deadline and lose the + // final flush — exactly the window worth having logs for. + telemetryCtx, cancelTelemetry := context.WithTimeout(context.Background(), 10*time.Second) + defer cancelTelemetry() + if meterProvider != nil { - if err := meterProvider.Shutdown(shutdownCtx); err != nil { + if err := meterProvider.Shutdown(telemetryCtx); err != nil { log.Error().Err(err).Msg("meter provider shutdown error") } } if otelLogWriter != nil { - if err := otelLogWriter.Shutdown(shutdownCtx); err != nil { + if err := otelLogWriter.Shutdown(telemetryCtx); err != nil { log.Error().Err(err).Msg("OTLP log provider shutdown error") } } @@ -164,6 +172,24 @@ func main() { log.Info().Msg("shutdown complete") } +// drain retires this instance in two phases. Phase one fails /readyz while +// serving exactly as before, giving the readiness probe, the EndpointSlice and +// Envoy's EDS push time to move new restic requests onto another replica. Only +// then does phase two close the listener and wait out whatever is still in +// flight — closing it up front would reset live blob transfers and fail the +// backup that was running. +func drain(srv *handlers.Server, httpSrv *http.Server, delay, timeout time.Duration) error { + srv.BeginDrain() + log.Info().Dur("drain_delay", delay).Msg("draining: readiness failing, still serving") + time.Sleep(delay) + + ctx, cancel := context.WithTimeout(context.Background(), timeout) + defer cancel() + + log.Info().Dur("timeout", timeout).Msg("closing listener, waiting for in-flight requests") + return httpSrv.Shutdown(ctx) +} + // buildClusters builds one Storage per storage cluster michael fronts, keyed by // cluster code, plus the subset of those that are load-balancing pools (returned // concretely so the caller can register their metrics and run their reconcile diff --git a/packages/michael/main_test.go b/packages/michael/main_test.go new file mode 100644 index 00000000..3fb9f691 --- /dev/null +++ b/packages/michael/main_test.go @@ -0,0 +1,343 @@ +package main + +import ( + "bufio" + "context" + "crypto/ecdsa" + "crypto/elliptic" + "crypto/rand" + "encoding/base64" + "errors" + "io" + "net" + "net/http" + "strconv" + "strings" + "sync" + "testing" + "time" + + "michael/internal/handlers" + "michael/internal/storage" + + "github.com/golang-jwt/jwt/v5" +) + +var drainKey, _ = ecdsa.GenerateKey(elliptic.P256(), rand.Reader) + +const ( + drainUser = "00000000-0000-0000-0000-0000000000d1" + drainRepo = "00000000-0000-0000-0000-0000000000d2" + drainBlob = "1111111111111111111111111111111111111111111111111111111111111111" +) + +// drainStorage stands in for the S3 backend. PutObject drains the request body +// and then blocks until release is closed, so a test can hold one blob upload +// open across an entire drain the way a slow restic client would. +type drainStorage struct { + started chan string + release chan struct{} + releaseOnce sync.Once + mu sync.Mutex + stored map[string]int +} + +func (d *drainStorage) releaseAll() { d.releaseOnce.Do(func() { close(d.release) }) } + +func newDrainStorage() *drainStorage { + return &drainStorage{ + started: make(chan string, 8), + release: make(chan struct{}), + stored: map[string]int{}, + } +} + +func (d *drainStorage) PutObject(ctx context.Context, bucket, key string, body io.Reader, _ int64, _ bool, _ string) error { + n, err := io.Copy(io.Discard, body) + if err != nil { + return err + } + d.started <- key + // "locks" blobs return immediately: a test needs some request that completes + // on its own while another is still pinned open. + if !strings.HasPrefix(key, "locks/") { + <-d.release + } + d.mu.Lock() + d.stored[key] = int(n) + d.mu.Unlock() + return nil +} + +func (d *drainStorage) storedBytes(key string) (int, bool) { + d.mu.Lock() + defer d.mu.Unlock() + n, ok := d.stored[key] + return n, ok +} + +func (d *drainStorage) CheckBucket(context.Context, string) (bool, error) { return true, nil } +func (d *drainStorage) CreateBucket(context.Context, string) error { return nil } +func (d *drainStorage) HeadObject(context.Context, string, string) (int64, error) { + return 0, nil +} +func (d *drainStorage) GetObject(context.Context, string, string, string) (*storage.S3Object, error) { + return nil, errors.New("not used") +} +func (d *drainStorage) ListObjects(context.Context, string, string, func(storage.BlobInfo) error) error { + return nil +} +func (d *drainStorage) DeleteObject(context.Context, string, string) error { return nil } + +func drainAuthHeader(t *testing.T) string { + t.Helper() + token := jwt.NewWithClaims(jwt.SigningMethodES256, jwt.MapClaims{ + "user": drainUser, + "repository": drainRepo, + "writeOnce": false, + }) + signed, err := token.SignedString(drainKey) + if err != nil { + t.Fatalf("sign token: %v", err) + } + return "Basic " + base64.StdEncoding.EncodeToString([]byte("restic:"+signed)) +} + +// drainHarness is michael's real handler stack on a real TCP listener, reached +// over a real HTTP client — the drain is about socket lifecycle, so httptest's +// in-memory shortcuts would test the wrong thing. +type drainHarness struct { + store *drainStorage + srv *handlers.Server + httpSrv *http.Server + addr string + client *http.Client +} + +func newDrainHarness(t *testing.T) *drainHarness { + t.Helper() + store := newDrainStorage() + srv := handlers.NewServer(store, &drainKey.PublicKey, nil) + + ln, err := net.Listen("tcp", "127.0.0.1:0") + if err != nil { + t.Fatalf("listen: %v", err) + } + httpSrv := &http.Server{Handler: srv.Handler()} + go func() { _ = httpSrv.Serve(ln) }() + + h := &drainHarness{ + store: store, + srv: srv, + httpSrv: httpSrv, + addr: ln.Addr().String(), + client: &http.Client{Transport: &http.Transport{}}, + } + t.Cleanup(func() { + store.releaseAll() + _ = httpSrv.Close() + }) + return h +} + +func (h *drainHarness) url(path string) string { return "http://" + h.addr + path } + +func (h *drainHarness) putBlob(t *testing.T, blobType, name, body string) (*http.Response, error) { + t.Helper() + req, err := http.NewRequest(http.MethodPost, h.url("/"+drainRepo+"/"+blobType+"/"+name), strings.NewReader(body)) + if err != nil { + t.Fatalf("build request: %v", err) + } + req.Header.Set("Authorization", drainAuthHeader(t)) + return h.client.Do(req) +} + +func (h *drainHarness) readyStatus(t *testing.T) int { + t.Helper() + resp, err := h.client.Get(h.url("/readyz")) + if err != nil { + t.Fatalf("GET /readyz: %v", err) + } + defer resp.Body.Close() + _, _ = io.Copy(io.Discard, resp.Body) + return resp.StatusCode +} + +// The whole point of the two phases: while the drain delay runs the instance is +// unready but fully functional, and only once the listener closes does it stop +// taking connections — with the upload that was already running still finishing. +func TestDrainServesThroughTheDelayThenFinishesInFlight(t *testing.T) { + h := newDrainHarness(t) + payload := strings.Repeat("x", 4096) + + if got := h.readyStatus(t); got != http.StatusOK { + t.Fatalf("ready before drain: got %d, want 200", got) + } + + uploadDone := make(chan *http.Response, 1) + uploadErr := make(chan error, 1) + go func() { + resp, err := h.putBlob(t, "data", drainBlob, payload) + if err != nil { + uploadErr <- err + return + } + uploadDone <- resp + }() + + select { + case key := <-h.store.started: + if key != "data/"+drainBlob { + t.Fatalf("unexpected in-flight key %q", key) + } + case err := <-uploadErr: + t.Fatalf("upload failed before drain started: %v", err) + case <-time.After(5 * time.Second): + t.Fatal("upload never reached the storage layer") + } + + drainErr := make(chan error, 1) + drainStarted := time.Now() + go func() { drainErr <- drain(h.srv, h.httpSrv, 750*time.Millisecond, 10*time.Second) }() + + // Phase one: unready, but new work is still accepted and served in full. + waitFor(t, 2*time.Second, "readyz to report draining", func() bool { + return h.readyStatus(t) == http.StatusServiceUnavailable + }) + if elapsed := time.Since(drainStarted); elapsed > 700*time.Millisecond { + t.Fatalf("readiness only flipped after %s — the delay had already elapsed", elapsed) + } + + lockName := strings.Repeat("a", 64) + resp, err := h.putBlob(t, "locks", lockName, "lock") + if err != nil { + t.Fatalf("request during drain delay failed: %v", err) + } + resp.Body.Close() + if resp.StatusCode != http.StatusOK { + t.Fatalf("request during drain delay: got %d, want 200", resp.StatusCode) + } + <-h.store.started + + // Phase two: the listener closes, so nothing new can connect. + waitFor(t, 5*time.Second, "listener to stop accepting", func() bool { + c, err := net.DialTimeout("tcp", h.addr, 500*time.Millisecond) + if err != nil { + return true + } + _ = c.Close() + return false + }) + + select { + case <-uploadDone: + t.Fatal("the in-flight upload completed before it was released") + default: + } + + h.store.releaseAll() + + select { + case resp := <-uploadDone: + defer resp.Body.Close() + if resp.StatusCode != http.StatusOK { + t.Fatalf("in-flight upload finished with %d, want 200", resp.StatusCode) + } + case err := <-uploadErr: + t.Fatalf("in-flight upload was cut off by the shutdown: %v", err) + case <-time.After(5 * time.Second): + t.Fatal("in-flight upload never completed") + } + + if n, ok := h.store.storedBytes("data/" + drainBlob); !ok || n != len(payload) { + t.Fatalf("blob stored as %d bytes (present=%v), want %d", n, ok, len(payload)) + } + + select { + case err := <-drainErr: + if err != nil { + t.Fatalf("drain returned %v, want nil once the upload finished", err) + } + case <-time.After(5 * time.Second): + t.Fatal("drain never returned after the upload finished") + } +} + +// A pooled-but-idle keep-alive is what restic leaves behind between blobs. It +// must not hold the pod open, and the server has to close it rather than leave +// restic holding a socket it will never answer again. Driven over a raw +// connection so the assertion is about the server's FIN and not about whatever +// http.Transport decides to do with its pool. +func TestDrainClosesIdleKeepAlives(t *testing.T) { + h := newDrainHarness(t) + + conn, err := net.Dial("tcp", h.addr) + if err != nil { + t.Fatalf("dial: %v", err) + } + defer func() { _ = conn.Close() }() + + body := "lock" + req := "POST /" + drainRepo + "/locks/" + strings.Repeat("b", 64) + " HTTP/1.1\r\n" + + "Host: " + h.addr + "\r\n" + + "Authorization: " + drainAuthHeader(t) + "\r\n" + + "Content-Length: " + strconv.Itoa(len(body)) + "\r\n\r\n" + body + if _, err := conn.Write([]byte(req)); err != nil { + t.Fatalf("write request: %v", err) + } + + br := bufio.NewReader(conn) + resp, err := http.ReadResponse(br, nil) + if err != nil { + t.Fatalf("read response: %v", err) + } + _, _ = io.Copy(io.Discard, resp.Body) + resp.Body.Close() + if resp.StatusCode != http.StatusOK { + t.Fatalf("warm-up request: got %d, want 200", resp.StatusCode) + } + if resp.Close { + t.Fatal("server closed the connection before the drain; the test proves nothing") + } + <-h.store.started + + started := time.Now() + if err := drain(h.srv, h.httpSrv, 0, 10*time.Second); err != nil { + t.Fatalf("drain returned %v, want nil with only idle connections open", err) + } + if elapsed := time.Since(started); elapsed > 3*time.Second { + t.Fatalf("drain waited %s on an idle keep-alive", elapsed) + } + + _ = conn.SetReadDeadline(time.Now().Add(2 * time.Second)) + if _, err := br.Read(make([]byte, 1)); !errors.Is(err, io.EOF) { + t.Fatalf("idle keep-alive read returned %v, want EOF (server should have closed it)", err) + } +} + +// The shutdown timeout is a ceiling, not a promise: a transfer that outlives it +// is reported rather than waited on forever, so the kubelet's grace period is +// never what ends the process. +func TestDrainReportsRequestsPastTheDeadline(t *testing.T) { + h := newDrainHarness(t) + + go func() { _, _ = h.putBlob(t, "data", drainBlob, "payload") }() + <-h.store.started + + err := drain(h.srv, h.httpSrv, 0, 250*time.Millisecond) + if !errors.Is(err, context.DeadlineExceeded) { + t.Fatalf("drain returned %v, want context.DeadlineExceeded", err) + } +} + +func waitFor(t *testing.T, limit time.Duration, what string, cond func() bool) { + t.Helper() + deadline := time.Now().Add(limit) + for time.Now().Before(deadline) { + if cond() { + return + } + time.Sleep(10 * time.Millisecond) + } + t.Fatalf("timed out after %s waiting for %s", limit, what) +}