Merge pull request 'A build waits out a registry held still, and a retry asks every failed build of the tier (issue 457)' (#225) from fix/457-a-build-waits-out-a-registry-pause into main

This commit was merged in pull request #225.
This commit is contained in:
2026-10-11 09:17:09 +00:00
7 changed files with 419 additions and 26 deletions
+5
View File
@@ -429,6 +429,11 @@ func (r Registry) copyBlob(ctx context.Context, src *source, where upstream, dig
if response.ContentLength > 0 {
put.ContentLength = response.ContentLength
}
// **Not waited for if refused** (novox/hq issue 457): the body streams from upstream and cannot be
// read twice, so a registry held still between the POST above and this PUT fails the copy with
// "its body cannot be read twice" rather than waiting. Accepted: the POST a moment before already
// waited the registry out, so the window is the length of one upstream fetch, and the build fails
// loudly, to be asked again, rather than buffering every base blob in memory.
done, err := r.client().Do(put)
if err != nil {
return fmt.Errorf("cannot upload blob %s: %w", digest, err)
+9 -3
View File
@@ -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.
+137
View File
@@ -0,0 +1,137 @@
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) {
// Said as it is: the registry answering is only the build going on when the call worked.
if err == nil {
tell("registry", "%s: the registry at %s answers again after %s; the build goes on", what, address,
time.Since(started).Round(time.Millisecond))
} else {
tell("registry", "%s: the registry at %s no longer refuses after %s, and answered with: %v", what,
address, time.Since(started).Round(time.Millisecond), err)
}
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
}
+154
View File
@@ -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) }