fix(build): retry the context fetch through a control-plane restart (#68)
This commit is contained in:
2 files changed
+208
-12
No files matched your search
+64
-12
@@ -18,9 +18,10 @@ import (
|
|||||||
)
|
)
|
||||||
|
|
||||||
// cmdFetchContext is the in-Pod entrypoint the build Job's context-fetch
|
// cmdFetchContext is the in-Pod entrypoint the build Job's context-fetch
|
||||||
// initContainer runs. It performs one read against the felis-api INTERNAL face —
|
// initContainer runs. It reads the blob the platform stored for a submission
|
||||||
// the blob the platform stored for a submission — and extracts it into the shared
|
// from the felis-api INTERNAL face (with a bounded retry — see
|
||||||
// emptyDir the Kaniko container then builds from.
|
// fetchContextWithRetry) and extracts it into the shared emptyDir the Kaniko
|
||||||
|
// container then builds from.
|
||||||
//
|
//
|
||||||
// Why this exists: the build Pod runs in the build namespace, where it can neither
|
// Why this exists: the build Pod runs in the build namespace, where it can neither
|
||||||
// mount the control-plane uploads PVC (a PVC does not cross namespaces) nor hold
|
// mount the control-plane uploads PVC (a PVC does not cross namespaces) nor hold
|
||||||
@@ -56,26 +57,22 @@ func cmdFetchContext(args []string, _, stderr io.Writer) int {
|
|||||||
ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM)
|
ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM)
|
||||||
defer stop()
|
defer stop()
|
||||||
|
|
||||||
req, err := http.NewRequestWithContext(ctx, http.MethodGet, *url, nil)
|
// Validate the URL once up front: a bad one is a usage error (2), not
|
||||||
if err != nil {
|
// something to sit in the retry loop.
|
||||||
|
if _, err := http.NewRequest(http.MethodGet, *url, nil); err != nil {
|
||||||
fmt.Fprintf(stderr, "felis fetch-context: bad --url: %v\n", err)
|
fmt.Fprintf(stderr, "felis fetch-context: bad --url: %v\n", err)
|
||||||
return 2
|
return 2
|
||||||
}
|
}
|
||||||
req.Header.Set("Authorization", "Bearer "+token)
|
|
||||||
// No overall client timeout: a legitimate modpack context can be large and the
|
// No overall client timeout: a legitimate modpack context can be large and the
|
||||||
// Job's activeDeadlineSeconds is the real bound. The header timeout catches a
|
// Job's activeDeadlineSeconds is the real bound. The header timeout catches a
|
||||||
// wedged endpoint without capping a healthy download.
|
// wedged endpoint without capping a healthy download.
|
||||||
client := &http.Client{Transport: &http.Transport{ResponseHeaderTimeout: time.Minute}}
|
client := &http.Client{Transport: &http.Transport{ResponseHeaderTimeout: time.Minute}}
|
||||||
resp, err := client.Do(req)
|
resp, err := fetchContextWithRetry(ctx, client, *url, token, stderr)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
fmt.Fprintf(stderr, "felis fetch-context: GET failed: %v\n", err)
|
fmt.Fprintf(stderr, "felis fetch-context: %v\n", err)
|
||||||
return 1
|
return 1
|
||||||
}
|
}
|
||||||
defer resp.Body.Close()
|
defer resp.Body.Close()
|
||||||
if resp.StatusCode != http.StatusOK {
|
|
||||||
fmt.Fprintf(stderr, "felis fetch-context: %s\n", resp.Status)
|
|
||||||
return 1
|
|
||||||
}
|
|
||||||
|
|
||||||
if err := extractTarGz(resp.Body, *out); err != nil {
|
if err := extractTarGz(resp.Body, *out); err != nil {
|
||||||
fmt.Fprintf(stderr, "felis fetch-context: %v\n", err)
|
fmt.Fprintf(stderr, "felis fetch-context: %v\n", err)
|
||||||
@@ -84,6 +81,61 @@ func cmdFetchContext(args []string, _, stderr io.Writer) int {
|
|||||||
return 0
|
return 0
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// fetchRetryInterval/fetchRetryWindow bound how long the fetch waits out a
|
||||||
|
// control-plane blip before giving up. The api pod being replaced is a normal
|
||||||
|
// event (rollout, eviction, a chaos drill), and without a retry one refused
|
||||||
|
// dial turns it into a failed build: BackoffLimit=0 gives the Job no second
|
||||||
|
// Pod, so the terminal verdict costs a manual re-approval — the live drill hit
|
||||||
|
// exactly this (context-fetch exit 1 on `connect: connection refused` while
|
||||||
|
// the api pod rolled; the new pod was serving 11 seconds later and the same
|
||||||
|
// 198-byte blob). The window is tiny next to the Job's 30-minute
|
||||||
|
// activeDeadline; a 4xx (missing blob, rejected token) still fails fast.
|
||||||
|
//
|
||||||
|
// Vars, not consts, so tests can shrink the window.
|
||||||
|
var (
|
||||||
|
fetchRetryInterval = 3 * time.Second
|
||||||
|
fetchRetryWindow = 45 * time.Second
|
||||||
|
)
|
||||||
|
|
||||||
|
// fetchContextWithRetry GETs the context tarball, retrying transport failures
|
||||||
|
// and 5xx responses until fetchRetryWindow runs out. A 4xx is an answer, not a
|
||||||
|
// blip — retrying it only delays the honest error.
|
||||||
|
func fetchContextWithRetry(ctx context.Context, client *http.Client, url, token string, stderr io.Writer) (*http.Response, error) {
|
||||||
|
deadline := time.Now().Add(fetchRetryWindow)
|
||||||
|
for {
|
||||||
|
req, err := http.NewRequestWithContext(ctx, http.MethodGet, url, nil)
|
||||||
|
if err != nil {
|
||||||
|
return nil, fmt.Errorf("bad --url: %w", err)
|
||||||
|
}
|
||||||
|
req.Header.Set("Authorization", "Bearer "+token)
|
||||||
|
|
||||||
|
resp, err := client.Do(req)
|
||||||
|
if err == nil && resp.StatusCode == http.StatusOK {
|
||||||
|
return resp, nil
|
||||||
|
}
|
||||||
|
if err == nil {
|
||||||
|
status := resp.Status
|
||||||
|
_ = resp.Body.Close()
|
||||||
|
err = fmt.Errorf("GET returned %s", status)
|
||||||
|
if resp.StatusCode < 500 {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if ctx.Err() != nil {
|
||||||
|
return nil, fmt.Errorf("GET failed: %w", err)
|
||||||
|
}
|
||||||
|
if time.Now().After(deadline) {
|
||||||
|
return nil, fmt.Errorf("GET failed (retried for %s): %w", fetchRetryWindow, err)
|
||||||
|
}
|
||||||
|
fmt.Fprintf(stderr, "felis fetch-context: %v; retrying (the internal face may be restarting)\n", err)
|
||||||
|
select {
|
||||||
|
case <-ctx.Done():
|
||||||
|
return nil, fmt.Errorf("GET failed: %w", err)
|
||||||
|
case <-time.After(fetchRetryInterval):
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
// extractTarGz streams a gzip'd tarball into root, creating directories as
|
// extractTarGz streams a gzip'd tarball into root, creating directories as
|
||||||
// needed. Every entry is vetted BEFORE anything is written: a path that is
|
// needed. Every entry is vetted BEFORE anything is written: a path that is
|
||||||
// absolute or escapes root (via ".."), a link (symlink or hardlink), or any
|
// absolute or escapes root (via ".."), a link (symlink or hardlink), or any
|
||||||
|
|||||||
@@ -5,12 +5,15 @@ import (
|
|||||||
"bytes"
|
"bytes"
|
||||||
"compress/gzip"
|
"compress/gzip"
|
||||||
"io"
|
"io"
|
||||||
|
"net"
|
||||||
"net/http"
|
"net/http"
|
||||||
"net/http/httptest"
|
"net/http/httptest"
|
||||||
"os"
|
"os"
|
||||||
"path/filepath"
|
"path/filepath"
|
||||||
"strings"
|
"strings"
|
||||||
|
"sync/atomic"
|
||||||
"testing"
|
"testing"
|
||||||
|
"time"
|
||||||
)
|
)
|
||||||
|
|
||||||
type tarEntry struct {
|
type tarEntry struct {
|
||||||
@@ -169,3 +172,144 @@ func TestExtractTarGzRejectsNonGzip(t *testing.T) {
|
|||||||
t.Fatalf("err = %v, want a gzip complaint", err)
|
t.Fatalf("err = %v, want a gzip complaint", err)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// shrinkFetchWindow swaps the retry knobs for a faster test and restores them
|
||||||
|
// afterwards, so no test leaks a tiny window into another.
|
||||||
|
func shrinkFetchWindow(t *testing.T, interval, window time.Duration) {
|
||||||
|
t.Helper()
|
||||||
|
oldInterval, oldWindow := fetchRetryInterval, fetchRetryWindow
|
||||||
|
fetchRetryInterval, fetchRetryWindow = interval, window
|
||||||
|
t.Cleanup(func() { fetchRetryInterval, fetchRetryWindow = oldInterval, oldWindow })
|
||||||
|
}
|
||||||
|
|
||||||
|
// A control-plane blip mid-fetch is survived: a 5xx on the first attempt is
|
||||||
|
// retried and the second attempt's tarball extracts. This walks back the live
|
||||||
|
// drill's failure, where the api pod rolled mid-fetch and the single attempt
|
||||||
|
// died, failing the build Job.
|
||||||
|
func TestFetchContextRetriesThroughBlip(t *testing.T) {
|
||||||
|
shrinkFetchWindow(t, 10*time.Millisecond, time.Second)
|
||||||
|
body := tgzBody(t, tarEntry{name: "Dockerfile", body: "FROM scratch\n"})
|
||||||
|
var calls int32
|
||||||
|
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
|
||||||
|
if atomic.AddInt32(&calls, 1) == 1 {
|
||||||
|
w.WriteHeader(http.StatusBadGateway) // the port is up, the API is not
|
||||||
|
return
|
||||||
|
}
|
||||||
|
_, _ = w.Write(body)
|
||||||
|
}))
|
||||||
|
defer srv.Close()
|
||||||
|
|
||||||
|
dir := t.TempDir()
|
||||||
|
t.Setenv("FELIS_SERVICE_TOKEN", "test-token")
|
||||||
|
var stderr bytes.Buffer
|
||||||
|
if code := cmdFetchContext([]string{"--url=" + srv.URL + "/sub-1/context", "--out=" + dir}, io.Discard, &stderr); code != 0 {
|
||||||
|
t.Fatalf("exit = %d, want 0 (stderr %q)", code, stderr.String())
|
||||||
|
}
|
||||||
|
if got, err := os.ReadFile(filepath.Join(dir, "Dockerfile")); err != nil || string(got) != "FROM scratch\n" {
|
||||||
|
t.Fatalf("extracted Dockerfile = (%q, %v)", got, err)
|
||||||
|
}
|
||||||
|
if !strings.Contains(stderr.String(), "retrying") {
|
||||||
|
t.Fatalf("stderr %q does not mention the retry", stderr.String())
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// The live drill's exact shape: the dial itself is refused (the api pod is
|
||||||
|
// gone and no endpoint answers). A refused dial is retried like any other
|
||||||
|
// transport failure, and once the face is back the fetch completes.
|
||||||
|
func TestFetchContextRetriesRefusedDial(t *testing.T) {
|
||||||
|
shrinkFetchWindow(t, 10*time.Millisecond, 5*time.Second)
|
||||||
|
body := tgzBody(t, tarEntry{name: "Dockerfile", body: "FROM scratch\n"})
|
||||||
|
|
||||||
|
// Borrow a listen address, then close it: the first attempts dial into a
|
||||||
|
// refused connection, exactly like a restarting control plane.
|
||||||
|
probe := httptest.NewServer(http.HandlerFunc(func(http.ResponseWriter, *http.Request) {}))
|
||||||
|
addr := strings.TrimPrefix(probe.URL, "http://")
|
||||||
|
probe.Close()
|
||||||
|
|
||||||
|
dir := t.TempDir()
|
||||||
|
t.Setenv("FELIS_SERVICE_TOKEN", "test-token")
|
||||||
|
var stderr bytes.Buffer
|
||||||
|
// Start the fetch; while the retry loop burns refused dials, bring the same
|
||||||
|
// address back.
|
||||||
|
result := make(chan int, 1)
|
||||||
|
go func() {
|
||||||
|
result <- cmdFetchContext([]string{"--url=http://" + addr + "/sub-1/context", "--out=" + dir}, io.Discard, &stderr)
|
||||||
|
}()
|
||||||
|
time.Sleep(100 * time.Millisecond) // let a handful of dials be refused
|
||||||
|
ln, err := net.Listen("tcp", addr)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("rebind %s: %v", addr, err)
|
||||||
|
}
|
||||||
|
back := &http.Server{Handler: http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||||
|
if r.Header.Get("Authorization") != "Bearer test-token" {
|
||||||
|
w.WriteHeader(http.StatusUnauthorized)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
_, _ = w.Write(body)
|
||||||
|
})}
|
||||||
|
defer back.Close()
|
||||||
|
go func() { _ = back.Serve(ln) }()
|
||||||
|
|
||||||
|
code := <-result
|
||||||
|
if code != 0 {
|
||||||
|
t.Fatalf("exit = %d, want 0 (stderr %q)", code, stderr.String())
|
||||||
|
}
|
||||||
|
if got, err := os.ReadFile(filepath.Join(dir, "Dockerfile")); err != nil || string(got) != "FROM scratch\n" {
|
||||||
|
t.Fatalf("extracted Dockerfile = (%q, %v)", got, err)
|
||||||
|
}
|
||||||
|
if !strings.Contains(stderr.String(), "retrying") {
|
||||||
|
t.Fatalf("stderr %q does not mention the retry", stderr.String())
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// A 4xx is an answer, not a blip: a missing/never-uploaded context fails
|
||||||
|
// immediately — no retry loop burns the build's deadline on a terminal error.
|
||||||
|
func TestFetchContextDoesNotRetry4xx(t *testing.T) {
|
||||||
|
shrinkFetchWindow(t, 5*time.Millisecond, 200*time.Millisecond)
|
||||||
|
var calls int32
|
||||||
|
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
|
||||||
|
atomic.AddInt32(&calls, 1)
|
||||||
|
w.WriteHeader(http.StatusNotFound)
|
||||||
|
}))
|
||||||
|
defer srv.Close()
|
||||||
|
|
||||||
|
t.Setenv("FELIS_SERVICE_TOKEN", "test-token")
|
||||||
|
var stderr bytes.Buffer
|
||||||
|
if code := cmdFetchContext([]string{"--url=" + srv.URL + "/sub-1/context", "--out=" + t.TempDir()}, io.Discard, &stderr); code != 1 {
|
||||||
|
t.Fatalf("exit = %d, want 1 (stderr %q)", code, stderr.String())
|
||||||
|
}
|
||||||
|
if got := atomic.LoadInt32(&calls); got != 1 {
|
||||||
|
t.Fatalf("server saw %d attempts, want exactly 1", got)
|
||||||
|
}
|
||||||
|
if strings.Contains(stderr.String(), "retrying") {
|
||||||
|
t.Fatalf("stderr %q mentions a retry for a terminal 4xx", stderr.String())
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// The retry is bounded: an internal face that stays down does not hang the
|
||||||
|
// build pod; the window runs out and the fetch reports the exhausted retries.
|
||||||
|
func TestFetchContextGivesUpAfterWindow(t *testing.T) {
|
||||||
|
shrinkFetchWindow(t, 5*time.Millisecond, 60*time.Millisecond)
|
||||||
|
var calls int32
|
||||||
|
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
|
||||||
|
atomic.AddInt32(&calls, 1)
|
||||||
|
w.WriteHeader(http.StatusServiceUnavailable)
|
||||||
|
}))
|
||||||
|
defer srv.Close() // the face is up but never healthy: 503 forever
|
||||||
|
|
||||||
|
t.Setenv("FELIS_SERVICE_TOKEN", "test-token")
|
||||||
|
var stderr bytes.Buffer
|
||||||
|
start := time.Now()
|
||||||
|
if code := cmdFetchContext([]string{"--url=" + srv.URL + "/sub-1/context", "--out=" + t.TempDir()}, io.Discard, &stderr); code != 1 {
|
||||||
|
t.Fatalf("exit = %d, want 1 (stderr %q)", code, stderr.String())
|
||||||
|
}
|
||||||
|
if elapsed := time.Since(start); elapsed > 5*time.Second {
|
||||||
|
t.Fatalf("gave up after %v; the window is supposed to bound it", elapsed)
|
||||||
|
}
|
||||||
|
if got := atomic.LoadInt32(&calls); got < 2 {
|
||||||
|
t.Fatalf("server saw %d attempts, want at least one retry", got)
|
||||||
|
}
|
||||||
|
if !strings.Contains(stderr.String(), "retried for") {
|
||||||
|
t.Fatalf("stderr %q does not report the exhausted retry window", stderr.String())
|
||||||
|
}
|
||||||
|
}
|
||||||
Reference in new issue
Block a user