Files
mesh-controller/cmd/mesh-controller/mirrors_test.go
T
jochen c6e372896b
mesh/merge-gate pass: builds build-agent, mesh-controller, route-proxy → ace, g14, novox, shanks; no bus step; every machine composes with the change as it…
mesh/repo-check pass: its merge-check.sh passed
mesh/delivery superseded: a newer delivery to the same trunk took over its walk
Leave an eligible index for a confirmed collect instead of letting it go alone
Let go of alone after a build, an index's platforms stayed for ever under a record that
said collected, so no later collect could reach them (re-review of #144).
2026-10-08 15:47:42 +02:00

299 lines
12 KiB
Go

package main
import (
"fmt"
"net/http"
"net/http/httptest"
"slices"
"strings"
"sync"
"testing"
"time"
"github.com/novox/mesh-controller/internal/artifacts"
"github.com/novox/mesh-controller/internal/catalogue"
"github.com/novox/mesh-controller/internal/inventory"
"github.com/novox/mesh-controller/internal/link"
)
// The bases the mesh copied in, and recording a copy no record names (novox/hq ADR 0257).
func TestTheMirrorsVerbComposesItsCommandAndRefusesWhatItDoesNotTake(t *testing.T) {
r := ref("route-proxy", "on-go_base", 1)
for _, c := range []struct {
args map[string]any
want string
}{
{map[string]any{}, "mirrors --json"},
{map[string]any{"record": r}, "mirrors --json --record " + r},
{map[string]any{"record": r, "confirm": "true", "why": "copied before copies were recorded"},
"mirrors --json --record " + r + " --confirm --why copied before copies were recorded"},
} {
argv, err := argvFor("mirrors", c.args)
if err != nil || strings.Join(argv, " ") != c.want {
t.Errorf("mirrors %v: %q %v, want %q", c.args, argv, err, c.want)
}
}
for _, args := range []map[string]any{
{"record": r, "confirm": "true"}, // recorded without a why
{"record": r, "why": "x"}, // a why for a dry run
{"confirm": "true", "why": "x"}, // nothing to record
{"record": r, "confirm": "yes", "why": "x"}, // a switch is true or false
{"delete": "true"},
} {
if argv, err := argvFor("mirrors", args); err == nil {
t.Errorf("mirrors %v was composed as %q", args, argv)
}
}
if repairingCommand([]string{"mirrors", "--json", "--record", r, "--confirm", "--why", "x"}) != "mirrors" {
t.Error("recording a copy through the generic verb would go unrecorded")
}
if repairingCommand([]string{"mirrors", "--json", "--record", r}) != "" {
t.Error("a dry run was taken for an act by hand")
}
if !personsDecision(link.HandAct{Verb: "mirrors"}) {
t.Error("recording a copy would count toward a healer the mesh lacks")
}
}
func TestOnlyACopyInARepositoryTheMirrorWritesIsTaken(t *testing.T) {
for given, ok := range map[string]bool{
ref("route-proxy", "on-go_base", 1): true,
catalogue.ArtifactStoreScheme + "upstream/docker.io/library/golang@sha256:" + strings.Repeat("a", 64): true,
ref("route-proxy", "server", 1): false, // a module's own image
catalogue.ArtifactStoreScheme + "novox/invoicing-api@sha256:" + strings.Repeat("a", 64): false,
archiveRef("web", "on-tools", 1): false, // a blob
catalogue.ArtifactStoreScheme + "web/on-x@sha256:abc": false, // not a whole digest
"docker.io/library/golang@sha256:" + strings.Repeat("a", 64): false,
} {
if _, err := mirrorReference(given); (err == nil) != ok {
t.Errorf("%s: taken %v, want %v (%v)", given, err == nil, ok, err)
}
}
}
// heldStore answers a manifest HEAD with 200 for the digests it holds, and records every request.
func heldStore(t *testing.T, holds ...string) (artifacts.Store, *[]string) {
t.Helper()
var asked []string
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
asked = append(asked, r.Method+" "+r.URL.Path)
if r.Method == http.MethodHead && slices.ContainsFunc(holds, func(d string) bool { return strings.HasSuffix(r.URL.Path, d) }) {
w.WriteHeader(http.StatusOK)
return
}
w.WriteHeader(http.StatusNotFound)
}))
t.Cleanup(server.Close)
return artifacts.Store{Address: strings.TrimPrefix(server.URL, "http://")}, &asked
}
func TestRecordingACopyTakesWhatTheStoreHoldsAndDeletesNothing(t *testing.T) {
inv := inventory.ForTest(t)
held := ref("route-proxy", "on-go_base", 1)
absent := ref("route-proxy", "on-go_base", 2)
known := ref("nats", "on-nats_base", 3)
notACopy := ref("route-proxy", "server", 4)
if err := inv.RecordMirrored(t.Context(), "b1", "", []string{known}); err != nil {
t.Fatal(err)
}
store, asked := heldStore(t, fmt.Sprintf("%064x", 1))
given := []string{held, absent, known, notACopy}
mirrored, err := inv.Mirrored(t.Context())
if err != nil {
t.Fatal(err)
}
// A dry run says, and records nothing.
dry, err := recordMirrors(t.Context(), inv, store, given, mirrored, "", false)
if err != nil {
t.Fatal(err)
}
if !dry.DryRun || !slices.Equal(dry.Recorded, []string{held}) || !slices.Equal(dry.AlreadyRecorded, []string{known}) ||
len(dry.Refused) != 2 || dry.Refused[absent] == "" || dry.Refused[notACopy] == "" {
t.Fatalf("a dry run answered %+v", dry)
}
if !strings.Contains(dry.Then, "would become eligible") || !strings.Contains(dry.Then, "next sweep") {
t.Fatalf("a dry run did not say what recording does: %q", dry.Then)
}
if again, _ := inv.Mirrored(t.Context()); again[held] {
t.Fatal("a dry run recorded a copy")
}
// A real one records what the store holds, with the person's why, and the sweep may now decide.
real, err := recordMirrors(t.Context(), inv, store, given, mirrored, "copied before copies were recorded", true)
if err != nil {
t.Fatal(err)
}
if !slices.Equal(real.Recorded, []string{held}) {
t.Fatalf("recorded %v", real.Recorded)
}
left, err := inv.ToCollect(t.Context())
if err != nil {
t.Fatal(err)
}
if !slices.Contains(left, held) {
t.Fatalf("a recorded copy no kept build stood on is not offered to collect: %v", left)
}
for _, r := range *asked {
if !strings.HasPrefix(r, "HEAD ") {
t.Fatalf("recording a copy asked the store %s", r)
}
}
}
// indexRegistry holds indexes and platforms of one image's repository, answers GETs with their
// documents, refuses GETs for the digests in fail, and records every delete.
type indexRegistry struct {
mu sync.Mutex
indexes map[string][]string
held map[string]bool
fail map[string]bool
deleted []string
}
func (f *indexRegistry) serve(t *testing.T) artifacts.Store {
t.Helper()
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
f.mu.Lock()
defer f.mu.Unlock()
digest := r.URL.Path[strings.LastIndex(r.URL.Path, "/")+1:]
switch r.Method {
case http.MethodGet:
if f.fail[digest] {
w.WriteHeader(http.StatusInternalServerError)
return
}
if children, ok := f.indexes[digest]; ok && f.held[digest] {
var named []string
for _, c := range children {
named = append(named, `{"mediaType":"application/vnd.oci.image.manifest.v1+json","digest":"`+c+`"}`)
}
_, _ = w.Write([]byte(`{"schemaVersion":2,"mediaType":"application/vnd.oci.image.index.v1+json","manifests":[` +
strings.Join(named, ",") + `]}`))
return
}
if f.held[digest] {
_, _ = w.Write([]byte(`{"schemaVersion":2,"layers":[]}`))
return
}
w.WriteHeader(http.StatusNotFound)
case http.MethodHead:
if f.held[digest] {
w.WriteHeader(http.StatusOK)
return
}
w.WriteHeader(http.StatusNotFound)
case http.MethodDelete:
if !f.held[digest] {
w.WriteHeader(http.StatusNotFound)
return
}
delete(f.held, digest)
f.deleted = append(f.deleted, digest)
w.WriteHeader(http.StatusAccepted)
default:
w.WriteHeader(http.StatusBadRequest)
}
}))
t.Cleanup(server.Close)
return artifacts.Store{Address: strings.TrimPrefix(server.URL, "http://")}
}
func copyDigest(n int) string { return fmt.Sprintf("sha256:%064x", n) }
// twoCopies records, in one image's repository, a kept copy (stood on by a held module's build) naming
// platforms 11 and 12, and an eligible one (a failed build's) naming 12 and 13.
func twoCopies(t *testing.T, inv *inventory.Inventory) (kept, eligible string, reg *indexRegistry) {
t.Helper()
const repository = "upstream/docker.io/library/golang"
kept = catalogue.ArtifactStoreScheme + repository + "@" + copyDigest(1)
eligible = catalogue.ArtifactStoreScheme + repository + "@" + copyDigest(2)
if err := inv.RegisterModule(t.Context(), catalogue.Manifest{Module: "proxy", Version: "1"},
inventory.Source{Repository: "https://forge.invalid/proxy.git"}); err != nil {
t.Fatal(err)
}
failed := inventory.Build{ID: "f1", Repository: "https://forge.invalid/proxy.git", On: "a-build-machine",
Failed: "a recipe refused", Mirrored: []string{eligible}}
if err := inv.RecordBuild(t.Context(), failed); err != nil {
t.Fatal(err)
}
ok := inventory.Build{ID: "b1", Repository: "https://forge.invalid/proxy.git", Module: "proxy", On: "a-build-machine",
Commit: "c0ffee", Made: []inventory.Artifact{{Name: "app", Kind: "image", Reference: ref("proxy", "app", 7)}},
Against: []string{kept}, Mirrored: []string{kept}}
if err := inv.RecordBuild(t.Context(), ok); err != nil {
t.Fatal(err)
}
reg = &indexRegistry{
indexes: map[string][]string{copyDigest(1): {copyDigest(11), copyDigest(12)}, copyDigest(2): {copyDigest(12), copyDigest(13)}},
held: map[string]bool{copyDigest(1): true, copyDigest(2): true, copyDigest(11): true, copyDigest(12): true, copyDigest(13): true},
fail: map[string]bool{},
}
return kept, eligible, reg
}
// Through the records: a confirmed collect lets the eligible index go with the platform only it names,
// and leaves the platform the kept index names.
func TestAConfirmedCollectLetsAnIndexGoWithItsOwnPlatformsOnly(t *testing.T) {
inv := inventory.ForTest(t)
_, eligible, reg := twoCopies(t, inv)
store := reg.serve(t)
references, err := inv.ToCollect(t.Context())
if err != nil {
t.Fatal(err)
}
if !slices.Equal(references, []string{eligible}) {
t.Fatalf("offered %v, want the failed build's copy", references)
}
a := runCollect(t.Context(), inv, store, references, nil, true,
sweepBounds{most: 10, budget: 5 * time.Second, platforms: true})
if !slices.Equal(a.LetGo, []string{eligible}) || a.Stopped != "" {
t.Fatalf("let go %v, stopped %q", a.LetGo, a.Stopped)
}
if !slices.Equal(reg.deleted, []string{copyDigest(13), copyDigest(2)}) {
t.Fatalf("deleted %v; want its own platform, then the index", reg.deleted)
}
}
// The sweep after a build leaves an eligible index in place and eligible: let go of alone, its
// platforms would stay for ever under a record that says collected. A later collect a person confirms
// takes it with the platform only it names. An image the same sweep reaches is let go of as before.
func TestTheSweepAfterABuildLeavesAnIndexForAConfirmedCollect(t *testing.T) {
inv := inventory.ForTest(t)
_, eligible, reg := twoCopies(t, inv)
reg.held[copyDigest(5)] = true
image := catalogue.ArtifactStoreScheme + "upstream/docker.io/library/golang@" + copyDigest(5)
store := reg.serve(t)
r := sweep(t.Context(), inv, store, []string{eligible, image}, nil, afterBuild)
if r.Indexes != 1 || !slices.Equal(r.LetGo, []string{image}) || !slices.Equal(reg.deleted, []string{copyDigest(5)}) {
t.Fatalf("after a build: %d indexes left, let go %v, deleted %v; want the index left and the image gone",
r.Indexes, r.LetGo, reg.deleted)
}
left, err := inv.ToCollect(t.Context())
if err != nil {
t.Fatal(err)
}
if !slices.Equal(left, []string{eligible}) {
t.Fatalf("after the sweep the records offer %v; want the index still eligible", left)
}
a := runCollect(t.Context(), inv, store, left, nil, true,
sweepBounds{most: 10, budget: 5 * time.Second, platforms: true})
if !slices.Equal(a.LetGo, []string{eligible}) ||
!slices.Equal(reg.deleted, []string{copyDigest(5), copyDigest(13), copyDigest(2)}) {
t.Fatalf("a confirmed collect let go %v, deleted %v; want the index with its own platform", a.LetGo, reg.deleted)
}
}
// A kept index the store will not answer for stops a confirmed collect before its first delete.
func TestASpareListThatCannotBeReadStopsTheSweepBeforeAnyDelete(t *testing.T) {
inv := inventory.ForTest(t)
_, eligible, reg := twoCopies(t, inv)
reg.fail[copyDigest(1)] = true
store := reg.serve(t)
a := runCollect(t.Context(), inv, store, []string{eligible}, nil, true,
sweepBounds{most: 10, budget: 5 * time.Second, platforms: true})
if len(a.LetGo) != 0 || len(reg.deleted) != 0 || !strings.Contains(a.Stopped, "nothing was let go") {
t.Fatalf("let go %v, deleted %v, stopped %q", a.LetGo, reg.deleted, a.Stopped)
}
}