From 094d3d5bc6c704bcaa451c44df90113544f1921b Mon Sep 17 00:00:00 2001 From: jochen Date: Sun, 11 Oct 2026 10:51:41 +0200 Subject: [PATCH 1/3] 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) } From 7758301444c99eb1163377ae53fccd13f75d393d Mon Sep 17 00:00:00 2001 From: jochen Date: Sun, 11 Oct 2026 10:51:41 +0200 Subject: [PATCH 2/3] Ask every failed build of the tier on a retry, not only the first (issue 457) An outcome arriving after its plan failed is kept only in the build records, so a retry read from the plan alone asked one failed build and the records then failed the plan again on the next. Settle the tier from the records first. --- cmd/mesh-controller/plan_retry.go | 22 ++++++++ cmd/mesh-controller/plan_retry_tier_test.go | 57 +++++++++++++++++++++ cmd/mesh-controller/release_plan.go | 51 ++++++++++-------- 3 files changed, 110 insertions(+), 20 deletions(-) create mode 100644 cmd/mesh-controller/plan_retry_tier_test.go diff --git a/cmd/mesh-controller/plan_retry.go b/cmd/mesh-controller/plan_retry.go index a8377bc1..10ed9c45 100644 --- a/cmd/mesh-controller/plan_retry.go +++ b/cmd/mesh-controller/plan_retry.go @@ -288,6 +288,15 @@ func retryPlan(ctx context.Context, open *stores, id string) (string, error) { "those walks answer them; a newer merge, or `rebuild `, builds again", p.ID) } } + // Settled from the build records before anything is judged (novox/hq issue 457): a build of the tier + // that failed after the plan did is as failed as the one that failed it. + if p.Tier < len(p.Tiers) { + recorded, byID, err := recordsOfAsked(ctx, inv, &p, p.Tiers[p.Tier]) + if err != nil { + return "", err + } + toRetry(&p, recorded, byID) + } if err := retryRefusal(p, plans); err != nil { return "", err } @@ -525,3 +534,16 @@ func retryTierWhole(ctx context.Context, open *stores, p *inventory.Plan, again func sendAgain(s *inventory.PlanModule) { s.First, s.FirstAt, s.Gate, s.GatedBy, s.Previous, s.Why = nil, nil, nil, "", "", "" } + +// toRetry is the modules a retry asks again: every one of the tier that failed, the plan's state +// settled from the build records first (novox/hq issue 457). A plan fails on the first failure in its +// tier, and an outcome arriving after that finds no open plan to answer — it is kept only in the build +// records. Read from the plan alone, a retry asked only the build that failed first, and the records +// then failed the plan again on the next: each failed build of a tier took a retry of its own. What +// still runs is left asked, and its outcome is the plan's once the retry sets it building. +func toRetry(p *inventory.Plan, recorded map[string][]inventory.Build, byID map[string]inventory.Build) []string { + if p.Tier < len(p.Tiers) { + settleFromRecords(p, p.Tiers[p.Tier], recorded, byID) + } + return failedIn(*p) +} diff --git a/cmd/mesh-controller/plan_retry_tier_test.go b/cmd/mesh-controller/plan_retry_tier_test.go new file mode 100644 index 00000000..7efd44a8 --- /dev/null +++ b/cmd/mesh-controller/plan_retry_tier_test.go @@ -0,0 +1,57 @@ +package main + +import ( + "reflect" + "testing" + "time" + + "github.com/novox/mesh-controller/internal/inventory" +) + +// A retry asks every failed build of the tier, not one (novox/hq issue 457). Seen 2026-10-11: gitea +// and plex both failed in tier 0 while the registry was held still; gitea's failure failed the plan, +// and plex's, arriving after, found no open plan and was kept only in the build records. The first +// retry asked gitea alone, the records then failed the plan again on plex, and a second retry asked +// plex. +func TestARetryAsksEveryFailedBuildOfTheTier(t *testing.T) { + asked := time.Date(2026, 10, 11, 1, 29, 0, 0, time.UTC) + failedAt := asked.Add(2 * time.Minute) + p := inventory.Plan{ID: "plan-457", State: inventory.PlanFailed, Tier: 0, + Tiers: [][]string{{"gitea", "plex"}}, Note: "gitea failed to build in tier 0", + Modules: map[string]*inventory.PlanModule{ + "gitea": {State: "failed", AskedAt: &asked, Build: "build-gitea", Why: "cannot reach the registry"}, + "plex": {State: "asked", AskedAt: &asked, Build: "build-plex"}, + }} + byID := map[string]inventory.Build{ + "build-plex": {ID: "build-plex", Module: "plex", At: failedAt, Failed: "cannot reach the registry"}, + } + + if got := toRetry(&p, nil, byID); !reflect.DeepEqual(got, []string{"gitea", "plex"}) { + t.Fatalf("a retry of a tier where gitea and plex failed asks %v", got) + } +} + +// A build of the tier still running when the plan failed is not asked again: its outcome is the plan's +// once the retry sets it building, and one that is recorded built is taken as built. +func TestARetryLeavesABuildThatRunsOrWorked(t *testing.T) { + asked := time.Date(2026, 10, 11, 1, 29, 0, 0, time.UTC) + builtAt := asked.Add(3 * time.Minute) + p := inventory.Plan{ID: "plan-457", State: inventory.PlanFailed, Tier: 0, + Tiers: [][]string{{"a", "b", "c"}}, + Modules: map[string]*inventory.PlanModule{ + "a": {State: "failed", AskedAt: &asked, Build: "build-a", Why: "broken"}, + "b": {State: "asked", AskedAt: &asked, Build: "build-b"}, + "c": {State: "asked", AskedAt: &asked, Build: "build-c"}, + }} + byID := map[string]inventory.Build{"build-c": {ID: "build-c", Module: "c", At: builtAt, Commit: "c0ffee"}} + + if got := toRetry(&p, nil, byID); !reflect.DeepEqual(got, []string{"a"}) { + t.Fatalf("asks %v, want a alone", got) + } + if s := p.Modules["c"]; s.State != "built" || s.Commit != "c0ffee" { + t.Errorf("c, recorded built, is %+v", s) + } + if s := p.Modules["b"]; s.State != "asked" { + t.Errorf("b, still building, is %+v", s) + } +} diff --git a/cmd/mesh-controller/release_plan.go b/cmd/mesh-controller/release_plan.go index 76fbd95b..7adc089c 100644 --- a/cmd/mesh-controller/release_plan.go +++ b/cmd/mesh-controller/release_plan.go @@ -649,26 +649,9 @@ func advanceOnce(ctx context.Context, open *stores, p *inventory.Plan, // the controller in its first tier: the build that produced the new one is recorded, and the // plan never hears it. The record is the fact; a build recorded after the ask is that tier's // outcome, whoever was listening. - recorded := map[string][]inventory.Build{} - byID := map[string]inventory.Build{} - for _, m := range tier { - if s := p.Modules[m]; s != nil && s.State == "asked" { - builds, err := inv.Builds(ctx, m, 5) - if err != nil { - return false, err - } - recorded[m] = builds - // Its own ask's record, by id — found even when the outcome named no module (ADR 0219). - if s.Build != "" { - b, found, err := inv.BuildByID(ctx, s.Build) - if err != nil { - return false, err - } - if found { - byID[s.Build] = b - } - } - } + recorded, byID, err := recordsOfAsked(ctx, inv, p, tier) + if err != nil { + return false, err } if settleFromRecords(p, tier, recorded, byID) { return true, nil @@ -1776,6 +1759,34 @@ func splitList(s string) []string { return out } +// recordsOfAsked is what settleFromRecords reads: for every module of the tier still `asked`, its last +// builds and the record of its own ask, by id. +func recordsOfAsked(ctx context.Context, inv *inventory.Inventory, p *inventory.Plan, tier []string) ( + map[string][]inventory.Build, map[string]inventory.Build, error) { + recorded := map[string][]inventory.Build{} + byID := map[string]inventory.Build{} + for _, m := range tier { + if s := p.Modules[m]; s != nil && s.State == "asked" { + builds, err := inv.Builds(ctx, m, 5) + if err != nil { + return nil, nil, err + } + recorded[m] = builds + // Its own ask's record, by id — found even when the outcome named no module (ADR 0219). + if s.Build != "" { + b, found, err := inv.BuildByID(ctx, s.Build) + if err != nil { + return nil, nil, err + } + if found { + byID[s.Build] = b + } + } + } + } + return recorded, byID, nil +} + // settleFromRecords marks every module of the tier still `asked` built — or failed — from a build // recorded after it was asked, and says whether it changed anything (novox/hq 04-ISSUES/214). // Newest first, as Builds answers: the first record after the ask is the outcome of that ask. From 83ce19b9c0141624445728d43a2cd969227a56e6 Mon Sep 17 00:00:00 2001 From: jochen Date: Sun, 11 Oct 2026 11:05:30 +0200 Subject: [PATCH 3/3] Say a registry came back only when the call worked, and retry from toRetry's set (issue 457 review) The streaming blob PUT is left unwaited, with why, since its body cannot be read twice and the POST before it already waited. --- cmd/mesh-controller/plan_retry.go | 11 ++++++----- internal/builder/mirror.go | 5 +++++ internal/builder/registry_wait.go | 10 ++++++++-- 3 files changed, 19 insertions(+), 7 deletions(-) diff --git a/cmd/mesh-controller/plan_retry.go b/cmd/mesh-controller/plan_retry.go index 10ed9c45..df28926e 100644 --- a/cmd/mesh-controller/plan_retry.go +++ b/cmd/mesh-controller/plan_retry.go @@ -289,21 +289,23 @@ func retryPlan(ctx context.Context, open *stores, id string) (string, error) { } } // Settled from the build records before anything is judged (novox/hq issue 457): a build of the tier - // that failed after the plan did is as failed as the one that failed it. + // that failed after the plan did is as failed as the one that failed it. What toRetry says is the + // failed set every step below works from. + var failed []string if p.Tier < len(p.Tiers) { recorded, byID, err := recordsOfAsked(ctx, inv, &p, p.Tiers[p.Tier]) if err != nil { return "", err } - toRetry(&p, recorded, byID) + failed = toRetry(&p, recorded, byID) } if err := retryRefusal(p, plans); err != nil { return "", err } - if again := unjudgedAtGate(p); len(failedIn(p)) == 0 && len(again) > 0 { + if again := unjudgedAtGate(p); len(failed) == 0 && len(again) > 0 { return retryTierWhole(ctx, open, &p, again) } - if len(failedIn(p)) == 0 { + if len(failed) == 0 { return retryRollouts(ctx, open, &p) } entries, err := inv.Catalogued(ctx) @@ -314,7 +316,6 @@ func retryPlan(ctx context.Context, open *stores, id string) (string, error) { for _, e := range entries { byName[e.Manifest.Module] = e } - failed := failedIn(p) var asked []string for _, m := range failed { askModule(ctx, &p, m, byName) diff --git a/internal/builder/mirror.go b/internal/builder/mirror.go index 8840e3d5..4a2da87e 100644 --- a/internal/builder/mirror.go +++ b/internal/builder/mirror.go @@ -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) diff --git a/internal/builder/registry_wait.go b/internal/builder/registry_wait.go index 010622a0..f2b10bae 100644 --- a/internal/builder/registry_wait.go +++ b/internal/builder/registry_wait.go @@ -74,8 +74,14 @@ func waitForRegistry(ctx context.Context, address, what string, try func() error } 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)) + // 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 } }