From 094d3d5bc6c704bcaa451c44df90113544f1921b Mon Sep 17 00:00:00 2001 From: jochen Date: Sun, 11 Oct 2026 10:51:41 +0200 Subject: [PATCH] Wait out a registry held still, for up to five minutes, instead of failing the build (issue 457) The store's nightly collection stops the registry for about a minute and a half; builds in that window failed on connection refused and failed their whole plans. A refusal is now waited for with a growing pause, said in the build's log, and still fails the build past the bound. --- internal/builder/registry.go | 12 +- internal/builder/registry_wait.go | 131 +++++++++++++++++++++ internal/builder/registry_wait_test.go | 154 +++++++++++++++++++++++++ 3 files changed, 294 insertions(+), 3 deletions(-) create mode 100644 internal/builder/registry_wait.go create mode 100644 internal/builder/registry_wait_test.go diff --git a/internal/builder/registry.go b/internal/builder/registry.go index 219b4fe5..4f7ecd11 100644 --- a/internal/builder/registry.go +++ b/internal/builder/registry.go @@ -52,7 +52,11 @@ func (r Registry) PublishImage(ctx context.Context, localTag, repository string) if _, err := r.Run(ctx, "", "docker", "tag", localTag, remote); err != nil { return "", err } - if _, err := r.Run(ctx, "", "docker", "push", remote); err != nil { + // The push waits for a registry held still, as every request to it does (novox/hq issue 457). + if err := waitForRegistry(ctx, r.Address, "docker push "+remote, func() error { + _, err := r.Run(ctx, "", "docker", "push", remote) + return err + }); err != nil { return "", err } out, err := r.Run(ctx, "", "docker", "inspect", "--format", "{{index .RepoDigests 0}}", remote) @@ -169,11 +173,13 @@ func (r Registry) has(ctx context.Context, url string, accept ...string) (bool, return response.StatusCode == http.StatusOK, nil } +// client is the client for the registry and for upstream, whose requests to the registry wait out a +// registry held still (novox/hq issue 457). func (r Registry) client() *http.Client { if r.HTTP != nil { - return r.HTTP + return waiting(r.HTTP, r.Address) } - return http.DefaultClient + return waiting(http.DefaultClient, r.Address) } // separator is whether the upload location already carries a query. diff --git a/internal/builder/registry_wait.go b/internal/builder/registry_wait.go new file mode 100644 index 00000000..010622a0 --- /dev/null +++ b/internal/builder/registry_wait.go @@ -0,0 +1,131 @@ +package builder + +import ( + "context" + "errors" + "fmt" + "net/http" + "strings" + "syscall" + "time" +) + +// A build waits out a registry that refuses connections, for a bounded time (novox/hq issue 457). +// +// **Why waiting, and why here.** The store's nightly collection holds the registry still for its run — +// about a minute and a half, measured on 2026-10-11 — and a build that reached the registry in that +// window failed on "connection refused", and its whole delivery plan with it: a plan failed for a +// pause the mesh itself scheduled. The other design weighed was the collection telling the controller +// it holds the registry, and the controller holding build asks while it runs. Waiting here is smaller +// and covers more: it is local to the one place that talks to the registry, needs no new message +// between modules, and also carries a build over any other short outage — a registry restarted by its +// own update, say. A refusal is the one error waited for: nothing was sent, so trying again cannot +// do anything twice, and it is what a registry that is stopped answers. +// +// **Bounded, and loud past the bound.** registryWait is longer than the collection holds the registry +// (five minutes against about one and a half), so the pause the mesh schedules is always waited out, +// and a registry that is really down still fails the build — saying how long it was refused — rather +// than holding a build machine for ever. Every wait is said in the build's log, with why, and so is +// the registry answering again. + +var ( + // registryWait is how long a build waits for a registry that refuses, per call that found it so. + registryWait = 5 * time.Minute + // registryFirstPause is the first pause between tries; each pause doubles, up to registryMostPause. + registryFirstPause = time.Second +) + +// registryMostPause is the longest pause between two tries: short enough that a build goes on within +// seconds of the registry answering again. +const registryMostPause = 10 * time.Second + +// refused is whether an error is a connection refused: from a dial here, or as a command such as docker +// said it in its output. +func refused(err error) bool { + return err != nil && (errors.Is(err, syscall.ECONNREFUSED) || strings.Contains(err.Error(), "connection refused")) +} + +// waitForRegistry runs try, and while it fails because the registry at address refuses connections, +// tries again with a growing pause until registryWait has passed. what names the call, for the log. +func waitForRegistry(ctx context.Context, address, what string, try func() error) error { + err := try() + if !refused(err) { + return err + } + started := time.Now() + pause := registryFirstPause + tell("registry", "%s: refused; the build waits for the registry at %s, for up to %s — it is held still while "+ + "the store's nightly collection runs, about a minute and a half (novox/hq issue 457)", + what, address, registryWait) + for { + left := registryWait - time.Since(started) + if left <= 0 { + tell("registry", "%s: the registry at %s still refuses after %s; the build fails", what, address, + time.Since(started).Round(time.Second)) + return fmt.Errorf("the registry at %s refused every connection for %s, longer than its nightly "+ + "collection holds it still, so it is down, not paused: %w", + address, time.Since(started).Round(time.Millisecond), err) + } + wait := min(pause, left) + select { + case <-ctx.Done(): + return fmt.Errorf("stopped while waiting for the registry at %s: %w (last: %v)", address, ctx.Err(), err) + case <-time.After(wait): + } + pause = min(pause*2, registryMostPause) + if err = try(); !refused(err) { + tell("registry", "%s: the registry at %s answers again after %s; the build goes on", what, address, + time.Since(started).Round(time.Millisecond)) + return err + } + } +} + +// waitingTransport waits for the registry on every request to it, and on none to anywhere else: the +// same client copies from upstream registries, whose refusals are theirs to answer. +type waitingTransport struct { + base http.RoundTripper + address string +} + +func (t waitingTransport) RoundTrip(request *http.Request) (*http.Response, error) { + if request.URL.Host != t.address { + return t.base.RoundTrip(request) + } + var response *http.Response + tries := 0 + err := waitForRegistry(request.Context(), t.address, request.Method+" "+request.URL.Path, func() error { + attempt := request + if tries > 0 && request.Body != nil && request.Body != http.NoBody { + // A body is sent again only when it can be read again; one that cannot is not retried. + if request.GetBody == nil { + return fmt.Errorf("%s %s cannot be sent again: its body cannot be read twice", request.Method, request.URL) + } + body, err := request.GetBody() + if err != nil { + return err + } + attempt = request.Clone(request.Context()) + attempt.Body = body + } + tries++ + var err error + response, err = t.base.RoundTrip(attempt) + return err + }) + return response, err +} + +// waiting is a client like c whose requests to the registry wait for it. +func waiting(c *http.Client, address string) *http.Client { + if _, already := c.Transport.(waitingTransport); already { + return c + } + copied := *c + base := c.Transport + if base == nil { + base = http.DefaultTransport + } + copied.Transport = waitingTransport{base: base, address: address} + return &copied +} diff --git a/internal/builder/registry_wait_test.go b/internal/builder/registry_wait_test.go new file mode 100644 index 00000000..8bb9be85 --- /dev/null +++ b/internal/builder/registry_wait_test.go @@ -0,0 +1,154 @@ +package builder + +import ( + "context" + "crypto/sha256" + "encoding/hex" + "net" + "net/http" + "strings" + "sync" + "testing" + "time" +) + +// A registry held still for a while (novox/hq issue 457): the store's nightly collection stops it for +// about a minute and a half, and a build in that window was failed, and its whole plan with it, for a +// pause the mesh itself scheduled. A build waits it out — for a bound longer than the collection +// holds it — says so in its log, and still fails, loudly, past the bound. + +// heldStill is a registry address that refuses every connection until it starts answering after +// pause, or never when pause is negative. +func heldStill(t *testing.T, f *fakeRegistry, pause time.Duration) string { + t.Helper() + handler := f.serve(t).Config.Handler + reserved, err := net.Listen("tcp", "127.0.0.1:0") + if err != nil { + t.Fatal(err) + } + address := reserved.Addr().String() + reserved.Close() // refused from here on: nothing listens + if pause < 0 { + return address + } + server := &http.Server{Handler: handler} + go func() { + time.Sleep(pause) + l, err := net.Listen("tcp", address) + if err != nil { + t.Errorf("cannot answer at %s again: %v", address, err) + return + } + _ = server.Serve(l) + }() + t.Cleanup(func() { _ = server.Close() }) + return address +} + +// saying collects what a build says, as the build machine's per-build Said does. +func saying(t *testing.T) func() []string { + t.Helper() + var mu sync.Mutex + var lines []string + was := Said + Said = func(step, message string) { + mu.Lock() + defer mu.Unlock() + lines = append(lines, step+": "+message) + } + t.Cleanup(func() { Said = was }) + return func() []string { + mu.Lock() + defer mu.Unlock() + return append([]string(nil), lines...) + } +} + +// waitingFor shortens the bound and the first pause, so a test waits for milliseconds. +func waitingFor(t *testing.T, bound time.Duration) { + t.Helper() + wasBound, wasFirst := registryWait, registryFirstPause + registryWait, registryFirstPause = bound, 20*time.Millisecond + t.Cleanup(func() { registryWait, registryFirstPause = wasBound, wasFirst }) +} + +func TestABuildWaitsForARegistryHeldStillAndGoesOn(t *testing.T) { + waitingFor(t, 5*time.Second) + said := saying(t) + f := &fakeRegistry{} + r := Registry{Address: heldStill(t, f, 400*time.Millisecond)} + body := []byte("a theme") + sum := sha256.Sum256(body) + digest := "sha256:" + hex.EncodeToString(sum[:]) + + if _, err := r.PublishArchive(context.Background(), "shell/config", body, digest); err != nil { + t.Fatalf("a registry refusing for 400ms failed the build: %v", err) + } + if string(f.blobs[digest]) != "a theme" { + t.Fatalf("the registry holds %q", f.blobs[digest]) + } + log := strings.Join(said(), "\n") + if !strings.Contains(log, "waits for the registry at "+r.Address) || !strings.Contains(log, "nightly collection") { + t.Errorf("the build's log does not say it waited for the registry, and why:\n%s", log) + } + if !strings.Contains(log, "answers again") { + t.Errorf("the build's log does not say the registry came back:\n%s", log) + } +} + +func TestABuildFailsLoudlyOnARegistryRefusingPastTheBound(t *testing.T) { + waitingFor(t, 300*time.Millisecond) + said := saying(t) + f := &fakeRegistry{} + r := Registry{Address: heldStill(t, f, -1)} + body := []byte("a theme") + sum := sha256.Sum256(body) + digest := "sha256:" + hex.EncodeToString(sum[:]) + + started := time.Now() + _, err := r.PublishArchive(context.Background(), "shell/config", body, digest) + if err == nil { + t.Fatal("a registry that never answered published the archive") + } + if !strings.Contains(err.Error(), "refused every connection for") || !strings.Contains(err.Error(), "connection refused") { + t.Errorf("the failure does not say it waited and was refused throughout: %v", err) + } + if waited := time.Since(started); waited < 300*time.Millisecond { + t.Errorf("failed after %s, before the bound", waited) + } + if log := strings.Join(said(), "\n"); !strings.Contains(log, "waits for the registry") { + t.Errorf("the build's log does not say it waited:\n%s", log) + } +} + +func TestAnImagePushWaitsForARegistryHeldStill(t *testing.T) { + waitingFor(t, 5*time.Second) + said := saying(t) + pushes := 0 + run := func(_ context.Context, _ string, name string, args ...string) (string, error) { + switch args[0] { + case "push": + pushes++ + if pushes < 3 { + return "", errorString("docker push: dial tcp 127.0.0.1:5000: connect: connection refused") + } + case "inspect": + return "127.0.0.1:5000/m/server@sha256:abc\n", nil + } + return "", nil + } + r := Registry{Address: "127.0.0.1:5000", Run: run} + if _, err := r.PublishImage(context.Background(), "local", "m/server"); err != nil { + t.Fatalf("a push refused twice failed the build: %v", err) + } + if pushes != 3 { + t.Errorf("pushed %d times, want 3", pushes) + } + if log := strings.Join(said(), "\n"); !strings.Contains(log, "waits for the registry") { + t.Errorf("the build's log does not say it waited:\n%s", log) + } +} + +type errorString string + +func (e errorString) Error() string { return string(e) }